|
|
@@ -26,6 +26,9 @@ import { fromRow } from "./session/info"
|
|
|
import { SessionRunner } from "./session/runner/index"
|
|
|
import { SessionStore } from "./session/store"
|
|
|
import { SessionExecution } from "./session/execution"
|
|
|
+import { SessionExecutionLocal } from "./session/execution/local"
|
|
|
+import { makeGlobalNode } from "./effect/scoped-node"
|
|
|
+import { LocationServiceMap, node as locationServiceMapNode } from "./location-service-map"
|
|
|
import { MessageDecodeError } from "./session/error"
|
|
|
import { SessionEvent } from "./session/event"
|
|
|
import { SessionInput } from "./session/input"
|
|
|
@@ -170,279 +173,267 @@ export interface Interface {
|
|
|
|
|
|
export class Service extends Context.Service<Service, Interface>()("@opencode/v2/Session") {}
|
|
|
|
|
|
-export const layer = Layer.unwrap(
|
|
|
- Effect.promise(() => import("./location-layer")).pipe(
|
|
|
- Effect.map(({ LocationServiceMap }) =>
|
|
|
- Layer.effect(
|
|
|
- Service,
|
|
|
- Effect.gen(function* () {
|
|
|
- const database = yield* Database.Service
|
|
|
- const db = database.db
|
|
|
- const events = yield* EventV2.Service
|
|
|
- const projects = yield* ProjectV2.Service
|
|
|
- const execution = yield* SessionExecution.Service
|
|
|
- const store = yield* SessionStore.Service
|
|
|
- const locations = yield* LocationServiceMap
|
|
|
- const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Message)
|
|
|
- const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
|
|
|
- const decode = (row: typeof SessionMessageTable.$inferSelect) =>
|
|
|
- decodeMessage({ ...row.data, id: row.id, type: row.type }).pipe(
|
|
|
- Effect.mapError(
|
|
|
- () =>
|
|
|
- new MessageDecodeError({
|
|
|
- sessionID: SessionSchema.ID.make(row.session_id),
|
|
|
- messageID: SessionMessage.ID.make(row.id),
|
|
|
- }),
|
|
|
- ),
|
|
|
- )
|
|
|
-
|
|
|
- const result = Service.of({
|
|
|
- create: Effect.fn("V2Session.create")(function* (input) {
|
|
|
- const sessionID = input.id ?? SessionSchema.ID.create()
|
|
|
- const recorded = yield* store.get(sessionID)
|
|
|
- if (recorded) return recorded
|
|
|
- const project = yield* projects.resolve(input.location.directory)
|
|
|
- yield* db
|
|
|
- .insert(ProjectTable)
|
|
|
- .values({ id: project.id, worktree: project.directory, vcs: project.vcs?.type, sandboxes: [] })
|
|
|
- .onConflictDoNothing()
|
|
|
- .run()
|
|
|
- .pipe(Effect.orDie)
|
|
|
- const now = Date.now()
|
|
|
- const info = SessionV1.SessionInfo.make({
|
|
|
- id: sessionID,
|
|
|
- slug: Slug.create(),
|
|
|
- version: InstallationVersion,
|
|
|
- projectID: project.id,
|
|
|
- directory: input.location.directory,
|
|
|
- path: path.relative(project.directory, input.location.directory).replaceAll("\\", "/"),
|
|
|
- workspaceID: input.location.workspaceID ? WorkspaceV2.ID.make(input.location.workspaceID) : undefined,
|
|
|
- title: `New session - ${new Date(now).toISOString()}`,
|
|
|
- agent: input.agent,
|
|
|
- model: input.model
|
|
|
- ? {
|
|
|
- id: ModelV2.ID.make(input.model.id),
|
|
|
- providerID: input.model.providerID,
|
|
|
- variant: input.model.variant,
|
|
|
- }
|
|
|
- : undefined,
|
|
|
- cost: 0,
|
|
|
- tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
|
|
- time: { created: now, updated: now },
|
|
|
- })
|
|
|
- const projected = yield* events
|
|
|
- .publish(SessionV1.Event.Created, { sessionID, info }, { location: input.location })
|
|
|
- .pipe(
|
|
|
- Effect.as({ type: "created" } as const),
|
|
|
- Effect.catchDefect((defect) => {
|
|
|
- if (!(defect instanceof SessionProjector.SessionAlreadyProjected)) {
|
|
|
- return Effect.die(defect)
|
|
|
- }
|
|
|
- // Concurrent creation lost the projection race. The existing Session identity wins.
|
|
|
- return store
|
|
|
- .get(sessionID)
|
|
|
- .pipe(
|
|
|
- Effect.flatMap((session) =>
|
|
|
- session ? Effect.succeed({ type: "existing", session } as const) : Effect.die(defect),
|
|
|
- ),
|
|
|
- )
|
|
|
- }),
|
|
|
- )
|
|
|
- if (projected.type === "existing") return projected.session
|
|
|
- // TODO: Restore recorded sessions onto replacement synchronized workspaces in a future API slice.
|
|
|
- return yield* result.get(sessionID).pipe(Effect.orDie)
|
|
|
+export const layer = Layer.effect(
|
|
|
+ Service,
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const database = yield* Database.Service
|
|
|
+ const db = database.db
|
|
|
+ const events = yield* EventV2.Service
|
|
|
+ const projects = yield* ProjectV2.Service
|
|
|
+ const execution = yield* SessionExecution.Service
|
|
|
+ const store = yield* SessionStore.Service
|
|
|
+ const locations = yield* LocationServiceMap
|
|
|
+ const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Message)
|
|
|
+ const isDurableSessionEvent = Schema.is(SessionEvent.Durable)
|
|
|
+ const decode = (row: typeof SessionMessageTable.$inferSelect) =>
|
|
|
+ decodeMessage({ ...row.data, id: row.id, type: row.type }).pipe(
|
|
|
+ Effect.mapError(
|
|
|
+ () =>
|
|
|
+ new MessageDecodeError({
|
|
|
+ sessionID: SessionSchema.ID.make(row.session_id),
|
|
|
+ messageID: SessionMessage.ID.make(row.id),
|
|
|
}),
|
|
|
- get: Effect.fn("V2Session.get")(function* (sessionID) {
|
|
|
- const session = yield* store.get(sessionID)
|
|
|
- if (!session) return yield* new NotFoundError({ sessionID })
|
|
|
- return session
|
|
|
- }),
|
|
|
- list: Effect.fn("V2Session.list")(function* (input = {}) {
|
|
|
- const direction = input.anchor?.direction ?? "next"
|
|
|
- const requestedOrder = input.order ?? "desc"
|
|
|
- const order = direction === "previous" ? (requestedOrder === "asc" ? "desc" : "asc") : requestedOrder
|
|
|
- const sortColumn = SessionTable.time_created
|
|
|
- const conditions: SQL[] = []
|
|
|
- if ("directory" in input) conditions.push(eq(SessionTable.directory, input.directory))
|
|
|
- if (input.workspaceID) conditions.push(eq(SessionTable.workspace_id, input.workspaceID))
|
|
|
- if ("project" in input) conditions.push(eq(SessionTable.project_id, input.project))
|
|
|
- if (input.search) conditions.push(like(SessionTable.title, `%${input.search}%`))
|
|
|
- if (input.anchor) {
|
|
|
- conditions.push(
|
|
|
- order === "asc"
|
|
|
- ? or(
|
|
|
- gt(sortColumn, input.anchor.time),
|
|
|
- and(eq(sortColumn, input.anchor.time), gt(SessionTable.id, input.anchor.id)),
|
|
|
- )!
|
|
|
- : or(
|
|
|
- lt(sortColumn, input.anchor.time),
|
|
|
- and(eq(sortColumn, input.anchor.time), lt(SessionTable.id, input.anchor.id)),
|
|
|
- )!,
|
|
|
- )
|
|
|
+ ),
|
|
|
+ )
|
|
|
+
|
|
|
+ const result = Service.of({
|
|
|
+ create: Effect.fn("V2Session.create")(function* (input) {
|
|
|
+ const sessionID = input.id ?? SessionSchema.ID.create()
|
|
|
+ const recorded = yield* store.get(sessionID)
|
|
|
+ if (recorded) return recorded
|
|
|
+ const project = yield* projects.resolve(input.location.directory)
|
|
|
+ yield* db
|
|
|
+ .insert(ProjectTable)
|
|
|
+ .values({ id: project.id, worktree: project.directory, vcs: project.vcs?.type, sandboxes: [] })
|
|
|
+ .onConflictDoNothing()
|
|
|
+ .run()
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ const now = Date.now()
|
|
|
+ const info = SessionV1.SessionInfo.make({
|
|
|
+ id: sessionID,
|
|
|
+ slug: Slug.create(),
|
|
|
+ version: InstallationVersion,
|
|
|
+ projectID: project.id,
|
|
|
+ directory: input.location.directory,
|
|
|
+ path: path.relative(project.directory, input.location.directory).replaceAll("\\", "/"),
|
|
|
+ workspaceID: input.location.workspaceID ? WorkspaceV2.ID.make(input.location.workspaceID) : undefined,
|
|
|
+ title: `New session - ${new Date(now).toISOString()}`,
|
|
|
+ agent: input.agent,
|
|
|
+ model: input.model
|
|
|
+ ? {
|
|
|
+ id: ModelV2.ID.make(input.model.id),
|
|
|
+ providerID: input.model.providerID,
|
|
|
+ variant: input.model.variant,
|
|
|
+ }
|
|
|
+ : undefined,
|
|
|
+ cost: 0,
|
|
|
+ tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
|
|
+ time: { created: now, updated: now },
|
|
|
+ })
|
|
|
+ const projected = yield* events
|
|
|
+ .publish(SessionV1.Event.Created, { sessionID, info }, { location: input.location })
|
|
|
+ .pipe(
|
|
|
+ Effect.as({ type: "created" } as const),
|
|
|
+ Effect.catchDefect((defect) => {
|
|
|
+ if (!(defect instanceof SessionProjector.SessionAlreadyProjected)) {
|
|
|
+ return Effect.die(defect)
|
|
|
}
|
|
|
- const query = db
|
|
|
- .select()
|
|
|
- .from(SessionTable)
|
|
|
- .where(conditions.length > 0 ? and(...conditions) : undefined)
|
|
|
- .orderBy(
|
|
|
- order === "asc" ? asc(sortColumn) : desc(sortColumn),
|
|
|
- order === "asc" ? asc(SessionTable.id) : desc(SessionTable.id),
|
|
|
+ // Concurrent creation lost the projection race. The existing Session identity wins.
|
|
|
+ return store
|
|
|
+ .get(sessionID)
|
|
|
+ .pipe(
|
|
|
+ Effect.flatMap((session) =>
|
|
|
+ session ? Effect.succeed({ type: "existing", session } as const) : Effect.die(defect),
|
|
|
+ ),
|
|
|
)
|
|
|
- const rows = yield* (input.limit === undefined ? query.all() : query.limit(input.limit).all()).pipe(
|
|
|
- Effect.orDie,
|
|
|
- )
|
|
|
- return (direction === "previous" ? rows.toReversed() : rows).map((row) => fromRow(row))
|
|
|
}),
|
|
|
- messages: Effect.fn("V2Session.messages")(function* (input) {
|
|
|
- yield* result.get(input.sessionID)
|
|
|
- const direction = input.cursor?.direction ?? "next"
|
|
|
- const requestedOrder = input.order ?? "desc"
|
|
|
- const order = direction === "previous" ? (requestedOrder === "asc" ? "desc" : "asc") : requestedOrder
|
|
|
- const anchor = input.cursor
|
|
|
- ? yield* db
|
|
|
- .select({ seq: SessionMessageTable.seq })
|
|
|
- .from(SessionMessageTable)
|
|
|
- .where(
|
|
|
- and(
|
|
|
- eq(SessionMessageTable.session_id, input.sessionID),
|
|
|
- eq(SessionMessageTable.id, input.cursor.id),
|
|
|
- ),
|
|
|
- )
|
|
|
- .get()
|
|
|
- .pipe(Effect.orDie)
|
|
|
- : undefined
|
|
|
- if (input.cursor && !anchor) return []
|
|
|
- const boundary = anchor
|
|
|
- ? order === "asc"
|
|
|
- ? gt(SessionMessageTable.seq, anchor.seq)
|
|
|
- : lt(SessionMessageTable.seq, anchor.seq)
|
|
|
- : undefined
|
|
|
- const where = boundary
|
|
|
- ? and(eq(SessionMessageTable.session_id, input.sessionID), boundary)
|
|
|
- : eq(SessionMessageTable.session_id, input.sessionID)
|
|
|
- const query = db
|
|
|
- .select()
|
|
|
- .from(SessionMessageTable)
|
|
|
- .where(where)
|
|
|
- .orderBy(order === "asc" ? asc(SessionMessageTable.seq) : desc(SessionMessageTable.seq))
|
|
|
- const rows = yield* (input.limit === undefined ? query.all() : query.limit(input.limit).all()).pipe(
|
|
|
- Effect.orDie,
|
|
|
+ )
|
|
|
+ if (projected.type === "existing") return projected.session
|
|
|
+ // TODO: Restore recorded sessions onto replacement synchronized workspaces in a future API slice.
|
|
|
+ return yield* result.get(sessionID).pipe(Effect.orDie)
|
|
|
+ }),
|
|
|
+ get: Effect.fn("V2Session.get")(function* (sessionID) {
|
|
|
+ const session = yield* store.get(sessionID)
|
|
|
+ if (!session) return yield* new NotFoundError({ sessionID })
|
|
|
+ return session
|
|
|
+ }),
|
|
|
+ list: Effect.fn("V2Session.list")(function* (input = {}) {
|
|
|
+ const direction = input.anchor?.direction ?? "next"
|
|
|
+ const requestedOrder = input.order ?? "desc"
|
|
|
+ const order = direction === "previous" ? (requestedOrder === "asc" ? "desc" : "asc") : requestedOrder
|
|
|
+ const sortColumn = SessionTable.time_created
|
|
|
+ const conditions: SQL[] = []
|
|
|
+ if ("directory" in input) conditions.push(eq(SessionTable.directory, input.directory))
|
|
|
+ if (input.workspaceID) conditions.push(eq(SessionTable.workspace_id, input.workspaceID))
|
|
|
+ if ("project" in input) conditions.push(eq(SessionTable.project_id, input.project))
|
|
|
+ if (input.search) conditions.push(like(SessionTable.title, `%${input.search}%`))
|
|
|
+ if (input.anchor) {
|
|
|
+ conditions.push(
|
|
|
+ order === "asc"
|
|
|
+ ? or(
|
|
|
+ gt(sortColumn, input.anchor.time),
|
|
|
+ and(eq(sortColumn, input.anchor.time), gt(SessionTable.id, input.anchor.id)),
|
|
|
+ )!
|
|
|
+ : or(
|
|
|
+ lt(sortColumn, input.anchor.time),
|
|
|
+ and(eq(sortColumn, input.anchor.time), lt(SessionTable.id, input.anchor.id)),
|
|
|
+ )!,
|
|
|
+ )
|
|
|
+ }
|
|
|
+ const query = db
|
|
|
+ .select()
|
|
|
+ .from(SessionTable)
|
|
|
+ .where(conditions.length > 0 ? and(...conditions) : undefined)
|
|
|
+ .orderBy(
|
|
|
+ order === "asc" ? asc(sortColumn) : desc(sortColumn),
|
|
|
+ order === "asc" ? asc(SessionTable.id) : desc(SessionTable.id),
|
|
|
+ )
|
|
|
+ const rows = yield* (input.limit === undefined ? query.all() : query.limit(input.limit).all()).pipe(
|
|
|
+ Effect.orDie,
|
|
|
+ )
|
|
|
+ return (direction === "previous" ? rows.toReversed() : rows).map((row) => fromRow(row))
|
|
|
+ }),
|
|
|
+ messages: Effect.fn("V2Session.messages")(function* (input) {
|
|
|
+ yield* result.get(input.sessionID)
|
|
|
+ const direction = input.cursor?.direction ?? "next"
|
|
|
+ const requestedOrder = input.order ?? "desc"
|
|
|
+ const order = direction === "previous" ? (requestedOrder === "asc" ? "desc" : "asc") : requestedOrder
|
|
|
+ const anchor = input.cursor
|
|
|
+ ? yield* db
|
|
|
+ .select({ seq: SessionMessageTable.seq })
|
|
|
+ .from(SessionMessageTable)
|
|
|
+ .where(
|
|
|
+ and(eq(SessionMessageTable.session_id, input.sessionID), eq(SessionMessageTable.id, input.cursor.id)),
|
|
|
)
|
|
|
- return yield* Effect.forEach(direction === "previous" ? rows.toReversed() : rows, decode)
|
|
|
- }),
|
|
|
- message: Effect.fn("V2Session.message")(function* (input) {
|
|
|
- const stored = yield* store.message(input.messageID)
|
|
|
- return stored?.sessionID === input.sessionID ? stored.message : undefined
|
|
|
- }),
|
|
|
- context: Effect.fn("V2Session.context")(function* (sessionID) {
|
|
|
- yield* result.get(sessionID)
|
|
|
- return yield* store.context(sessionID)
|
|
|
- }),
|
|
|
- events: (input) =>
|
|
|
- Stream.unwrap(
|
|
|
- result
|
|
|
- .get(input.sessionID)
|
|
|
- .pipe(Effect.as(events.durable({ aggregateID: input.sessionID, after: input.after }))),
|
|
|
- ).pipe(Stream.filter((event): event is SessionEvent.DurableEvent => isDurableSessionEvent(event))),
|
|
|
- prompt: Effect.fn("V2Session.prompt")((input) =>
|
|
|
- Effect.uninterruptible(
|
|
|
- Effect.gen(function* () {
|
|
|
- yield* result.get(input.sessionID)
|
|
|
- const messageID = input.id ?? SessionMessage.ID.create()
|
|
|
- const delivery = input.delivery ?? "steer"
|
|
|
- const expected = { sessionID: input.sessionID, messageID, prompt: input.prompt, delivery }
|
|
|
- const admitted = yield* SessionInput.admit(db, events, {
|
|
|
- id: messageID,
|
|
|
- sessionID: input.sessionID,
|
|
|
- prompt: input.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* execution.wake(admitted.sessionID)
|
|
|
- return admitted
|
|
|
- }),
|
|
|
+ .get()
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ : undefined
|
|
|
+ if (input.cursor && !anchor) return []
|
|
|
+ const boundary = anchor
|
|
|
+ ? order === "asc"
|
|
|
+ ? gt(SessionMessageTable.seq, anchor.seq)
|
|
|
+ : lt(SessionMessageTable.seq, anchor.seq)
|
|
|
+ : undefined
|
|
|
+ const where = boundary
|
|
|
+ ? and(eq(SessionMessageTable.session_id, input.sessionID), boundary)
|
|
|
+ : eq(SessionMessageTable.session_id, input.sessionID)
|
|
|
+ const query = db
|
|
|
+ .select()
|
|
|
+ .from(SessionMessageTable)
|
|
|
+ .where(where)
|
|
|
+ .orderBy(order === "asc" ? asc(SessionMessageTable.seq) : desc(SessionMessageTable.seq))
|
|
|
+ const rows = yield* (input.limit === undefined ? query.all() : query.limit(input.limit).all()).pipe(
|
|
|
+ Effect.orDie,
|
|
|
+ )
|
|
|
+ return yield* Effect.forEach(direction === "previous" ? rows.toReversed() : rows, decode)
|
|
|
+ }),
|
|
|
+ message: Effect.fn("V2Session.message")(function* (input) {
|
|
|
+ const stored = yield* store.message(input.messageID)
|
|
|
+ return stored?.sessionID === input.sessionID ? stored.message : undefined
|
|
|
+ }),
|
|
|
+ context: Effect.fn("V2Session.context")(function* (sessionID) {
|
|
|
+ yield* result.get(sessionID)
|
|
|
+ return yield* store.context(sessionID)
|
|
|
+ }),
|
|
|
+ events: (input) =>
|
|
|
+ Stream.unwrap(
|
|
|
+ result
|
|
|
+ .get(input.sessionID)
|
|
|
+ .pipe(Effect.as(events.durable({ aggregateID: input.sessionID, after: input.after }))),
|
|
|
+ ).pipe(Stream.filter((event): event is SessionEvent.DurableEvent => isDurableSessionEvent(event))),
|
|
|
+ prompt: Effect.fn("V2Session.prompt")((input) =>
|
|
|
+ Effect.uninterruptible(
|
|
|
+ Effect.gen(function* () {
|
|
|
+ yield* result.get(input.sessionID)
|
|
|
+ const messageID = input.id ?? SessionMessage.ID.create()
|
|
|
+ const delivery = input.delivery ?? "steer"
|
|
|
+ const expected = { sessionID: input.sessionID, messageID, prompt: input.prompt, delivery }
|
|
|
+ const admitted = yield* SessionInput.admit(db, events, {
|
|
|
+ id: messageID,
|
|
|
+ sessionID: input.sessionID,
|
|
|
+ prompt: input.prompt,
|
|
|
+ delivery,
|
|
|
+ }).pipe(
|
|
|
+ Effect.catchDefect((defect) =>
|
|
|
+ defect instanceof SessionInput.LifecycleConflict
|
|
|
+ ? new PromptConflictError({ sessionID: input.sessionID, messageID })
|
|
|
+ : Effect.die(defect),
|
|
|
),
|
|
|
- ),
|
|
|
- shell: Effect.fn("V2Session.shell")(function* () {
|
|
|
- return yield* new OperationUnavailableError({ operation: "shell" })
|
|
|
- }),
|
|
|
- skill: Effect.fn("V2Session.skill")(function* () {
|
|
|
- return yield* new OperationUnavailableError({ operation: "skill" })
|
|
|
- }),
|
|
|
- switchAgent: Effect.fn("V2Session.switchAgent")(function* (input) {
|
|
|
- yield* result.get(input.sessionID)
|
|
|
- yield* events.publish(SessionEvent.AgentSwitched, {
|
|
|
- sessionID: input.sessionID,
|
|
|
- messageID: SessionMessage.ID.create(),
|
|
|
- timestamp: yield* DateTime.now,
|
|
|
- agent: input.agent,
|
|
|
- })
|
|
|
- }),
|
|
|
- switchModel: Effect.fn("V2Session.switchModel")(function* (input) {
|
|
|
- yield* result.get(input.sessionID)
|
|
|
- yield* events.publish(SessionEvent.ModelSwitched, {
|
|
|
- sessionID: input.sessionID,
|
|
|
- messageID: SessionMessage.ID.create(),
|
|
|
- timestamp: yield* DateTime.now,
|
|
|
- model: input.model,
|
|
|
- })
|
|
|
- }),
|
|
|
- compact: Effect.fn("V2Session.compact")(function* (input) {
|
|
|
- yield* result.get(input.sessionID)
|
|
|
- return yield* new OperationUnavailableError({ operation: "compact" })
|
|
|
- }),
|
|
|
- wait: Effect.fn("V2Session.wait")(function* (sessionID) {
|
|
|
- yield* result.get(sessionID)
|
|
|
- return yield* new OperationUnavailableError({ operation: "wait" })
|
|
|
- }),
|
|
|
- resume: Effect.fn("V2Session.resume")(function* (sessionID) {
|
|
|
- yield* result.get(sessionID)
|
|
|
- yield* execution.resume(sessionID)
|
|
|
- }),
|
|
|
- interrupt: Effect.fn("V2Session.interrupt")((sessionID) =>
|
|
|
- Effect.uninterruptible(execution.interrupt(sessionID)),
|
|
|
- ),
|
|
|
- revert: {
|
|
|
- stage: Effect.fn("V2Session.revert.stage")(function* (input) {
|
|
|
- const session = yield* result.get(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) {
|
|
|
- const session = yield* result.get(sessionID)
|
|
|
- yield* SessionRevert.clear(session).pipe(
|
|
|
- Effect.provideService(EventV2.Service, events),
|
|
|
- Effect.provide(locations.get(session.location)),
|
|
|
- )
|
|
|
- }),
|
|
|
- commit: Effect.fn("V2Session.revert.commit")(function* (sessionID) {
|
|
|
- const session = yield* result.get(sessionID)
|
|
|
- yield* SessionRevert.commit(session).pipe(Effect.provideService(EventV2.Service, events))
|
|
|
- }),
|
|
|
- },
|
|
|
- })
|
|
|
-
|
|
|
- return result
|
|
|
- }),
|
|
|
+ )
|
|
|
+ if (!SessionInput.equivalent(admitted, expected))
|
|
|
+ return yield* new PromptConflictError({ sessionID: input.sessionID, messageID })
|
|
|
+ if (input.resume !== false) yield* execution.wake(admitted.sessionID)
|
|
|
+ return admitted
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ ),
|
|
|
+ shell: Effect.fn("V2Session.shell")(function* () {
|
|
|
+ return yield* new OperationUnavailableError({ operation: "shell" })
|
|
|
+ }),
|
|
|
+ skill: Effect.fn("V2Session.skill")(function* () {
|
|
|
+ return yield* new OperationUnavailableError({ operation: "skill" })
|
|
|
+ }),
|
|
|
+ switchAgent: Effect.fn("V2Session.switchAgent")(function* (input) {
|
|
|
+ yield* result.get(input.sessionID)
|
|
|
+ yield* events.publish(SessionEvent.AgentSwitched, {
|
|
|
+ sessionID: input.sessionID,
|
|
|
+ messageID: SessionMessage.ID.create(),
|
|
|
+ timestamp: yield* DateTime.now,
|
|
|
+ agent: input.agent,
|
|
|
+ })
|
|
|
+ }),
|
|
|
+ switchModel: Effect.fn("V2Session.switchModel")(function* (input) {
|
|
|
+ yield* result.get(input.sessionID)
|
|
|
+ yield* events.publish(SessionEvent.ModelSwitched, {
|
|
|
+ sessionID: input.sessionID,
|
|
|
+ messageID: SessionMessage.ID.create(),
|
|
|
+ timestamp: yield* DateTime.now,
|
|
|
+ model: input.model,
|
|
|
+ })
|
|
|
+ }),
|
|
|
+ compact: Effect.fn("V2Session.compact")(function* (input) {
|
|
|
+ yield* result.get(input.sessionID)
|
|
|
+ return yield* new OperationUnavailableError({ operation: "compact" })
|
|
|
+ }),
|
|
|
+ wait: Effect.fn("V2Session.wait")(function* (sessionID) {
|
|
|
+ yield* result.get(sessionID)
|
|
|
+ return yield* new OperationUnavailableError({ operation: "wait" })
|
|
|
+ }),
|
|
|
+ resume: Effect.fn("V2Session.resume")(function* (sessionID) {
|
|
|
+ yield* result.get(sessionID)
|
|
|
+ yield* execution.resume(sessionID)
|
|
|
+ }),
|
|
|
+ interrupt: Effect.fn("V2Session.interrupt")((sessionID) =>
|
|
|
+ Effect.uninterruptible(execution.interrupt(sessionID)),
|
|
|
),
|
|
|
- ),
|
|
|
- ),
|
|
|
+ revert: {
|
|
|
+ stage: Effect.fn("V2Session.revert.stage")(function* (input) {
|
|
|
+ const session = yield* result.get(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) {
|
|
|
+ const session = yield* result.get(sessionID)
|
|
|
+ yield* SessionRevert.clear(session).pipe(
|
|
|
+ Effect.provideService(EventV2.Service, events),
|
|
|
+ Effect.provide(locations.get(session.location)),
|
|
|
+ )
|
|
|
+ }),
|
|
|
+ commit: Effect.fn("V2Session.revert.commit")(function* (sessionID) {
|
|
|
+ const session = yield* result.get(sessionID)
|
|
|
+ yield* SessionRevert.commit(session).pipe(Effect.provideService(EventV2.Service, events))
|
|
|
+ }),
|
|
|
+ },
|
|
|
+ })
|
|
|
+
|
|
|
+ return result
|
|
|
+ }),
|
|
|
)
|
|
|
|
|
|
export const defaultLayer = layer.pipe(
|
|
|
- Layer.provide(
|
|
|
- Layer.unwrap(Effect.promise(() => import("./location-layer")).pipe(Effect.map((m) => m.LocationServiceMap.layer))),
|
|
|
- ),
|
|
|
Layer.provide(SessionExecution.noopLayer),
|
|
|
Layer.provide(SessionStore.defaultLayer),
|
|
|
Layer.provide(SessionProjector.defaultLayer),
|
|
|
@@ -451,3 +442,17 @@ export const defaultLayer = layer.pipe(
|
|
|
Layer.provide(ProjectV2.defaultLayer),
|
|
|
Layer.orDie,
|
|
|
)
|
|
|
+
|
|
|
+export const node = makeGlobalNode({
|
|
|
+ service: Service,
|
|
|
+ layer: layer.pipe(Layer.orDie),
|
|
|
+ deps: [
|
|
|
+ Database.node,
|
|
|
+ EventV2.node,
|
|
|
+ ProjectV2.node,
|
|
|
+ SessionExecutionLocal.node,
|
|
|
+ SessionStore.node,
|
|
|
+ locationServiceMapNode,
|
|
|
+ SessionProjector.node,
|
|
|
+ ],
|
|
|
+})
|