Compare commits

...
1 Commits
Author SHA1 Message Date
Kit Langton da59dcf592 fix(server): decouple session controls from Location startup 2026-08-26 13:26:20 -04:00
8 changed files with 275 additions and 30 deletions
+6
View File
@@ -0,0 +1,6 @@
---
"@opencode-ai/protocol": patch
"@opencode-ai/server": patch
---
Allow session wait and interrupt requests to reach process-local execution without booting the session's Location services. Preserve session validation errors, immediate interrupt acceptance, and waiting for execution cleanup to finish.
+3 -2
View File
@@ -30,6 +30,7 @@ import { Model } from "@opencode-ai/schema/model"
import { Location } from "@opencode-ai/schema/location"
import { SessionEvent } from "@opencode-ai/schema/session-event"
import { EventLog } from "@opencode-ai/schema/event-log"
import { SessionValidationMiddleware } from "../middleware/session-validation.js"
const ParentIDFilter = Schema.Union([
Session.ID,
@@ -463,7 +464,7 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
success: HttpApiSchema.NoContent,
error: [SessionNotFoundError, ServiceUnavailableError],
})
.middleware(sessionLocationMiddleware)
.middleware(SessionValidationMiddleware)
.annotateMerge(
OpenApi.annotations({
identifier: "v2.session.wait",
@@ -671,7 +672,7 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
}).annotate({ identifier: "SessionInterruptResponse" }),
error: SessionNotFoundError,
})
.middleware(sessionLocationMiddleware)
.middleware(SessionValidationMiddleware)
.annotateMerge(
OpenApi.annotations({
identifier: "v2.session.interrupt",
@@ -0,0 +1,7 @@
import { HttpApiMiddleware } from "effect/unstable/httpapi"
import { InvalidRequestError, SessionNotFoundError } from "../errors.js"
export class SessionValidationMiddleware extends HttpApiMiddleware.Service<SessionValidationMiddleware>()(
"@opencode/HttpApiSessionValidation",
{ error: [InvalidRequestError, SessionNotFoundError] },
) {}
+103
View File
@@ -444,6 +444,109 @@ it.live("embedded client exposes plugin-backed web search", () =>
),
)
for (const continuation of [false, true]) {
it.live(
`session controls bypass a cold Location during interrupt cleanup (continue=${continuation})`,
() =>
withEmbedded("opencode-embedded-session-controls-", (fixture) =>
Effect.gen(function* () {
const llm = yield* TestLLM.Service
const started = yield* Deferred.make<void>()
const cleanupStarted = yield* Deferred.make<void>()
const cleanupGate = yield* Deferred.make<void>()
const unavailable = yield* Ref.make(false)
const boots = yield* Ref.make(0)
const model = LanguageModel.make({ id: "session-controls", provider: "test", route: OpenAIChat.route })
yield* llm.push(
Stream.fromEffect(
Deferred.succeed(started, undefined).pipe(
Effect.andThen(Effect.never),
Effect.onInterrupt(() =>
Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))),
),
),
),
)
const models = Layer.effect(
SessionRunnerModel.Service,
Effect.gen(function* () {
yield* Ref.update(boots, (count) => count + 1)
if (yield* Ref.get(unavailable)) return yield* Effect.die("Location is unavailable during cleanup")
return SessionRunnerModel.Service.of({
resolve: () =>
Effect.succeed(
SessionRunnerModel.resolved(model, {
capabilities: { tools: true, input: ["text"], output: ["text"] },
cost: [],
limit: { context: 100_000, output: 1_000 },
}),
),
})
}),
)
const opencode = yield* fixture.sdk.OpenCode.create(
{
config: { directory: fixture.directory, project: false, content: "{}" },
fs: { filewatcher: false },
},
{
overrides: [
[llmClient, Layer.succeed(LLMClient.Service, llm.client)],
[SessionRunnerModel.node, models],
],
},
)
// Release blocked cleanup before the embedded host's finalizer on assertion failure.
yield* Effect.addFinalizer(() => Deferred.succeed(cleanupGate, undefined).pipe(Effect.asVoid))
const session = yield* opencode.sessions.create({ title: "Session controls", location: location(fixture) })
yield* opencode.sessions.wait({ sessionID: session.id })
expect(yield* opencode.sessions.interrupt({ sessionID: session.id })).toEqual({ interrupted: false })
expect(yield* Ref.get(boots)).toBe(0)
yield* opencode.sessions.prompt({ sessionID: session.id, text: "Start the model" })
yield* Deferred.await(started).pipe(Effect.timeout("5 seconds"))
expect(yield* Ref.get(boots)).toBe(1)
expect(llm.requests).toHaveLength(1)
const steer = yield* opencode.sessions.prompt({ sessionID: session.id, text: "Continue here", resume: false })
const queued = yield* opencode.sessions.prompt({
sessionID: session.id,
text: "Keep this queued",
delivery: "queue",
resume: false,
})
yield* Ref.set(unavailable, true)
yield* opencode.debug.location.evict({ location: location(fixture) })
expect(yield* opencode.debug.location.list()).toEqual([])
const waiting = yield* opencode.sessions.wait({ sessionID: session.id }).pipe(Effect.forkScoped)
expect(
yield* opencode.sessions
.interrupt({ sessionID: session.id, continue: continuation })
.pipe(Effect.timeout("2 seconds")),
).toEqual({ interrupted: true })
yield* Deferred.await(cleanupStarted).pipe(Effect.timeout("2 seconds"))
expect(waiting.pollUnsafe()).toBeUndefined()
expect(yield* opencode.sessions.active()).toEqual({ [session.id]: { type: "running" } })
expect(yield* opencode.sessions.interrupt({ sessionID: session.id })).toEqual({ interrupted: false })
expect(yield* Ref.get(boots)).toBe(1)
expect(yield* opencode.debug.location.list()).toEqual([])
// Only real continuation may acquire a fresh graph, after the interrupted drain settles.
yield* Ref.set(unavailable, false)
yield* Deferred.succeed(cleanupGate, undefined)
yield* Fiber.join(waiting).pipe(Effect.timeout("5 seconds"))
expect(yield* opencode.sessions.active()).toEqual({})
expect(yield* Ref.get(boots)).toBe(continuation ? 2 : 1)
expect(llm.requests).toHaveLength(continuation ? 2 : 1)
expect((yield* opencode.sessions.inbox.list({ sessionID: session.id })).map((item) => item.id)).toEqual(
continuation ? [queued.id] : [steer.id, queued.id],
)
}),
).pipe(Effect.provide(TestLLM.layer({ fallback: TestLLM.text("Finished", "answer") }))),
15_000,
)
}
it.live(
"Location-owned runner events reach the ready global client",
() =>
@@ -2,15 +2,12 @@ import { Database } from "@opencode-ai/core/database/database"
import { LocationServiceMap } from "@opencode-ai/core/location-services"
import { Location } from "@opencode-ai/core/location"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { Session } from "@opencode-ai/core/session"
import { SessionTable } from "@opencode-ai/core/session/sql"
import { Workspace } from "@opencode-ai/core/workspace"
import { eq } from "drizzle-orm"
import { Effect, Layer, Schema } from "effect"
import { HttpRouter } from "effect/unstable/http"
import { Effect, Layer } from "effect"
import { HttpApiMiddleware } from "effect/unstable/httpapi"
import { InvalidRequestError, SessionNotFoundError } from "@opencode-ai/protocol/errors"
import type { LocationServices } from "../location"
import { requireSession } from "./session-validation"
export class SessionLocationMiddleware extends HttpApiMiddleware.Service<
SessionLocationMiddleware,
@@ -19,8 +16,6 @@ export class SessionLocationMiddleware extends HttpApiMiddleware.Service<
error: [InvalidRequestError, SessionNotFoundError],
}) {}
const decodeSessionID = Schema.decodeUnknownEffect(Session.ID)
export const sessionLocationLayer = Layer.effect(
SessionLocationMiddleware,
Effect.gen(function* () {
@@ -29,27 +24,7 @@ export const sessionLocationLayer = Layer.effect(
return SessionLocationMiddleware.of((effect) =>
Effect.gen(function* () {
const route = yield* HttpRouter.RouteContext
const sessionID = yield* decodeSessionID(route.params.sessionID).pipe(
Effect.mapError(
() =>
new InvalidRequestError({
message: "Invalid session ID",
field: "sessionID",
}),
),
)
const row = yield* db
.select({ directory: SessionTable.directory, workspaceID: SessionTable.workspace_id })
.from(SessionTable)
.where(eq(SessionTable.id, sessionID))
.get()
.pipe(Effect.orDie)
if (!row)
return yield* new SessionNotFoundError({
sessionID,
message: `Session not found: ${sessionID}`,
})
const row = yield* requireSession(db)
return yield* effect.pipe(
Effect.provide(
@@ -0,0 +1,34 @@
import { Database } from "@opencode-ai/core/database/database"
import { Session } from "@opencode-ai/core/session"
import { SessionTable } from "@opencode-ai/core/session/sql"
import { InvalidRequestError, SessionNotFoundError } from "@opencode-ai/protocol/errors"
import { SessionValidationMiddleware } from "@opencode-ai/protocol/middleware/session-validation"
import { eq } from "drizzle-orm"
import { Effect, Layer, Schema } from "effect"
import { HttpRouter } from "effect/unstable/http"
const decodeSessionID = Schema.decodeUnknownEffect(Session.ID)
export const sessionValidationLayer = Layer.effect(
SessionValidationMiddleware,
Effect.gen(function* () {
const database = yield* Database.Service
return SessionValidationMiddleware.of((effect) => requireSession(database.db).pipe(Effect.andThen(effect)))
}),
)
// Middleware validates before query decoding, preserving the public session error precedence.
export const requireSession = Effect.fn("HttpApi.requireSession")(function* (db: Database.Interface["db"]) {
const route = yield* HttpRouter.RouteContext
const sessionID = yield* decodeSessionID(route.params.sessionID).pipe(
Effect.mapError(() => new InvalidRequestError({ message: "Invalid session ID", field: "sessionID" })),
)
const row = yield* db
.select({ directory: SessionTable.directory, workspaceID: SessionTable.workspace_id })
.from(SessionTable)
.where(eq(SessionTable.id, sessionID))
.get()
.pipe(Effect.orDie)
if (!row) return yield* new SessionNotFoundError({ sessionID, message: `Session not found: ${sessionID}` })
return row
})
+2
View File
@@ -42,6 +42,7 @@ import { PtyEnvironment } from "./pty-environment"
import { layer } from "./location"
import { formLocationLayer } from "./middleware/form-location"
import { sessionLocationLayer } from "./middleware/session-location"
import { sessionValidationLayer } from "./middleware/session-validation"
import { ServerInfo } from "./server-info"
import type { ServerOptions } from "./options"
@@ -151,6 +152,7 @@ function makeRoutes<AuthError, AuthServices>(
Layer.provide(handlers.pipe(Layer.provide(services))),
Layer.provide(formLocationLayer),
Layer.provide(sessionLocationLayer),
Layer.provide(sessionValidationLayer),
Layer.provide(layer),
Layer.provide(authorizationLayer),
Layer.provide(schemaErrorLayer),
@@ -0,0 +1,117 @@
import { expect } from "bun:test"
import { Location } from "@opencode-ai/core/location"
import { LocationServiceMap } from "@opencode-ai/core/location-service-map"
import type { LocationError, LocationServices } from "@opencode-ai/core/location-services"
import { Session } from "@opencode-ai/core/session"
import { Effect, Layer, LayerMap } from "effect"
import { it } from "../../core/test/lib/effect"
import { ServerFetch } from "../src/fetch"
const fixture = Effect.gen(function* () {
const acquisitions: Location.Ref[] = []
const locations = Layer.effect(
LocationServiceMap.Service,
LayerMap.make((ref: Location.Ref) => {
acquisitions.push(ref)
return Layer.effectContext<LocationServices, LocationError, never>(Effect.die("Location must not be acquired"))
}),
)
const handler = yield* ServerFetch.make(
{
app: { version: "test-version" },
database: { path: ":memory:" },
fs: { filewatcher: false },
config: { project: false, content: "{}" },
},
{ overrides: [[LocationServiceMap.node, locations]] },
)
const post = (pathname: string, body?: unknown) =>
Effect.promise(() =>
handler(
new Request(`http://opencode.local${pathname}`, {
method: "POST",
headers: { "content-type": "application/json" },
body: body === undefined ? undefined : JSON.stringify(body),
}),
),
)
return { acquisitions, post }
})
it.live("session controls preserve malformed and unknown session errors without acquiring a Location", () =>
Effect.gen(function* () {
const server = yield* fixture
yield* Effect.forEach(
["wait", "interrupt", "interrupt?continue=invalid", "prompt", "synthetic", "compact"],
(operation) =>
Effect.gen(function* () {
yield* Effect.forEach(["invalid", "msg_invalid", "SES_invalid"], (id) =>
Effect.gen(function* () {
const response = yield* server.post(`/api/session/${id}/${operation}`)
expect(response.status).toBe(400)
expect(yield* Effect.promise(() => response.json())).toEqual({
_tag: "InvalidRequestError",
message: "Invalid session ID",
field: "sessionID",
})
}),
)
// Session IDs retain their existing loose prefix validation.
yield* Effect.forEach(["ses", Session.ID.create()], (id) =>
Effect.gen(function* () {
const response = yield* server.post(`/api/session/${id}/${operation}`)
expect(response.status).toBe(404)
expect(yield* Effect.promise(() => response.json())).toEqual({
_tag: "SessionNotFoundError",
sessionID: id,
message: `Session not found: ${id}`,
})
}),
)
}),
)
expect(server.acquisitions).toEqual([])
}),
)
it.live("idle session controls do not acquire an unavailable Location", () =>
Effect.gen(function* () {
const server = yield* fixture
const id = Session.ID.create()
expect((yield* server.post("/api/session", { id })).status).toBe(200)
expect(server.acquisitions).toEqual([])
const waited = yield* server.post(`/api/session/${id}/wait`)
expect(waited.status).toBe(204)
expect(yield* Effect.promise(() => waited.text())).toBe("")
yield* Effect.forEach(["", "?continue=false", "?continue=true"], (query) =>
Effect.gen(function* () {
const interrupted = yield* server.post(`/api/session/${id}/interrupt${query}`)
expect(interrupted.status).toBe(200)
expect(yield* Effect.promise(() => interrupted.json())).toEqual({ interrupted: false })
}),
)
const invalidQuery = yield* server.post(`/api/session/${id}/interrupt?continue=invalid`)
expect(invalidQuery.status).toBe(400)
expect(yield* Effect.promise(() => invalidQuery.json())).toMatchObject({
_tag: "InvalidRequestError",
kind: "Query",
})
expect(server.acquisitions).toEqual([])
}),
)
it.live("session admission endpoints still require the Location graph", () =>
Effect.gen(function* () {
const server = yield* fixture
const id = Session.ID.create()
expect((yield* server.post("/api/session", { id })).status).toBe(200)
yield* Effect.forEach(["prompt", "synthetic", "compact"], (operation) =>
Effect.gen(function* () {
const response = yield* server.post(`/api/session/${id}/${operation}`, { text: "input", resume: false })
expect(response.status).toBe(500)
}),
)
expect(server.acquisitions.length).toBeGreaterThan(0)
}),
)