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): fence provider dispatch by epoch
  • Loading branch information
kitlangton committed Jun 5, 2026
commit 608a203949a4ae1bb0bdd4994fdad0985443d04a
27 changes: 22 additions & 5 deletions packages/core/src/session/context-epoch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ const retryRevisionMismatch = <A, E>(attempt: () => Effect.Effect<A, E>): Effect
interface Prepared {
readonly baseline: string
readonly baselineSeq: number
readonly revision: number
}

export function initialize(
Expand Down Expand Up @@ -78,7 +79,7 @@ const prepareOnce = Effect.fnUntraced(function* (
if (!stored) {
const generation = yield* SystemContext.initialize(value)
const baselineSeq = yield* insert(db, sessionID, location, agent, generation)
return { baseline: generation.baseline, baselineSeq }
return { baseline: generation.baseline, baselineSeq, revision: 0 }
}

const snapshot = yield* Schema.decodeUnknownEffect(SystemContext.Snapshot)(stored.snapshot).pipe(
Expand All @@ -95,20 +96,20 @@ const prepareOnce = Effect.fnUntraced(function* (
}
if (result._tag === "Unchanged" || result._tag === "ReplacementBlocked") {
yield* fence(db, sessionID, agent, stored.revision)
return { baseline: stored.baseline, baselineSeq: stored.baseline_seq }
return { baseline: stored.baseline, baselineSeq: stored.baseline_seq, revision: stored.revision }
}
if (result._tag === "ReplacementReady") {
const replacementSeq = stored.replacement_seq ?? (yield* SessionInput.latestSeq(db, sessionID))
yield* replace(db, sessionID, agent, stored.revision, replacementSeq, result.generation)
return { baseline: result.generation.baseline, baselineSeq: replacementSeq }
return { baseline: result.generation.baseline, baselineSeq: replacementSeq, revision: stored.revision + 1 }
}

yield* events.publish(
SessionEvent.ContextUpdated,
{ sessionID, messageID: SessionMessageID.ID.create(), timestamp: yield* DateTime.now, text: result.text },
{ commit: () => advance(db, sessionID, stored.revision, result.snapshot).pipe(Effect.orDie) },
)
return { baseline: stored.baseline, baselineSeq: stored.baseline_seq }
return { baseline: stored.baseline, baselineSeq: stored.baseline_seq, revision: stored.revision + 1 }
})

const initializeOnce = Effect.fnUntraced(function* (
Expand All @@ -121,7 +122,7 @@ const initializeOnce = Effect.fnUntraced(function* (
if (yield* exists(db, sessionID)) return
const generation = yield* context.pipe(Effect.flatMap(SystemContext.initialize))
const baselineSeq = yield* insert(db, sessionID, location, agent, generation)
return { baseline: generation.baseline, baselineSeq }
return { baseline: generation.baseline, baselineSeq, revision: 0 }
})

const exists = Effect.fn("SessionContextEpoch.exists")(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
Expand Down Expand Up @@ -296,6 +297,22 @@ const fence = Effect.fnUntraced(function* (
if (current.revision !== expectedRevision) return yield* Effect.die(new RevisionMismatch())
})

export const current = Effect.fn("SessionContextEpoch.current")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
agent: AgentV2.ID,
revision: number,
) {
const value = yield* db
.select({ agent: SessionTable.agent, revision: SessionContextEpochTable.revision })
.from(SessionContextEpochTable)
.innerJoin(SessionTable, eq(SessionTable.id, SessionContextEpochTable.session_id))
.where(eq(SessionContextEpochTable.session_id, sessionID))
.get()
.pipe(Effect.orDie)
return value !== undefined && AgentV2.effectiveID(value.agent) === agent && value.revision === revision
})

const advance = Effect.fnUntraced(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
Expand Down
12 changes: 11 additions & 1 deletion packages/core/src/session/runner/llm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,13 @@ export const layer = Layer.effect(
defect instanceof SessionContextEpoch.AgentMismatch ? Effect.die(new RetryTurn(promotion)) : Effect.die(defect),
)

const sameModel = (left: ModelV2.Ref | undefined, right: ModelV2.Ref | undefined) =>
left === right ||
(left !== undefined &&
right !== undefined &&
left.id === right.id &&
left.providerID === right.providerID &&
left.variant === right.variant)
const loadSystemContext = (agent: AgentV2.ID) =>
Effect.all([systemContext.load(), skillGuidance.load(agent)], { concurrency: "unbounded" }).pipe(
Effect.map(SystemContext.combine),
Expand Down Expand Up @@ -173,7 +180,8 @@ export const layer = Layer.effect(
agentID,
).pipe(retryAgentMismatch(undefined)))
const current = yield* getSession(sessionID)
if ((yield* agents.resolve(current.agent))?.id !== agent?.id) return yield* Effect.die(new RetryTurn(undefined))
if ((yield* agents.resolve(current.agent))?.id !== agent?.id || !sameModel(current.model, session.model))
return yield* Effect.die(new RetryTurn(undefined))
const model = yield* models.resolve(session)
const context = yield* store.runnerContext(session.id, system.baselineSeq)
const request = LLM.request({
Expand All @@ -195,6 +203,8 @@ export const layer = Layer.effect(
})
const withPublication = Semaphore.makeUnsafe(1).withPermit
const publish = (event: LLMEvent) => withPublication(publisher.publish(event))
if (!(yield* SessionContextEpoch.current(db, session.id, agentID, system.revision)))
return yield* Effect.die(new RetryTurn(undefined))
const providerStream = llm.stream(request).pipe(
Stream.runForEach((event) =>
Effect.gen(function* () {
Expand Down
21 changes: 4 additions & 17 deletions packages/core/test/session-runner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -925,7 +925,7 @@ describe("SessionRunnerLLM", () => {
}),
)

it.effect("applies an agent switch after the safe boundary to the next provider turn", () =>
it.effect("retries an agent switch before the final provider-dispatch boundary", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
Expand All @@ -951,12 +951,7 @@ describe("SessionRunnerLLM", () => {
requests.length = 0
response = []
yield* session.resume(sessionID)
modelResolveHook = Effect.void
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
yield* session.resume(sessionID)

expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
["Initial context\n\nBuild skills"],
["Initial context\n\nReviewer skills"],
])
expect(
Expand All @@ -970,7 +965,7 @@ describe("SessionRunnerLLM", () => {
}),
)

it.effect("applies a model switch after the safe boundary to the next provider turn", () =>
it.effect("retries a model switch before the final provider-dispatch boundary", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
Expand All @@ -993,16 +988,8 @@ describe("SessionRunnerLLM", () => {
requests.length = 0
response = []
yield* session.resume(sessionID)
modelResolveHook = Effect.void
systemBaseline = "Replacement context"
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
yield* session.resume(sessionID)

expect(requests.map((request) => request.model)).toEqual([model, replacementModel])
expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
["Initial context"],
["Replacement context"],
])
expect(requests.map((request) => request.model)).toEqual([replacementModel])
expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([["Initial context"]])
}),
)

Expand Down