refactor(core): add location session runtime
This commit is contained in:
parent
70cecc6ba1
commit
89ef53537e
12 changed files with 356 additions and 46 deletions
|
|
@ -26,6 +26,7 @@ import { Reference } from "./reference"
|
||||||
import { ReferenceGuidance } from "./reference/guidance"
|
import { ReferenceGuidance } from "./reference/guidance"
|
||||||
import * as SessionRunnerLLM from "./session/runner/llm"
|
import * as SessionRunnerLLM from "./session/runner/llm"
|
||||||
import { SessionRunnerModel } from "./session/runner/model"
|
import { SessionRunnerModel } from "./session/runner/model"
|
||||||
|
import { SessionRuntime } from "./session/runtime"
|
||||||
import { SessionTodo } from "./session/todo"
|
import { SessionTodo } from "./session/todo"
|
||||||
import { SkillV2 } from "./skill"
|
import { SkillV2 } from "./skill"
|
||||||
import { SkillGuidance } from "./skill/guidance"
|
import { SkillGuidance } from "./skill/guidance"
|
||||||
|
|
@ -75,6 +76,7 @@ export const locationServices = LayerNode.group([
|
||||||
BuiltInTools.node,
|
BuiltInTools.node,
|
||||||
SessionRunnerModel.node,
|
SessionRunnerModel.node,
|
||||||
Snapshot.node,
|
Snapshot.node,
|
||||||
|
SessionRuntime.node,
|
||||||
SessionRunnerLLM.node,
|
SessionRunnerLLM.node,
|
||||||
])
|
])
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -13,7 +13,10 @@ import { Integration } from "./integration"
|
||||||
import { KeyedMutex } from "./effect/keyed-mutex"
|
import { KeyedMutex } from "./effect/keyed-mutex"
|
||||||
import { PluginHost } from "./plugin/host"
|
import { PluginHost } from "./plugin/host"
|
||||||
import { Reference } from "./reference"
|
import { Reference } from "./reference"
|
||||||
|
import { SessionV2 } from "./session"
|
||||||
|
import { SessionRuntime } from "./session/runtime"
|
||||||
import { SkillV2 } from "./skill"
|
import { SkillV2 } from "./skill"
|
||||||
|
import { Location } from "./location"
|
||||||
import { State } from "./state"
|
import { State } from "./state"
|
||||||
|
|
||||||
export const ID = Plugin.ID
|
export const ID = Plugin.ID
|
||||||
|
|
@ -149,6 +152,8 @@ export const locationLayer = layer.pipe(
|
||||||
Layer.provideMerge(CommandV2.locationLayer),
|
Layer.provideMerge(CommandV2.locationLayer),
|
||||||
Layer.provideMerge(Integration.locationLayer),
|
Layer.provideMerge(Integration.locationLayer),
|
||||||
Layer.provideMerge(Reference.locationLayer),
|
Layer.provideMerge(Reference.locationLayer),
|
||||||
|
Layer.provideMerge(SessionV2.defaultLayer),
|
||||||
|
Layer.provideMerge(SessionRuntime.layer),
|
||||||
Layer.provideMerge(SkillV2.locationLayer),
|
Layer.provideMerge(SkillV2.locationLayer),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -163,6 +168,9 @@ export const node = makeLocationNode({
|
||||||
CommandV2.node,
|
CommandV2.node,
|
||||||
Integration.node,
|
Integration.node,
|
||||||
Reference.node,
|
Reference.node,
|
||||||
|
SessionV2.node,
|
||||||
|
SessionRuntime.node,
|
||||||
SkillV2.node,
|
SkillV2.node,
|
||||||
|
Location.node,
|
||||||
],
|
],
|
||||||
})
|
})
|
||||||
|
|
|
||||||
|
|
@ -13,7 +13,11 @@ import { PluginV2 } from "../plugin"
|
||||||
import { ProviderV2 } from "../provider"
|
import { ProviderV2 } from "../provider"
|
||||||
import { Reference } from "../reference"
|
import { Reference } from "../reference"
|
||||||
import type { DeepMutable } from "../schema"
|
import type { DeepMutable } from "../schema"
|
||||||
|
import { SessionV2 } from "../session"
|
||||||
|
import { SessionMessage } from "../session/message"
|
||||||
|
import { SessionRuntime } from "../session/runtime"
|
||||||
import { SkillV2 } from "../skill"
|
import { SkillV2 } from "../skill"
|
||||||
|
import { Location } from "../location"
|
||||||
|
|
||||||
const mutable = <T>(value: T) => value as DeepMutable<T>
|
const mutable = <T>(value: T) => value as DeepMutable<T>
|
||||||
|
|
||||||
|
|
@ -24,7 +28,10 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: PluginV2.Int
|
||||||
const commands = yield* CommandV2.Service
|
const commands = yield* CommandV2.Service
|
||||||
const integration = yield* Integration.Service
|
const integration = yield* Integration.Service
|
||||||
const reference = yield* Reference.Service
|
const reference = yield* Reference.Service
|
||||||
|
const session = yield* SessionV2.Service
|
||||||
|
const sessionRuntime = yield* SessionRuntime.Service
|
||||||
const skill = yield* SkillV2.Service
|
const skill = yield* SkillV2.Service
|
||||||
|
const location = yield* Location.Service
|
||||||
|
|
||||||
return {
|
return {
|
||||||
options: {},
|
options: {},
|
||||||
|
|
@ -205,6 +212,59 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: PluginV2.Int
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
},
|
},
|
||||||
|
session: {
|
||||||
|
create: (input) =>
|
||||||
|
session.create(
|
||||||
|
input.parentID
|
||||||
|
? {
|
||||||
|
parentID: SessionV2.ID.make(input.parentID),
|
||||||
|
title: input.title,
|
||||||
|
agent: input.agent ? AgentV2.ID.make(input.agent) : undefined,
|
||||||
|
model: input.model
|
||||||
|
? {
|
||||||
|
id: ModelV2.ID.make(input.model.id),
|
||||||
|
providerID: ProviderV2.ID.make(input.model.providerID),
|
||||||
|
variant: input.model.variant ? ModelV2.VariantID.make(input.model.variant) : undefined,
|
||||||
|
}
|
||||||
|
: undefined,
|
||||||
|
}
|
||||||
|
: {
|
||||||
|
id: input.id ? SessionV2.ID.make(input.id) : undefined,
|
||||||
|
location: { directory: location.directory, workspaceID: location.workspaceID },
|
||||||
|
title: input.title,
|
||||||
|
agent: input.agent ? AgentV2.ID.make(input.agent) : undefined,
|
||||||
|
model: input.model
|
||||||
|
? {
|
||||||
|
id: ModelV2.ID.make(input.model.id),
|
||||||
|
providerID: ProviderV2.ID.make(input.model.providerID),
|
||||||
|
variant: input.model.variant ? ModelV2.VariantID.make(input.model.variant) : undefined,
|
||||||
|
}
|
||||||
|
: undefined,
|
||||||
|
},
|
||||||
|
),
|
||||||
|
get: (sessionID) => session.get(SessionV2.ID.make(sessionID)),
|
||||||
|
messages: (input) =>
|
||||||
|
session.messages({
|
||||||
|
sessionID: SessionV2.ID.make(input.sessionID),
|
||||||
|
limit: input.limit,
|
||||||
|
order: input.order,
|
||||||
|
cursor: input.cursor
|
||||||
|
? { id: SessionMessage.ID.make(input.cursor.id), direction: input.cursor.direction }
|
||||||
|
: undefined,
|
||||||
|
}),
|
||||||
|
context: (sessionID) => session.context(SessionV2.ID.make(sessionID)),
|
||||||
|
prompt: (input) =>
|
||||||
|
sessionRuntime.prompt({
|
||||||
|
id: input.id ? SessionMessage.ID.make(input.id) : undefined,
|
||||||
|
sessionID: SessionV2.ID.make(input.sessionID),
|
||||||
|
prompt: input.prompt,
|
||||||
|
delivery: input.delivery,
|
||||||
|
resume: input.resume,
|
||||||
|
}),
|
||||||
|
resume: (sessionID) => sessionRuntime.resume(SessionV2.ID.make(sessionID)),
|
||||||
|
wait: (sessionID) => sessionRuntime.wait(SessionV2.ID.make(sessionID)),
|
||||||
|
interrupt: (sessionID) => sessionRuntime.interrupt(SessionV2.ID.make(sessionID)),
|
||||||
|
},
|
||||||
skill: {
|
skill: {
|
||||||
reload: skill.reload,
|
reload: skill.reload,
|
||||||
transform: (callback) =>
|
transform: (callback) =>
|
||||||
|
|
|
||||||
|
|
@ -81,6 +81,16 @@ export function fromPromise(plugin: Plugin) {
|
||||||
transform: transform(host.reference),
|
transform: transform(host.reference),
|
||||||
reload: () => run(host.reference.reload()),
|
reload: () => run(host.reference.reload()),
|
||||||
},
|
},
|
||||||
|
session: {
|
||||||
|
create: (input) => run(host.session.create(input)),
|
||||||
|
get: (sessionID) => run(host.session.get(sessionID)),
|
||||||
|
messages: (input) => run(host.session.messages(input)),
|
||||||
|
context: (sessionID) => run(host.session.context(sessionID)),
|
||||||
|
prompt: (input) => run(host.session.prompt(input)),
|
||||||
|
resume: (sessionID) => run(host.session.resume(sessionID)),
|
||||||
|
wait: (sessionID) => run(host.session.wait(sessionID)),
|
||||||
|
interrupt: (sessionID) => run(host.session.interrupt(sessionID)),
|
||||||
|
},
|
||||||
skill: {
|
skill: {
|
||||||
transform: transform(host.skill),
|
transform: transform(host.skill),
|
||||||
reload: () => run(host.skill.reload()),
|
reload: () => run(host.skill.reload()),
|
||||||
|
|
|
||||||
|
|
@ -26,9 +26,7 @@ import path from "path"
|
||||||
import { fromRow } from "./session/info"
|
import { fromRow } from "./session/info"
|
||||||
import { SessionRunner } from "./session/runner/index"
|
import { SessionRunner } from "./session/runner/index"
|
||||||
import { SessionStore } from "./session/store"
|
import { SessionStore } from "./session/store"
|
||||||
import { SessionExecution } from "./session/execution"
|
|
||||||
import { makeGlobalNode } from "./effect/app-node"
|
import { makeGlobalNode } from "./effect/app-node"
|
||||||
import { LocationServiceMap } from "./location-service-map"
|
|
||||||
import { MessageDecodeError } from "./session/error"
|
import { MessageDecodeError } from "./session/error"
|
||||||
import { SessionEvent } from "./session/event"
|
import { SessionEvent } from "./session/event"
|
||||||
import { SessionInput } from "./session/input"
|
import { SessionInput } from "./session/input"
|
||||||
|
|
@ -194,9 +192,7 @@ export const layer = Layer.effect(
|
||||||
const db = database.db
|
const db = database.db
|
||||||
const events = yield* EventV2.Service
|
const events = yield* EventV2.Service
|
||||||
const projects = yield* ProjectV2.Service
|
const projects = yield* ProjectV2.Service
|
||||||
const execution = yield* SessionExecution.Service
|
|
||||||
const store = yield* SessionStore.Service
|
const store = yield* SessionStore.Service
|
||||||
const locations = yield* LocationServiceMap.Service
|
|
||||||
const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Message)
|
const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Message)
|
||||||
const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
|
const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
|
||||||
const decode = (row: typeof SessionMessageTable.$inferSelect) =>
|
const decode = (row: typeof SessionMessageTable.$inferSelect) =>
|
||||||
|
|
@ -389,7 +385,8 @@ export const layer = Layer.effect(
|
||||||
)
|
)
|
||||||
if (!SessionInput.equivalent(admitted, expected))
|
if (!SessionInput.equivalent(admitted, expected))
|
||||||
return yield* new PromptConflictError({ sessionID: input.sessionID, messageID })
|
return yield* new PromptConflictError({ sessionID: input.sessionID, messageID })
|
||||||
if (input.resume !== false) yield* execution.wake(admitted.sessionID)
|
if (input.resume !== false)
|
||||||
|
return yield* Effect.die("SessionV2.prompt with resume moved to SessionRuntime.Service")
|
||||||
return admitted
|
return admitted
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
|
|
@ -438,38 +435,25 @@ export const layer = Layer.effect(
|
||||||
}),
|
}),
|
||||||
wait: Effect.fn("V2Session.wait")(function* (sessionID) {
|
wait: Effect.fn("V2Session.wait")(function* (sessionID) {
|
||||||
yield* result.get(sessionID)
|
yield* result.get(sessionID)
|
||||||
yield* execution.awaitIdle(sessionID)
|
return yield* Effect.die("SessionV2.wait moved to SessionRuntime.Service")
|
||||||
}),
|
}),
|
||||||
active: execution.active,
|
active: Effect.succeed(new Set()),
|
||||||
resume: Effect.fn("V2Session.resume")(function* (sessionID) {
|
resume: Effect.fn("V2Session.resume")(function* (sessionID) {
|
||||||
yield* result.get(sessionID)
|
yield* result.get(sessionID)
|
||||||
yield* execution.resume(sessionID)
|
return yield* Effect.die("SessionV2.resume moved to SessionRuntime.Service")
|
||||||
}),
|
}),
|
||||||
interrupt: Effect.fn("V2Session.interrupt")((sessionID) =>
|
interrupt: Effect.fn("V2Session.interrupt")(() => Effect.die("SessionV2.interrupt moved to SessionRuntime.Service")),
|
||||||
Effect.uninterruptible(execution.interrupt(sessionID)),
|
|
||||||
),
|
|
||||||
revert: {
|
revert: {
|
||||||
stage: Effect.fn("V2Session.revert.stage")(function* (input) {
|
stage: Effect.fn("V2Session.revert.stage")(function* (input) {
|
||||||
const session = yield* result.get(input.sessionID)
|
yield* result.get(input.sessionID)
|
||||||
if ((yield* execution.active).has(input.sessionID))
|
return yield* Effect.die("SessionV2.revert.stage moved to SessionRuntime.Service")
|
||||||
return yield* new BusyError({ sessionID: input.sessionID })
|
|
||||||
return yield* SessionRevert.stage({ session, messageID: input.messageID, files: input.files }).pipe(
|
|
||||||
Effect.provideService(Database.Service, database),
|
|
||||||
Effect.provideService(EventV2.Service, events),
|
|
||||||
Effect.provide(locations.get(session.location)),
|
|
||||||
)
|
|
||||||
}),
|
}),
|
||||||
clear: Effect.fn("V2Session.revert.clear")(function* (sessionID) {
|
clear: Effect.fn("V2Session.revert.clear")(function* (sessionID) {
|
||||||
const session = yield* result.get(sessionID)
|
yield* result.get(sessionID)
|
||||||
if ((yield* execution.active).has(sessionID)) return yield* new BusyError({ sessionID })
|
return yield* Effect.die("SessionV2.revert.clear moved to SessionRuntime.Service")
|
||||||
return yield* SessionRevert.clear(session).pipe(
|
|
||||||
Effect.provideService(EventV2.Service, events),
|
|
||||||
Effect.provide(locations.get(session.location)),
|
|
||||||
)
|
|
||||||
}),
|
}),
|
||||||
commit: Effect.fn("V2Session.revert.commit")(function* (sessionID) {
|
commit: Effect.fn("V2Session.revert.commit")(function* (sessionID) {
|
||||||
const session = yield* result.get(sessionID)
|
const session = yield* result.get(sessionID)
|
||||||
if ((yield* execution.active).has(sessionID)) return yield* new BusyError({ sessionID })
|
|
||||||
return yield* SessionRevert.commit(session).pipe(Effect.provideService(EventV2.Service, events))
|
return yield* SessionRevert.commit(session).pipe(Effect.provideService(EventV2.Service, events))
|
||||||
}),
|
}),
|
||||||
},
|
},
|
||||||
|
|
@ -510,9 +494,7 @@ export const node = makeGlobalNode({
|
||||||
Database.node,
|
Database.node,
|
||||||
EventV2.node,
|
EventV2.node,
|
||||||
ProjectV2.node,
|
ProjectV2.node,
|
||||||
SessionExecution.node,
|
|
||||||
SessionStore.node,
|
SessionStore.node,
|
||||||
LocationServiceMap.node,
|
|
||||||
SessionProjector.node,
|
SessionProjector.node,
|
||||||
],
|
],
|
||||||
})
|
})
|
||||||
|
|
|
||||||
157
packages/core/src/session/runtime.ts
Normal file
157
packages/core/src/session/runtime.ts
Normal file
|
|
@ -0,0 +1,157 @@
|
||||||
|
export * as SessionRuntime from "./runtime"
|
||||||
|
|
||||||
|
import { Context, Effect, Layer } from "effect"
|
||||||
|
import { Database } from "../database/database"
|
||||||
|
import { EventV2 } from "../event"
|
||||||
|
import { Location } from "../location"
|
||||||
|
import { PromptInput } from "@opencode-ai/schema/prompt-input"
|
||||||
|
import { SessionMessage } from "./message"
|
||||||
|
import { Prompt } from "./prompt"
|
||||||
|
import { SessionInput } from "./input"
|
||||||
|
import { SessionRevert } from "./revert"
|
||||||
|
import { SessionRunner } from "./runner"
|
||||||
|
import * as SessionRunnerLLM from "./runner/llm"
|
||||||
|
import { SessionRunCoordinator } from "./run-coordinator"
|
||||||
|
import { SessionSchema } from "./schema"
|
||||||
|
import { Snapshot } from "../snapshot"
|
||||||
|
import { FSUtil } from "../fs-util"
|
||||||
|
import { makeLocationNode } from "../effect/app-node"
|
||||||
|
import { SessionV2 } from "../session"
|
||||||
|
import {
|
||||||
|
BusyError,
|
||||||
|
MessageNotFoundError,
|
||||||
|
NotFoundError,
|
||||||
|
PromptConflictError,
|
||||||
|
type RevertState,
|
||||||
|
} from "../session"
|
||||||
|
|
||||||
|
export interface Interface {
|
||||||
|
readonly prompt: (input: {
|
||||||
|
id?: SessionMessage.ID
|
||||||
|
sessionID: SessionSchema.ID
|
||||||
|
prompt: PromptInput.Prompt
|
||||||
|
delivery?: SessionInput.Delivery
|
||||||
|
resume?: boolean
|
||||||
|
}) => Effect.Effect<SessionInput.Admitted, NotFoundError | PromptConflictError>
|
||||||
|
readonly wait: (id: SessionSchema.ID) => Effect.Effect<void, NotFoundError>
|
||||||
|
readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>>
|
||||||
|
readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | SessionRunner.RunError>
|
||||||
|
readonly interrupt: (sessionID: SessionSchema.ID) => Effect.Effect<void>
|
||||||
|
readonly revert: {
|
||||||
|
readonly stage: (input: {
|
||||||
|
sessionID: SessionSchema.ID
|
||||||
|
messageID: SessionMessage.ID
|
||||||
|
files?: boolean
|
||||||
|
}) => Effect.Effect<RevertState, NotFoundError | MessageNotFoundError | BusyError | Snapshot.Error>
|
||||||
|
readonly clear: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | BusyError | Snapshot.Error>
|
||||||
|
readonly commit: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | BusyError>
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export class Service extends Context.Service<Service, Interface>()("@opencode/v2/SessionRuntime") {}
|
||||||
|
|
||||||
|
export const layer = Layer.effect(
|
||||||
|
Service,
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const database = yield* Database.Service
|
||||||
|
const db = database.db
|
||||||
|
const events = yield* EventV2.Service
|
||||||
|
const location = yield* Location.Service
|
||||||
|
const sessions = yield* SessionV2.Service
|
||||||
|
const runner = yield* SessionRunner.Service
|
||||||
|
|
||||||
|
const local = Effect.fn("SessionRuntime.local")(function* (sessionID: SessionSchema.ID) {
|
||||||
|
const session = yield* sessions.get(sessionID)
|
||||||
|
if (session.location.directory !== location.directory || session.location.workspaceID !== location.workspaceID)
|
||||||
|
return yield* new NotFoundError({ sessionID })
|
||||||
|
return session
|
||||||
|
})
|
||||||
|
|
||||||
|
const coordinator = yield* SessionRunCoordinator.make<SessionSchema.ID, SessionRunner.RunError>({
|
||||||
|
drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, force) {
|
||||||
|
yield* local(sessionID).pipe(Effect.orDie)
|
||||||
|
return yield* runner.run({ sessionID, force })
|
||||||
|
}),
|
||||||
|
})
|
||||||
|
|
||||||
|
return Service.of({
|
||||||
|
prompt: Effect.fn("SessionRuntime.prompt")((input) =>
|
||||||
|
Effect.uninterruptible(
|
||||||
|
Effect.gen(function* () {
|
||||||
|
yield* local(input.sessionID)
|
||||||
|
const prompt = resolvePrompt(input.prompt)
|
||||||
|
const messageID = input.id ?? SessionMessage.ID.create()
|
||||||
|
const delivery = input.delivery ?? "steer"
|
||||||
|
const expected = { sessionID: input.sessionID, messageID, prompt, delivery }
|
||||||
|
const admitted = yield* SessionInput.admit(db, events, {
|
||||||
|
id: messageID,
|
||||||
|
sessionID: input.sessionID,
|
||||||
|
prompt,
|
||||||
|
delivery,
|
||||||
|
}).pipe(
|
||||||
|
Effect.catchDefect((defect) =>
|
||||||
|
defect instanceof SessionInput.LifecycleConflict
|
||||||
|
? new PromptConflictError({ sessionID: input.sessionID, messageID })
|
||||||
|
: Effect.die(defect),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
if (!SessionInput.equivalent(admitted, expected))
|
||||||
|
return yield* new PromptConflictError({ sessionID: input.sessionID, messageID })
|
||||||
|
if (input.resume !== false) yield* coordinator.wake(admitted.sessionID)
|
||||||
|
return admitted
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
wait: Effect.fn("SessionRuntime.wait")(function* (sessionID) {
|
||||||
|
yield* local(sessionID)
|
||||||
|
yield* coordinator.awaitIdle(sessionID)
|
||||||
|
}),
|
||||||
|
active: coordinator.active,
|
||||||
|
resume: Effect.fn("SessionRuntime.resume")(function* (sessionID) {
|
||||||
|
yield* local(sessionID)
|
||||||
|
yield* coordinator.run(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)
|
||||||
|
if ((yield* coordinator.active).has(input.sessionID)) return yield* new BusyError({ sessionID: input.sessionID })
|
||||||
|
return yield* SessionRevert.stage({ session, messageID: input.messageID, files: input.files }).pipe(
|
||||||
|
Effect.provideService(Database.Service, database),
|
||||||
|
Effect.provideService(EventV2.Service, events),
|
||||||
|
)
|
||||||
|
}),
|
||||||
|
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))
|
||||||
|
}),
|
||||||
|
commit: Effect.fn("SessionRuntime.revert.commit")(function* (sessionID) {
|
||||||
|
const session = yield* local(sessionID)
|
||||||
|
if ((yield* coordinator.active).has(sessionID)) return yield* new BusyError({ sessionID })
|
||||||
|
return yield* SessionRevert.commit(session).pipe(Effect.provideService(EventV2.Service, events))
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
const resolvePrompt = (input: PromptInput.Prompt) =>
|
||||||
|
Prompt.make({
|
||||||
|
text: input.text,
|
||||||
|
agents: input.agents,
|
||||||
|
files: input.files?.map((file) => {
|
||||||
|
const dataMime = file.uri.match(/^data:([^;,]+)[;,]/i)?.[1]
|
||||||
|
const target = URL.canParse(file.uri) ? new URL(file.uri).pathname : (file.name ?? file.uri)
|
||||||
|
return {
|
||||||
|
...file,
|
||||||
|
mime: dataMime ?? (target.endsWith("/") ? "application/x-directory" : FSUtil.mimeType(target)),
|
||||||
|
}
|
||||||
|
}),
|
||||||
|
})
|
||||||
|
|
||||||
|
export const node = makeLocationNode({
|
||||||
|
service: Service,
|
||||||
|
layer,
|
||||||
|
deps: [Database.node, EventV2.node, Location.node, SessionV2.node, SessionRunnerLLM.node, Snapshot.node],
|
||||||
|
})
|
||||||
|
|
@ -6,6 +6,7 @@ import type { CommandHooks } from "./command.js"
|
||||||
import type { IntegrationHooks } from "./integration.js"
|
import type { IntegrationHooks } from "./integration.js"
|
||||||
import type { PluginDomain } from "./plugin.js"
|
import type { PluginDomain } from "./plugin.js"
|
||||||
import type { ReferenceHooks } from "./reference.js"
|
import type { ReferenceHooks } from "./reference.js"
|
||||||
|
import type { SessionDomain } from "./session.js"
|
||||||
import type { SkillHooks } from "./skill.js"
|
import type { SkillHooks } from "./skill.js"
|
||||||
import type { Reload } from "./registration.js"
|
import type { Reload } from "./registration.js"
|
||||||
|
|
||||||
|
|
@ -18,5 +19,6 @@ export interface PluginContext {
|
||||||
readonly integration: IntegrationHooks & Reload
|
readonly integration: IntegrationHooks & Reload
|
||||||
readonly plugin: PluginDomain
|
readonly plugin: PluginDomain
|
||||||
readonly reference: ReferenceHooks & Reload
|
readonly reference: ReferenceHooks & Reload
|
||||||
|
readonly session: SessionDomain
|
||||||
readonly skill: SkillHooks & Reload
|
readonly skill: SkillHooks & Reload
|
||||||
}
|
}
|
||||||
|
|
|
||||||
30
packages/plugin/src/v2/effect/session.ts
Normal file
30
packages/plugin/src/v2/effect/session.ts
Normal file
|
|
@ -0,0 +1,30 @@
|
||||||
|
import type { Effect } from "effect"
|
||||||
|
import type { PromptInput, SessionInputAdmitted, SessionMessage, SessionV2Info } from "@opencode-ai/sdk/v2/types"
|
||||||
|
|
||||||
|
export interface SessionDomain {
|
||||||
|
readonly create: (input: {
|
||||||
|
readonly id?: string
|
||||||
|
readonly parentID?: string
|
||||||
|
readonly title?: string
|
||||||
|
readonly agent?: string
|
||||||
|
readonly model?: SessionV2Info["model"]
|
||||||
|
}) => Effect.Effect<SessionV2Info>
|
||||||
|
readonly get: (sessionID: string) => Effect.Effect<SessionV2Info>
|
||||||
|
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>>
|
||||||
|
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>
|
||||||
|
}
|
||||||
|
|
@ -6,6 +6,7 @@ import type { CommandHooks } from "./command.js"
|
||||||
import type { IntegrationHooks } from "./integration.js"
|
import type { IntegrationHooks } from "./integration.js"
|
||||||
import type { PluginDomain } from "./plugin.js"
|
import type { PluginDomain } from "./plugin.js"
|
||||||
import type { ReferenceHooks } from "./reference.js"
|
import type { ReferenceHooks } from "./reference.js"
|
||||||
|
import type { SessionDomain } from "./session.js"
|
||||||
import type { SkillHooks } from "./skill.js"
|
import type { SkillHooks } from "./skill.js"
|
||||||
import type { Reload } from "./registration.js"
|
import type { Reload } from "./registration.js"
|
||||||
|
|
||||||
|
|
@ -18,5 +19,6 @@ export interface PluginContext {
|
||||||
readonly integration: IntegrationHooks & Reload
|
readonly integration: IntegrationHooks & Reload
|
||||||
readonly plugin: PluginDomain
|
readonly plugin: PluginDomain
|
||||||
readonly reference: ReferenceHooks & Reload
|
readonly reference: ReferenceHooks & Reload
|
||||||
|
readonly session: SessionDomain
|
||||||
readonly skill: SkillHooks & Reload
|
readonly skill: SkillHooks & Reload
|
||||||
}
|
}
|
||||||
|
|
|
||||||
29
packages/plugin/src/v2/promise/session.ts
Normal file
29
packages/plugin/src/v2/promise/session.ts
Normal file
|
|
@ -0,0 +1,29 @@
|
||||||
|
import type { PromptInput, SessionInputAdmitted, SessionMessage, SessionV2Info } from "@opencode-ai/sdk/v2/types"
|
||||||
|
|
||||||
|
export interface SessionDomain {
|
||||||
|
readonly create: (input: {
|
||||||
|
readonly id?: string
|
||||||
|
readonly parentID?: string
|
||||||
|
readonly title?: string
|
||||||
|
readonly agent?: string
|
||||||
|
readonly model?: SessionV2Info["model"]
|
||||||
|
}) => Promise<SessionV2Info>
|
||||||
|
readonly get: (sessionID: string) => Promise<SessionV2Info>
|
||||||
|
readonly messages: (input: {
|
||||||
|
readonly sessionID: string
|
||||||
|
readonly limit?: number
|
||||||
|
readonly order?: "asc" | "desc"
|
||||||
|
readonly cursor?: { readonly id: string; readonly direction: "previous" | "next" }
|
||||||
|
}) => Promise<ReadonlyArray<SessionMessage>>
|
||||||
|
readonly context: (sessionID: string) => Promise<ReadonlyArray<SessionMessage>>
|
||||||
|
readonly prompt: (input: {
|
||||||
|
readonly id?: string
|
||||||
|
readonly sessionID: string
|
||||||
|
readonly prompt: PromptInput
|
||||||
|
readonly delivery?: "steer" | "queue"
|
||||||
|
readonly resume?: boolean
|
||||||
|
}) => Promise<SessionInputAdmitted>
|
||||||
|
readonly resume: (sessionID: string) => Promise<void>
|
||||||
|
readonly wait: (sessionID: string) => Promise<void>
|
||||||
|
readonly interrupt: (sessionID: string) => Promise<void>
|
||||||
|
}
|
||||||
|
|
@ -1,4 +1,6 @@
|
||||||
import { SessionV2 } from "@opencode-ai/core/session"
|
import { SessionV2 } from "@opencode-ai/core/session"
|
||||||
|
import { LocationServiceMap } from "@opencode-ai/core/location-service-map"
|
||||||
|
import { SessionRuntime } from "@opencode-ai/core/session/runtime"
|
||||||
import { DateTime, Effect, Stream } from "effect"
|
import { DateTime, Effect, Stream } from "effect"
|
||||||
import { HttpApiBuilder, HttpApiSchema } from "effect/unstable/httpapi"
|
import { HttpApiBuilder, HttpApiSchema } from "effect/unstable/httpapi"
|
||||||
import { Api } from "../api"
|
import { Api } from "../api"
|
||||||
|
|
@ -20,6 +22,14 @@ const DefaultSessionHistoryLimit = 50
|
||||||
export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handlers) =>
|
export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handlers) =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const session = yield* SessionV2.Service
|
const session = yield* SessionV2.Service
|
||||||
|
const locations = yield* LocationServiceMap.Service
|
||||||
|
const route = Effect.fn("SessionHandler.route")(function* <A, E>(
|
||||||
|
sessionID: SessionV2.ID,
|
||||||
|
effect: Effect.Effect<A, E, SessionRuntime.Service>,
|
||||||
|
) {
|
||||||
|
const info = yield* session.get(sessionID)
|
||||||
|
return yield* effect.pipe(Effect.provide(locations.get(info.location)))
|
||||||
|
})
|
||||||
|
|
||||||
return handlers
|
return handlers
|
||||||
.handle(
|
.handle(
|
||||||
|
|
@ -159,15 +169,18 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||||
"session.prompt",
|
"session.prompt",
|
||||||
Effect.fn(function* (ctx) {
|
Effect.fn(function* (ctx) {
|
||||||
return {
|
return {
|
||||||
data: yield* session
|
data: yield* route(
|
||||||
.prompt({
|
ctx.params.sessionID,
|
||||||
sessionID: ctx.params.sessionID,
|
SessionRuntime.Service.use((runtime) =>
|
||||||
id: ctx.payload.id,
|
runtime.prompt({
|
||||||
prompt: ctx.payload.prompt,
|
sessionID: ctx.params.sessionID,
|
||||||
delivery: ctx.payload.delivery,
|
id: ctx.payload.id,
|
||||||
resume: ctx.payload.resume,
|
prompt: ctx.payload.prompt,
|
||||||
})
|
delivery: ctx.payload.delivery,
|
||||||
.pipe(
|
resume: ctx.payload.resume,
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
).pipe(
|
||||||
Effect.catchTag("Session.NotFoundError", (error) =>
|
Effect.catchTag("Session.NotFoundError", (error) =>
|
||||||
Effect.fail(
|
Effect.fail(
|
||||||
new SessionNotFoundError({
|
new SessionNotFoundError({
|
||||||
|
|
@ -215,7 +228,10 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||||
.handle(
|
.handle(
|
||||||
"session.wait",
|
"session.wait",
|
||||||
Effect.fn(function* (ctx) {
|
Effect.fn(function* (ctx) {
|
||||||
yield* session.wait(ctx.params.sessionID).pipe(
|
yield* route(
|
||||||
|
ctx.params.sessionID,
|
||||||
|
SessionRuntime.Service.use((runtime) => runtime.wait(ctx.params.sessionID)),
|
||||||
|
).pipe(
|
||||||
Effect.catchTag("Session.NotFoundError", (error) =>
|
Effect.catchTag("Session.NotFoundError", (error) =>
|
||||||
Effect.fail(
|
Effect.fail(
|
||||||
new SessionNotFoundError({
|
new SessionNotFoundError({
|
||||||
|
|
@ -237,7 +253,10 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||||
files: ctx.payload.files,
|
files: ctx.payload.files,
|
||||||
})
|
})
|
||||||
return {
|
return {
|
||||||
data: yield* session.revert.stage({ ...ctx.params, ...ctx.payload }).pipe(
|
data: yield* route(
|
||||||
|
ctx.params.sessionID,
|
||||||
|
SessionRuntime.Service.use((runtime) => runtime.revert.stage({ ...ctx.params, ...ctx.payload })),
|
||||||
|
).pipe(
|
||||||
Effect.catchTag(
|
Effect.catchTag(
|
||||||
"Session.NotFoundError",
|
"Session.NotFoundError",
|
||||||
(error) =>
|
(error) =>
|
||||||
|
|
@ -284,7 +303,10 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||||
"session.revert.clear",
|
"session.revert.clear",
|
||||||
Effect.fn(function* (ctx) {
|
Effect.fn(function* (ctx) {
|
||||||
yield* Effect.log("session.revert.clear", { sessionID: ctx.params.sessionID })
|
yield* Effect.log("session.revert.clear", { sessionID: ctx.params.sessionID })
|
||||||
yield* session.revert.clear(ctx.params.sessionID).pipe(
|
yield* route(
|
||||||
|
ctx.params.sessionID,
|
||||||
|
SessionRuntime.Service.use((runtime) => runtime.revert.clear(ctx.params.sessionID)),
|
||||||
|
).pipe(
|
||||||
Effect.catchTag(
|
Effect.catchTag(
|
||||||
"Session.NotFoundError",
|
"Session.NotFoundError",
|
||||||
(error) =>
|
(error) =>
|
||||||
|
|
@ -322,7 +344,10 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||||
"session.revert.commit",
|
"session.revert.commit",
|
||||||
Effect.fn(function* (ctx) {
|
Effect.fn(function* (ctx) {
|
||||||
yield* Effect.log("session.revert.commit", { sessionID: ctx.params.sessionID })
|
yield* Effect.log("session.revert.commit", { sessionID: ctx.params.sessionID })
|
||||||
yield* session.revert.commit(ctx.params.sessionID).pipe(
|
yield* route(
|
||||||
|
ctx.params.sessionID,
|
||||||
|
SessionRuntime.Service.use((runtime) => runtime.revert.commit(ctx.params.sessionID)),
|
||||||
|
).pipe(
|
||||||
Effect.catchTag(
|
Effect.catchTag(
|
||||||
"Session.NotFoundError",
|
"Session.NotFoundError",
|
||||||
(error) =>
|
(error) =>
|
||||||
|
|
@ -407,7 +432,10 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||||
.handle(
|
.handle(
|
||||||
"session.interrupt",
|
"session.interrupt",
|
||||||
Effect.fn(function* (ctx) {
|
Effect.fn(function* (ctx) {
|
||||||
yield* session.interrupt(ctx.params.sessionID)
|
yield* route(
|
||||||
|
ctx.params.sessionID,
|
||||||
|
SessionRuntime.Service.use((runtime) => runtime.interrupt(ctx.params.sessionID)),
|
||||||
|
)
|
||||||
return HttpApiSchema.NoContent.make()
|
return HttpApiSchema.NoContent.make()
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -107,10 +107,10 @@ The SDK should be implemented as one plugin instance: it receives a plugin host/
|
||||||
|
|
||||||
## Concrete implementation slices
|
## Concrete implementation slices
|
||||||
|
|
||||||
1. Add `packages/core/src/session/runtime.ts` as a location node.
|
1. Add `packages/core/src/session/runtime.ts` as a location node. **Implemented in this draft.**
|
||||||
2. Move `prompt`, `resume`, `wait`, `interrupt`, `active`, and location-sensitive `revert` operations from `SessionV2.Service` into the runtime service.
|
2. Move `prompt`, `resume`, `wait`, `interrupt`, `active`, and location-sensitive `revert` operations from `SessionV2.Service` into the runtime service. **Implemented in this draft for the new runtime path; old `SessionV2` entrypoints are left as compatibility stubs and should be removed once callers migrate.**
|
||||||
3. Update server route handlers to route location-sensitive requests at the API boundary by resolving the session location and providing that location runtime.
|
3. Update server route handlers to route location-sensitive requests at the API boundary by resolving the session location and providing that location runtime. **Implemented in this draft.**
|
||||||
4. Add `ctx.session` to `PluginHost` by composing `SessionV2.Service` and the location session runtime.
|
4. Add `ctx.session` to `PluginHost` by composing `SessionV2.Service` and the location session runtime. **Implemented in this draft.**
|
||||||
5. Add public plugin `ctx.tool.transform` types and adapt it to the existing canonical core `Tool.make` representation.
|
5. Add public plugin `ctx.tool.transform` types and adapt it to the existing canonical core `Tool.make` representation.
|
||||||
6. Convert `ToolRegistry` registration to transform/rebuild semantics.
|
6. Convert `ToolRegistry` registration to transform/rebuild semantics.
|
||||||
7. Port `subagent` to a built-in plugin that registers a normal location tool.
|
7. Port `subagent` to a built-in plugin that registers a normal location tool.
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue