From 3c59d3c1a42e77ebcd5c45576512515ff2897ee2 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Fri, 21 Aug 2026 13:26:56 -0400 Subject: [PATCH 1/3] refactor(session): clarify execution state --- packages/core/src/session/inbox.ts | 10 +++++----- packages/core/src/session/model-request.ts | 14 +++++++------- .../core/src/session/runner/publish-llm-event.ts | 14 +++++++------- 3 files changed, 19 insertions(+), 19 deletions(-) diff --git a/packages/core/src/session/inbox.ts b/packages/core/src/session/inbox.ts index 11d1a4ccec77..1ebf86712d70 100644 --- a/packages/core/src/session/inbox.ts +++ b/packages/core/src/session/inbox.ts @@ -157,14 +157,14 @@ export const admit = Effect.fn("SessionInbox.admit")(function* ( item: request.item, }) .pipe( - Effect.flatMap((event) => { - const base = { + Effect.map((event) => + Info.make({ id: request.id, sessionID: request.sessionID, timeCreated: DateTime.makeUnsafe(event.created), - } - return Effect.succeed(Info.make({ ...base, ...request.item })) - }), + ...request.item, + }), + ), Effect.catchDefect((defect) => find(db, request.id).pipe( Effect.flatMap((stored) => diff --git a/packages/core/src/session/model-request.ts b/packages/core/src/session/model-request.ts index 722ac6023106..2adba5e52ca2 100644 --- a/packages/core/src/session/model-request.ts +++ b/packages/core/src/session/model-request.ts @@ -266,9 +266,9 @@ export const layer = Layer.effect( toolChoice: input.toolChoice, }), ) - const webSocketEligible = - !(yield* hooks.has("session", "http.request", resolved.ref.providerID)) && - !(yield* hooks.has("session", "http.response", resolved.ref.providerID)) + const hasHttpHooks = + (yield* hooks.has("session", "http.request", resolved.ref.providerID)) || + (yield* hooks.has("session", "http.response", resolved.ref.providerID)) const webSocket = resolved.capabilities.responsesWebsockets === true ? yield* Config.boolean(responsesWebSocketFlag(resolved.ref.providerID)).pipe( @@ -276,18 +276,18 @@ export const layer = Layer.effect( Effect.orDie, ) : false - const http = webSocketEligible - ? undefined - : SessionModelHttp.middleware(hooks, { + const http = hasHttpHooks + ? SessionModelHttp.middleware(hooks, { sessionID: session.id, agent: input.scope.agentID, model: resolved.ref, }) + : undefined const options: StreamOptions = { ...(http ? { http } : {}), ...(input.webSocket === "session" && webSocket && - webSocketEligible && + !hasHttpHooks && resolved.capabilities.responsesWebsockets === true ? { webSocket: transport.bind(session.id) } : {}), diff --git a/packages/core/src/session/runner/publish-llm-event.ts b/packages/core/src/session/runner/publish-llm-event.ts index f569f686fb98..1d25a5330c60 100644 --- a/packages/core/src/session/runner/publish-llm-event.ts +++ b/packages/core/src/session/runner/publish-llm-event.ts @@ -263,13 +263,14 @@ export const createLLMEventPublisher = (bus: Pick, inp }) { if (tools.has(event.id)) return yield* Effect.die(new Error(`Duplicate tool input start: ${event.id}`)) const assistantMessageID = yield* startAssistant() - tools.set(event.id, { + const tool: NonNullable> = { assistantMessageID, name: event.name, called: false, settled: false, providerExecuted: event.providerExecuted === true, - }) + } + tools.set(event.id, tool) yield* toolInput.start(event.id) yield* bus.publish(SessionEvent.Tool.Input.Started, { sessionID: input.sessionID, @@ -277,6 +278,7 @@ export const createLLMEventPublisher = (bus: Pick, inp id: event.id, name: event.name, }) + return tool }) const endToolInput = Effect.fnUntraced(function* ( @@ -296,9 +298,8 @@ export const createLLMEventPublisher = (bus: Pick, inp readonly name: string readonly raw: string }) { - if (!tools.has(event.id)) yield* startToolInput(event) - const tool = tools.get(event.id) - if (!tool || tool.called || tool.settled) + const tool = tools.get(event.id) ?? (yield* startToolInput(event)) + if (tool.called || tool.settled) return yield* Effect.die(new Error(`Malformed tool input after call settlement: ${event.id}`)) if (tool.name !== event.name) return yield* Effect.die(new Error(`Tool input name changed for ${event.id}: ${tool.name} -> ${event.name}`)) @@ -443,8 +444,7 @@ export const createLLMEventPublisher = (bus: Pick, inp return case "tool-call": { outputStarted = true - if (!tools.has(event.id)) yield* startToolInput(event) - const tool = tools.get(event.id)! + const tool = tools.get(event.id) ?? (yield* startToolInput(event)) if (toolInput.has(event.id)) yield* endToolInput(event) if (tool.name !== event.name) return yield* Effect.die(new Error(`Tool call name changed for ${event.id}: ${tool.name} -> ${event.name}`)) From d4c1e513e93b923b2a46a6f98d2afecf6d83f81c Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Fri, 21 Aug 2026 13:25:42 -0400 Subject: [PATCH 2/3] refactor(core): flatten state transformations --- packages/core/src/form.ts | 25 ++++++++------------ packages/core/src/session/message-updater.ts | 16 ++++++------- packages/core/src/state.ts | 14 +++++------ 3 files changed, 24 insertions(+), 31 deletions(-) diff --git a/packages/core/src/form.ts b/packages/core/src/form.ts index b8098733e3fe..9c54ddd66304 100644 --- a/packages/core/src/form.ts +++ b/packages/core/src/form.ts @@ -109,16 +109,11 @@ export const layer = Layer.effect( }, ) - const find = Effect.fn("Form.find")(function* (id: ID) { - return yield* Cache.getSuccess(forms, id).pipe( - Effect.flatMap((entry) => - Option.match(entry, { - onNone: () => Effect.fail(new NotFoundError({ id })), - onSome: Effect.succeed, - }), - ), - ) - }) + const requireEntry = Effect.fn("Form.requireEntry")((id: ID) => + Cache.getSuccess(forms, id).pipe( + Effect.flatMap((entry) => Effect.fromOption(entry, () => new NotFoundError({ id }))), + ), + ) const create = Effect.fn("Form.create")((input: CreateInput) => Effect.uninterruptible( @@ -151,7 +146,7 @@ export const layer = Layer.effect( Effect.uninterruptibleMask((restore) => Effect.gen(function* () { const form = yield* create(input) - const entry = yield* find(form.id).pipe(Effect.orDie) + const entry = yield* requireEntry(form.id).pipe(Effect.orDie) return yield* restore(Deferred.await(entry.deferred)).pipe( Effect.onInterrupt(() => Effect.ignore(cancel(form.id))), ) @@ -160,7 +155,7 @@ export const layer = Layer.effect( ) const get = Effect.fn("Form.get")(function* (id: ID) { - return (yield* find(id)).form + return (yield* requireEntry(id)).form }) const list = Effect.fn("Form.list")(function* (input?: ListInput) { @@ -172,13 +167,13 @@ export const layer = Layer.effect( }) const state = Effect.fn("Form.state")(function* (id: ID) { - return (yield* find(id)).state + return (yield* requireEntry(id)).state }) const reply = Effect.fn("Form.reply")((input: ReplyInput) => Effect.uninterruptible( Effect.gen(function* () { - const entry = yield* find(input.id) + const entry = yield* requireEntry(input.id) if (entry.state.status !== "pending") return yield* new AlreadySettledError({ id: input.id }) const invalid = validateAnswer(entry.form.fields, input.answer) if (invalid) return yield* new InvalidAnswerError({ id: input.id, message: invalid }) @@ -197,7 +192,7 @@ export const layer = Layer.effect( const cancel = Effect.fn("Form.cancel")((id: ID) => Effect.uninterruptible( Effect.gen(function* () { - const entry = yield* find(id) + const entry = yield* requireEntry(id) if (entry.state.status !== "pending") return yield* new AlreadySettledError({ id }) const next: TerminalState = { status: "cancelled" } yield* bus.publish(Form.Event.Cancelled, { id, sessionID: entry.form.sessionID }) diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index a493fadb317b..ba3a7e0eb4c0 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -42,18 +42,18 @@ export function update(adapter: Adapter, event: SessionEvent.DurableEvent) { const updateOwnedAssistant = (messageID: SessionMessage.ID, recipe: (draft: DraftAssistant) => void) => Effect.gen(function* () { const assistant = yield* adapter.getAssistant(messageID) - if (assistant) yield* adapter.updateAssistant(produce(assistant, recipe)) + if (!assistant) return + yield* adapter.updateAssistant(produce(assistant, recipe)) }) const clearCurrentRetry = Effect.gen(function* () { const assistant = yield* adapter.getCurrentAssistant() - if (assistant?.retry) { - yield* adapter.updateAssistant( - produce(assistant, (draft) => { - draft.retry = undefined - }), - ) - } + if (!assistant?.retry) return + yield* adapter.updateAssistant( + produce(assistant, (draft) => { + draft.retry = undefined + }), + ) }) const project = pipe( diff --git a/packages/core/src/state.ts b/packages/core/src/state.ts index dd64809f1084..432052c482e0 100644 --- a/packages/core/src/state.ts +++ b/packages/core/src/state.ts @@ -89,15 +89,14 @@ export function create(options: Options): Inte if (options.finalize) yield* options.finalize(options.draft(next)) }) - const apply = (transform: TransformCallback, draft: DraftApi) => - Effect.sync(() => { - transform(draft) - }) - const materialize = Effect.fnUntraced(function* () { const next = options.initial() const api = options.draft(next) - for (const transform of transforms) yield* apply(transform.run, api) + for (const transform of transforms) { + yield* Effect.sync(() => { + transform.run(api) + }) + } yield* commit(next) }) @@ -135,7 +134,7 @@ export function create(options: Options): Inte return yield* Deferred.await(done) }) - const result: Interface = { + return { get: () => state, transform: Effect.fn("State.transform")(function* (update) { yield* Effect.annotateCurrentSpan("state", options.name ?? "anonymous") @@ -176,5 +175,4 @@ export function create(options: Options): Inte }), reload, } - return result } From cd3e1ec24018979af5237f065dacd54ff0f243ec Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Fri, 21 Aug 2026 13:36:20 -0400 Subject: [PATCH 3/3] refactor(session): name tool state --- .../src/session/runner/publish-llm-event.ts | 22 +++++++++---------- 1 file changed, 10 insertions(+), 12 deletions(-) diff --git a/packages/core/src/session/runner/publish-llm-event.ts b/packages/core/src/session/runner/publish-llm-event.ts index 1d25a5330c60..30cd003bcc6f 100644 --- a/packages/core/src/session/runner/publish-llm-event.ts +++ b/packages/core/src/session/runner/publish-llm-event.ts @@ -84,17 +84,15 @@ const hostedContent = (result: ToolResultValue): NonEmptyContent => { */ export const createLLMEventPublisher = (bus: Pick, input: Input) => { const deltaBatchInterval = 100 - const tools = new Map< - string, - { - readonly assistantMessageID: SessionMessage.ID - readonly name: string - called: boolean - settled: boolean - providerExecuted: boolean - progress?: Tool.Metadata - } - >() + type ToolState = { + readonly assistantMessageID: SessionMessage.ID + readonly name: string + called: boolean + settled: boolean + providerExecuted: boolean + progress?: Tool.Metadata + } + const tools = new Map() const failureSnapshot = (tool: { readonly progress?: Tool.Metadata }, metadata?: Tool.Metadata) => { if (tool.progress === undefined) return metadata === undefined ? {} : { metadata } if (metadata === undefined) return { metadata: tool.progress } @@ -263,7 +261,7 @@ export const createLLMEventPublisher = (bus: Pick, inp }) { if (tools.has(event.id)) return yield* Effect.die(new Error(`Duplicate tool input start: ${event.id}`)) const assistantMessageID = yield* startAssistant() - const tool: NonNullable> = { + const tool: ToolState = { assistantMessageID, name: event.name, called: false,