fix(core): share session runtime coordination
This commit is contained in:
parent
8a819da177
commit
d3ace428d6
5 changed files with 80 additions and 20 deletions
|
|
@ -217,6 +217,7 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: PluginV2.Int
|
|||
session.create(
|
||||
input.parentID
|
||||
? {
|
||||
id: input.id ? SessionV2.ID.make(input.id) : undefined,
|
||||
parentID: SessionV2.ID.make(input.parentID),
|
||||
title: input.title,
|
||||
agent: input.agent ? AgentV2.ID.make(input.agent) : undefined,
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ import { ProjectTable } from "./project/sql"
|
|||
import path from "path"
|
||||
import { fromRow } from "./session/info"
|
||||
import { SessionRunner } from "./session/runner/index"
|
||||
import { SessionRuntimeCoordinator } from "./session/runtime-coordinator"
|
||||
import { SessionStore } from "./session/store"
|
||||
import { makeGlobalNode } from "./effect/app-node"
|
||||
import { MessageDecodeError } from "./session/error"
|
||||
|
|
@ -192,6 +193,7 @@ export const layer = Layer.effect(
|
|||
const db = database.db
|
||||
const events = yield* EventV2.Service
|
||||
const projects = yield* ProjectV2.Service
|
||||
const runtimeCoordinator = yield* SessionRuntimeCoordinator.Service
|
||||
const store = yield* SessionStore.Service
|
||||
const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Message)
|
||||
const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
|
||||
|
|
@ -437,7 +439,7 @@ export const layer = Layer.effect(
|
|||
yield* result.get(sessionID)
|
||||
return yield* Effect.die("SessionV2.wait moved to SessionRuntime.Service")
|
||||
}),
|
||||
active: Effect.succeed(new Set()),
|
||||
active: runtimeCoordinator.active,
|
||||
resume: Effect.fn("V2Session.resume")(function* (sessionID) {
|
||||
yield* result.get(sessionID)
|
||||
return yield* Effect.die("SessionV2.resume moved to SessionRuntime.Service")
|
||||
|
|
@ -469,6 +471,7 @@ export const defaultLayer = layer.pipe(
|
|||
Layer.provide(EventV2.defaultLayer),
|
||||
Layer.provide(Database.defaultLayer),
|
||||
Layer.provide(ProjectV2.defaultLayer),
|
||||
Layer.provide(SessionRuntimeCoordinator.defaultLayer),
|
||||
Layer.orDie,
|
||||
)
|
||||
|
||||
|
|
@ -494,6 +497,7 @@ export const node = makeGlobalNode({
|
|||
Database.node,
|
||||
EventV2.node,
|
||||
ProjectV2.node,
|
||||
SessionRuntimeCoordinator.node,
|
||||
SessionStore.node,
|
||||
SessionProjector.node,
|
||||
],
|
||||
|
|
|
|||
47
packages/core/src/session/runtime-coordinator.ts
Normal file
47
packages/core/src/session/runtime-coordinator.ts
Normal file
|
|
@ -0,0 +1,47 @@
|
|||
export * as SessionRuntimeCoordinator from "./runtime-coordinator"
|
||||
|
||||
import { Context, Effect, Layer } from "effect"
|
||||
import { makeGlobalNode } from "../effect/app-node"
|
||||
import { SessionRunner } from "./runner"
|
||||
import { SessionRunCoordinator } from "./run-coordinator"
|
||||
import { SessionSchema } from "./schema"
|
||||
|
||||
type Drain = (force: boolean) => Effect.Effect<void, SessionRunner.RunError>
|
||||
|
||||
export interface Interface {
|
||||
readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>>
|
||||
readonly run: (sessionID: SessionSchema.ID, drain: Drain) => Effect.Effect<void, SessionRunner.RunError>
|
||||
readonly wake: (sessionID: SessionSchema.ID, drain: Drain) => Effect.Effect<void>
|
||||
readonly interrupt: (sessionID: SessionSchema.ID) => Effect.Effect<void>
|
||||
readonly wait: (sessionID: SessionSchema.ID) => Effect.Effect<void>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/v2/SessionRuntimeCoordinator") {}
|
||||
|
||||
export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const drains = new Map<SessionSchema.ID, Drain>()
|
||||
const coordinator = yield* SessionRunCoordinator.make<SessionSchema.ID, SessionRunner.RunError>({
|
||||
drain: Effect.fnUntraced(function* (sessionID, force) {
|
||||
const drain = drains.get(sessionID)
|
||||
if (!drain) return yield* Effect.die(`No SessionRuntime drain registered for ${sessionID}`)
|
||||
return yield* drain(force)
|
||||
}),
|
||||
})
|
||||
|
||||
return Service.of({
|
||||
active: coordinator.active,
|
||||
run: (sessionID, drain) =>
|
||||
Effect.sync(() => drains.set(sessionID, drain)).pipe(Effect.andThen(coordinator.run(sessionID))),
|
||||
wake: (sessionID, drain) =>
|
||||
Effect.sync(() => drains.set(sessionID, drain)).pipe(Effect.andThen(coordinator.wake(sessionID))),
|
||||
interrupt: coordinator.interrupt,
|
||||
wait: coordinator.awaitIdle,
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export const defaultLayer = layer
|
||||
|
||||
export const node = makeGlobalNode({ service: Service, layer, deps: [] })
|
||||
|
|
@ -11,7 +11,7 @@ import { SessionInput } from "./input"
|
|||
import { SessionRevert } from "./revert"
|
||||
import { SessionRunner } from "./runner"
|
||||
import * as SessionRunnerLLM from "./runner/llm"
|
||||
import { SessionRunCoordinator } from "./run-coordinator"
|
||||
import { SessionRuntimeCoordinator } from "./runtime-coordinator"
|
||||
import { SessionSchema } from "./schema"
|
||||
import { Snapshot } from "../snapshot"
|
||||
import { FSUtil } from "../fs-util"
|
||||
|
|
@ -59,6 +59,8 @@ export const layer = Layer.effect(
|
|||
const location = yield* Location.Service
|
||||
const sessions = yield* SessionV2.Service
|
||||
const runner = yield* SessionRunner.Service
|
||||
const snapshot = yield* Snapshot.Service
|
||||
const coordinator = yield* SessionRuntimeCoordinator.Service
|
||||
|
||||
const local = Effect.fn("SessionRuntime.local")(function* (sessionID: SessionSchema.ID) {
|
||||
const session = yield* sessions.get(sessionID)
|
||||
|
|
@ -67,12 +69,11 @@ export const layer = Layer.effect(
|
|||
return session
|
||||
})
|
||||
|
||||
const coordinator = yield* SessionRunCoordinator.make<SessionSchema.ID, SessionRunner.RunError>({
|
||||
drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, force) {
|
||||
const drain = (sessionID: SessionSchema.ID) =>
|
||||
Effect.fnUntraced(function* (force: boolean) {
|
||||
yield* local(sessionID).pipe(Effect.orDie)
|
||||
return yield* runner.run({ sessionID, force })
|
||||
}),
|
||||
})
|
||||
})
|
||||
|
||||
return Service.of({
|
||||
prompt: Effect.fn("SessionRuntime.prompt")((input) =>
|
||||
|
|
@ -97,21 +98,24 @@ export const layer = Layer.effect(
|
|||
)
|
||||
if (!SessionInput.equivalent(admitted, expected))
|
||||
return yield* new PromptConflictError({ sessionID: input.sessionID, messageID })
|
||||
if (input.resume !== false) yield* coordinator.wake(admitted.sessionID)
|
||||
if (input.resume !== false) yield* coordinator.wake(admitted.sessionID, drain(admitted.sessionID))
|
||||
return admitted
|
||||
}),
|
||||
),
|
||||
),
|
||||
wait: Effect.fn("SessionRuntime.wait")(function* (sessionID) {
|
||||
yield* local(sessionID)
|
||||
yield* coordinator.awaitIdle(sessionID)
|
||||
yield* coordinator.wait(sessionID)
|
||||
}),
|
||||
active: coordinator.active,
|
||||
resume: Effect.fn("SessionRuntime.resume")(function* (sessionID) {
|
||||
yield* local(sessionID)
|
||||
yield* coordinator.run(sessionID)
|
||||
yield* coordinator.run(sessionID, drain(sessionID))
|
||||
}),
|
||||
interrupt: Effect.fn("SessionRuntime.interrupt")(function* (sessionID) {
|
||||
yield* local(sessionID)
|
||||
yield* Effect.uninterruptible(coordinator.interrupt(sessionID))
|
||||
}),
|
||||
interrupt: Effect.fn("SessionRuntime.interrupt")((sessionID) => Effect.uninterruptible(coordinator.interrupt(sessionID))),
|
||||
revert: {
|
||||
stage: Effect.fn("SessionRuntime.revert.stage")(function* (input) {
|
||||
const session = yield* local(input.sessionID)
|
||||
|
|
@ -119,12 +123,16 @@ export const layer = Layer.effect(
|
|||
return yield* SessionRevert.stage({ session, messageID: input.messageID, files: input.files }).pipe(
|
||||
Effect.provideService(Database.Service, database),
|
||||
Effect.provideService(EventV2.Service, events),
|
||||
Effect.provideService(Snapshot.Service, snapshot),
|
||||
)
|
||||
}),
|
||||
clear: Effect.fn("SessionRuntime.revert.clear")(function* (sessionID) {
|
||||
const session = yield* local(sessionID)
|
||||
if ((yield* coordinator.active).has(sessionID)) return yield* new BusyError({ sessionID })
|
||||
return yield* SessionRevert.clear(session).pipe(Effect.provideService(EventV2.Service, events))
|
||||
return yield* SessionRevert.clear(session).pipe(
|
||||
Effect.provideService(EventV2.Service, events),
|
||||
Effect.provideService(Snapshot.Service, snapshot),
|
||||
)
|
||||
}),
|
||||
commit: Effect.fn("SessionRuntime.revert.commit")(function* (sessionID) {
|
||||
const session = yield* local(sessionID)
|
||||
|
|
@ -153,5 +161,5 @@ const resolvePrompt = (input: PromptInput.Prompt) =>
|
|||
export const node = makeLocationNode({
|
||||
service: Service,
|
||||
layer,
|
||||
deps: [Database.node, EventV2.node, Location.node, SessionV2.node, SessionRunnerLLM.node, Snapshot.node],
|
||||
deps: [Database.node, EventV2.node, Location.node, SessionV2.node, SessionRunnerLLM.node, Snapshot.node, SessionRuntimeCoordinator.node],
|
||||
})
|
||||
|
|
|
|||
|
|
@ -8,23 +8,23 @@ export interface SessionDomain {
|
|||
readonly title?: string
|
||||
readonly agent?: string
|
||||
readonly model?: SessionV2Info["model"]
|
||||
}) => Effect.Effect<SessionV2Info>
|
||||
readonly get: (sessionID: string) => Effect.Effect<SessionV2Info>
|
||||
}) => Effect.Effect<SessionV2Info, unknown>
|
||||
readonly get: (sessionID: string) => Effect.Effect<SessionV2Info, unknown>
|
||||
readonly messages: (input: {
|
||||
readonly sessionID: string
|
||||
readonly limit?: number
|
||||
readonly order?: "asc" | "desc"
|
||||
readonly cursor?: { readonly id: string; readonly direction: "previous" | "next" }
|
||||
}) => Effect.Effect<ReadonlyArray<SessionMessage>>
|
||||
readonly context: (sessionID: string) => Effect.Effect<ReadonlyArray<SessionMessage>>
|
||||
}) => Effect.Effect<ReadonlyArray<SessionMessage>, unknown>
|
||||
readonly context: (sessionID: string) => Effect.Effect<ReadonlyArray<SessionMessage>, unknown>
|
||||
readonly prompt: (input: {
|
||||
readonly id?: string
|
||||
readonly sessionID: string
|
||||
readonly prompt: PromptInput
|
||||
readonly delivery?: "steer" | "queue"
|
||||
readonly resume?: boolean
|
||||
}) => Effect.Effect<SessionInputAdmitted>
|
||||
readonly resume: (sessionID: string) => Effect.Effect<void>
|
||||
readonly wait: (sessionID: string) => Effect.Effect<void>
|
||||
readonly interrupt: (sessionID: string) => Effect.Effect<void>
|
||||
}) => Effect.Effect<SessionInputAdmitted, unknown>
|
||||
readonly resume: (sessionID: string) => Effect.Effect<void, unknown>
|
||||
readonly wait: (sessionID: string) => Effect.Effect<void, unknown>
|
||||
readonly interrupt: (sessionID: string) => Effect.Effect<void, unknown>
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue