Compare commits

...
7 changed files with 514 additions and 703 deletions

No files matched your search

+335 -46
View File
@@ -1,7 +1,7 @@
export * as Session from "./session.js"
export * from "./session/schema.js"
import { Effect, Layer, Schema, Context, Stream } from "effect"
import { DateTime, Effect, Fiber, Layer, Schema, Scope, Context, Stream } from "effect"
import { LLMClient } from "@opencode/ai"
import { ListAnchor } from "@opencode/schema/session"
import { and, desc, eq } from "drizzle-orm"
@@ -44,7 +44,6 @@ import { SessionEvent } from "./session/event.js"
import { SessionInbox } from "./session/inbox.js"
import { InstructionState } from "./session/instruction-state.js"
import { SessionGenerate } from "./session/generate.js"
import { SessionCommand } from "./session/command.js"
import {
SessionMove,
DestinationNotFoundError,
@@ -54,16 +53,22 @@ import {
import { SessionModelTransport } from "./session/model-transport.js"
import { llmClient } from "./effect/app-node-platform.js"
import { Snapshot } from "./snapshot.js"
import { Session } from "./session/session.js"
import { SessionDiff, TurnRangeError } from "./session/diff.js"
import { LocationServiceMap } from "./location-service-map.js"
import { FSUtil } from "@opencode/util/fs-util"
import type { EventLog } from "@opencode/schema/event-log"
import type { FileDiff } from "@opencode/schema/file-diff"
import { Job } from "./job.js"
import type { Command } from "./command.js"
import { Command } from "./command.js"
import { SessionEnvironment } from "./session/environment.js"
import { InstructionEntry } from "./session/instruction-entry.js"
import { SessionPrompt } from "./session/prompt.js"
import { SessionRevert } from "./session/revert.js"
import { Plugin } from "./plugin/service.js"
import { Shell } from "./shell.js"
import { ShellResult } from "./shell/result.js"
import { Skill } from "./skill.js"
import { Event } from "@opencode/schema/event"
// get project -> project.locations
//
@@ -90,8 +95,6 @@ type CreateBaseInput = {
type CreateInput = CreateBaseInput &
({ location: Location.Ref; parentID?: never } | { parentID: SessionSchema.ID; location?: never })
type CompactInput = Parameters<Session.Handle["compact"]>[0] & { sessionID: SessionSchema.ID }
type ForkInput = {
sessionID: SessionSchema.ID
before?: SessionMessage.ID
@@ -180,8 +183,8 @@ export interface Interface {
}) => Effect.Effect<void, NotFoundError>
readonly move: SessionMove.Interface["move"]
readonly prompt: (
input: Parameters<Session.Handle["prompt"]>[0] & { sessionID: SessionSchema.ID },
) => ReturnType<Session.Handle["prompt"]>
input: SessionPrompt.Input & { sessionID: SessionSchema.ID; id?: SessionMessage.ID; resume?: boolean },
) => Effect.Effect<SessionInbox.User, NotFoundError | PromptConflictError | AttachmentError | SkillNotFoundError>
/** Generates text from current Session context without admitting input or mutating history. */
readonly generate: (input: {
sessionID: SessionSchema.ID
@@ -196,23 +199,36 @@ export interface Interface {
skills?: PromptInput.Prompt["skills"]
delivery?: SessionInbox.Delivery
}) => Effect.Effect<void, NotFoundError | Command.NotFoundError | Command.ExecutionError>
readonly shell: (
input: Parameters<Session.Handle["shell"]>[0] & { sessionID: SessionSchema.ID },
) => ReturnType<Session.Handle["shell"]>
readonly skill: (
input: Parameters<Session.Handle["skill"]>[0] & { sessionID: SessionSchema.ID },
) => ReturnType<Session.Handle["skill"]>
readonly compact: (
input: CompactInput,
) => Effect.Effect<SessionInbox.Compaction, NotFoundError | CompactionConflictError>
readonly shell: (input: {
sessionID: SessionSchema.ID
id?: SessionMessage.ID
command: string
}) => Effect.Effect<void, NotFoundError>
readonly skill: (input: {
sessionID: SessionSchema.ID
messageID?: SessionMessage.ID
skill: Skill.ID
resume?: boolean
}) => Effect.Effect<void, NotFoundError | SkillNotFoundError>
readonly compact: (input: {
sessionID: SessionSchema.ID
id?: SessionMessage.ID
delivery?: SessionInbox.Delivery
}) => Effect.Effect<SessionInbox.Compaction, NotFoundError | CompactionConflictError>
readonly wait: (id: SessionSchema.ID) => Effect.Effect<void, NotFoundError>
readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>>
readonly background: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError>
readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | SessionRunner.RunError>
readonly interrupt: (sessionID: SessionSchema.ID, options?: { readonly resume?: boolean }) => Effect.Effect<boolean>
readonly synthetic: (
input: Parameters<Session.Handle["synthetic"]>[0] & { sessionID: SessionSchema.ID },
) => ReturnType<Session.Handle["synthetic"]>
readonly synthetic: (input: {
sessionID: SessionSchema.ID
id?: SessionMessage.ID
text: string
description?: string
metadata?: Record<string, unknown>
delivery?: SessionInbox.Delivery
resume?: boolean
}) => Effect.Effect<SessionInbox.Synthetic, NotFoundError | SyntheticConflictError>
readonly revert: {
readonly stage: (input: {
sessionID: SessionSchema.ID
@@ -243,9 +259,27 @@ const layer = Layer.effect(
const jobs = yield* Job.Service
const environments = yield* SessionEnvironment.Service
const locations = yield* LocationServiceMap.Service
const sessions = yield* Session.make()
const admission = yield* SessionInbox.Service
const fs = yield* FSUtil.Service
const scope = yield* Scope.Scope
const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
const mutatePending = (
input: InboxItemRef,
mutation: (input: {
readonly id: SessionMessage.ID
readonly sessionID: SessionSchema.ID
}) => Effect.Effect<void, SessionInbox.LifecycleConflict>,
) =>
mutation({ sessionID: input.sessionID, id: input.inboxID }).pipe(
Effect.catchTag("SessionInbox.LifecycleConflict", () =>
Effect.gen(function* () {
yield* result.get(input.sessionID)
return yield* new InboxConflictError({ sessionID: input.sessionID, inboxID: input.inboxID })
}),
),
)
const result = Service.of({
create: Effect.fn("Session.create")(function* (input) {
const sessionID = input.id ?? SessionSchema.ID.create()
@@ -345,13 +379,26 @@ const layer = Layer.effect(
})
return yield* result.get(sessionID).pipe(Effect.orDie)
}),
get: (sessionID) => sessions.forSession(sessionID).get(),
get: Effect.fn("Session.get")(function* (sessionID) {
const session = yield* store.get(sessionID)
if (!session) return yield* new NotFoundError({ sessionID })
return session
}),
environment: Effect.fn("Session.environment")(function* (input) {
yield* result.get(input.sessionID)
if (input.variables !== undefined) yield* environments.set(input.sessionID, input.variables)
return yield* environments.get(input.sessionID)
}),
view: (input) => sessions.forSession(input.sessionID).view(input),
view: Effect.fn("Session.view")(function* (input) {
const session = yield* result.get(input.sessionID)
if (
session.time.idle === undefined ||
input.idle > DateTime.toEpochMillis(session.time.idle) ||
(session.time.viewed !== undefined && DateTime.toEpochMillis(session.time.viewed) >= input.idle)
)
return
yield* bus.publish(SessionEvent.Viewed, { sessionID: input.sessionID, idle: input.idle })
}),
remove: Effect.fn("Session.remove")(function* (sessionID) {
yield* result.get(sessionID)
yield* execution.interrupt(sessionID)
@@ -370,7 +417,10 @@ const layer = Layer.effect(
yield* result.get(input.sessionID)
return yield* store.messages(input)
}),
message: (input) => sessions.forSession(input.sessionID).message(input.messageID),
message: Effect.fn("Session.message")(function* (input) {
const stored = yield* store.message(input.messageID)
return stored?.sessionID === input.sessionID ? stored.message : undefined
}),
context: Effect.fn("Session.context")(function* (sessionID) {
yield* result.get(sessionID)
return yield* store.context(sessionID)
@@ -386,10 +436,22 @@ const layer = Layer.effect(
context: input.context,
})
}),
inbox: (sessionID) => sessions.forSession(sessionID).inbox(),
cancelInbox: (input) => sessions.forSession(input.sessionID).cancelInbox(input.inboxID),
steerInbox: (input) => sessions.forSession(input.sessionID).steerInbox(input.inboxID),
queueInbox: (input) => sessions.forSession(input.sessionID).queueInbox(input.inboxID),
inbox: Effect.fn("Session.inbox")(function* (sessionID) {
yield* result.get(sessionID)
return yield* admission.list(sessionID)
}),
cancelInbox: Effect.fn("Session.cancelInbox")(
(input) => mutatePending(input, admission.cancel),
Effect.uninterruptible,
),
steerInbox: Effect.fn("Session.steerInbox")(function* (input) {
yield* mutatePending(input, admission.steer)
yield* execution.wake(input.sessionID)
}, Effect.uninterruptible),
queueInbox: Effect.fn("Session.queueInbox")(
(input) => mutatePending(input, admission.queue),
Effect.uninterruptible,
),
log: (input) =>
Stream.unwrap(
result
@@ -401,7 +463,43 @@ const layer = Layer.effect(
Bus.isSynced(item) || isDurableSessionEvent(item),
),
),
prompt: (input) => sessions.forSession(input.sessionID).prompt(input),
prompt: Effect.fn("Session.prompt")((input) =>
Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const session = yield* result.get(input.sessionID)
const messageID = input.id ?? SessionMessage.ID.create()
const admitted = yield* Effect.gen(function* () {
const existing = yield* admission.reconcile({
id: messageID,
sessionID: session.id,
type: "user",
delivery: input.delivery ?? "steer",
})
if (existing) return existing
const item = yield* restore(
SessionPrompt.prepare({ session, messageID, input }).pipe(
Effect.provideService(Instance.Service, instances),
Effect.provideService(FSUtil.Service, fs),
),
)
// Commit a staged revert only after preparation succeeds, before admitting new work.
if (session.revert) yield* SessionRevert.commit(bus, session)
return yield* admission.admit({
id: messageID,
sessionID: session.id,
item,
})
}).pipe(
Effect.catchTag(
"SessionInbox.LifecycleConflict",
() => new PromptConflictError({ sessionID: input.sessionID, messageID }),
),
)
if (input.resume !== false) yield* execution.wake(input.sessionID)
return admitted
}),
),
),
generate: Effect.fn("Session.generate")(function* (input) {
const session = yield* result.get(input.sessionID)
return yield* SessionGenerate.generate({ session, prompt: input.prompt }).pipe(
@@ -412,20 +510,155 @@ const layer = Layer.effect(
}),
command: Effect.fn("Session.command")(function* (input) {
const session = yield* result.get(input.sessionID)
return yield* SessionCommand.execute({ ...input, session }).pipe(
Effect.provideService(Instance.Service, instances),
const commands = yield* Plugin.awaitActivation.pipe(Effect.andThen(Command.Service), instances.provide(session))
yield* commands.execute({
name: input.command,
invocation: {
sessionID: session.id,
prompt: {
text: input.text,
files: input.files,
agents: input.agents,
skills: input.skills,
},
delivery: input.delivery ?? "steer",
},
})
}),
shell: Effect.fn("Session.shell")(function* (input) {
const session = yield* result.get(input.sessionID)
// The server owns completion recording even if the submitting client disconnects.
const running = yield* Effect.gen(function* () {
const started = yield* Effect.gen(function* () {
const shells = yield* Plugin.awaitActivation.pipe(Effect.andThen(Shell.Service), instances.provide(session))
const info = yield* shells.create({
command: input.command,
cwd: session.location.directory,
timeout: 0,
metadata: { sessionID: session.id, background: true },
})
return { shells, info }
}).pipe(
Effect.tapError((error) =>
result.synthetic({
sessionID: input.sessionID,
text: `User shell command failed to start:\n${input.command}\n\n${error.message}`,
description: input.command,
metadata: { source: "shell", state: "error" },
resume: false,
}),
),
Effect.orDie,
)
yield* bus.publish(
SessionEvent.Shell.Started,
{
sessionID: input.sessionID,
shell: started.info,
},
{ id: input.id ? Event.ID.make(input.id.replace(/^msg_/, "evt_")) : undefined },
)
// Keep completion tied to the original shell even if the Session moves.
const terminal = yield* started.shells.result(started.info)
const preview = yield* started.shells
.output(started.info.id, { limit: 1024 * 1024 })
.pipe(Effect.catchTag("Shell.NotFoundError", () => Effect.succeed(ShellResult.unavailable)))
yield* bus.publish(SessionEvent.Shell.Ended, {
sessionID: input.sessionID,
shell: terminal.info,
output: preview,
})
yield* result
.synthetic({
sessionID: input.sessionID,
...ShellResult.userNotification(terminal),
resume: false,
})
.pipe(
Effect.catchTag("Session.NotFoundError", () => Effect.void),
Effect.orDie,
)
}).pipe(Effect.forkIn(scope, { startImmediately: true }))
yield* Fiber.join(running)
}),
skill: Effect.fn("Session.skill")(function* (input) {
const session = yield* result.get(input.sessionID)
const skills = yield* Plugin.awaitActivation.pipe(Effect.andThen(Skill.Service), instances.provide(session))
const skill = yield* skills.get(input.skill)
if (!skill) return yield* new SkillNotFoundError({ skill: input.skill })
yield* bus.publish(
SessionEvent.Skill.Activated,
{
sessionID: input.sessionID,
id: skill.id,
name: skill.name,
text: skill.content,
},
{ id: input.messageID ? Event.ID.make(input.messageID.replace(/^msg_/, "evt_")) : undefined },
)
if (input.resume !== false)
yield* execution
.resume(input.sessionID)
.pipe(Effect.ignore, Effect.forkIn(scope, { startImmediately: true }), Effect.asVoid)
}),
switchAgent: Effect.fn("Session.switchAgent")(function* (input) {
const session = yield* result.get(input.sessionID)
yield* bus.publish(SessionEvent.AgentSelected, {
sessionID: input.sessionID,
agent: input.agent,
previous: session.agent,
})
}),
switchModel: Effect.fn("Session.switchModel")(function* (input) {
const session = yield* result.get(input.sessionID)
if (
session.model?.providerID === input.model.providerID &&
session.model.id === input.model.id &&
(session.model.variant ?? "default") === (input.model.variant ?? "default")
)
return
yield* bus.publish(SessionEvent.ModelSelected, {
sessionID: input.sessionID,
model: input.model,
previous: session.model,
})
}),
rename: Effect.fn("Session.rename")(function* (input) {
yield* result.get(input.sessionID)
yield* bus.publish(SessionEvent.Renamed, { sessionID: input.sessionID, title: input.title })
}),
setMetadata: Effect.fn("Session.setMetadata")(function* (input) {
yield* result.get(input.sessionID)
yield* bus.publish(SessionEvent.MetadataUpdated, { sessionID: input.sessionID, metadata: input.metadata })
}),
setPermissions: Effect.fn("Session.setPermissions")(function* (input) {
yield* result.get(input.sessionID)
yield* bus.publish(SessionEvent.Permissions, { sessionID: input.sessionID, permissions: input.permissions })
}),
shell: (input) => sessions.forSession(input.sessionID).shell(input),
skill: (input) => sessions.forSession(input.sessionID).skill(input),
switchAgent: (input) => sessions.forSession(input.sessionID).switchAgent(input),
switchModel: (input) => sessions.forSession(input.sessionID).switchModel(input),
rename: (input) => sessions.forSession(input.sessionID).rename(input),
setMetadata: (input) => sessions.forSession(input.sessionID).setMetadata(input),
setPermissions: (input) => sessions.forSession(input.sessionID).setPermissions(input),
move: moves.move,
compact: (input) => sessions.forSession(input.sessionID).compact(input),
wait: (sessionID) => sessions.forSession(sessionID).wait(),
compact: Effect.fn("Session.compact")(function* (input) {
const session = yield* result.get(input.sessionID)
if (session.revert) yield* SessionRevert.commit(bus, session)
const inputID = input.id ?? SessionMessage.ID.create()
const admitted = yield* admission
.admitCompaction({
id: inputID,
sessionID: input.sessionID,
delivery: input.delivery ?? "steer",
})
.pipe(
Effect.catchTag(
"SessionInbox.LifecycleConflict",
() => new CompactionConflictError({ sessionID: input.sessionID, inputID }),
),
)
yield* execution.wake(input.sessionID)
return admitted
}),
wait: Effect.fn("Session.wait")(function* (sessionID) {
yield* result.get(sessionID)
yield* execution.awaitIdle(sessionID)
}),
active: execution.active,
background: Effect.fn("Session.background")(function* (sessionID) {
yield* result.get(sessionID)
@@ -445,13 +678,69 @@ const layer = Layer.effect(
})
.pipe(Effect.catchTag("Session.SyntheticConflictError", Effect.die))
}),
resume: (sessionID) => sessions.forSession(sessionID).resume(),
synthetic: (input) => sessions.forSession(input.sessionID).synthetic(input),
interrupt: (sessionID, options) => sessions.forSession(sessionID).interrupt(options),
resume: Effect.fn("Session.resume")(function* (sessionID) {
yield* result.get(sessionID)
yield* execution.resume(sessionID)
}),
synthetic: Effect.fn("Session.synthetic")((input) =>
Effect.uninterruptible(
Effect.gen(function* () {
yield* result.get(input.sessionID)
const inputID = input.id ?? SessionMessage.ID.create()
const item = {
type: "synthetic",
payload: SessionInbox.SyntheticPayload.make({
text: input.text,
description: input.description,
metadata: input.metadata,
}),
delivery: SessionInbox.Delivery.make(input.delivery ?? "steer"),
} satisfies SessionInbox.Item
const admitted = yield* admission
.admit({
id: inputID,
sessionID: input.sessionID,
item,
})
.pipe(
Effect.catchTag(
"SessionInbox.LifecycleConflict",
() => new SyntheticConflictError({ sessionID: input.sessionID, inputID }),
),
)
if (input.resume !== false && !(yield* result.get(input.sessionID)).revert)
yield* execution.wake(input.sessionID)
return admitted
}),
),
),
interrupt: Effect.fn("Session.interrupt")((sessionID, options) =>
Effect.uninterruptible(execution.interrupt(sessionID, options)),
),
revert: {
stage: (input) => sessions.forSession(input.sessionID).revert.stage(input),
clear: (sessionID) => sessions.forSession(sessionID).revert.clear(),
commit: (sessionID) => sessions.forSession(sessionID).revert.commit(),
stage: Effect.fn("Session.revert.stage")(function* (input) {
const session = yield* result.get(input.sessionID)
if (yield* execution.isActive(input.sessionID)) return yield* new BusyError({ sessionID: input.sessionID })
return yield* SessionRevert.stage({ session, messageID: input.messageID, files: input.files }).pipe(
Effect.provideService(Instance.Service, instances),
Effect.provideService(Database.Service, database),
Effect.provideService(Bus.Service, bus),
)
}),
clear: Effect.fn("Session.revert.clear")(function* (sessionID) {
const session = yield* result.get(sessionID)
if (yield* execution.isActive(sessionID)) return yield* new BusyError({ sessionID })
yield* SessionRevert.clear(session).pipe(
Effect.provideService(Instance.Service, instances),
Effect.provideService(Bus.Service, bus),
)
return yield* execution.wake(sessionID)
}),
commit: Effect.fn("Session.revert.commit")(function* (sessionID) {
const session = yield* result.get(sessionID)
if (yield* execution.isActive(sessionID)) return yield* new BusyError({ sessionID })
return yield* SessionRevert.commit(bus, session)
}),
},
})
-35
View File
@@ -1,35 +0,0 @@
export * as SessionCommand from "./command.js"
import type { PromptInput } from "@opencode/schema/prompt-input"
import type { Session } from "@opencode/schema/session"
import type { SessionInbox } from "@opencode/schema/session-inbox"
import { Effect } from "effect"
import { Command } from "../command.js"
import { Instance } from "../instance/service.js"
import { Plugin } from "../plugin/service.js"
export const execute = Effect.fn("SessionCommand.execute")(function* (input: {
session: Session.Info
command: string
text: string
files?: PromptInput.Prompt["files"]
agents?: PromptInput.Prompt["agents"]
skills?: PromptInput.Prompt["skills"]
delivery?: SessionInbox.Delivery
}) {
const instances = yield* Instance.Service
const commands = yield* Plugin.awaitActivation.pipe(Effect.andThen(Command.Service), instances.provide(input.session))
yield* commands.execute({
name: input.command,
invocation: {
sessionID: input.session.id,
prompt: {
text: input.text,
files: input.files,
agents: input.agents,
skills: input.skills,
},
delivery: input.delivery ?? "steer",
},
})
})
-427
View File
@@ -1,427 +0,0 @@
export * as Session from "./session.js"
import { DateTime, Effect, Fiber, Scope } from "effect"
import type { Agent } from "@opencode/schema/agent"
import type { Model } from "@opencode/schema/model"
import type { Permission } from "@opencode/schema/permission"
import { Event } from "@opencode/schema/event"
import { FSUtil } from "@opencode/util/fs-util"
import { Bus } from "../bus.js"
import { Database } from "../database/database.js"
import { Instance } from "../instance/service.js"
import { ShellResult } from "../shell/result.js"
import type { Skill } from "../skill.js"
import {
BusyError,
CompactionConflictError,
InboxConflictError,
MessageNotFoundError,
NotFoundError,
PromptConflictError,
SyntheticConflictError,
} from "./error.js"
import { SessionEvent } from "./event.js"
import { SessionExecution } from "./execution.js"
import { SessionInbox } from "./inbox.js"
import { SessionMessage } from "./message.js"
import { SessionPrompt } from "./prompt.js"
import { SessionRevert } from "./revert.js"
import { SessionShell } from "./shell.js"
import { SessionSkill } from "./skill.js"
import { SessionSchema } from "./schema.js"
import { SessionStore } from "./store.js"
type PromptRequest = SessionPrompt.Input & {
id?: SessionMessage.ID
resume?: boolean
}
/**
* Build once in the host Scope: `const sessions = yield* Session.make()`.
* Use `sessions.forSession(id)` for handles that share host services and reload current state.
*/
export const make = Effect.fn("Session.make")(function* () {
const bus = yield* Bus.Service
const database = yield* Database.Service
const store = yield* SessionStore.Service
const instances = yield* Instance.Service
const execution = yield* SessionExecution.Service
const admission = yield* SessionInbox.Service
const fs = yield* FSUtil.Service
const scope = yield* Scope.Scope
const get = Effect.fn("Session.get")(function* (sessionID: SessionSchema.ID) {
const session = yield* store.get(sessionID)
if (!session) return yield* new NotFoundError({ sessionID })
return session
})
const message = Effect.fn("Session.message")(function* (sessionID: SessionSchema.ID, messageID: SessionMessage.ID) {
const stored = yield* store.message(messageID)
return stored?.sessionID === sessionID ? stored.message : undefined
})
const view = Effect.fn("Session.view")(function* (sessionID: SessionSchema.ID, input: { idle: number }) {
const session = yield* get(sessionID)
if (
session.time.idle === undefined ||
input.idle > DateTime.toEpochMillis(session.time.idle) ||
(session.time.viewed !== undefined && DateTime.toEpochMillis(session.time.viewed) >= input.idle)
)
return
yield* bus.publish(SessionEvent.Viewed, { sessionID, idle: input.idle })
})
const rename = Effect.fn("Session.rename")(function* (sessionID: SessionSchema.ID, input: { title: string }) {
yield* get(sessionID)
yield* bus.publish(SessionEvent.Renamed, { sessionID, title: input.title })
})
const setMetadata = Effect.fn("Session.setMetadata")(function* (
sessionID: SessionSchema.ID,
input: { metadata: SessionSchema.Metadata },
) {
yield* get(sessionID)
yield* bus.publish(SessionEvent.MetadataUpdated, { sessionID, metadata: input.metadata })
})
const setPermissions = Effect.fn("Session.setPermissions")(function* (
sessionID: SessionSchema.ID,
input: { permissions: Permission.Ruleset },
) {
yield* get(sessionID)
yield* bus.publish(SessionEvent.Permissions, { sessionID, permissions: input.permissions })
})
const switchAgent = Effect.fn("Session.switchAgent")(function* (
sessionID: SessionSchema.ID,
input: { agent: Agent.ID },
) {
const session = yield* get(sessionID)
yield* bus.publish(SessionEvent.AgentSelected, { sessionID, agent: input.agent, previous: session.agent })
})
const switchModel = Effect.fn("Session.switchModel")(function* (
sessionID: SessionSchema.ID,
input: { model: Model.Ref },
) {
const session = yield* get(sessionID)
if (
session.model?.providerID === input.model.providerID &&
session.model.id === input.model.id &&
(session.model.variant ?? "default") === (input.model.variant ?? "default")
)
return
yield* bus.publish(SessionEvent.ModelSelected, { sessionID, model: input.model, previous: session.model })
})
const mutatePending = (
sessionID: SessionSchema.ID,
inboxID: SessionMessage.ID,
mutation: (input: {
readonly id: SessionMessage.ID
readonly sessionID: SessionSchema.ID
}) => Effect.Effect<void, SessionInbox.LifecycleConflict>,
) =>
mutation({ sessionID, id: inboxID }).pipe(
Effect.catchTag("SessionInbox.LifecycleConflict", () =>
Effect.gen(function* () {
yield* get(sessionID)
return yield* new InboxConflictError({ sessionID, inboxID })
}),
),
)
const inbox = Effect.fn("Session.inbox")(function* (sessionID: SessionSchema.ID) {
yield* get(sessionID)
return yield* admission.list(sessionID)
})
const cancelInbox = Effect.fn("Session.cancelInbox")(
(sessionID: SessionSchema.ID, inboxID: SessionMessage.ID) => mutatePending(sessionID, inboxID, admission.cancel),
Effect.uninterruptible,
)
const steerInbox = Effect.fn("Session.steerInbox")(function* (
sessionID: SessionSchema.ID,
inboxID: SessionMessage.ID,
) {
yield* mutatePending(sessionID, inboxID, admission.steer)
yield* execution.wake(sessionID)
}, Effect.uninterruptible)
const queueInbox = Effect.fn("Session.queueInbox")(
(sessionID: SessionSchema.ID, inboxID: SessionMessage.ID) => mutatePending(sessionID, inboxID, admission.queue),
Effect.uninterruptible,
)
const prompt = Effect.fn("Session.prompt")((sessionID: SessionSchema.ID, input: PromptRequest) =>
Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const session = yield* get(sessionID)
const messageID = input.id ?? SessionMessage.ID.create()
const admitted = yield* Effect.gen(function* () {
const existing = yield* admission.reconcile({
id: messageID,
sessionID: session.id,
type: "user",
delivery: input.delivery ?? "steer",
})
if (existing) return existing
const item = yield* restore(
SessionPrompt.prepare({ session, messageID, input }).pipe(
Effect.provideService(Instance.Service, instances),
Effect.provideService(FSUtil.Service, fs),
),
)
// Commit a staged revert only after preparation succeeds, before admitting new work.
if (session.revert) yield* SessionRevert.commit(bus, session)
return yield* admission.admit({
id: messageID,
sessionID: session.id,
item,
})
}).pipe(
Effect.catchTag("SessionInbox.LifecycleConflict", () => new PromptConflictError({ sessionID, messageID })),
)
if (input.resume !== false) yield* execution.wake(sessionID)
return admitted
}),
),
)
const shell = Effect.fn("Session.shell")(function* (
sessionID: SessionSchema.ID,
input: { id?: SessionMessage.ID; command: string },
) {
const session = yield* get(sessionID)
// The server owns completion recording even if the submitting client disconnects.
const running = yield* Effect.gen(function* () {
const started = yield* SessionShell.start({ session, command: input.command }).pipe(
Effect.provideService(Instance.Service, instances),
Effect.tapError((error) =>
synthetic(sessionID, {
text: `User shell command failed to start:\n${input.command}\n\n${error.message}`,
description: input.command,
metadata: { source: "shell", state: "error" },
resume: false,
}),
),
Effect.orDie,
)
yield* bus.publish(
SessionEvent.Shell.Started,
{
sessionID,
shell: started.info,
},
{ id: input.id ? Event.ID.make(input.id.replace(/^msg_/, "evt_")) : undefined },
)
const terminal = yield* started.result
const preview = yield* started.output
yield* bus.publish(SessionEvent.Shell.Ended, {
sessionID,
shell: terminal.info,
output: preview,
})
yield* synthetic(sessionID, {
...ShellResult.userNotification(terminal),
resume: false,
}).pipe(
Effect.catchTag("Session.NotFoundError", () => Effect.void),
Effect.orDie,
)
}).pipe(Effect.forkIn(scope, { startImmediately: true }))
yield* Fiber.join(running)
})
const skill = Effect.fn("Session.skill")(function* (
sessionID: SessionSchema.ID,
input: { messageID?: SessionMessage.ID; skill: Skill.ID; resume?: boolean },
) {
const session = yield* get(sessionID)
const skill = yield* SessionSkill.get({ session, skill: input.skill }).pipe(
Effect.provideService(Instance.Service, instances),
)
yield* bus.publish(
SessionEvent.Skill.Activated,
{
sessionID,
id: skill.id,
name: skill.name,
text: skill.content,
},
{ id: input.messageID ? Event.ID.make(input.messageID.replace(/^msg_/, "evt_")) : undefined },
)
if (input.resume !== false)
yield* execution
.resume(sessionID)
.pipe(Effect.ignore, Effect.forkIn(scope, { startImmediately: true }), Effect.asVoid)
})
const compact = Effect.fn("Session.compact")(function* (
sessionID: SessionSchema.ID,
input: { id?: SessionMessage.ID; delivery?: SessionInbox.Delivery },
) {
const session = yield* get(sessionID)
if (session.revert) yield* SessionRevert.commit(bus, session)
const inputID = input.id ?? SessionMessage.ID.create()
const admitted = yield* admission
.admitCompaction({
id: inputID,
sessionID,
delivery: input.delivery ?? "steer",
})
.pipe(
Effect.catchTag("SessionInbox.LifecycleConflict", () => new CompactionConflictError({ sessionID, inputID })),
)
yield* execution.wake(sessionID)
return admitted
})
const wait = Effect.fn("Session.wait")(function* (sessionID: SessionSchema.ID) {
yield* get(sessionID)
yield* execution.awaitIdle(sessionID)
})
const resume = Effect.fn("Session.resume")(function* (sessionID: SessionSchema.ID) {
yield* get(sessionID)
yield* execution.resume(sessionID)
})
const synthetic = Effect.fn("Session.synthetic")(
(
sessionID: SessionSchema.ID,
input: {
id?: SessionMessage.ID
text: string
description?: string
metadata?: Record<string, unknown>
delivery?: SessionInbox.Delivery
resume?: boolean
},
) =>
Effect.uninterruptible(
Effect.gen(function* () {
yield* get(sessionID)
const inputID = input.id ?? SessionMessage.ID.create()
const admittedInput = {
type: "synthetic",
payload: SessionInbox.SyntheticPayload.make({
text: input.text,
description: input.description,
metadata: input.metadata,
}),
delivery: SessionInbox.Delivery.make(input.delivery ?? "steer"),
} satisfies SessionInbox.Item
const admitted = yield* admission
.admit({
id: inputID,
sessionID,
item: admittedInput,
})
.pipe(
Effect.catchTag(
"SessionInbox.LifecycleConflict",
() => new SyntheticConflictError({ sessionID, inputID }),
),
)
if (input.resume !== false && !(yield* get(sessionID)).revert) yield* execution.wake(sessionID)
return admitted
}),
),
)
const interrupt = Effect.fn("Session.interrupt")(
(sessionID: SessionSchema.ID, options?: { readonly resume?: boolean }) =>
Effect.uninterruptible(execution.interrupt(sessionID, options)),
)
const stage = Effect.fn("Session.revert.stage")(function* (
sessionID: SessionSchema.ID,
input: { messageID: SessionMessage.ID; files?: boolean },
) {
const session = yield* get(sessionID)
if (yield* execution.isActive(sessionID)) return yield* new BusyError({ sessionID })
return yield* SessionRevert.stage({ session, messageID: input.messageID, files: input.files }).pipe(
Effect.provideService(Instance.Service, instances),
Effect.provideService(Database.Service, database),
Effect.provideService(Bus.Service, bus),
)
})
const clear = Effect.fn("Session.revert.clear")(function* (sessionID: SessionSchema.ID) {
const session = yield* get(sessionID)
if (yield* execution.isActive(sessionID)) return yield* new BusyError({ sessionID })
yield* SessionRevert.clear(session).pipe(
Effect.provideService(Instance.Service, instances),
Effect.provideService(Bus.Service, bus),
)
return yield* execution.wake(sessionID)
})
const commit = Effect.fn("Session.revert.commit")(function* (sessionID: SessionSchema.ID) {
const session = yield* get(sessionID)
if (yield* execution.isActive(sessionID)) return yield* new BusyError({ sessionID })
return yield* SessionRevert.commit(bus, session)
})
const revert = { stage, clear, commit }
const operations = {
get,
message,
view,
rename,
setMetadata,
setPermissions,
switchAgent,
switchModel,
inbox,
prompt,
synthetic,
shell,
skill,
compact,
wait,
resume,
interrupt,
cancelInbox,
steerInbox,
queueInbox,
revert,
}
const forSession = (sessionID: SessionSchema.ID) => {
const get = operations.get.bind(undefined, sessionID)
const message = operations.message.bind(undefined, sessionID)
const view = operations.view.bind(undefined, sessionID)
const rename = operations.rename.bind(undefined, sessionID)
const setMetadata = operations.setMetadata.bind(undefined, sessionID)
const setPermissions = operations.setPermissions.bind(undefined, sessionID)
const switchAgent = operations.switchAgent.bind(undefined, sessionID)
const switchModel = operations.switchModel.bind(undefined, sessionID)
const inbox = operations.inbox.bind(undefined, sessionID)
const prompt = operations.prompt.bind(undefined, sessionID)
const synthetic = operations.synthetic.bind(undefined, sessionID)
const shell = operations.shell.bind(undefined, sessionID)
const skill = operations.skill.bind(undefined, sessionID)
const compact = operations.compact.bind(undefined, sessionID)
const wait = operations.wait.bind(undefined, sessionID)
const resume = operations.resume.bind(undefined, sessionID)
const interrupt = operations.interrupt.bind(undefined, sessionID)
const cancelInbox = operations.cancelInbox.bind(undefined, sessionID)
const steerInbox = operations.steerInbox.bind(undefined, sessionID)
const queueInbox = operations.queueInbox.bind(undefined, sessionID)
const stage = operations.revert.stage.bind(undefined, sessionID)
const clear = operations.revert.clear.bind(undefined, sessionID)
const commit = operations.revert.commit.bind(undefined, sessionID)
const revert = { stage, clear, commit }
return {
id: sessionID,
get,
message,
view,
rename,
setMetadata,
setPermissions,
switchAgent,
switchModel,
inbox,
prompt,
synthetic,
shell,
skill,
compact,
wait,
resume,
interrupt,
cancelInbox,
steerInbox,
queueInbox,
revert,
}
}
return { forSession }
})
export type Handle = ReturnType<Effect.Success<ReturnType<typeof make>>["forSession"]>
// Mirrors the shell tool's in-memory preview safety limit.
-27
View File
@@ -1,27 +0,0 @@
export * as SessionShell from "./shell.js"
import type { Session } from "@opencode/schema/session"
import { Effect } from "effect"
import { Instance } from "../instance/service.js"
import { Plugin } from "../plugin/service.js"
import { Shell } from "../shell.js"
import { ShellResult } from "../shell/result.js"
export const start = Effect.fn("SessionShell.start")(function* (input: { session: Session.Info; command: string }) {
const instances = yield* Instance.Service
const shell = yield* Plugin.awaitActivation.pipe(Effect.andThen(Shell.Service), instances.provide(input.session))
const info = yield* shell.create({
command: input.command,
cwd: input.session.location.directory,
timeout: 0,
metadata: { sessionID: input.session.id, background: true },
})
// Keep completion tied to the original shell even if the Session moves.
return {
info,
result: shell.result(info),
output: shell
.output(info.id, { limit: 1024 * 1024 })
.pipe(Effect.catchTag("Shell.NotFoundError", () => Effect.succeed(ShellResult.unavailable))),
}
})
-16
View File
@@ -1,16 +0,0 @@
export * as SessionSkill from "./skill.js"
import type { Session } from "@opencode/schema/session"
import { Effect } from "effect"
import { Instance } from "../instance/service.js"
import { Plugin } from "../plugin/service.js"
import { Skill } from "../skill.js"
import { SkillNotFoundError } from "./error.js"
export const get = Effect.fn("SessionSkill.get")(function* (input: { session: Session.Info; skill: Skill.ID }) {
const instances = yield* Instance.Service
const skills = yield* Plugin.awaitActivation.pipe(Effect.andThen(Skill.Service), instances.provide(input.session))
const skill = yield* skills.get(input.skill)
if (!skill) return yield* new SkillNotFoundError({ skill: input.skill })
return skill
})
+171 -141
View File
@@ -1,11 +1,11 @@
import { describe, expect } from "bun:test"
import { and, eq } from "drizzle-orm"
import { Cause, Context, DateTime, Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect"
import { LLMClient } from "@opencode/ai"
import { Agent } from "@opencode/schema/agent"
import { Event } from "@opencode/schema/event"
import { Model } from "@opencode/schema/model"
import { Money } from "@opencode/schema/money"
import { Project } from "@opencode/schema/project"
import { Provider } from "@opencode/schema/provider"
import { ID, Info, Output } from "@opencode/schema/shell"
import { LayerNode } from "@opencode/util/effect/layer-node"
@@ -13,14 +13,19 @@ import { FSUtil } from "@opencode/util/fs-util"
import { Global } from "@opencode/util/global"
import { Bus } from "../src/bus.js"
import { Database } from "../src/database/database.js"
import { AppNodeBuilder } from "../src/effect/app-node-builder.js"
import { llmClient } from "../src/effect/app-node-platform.js"
import { EventTable } from "../src/event/sql.js"
import { Image } from "../src/image.js"
import { Instance } from "../src/instance/service.js"
import { Job } from "../src/job.js"
import { Location } from "../src/location.js"
import { Plugin } from "../src/plugin.js"
import { PluginHooks } from "../src/plugin/hooks.js"
import { Project } from "../src/project.js"
import { ProjectTable } from "../src/project/sql.js"
import { AbsolutePath, RelativePath } from "../src/schema.js"
import { Session } from "../src/session.js"
import { InboxConflictError, NotFoundError, PromptConflictError } from "../src/session/error.js"
import { SessionEvent } from "../src/session/event.js"
import { SessionExecution } from "../src/session/execution.js"
@@ -31,7 +36,8 @@ import { SessionProjector } from "../src/session/projector.js"
import { SessionRevert } from "../src/session/revert.js"
import { SessionRunCoordinator } from "../src/session/run-coordinator.js"
import { SessionSchema } from "../src/session/schema.js"
import { Session } from "../src/session/session.js"
import { SessionEnvironment } from "../src/session/environment.js"
import { SessionMove } from "../src/session/move.js"
import { SessionTable } from "../src/session/sql.js"
import { SessionStore } from "../src/session/store.js"
import { Shell } from "../src/shell.js"
@@ -148,28 +154,33 @@ const setup = Effect.fnUntraced(function* (options?: {
// This fixture supplies only the instance services exercised by Session.
provide: (session) => Effect.provide(servicesFor(session.location) as Layer.Layer<Instance.Services>),
})
const sessions = yield* Session.make().pipe(
Effect.satisfiesServicesType<
| Bus.Service
| Database.Service
| FSUtil.Service
| SessionStore.Service
| Instance.Service
| SessionExecution.Service
| SessionInbox.Service
| Scope.Scope
>(),
Effect.provideService(Instance.Service, instances),
Effect.provideService(SessionExecution.Service, options?.execution ?? execution),
const context = yield* Layer.build(
AppNodeBuilder.build(Session.node, [
Database.node.replace(Layer.succeed(Database.Service, database)),
Bus.node.replace(Layer.succeed(Bus.Service, bus)),
SessionStore.node.replace(Layer.succeed(SessionStore.Service, store)),
SessionInbox.node.replace(Layer.succeed(SessionInbox.Service, yield* SessionInbox.Service)),
FSUtil.node.replace(Layer.succeed(FSUtil.Service, yield* FSUtil.Service)),
SessionProjector.node.replace(Layer.empty),
Instance.node.replace(Layer.succeed(Instance.Service, instances)),
SessionExecution.node.replace(Layer.succeed(SessionExecution.Service, options?.execution ?? execution)),
// Session uses these only for collection, move, and host-routing operations this suite does not exercise.
Project.node.replace(Layer.mock(Project.Service, {})),
Job.node.replace(Layer.mock(Job.Service, {})),
SessionEnvironment.node.replace(Layer.mock(SessionEnvironment.Service, {})),
SessionMove.node.replace(Layer.mock(SessionMove.Service, {})),
llmClient.replace(Layer.mock(LLMClient.Service, {})),
]),
)
const sessions = Context.get(context, Session.Service)
return { sessions, instances, hooks, locations, activationWaits, resumes, wakes, db: database.db, bus, store }
})
describe("Session-owned handles", () => {
describe("Session-owned operations", () => {
it.live("owns state changes and message reads without caller services or Location acquisition", () =>
Effect.gen(function* () {
const fixture = yield* setup()
const handle = fixture.sessions.forSession(sessionID)
const sessions = fixture.sessions
const model = { id: Model.ID.make("test-model"), providerID: Provider.ID.make("test-provider") }
const messageID = SessionMessage.ID.create()
yield* fixture.bus.publish(SessionEvent.Step.Started, {
@@ -192,23 +203,21 @@ describe("Session-owned handles", () => {
.where(eq(SessionTable.id, sessionID))
.run()
.pipe(Effect.orDie)
const { rename, switchAgent, switchModel, view, message } = handle
yield* Effect.gen(function* () {
yield* rename({ title: "Renamed" })
yield* switchAgent({ agent: Agent.ID.make("review") })
yield* switchModel({ model })
yield* switchModel({ model })
yield* view({ idle: 0 })
yield* view({ idle: 0 })
expect(yield* message(messageID)).toMatchObject({ type: "assistant", content: [] })
yield* sessions.rename({ sessionID, title: "Renamed" })
yield* sessions.switchAgent({ sessionID, agent: Agent.ID.make("review") })
yield* sessions.switchModel({ sessionID, model })
yield* sessions.switchModel({ sessionID, model })
yield* sessions.view({ sessionID, idle: 0 })
yield* sessions.view({ sessionID, idle: 0 })
expect(yield* sessions.message({ sessionID, messageID })).toMatchObject({ type: "assistant", content: [] })
}).pipe(Effect.satisfiesServicesType<never>(), Effect.setContext(Context.empty()))
const session = yield* handle.get()
const session = yield* sessions.get(sessionID)
expect(session).toMatchObject({ title: "Renamed", agent: "review", model })
expect(session.time.viewed && DateTime.toEpochMillis(session.time.viewed)).toBe(0)
expect(yield* fixture.sessions.forSession(otherID).message(messageID)).toBeUndefined()
expect((yield* fixture.sessions.forSession(otherID).get()).title).toBe("Owned session")
expect(yield* sessions.message({ sessionID: otherID, messageID })).toBeUndefined()
expect((yield* sessions.get(otherID)).title).toBe("Owned session")
expect(fixture.locations).toEqual([])
expect(fixture.wakes).toEqual([])
const events = yield* fixture.db
@@ -227,11 +236,9 @@ describe("Session-owned handles", () => {
it.live("acquires Location only for new prompt preparation and persists before waking", () =>
Effect.gen(function* () {
const fixture = yield* setup()
const handle = fixture.sessions.forSession(sessionID)
const { get, prompt } = handle
expect(handle.id).toBe(sessionID)
expect((yield* get().pipe(Effect.satisfiesServicesType<never>())).location).toEqual(source)
const synthetic = yield* handle.synthetic({ text: "Background result", resume: false })
const sessions = fixture.sessions
expect((yield* sessions.get(sessionID).pipe(Effect.satisfiesServicesType<never>())).location).toEqual(source)
const synthetic = yield* sessions.synthetic({ sessionID, text: "Background result", resume: false })
expect(fixture.locations).toEqual([])
expect(fixture.wakes).toEqual([])
@@ -243,12 +250,14 @@ describe("Session-owned handles", () => {
event.prompt.text += " prepared"
}),
)
const first = yield* prompt({
const first = yield* sessions.prompt({
sessionID,
id: SessionMessage.ID.make("msg_owned_prepared"),
text: "Original",
files: [{ uri: new URL("./session-owned.test.ts", import.meta.url).href }],
})
const retried = yield* fixture.sessions.forSession(sessionID).prompt({
const retried = yield* sessions.prompt({
sessionID,
id: first.id,
text: "Ignored retry",
files: [{ uri: "file:///missing-owned-retry" }],
@@ -273,26 +282,34 @@ describe("Session-owned handles", () => {
}),
)
it.live("keeps the first admission across handles, including delivered retries and identity conflicts", () =>
it.live("keeps the first admission, including delivered retries and identity conflicts", () =>
Effect.gen(function* () {
const fixture = yield* setup()
const first = fixture.sessions.forSession(sessionID)
const second = fixture.sessions.forSession(sessionID)
const other = fixture.sessions.forSession(otherID)
const prompt = yield* first.prompt({ text: "Keep this", metadata: { source: "first" }, resume: false })
const retry = { id: prompt.id, text: "Ignore this", metadata: { source: "retry" }, resume: false }
expect(yield* second.prompt({ ...retry, delivery: "queue" })).toEqual(prompt)
const conflict = yield* other.prompt(retry).pipe(Effect.flip)
const sessions = fixture.sessions
const prompt = yield* sessions.prompt({
sessionID,
text: "Keep this",
metadata: { source: "first" },
resume: false,
})
const retry = { sessionID, id: prompt.id, text: "Ignore this", metadata: { source: "retry" }, resume: false }
expect(yield* sessions.prompt({ ...retry, delivery: "queue" })).toEqual(prompt)
const conflict = yield* sessions.prompt({ ...retry, sessionID: otherID }).pipe(Effect.flip)
expect(conflict).toBeInstanceOf(PromptConflictError)
expect(conflict).toMatchObject({ _tag: "Session.PromptConflictError", sessionID: otherID, messageID: prompt.id })
expect(yield* second.synthetic(retry).pipe(Effect.flip)).toMatchObject({
expect(yield* sessions.synthetic(retry).pipe(Effect.flip)).toMatchObject({
_tag: "Session.SyntheticConflictError",
sessionID,
inputID: prompt.id,
})
const synthetic = yield* first.synthetic({ text: "Original completion", description: "Job", resume: false })
expect(yield* second.synthetic({ ...retry, id: synthetic.id })).toEqual(synthetic)
expect(yield* first.inbox()).toEqual([prompt, synthetic])
const synthetic = yield* sessions.synthetic({
sessionID,
text: "Original completion",
description: "Job",
resume: false,
})
expect(yield* sessions.synthetic({ ...retry, id: synthetic.id })).toEqual(synthetic)
expect(yield* sessions.inbox(sessionID)).toEqual([prompt, synthetic])
yield* SessionInbox.promote(fixture.db, fixture.bus, sessionID, "steer")
// Delivered identity must be recoverable from the message, without retained enqueue history.
@@ -306,21 +323,23 @@ describe("Session-owned handles", () => {
)
.run()
.pipe(Effect.orDie)
expect((yield* second.prompt({ ...retry, files: [{ uri: "file:///missing-owned-retry" }] })).payload).toEqual(
expect((yield* sessions.prompt({ ...retry, files: [{ uri: "file:///missing-owned-retry" }] })).payload).toEqual(
prompt.payload,
)
expect((yield* second.synthetic({ ...retry, id: synthetic.id })).payload).toEqual(synthetic.payload)
expect(yield* other.synthetic({ ...retry, id: synthetic.id }).pipe(Effect.flip)).toMatchObject({
expect((yield* sessions.synthetic({ ...retry, id: synthetic.id })).payload).toEqual(synthetic.payload)
expect(
yield* sessions.synthetic({ ...retry, sessionID: otherID, id: synthetic.id }).pipe(Effect.flip),
).toMatchObject({
_tag: "Session.SyntheticConflictError",
sessionID: otherID,
inputID: synthetic.id,
})
expect(yield* second.prompt({ ...retry, id: synthetic.id }).pipe(Effect.flip)).toMatchObject({
expect(yield* sessions.prompt({ ...retry, id: synthetic.id }).pipe(Effect.flip)).toMatchObject({
_tag: "Session.PromptConflictError",
sessionID,
messageID: synthetic.id,
})
expect(yield* second.inbox()).toEqual([])
expect(yield* sessions.inbox(sessionID)).toEqual([])
expect(yield* fixture.store.context(sessionID)).toMatchObject([
{ id: prompt.id, text: "Keep this", metadata: { source: "first" } },
{ id: synthetic.id, text: "Original completion", description: "Job" },
@@ -329,13 +348,12 @@ describe("Session-owned handles", () => {
}),
)
it.live("reads fresh placement through an existing handle after a projected move", () =>
it.live("reads fresh placement through effects constructed before a projected move", () =>
Effect.gen(function* () {
const fixture = yield* setup()
const handle = fixture.sessions.forSession(sessionID)
yield* handle.prompt({ text: "Before move", resume: false })
const get = handle.get()
const prompt = handle.prompt({ text: "After move", resume: false })
yield* fixture.sessions.prompt({ sessionID, text: "Before move", resume: false })
const get = fixture.sessions.get(sessionID)
const prompt = fixture.sessions.prompt({ sessionID, text: "After move", resume: false })
const destination = Location.Ref.make({ directory: AbsolutePath.make("/project/moved") })
yield* fixture.bus.publish(SessionEvent.Moved, {
@@ -350,18 +368,17 @@ describe("Session-owned handles", () => {
yield* prompt
expect(fixture.locations).toEqual([source, destination])
expect(fixture.activationWaits).toEqual([source, destination])
expect((yield* fixture.sessions.forSession(otherID).get()).location).toEqual(source)
expect((yield* fixture.sessions.get(otherID)).location).toEqual(source)
}),
)
it.live("activates skills through detached handles using fresh placement and ambient publication context", () =>
it.live("activates skills through detached effects using fresh placement and ambient publication context", () =>
Effect.gen(function* () {
const fixture = yield* setup({
skills: (ref) =>
Layer.mock(Skill.Service, { get: () => Effect.succeed({ ...skillInfo, content: ref.directory }) }),
})
const handle = fixture.sessions.forSession(sessionID)
const { skill } = handle
const sessions = fixture.sessions
const events: Event.Payload[] = []
yield* fixture.bus.listen((event) =>
Effect.sync(() => {
@@ -369,12 +386,11 @@ describe("Session-owned handles", () => {
}),
)
const initial = SessionMessage.ID.make("msg_owned_skill_initial")
yield* skill({ messageID: initial, skill: skillInfo.id, resume: false }).pipe(
Effect.satisfiesServicesType<never>(),
Effect.setContext(Context.empty()),
)
yield* sessions
.skill({ sessionID, messageID: initial, skill: skillInfo.id, resume: false })
.pipe(Effect.satisfiesServicesType<never>(), Effect.setContext(Context.empty()))
const moved = SessionMessage.ID.make("msg_owned_skill_moved")
const activation = skill({ messageID: moved, skill: skillInfo.id, resume: false })
const activation = sessions.skill({ sessionID, messageID: moved, skill: skillInfo.id, resume: false })
const destination = Location.Ref.make({ directory: AbsolutePath.make("/project/moved") })
yield* fixture.bus.publish(SessionEvent.Moved, {
sessionID,
@@ -384,13 +400,19 @@ describe("Session-owned handles", () => {
})
yield* activation.pipe(Effect.satisfiesServicesType<never>(), Effect.setContext(Context.empty()))
yield* skill({ skill: skillInfo.id, resume: false }).pipe(
Effect.provideService(Location.Service, location(source)),
)
yield* sessions
.skill({ sessionID, skill: skillInfo.id, resume: false })
.pipe(Effect.provideService(Location.Service, location(source)))
expect(fixture.locations).toEqual([source, destination, destination])
expect(yield* handle.message(initial)).toMatchObject({ type: "skill", text: source.directory })
expect(yield* handle.message(moved)).toMatchObject({ type: "skill", text: destination.directory })
expect(yield* sessions.message({ sessionID, messageID: initial })).toMatchObject({
type: "skill",
text: source.directory,
})
expect(yield* sessions.message({ sessionID, messageID: moved })).toMatchObject({
type: "skill",
text: destination.directory,
})
expect(
events.filter((event) => event.type === SessionEvent.Skill.Activated.type).map((event) => event.location),
).toEqual([undefined, undefined, source])
@@ -411,22 +433,21 @@ describe("Session-owned handles", () => {
)
const missingID = SessionSchema.ID.make("ses_missing_skill")
expect(
yield* fixture.sessions.forSession(missingID).skill({ skill: skillInfo.id }).pipe(Effect.flip),
yield* fixture.sessions.skill({ sessionID: missingID, skill: skillInfo.id }).pipe(Effect.flip),
).toMatchObject({ _tag: "Session.NotFoundError", sessionID: missingID })
expect(fixture.locations).toEqual([])
const handle = fixture.sessions.forSession(sessionID)
const before = yield* handle.get()
const before = yield* fixture.sessions.get(sessionID)
const missing = Skill.ID.make("missing")
expect(yield* handle.skill({ skill: missing }).pipe(Effect.flip)).toMatchObject({
expect(yield* fixture.sessions.skill({ sessionID, skill: missing }).pipe(Effect.flip)).toMatchObject({
_tag: "Session.SkillNotFoundError",
skill: missing,
})
expect(fixture.locations).toEqual([source])
expect(events).toEqual([])
expect(yield* handle.get()).toEqual(before)
expect(yield* handle.inbox()).toEqual([])
expect(yield* fixture.sessions.get(sessionID)).toEqual(before)
expect(yield* fixture.sessions.inbox(sessionID)).toEqual([])
expect(yield* fixture.store.context(sessionID)).toEqual([])
expect(fixture.activationWaits).toEqual([source])
expect(fixture.resumes).toEqual([])
@@ -466,16 +487,26 @@ describe("Session-owned handles", () => {
if (event.type === SessionEvent.Skill.Activated.type) calls.push(`published:${event.id}`)
}),
)
const { skill } = fixture.sessions.forSession(sessionID)
const sessions = fixture.sessions
yield* skill({ messageID: SessionMessage.ID.make("msg_skill_no_resume"), skill: skillInfo.id, resume: false })
yield* sessions.skill({
sessionID,
messageID: SessionMessage.ID.make("msg_skill_no_resume"),
skill: skillInfo.id,
resume: false,
})
expect(calls).toEqual(["published:evt_skill_no_resume"])
yield* Effect.forEach(
[
{ messageID: SessionMessage.ID.make("msg_skill_default_resume"), skill: skillInfo.id },
{ messageID: SessionMessage.ID.make("msg_skill_explicit_resume"), skill: skillInfo.id, resume: true },
{ sessionID, messageID: SessionMessage.ID.make("msg_skill_default_resume"), skill: skillInfo.id },
{
sessionID,
messageID: SessionMessage.ID.make("msg_skill_explicit_resume"),
skill: skillInfo.id,
resume: true,
},
],
(input) => skill(input).pipe(Effect.scoped, Effect.forkScoped, Effect.flatMap(Fiber.join)),
(input) => sessions.skill(input).pipe(Effect.scoped, Effect.forkScoped, Effect.flatMap(Fiber.join)),
)
expect(calls).toEqual([
@@ -488,11 +519,11 @@ describe("Session-owned handles", () => {
expect(stopped).toEqual([])
yield* Scope.close(host, Exit.void)
expect(stopped).toEqual([sessionID, sessionID])
expect(yield* fixture.sessions.forSession(sessionID).inbox()).toEqual([])
expect(yield* sessions.inbox(sessionID)).toEqual([])
}),
)
it.live("keeps prompt wakes independent of shell work across handles", () =>
it.live("keeps prompt wakes independent of shell work", () =>
Effect.gen(function* () {
const blocked = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
@@ -530,15 +561,14 @@ describe("Session-owned handles", () => {
}),
})
const shell = yield* fixture.sessions
.forSession(sessionID)
.shell({ id: SessionMessage.ID.make("msg_owned_shell"), command: started.command })
.shell({ sessionID, id: SessionMessage.ID.make("msg_owned_shell"), command: started.command })
.pipe(Effect.forkScoped)
yield* Deferred.await(blocked)
const admitted = yield* fixture.sessions.forSession(sessionID).prompt({ text: "Admit while the shell runs" })
const admitted = yield* fixture.sessions.prompt({ sessionID, text: "Admit while the shell runs" })
expect(yield* SessionInbox.find(fixture.db, admitted.id)).toEqual(admitted)
expect(fixture.wakes).toEqual([{ sessionID, pending: [admitted.id], enqueued: 1 }])
const other = yield* fixture.sessions.forSession(otherID).prompt({ text: "Independent Session" })
const other = yield* fixture.sessions.prompt({ sessionID: otherID, text: "Independent Session" })
expect(fixture.wakes).toEqual([
{ sessionID, pending: [admitted.id], enqueued: 1 },
{ sessionID: otherID, pending: [other.id], enqueued: 1 },
@@ -553,56 +583,55 @@ describe("Session-owned handles", () => {
expect(yield* fixture.store.context(sessionID)).toMatchObject([
{ type: "shell", shellID: started.id, status: "exited", output: { output: "owned" } },
])
expect(yield* fixture.sessions.forSession(sessionID).inbox()).toMatchObject([
expect(yield* fixture.sessions.inbox(sessionID)).toMatchObject([
{ id: admitted.id, type: "user" },
{ type: "synthetic", payload: { metadata: { source: "shell", shellID: started.id, state: "completed" } } },
])
}),
)
it.live("mutates only this handle's pending inbox and preserves public conflict tags", () =>
it.live("mutates only the addressed Session's pending inbox and preserves public conflict tags", () =>
Effect.gen(function* () {
const fixture = yield* setup()
const handle = fixture.sessions.forSession(sessionID)
const second = fixture.sessions.forSession(sessionID)
const queued = yield* handle.synthetic({ text: "Queued", delivery: "queue", resume: false })
const steer = yield* handle.prompt({ text: "Steer", resume: false })
const compact = yield* handle.compact({ delivery: "queue" })
const sessions = fixture.sessions
const queued = yield* sessions.synthetic({ sessionID, text: "Queued", delivery: "queue", resume: false })
const steer = yield* sessions.prompt({ sessionID, text: "Steer", resume: false })
const compact = yield* sessions.compact({ sessionID, delivery: "queue" })
yield* second.steerInbox(queued.id)
yield* second.queueInbox(steer.id)
expect(yield* handle.inbox()).toMatchObject([
yield* sessions.steerInbox({ sessionID, inboxID: queued.id })
yield* sessions.queueInbox({ sessionID, inboxID: steer.id })
expect(yield* sessions.inbox(sessionID)).toMatchObject([
{ id: queued.id, delivery: "steer" },
{ id: steer.id, delivery: "queue" },
{ id: compact.id, type: "compaction", delivery: "queue" },
])
expect(fixture.wakes).toHaveLength(2)
expect(yield* fixture.sessions.forSession(otherID).cancelInbox(queued.id).pipe(Effect.flip)).toMatchObject({
expect(yield* sessions.cancelInbox({ sessionID: otherID, inboxID: queued.id }).pipe(Effect.flip)).toMatchObject({
_tag: "Session.InboxConflictError",
sessionID: otherID,
inboxID: queued.id,
})
yield* second.cancelInbox(compact.id)
const cancelled = yield* handle.cancelInbox(compact.id).pipe(Effect.flip)
yield* sessions.cancelInbox({ sessionID, inboxID: compact.id })
const cancelled = yield* sessions.cancelInbox({ sessionID, inboxID: compact.id }).pipe(Effect.flip)
expect(cancelled).toBeInstanceOf(InboxConflictError)
expect(cancelled).toMatchObject({ _tag: "Session.InboxConflictError", sessionID, inboxID: compact.id })
expect(yield* handle.compact({ id: steer.id }).pipe(Effect.flip)).toMatchObject({
expect(yield* sessions.compact({ sessionID, id: steer.id }).pipe(Effect.flip)).toMatchObject({
_tag: "Session.CompactionConflictError",
sessionID,
inputID: steer.id,
})
expect(yield* SessionInbox.promote(fixture.db, fixture.bus, sessionID, "steer")).toBe(1)
expect(yield* second.queueInbox(queued.id).pipe(Effect.flip)).toMatchObject({
expect(yield* sessions.queueInbox({ sessionID, inboxID: queued.id }).pipe(Effect.flip)).toMatchObject({
_tag: "Session.InboxConflictError",
sessionID,
inboxID: queued.id,
})
expect(yield* handle.inbox()).toMatchObject([{ id: steer.id, delivery: "queue" }])
yield* second.cancelInbox(steer.id)
expect(yield* handle.inbox()).toEqual([])
expect(yield* sessions.inbox(sessionID)).toMatchObject([{ id: steer.id, delivery: "queue" }])
yield* sessions.cancelInbox({ sessionID, inboxID: steer.id })
expect(yield* sessions.inbox(sessionID)).toEqual([])
const missingID = SessionSchema.ID.make("ses_owned_missing")
const missing = yield* fixture.sessions.forSession(missingID).inbox().pipe(Effect.flip)
const missing = yield* sessions.inbox(missingID).pipe(Effect.flip)
expect(missing).toBeInstanceOf(NotFoundError)
expect(missing).toMatchObject({ _tag: "Session.NotFoundError", sessionID: missingID })
expect(fixture.locations).toEqual([source])
@@ -642,9 +671,9 @@ describe("Session-owned handles", () => {
),
}),
})
const first = yield* fixture.sessions.forSession(sessionID).resume().pipe(Effect.forkScoped)
const first = yield* fixture.sessions.resume(sessionID).pipe(Effect.forkScoped)
yield* Deferred.await(started)
const second = yield* fixture.sessions.forSession(sessionID).resume().pipe(Effect.forkScoped)
const second = yield* fixture.sessions.resume(sessionID).pipe(Effect.forkScoped)
yield* Deferred.await(joining)
yield* Fiber.interrupt(second)
@@ -654,11 +683,11 @@ describe("Session-owned handles", () => {
expect(drains).toEqual([sessionID])
yield* Deferred.succeed(release, undefined)
yield* Fiber.join(first)
yield* fixture.sessions.forSession(sessionID).wait()
yield* fixture.sessions.wait(sessionID)
expect(drains).toEqual([sessionID])
expect(yield* coordinator.active).toEqual(new Set())
expect(yield* fixture.sessions.forSession(sessionID).interrupt({ resume: true })).toBe(false)
expect(yield* fixture.sessions.forSession(sessionID).interrupt()).toBe(false)
expect(yield* fixture.sessions.interrupt(sessionID, { resume: true })).toBe(false)
expect(yield* fixture.sessions.interrupt(sessionID)).toBe(false)
expect(interrupts).toEqual([
{ sessionID, options: { resume: true } },
{ sessionID, options: undefined },
@@ -672,33 +701,35 @@ describe("Session-owned handles", () => {
const fixture = yield* setup({
snapshot: () => Layer.mock(Snapshot.Service, { capture: () => Effect.undefined }),
})
const handle = fixture.sessions.forSession(sessionID)
const boundary = yield* handle.synthetic({ text: "Revert boundary", resume: false })
const sessions = fixture.sessions
const boundary = yield* sessions.synthetic({ sessionID, text: "Revert boundary", resume: false })
yield* SessionInbox.promote(fixture.db, fixture.bus, sessionID, "steer")
yield* handle.revert.stage({ messageID: boundary.id, files: false })
yield* sessions.revert.stage({ sessionID, messageID: boundary.id, files: false })
const entered = yield* Deferred.make<void>()
const hook = yield* fixture.hooks.register("session", "prompt", () =>
Deferred.succeed(entered, undefined).pipe(Effect.andThen(Effect.never)),
)
const submission = yield* handle.prompt({ text: "Cancelled before admission" }).pipe(Effect.forkScoped)
const submission = yield* sessions
.prompt({ sessionID, text: "Cancelled before admission" })
.pipe(Effect.forkScoped)
yield* Deferred.await(entered)
yield* Fiber.interrupt(submission)
const cancelled = yield* Fiber.await(submission)
expect(Exit.isFailure(cancelled) && Cause.hasInterruptsOnly(cancelled.cause)).toBe(true)
expect(yield* handle.inbox()).toEqual([])
expect((yield* handle.get()).revert?.messageID).toBe(boundary.id)
expect(yield* sessions.inbox(sessionID)).toEqual([])
expect((yield* sessions.get(sessionID)).revert?.messageID).toBe(boundary.id)
expect(yield* fixture.store.context(sessionID)).toMatchObject([{ id: boundary.id }])
expect(fixture.wakes).toEqual([])
yield* hook.dispose
yield* handle.revert.clear()
expect((yield* handle.get()).revert).toBeUndefined()
yield* sessions.revert.clear(sessionID)
expect((yield* sessions.get(sessionID)).revert).toBeUndefined()
expect(yield* fixture.store.context(sessionID)).toMatchObject([{ id: boundary.id }])
yield* handle.revert.stage({ messageID: boundary.id, files: false })
yield* sessions.revert.stage({ sessionID, messageID: boundary.id, files: false })
const acquisitions = fixture.locations.length
yield* fixture.sessions.forSession(sessionID).revert.commit()
expect((yield* handle.get()).revert).toBeUndefined()
yield* sessions.revert.commit(sessionID)
expect((yield* sessions.get(sessionID)).revert).toBeUndefined()
expect(yield* fixture.store.context(sessionID)).toEqual([])
expect(fixture.locations).toHaveLength(acquisitions)
}),
@@ -717,10 +748,10 @@ describe("Session-owned handles", () => {
}),
}),
})
const handle = fixture.sessions.forSession(sessionID)
const boundary = yield* handle.synthetic({ text: "Revert boundary", resume: false })
const sessions = fixture.sessions
const boundary = yield* sessions.synthetic({ sessionID, text: "Revert boundary", resume: false })
yield* SessionInbox.promote(fixture.db, fixture.bus, sessionID, "steer")
yield* handle.revert.stage({ messageID: boundary.id, files: false })
yield* sessions.revert.stage({ sessionID, messageID: boundary.id, files: false })
const destination = Location.Ref.make({ directory: AbsolutePath.make("/project/moved") })
yield* fixture.bus.publish(SessionEvent.Moved, {
sessionID,
@@ -729,13 +760,13 @@ describe("Session-owned handles", () => {
subpath: RelativePath.make("moved"),
})
yield* handle.revert.stage({ messageID: boundary.id, files: false })
yield* handle.revert.clear()
yield* sessions.revert.stage({ sessionID, messageID: boundary.id, files: false })
yield* sessions.revert.clear(sessionID)
expect(captures).toEqual([source, destination])
expect(fixture.locations).toEqual([source, destination, destination])
expect(fixture.activationWaits).toEqual([])
expect((yield* handle.get()).revert).toBeUndefined()
expect((yield* sessions.get(sessionID)).revert).toBeUndefined()
}),
)
})
@@ -746,7 +777,7 @@ describe("SessionPrompt preparation", () => {
const fixture = yield* setup()
const input = { text: "Original", files: [{ uri: new URL("./session-owned.test.ts", import.meta.url).href }] }
const request = {
session: yield* fixture.sessions.forSession(sessionID).get(),
session: yield* fixture.sessions.get(sessionID),
messageID: SessionMessage.ID.create(),
input,
}
@@ -759,7 +790,7 @@ describe("SessionPrompt preparation", () => {
expect(items[0]).toMatchObject({ type: "user", payload: { text: "Original" }, delivery: "steer" })
expect(items[0]?.payload.files?.[0]?.mime).toBe("text/plain")
expect(input.text).toBe("Original")
expect(yield* fixture.sessions.forSession(sessionID).inbox()).toEqual([])
expect(yield* fixture.sessions.inbox(sessionID)).toEqual([])
expect(fixture.wakes).toEqual([])
}),
)
@@ -788,20 +819,20 @@ describe("SessionRevert operations", () => {
}),
}),
})
const handle = fixture.sessions.forSession(sessionID)
const boundary = yield* handle.synthetic({ text: "Revert boundary", resume: false })
const sessions = fixture.sessions
const boundary = yield* sessions.synthetic({ sessionID, text: "Revert boundary", resume: false })
yield* SessionInbox.promote(fixture.db, fixture.bus, sessionID, "steer")
expect(calls).toEqual([])
const session = yield* handle.get()
const session = yield* sessions.get(sessionID)
yield* SessionRevert.stage({ session, messageID: boundary.id, files: false }).pipe(
Effect.provideService(Instance.Service, fixture.instances),
)
expect(calls).toEqual(["capture", "capture", "diff"])
const staged = yield* handle.get()
const staged = yield* sessions.get(sessionID)
expect(staged.revert?.snapshot).toBe(Snapshot.ID.make("captured-tree"))
yield* SessionRevert.clear(staged).pipe(Effect.provideService(Instance.Service, fixture.instances))
const cleared = yield* handle.get()
const cleared = yield* sessions.get(sessionID)
expect(cleared.revert).toBeUndefined()
yield* SessionRevert.clear(cleared).pipe(Effect.provideService(Instance.Service, fixture.instances))
expect(calls).toEqual(["capture", "capture", "diff", "restore"])
@@ -824,13 +855,12 @@ describe("SessionInbox command contracts", () => {
}),
),
)
const handle = fixture.sessions.forSession(sessionID)
const pending = yield* handle.synthetic({ text: "Pending", resume: false })
const pending = yield* fixture.sessions.synthetic({ sessionID, text: "Pending", resume: false })
yield* handle.cancelInbox(pending.id).pipe(Effect.setContext(Context.empty()))
yield* fixture.sessions.cancelInbox({ sessionID, inboxID: pending.id }).pipe(Effect.setContext(Context.empty()))
expect(cancelled).toEqual([pending.id])
expect(yield* handle.inbox()).toEqual([])
expect(yield* fixture.sessions.inbox(sessionID)).toEqual([])
expect(fixture.wakes).toEqual([])
}),
)
+8 -11
View File
@@ -24,29 +24,26 @@ Manual compaction and Session movement use the same inbox as control items. Each
## Session Operations Own Admission Policy
Core's ID-addressed Session facade delegates prompt, synthetic, compaction, pending-input, shell, revert, execution controls, renaming, agent/model selection, viewed acknowledgements, and message lookup/editing to the ID-bound Session values in `session/session.ts`. Those operations own request idempotency, admission ordering, wake policy, and per-Session state-change rules. The facade supplies lazy Location services; the lower operations do not depend on `LocationServiceMap` or call back into the facade. Collection operations and host-infrastructure routing remain in the facade. Project resolution persists the owning Project; Session does not repeat that write.
The host acquires `Session.make` once within its host `Scope`, then selects ID-bound values through `forSession`. The factory captures that Scope for shell completion recording; it must outlive individual callers. The existing facade remains the Effect service; the lower factory does not introduce a second service or a Layer per Session ID.
`Session.Service` implements prompt, synthetic, compaction, pending-input, shell, skill, command, revert, execution controls, renaming, agent/model selection, viewed acknowledgements, and message lookup directly as ID-addressed operations. Those operations own request idempotency, admission ordering, wake policy, and per-Session state-change rules. They acquire Location services lazily through `Instance.Service`; they do not use `LocationServiceMap`. Project resolution persists the owning Project; Session does not repeat that write.
```ts
Effect.gen(function* () {
const sessions = yield* Session.make((ref) => locations.get(ref))
const session = sessions.forSession(sessionID)
yield* session.prompt({ text: "Inspect the failing tests", resume: false })
const sessions = yield* Session.Service
yield* sessions.prompt({ sessionID, text: "Inspect the failing tests", resume: false })
})
```
References share operation implementations, host services, and the existing execution coordinator. A Session value retains its ID, not a cached projection or permanently selected runner. Obtaining or discarding a value does not start or interrupt execution. Admission stays independently callable while execution is active.
The service layer captures its Scope for shell completion recording and skill-triggered resumes; it must outlive individual callers. There is no per-Session value or Layer: every operation takes the Session ID and reloads the current projection, so no cached state or permanently selected runner survives movement. Admission stays independently callable while execution is active.
User shell commands start immediately as background work without waiting for model execution or other user shells. They do not suppress prompt wakeups. The shell operation forks completion recording into the captured host Scope with `startImmediately: true` and joins that fiber, so caller cancellation does not cancel recording; closing the host Scope does. Shell started and ended events retain output in one shell entry. Completion and startup-failure notifications are admitted as synthetic input with `resume: false`, without waking execution. Shell services are resolved at startup, while Session events remain outside that Location context so movement does not pin them to the old Location.
User shell commands start immediately as background work without waiting for model execution or other user shells. They do not suppress prompt wakeups. The shell operation forks completion recording into the Session layer Scope with `startImmediately: true` and joins that fiber, so caller cancellation does not cancel recording; closing that Scope does. Shell started and ended events retain output in one shell entry. Completion and startup-failure notifications are admitted as synthetic input with `resume: false`, without waking execution. Shell services are resolved at startup, while Session events remain outside that Location context so movement does not pin them to the old Location.
Location services are acquired only when an operation needs them. In particular, retry reconciliation happens before prompt preparation, so an already-admitted input skips hooks and attachment resolution. Execution continues to resolve placement independently at drain start and after movement.
`servicesFor` selects instance services from the saved placement. Each instance constructs `SessionPrompt.Service`, whose `prepare` method turns submitted input into a user inbox item without admitting it, committing a revert, or waking execution. It captures FSUtil, PluginSupervisor, PluginHooks, Image, and Skill; readiness is checked before hooks on every call. Prompt keeps early retry reconciliation outside preparation and invokes preparation interruptibly in the current instance. Lower Session does not depend directly on Database or FSUtil; its Location requirements are SessionPrompt, SessionRevert, Shell, and the PluginSupervisor still used by manual shell startup.
`servicesFor` selects instance services from the saved placement. Each instance constructs `SessionPrompt.Service`, whose `prepare` method turns submitted input into a user inbox item without admitting it, committing a revert, or waking execution. It captures FSUtil, PluginSupervisor, PluginHooks, Image, and Skill; readiness is checked before hooks on every call. Prompt keeps early retry reconciliation outside preparation and invokes preparation interruptibly in the current instance.
Each instance constructs its `SessionRevert.Service` through `SessionRevert.make`, capturing Database, Bus, PluginSupervisor, and Snapshot. Stage and clear check plugin readiness on each invocation and require no service provisioning inside their implementations. Session methods select the current instance for each operation, so an ID-bound Session does not retain a previous instance's snapshots after movement. Commit uses only the host's captured Bus and does not acquire an instance.
Each instance constructs its `SessionRevert.Service` through `SessionRevert.make`, capturing Database, Bus, PluginSupervisor, and Snapshot. Stage and clear check plugin readiness on each invocation and require no service provisioning inside their implementations. Session methods select the current instance for each operation, so a Session operation does not retain a previous instance's snapshots after movement. Commit uses only the host's captured Bus and does not acquire an instance.
`SessionInbox.Service` is host-scoped. Its node depends on Database and Bus, and its layer uses `SessionInbox.make` to capture those services and construct admission and pending-input commands that take only domain inputs. Its `list` method supplies the normal pending-input read without exposing Database to Session. `Session.make` and the facade's move admission consume the registered service rather than constructing separate command objects. The service is not Session-ID scoped and does not own execution. Its commands retain the existing shared inbox serialization lock. Standalone query helpers, transaction-facing projectors, and runner delivery retain their explicit database/Bus inputs. The layer compiler is unchanged.
`SessionInbox.Service` is host-scoped. Its node depends on Database and Bus, and its layer uses `SessionInbox.make` to capture those services and construct admission and pending-input commands that take only domain inputs. Its `list` method supplies the normal pending-input read without exposing Database to Session. Session operations and move admission consume the registered service rather than constructing separate command objects. The service is not Session-ID scoped and does not own execution. Its commands retain the existing shared inbox serialization lock. Standalone query helpers, transaction-facing projectors, and runner delivery retain their explicit database/Bus inputs. The layer compiler is unchanged.
Inbox commands own identity and type checks and return typed `SessionInbox.LifecycleConflict` errors. Session operations translate these into their public operation-specific errors and decide whether to wake execution. The Bus/projector boundary still uses defects to abort invalid projections; Inbox translates only lifecycle conflicts, not unrelated defects. Pending-input mutation does not schedule execution itself: steering wakes after a successful mutation, while queueing and cancellation do not.