Compare commits

...
Author SHA1 Message Date
Kit Langton b6db6631d6 chore(core): merge v2 into source recovery
Integrate opaque layer graphs and session-aware instance selection while preserving idle-source recovery and active handoff guards. Migrate recovery fixtures and the base execution and instance-selection fixtures to checked replacements.
2026-08-31 14:09:45 -04:00
Kit Langton 6fb36773a5 refactor(core): reuse directory checks for move recovery 2026-08-31 13:30:51 -04:00
Kit Langton 3b8283789d fix(core): recover moves from unavailable source locations 2026-08-31 12:50:51 -04:00
4 changed files with 174 additions and 17 deletions
+16 -2
View File
@@ -455,6 +455,16 @@ const layer = Layer.effect(
)
}),
)
const unavailable =
!(yield* execution.isActive(input.sessionID)) &&
(!(yield* fs.isDir(current.location.directory)) ||
!(yield* locations.contextEffect(current.location).pipe(
Effect.scoped,
Effect.as(true),
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause) ? Effect.failCause(cause) : Effect.succeed(false),
),
)))
const item = SessionInbox.Item.make({
type: "move",
payload,
@@ -464,9 +474,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 === current.location.directory &&
latest.location.workspaceID === current.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,
)
+3 -1
View File
@@ -1374,7 +1374,9 @@ function buildExecution(
Layer.provide(Layer.succeed(Job.Service, jobs)),
// Do not reuse the outer harness's selector with its already-captured Location map.
Layer.provide(
LayerNode.compile(Instance.byLocationNode, [[LocationServiceMap.node, locations]]).pipe(Layer.fresh),
LayerNode.compile(Instance.byLocationNode, {
replacements: [LocationServiceMap.node.replace(locations)],
}).pipe(Layer.fresh),
),
),
scope,
+141 -1
View File
@@ -1,11 +1,13 @@
import { describe, expect } from "bun:test"
import path from "path"
import { mkdir, rm } from "fs/promises"
import { Effect, Layer, LayerMap } from "effect"
import { Cause, Context, Deferred, Duration, Effect, Exit, Fiber, Layer, LayerMap } from "effect"
import { Worktree } from "@opencode-ai/schema/worktree"
import { Workspace } from "@opencode-ai/schema/workspace"
import { Bus } from "@opencode-ai/core/bus"
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 +20,8 @@ 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 { Global } from "@opencode-ai/util/global"
import { tempGlobalLayer } from "./fixture/global"
import { tmpdirScoped } from "./fixture/tmpdir"
import { testEffect } from "./lib/effect"
import { globalProjectNode } from "./lib/project"
@@ -57,6 +61,9 @@ const itWithActiveExecution = testEffect(
],
),
)
const itWithExecution = testEffect(
AppNodeBuilder.build(LayerNode.group([Session.node, SessionExecution.node]), [Global.node.replace(tempGlobalLayer)]),
)
const unavailableLocations = Layer.effect(
LocationServiceMap.Service,
LayerMap.make(
@@ -73,8 +80,141 @@ const itWithUnavailableDestination = testEffect(
],
),
)
const itWithSourceProbe = testEffect(Layer.empty)
const sourceProbe = 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 entered = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
const replacements: LayerNode.Replacements = [
Project.node.replace(globalProjectNode),
SessionExecution.node.replace(SessionExecution.noopLayer),
Global.node.replace(tempGlobalLayer),
]
const context = yield* Layer.build(
AppNodeBuilder.build(LayerNode.group([Session.node, Bus.node]), [
...replacements,
LocationServiceMap.node.replace(
Layer.effect(
LocationServiceMap.Service,
LayerMap.make(
(ref: Location.Ref) =>
Layer.unwrap(
Effect.gen(function* () {
if (ref.directory === source) {
yield* Deferred.succeed(entered, undefined)
yield* Deferred.await(release)
}
return Instance.layer(ref, { replacements })
}),
),
{ idleTimeToLive: Duration.infinity },
),
),
),
]),
)
return {
source,
destination,
entered,
release,
session: Context.get(context, Session.Service),
bus: Context.get(context, Bus.Service),
}
})
describe("Session.move", () => {
itWithExecution.live(
"recovers an idle session whose source configuration cannot load",
() =>
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)]))
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 },
)
itWithSourceProbe.live("does not recover or enqueue a move when source initialization is interrupted", () =>
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.exit, Effect.forkScoped)
yield* Deferred.await(fixture.entered)
yield* Deferred.interrupt(fixture.release)
const exit = yield* Fiber.join(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")),
)
for (const changed of ["directory", "workspace"] as const) {
itWithSourceProbe.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)
yield* Deferred.await(fixture.entered)
yield* fixture.session.cancelInbox({ sessionID: created.id, inboxID: pending.id }).pipe(
Effect.timeout("2 seconds"),
Effect.onError(() => Deferred.interrupt(fixture.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,
})
// Apply a concurrent placement change while the original source probe is suspended.
yield* fixture.bus.publish(SessionEvent.Moved, {
sessionID: created.id,
location,
projectID: Project.ID.global,
})
yield* Deferred.die(fixture.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) =>
+14 -13
View File
@@ -52,18 +52,19 @@ it.live(
const cell = PluginRuntime.makeCell()
// Host and private instances must reuse the same global layer identities.
const replacements: LayerNode.Replacements = [
[Global.node, tempGlobalLayer],
[Database.node, Database.node],
[Bus.node, Bus.node],
[App.node, App.node],
[ModelsDev.node, ModelsDev.configured({ fetch: false })],
[Watcher.node, Watcher.configured({ enabled: false })],
[PluginRuntime.node, PluginRuntime.layerWithCell(cell)],
[PluginRuntime.providerNode, PluginRuntime.providerNodeWithCell(cell)],
[llmClient, Layer.succeed(LLMClient.Service, llm)],
[SessionRunnerModel.node, Layer.succeed(SessionRunnerModel.Service, { resolve: () => Effect.succeed(model) })],
[
Instance.byLocationNode,
Global.node.replace(tempGlobalLayer),
Database.node.replace(Database.node),
Bus.node.replace(Bus.node),
App.node.replace(App.node),
ModelsDev.node.replace(ModelsDev.configured({ fetch: false })),
Watcher.node.replace(Watcher.configured({ enabled: false })),
PluginRuntime.node.replace(PluginRuntime.layerWithCell(cell)),
PluginRuntime.providerNode.replace(PluginRuntime.providerNodeWithCell(cell)),
llmClient.replace(Layer.succeed(LLMClient.Service, llm)),
SessionRunnerModel.node.replace(
Layer.succeed(SessionRunnerModel.Service, { resolve: () => Effect.succeed(model) }),
),
Instance.byLocationNode.replace(
Layer.effect(
Instance.Service,
Effect.gen(function* () {
@@ -151,7 +152,7 @@ it.live(
})
}),
),
],
),
]
const context = yield* Layer.build(
createEmbeddedRoutes({}, replacements).pipe(Layer.provide(HttpServer.layerServices)),