Compare commits

...

3 Commits

Author SHA1 Message Date
Kit Langton f1ff8cdebe refactor(core): tighten restart recovery scan 2026-07-07 20:58:22 -04:00
Kit Langton 6c16d4da0f refactor(core): simplify session recovery ownership 2026-07-07 20:53:42 -04:00
Kit Langton c05387ec9c fix(core): resume sessions after restart 2026-07-07 20:45:43 -04:00
2 changed files with 102 additions and 2 deletions
+70 -1
View File
@@ -1,5 +1,8 @@
import { sql } from "drizzle-orm"
import { Cause, Effect, Exit, Layer } from "effect"
import { Database } from "../../database/database"
import { EventV2 } from "../../event"
import { EventTable } from "../../event/sql"
import { LocationServiceMap } from "../../location-service-map"
import { makeGlobalNode } from "../../effect/app-node"
import { SessionEvent } from "../event"
@@ -10,6 +13,7 @@ import { SessionStore } from "../store"
import { SessionExecution } from "../execution"
import { toSessionError } from "../to-session-error"
import { UserInterruptedError } from "../error"
import { EffectFlock } from "../../util/effect-flock"
export function terminal(exit: Exit.Exit<void, SessionRunner.RunError>, reason?: "user" | "shutdown" | "superseded") {
if (Exit.isSuccess(exit)) return { type: "succeeded" as const }
@@ -19,6 +23,43 @@ export function terminal(exit: Exit.Exit<void, SessionRunner.RunError>, reason?:
return { type: "failed" as const, error: toSessionError(failure) }
}
export const sessionsInterruptedByShutdown = Effect.fn("SessionExecutionLocal.sessionsInterruptedByShutdown")(
function* (db: Database.Interface["db"]) {
const started = EventV2.versionedType(
SessionEvent.Execution.Started.type,
SessionEvent.Execution.Started.durable.version,
)
const succeeded = EventV2.versionedType(
SessionEvent.Execution.Succeeded.type,
SessionEvent.Execution.Succeeded.durable.version,
)
const failed = EventV2.versionedType(
SessionEvent.Execution.Failed.type,
SessionEvent.Execution.Failed.durable.version,
)
const interrupted = EventV2.versionedType(
SessionEvent.Execution.Interrupted.type,
SessionEvent.Execution.Interrupted.durable.version,
)
const latest = yield* db.all<{ sessionID: string }>(sql`
SELECT aggregate_id AS sessionID
FROM (
SELECT aggregate_id, type, data,
row_number() OVER (PARTITION BY aggregate_id ORDER BY seq DESC) AS rank
FROM ${EventTable}
WHERE type IN (
${started},
${succeeded},
${failed},
${interrupted}
)
)
WHERE rank = 1 AND type = ${interrupted} AND json_extract(data, '$.reason') = 'shutdown'
`)
return latest.map((event) => SessionSchema.ID.make(event.sessionID))
},
)
/** Current-process routing for implicit-local Locations. Future remote placement belongs here. */
const layer = Layer.effect(
SessionExecution.Service,
@@ -26,6 +67,8 @@ const layer = Layer.effect(
const store = yield* SessionStore.Service
const locations = yield* LocationServiceMap.Service
const events = yield* EventV2.Service
const flock = yield* EffectFlock.Service
const { db } = yield* Database.Service
const reportLifecycle = <A>(sessionID: SessionSchema.ID, effect: Effect.Effect<A>) =>
effect.pipe(
Effect.tapCause((cause) =>
@@ -77,6 +120,32 @@ const layer = Layer.effect(
),
})
yield* flock
.withLock(
Effect.gen(function* () {
const interrupted = yield* sessionsInterruptedByShutdown(db)
yield* Effect.forEach(
interrupted,
(sessionID) =>
coordinator.run(sessionID).pipe(
Effect.tapCause((cause) =>
Effect.logError("Failed to recover Session after shutdown", cause).pipe(
Effect.annotateLogs({ sessionID }),
),
),
Effect.ignore,
),
{ concurrency: "unbounded", discard: true },
)
}),
`session-shutdown-recovery:${Database.path()}`,
)
.pipe(
Effect.tapCause((cause) => Effect.logError("Failed to recover Sessions after shutdown", cause)),
Effect.ignore,
Effect.forkScoped,
)
return SessionExecution.Service.of({
active: coordinator.active,
interrupt: (sessionID) => coordinator.interrupt(sessionID, "user"),
@@ -90,7 +159,7 @@ const layer = Layer.effect(
export const node = makeGlobalNode({
service: SessionExecution.Service,
layer,
deps: [SessionStore.node, LocationServiceMap.node, EventV2.node],
deps: [SessionStore.node, LocationServiceMap.node, EventV2.node, Database.node, EffectFlock.node],
})
export * as SessionExecutionLocal from "./local"
@@ -1,9 +1,18 @@
import { describe, expect, test } from "bun:test"
import { LLMError, TransportReason } from "@opencode-ai/llm"
import { terminal } from "@opencode-ai/core/session/execution/local"
import { Database } from "@opencode-ai/core/database/database"
import { EventV2 } from "@opencode-ai/core/event"
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { SessionEvent } from "@opencode-ai/core/session/event"
import { sessionsInterruptedByShutdown, terminal } from "@opencode-ai/core/session/execution/local"
import { SessionV2 } from "@opencode-ai/core/session"
import { UserInterruptedError } from "@opencode-ai/core/session/error"
import { ToolOutputStore } from "@opencode-ai/core/tool-output-store"
import { Effect, Exit } from "effect"
import { testEffect } from "./lib/effect"
const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node])))
describe("SessionExecutionLocal lifecycle", () => {
test("classifies success and typed failure terminals", () => {
@@ -33,4 +42,26 @@ describe("SessionExecutionLocal lifecycle", () => {
expect(terminal(interrupted, "superseded")).toEqual({ type: "interrupted", reason: "superseded" })
expect(terminal(Exit.fail(new UserInterruptedError()))).toEqual({ type: "interrupted", reason: "user" })
})
it.effect("selects only sessions whose latest execution ended for shutdown", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const { db } = yield* Database.Service
const recover = SessionV2.ID.make("ses_recover")
const completed = SessionV2.ID.make("ses_completed")
const user = SessionV2.ID.make("ses_user")
const crashed = SessionV2.ID.make("ses_crashed")
yield* events.publish(SessionEvent.Execution.Started, { sessionID: recover })
yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID: recover, reason: "shutdown" })
yield* events.publish(SessionEvent.InstructionsUpdated, { sessionID: recover, text: "later non-lifecycle event" })
yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID: completed, reason: "shutdown" })
yield* events.publish(SessionEvent.Execution.Started, { sessionID: completed })
yield* events.publish(SessionEvent.Execution.Succeeded, { sessionID: completed })
yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID: user, reason: "user" })
yield* events.publish(SessionEvent.Execution.Started, { sessionID: crashed })
expect(yield* sessionsInterruptedByShutdown(db)).toEqual([recover])
}),
)
})