Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
fix(core): harden v2 session input cutover
  • Loading branch information
kitlangton committed Jun 4, 2026
commit c9ac2d04824ef9e8c54c81abbf90e92a339317f6

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 6 additions & 4 deletions packages/core/src/event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -370,11 +370,13 @@ export const layerWith = (options?: LayerOptions) =>
return Effect.gen(function* () {
const durable = registry.get(event.type)?.sync !== undefined
if (durable) {
for (const sync of syncHandlers) {
yield* sync(event as Payload)
}
const committed = yield* commitSyncEvent(event as Payload)
if (committed) event = { ...event, seq: committed.seq }
if (committed) {
event = { ...event, seq: committed.seq }
for (const sync of syncHandlers) {
yield* sync(event as Payload)
}
}
}
for (const listener of listeners) {
yield* listener(event as Payload)
Expand Down
6 changes: 3 additions & 3 deletions packages/core/src/session/event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,7 @@ export namespace Step {
...stepSettlementOptions,
schema: {
...Base,
assistantMessageID: EventV2.ID,
assistantCreatorEventID: EventV2.ID,
finish: Schema.String,
cost: Schema.Finite,
tokens: Schema.Struct({
Expand All @@ -188,7 +188,7 @@ export namespace Step {
...stepSettlementOptions,
schema: {
...Base,
assistantMessageID: EventV2.ID,
assistantCreatorEventID: EventV2.ID,
error: UnknownError,
},
})
Expand Down Expand Up @@ -268,7 +268,7 @@ export namespace Reasoning {
export namespace Tool {
const ToolBase = {
...Base,
assistantMessageID: EventV2.ID,
assistantCreatorEventID: EventV2.ID,
callID: Schema.String,
}

Expand Down
22 changes: 20 additions & 2 deletions packages/core/src/session/input.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
export * as SessionInput from "./input"

import { and, asc, eq, inArray, isNull } from "drizzle-orm"
import { and, asc, eq, inArray, isNull, lte } from "drizzle-orm"
import { DateTime, Effect, Schema } from "effect"
import type { Database } from "../database/database"
import type { EventV2 } from "../event"
import { EventSequenceTable } from "../event/sql"
import { NonNegativeInt } from "../schema"
import { V2Schema } from "../v2-schema"
import { SessionEvent } from "./event"
Expand All @@ -24,6 +25,7 @@ export class Admitted extends Schema.Class<Admitted>("SessionInput.Admitted")({
prompt: Prompt,
delivery: Delivery,
timeCreated: V2Schema.DateTimeUtcFromMillis,
state: Schema.Literals(["pending", "promoted"]),
promotedSeq: NonNegativeInt.pipe(Schema.optional),
}) {}

Expand All @@ -39,6 +41,7 @@ const fromRow = (row: typeof SessionInputTable.$inferSelect): Admitted =>
prompt: decodePrompt(row.prompt),
delivery: row.delivery,
timeCreated: DateTime.makeUnsafe(row.time_created),
state: row.promoted_seq === null ? "pending" : "promoted",
...(row.promoted_seq === null ? {} : { promotedSeq: row.promoted_seq }),
})

Expand Down Expand Up @@ -91,6 +94,19 @@ export const pending = Effect.fn("SessionInput.pending")(function* (db: Database
.pipe(Effect.orDie)).map(fromRow)
})

export const latestSeq = Effect.fn("SessionInput.latestSeq")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
) {
const row = yield* db
.select({ seq: EventSequenceTable.seq })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, sessionID))
.get()
.pipe(Effect.orDie)
return row?.seq ?? -1
})

export const projectAdmitted = Effect.fn("SessionInput.projectAdmitted")(function* (
db: DatabaseService,
input: {
Expand Down Expand Up @@ -192,7 +208,7 @@ export const guardReservedID = Effect.fn("SessionInput.guardReservedID")(functio
if (Schema.is(SessionEvent.PromptLifecycle.Admitted)(event)) return
const id = Schema.is(SessionEvent.PromptLifecycle.Promoted)(event)
? event.data.messageID
: SessionMessage.ID.fromEvent(event.id)
: SessionMessage.ID.fromCreatorEvent(event.id)
const admitted = yield* find(db, id)
if (admitted === undefined) return
if (Schema.is(SessionEvent.PromptLifecycle.Promoted)(event)) return
Expand Down Expand Up @@ -301,6 +317,7 @@ export const promoteSteers = Effect.fn("SessionInput.promoteSteers")(function* (
db: DatabaseService,
events: EventV2.Interface,
sessionID: SessionSchema.ID,
cutoff: number,
) {
const rows = yield* db
.select()
Expand All @@ -310,6 +327,7 @@ export const promoteSteers = Effect.fn("SessionInput.promoteSteers")(function* (
eq(SessionInputTable.session_id, sessionID),
isNull(SessionInputTable.promoted_seq),
eq(SessionInputTable.delivery, "steer"),
lte(SessionInputTable.admitted_seq, cutoff),
),
)
.orderBy(asc(SessionInputTable.admitted_seq))
Expand Down
4 changes: 2 additions & 2 deletions packages/core/src/session/message-id.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,8 @@ export const ID = Schema.String.check(Schema.isStartsWith("msg_")).pipe(
Schema.brand("Session.Message.ID"),
withStatics((schema) => ({
create: () => schema.make("msg_" + Identifier.ascending()),
fromEvent: (id: EventV2.ID) => schema.make("msg" + id.slice(3)),
toEvent: (id: ID) => EventV2.ID.make("evt" + id.slice(3)),
fromCreatorEvent: (id: EventV2.ID) => schema.make("msg" + id.slice(3)),
toCreatorEvent: (id: ID) => EventV2.ID.make("evt" + id.slice(3)),
})),
)
export type ID = typeof ID.Type
30 changes: 15 additions & 15 deletions packages/core/src/session/message-updater.ts
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.agent.switched": (event) => {
return adapter.appendMessage(
new SessionMessage.AgentSwitched({
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "agent-switched",
metadata: event.metadata,
agent: event.data.agent,
Expand All @@ -134,7 +134,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.model.switched": (event) => {
return adapter.appendMessage(
new SessionMessage.ModelSwitched({
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "model-switched",
metadata: event.metadata,
model: event.data.model,
Expand All @@ -146,7 +146,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.prompted": (event) => {
return adapter.appendMessage(
new SessionMessage.User({
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "user",
metadata: event.metadata,
text: event.data.prompt.text,
Expand All @@ -164,7 +164,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
new SessionMessage.Synthetic({
sessionID: event.data.sessionID,
text: event.data.text,
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "synthetic",
time: { created: event.data.timestamp },
}),
Expand All @@ -173,7 +173,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.shell.started": (event) => {
return adapter.appendMessage(
new SessionMessage.Shell({
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "shell",
metadata: event.metadata,
callID: event.data.callID,
Expand Down Expand Up @@ -208,7 +208,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
yield* adapter.appendMessage(
new SessionMessage.Assistant({
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "assistant",
agent: event.data.agent,
model: event.data.model,
Expand All @@ -220,7 +220,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
})
},
"session.next.step.ended": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
draft.time.completed = event.data.timestamp
draft.finish = event.data.finish
draft.cost = event.data.cost
Expand All @@ -229,7 +229,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
})
},
"session.next.step.failed": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
draft.time.completed = event.data.timestamp
draft.finish = "error"
draft.error = event.data.error
Expand Down Expand Up @@ -276,7 +276,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
})
},
"session.next.tool.input.started": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
draft.content.push(
castDraft(
new SessionMessage.AssistantTool({
Expand All @@ -292,13 +292,13 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
},
"session.next.tool.input.delta": () => Effect.void,
"session.next.tool.input.ended": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "pending") match.state.input = event.data.text
})
},
"session.next.tool.called": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
const match = latestTool(draft, event.data.callID)
if (match) {
match.provider = event.data.provider
Expand All @@ -315,7 +315,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
})
},
"session.next.tool.progress": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
match.state.structured = event.data.structured
Expand All @@ -324,7 +324,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
})
},
"session.next.tool.success": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
match.provider = {
Expand All @@ -346,7 +346,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
})
},
"session.next.tool.failed": (event) => {
return updateOwnedAssistant(SessionMessage.ID.fromEvent(event.data.assistantMessageID), (draft) => {
return updateOwnedAssistant(SessionMessage.ID.fromCreatorEvent(event.data.assistantCreatorEventID), (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && (match.state.status === "pending" || match.state.status === "running")) {
match.provider = {
Expand Down Expand Up @@ -422,7 +422,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.next.compaction.started": (event) => {
return adapter.appendMessage(
new SessionMessage.Compaction({
id: SessionMessage.ID.fromEvent(event.id),
id: SessionMessage.ID.fromCreatorEvent(event.id),
type: "compaction",
metadata: event.metadata,
reason: event.data.reason,
Expand Down
2 changes: 1 addition & 1 deletion packages/core/src/session/projector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -367,7 +367,7 @@ export const layer = Layer.effectDiscard(
)
yield* events.project(SessionEvent.Prompted, (event) =>
Effect.gen(function* () {
const messageID = SessionMessage.ID.fromEvent(event.id)
const messageID = SessionMessage.ID.fromCreatorEvent(event.id)
const existing = yield* db
.select({ id: SessionMessageTable.id })
.from(SessionMessageTable)
Expand Down
15 changes: 9 additions & 6 deletions packages/core/src/session/runner/llm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -107,7 +107,7 @@ export const layer = Layer.effect(
yield* events.publish(SessionEvent.Tool.Failed, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: SessionMessage.ID.toEvent(message.id),
assistantCreatorEventID: SessionMessage.ID.toCreatorEvent(message.id),
callID: tool.id,
error: { type: "unknown", message: "Tool execution interrupted" },
provider: {
Expand All @@ -133,10 +133,13 @@ export const layer = Layer.effect(
const model = yield* models.resolve(session)
const toolFibers = yield* FiberSet.make<void, never>()
let needsContinuation = false
if (promotion === "steer") yield* SessionInput.promoteSteers(db, events, session.id)
if (promotion === "queue") {
yield* SessionInput.promoteNextQueued(db, events, session.id)
yield* SessionInput.promoteSteers(db, events, session.id)
if (promotion) {
const cutoff = yield* SessionInput.latestSeq(db, session.id)
if (promotion === "steer") yield* SessionInput.promoteSteers(db, events, session.id, cutoff)
if (promotion === "queue") {
yield* SessionInput.promoteNextQueued(db, events, session.id)
yield* SessionInput.promoteSteers(db, events, session.id, cutoff)
}
}
yield* failInterruptedTools(session.id)
const context = yield* getContext(session.id)
Expand Down Expand Up @@ -199,7 +202,7 @@ export const layer = Layer.effect(
events.publish(SessionEvent.Step.Failed, {
sessionID: session.id,
timestamp: yield* DateTime.now,
assistantMessageID: yield* publisher.startAssistant(),
assistantCreatorEventID: yield* publisher.startAssistant(),
error: { type: "unknown", message: llmFailure.reason.message },
}),
)
Expand Down
Loading