From c23cb367adf99fb4cd201f13f1a09d45fb3dda65 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Wed, 2 Sep 2026 21:48:20 -0400 Subject: [PATCH 1/3] fix(core): recover idle moves through the selected instance --- packages/core/src/session.ts | 21 +- packages/core/test/session-move.test.ts | 326 +++++++++++++++++++++++- 2 files changed, 340 insertions(+), 7 deletions(-) diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index df1e4b223684..1b85d435a977 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -1,7 +1,7 @@ export * as Session from "./session.js" export * from "./session/schema.js" -import { Effect, Layer, Schema, Context, Stream } from "effect" +import { Cause, Effect, Layer, Schema, Context, Stream } from "effect" import { LLMClient } from "@opencode-ai/ai" import { ListAnchor } from "@opencode-ai/schema/session" import { and, desc, eq } from "drizzle-orm" @@ -418,6 +418,17 @@ const layer = Layer.effect( Effect.provideService(Project.Service, projects), Effect.provideService(LocationServiceMap.Service, locations), ) + // Probe the same instance execution would use, without holding up inbox cancellation. + const unavailable = + !(yield* execution.isActive(input.sessionID)) && + (!(yield* fs.isDir(session.location.directory)) || + !(yield* SessionRunner.Service.pipe( + instances.provide(session), + Effect.as(true), + Effect.catchCause((cause) => + Cause.hasInterrupts(cause) ? Effect.failCause(cause) : Effect.succeed(false), + ), + ))) const item = SessionInbox.Item.make({ type: "move", payload, @@ -427,9 +438,13 @@ const layer = Layer.effect( input.sessionID, Effect.gen(function* () { const latest = yield* result.get(input.sessionID) - const source = yield* fs.stat(latest.location.directory).pipe(Effect.orElseSucceed(() => undefined)) // Active runners must hand off at a step boundary to retain their continuation. - if ((!source || source.type !== "Directory") && !(yield* execution.isActive(input.sessionID))) { + if ( + unavailable && + latest.location.directory === session.location.directory && + latest.location.workspaceID === session.location.workspaceID && + !(yield* execution.isActive(input.sessionID)) + ) { const cancellations = (yield* SessionInbox.moveIDs(db, input.sessionID)).map( (item) => [SessionEvent.InboxCancelled, { sessionID: input.sessionID, inboxID: item.id }] as const, ) diff --git a/packages/core/test/session-move.test.ts b/packages/core/test/session-move.test.ts index 2f297fbe5fa6..ab415fc856d1 100644 --- a/packages/core/test/session-move.test.ts +++ b/packages/core/test/session-move.test.ts @@ -1,11 +1,14 @@ import { describe, expect } from "bun:test" import path from "path" -import { mkdir, rm } from "fs/promises" -import { Effect, Layer, LayerMap } from "effect" +import { chmod, mkdir, readdir, rm } from "fs/promises" +import { Cause, Context, Deferred, Duration, Effect, Exit, Fiber, Layer, LayerMap, Queue } from "effect" import { Worktree } from "@opencode-ai/schema/worktree" +import { Workspace } from "@opencode-ai/schema/workspace" import { Bus } from "@opencode-ai/core/bus" +import { Config } from "@opencode-ai/core/config" import { Database } from "@opencode-ai/core/database/database" import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" +import { Instance } from "@opencode-ai/core/instance" import { Location } from "@opencode-ai/core/location" import { LocationServiceMap } from "@opencode-ai/core/location-service-map" import type { LocationServices } from "@opencode-ai/core/location-services" @@ -18,6 +21,9 @@ import { SessionProjector } from "@opencode-ai/core/session/projector" import { SessionRunner } from "@opencode-ai/core/session/runner/index" import { SessionStore } from "@opencode-ai/core/session/store" import { LayerNode } from "@opencode-ai/util/effect/layer-node" +import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" +import { Global } from "@opencode-ai/util/global" +import { tempGlobalLayer } from "./fixture/global" import { offlineModels } from "./fixture/models" import { tmpdirScoped } from "./fixture/tmpdir" import { testEffect } from "./lib/effect" @@ -74,8 +80,312 @@ const itWithUnavailableDestination = testEffect( ], ), ) +const itWithExecution = testEffect( + AppNodeBuilder.build(LayerNode.group([Session.node, SessionExecution.node]), [ + Global.node.replace(tempGlobalLayer), + offlineModels, + ]), +) +// Windows does not enforce POSIX mode bits, and root can traverse mode-000 directories. +const itWithPermissions = + process.platform === "win32" || process.getuid?.() === 0 ? itWithExecution.live.skip : itWithExecution.live +const itWithInstance = testEffect(Layer.empty) +const sourceProbe = (options: { execution?: boolean } = {}) => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const source = AbsolutePath.make(path.join(tmp.path, "source")) + const destination = AbsolutePath.make(tmp.path) + yield* Effect.promise(() => mkdir(source)) + const probes = yield* Queue.unbounded>() + const context = yield* Layer.build( + AppNodeBuilder.build(LayerNode.group([Session.node, Bus.node, SessionExecution.node]), [ + Global.node.replace(tempGlobalLayer), + ...(options.execution ? [] : [SessionExecution.node.replace(SessionExecution.noopLayer)]), + offlineModels, + Instance.node.replace( + makeGlobalNode({ + service: Instance.Service, + deps: [LocationServiceMap.node], + layer: Layer.effect( + Instance.Service, + Effect.gen(function* () { + const locations = yield* LocationServiceMap.Service + return Instance.Service.of({ + provide: (session) => (effect) => + Effect.gen(function* () { + if (session.location.directory === source) { + const release = yield* Deferred.make() + yield* Queue.offer(probes, release) + yield* Deferred.await(release) + } + return yield* effect.pipe(Effect.provide(locations.get(session.location))) + }), + }) + }), + ), + }), + ), + ]), + ) + return { + source, + destination, + probes, + session: Context.get(context, Session.Service), + bus: Context.get(context, Bus.Service), + execution: Context.get(context, SessionExecution.Service), + } + }) describe("Session.move", () => { + for (const broken of [false, true]) { + itWithExecution.live( + `moves an idle session from ${broken ? "broken" : "healthy"} source configuration`, + () => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const source = AbsolutePath.make(path.join(tmp.path, "source")) + const destination = AbsolutePath.make(path.join(tmp.path, "destination")) + yield* Effect.promise(() => Promise.all([mkdir(source), mkdir(destination)])) + if (broken) + yield* Effect.promise(() => + Bun.write(path.join(source, "opencode.json"), JSON.stringify({ instructions: ["{file:./missing.txt}"] })), + ) + const session = yield* Session.Service + const execution = yield* SessionExecution.Service + const created = yield* session.create({ location: Location.Ref.make({ directory: source }) }) + + yield* session.move({ sessionID: created.id, directory: destination }) + yield* execution.awaitIdle(created.id) + + expect((yield* session.get(created.id)).location.directory).toBe(destination) + expect(yield* session.inbox(created.id)).toEqual([]) + }), + { timeout: 15_000 }, + ) + } + + itWithPermissions( + "recovers an idle session from an unreadable source directory", + () => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const source = AbsolutePath.make(path.join(tmp.path, "source")) + const destination = AbsolutePath.make(path.join(tmp.path, "destination")) + yield* Effect.promise(() => Promise.all([mkdir(source), mkdir(destination)])) + const session = yield* Session.Service + const execution = yield* SessionExecution.Service + const created = yield* session.create({ location: Location.Ref.make({ directory: source }) }) + yield* Effect.addFinalizer(() => Effect.promise(() => chmod(source, 0o755))) + yield* Effect.promise(() => chmod(source, 0o000)) + expect( + yield* Effect.promise(() => + readdir(source).then( + () => false, + () => true, + ), + ), + ).toBe(true) + + yield* session.move({ sessionID: created.id, directory: destination }) + yield* execution.awaitIdle(created.id) + + expect((yield* session.get(created.id)).location.directory).toBe(destination) + expect(yield* session.inbox(created.id)).toEqual([]) + }), + { timeout: 15_000 }, + ) + + for (const broken of [false, true]) { + itWithInstance.live( + `uses the ${broken ? "broken" : "healthy discovery-disabled"} selected instance rather than the default Location`, + () => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const source = Location.Ref.make({ directory: AbsolutePath.make(path.join(tmp.path, "source")) }) + const destination = AbsolutePath.make(path.join(tmp.path, "destination")) + yield* Effect.promise(() => Promise.all([mkdir(source.directory), mkdir(destination)])) + const config = JSON.stringify({ instructions: ["{file:./missing.txt}"] }) + if (!broken) yield* Effect.promise(() => Bun.write(path.join(source.directory, "opencode.json"), config)) + const selectedID = Session.ID.create() + const replacements: LayerNode.Replacements = [ + Global.node.replace(tempGlobalLayer), + SessionExecution.node.replace(SessionExecution.noopLayer), + offlineModels, + Instance.node.replace( + makeGlobalNode({ + service: Instance.Service, + deps: [LocationServiceMap.node], + layer: Layer.effect( + Instance.Service, + Effect.gen(function* () { + const locations = yield* LocationServiceMap.Service + const privateInstances = yield* LayerMap.make( + () => + Instance.layer(source, { + discovery: false, + replacements: [ + ...bindings, + ...(broken + ? [ + Config.node.replace( + Config.configured({ project: false, global: false, content: config }), + ), + ] + : []), + ], + }), + { idleTimeToLive: Duration.infinity }, + ) + const selector = Instance.Service.of({ + provide: (session) => + Effect.provide( + session.id === selectedID && session.location.directory === source.directory + ? privateInstances.get(session.id) + : locations.get(session.location), + ), + }) + const bindings: LayerNode.Replacements = [ + ...replacements, + Instance.node.replace(Layer.succeed(Instance.Service, selector)), + LocationServiceMap.node.replace(Layer.succeed(LocationServiceMap.Service, locations)), + ] + return selector + }), + ), + }), + ), + ] + const context = yield* Layer.build(AppNodeBuilder.build(Session.node, replacements)) + const session = Context.get(context, Session.Service) + const created = yield* session.create({ id: selectedID, location: source }) + const pending = yield* session.synthetic({ + sessionID: created.id, + text: "Keep pending", + delivery: "queue", + resume: false, + }) + + yield* session.move({ sessionID: created.id, directory: destination, delivery: "queue" }) + + expect((yield* session.get(created.id)).location.directory).toBe(broken ? destination : source.directory) + const inbox = yield* session.inbox(created.id) + expect(inbox[0]).toEqual(pending) + if (broken) expect(inbox).toEqual([pending]) + if (!broken) expect(inbox.slice(1)).toMatchObject([{ type: "move", delivery: "queue" }]) + }), + { timeout: 15_000 }, + ) + } + + for (const interrupt of ["source", "caller"] as const) { + itWithInstance.live(`does not recover or enqueue a move when the ${interrupt} interrupts the probe`, () => + Effect.gen(function* () { + const fixture = yield* sourceProbe() + const created = yield* fixture.session.create({ location: Location.Ref.make({ directory: fixture.source }) }) + const pending = yield* fixture.session.synthetic({ sessionID: created.id, text: "Keep pending", resume: false }) + const moving = yield* fixture.session + .move({ sessionID: created.id, directory: fixture.destination }) + .pipe(Effect.forkScoped) + const release = yield* Queue.take(fixture.probes) + + if (interrupt === "source") yield* Deferred.interrupt(release) + if (interrupt === "caller") yield* Fiber.interrupt(moving) + const exit = yield* Fiber.await(moving) + + expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true) + expect((yield* fixture.session.get(created.id)).location.directory).toBe(fixture.source) + expect(yield* fixture.session.inbox(created.id)).toEqual([pending]) + }).pipe(Effect.timeout("5 seconds")), + ) + } + + itWithInstance.live("does not recover if execution starts during the source probe", () => + Effect.gen(function* () { + const fixture = yield* sourceProbe({ execution: true }) + const created = yield* fixture.session.create({ location: Location.Ref.make({ directory: fixture.source }) }) + const moving = yield* fixture.session + .move({ sessionID: created.id, directory: fixture.destination }) + .pipe(Effect.forkScoped) + const release = yield* Queue.take(fixture.probes) + + yield* fixture.execution.wake(created.id) + // The real coordinator now owns execution; its separate instance acquisition stays suspended. + yield* Queue.take(fixture.probes) + expect(yield* fixture.execution.isActive(created.id)).toBe(true) + yield* Deferred.die(release, new Error("source unavailable")) + yield* Fiber.join(moving) + + expect((yield* fixture.session.get(created.id)).location.directory).toBe(fixture.source) + expect(yield* fixture.session.inbox(created.id)).toMatchObject([{ type: "move", delivery: "steer" }]) + expect(yield* fixture.execution.isActive(created.id)).toBe(true) + yield* fixture.execution.interrupt(created.id) + yield* fixture.execution.awaitIdle(created.id) + }).pipe(Effect.timeout("5 seconds")), + ) + + itWithInstance.live( + "recovers a missing source without initializing its instance and retains destination workspace identity", + () => + Effect.gen(function* () { + const fixture = yield* sourceProbe() + const created = yield* fixture.session.create({ + location: Location.Ref.make({ directory: fixture.source, workspaceID: Workspace.ID.create() }), + }) + yield* Effect.promise(() => rm(fixture.source, { recursive: true })) + const workspaceID = Workspace.ID.create() + + yield* fixture.session.move({ sessionID: created.id, directory: fixture.destination, workspaceID }) + + expect((yield* fixture.session.get(created.id)).location).toEqual( + Location.Ref.make({ directory: fixture.destination, workspaceID }), + ) + expect(yield* fixture.session.inbox(created.id)).toEqual([]) + expect(yield* Queue.size(fixture.probes)).toBe(0) + }).pipe(Effect.timeout("5 seconds")), + ) + + for (const changed of ["directory", "workspace"] as const) { + itWithInstance.live(`allows inbox cancellation during a source probe and rejects stale ${changed} recovery`, () => + Effect.gen(function* () { + const fixture = yield* sourceProbe() + const created = yield* fixture.session.create({ location: Location.Ref.make({ directory: fixture.source }) }) + const pending = yield* fixture.session.synthetic({ + sessionID: created.id, + text: "Cancel pending", + resume: false, + }) + const moving = yield* fixture.session + .move({ sessionID: created.id, directory: fixture.destination }) + .pipe(Effect.exit, Effect.forkScoped) + const release = yield* Queue.take(fixture.probes) + + yield* fixture.session.cancelInbox({ sessionID: created.id, inboxID: pending.id }).pipe( + Effect.timeout("2 seconds"), + Effect.onError(() => Deferred.interrupt(release)), + ) + expect(yield* fixture.session.inbox(created.id)).toEqual([]) + expect(moving.pollUnsafe()).toBeUndefined() + const location = Location.Ref.make({ + directory: changed === "directory" ? fixture.destination : fixture.source, + workspaceID: changed === "workspace" ? Workspace.ID.create() : undefined, + }) + yield* fixture.bus.publish(SessionEvent.Moved, { + sessionID: created.id, + location, + projectID: created.projectID, + }) + yield* Deferred.die(release, new Error("source unavailable")) + expect(Exit.isSuccess(yield* Fiber.join(moving))).toBe(true) + + expect((yield* fixture.session.get(created.id)).location).toEqual(location) + expect(yield* fixture.session.inbox(created.id)).toMatchObject([ + { type: "move", payload: { location: { directory: fixture.destination } } }, + ]) + }).pipe(Effect.timeout("5 seconds")), + ) + } + itWithUnavailableDestination.effect("rejects an unavailable destination before admitting the move", () => tmpdirScoped().pipe( Effect.flatMap((tmp) => @@ -111,15 +421,23 @@ describe("Session.move", () => { yield* session.move({ sessionID: created.id, directory: destination }) expect((yield* session.get(created.id)).location.directory).toBe(AbsolutePath.make(source)) expect(yield* session.inbox(created.id)).toHaveLength(1) + const pending = yield* session.synthetic({ + sessionID: created.id, + text: "Keep queued", + delivery: "queue", + resume: false, + }) + yield* session.move({ sessionID: created.id, directory: destination, delivery: "queue" }) + expect(yield* session.inbox(created.id)).toHaveLength(3) yield* Effect.promise(() => rm(source, { recursive: true })) yield* session.move({ sessionID: created.id, directory: destination }) expect((yield* session.get(created.id)).location.directory).toBe(destination) - expect(yield* session.inbox(created.id)).toEqual([]) + expect(yield* session.inbox(created.id)).toEqual([pending]) yield* session.move({ sessionID: created.id, directory: destination }) - expect(yield* session.inbox(created.id)).toHaveLength(1) + expect(yield* session.inbox(created.id)).toHaveLength(2) yield* Effect.promise(() => mkdir(path.join(tmp.path, "other"))) const steered = yield* session.create({ From 930f56bbd106cd8d909578dafcf474c7108cef2f Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 3 Sep 2026 12:15:21 -0400 Subject: [PATCH 2/3] refactor(core): bind move preparation in a service --- packages/core/src/session.ts | 15 +--- packages/core/src/session/move.ts | 91 ++++++++++++++++--------- packages/core/test/session-move.test.ts | 47 +++++++++++++ 3 files changed, 108 insertions(+), 45 deletions(-) diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index 1b85d435a977..f17ddd5d6950 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -42,7 +42,6 @@ import { } from "./session/error.js" import { Node } from "@opencode-ai/util/effect/app-node" import { LayerNode } from "@opencode-ai/util/effect/layer-node" -import { LocationServiceMap } from "./location-service-map.js" import { SessionEvent } from "./session/event.js" import { SessionInbox } from "./session/inbox.js" import { InstructionState } from "./session/instruction-state.js" @@ -62,7 +61,6 @@ import { FSUtil } from "@opencode-ai/util/fs-util" import type { EventLog } from "@opencode-ai/schema/event-log" import { Job } from "./job.js" import type { Command } from "./command.js" -import { Global } from "@opencode-ai/util/global" import { SessionEnvironment } from "./session/environment.js" import { InstructionEntry } from "./session/instruction-entry.js" @@ -232,13 +230,12 @@ const layer = Layer.effect( const db = database.db const bus = yield* Bus.Service const projects = yield* Project.Service - const global = yield* Global.Service const execution = yield* SessionExecution.Service const llm = yield* LLMClient.Service const transport = yield* SessionModelTransport.Service const store = yield* SessionStore.Service const instances = yield* Instance.Service - const locations = yield* LocationServiceMap.Service + const moves = yield* SessionMove.Service const fs = yield* FSUtil.Service const jobs = yield* Job.Service const environments = yield* SessionEnvironment.Service @@ -412,12 +409,7 @@ const layer = Layer.effect( rename: (input) => sessions.forSession(input.sessionID).rename(input), move: Effect.fn("Session.move")(function* (input) { const session = yield* result.get(input.sessionID) - const payload = yield* SessionMove.prepare({ ...input, session }).pipe( - Effect.provideService(FSUtil.Service, fs), - Effect.provideService(Global.Service, global), - Effect.provideService(Project.Service, projects), - Effect.provideService(LocationServiceMap.Service, locations), - ) + const payload = yield* moves.prepare({ ...input, session }) // Probe the same instance execution would use, without holding up inbox cancellation. const unavailable = !(yield* execution.isActive(input.sessionID)) && @@ -514,10 +506,9 @@ export const node: LayerNode.Provider undefined)) - if (!info) return yield* new DestinationNotFoundError({ directory }) - if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory }) - const project = yield* projects.resolve(directory) - const payload: SessionInbox.MovePayload = { - location: Location.Ref.make({ directory, workspaceID: input.workspaceID }), - projectID: project.id, - subpath: RelativePath.make(path.relative(project.directory, directory).replaceAll("\\", "/")), - } - yield* Location.Service.pipe( - Effect.provide(locations.get(payload.location)), - Effect.scoped, - Effect.catchCause((cause) => { - if (Cause.hasInterruptsOnly(cause)) return Effect.failCause(cause) - return Effect.logWarning("session move destination unavailable", { directory, cause }).pipe( - Effect.andThen(Effect.fail(new DestinationUnavailableError({ directory }))), - ) - }), - ) - return payload +export interface Interface { + readonly prepare: (input: { + session: Session.Info + directory: AbsolutePath + workspaceID?: Location.Ref["workspaceID"] + }) => Effect.Effect< + SessionInbox.MovePayload, + DestinationNotFoundError | DestinationNotDirectoryError | DestinationUnavailableError + > +} + +export class Service extends Context.Service()("@opencode/SessionMove") {} + +const layer = Layer.effect( + Service, + Effect.gen(function* () { + const fs = yield* FSUtil.Service + const global = yield* Global.Service + const projects = yield* Project.Service + const locations = yield* LocationServiceMap.Service + + return Service.of({ + prepare: Effect.fn("SessionMove.prepare")(function* (input) { + const value = input.directory.trim() + const expanded = + value === "~" ? global.home : value.startsWith("~/") ? path.join(global.home, value.slice(2)) : value + const directory = AbsolutePath.make(path.resolve(input.session.location.directory, expanded)) + const info = yield* fs.stat(directory).pipe(Effect.orElseSucceed(() => undefined)) + if (!info) return yield* new DestinationNotFoundError({ directory }) + if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory }) + const project = yield* projects.resolve(directory) + const payload: SessionInbox.MovePayload = { + location: Location.Ref.make({ directory, workspaceID: input.workspaceID }), + projectID: project.id, + subpath: RelativePath.make(path.relative(project.directory, directory).replaceAll("\\", "/")), + } + yield* Location.Service.pipe( + Effect.provide(locations.get(payload.location)), + Effect.scoped, + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) return Effect.failCause(cause) + return Effect.logWarning("session move destination unavailable", { directory, cause }).pipe( + Effect.andThen(Effect.fail(new DestinationUnavailableError({ directory }))), + ) + }), + ) + return payload + }), + }) + }), +) + +export const node = makeGlobalNode({ + service: Service, + layer, + deps: [FSUtil.node, Global.node, Project.node, LocationServiceMap.node], }) diff --git a/packages/core/test/session-move.test.ts b/packages/core/test/session-move.test.ts index ab415fc856d1..5b1bb34b8185 100644 --- a/packages/core/test/session-move.test.ts +++ b/packages/core/test/session-move.test.ts @@ -17,6 +17,7 @@ import { AbsolutePath } from "@opencode-ai/core/schema" import { Session } from "@opencode-ai/core/session" import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionExecution } from "@opencode-ai/core/session/execution" +import { SessionMove } from "@opencode-ai/core/session/move" import { SessionProjector } from "@opencode-ai/core/session/projector" import { SessionRunner } from "@opencode-ai/core/session/runner/index" import { SessionStore } from "@opencode-ai/core/session/store" @@ -138,6 +139,52 @@ const sourceProbe = (options: { execution?: boolean } = {}) => }) describe("Session.move", () => { + itWithInstance.live("binds move preparation dependencies in the service implementation", () => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const directory = AbsolutePath.make(tmp.path) + const context = yield* Layer.build( + AppNodeBuilder.build(LayerNode.group([Session.node, SessionMove.node]), [ + Global.node.replace(tempGlobalLayer), + Project.node.replace(globalProjectNode), + SessionExecution.node.replace(SessionExecution.noopLayer), + offlineModels, + ]), + ) + const sessions = Context.get(context, Session.Service) + const moves = Context.get(context, SessionMove.Service) + const session = yield* sessions.create({ location: Location.Ref.make({ directory }) }) + + // Call the bound service outside its construction context: no leaf services are provided here. + expect(yield* moves.prepare({ session, directory })).toMatchObject({ + location: { directory }, + projectID: Project.ID.global, + }) + }), + ) + + itWithInstance.live("uses the provided move preparation service", () => + Effect.gen(function* () { + const tmp = yield* tmpdirScoped() + const directory = AbsolutePath.make(tmp.path) + const rejection = new Session.DestinationUnavailableError({ directory }) + const context = yield* Layer.build( + AppNodeBuilder.build(Session.node, [ + Global.node.replace(tempGlobalLayer), + Project.node.replace(globalProjectNode), + SessionExecution.node.replace(SessionExecution.noopLayer), + SessionMove.node.replace(Layer.succeed(SessionMove.Service, { prepare: () => Effect.fail(rejection) })), + offlineModels, + ]), + ) + const sessions = Context.get(context, Session.Service) + const created = yield* sessions.create({ location: Location.Ref.make({ directory }) }) + + expect(yield* sessions.move({ sessionID: created.id, directory }).pipe(Effect.flip)).toBe(rejection) + expect(yield* sessions.inbox(created.id)).toEqual([]) + }), + ) + for (const broken of [false, true]) { itWithExecution.live( `moves an idle session from ${broken ? "broken" : "healthy"} source configuration`, From c08dfb71f743879f942ed28f08a887466c67a188 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 3 Sep 2026 12:27:14 -0400 Subject: [PATCH 3/3] refactor(core): give move service full ownership --- packages/core/src/session.ts | 66 +---------- packages/core/src/session/move.ts | 146 +++++++++++++++++++----- packages/core/test/session-move.test.ts | 32 ++++-- 3 files changed, 144 insertions(+), 100 deletions(-) diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index f17ddd5d6950..3e8fdb577779 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -1,7 +1,7 @@ export * as Session from "./session.js" export * from "./session/schema.js" -import { Cause, Effect, Layer, Schema, Context, Stream } from "effect" +import { Effect, Layer, Schema, Context, Stream } from "effect" import { LLMClient } from "@opencode-ai/ai" import { ListAnchor } from "@opencode-ai/schema/session" import { and, desc, eq } from "drizzle-orm" @@ -16,7 +16,7 @@ import { Database } from "./database/database.js" import { SessionProjector } from "./session/projector.js" import { SessionMessageTable } from "./session/sql.js" import { SessionSchema } from "./session/schema.js" -import { AbsolutePath, RelativePath } from "./schema.js" +import { RelativePath } from "./schema.js" import { Agent } from "@opencode-ai/schema/agent" import { App } from "./app.js" import { Slug } from "./util/slug.js" @@ -166,15 +166,7 @@ export interface Interface { readonly switchAgent: (input: { sessionID: SessionSchema.ID; agent: Agent.ID }) => Effect.Effect readonly switchModel: (input: { sessionID: SessionSchema.ID; model: Model.Ref }) => Effect.Effect readonly rename: (input: { sessionID: SessionSchema.ID; title: string }) => Effect.Effect - readonly move: (input: { - sessionID: SessionSchema.ID - directory: AbsolutePath - workspaceID?: Location.Ref["workspaceID"] - delivery?: SessionInbox.Delivery - }) => Effect.Effect< - void, - NotFoundError | DestinationNotFoundError | DestinationNotDirectoryError | DestinationUnavailableError - > + readonly move: SessionMove.Interface["move"] readonly prompt: ( input: Parameters[0] & { sessionID: SessionSchema.ID }, ) => ReturnType @@ -236,11 +228,9 @@ const layer = Layer.effect( const store = yield* SessionStore.Service const instances = yield* Instance.Service const moves = yield* SessionMove.Service - const fs = yield* FSUtil.Service const jobs = yield* Job.Service const environments = yield* SessionEnvironment.Service const sessions = yield* Session.make() - const admission = yield* SessionInbox.Service const isDurableSessionEvent = Schema.is(SessionEvent.Durable) const result = Service.of({ @@ -407,55 +397,7 @@ const layer = Layer.effect( switchAgent: (input) => sessions.forSession(input.sessionID).switchAgent(input), switchModel: (input) => sessions.forSession(input.sessionID).switchModel(input), rename: (input) => sessions.forSession(input.sessionID).rename(input), - move: Effect.fn("Session.move")(function* (input) { - const session = yield* result.get(input.sessionID) - const payload = yield* moves.prepare({ ...input, session }) - // Probe the same instance execution would use, without holding up inbox cancellation. - const unavailable = - !(yield* execution.isActive(input.sessionID)) && - (!(yield* fs.isDir(session.location.directory)) || - !(yield* SessionRunner.Service.pipe( - instances.provide(session), - Effect.as(true), - Effect.catchCause((cause) => - Cause.hasInterrupts(cause) ? Effect.failCause(cause) : Effect.succeed(false), - ), - ))) - const item = SessionInbox.Item.make({ - type: "move", - payload, - delivery: input.delivery ?? "steer", - }) - yield* SessionInbox.serialized( - input.sessionID, - Effect.gen(function* () { - const latest = yield* result.get(input.sessionID) - // Active runners must hand off at a step boundary to retain their continuation. - if ( - unavailable && - latest.location.directory === session.location.directory && - latest.location.workspaceID === session.location.workspaceID && - !(yield* execution.isActive(input.sessionID)) - ) { - const cancellations = (yield* SessionInbox.moveIDs(db, input.sessionID)).map( - (item) => [SessionEvent.InboxCancelled, { sessionID: input.sessionID, inboxID: item.id }] as const, - ) - const moved = [SessionEvent.Moved, { sessionID: input.sessionID, ...payload }] as const - const first = cancellations[0] - if (!first) return yield* bus.publish(...moved).pipe(Effect.asVoid) - return yield* bus.publishAll([first, ...cancellations.slice(1), moved]) - } - yield* admission - .admit({ - id: SessionMessage.ID.create(), - sessionID: input.sessionID, - item, - }) - .pipe(Effect.orDie) - }), - ) - yield* execution.wake(input.sessionID) - }), + move: moves.move, compact: (input) => sessions.forSession(input.sessionID).compact(input), wait: (sessionID) => sessions.forSession(sessionID).wait(), active: execution.active, diff --git a/packages/core/src/session/move.ts b/packages/core/src/session/move.ts index 40d0bff9ebc4..1a66df00cba6 100644 --- a/packages/core/src/session/move.ts +++ b/packages/core/src/session/move.ts @@ -1,16 +1,26 @@ export * as SessionMove from "./move.js" import type { Session } from "@opencode-ai/schema/session" -import type { SessionInbox } from "@opencode-ai/schema/session-inbox" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" import { FSUtil } from "@opencode-ai/util/fs-util" import { Global } from "@opencode-ai/util/global" import { Cause, Context, Effect, Layer, Schema } from "effect" import path from "path" +import { Bus } from "../bus.js" +import { Database } from "../database/database.js" +import { Instance } from "../instance/service.js" import { Location } from "../location.js" import { LocationServiceMap } from "../location-service-map.js" import { Project } from "../project.js" import { AbsolutePath, RelativePath } from "../schema.js" +import { NotFoundError } from "./error.js" +import { SessionEvent } from "./event.js" +import { SessionExecution } from "./execution.js" +import { SessionInbox } from "./inbox.js" +import { SessionMessage } from "./message.js" +import { SessionProjector } from "./projector.js" +import { SessionRunner } from "./runner/index.js" +import { SessionStore } from "./store.js" export class DestinationNotFoundError extends Schema.TaggedError()( "Session.DestinationNotFoundError", @@ -28,13 +38,14 @@ export class DestinationUnavailableError extends Schema.TaggedError Effect.Effect< - SessionInbox.MovePayload, - DestinationNotFoundError | DestinationNotDirectoryError | DestinationUnavailableError + void, + NotFoundError | DestinationNotFoundError | DestinationNotDirectoryError | DestinationUnavailableError > } @@ -47,33 +58,98 @@ const layer = Layer.effect( const global = yield* Global.Service const projects = yield* Project.Service const locations = yield* LocationServiceMap.Service + const store = yield* SessionStore.Service + const execution = yield* SessionExecution.Service + const instances = yield* Instance.Service + const admission = yield* SessionInbox.Service + const database = yield* Database.Service + const bus = yield* Bus.Service + + const get = Effect.fn("SessionMove.get")(function* (sessionID: Session.ID) { + const session = yield* store.get(sessionID) + if (!session) return yield* new NotFoundError({ sessionID }) + return session + }) + + const resolveDestination = Effect.fn("SessionMove.resolveDestination")(function* ( + session: Session.Info, + input: Parameters[0], + ) { + const value = input.directory.trim() + const expanded = + value === "~" ? global.home : value.startsWith("~/") ? path.join(global.home, value.slice(2)) : value + const directory = AbsolutePath.make(path.resolve(session.location.directory, expanded)) + const info = yield* fs.stat(directory).pipe(Effect.orElseSucceed(() => undefined)) + if (!info) return yield* new DestinationNotFoundError({ directory }) + if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory }) + const project = yield* projects.resolve(directory) + const destination: SessionInbox.MovePayload = { + location: Location.Ref.make({ directory, workspaceID: input.workspaceID }), + projectID: project.id, + subpath: RelativePath.make(path.relative(project.directory, directory).replaceAll("\\", "/")), + } + yield* locations.contextEffect(destination.location).pipe( + Effect.scoped, + Effect.catchCause((cause) => { + if (Cause.hasInterruptsOnly(cause)) return Effect.failCause(cause) + return Effect.logWarning("session move destination unavailable", { directory, cause }).pipe( + Effect.andThen(Effect.fail(new DestinationUnavailableError({ directory }))), + ) + }), + ) + return destination + }) + + const sourceUnavailable = Effect.fn("SessionMove.sourceUnavailable")(function* (session: Session.Info) { + if (yield* execution.isActive(session.id)) return false + if (!(yield* fs.isDir(session.location.directory))) return true + return yield* SessionRunner.Service.pipe( + instances.provide(session), + Effect.as(false), + Effect.catchCause((cause) => (Cause.hasInterrupts(cause) ? Effect.failCause(cause) : Effect.succeed(true))), + ) + }) return Service.of({ - prepare: Effect.fn("SessionMove.prepare")(function* (input) { - const value = input.directory.trim() - const expanded = - value === "~" ? global.home : value.startsWith("~/") ? path.join(global.home, value.slice(2)) : value - const directory = AbsolutePath.make(path.resolve(input.session.location.directory, expanded)) - const info = yield* fs.stat(directory).pipe(Effect.orElseSucceed(() => undefined)) - if (!info) return yield* new DestinationNotFoundError({ directory }) - if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory }) - const project = yield* projects.resolve(directory) - const payload: SessionInbox.MovePayload = { - location: Location.Ref.make({ directory, workspaceID: input.workspaceID }), - projectID: project.id, - subpath: RelativePath.make(path.relative(project.directory, directory).replaceAll("\\", "/")), - } - yield* Location.Service.pipe( - Effect.provide(locations.get(payload.location)), - Effect.scoped, - Effect.catchCause((cause) => { - if (Cause.hasInterruptsOnly(cause)) return Effect.failCause(cause) - return Effect.logWarning("session move destination unavailable", { directory, cause }).pipe( - Effect.andThen(Effect.fail(new DestinationUnavailableError({ directory }))), - ) + move: Effect.fn("SessionMove.move")(function* (input) { + const session = yield* get(input.sessionID) + const destination = yield* resolveDestination(session, input) + // Probe outside the inbox lock so cancellation remains available during initialization. + const unavailable = yield* sourceUnavailable(session) + const item = SessionInbox.Item.make({ + type: "move", + payload: destination, + delivery: input.delivery ?? "steer", + }) + yield* SessionInbox.serialized( + input.sessionID, + Effect.gen(function* () { + const latest = yield* get(input.sessionID) + // Only recover the placement we probed; active runners retain their step-boundary handoff. + if ( + unavailable && + latest.location.directory === session.location.directory && + latest.location.workspaceID === session.location.workspaceID && + !(yield* execution.isActive(input.sessionID)) + ) { + const cancellations = (yield* SessionInbox.moveIDs(database.db, input.sessionID)).map( + (item) => [SessionEvent.InboxCancelled, { sessionID: input.sessionID, inboxID: item.id }] as const, + ) + const moved = [SessionEvent.Moved, { sessionID: input.sessionID, ...destination }] as const + const first = cancellations[0] + if (!first) return yield* bus.publish(...moved).pipe(Effect.asVoid) + return yield* bus.publishAll([first, ...cancellations.slice(1), moved]) + } + yield* admission + .admit({ + id: SessionMessage.ID.create(), + sessionID: input.sessionID, + item, + }) + .pipe(Effect.orDie) }), ) - return payload + yield* execution.wake(input.sessionID) }), }) }), @@ -82,5 +158,17 @@ const layer = Layer.effect( export const node = makeGlobalNode({ service: Service, layer, - deps: [FSUtil.node, Global.node, Project.node, LocationServiceMap.node], + deps: [ + FSUtil.node, + Global.node, + Project.node, + LocationServiceMap.node, + SessionStore.node, + SessionExecution.node, + Instance.node, + SessionInbox.node, + Database.node, + Bus.node, + SessionProjector.node, + ], }) diff --git a/packages/core/test/session-move.test.ts b/packages/core/test/session-move.test.ts index 5b1bb34b8185..e235206f6ab6 100644 --- a/packages/core/test/session-move.test.ts +++ b/packages/core/test/session-move.test.ts @@ -139,31 +139,45 @@ const sourceProbe = (options: { execution?: boolean } = {}) => }) describe("Session.move", () => { - itWithInstance.live("binds move preparation dependencies in the service implementation", () => + itWithInstance.live("moves through the bound service without depending on the Session facade", () => Effect.gen(function* () { const tmp = yield* tmpdirScoped() const directory = AbsolutePath.make(tmp.path) const context = yield* Layer.build( - AppNodeBuilder.build(LayerNode.group([Session.node, SessionMove.node]), [ + AppNodeBuilder.build(LayerNode.group([SessionMove.node, SessionStore.node, Bus.node, Project.node]), [ Global.node.replace(tempGlobalLayer), Project.node.replace(globalProjectNode), SessionExecution.node.replace(SessionExecution.noopLayer), offlineModels, ]), ) - const sessions = Context.get(context, Session.Service) const moves = Context.get(context, SessionMove.Service) - const session = yield* sessions.create({ location: Location.Ref.make({ directory }) }) - - // Call the bound service outside its construction context: no leaf services are provided here. - expect(yield* moves.prepare({ session, directory })).toMatchObject({ + const store = Context.get(context, SessionStore.Service) + const bus = Context.get(context, Bus.Service) + const projects = Context.get(context, Project.Service) + const sessionID = Session.ID.create() + + // Call outside the construction context: the service owns all of its dependencies. + expect(yield* moves.move({ sessionID, directory }).pipe(Effect.flip)).toEqual( + new Session.NotFoundError({ sessionID }), + ) + yield* projects.resolve(directory) + yield* bus.publish(SessionEvent.Created, { + sessionID, + slug: "move-service", + version: "test", + projectID: Project.ID.global, + location: Location.Ref.make({ directory: AbsolutePath.make(path.join(tmp.path, "missing")) }), + }) + yield* moves.move({ sessionID, directory }) + expect(yield* store.get(sessionID)).toMatchObject({ location: { directory }, projectID: Project.ID.global, }) }), ) - itWithInstance.live("uses the provided move preparation service", () => + itWithInstance.live("delegates to the provided move service", () => Effect.gen(function* () { const tmp = yield* tmpdirScoped() const directory = AbsolutePath.make(tmp.path) @@ -173,7 +187,7 @@ describe("Session.move", () => { Global.node.replace(tempGlobalLayer), Project.node.replace(globalProjectNode), SessionExecution.node.replace(SessionExecution.noopLayer), - SessionMove.node.replace(Layer.succeed(SessionMove.Service, { prepare: () => Effect.fail(rejection) })), + SessionMove.node.replace(Layer.succeed(SessionMove.Service, { move: () => Effect.fail(rejection) })), offlineModels, ]), )