mirror of
https://github.com/anomalyco/opencode.git
synced 2026-10-07 07:48:16 +00:00
Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bb1d75ff5f | ||
|
|
9314014d86 | ||
|
|
c4417d707c | ||
|
|
cc502b3a48 | ||
|
|
7268ae7109 | ||
|
|
c64b231c98 |
No files matched your search
@@ -54,8 +54,12 @@ export type LocationGetInput = { readonly location?: { readonly directory?: stri
|
||||
export type LocationGetOutput = Location.PublicInfo
|
||||
export type LocationGetOperation<E = never> = (input?: LocationGetInput) => Effect.Effect<LocationGetOutput, E>
|
||||
|
||||
export type LocationReloadOutput = void
|
||||
export type LocationReloadOperation<E = never> = () => Effect.Effect<LocationReloadOutput, E>
|
||||
|
||||
export interface LocationApi<E = never> {
|
||||
readonly get: LocationGetOperation<E>
|
||||
readonly reload: LocationReloadOperation<E>
|
||||
}
|
||||
|
||||
export type AgentListInput = { readonly location?: { readonly directory?: string | undefined } | undefined }
|
||||
|
||||
@@ -8,6 +8,7 @@ import type {
|
||||
ServerStatusOutput,
|
||||
LocationGetInput,
|
||||
LocationGetOutput,
|
||||
LocationReloadOutput,
|
||||
AgentListInput,
|
||||
AgentListOutput,
|
||||
AgentGetInput,
|
||||
@@ -287,7 +288,13 @@ const EndpointLocationGet = (raw: RawClient["server.location"]) => (input?: Loca
|
||||
raw["location.get"]({ query: { location: input?.["location"] } }).pipe(Effect.mapError(mapClientError)),
|
||||
)
|
||||
|
||||
const adaptGroupLocation = (raw: RawClient["server.location"]) => ({ get: EndpointLocationGet(raw) })
|
||||
const EndpointLocationReload = (raw: RawClient["server.location"]) => () =>
|
||||
preserveEffect<LocationReloadOutput>()(raw["location.reload"]({}).pipe(Effect.mapError(mapClientError)))
|
||||
|
||||
const adaptGroupLocation = (raw: RawClient["server.location"]) => ({
|
||||
get: EndpointLocationGet(raw),
|
||||
reload: EndpointLocationReload(raw),
|
||||
})
|
||||
|
||||
const EndpointAgentList = (raw: RawClient["server.agent"]) => (input?: AgentListInput) =>
|
||||
preserveEffect<AgentListOutput>()(
|
||||
|
||||
@@ -2,6 +2,7 @@ import type {
|
||||
ServerStatusOutput,
|
||||
LocationGetInput,
|
||||
LocationGetOutput,
|
||||
LocationReloadOutput,
|
||||
AgentListInput,
|
||||
AgentListOutput,
|
||||
AgentGetInput,
|
||||
@@ -416,6 +417,17 @@ export function make(options: ClientOptions) {
|
||||
},
|
||||
requestOptions,
|
||||
),
|
||||
reload: (requestOptions?: RequestOptions) =>
|
||||
request<LocationReloadOutput>(
|
||||
{
|
||||
method: "POST",
|
||||
path: `/api/location/reload`,
|
||||
successStatus: 204,
|
||||
declaredStatuses: [400, 401, 503],
|
||||
empty: true,
|
||||
},
|
||||
requestOptions,
|
||||
),
|
||||
},
|
||||
agent: {
|
||||
list: (input?: AgentListInput, requestOptions?: RequestOptions) =>
|
||||
|
||||
@@ -820,6 +820,15 @@ export type SessionUsageRecorded = {
|
||||
data: { sessionID: string; source: "title" | "compaction"; cost: MoneyUSD; tokens: TokenUsageInfo }
|
||||
}
|
||||
|
||||
export type LocationShutdown = {
|
||||
id: string
|
||||
created: number
|
||||
metadata?: { [x: string]: any }
|
||||
type: "location.shutdown"
|
||||
location?: LocationRef
|
||||
data: {}
|
||||
}
|
||||
|
||||
export type ModelsDevRefreshed = {
|
||||
id: string
|
||||
created: number
|
||||
@@ -2321,6 +2330,7 @@ export type IntegrationInfo = {
|
||||
}
|
||||
|
||||
export type V2Event =
|
||||
| LocationShutdown
|
||||
| ModelsDevRefreshed
|
||||
| CredentialUpdated
|
||||
| CredentialSwitched
|
||||
@@ -2429,14 +2439,6 @@ export type UnauthorizedError = { readonly _tag: "UnauthorizedError"; readonly m
|
||||
export const isUnauthorizedError = (value: unknown): value is UnauthorizedError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "UnauthorizedError"
|
||||
|
||||
export type AgentNotFoundError = {
|
||||
readonly _tag: "AgentNotFoundError"
|
||||
readonly agentID: string
|
||||
readonly message: string
|
||||
}
|
||||
export const isAgentNotFoundError = (value: unknown): value is AgentNotFoundError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "AgentNotFoundError"
|
||||
|
||||
export type ServiceUnavailableError = {
|
||||
readonly _tag: "ServiceUnavailableError"
|
||||
readonly message: string
|
||||
@@ -2445,6 +2447,14 @@ export type ServiceUnavailableError = {
|
||||
export const isServiceUnavailableError = (value: unknown): value is ServiceUnavailableError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "ServiceUnavailableError"
|
||||
|
||||
export type AgentNotFoundError = {
|
||||
readonly _tag: "AgentNotFoundError"
|
||||
readonly agentID: string
|
||||
readonly message: string
|
||||
}
|
||||
export const isAgentNotFoundError = (value: unknown): value is AgentNotFoundError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "AgentNotFoundError"
|
||||
|
||||
export type InvalidCursorError = { readonly _tag: "InvalidCursorError"; readonly message: string }
|
||||
export const isInvalidCursorError = (value: unknown): value is InvalidCursorError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "InvalidCursorError"
|
||||
@@ -2653,6 +2663,8 @@ export type LocationGetInput = {
|
||||
|
||||
export type LocationGetOutput = LocationPublicInfo
|
||||
|
||||
export type LocationReloadOutput = void
|
||||
|
||||
export type AgentListInput = {
|
||||
readonly location?: { readonly location?: { readonly directory?: string | undefined } | undefined }["location"]
|
||||
}
|
||||
|
||||
@@ -79,6 +79,7 @@ export interface ListInput {
|
||||
}
|
||||
|
||||
export interface Interface {
|
||||
readonly close: Effect.Effect<void>
|
||||
readonly create: (input: CreateInput) => Effect.Effect<Info, AlreadyExistsError | InvalidFormError>
|
||||
readonly ask: (input: CreateInput) => Effect.Effect<TerminalState, AlreadyExistsError | InvalidFormError>
|
||||
readonly get: (id: ID) => Effect.Effect<Info, NotFoundError>
|
||||
@@ -100,6 +101,7 @@ export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const bus = yield* Bus.Service
|
||||
let closed = false
|
||||
const forms = yield* Cache.makeWith<ID, Entry>(
|
||||
() => Effect.die(new Error("Form cache must be used via set/getSuccess, never get")),
|
||||
{
|
||||
@@ -137,6 +139,7 @@ export const layer = Layer.effect(
|
||||
}
|
||||
yield* Cache.set(forms, id, entry)
|
||||
yield* bus.publish(Form.Event.Created, { form }).pipe(Effect.onError(() => Cache.invalidate(forms, id)))
|
||||
if (closed) yield* cancel(id).pipe(Effect.orDie)
|
||||
return form
|
||||
}),
|
||||
),
|
||||
@@ -202,19 +205,21 @@ export const layer = Layer.effect(
|
||||
),
|
||||
)
|
||||
|
||||
yield* Effect.addFinalizer(() =>
|
||||
Cache.values(forms).pipe(
|
||||
Effect.flatMap((entries) =>
|
||||
Effect.forEach(
|
||||
Array.from(entries).filter((entry) => entry.state.status === "pending"),
|
||||
(entry) => cancel(entry.form.id).pipe(Effect.ignore),
|
||||
{ discard: true },
|
||||
),
|
||||
const close = Effect.sync(() => {
|
||||
closed = true
|
||||
}).pipe(
|
||||
Effect.andThen(Cache.values(forms)),
|
||||
Effect.flatMap((entries) =>
|
||||
Effect.forEach(
|
||||
Array.from(entries).filter((entry) => entry.state.status === "pending"),
|
||||
(entry) => cancel(entry.form.id).pipe(Effect.ignore),
|
||||
{ discard: true },
|
||||
),
|
||||
),
|
||||
)
|
||||
yield* Effect.addFinalizer(() => close)
|
||||
|
||||
return Service.of({ create, ask, get, list, state, reply, cancel })
|
||||
return Service.of({ create, ask, get, list, state, reply, cancel, close })
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Effect, Layer } from "effect"
|
||||
import { Context, Effect, Layer } from "effect"
|
||||
import { Agent } from "./agent.js"
|
||||
import { AISDK } from "./aisdk.js"
|
||||
import { Model } from "./model.js"
|
||||
@@ -18,6 +18,7 @@ import { Image } from "./image.js"
|
||||
import { LocationWatcher } from "./filesystem/location-watcher.js"
|
||||
import { Integration } from "./integration.js"
|
||||
import { Location } from "./location.js"
|
||||
import { LocationLifecycle } from "./location-lifecycle.js"
|
||||
import { FileAccess } from "./file-access.js"
|
||||
import { ModelResolver } from "./model-resolver.js"
|
||||
import { Mcp } from "./mcp/index.js"
|
||||
@@ -57,6 +58,7 @@ export { Service, node, type Interface } from "./instance/service.js"
|
||||
|
||||
const nodes = [
|
||||
Location.node,
|
||||
LocationLifecycle.node,
|
||||
Environment.node,
|
||||
Config.node,
|
||||
Agent.node,
|
||||
@@ -157,6 +159,7 @@ export function layer(ref: Location.Ref, options: Options = {}): Layer.Layer<Ser
|
||||
return LayerNode.compile(graph, { replacements, shared: Node.tags.values.global }).pipe(
|
||||
// Instance boot failures are defects; provided operations retain their typed errors.
|
||||
Layer.orDie,
|
||||
Layer.tap((context) => Effect.addFinalizer(() => Context.get(context, LocationLifecycle.Service).shutdown)),
|
||||
Layer.tap(() =>
|
||||
Effect.logInfo("location services booted", {
|
||||
directory: ref.directory,
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
export * as LocationLifecycle from "./location-lifecycle.js"
|
||||
|
||||
import { Context, Effect, Layer } from "effect"
|
||||
import { makeLocationNode } from "@opencode/util/effect/app-node"
|
||||
import { LocationEvent } from "@opencode/schema/location-event"
|
||||
import { Bus } from "./bus.js"
|
||||
import { Form } from "./form.js"
|
||||
import { Location } from "./location.js"
|
||||
import { Permission } from "./permission.js"
|
||||
import { Rpc } from "./rpc.js"
|
||||
|
||||
export class Service extends Context.Service<
|
||||
Service,
|
||||
{ readonly isClosed: () => boolean; readonly shutdown: Effect.Effect<void> }
|
||||
>()("@opencode/LocationLifecycle") {}
|
||||
|
||||
const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const bus = yield* Bus.Service
|
||||
const location = yield* Location.Service
|
||||
const permission = yield* Permission.Service
|
||||
const forms = yield* Form.Service
|
||||
const rpc = yield* Rpc.Service
|
||||
let closed = false
|
||||
const shutdown = yield* Effect.cached(
|
||||
Effect.gen(function* () {
|
||||
closed = true
|
||||
yield* permission.close
|
||||
yield* forms.close
|
||||
yield* rpc.close
|
||||
yield* bus.publish(
|
||||
LocationEvent.Shutdown,
|
||||
{},
|
||||
{
|
||||
location: Location.Ref.make({ directory: location.directory, workspaceID: location.workspaceID }),
|
||||
},
|
||||
)
|
||||
}).pipe(Effect.uninterruptible),
|
||||
)
|
||||
return Service.of({
|
||||
isClosed: () => closed,
|
||||
shutdown,
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export const node = makeLocationNode({
|
||||
service: Service,
|
||||
layer,
|
||||
deps: [Bus.node, Location.node, Permission.node, Form.node, Rpc.node],
|
||||
})
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Context, Effect, Layer, LayerMap } from "effect"
|
||||
import { Context, Effect, Exit, Layer, LayerMap, RcMap } from "effect"
|
||||
import { LayerNode } from "@opencode/util/effect/layer-node"
|
||||
import { Node } from "@opencode/util/effect/app-node"
|
||||
import { AbsolutePath } from "@opencode/schema/schema"
|
||||
@@ -17,6 +17,24 @@ export class Service extends Context.Service<
|
||||
|
||||
export const node = LayerNode.unbound(Service, Node.tags.values.global)
|
||||
|
||||
export const reload = Effect.fn("LocationServiceMap.reload")(function* () {
|
||||
const locations = yield* Service
|
||||
const refs = Array.from(yield* RcMap.keys(locations.rcMap))
|
||||
yield* Effect.forEach(refs, (ref) => locations.invalidate(ref), {
|
||||
discard: true,
|
||||
concurrency: "unbounded",
|
||||
})
|
||||
// Boot every replacement now and let all builds settle even if one fails.
|
||||
const results = yield* Effect.forEach(
|
||||
refs,
|
||||
(ref) => Effect.scoped(locations.contextEffect(ref)).pipe(Effect.asVoid, Effect.exit),
|
||||
{ concurrency: "unbounded" },
|
||||
)
|
||||
const failure = results.find(Exit.isFailure)
|
||||
if (failure) return yield* Effect.failCause(failure.cause)
|
||||
yield* Effect.logInfo("location services reloaded", { count: refs.length })
|
||||
})
|
||||
|
||||
/** Normalize equivalent placements before they become resource-cache keys. */
|
||||
export function canonical(ref: Location.Ref) {
|
||||
return Location.Ref.make({
|
||||
|
||||
@@ -2,8 +2,8 @@ import { Context, Duration, Effect, Exit, Layer, LayerMap, MutableHashMap, Optio
|
||||
import { LayerNode } from "@opencode/util/effect/layer-node"
|
||||
import { Instance } from "./instance.js"
|
||||
import { Location } from "./location.js"
|
||||
import { LocationLifecycle } from "./location-lifecycle.js"
|
||||
import { LocationServiceMap } from "./location-service-map.js"
|
||||
import { Rpc } from "./rpc.js"
|
||||
|
||||
export { LocationServiceMap } from "./location-service-map.js"
|
||||
|
||||
@@ -28,12 +28,17 @@ export function buildLocationServiceMap(
|
||||
).pipe(
|
||||
Effect.onExit((exit) => {
|
||||
const finish = Effect.suspend(() => {
|
||||
if (Exit.isSuccess(exit)) {
|
||||
return Effect.gen(function* () {
|
||||
const lifecycle = Context.get(exit.value, LocationLifecycle.Service)
|
||||
// A boot detached while in flight still needs shutdown and cancellation.
|
||||
if (Option.getOrUndefined(MutableHashMap.get(builds, ref)) !== build)
|
||||
return yield* lifecycle.shutdown
|
||||
build.close = lifecycle.shutdown
|
||||
})
|
||||
}
|
||||
// An explicitly invalidated build must not evict its replacement.
|
||||
if (Option.getOrUndefined(MutableHashMap.get(builds, ref)) !== build) return Effect.void
|
||||
if (Exit.isSuccess(exit)) {
|
||||
build.close = Context.get(exit.value, Rpc.Service).close
|
||||
return Effect.void
|
||||
}
|
||||
MutableHashMap.remove(builds, ref)
|
||||
// Evict once per failed build, before its result reaches borrowers.
|
||||
return Exit.isFailure(exit) ? inner.invalidate(ref) : Effect.void
|
||||
@@ -61,10 +66,11 @@ export function buildLocationServiceMap(
|
||||
const key = LocationServiceMap.canonical(ref)
|
||||
const build = Option.getOrUndefined(MutableHashMap.get(builds, key))
|
||||
MutableHashMap.remove(builds, key)
|
||||
// Detach routing first, then end pending RPCs that still borrow the old graph.
|
||||
// Detach routing first, then cancel interactions and notify clients. Running
|
||||
// steps retain their borrowed graph until they can hand off at a boundary.
|
||||
// Do not await a boot here: failed/in-flight builds have their own cleanup path.
|
||||
return inner.invalidate(key).pipe(Effect.andThen(build?.close ?? Effect.void))
|
||||
}),
|
||||
}).pipe(Effect.uninterruptible),
|
||||
}
|
||||
// Cached instances borrow their owner instead of retaining its Layer scope.
|
||||
const bindings: LayerNode.Replacements = [
|
||||
|
||||
@@ -81,6 +81,8 @@ type ServerEntry = {
|
||||
// persisted session row, so their forms are owned by this opaque sentinel session identifier.
|
||||
const GLOBAL_ELICITATION_SESSION_ID = "global"
|
||||
const URL_ELICITATION_FIELD_KEY = "elicitation"
|
||||
// Connections remain Location-scoped, but shared remote endpoints should not receive concurrent startup bursts.
|
||||
const endpointLoads = KeyedMutex.makeUnsafe<string>()
|
||||
|
||||
type Data = {
|
||||
servers: Map<ServerName, Types.DeepMutable<Mcp.ServerConfig>>
|
||||
@@ -387,7 +389,7 @@ export const layer = (options?: Options) =>
|
||||
const { McpClient } = yield* Effect.promise(() => import("./client.js"))
|
||||
// List tools as part of connect so a failure here marks the server failed rather than
|
||||
// leaving it connected with a silently empty tool list and no path to recover.
|
||||
const result = yield* McpClient.connect(
|
||||
const load = McpClient.connect(
|
||||
name,
|
||||
entry.config,
|
||||
location.directory,
|
||||
@@ -399,8 +401,10 @@ export const layer = (options?: Options) =>
|
||||
// A stdio server is spawned on this location's execution plane, not the host's.
|
||||
Effect.provideService(Environment.Service, environment),
|
||||
Scope.provide(scope),
|
||||
Effect.exit,
|
||||
)
|
||||
const result = yield* (
|
||||
entry.config.type === "remote" ? endpointLoads.withLock(entry.config.url)(load) : load
|
||||
).pipe(Effect.exit)
|
||||
if (Exit.isSuccess(result)) {
|
||||
entry.client = result.value.connection
|
||||
entry.tools = result.value.tools.map((tool) => toTool(name, entry, tool))
|
||||
|
||||
@@ -101,6 +101,7 @@ export function merge(...rulesets: Permission.Ruleset[]): Permission.Ruleset {
|
||||
}
|
||||
|
||||
export interface Interface {
|
||||
readonly close: Effect.Effect<void>
|
||||
readonly ask: (input: AssertInput) => Effect.Effect<AskResult, SessionErrors.NotFoundError>
|
||||
readonly assert: (input: AssertInput) => Effect.Effect<void, Error | SessionErrors.NotFoundError>
|
||||
readonly reply: (input: ReplyInput) => Effect.Effect<void, NotFoundError>
|
||||
@@ -127,18 +128,22 @@ const layer = Layer.effect(
|
||||
const saved = yield* PermissionSaved.Service
|
||||
const hooks = yield* PluginHooks.Service
|
||||
const pending = new Map<ID, Pending>()
|
||||
let closed = false
|
||||
|
||||
yield* Effect.addFinalizer(() =>
|
||||
Effect.forEach(pending.values(), (item) => Deferred.fail(item.deferred, new DeclinedError()), {
|
||||
discard: true,
|
||||
}).pipe(
|
||||
Effect.ensuring(
|
||||
Effect.sync(() => {
|
||||
pending.clear()
|
||||
}),
|
||||
),
|
||||
),
|
||||
)
|
||||
const close = Effect.gen(function* () {
|
||||
closed = true
|
||||
yield* Effect.forEach(Array.from(pending.values()), (item) =>
|
||||
bus
|
||||
.publish(Permission.Event.Replied, {
|
||||
sessionID: item.request.sessionID,
|
||||
requestID: item.request.id,
|
||||
reply: "reject",
|
||||
})
|
||||
.pipe(Effect.ensuring(Deferred.fail(item.deferred, new DeclinedError()))),
|
||||
)
|
||||
pending.clear()
|
||||
}).pipe(Effect.uninterruptible)
|
||||
yield* Effect.addFinalizer(() => close)
|
||||
|
||||
const savedRules = Effect.fnUntraced(function* () {
|
||||
return (yield* saved.list({ projectID: location.project.id })).map(
|
||||
@@ -201,6 +206,10 @@ const layer = Layer.effect(
|
||||
Effect.gen(function* () {
|
||||
const deferred = yield* Deferred.make<void, DeclinedError | CorrectedError>()
|
||||
const item = { request, agent, deferred }
|
||||
if (closed) {
|
||||
yield* Deferred.fail(deferred, new DeclinedError())
|
||||
return item
|
||||
}
|
||||
if (pending.has(request.id))
|
||||
return yield* Effect.die(new Error(`Duplicate pending permission ID: ${request.id}`))
|
||||
pending.set(request.id, item)
|
||||
@@ -212,6 +221,7 @@ const layer = Layer.effect(
|
||||
)
|
||||
|
||||
const ask = Effect.fn("Permission.ask")(function* (input: AssertInput) {
|
||||
if (closed) return { id: input.id ?? ID.create(), effect: "deny" as const }
|
||||
const result = yield* evaluateInput(input)
|
||||
const value = request(input, result.message)
|
||||
if (result.effect === "ask") yield* create(value, input.agent)
|
||||
@@ -220,6 +230,7 @@ const layer = Layer.effect(
|
||||
|
||||
const assert = Effect.fn("Permission.assert")((input: AssertInput) =>
|
||||
Effect.gen(function* () {
|
||||
if (closed) return yield* Effect.die(new DeclinedError())
|
||||
const result = yield* evaluateInput(input)
|
||||
return yield* Effect.uninterruptibleMask((restore) =>
|
||||
Effect.gen(function* () {
|
||||
@@ -321,7 +332,7 @@ const layer = Layer.effect(
|
||||
return Array.from(pending.values(), (item) => item.request).filter((request) => request.sessionID === sessionID)
|
||||
})
|
||||
|
||||
return Service.of({ ask, assert, reply, get, forSession, list })
|
||||
return Service.of({ ask, assert, reply, get, forSession, list, close })
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -108,6 +108,7 @@ export const layer = Layer.effect(
|
||||
return yield* SessionRunner.DrainResult.$match(result, {
|
||||
Complete: () => Effect.void,
|
||||
Moved: (result) => drain(sessionID, false, result.continuation, promotable),
|
||||
Reloaded: (result) => drain(sessionID, result.force, result.continuation, promotable),
|
||||
})
|
||||
})
|
||||
const coordinator = yield* SessionRunCoordinator.make<SessionSchema.ID, SessionRunner.RunError, InterruptReason>({
|
||||
|
||||
@@ -22,6 +22,7 @@ export type Continuation = { readonly step: number }
|
||||
export type DrainResult = Data.TaggedEnum<{
|
||||
Complete: {}
|
||||
Moved: { readonly continuation?: Continuation }
|
||||
Reloaded: { readonly force: boolean; readonly continuation?: Continuation }
|
||||
}>
|
||||
export const DrainResult = Data.taggedEnum<DrainResult>()
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ import { and, desc, eq, sql } from "drizzle-orm"
|
||||
import { Cause, Effect, Exit, FiberMap, Layer } from "effect"
|
||||
import { Database } from "../../database/database.js"
|
||||
import { Bus } from "../../bus.js"
|
||||
import { LocationLifecycle } from "../../location-lifecycle.js"
|
||||
import { InstructionState } from "../instruction-state.js"
|
||||
import { SessionCompaction } from "../compaction.js"
|
||||
import { SessionContext } from "../context.js"
|
||||
@@ -37,6 +38,7 @@ const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const bus = yield* Bus.Service
|
||||
const lifecycle = yield* LocationLifecycle.Service
|
||||
const store = yield* SessionStore.Service
|
||||
const context = yield* SessionContext.Service
|
||||
const modelTransport = yield* SessionModelTransport.Service
|
||||
@@ -69,6 +71,10 @@ const layer = Layer.effect(
|
||||
Effect.uninterruptibleMask((restore) =>
|
||||
Effect.gen(function* () {
|
||||
while (true) {
|
||||
if (lifecycle.isClosed()) {
|
||||
yield* restore(modelTransport.close(sessionID))
|
||||
return DrainResult.Reloaded({ force, continuation: continuing ? { step } : undefined })
|
||||
}
|
||||
// Location entry and idle boundaries allow queued controls, not necessarily queued prompts.
|
||||
const pending = yield* SessionInbox.serialized(
|
||||
sessionID,
|
||||
@@ -240,6 +246,7 @@ const layer = Layer.effect(
|
||||
webSocket: "session",
|
||||
})
|
||||
const outcome = yield* steps.attempt({
|
||||
isLocationClosed: lifecycle.isClosed,
|
||||
sessionID,
|
||||
assistantMessageID,
|
||||
agent: loaded.agent.id,
|
||||
@@ -356,6 +363,7 @@ export const node = makeLocationNode({
|
||||
layer,
|
||||
deps: [
|
||||
Bus.node,
|
||||
LocationLifecycle.node,
|
||||
llmClient,
|
||||
SessionContext.node,
|
||||
SessionModelTransport.node,
|
||||
|
||||
@@ -42,6 +42,7 @@ export type Outcome = Data.TaggedEnum<{
|
||||
export const Outcome = Data.taggedEnum<Outcome>()
|
||||
|
||||
interface Input {
|
||||
readonly isLocationClosed: () => boolean
|
||||
readonly sessionID: SessionSchema.ID
|
||||
readonly assistantMessageID: SessionMessage.ID
|
||||
readonly agent: Agent.ID
|
||||
@@ -190,8 +191,9 @@ export const make = Effect.gen(function* () {
|
||||
for (const decline of tools.declines)
|
||||
yield* publisher.failTool(decline.call.id, {
|
||||
type: "aborted",
|
||||
message:
|
||||
decline.reason._tag === "QuestionTool.CancelledError"
|
||||
message: input.isLocationClosed()
|
||||
? "Interaction cancelled because the location shut down"
|
||||
: decline.reason._tag === "QuestionTool.CancelledError"
|
||||
? decline.reason.message
|
||||
: "The user declined this tool call",
|
||||
})
|
||||
@@ -251,7 +253,10 @@ export const make = Effect.gen(function* () {
|
||||
return Outcome.Continue({ error: llmError, decision: retry })
|
||||
|
||||
if (Exit.isFailure(stream)) return yield* Effect.failCause(stream.cause)
|
||||
if (tools.declines.length > 0) return yield* Effect.interrupt
|
||||
if (tools.declines.length > 0) {
|
||||
if (input.isLocationClosed()) return Outcome.Completed({ needsContinuation: true })
|
||||
return yield* Effect.interrupt
|
||||
}
|
||||
if (tools.interrupted && tools.failure) return yield* Effect.failCause(tools.failure)
|
||||
if (tools.interrupted && Exit.isFailure(joined)) return yield* Effect.failCause(joined.cause)
|
||||
if (record.failure) return yield* new StepFailedError({ error: record.failure })
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { Permission } from "@opencode/core/permission"
|
||||
import { Layer } from "effect"
|
||||
import { Effect, Layer } from "effect"
|
||||
|
||||
export const permissionLayer = (overrides: Partial<Permission.Interface> = {}) =>
|
||||
Layer.mock(Permission.Service, overrides)
|
||||
Layer.mock(Permission.Service, { close: Effect.void, ...overrides })
|
||||
@@ -55,6 +55,7 @@ const generateLayer = Layer.succeed(Generate.Service, Generate.Service.of({ text
|
||||
const permissionLayer = Layer.succeed(
|
||||
Permission.Service,
|
||||
Permission.Service.of({
|
||||
close: Effect.void,
|
||||
ask: (input) => Effect.succeed({ id: input.id ?? Permission.ID.create(), effect: "ask" }),
|
||||
assert: () => Effect.void,
|
||||
reply: () => Effect.void,
|
||||
|
||||
@@ -100,6 +100,7 @@ for (const fixture of [
|
||||
)
|
||||
const result = yield* steps
|
||||
.attempt({
|
||||
isLocationClosed: () => false,
|
||||
sessionID,
|
||||
assistantMessageID,
|
||||
agent: Agent.defaultID,
|
||||
|
||||
@@ -23,7 +23,7 @@ import { PersistentPtyGroup } from "./groups/persistent-pty.js"
|
||||
import { ShellGroup } from "./groups/shell.js"
|
||||
import { ReferenceGroup } from "./groups/reference.js"
|
||||
import { Authorization } from "./middleware/authorization.js"
|
||||
import { LocationGroup } from "./groups/location.js"
|
||||
import { makeLocationGroup } from "./groups/location.js"
|
||||
import { IntegrationGroup } from "./groups/integration.js"
|
||||
import { WebSearchGroup } from "./groups/websearch.js"
|
||||
import { McpGroup } from "./groups/mcp.js"
|
||||
@@ -35,7 +35,6 @@ import { MigrationGroup } from "./groups/migration.js"
|
||||
import { ConfigGroup } from "./groups/config.js"
|
||||
|
||||
type LocationGroups<LocationId extends HttpApiMiddleware.AnyId> =
|
||||
| HttpApiGroup.AddMiddleware<typeof LocationGroup, LocationId>
|
||||
| HttpApiGroup.AddMiddleware<typeof AgentGroup, LocationId>
|
||||
| HttpApiGroup.AddMiddleware<typeof PluginGroup, LocationId>
|
||||
| HttpApiGroup.AddMiddleware<typeof ModelGroup, LocationId>
|
||||
@@ -60,15 +59,17 @@ type SessionGroups<
|
||||
FormLocationId extends HttpApiMiddleware.AnyId,
|
||||
FormLocationService,
|
||||
> =
|
||||
| ReturnType<
|
||||
typeof makeSessionGroup<SessionLocationId, SessionLocationService, FormLocationId, FormLocationService>
|
||||
>
|
||||
| ReturnType<typeof makeSessionGroup<SessionLocationId, SessionLocationService, FormLocationId, FormLocationService>>
|
||||
| typeof MessageGroup
|
||||
|
||||
type FormGroups<LocationId extends HttpApiMiddleware.AnyId, LocationService> = ReturnType<
|
||||
typeof makeFormGroup<LocationId, LocationService>
|
||||
>
|
||||
|
||||
type LocationGroup<LocationId extends HttpApiMiddleware.AnyId, LocationService> = ReturnType<
|
||||
typeof makeLocationGroup<LocationId, LocationService>
|
||||
>
|
||||
|
||||
type MixedMiddlewareGroups<
|
||||
LocationId extends HttpApiMiddleware.AnyId,
|
||||
LocationService,
|
||||
@@ -93,6 +94,7 @@ type ApiGroups<
|
||||
| typeof PersistentPtyGroup
|
||||
| typeof CredentialGroup
|
||||
| LocationGroups<LocationId>
|
||||
| LocationGroup<LocationId, LocationService>
|
||||
| FormGroups<LocationId, LocationService>
|
||||
| SessionGroups<SessionLocationId, SessionLocationService, FormLocationId, FormLocationService>
|
||||
| MixedMiddlewareGroups<LocationId, LocationService, SessionLocationId, SessionLocationService>
|
||||
@@ -152,7 +154,7 @@ const makeApiFromGroup = <
|
||||
> =>
|
||||
HttpApi.make("server")
|
||||
.add(ServerGroup)
|
||||
.add(LocationGroup.middleware(locationMiddleware))
|
||||
.add(makeLocationGroup(locationMiddleware))
|
||||
.add(AgentGroup.middleware(locationMiddleware))
|
||||
.add(PluginGroup.middleware(locationMiddleware))
|
||||
.add(makeSessionGroup(sessionLocationMiddleware, formLocationMiddleware))
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { Location } from "@opencode/schema/location"
|
||||
import { Schema } from "effect"
|
||||
import { HttpApiEndpoint, HttpApiGroup, OpenApi } from "effect/unstable/httpapi"
|
||||
import { Context, Schema } from "effect"
|
||||
import { HttpApiEndpoint, HttpApiGroup, HttpApiMiddleware, HttpApiSchema, OpenApi } from "effect/unstable/httpapi"
|
||||
import { ServiceUnavailableError } from "../errors.js"
|
||||
|
||||
export const LocationQuery = Schema.Struct({
|
||||
location: Schema.optional(
|
||||
@@ -25,19 +26,38 @@ export const locationQueryOpenApi = OpenApi.annotations({
|
||||
},
|
||||
})
|
||||
|
||||
export const LocationGroup = HttpApiGroup.make("server.location")
|
||||
.add(
|
||||
HttpApiEndpoint.get("location.get", "/api/location", {
|
||||
query: LocationQuery,
|
||||
success: Location.PublicInfo,
|
||||
})
|
||||
.annotateMerge(locationQueryOpenApi)
|
||||
.annotateMerge(
|
||||
// Middleware is applied per endpoint: reload acts on every loaded location and
|
||||
// must not boot the caller's location first.
|
||||
export const makeLocationGroup = <LocationId extends HttpApiMiddleware.AnyId, LocationService>(
|
||||
locationMiddleware: Context.Key<LocationId, LocationService>,
|
||||
) =>
|
||||
HttpApiGroup.make("server.location")
|
||||
.add(
|
||||
HttpApiEndpoint.get("location.get", "/api/location", {
|
||||
query: LocationQuery,
|
||||
success: Location.PublicInfo,
|
||||
})
|
||||
.middleware(locationMiddleware)
|
||||
.annotateMerge(locationQueryOpenApi)
|
||||
.annotateMerge(
|
||||
OpenApi.annotations({
|
||||
identifier: "location.get",
|
||||
summary: "Get location",
|
||||
description: "Resolve the requested location or the server default location.",
|
||||
}),
|
||||
),
|
||||
)
|
||||
.add(
|
||||
HttpApiEndpoint.post("location.reload", "/api/location/reload", {
|
||||
success: HttpApiSchema.NoContent,
|
||||
error: ServiceUnavailableError,
|
||||
}).annotateMerge(
|
||||
OpenApi.annotations({
|
||||
identifier: "location.get",
|
||||
summary: "Get location",
|
||||
description: "Resolve the requested location or the server default location.",
|
||||
identifier: "location.reload",
|
||||
summary: "Reload locations",
|
||||
description:
|
||||
"Shut down and rebuild every loaded location. Pending permissions and forms are cancelled; running sessions continue with fresh services at the next step boundary. Emits location.shutdown for client recovery and responds once all replacement builds settle.",
|
||||
}),
|
||||
),
|
||||
)
|
||||
.annotateMerge(OpenApi.annotations({ title: "location" }))
|
||||
)
|
||||
.annotateMerge(OpenApi.annotations({ title: "location" }))
|
||||
@@ -14,6 +14,7 @@ import { InstallationEvent } from "./installation-event.js"
|
||||
import { Integration } from "./integration.js"
|
||||
import { LegacyEventV1 } from "./legacy-event.js"
|
||||
import { LspEvent } from "./lsp-event.js"
|
||||
import { LocationEvent } from "./location-event.js"
|
||||
import { McpEvent } from "./mcp-event.js"
|
||||
import { Model } from "./model.js"
|
||||
import { ModelsDev } from "./models-dev.js"
|
||||
@@ -40,6 +41,7 @@ import { WebSearch } from "./websearch.js"
|
||||
const coreDefinitions = Event.inventory(...SessionEvent.Definitions)
|
||||
|
||||
const foundationDefinitions = Event.inventory(
|
||||
...LocationEvent.Definitions,
|
||||
...ModelsDev.Event.Definitions,
|
||||
...Credential.Event.Definitions,
|
||||
...Integration.Event.Definitions,
|
||||
|
||||
@@ -0,0 +1,8 @@
|
||||
export * as LocationEvent from "./location-event.js"
|
||||
|
||||
import { ephemeral, inventory } from "./event.js"
|
||||
|
||||
/** The location's cached services were shut down; clients must revalidate its reads. */
|
||||
export const Shutdown = ephemeral({ type: "location.shutdown", schema: {} })
|
||||
|
||||
export const Definitions = inventory(Shutdown)
|
||||
@@ -1,17 +1,33 @@
|
||||
import { Location } from "@opencode/core/location"
|
||||
import { Effect } from "effect"
|
||||
import { LocationServiceMap } from "@opencode/core/location-service-map"
|
||||
import { ServiceUnavailableError } from "@opencode/protocol/errors"
|
||||
import { Cause, Effect } from "effect"
|
||||
import { HttpApiBuilder } from "effect/unstable/httpapi"
|
||||
import { Api } from "../api"
|
||||
|
||||
export const LocationHandler = HttpApiBuilder.group(Api, "server.location", (handlers) =>
|
||||
handlers.handle(
|
||||
"location.get",
|
||||
Effect.fn(function* () {
|
||||
const location = yield* Location.Service
|
||||
return new Location.Info({
|
||||
directory: location.directory,
|
||||
project: location.project,
|
||||
})
|
||||
}),
|
||||
),
|
||||
Effect.gen(function* () {
|
||||
const locations = yield* LocationServiceMap.Service
|
||||
return handlers
|
||||
.handle(
|
||||
"location.get",
|
||||
Effect.fn(function* () {
|
||||
const location = yield* Location.Service
|
||||
return new Location.Info({
|
||||
directory: location.directory,
|
||||
project: location.project,
|
||||
})
|
||||
}),
|
||||
)
|
||||
.handle("location.reload", () =>
|
||||
LocationServiceMap.reload().pipe(
|
||||
Effect.provideService(LocationServiceMap.Service, locations),
|
||||
Effect.catchCause((cause) =>
|
||||
Cause.hasInterruptsOnly(cause)
|
||||
? Effect.failCause(cause)
|
||||
: Effect.fail(new ServiceUnavailableError({ message: Cause.pretty(cause), service: "location" })),
|
||||
),
|
||||
),
|
||||
)
|
||||
}),
|
||||
)
|
||||
@@ -1003,6 +1003,22 @@ function App(props: { pair?: DialogPairCredentials }) {
|
||||
},
|
||||
]
|
||||
: []),
|
||||
{
|
||||
name: "location.reload",
|
||||
title: "Reload locations",
|
||||
slash: { name: "reload" },
|
||||
run: async () => {
|
||||
dialog.clear()
|
||||
toast.show({ variant: "info", message: "Reloading all locations…", duration: 30000 })
|
||||
await client.api.location
|
||||
.reload()
|
||||
.then(() => {
|
||||
toast.show({ variant: "success", message: "Locations reloaded" })
|
||||
})
|
||||
.catch(toast.error)
|
||||
},
|
||||
category: "System",
|
||||
},
|
||||
{
|
||||
name: "opencode.debug",
|
||||
title: "View debug info",
|
||||
|
||||
Reference in new issue
Block a user