diff --git a/apps/server/src/actionResume/ActionResume.test.ts b/apps/server/src/actionResume/ActionResume.test.ts index 91d36bf5866d..286366af68bb 100644 --- a/apps/server/src/actionResume/ActionResume.test.ts +++ b/apps/server/src/actionResume/ActionResume.test.ts @@ -27,6 +27,7 @@ import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSna import * as ThreadActionResume from "../orchestration/ThreadActionResume.ts"; import { makeProviderRegistryLayer } from "../provider/testUtils/providerRegistryMock.ts"; import * as TerminalManager from "../terminal/Manager.ts"; +import { UpdateDrainAdmission } from "../updateDrain/UpdateDrainAdmission.ts"; import * as ActionResume from "./ActionResume.ts"; const threadId = ThreadId.make("thread-action-resume"); @@ -114,6 +115,8 @@ it.effect("runs one opted-in Action and delivers exactly one automated follow-up const timeline: string[] = []; let terminalStatus: "running" | "exited" = "running"; let failWrite = false; + let admissionClosed = false; + const admittedKinds: string[] = []; let terminalListener: ((event: TerminalEvent) => Effect.Effect) | undefined; const dependencies = Layer.mergeAll( @@ -169,6 +172,14 @@ it.effect("runs one opted-in Action and delivers exactly one automated follow-up } as never, ]), ThreadActionResume.layer, + Layer.mock(UpdateDrainAdmission)({ + admit: (kind, effect) => + Effect.sync(() => admittedKinds.push(kind)).pipe( + Effect.andThen( + admissionClosed ? Effect.die("update drain is closed in this test") : effect, + ), + ), + }), NodeServices.layer, ); @@ -280,6 +291,35 @@ it.effect("runs one opted-in Action and delivers exactly one automated follow-up .pipe(Effect.flip); assert.equal(earlyExit.reason, "launch_failed"); assert.equal(registry.getLatest(threadId)?.runId, failedState?.runId); + + terminalStatus = "running"; + const blocked = yield* service.runProjectActionAndResume( + { threadId, providerInstanceId }, + "qa", + ); + admissionClosed = true; + yield* terminalListener!({ + type: "exited", + threadId, + terminalId: blocked.terminalId, + exitCode: 0, + exitSignal: null, + }); + assert.equal(dispatched.filter((command) => command.type === "thread.turn.start").length, 1); + assert.equal(admittedKinds.at(-1), "thread-turn"); + assert.deepInclude(registry.getLatest(threadId), { + outcome: "succeeded", + delivery: "pending", + }); + + admissionClosed = false; + yield* service.retryPendingFollowUps; + yield* service.retryPendingFollowUps; + assert.equal(dispatched.filter((command) => command.type === "thread.turn.start").length, 2); + assert.deepInclude(registry.getLatest(threadId), { + outcome: "succeeded", + delivery: "delivered", + }); }).pipe(Effect.provide(ActionResume.layer.pipe(Layer.provideMerge(dependencies))), Effect.scoped); }); @@ -420,6 +460,9 @@ it.effect("requires an explicit resume after a running Action is found on startu } as never, ]), ThreadActionResume.layer, + Layer.mock(UpdateDrainAdmission)({ + admit: (_kind, effect) => effect, + }), NodeServices.layer, ); diff --git a/apps/server/src/actionResume/ActionResume.ts b/apps/server/src/actionResume/ActionResume.ts index 291bdbfae67c..3688753d0b35 100644 --- a/apps/server/src/actionResume/ActionResume.ts +++ b/apps/server/src/actionResume/ActionResume.ts @@ -37,6 +37,7 @@ import { OrchestrationEngineService } from "../orchestration/Services/Orchestrat import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts"; import { ThreadActionResumeService } from "../orchestration/ThreadActionResume.ts"; import { ProviderRegistry } from "../provider/Services/ProviderRegistry.ts"; +import { UpdateDrainAdmission } from "../updateDrain/UpdateDrainAdmission.ts"; export const ACTION_RESUME_ACTIVITY_KIND = "action.resume.lifecycle"; @@ -74,6 +75,7 @@ export class ActionResume extends Context.Service< readonly cancelByArchive: (threadId: ThreadId) => Effect.Effect; readonly resumeInterrupted: (threadId: ThreadId) => Effect.Effect; readonly discardInterrupted: (threadId: ThreadId) => Effect.Effect; + readonly retryPendingFollowUps: Effect.Effect; readonly countRunning: Effect.Effect; } >()("t3/actionResume/ActionResume") {} @@ -262,6 +264,7 @@ const make = Effect.gen(function* () { const registry = yield* ThreadActionResumeService; const terminals = yield* TerminalManager.TerminalManager; const providers = yield* ProviderRegistry; + const admission = yield* UpdateDrainAdmission; const mutex = yield* Semaphore.make(1); const decodeState = Schema.decodeUnknownEffect(ActionResumeState); const outputCaptureByRunId = new Map(); @@ -376,7 +379,7 @@ const make = Effect.gen(function* () { mutex.withPermits(1)(deliverPendingUnlocked(threadId)); const deliverPending = (threadId: ThreadId) => - attemptDeliverPending(threadId).pipe( + admission.admit("thread-turn", attemptDeliverPending(threadId)).pipe( Effect.catchCause((cause) => Effect.logWarning("Action follow-up delivery failed; it remains pending", { threadId, @@ -385,6 +388,14 @@ const make = Effect.gen(function* () { ), ); + const retryPendingFollowUps = Effect.suspend(() => + Effect.forEach( + registry.listLatest().filter((state) => state.delivery === "pending"), + (state) => deliverPending(state.threadId), + { concurrency: 1, discard: true }, + ), + ); + const finishUnlocked = Effect.fn("ActionResume.finishUnlocked")(function* ( input: FinishActionInput, ) { @@ -778,6 +789,7 @@ const make = Effect.gen(function* () { resumeInterruptedImpl(threadId).pipe(mapActionResumeError("resume the interrupted Action")), discardInterrupted: (threadId) => discardInterruptedImpl(threadId).pipe(mapActionResumeError("discard the interrupted Action")), + retryPendingFollowUps, countRunning: Effect.sync(() => registry.countRunning()), }); }); diff --git a/apps/server/src/auth/RpcAuthorization.test.ts b/apps/server/src/auth/RpcAuthorization.test.ts index c065272918b4..b04b61d8b502 100644 --- a/apps/server/src/auth/RpcAuthorization.test.ts +++ b/apps/server/src/auth/RpcAuthorization.test.ts @@ -40,6 +40,9 @@ describe("RPC authorization scopes", () => { expect(requiredScopeForRpcMethod(WS_METHODS.serverCancelUpdateDrain)).toBe( AuthOrchestrationOperateScope, ); + expect(requiredScopeForRpcMethod(WS_METHODS.serverClaimUpdateActivation)).toBe( + AuthOrchestrationOperateScope, + ); }); it("allows relay status reads without granting relay installation access", () => { diff --git a/apps/server/src/auth/RpcAuthorization.ts b/apps/server/src/auth/RpcAuthorization.ts index 437a407dc966..f663a041078b 100644 --- a/apps/server/src/auth/RpcAuthorization.ts +++ b/apps/server/src/auth/RpcAuthorization.ts @@ -52,6 +52,7 @@ export const RPC_REQUIRED_SCOPES = { [WS_METHODS.serverGetBackgroundPolicy]: AuthOrchestrationReadScope, [WS_METHODS.serverStartUpdateDrain]: AuthOrchestrationOperateScope, [WS_METHODS.serverCancelUpdateDrain]: AuthOrchestrationOperateScope, + [WS_METHODS.serverClaimUpdateActivation]: AuthOrchestrationOperateScope, [WS_METHODS.serverGetUpdateDrainStatus]: AuthOrchestrationReadScope, [WS_METHODS.cloudGetRelayClientStatus]: AuthRelayReadScope, [WS_METHODS.cloudInstallRelayClient]: AuthRelayWriteScope, diff --git a/apps/server/src/mcp/McpHttpServer.test.ts b/apps/server/src/mcp/McpHttpServer.test.ts index fa2880f9c364..3da417af3a7d 100644 --- a/apps/server/src/mcp/McpHttpServer.test.ts +++ b/apps/server/src/mcp/McpHttpServer.test.ts @@ -1,7 +1,15 @@ import { expect, it } from "@effect/vitest"; import { NodeHttpServer } from "@effect/platform-node"; import * as NodeServices from "@effect/platform-node/NodeServices"; -import { EnvironmentId, PreviewTabId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; +import { + EnvironmentId, + PreviewTabId, + ProviderInstanceId, + ThreadId, + UpdateDrainAdmissionError, + UpdateDrainRequestId, + UpdateDrainTargetVersion, +} from "@t3tools/contracts"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Stream from "effect/Stream"; @@ -11,6 +19,8 @@ import { HttpBody, HttpClient, HttpRouter, HttpServerResponse } from "effect/uns import * as McpHttpServer from "./McpHttpServer.ts"; import * as McpInvocationContext from "./McpInvocationContext.ts"; import * as PreviewAutomationBroker from "./PreviewAutomationBroker.ts"; +import { ActionResume } from "../actionResume/ActionResume.ts"; +import { UpdateDrainAdmission } from "../updateDrain/UpdateDrainAdmission.ts"; const environmentId = EnvironmentId.make("environment-mcp-test"); const threadId = ThreadId.make("thread-mcp-test"); @@ -39,6 +49,51 @@ const TestLayer = McpHttpServer.PreviewToolkitRegistrationLive.pipe( Layer.provideMerge(PreviewAutomationBroker.layer.pipe(Layer.provide(NodeServices.layer))), ); +it.effect("rejects MCP action launch while update drain admission is closed", () => + Effect.gen(function* () { + let launched = false; + const maintenance = new UpdateDrainAdmissionError({ + reason: "update_draining", + requestId: UpdateDrainRequestId.make("mcp-drain"), + targetVersion: UpdateDrainTargetVersion.make("1.2.3"), + message: "LastCode is draining for an update.", + }); + const layer = McpHttpServer.ActionResumeToolkitRegistrationLive.pipe( + Layer.provideMerge(McpServer.McpServer.layer), + Layer.provideMerge( + Layer.mock(ActionResume)({ + runProjectActionAndResume: () => + Effect.sync(() => { + launched = true; + throw new Error("action launch should not run"); + }), + }), + ), + Layer.provideMerge( + Layer.mock(UpdateDrainAdmission)({ + admit: () => Effect.fail(maintenance), + }), + ), + ); + + const result = yield* Effect.gen(function* () { + const server = yield* McpServer.McpServer; + return yield* server + .callTool({ name: "run_project_action_and_resume", arguments: { actionId: "qa" } }) + .pipe( + Effect.provideService(McpInvocationContext.McpInvocationContext, { + ...invocation, + capabilities: new Set(["action-resume"] as const), + }), + Effect.provideService(McpSchema.McpServerClient, client), + ); + }).pipe(Effect.provide(layer)); + + expect(result.isError).toBe(true); + expect(launched).toBe(false); + }), +); + it("normalizes empty successful notification responses to accepted", () => { const notificationResponse = McpHttpServer.normalizeMcpHttpResponse( HttpServerResponse.text("", { status: 200, contentType: "application/json" }), diff --git a/apps/server/src/mcp/toolkits/actionResume/handlers.ts b/apps/server/src/mcp/toolkits/actionResume/handlers.ts index e295232481bc..39decfc6312a 100644 --- a/apps/server/src/mcp/toolkits/actionResume/handlers.ts +++ b/apps/server/src/mcp/toolkits/actionResume/handlers.ts @@ -1,46 +1,69 @@ import { ActionResumeError } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import { ActionResume } from "../../../actionResume/ActionResume.ts"; +import { UpdateDrainAdmission } from "../../../updateDrain/UpdateDrainAdmission.ts"; import * as McpInvocationContext from "../../McpInvocationContext.ts"; import { ActionResumeToolkit } from "./tools.ts"; -const handlers = { - list_project_actions: () => - Effect.gen(function* () { - const invocation = yield* McpInvocationContext.requireMcpCapability("action-resume"); - const service = yield* Effect.serviceOption(ActionResume); - if (Option.isNone(service)) { - return yield* new ActionResumeError({ - reason: "internal_error", - message: "Action resume is unavailable in this server runtime.", - }); - } - const actions = yield* service.value.listProjectActions({ - threadId: invocation.threadId, - providerInstanceId: invocation.providerInstanceId, - }); - return { actions }; - }), - run_project_action_and_resume: ({ actionId }) => - Effect.gen(function* () { - const invocation = yield* McpInvocationContext.requireMcpCapability("action-resume"); - const service = yield* Effect.serviceOption(ActionResume); - if (Option.isNone(service)) { - return yield* new ActionResumeError({ - reason: "internal_error", - message: "Action resume is unavailable in this server runtime.", - }); - } - return yield* service.value.runProjectActionAndResume( - { +const duringDrain = (message: string) => + new ActionResumeError({ + reason: "internal_error", + message, + }); + +const makeHandlers = (admission: UpdateDrainAdmission["Service"]) => + ({ + list_project_actions: () => + Effect.gen(function* () { + const invocation = yield* McpInvocationContext.requireMcpCapability("action-resume"); + const service = yield* Effect.serviceOption(ActionResume); + if (Option.isNone(service)) { + return yield* new ActionResumeError({ + reason: "internal_error", + message: "Action resume is unavailable in this server runtime.", + }); + } + const actions = yield* service.value.listProjectActions({ threadId: invocation.threadId, providerInstanceId: invocation.providerInstanceId, - }, - actionId, - ); - }), -} satisfies Parameters[0]; + }); + return { actions }; + }), + run_project_action_and_resume: ({ actionId }) => + Effect.gen(function* () { + const invocation = yield* McpInvocationContext.requireMcpCapability("action-resume"); + const service = yield* Effect.serviceOption(ActionResume); + if (Option.isNone(service)) { + return yield* new ActionResumeError({ + reason: "internal_error", + message: "Action resume is unavailable in this server runtime.", + }); + } + return yield* admission + .admit( + "action-resume", + service.value.runProjectActionAndResume( + { + threadId: invocation.threadId, + providerInstanceId: invocation.providerInstanceId, + }, + actionId, + ), + ) + .pipe( + Effect.catchTags({ + UpdateDrainAdmissionError: (error) => Effect.fail(duringDrain(error.message)), + UpdateDrainError: (error) => Effect.fail(duringDrain(error.message)), + }), + ); + }), + }) satisfies Parameters[0]; -export const ActionResumeToolkitHandlersLive = ActionResumeToolkit.toLayer(handlers); +export const ActionResumeToolkitHandlersLive = Layer.unwrap( + UpdateDrainAdmission.pipe( + Effect.map((admission) => ActionResumeToolkit.toLayer(makeHandlers(admission))), + ), +); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 5078f9e7c7e0..3b113085e9ae 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -59,6 +59,7 @@ import { import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; import { ProviderCommandReactor } from "../Services/ProviderCommandReactor.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; +import { ProjectionTurnRepository } from "../../persistence/Services/ProjectionTurns.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import * as Clock from "effect/Clock"; import { ServerSettingsService } from "../../serverSettings.ts"; @@ -94,7 +95,10 @@ async function waitFor( describe("ProviderCommandReactor", () => { let runtime: ManagedRuntime.ManagedRuntime< - OrchestrationEngineService | ProviderCommandReactor | ProjectionSnapshotQuery, + | OrchestrationEngineService + | ProviderCommandReactor + | ProjectionSnapshotQuery + | ProjectionTurnRepository, unknown > | null = null; let scope: Scope.Closeable | null = null; @@ -151,6 +155,7 @@ describe("ProviderCommandReactor", () => { readonly requiresNewThreadForModelChange?: boolean; readonly titleRegenerationCompletionDispatchFailures?: number; readonly titleRegenerationBeforeStart?: "one" | "two"; + readonly createSecondThread?: boolean; readonly startSessionEffect?: ( session: ProviderSession, ) => Effect.Effect; @@ -417,12 +422,14 @@ describe("ProviderCommandReactor", () => { Layer.provideMerge(ServerSettingsService.layerTest()), Layer.provideMerge(ServerConfig.layerTest(process.cwd(), baseDir)), Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(SqlitePersistenceMemory), ); runtime = ManagedRuntime.make(layer); const engine = await runtime.runPromise(Effect.service(OrchestrationEngineService)); const snapshotQuery = await runtime.runPromise(Effect.service(ProjectionSnapshotQuery)); const reactor = await runtime.runPromise(Effect.service(ProviderCommandReactor)); + const projectionTurns = await runtime.runPromise(Effect.service(ProjectionTurnRepository)); const runEffect = (effect: Effect.Effect) => runtime!.runPromise(effect); await Effect.runPromise( @@ -451,7 +458,7 @@ describe("ProviderCommandReactor", () => { createdAt: now, }), ); - if (input?.titleRegenerationBeforeStart === "two") { + if (input?.titleRegenerationBeforeStart === "two" || input?.createSecondThread) { await Effect.runPromise( engine.dispatch({ type: "thread.create", @@ -494,6 +501,7 @@ describe("ProviderCommandReactor", () => { return { engine, readModel: () => Effect.runPromise(snapshotQuery.getSnapshot()), + pendingTurnStarts: () => Effect.runPromise(projectionTurns.listPendingTurnStarts()), startSession, sendTurn, interruptTurn, @@ -554,6 +562,67 @@ describe("ProviderCommandReactor", () => { expect(thread?.session?.runtimeMode).toBe("approval-required"); }); + effectIt.effect("clears a queued pending start after its thread is deleted", () => + Effect.gen(function* () { + const releaseFirstStart = yield* Deferred.make(); + let startCount = 0; + const harness = yield* Effect.promise(() => + createHarness({ + createSecondThread: true, + startSessionEffect: (session) => { + startCount += 1; + return startCount === 1 + ? Deferred.await(releaseFirstStart).pipe(Effect.as(session)) + : Effect.succeed(session); + }, + }), + ); + const now = "2026-01-01T00:00:00.000Z"; + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-busy-thread"), + threadId: ThreadId.make("thread-2"), + message: { + messageId: asMessageId("user-message-busy-thread"), + role: "user", + text: "keep the provider worker busy", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }); + yield* Effect.promise(() => waitFor(() => harness.startSession.mock.calls.length === 1)); + + yield* harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-deleted-thread"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-deleted-thread"), + role: "user", + text: "delete before the worker reaches this turn", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }); + yield* harness.engine.dispatch({ + type: "thread.delete", + commandId: CommandId.make("cmd-delete-thread-with-queued-turn"), + threadId: ThreadId.make("thread-1"), + }); + + yield* Deferred.succeed(releaseFirstStart, undefined); + yield* Effect.promise(() => harness.drain()); + + const pendingStarts = yield* Effect.promise(() => harness.pendingTurnStarts()); + expect(pendingStarts.some((pending) => pending.threadId === "thread-1")).toBe(false); + }), + ); + effectIt.effect("projects starting before a slow provider session finishes", () => Effect.gen(function* () { const releaseStart = yield* Deferred.make(); @@ -2859,7 +2928,7 @@ describe("ProviderCommandReactor", () => { }), ); - await Effect.runPromise( + await harness.runEffect( harness.engine.dispatch({ type: "thread.user-input.respond", commandId: CommandId.make("cmd-user-input-respond-stale"), diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 223c78da3a37..a146480d4392 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -29,6 +29,8 @@ import { resolveThreadWorkspaceCwd } from "../../checkpointing/Utils.ts"; import { increment, orchestrationEventsProcessedTotal } from "../../observability/Metrics.ts"; import { ProviderAdapterRequestError } from "../../provider/Errors.ts"; import type { ProviderServiceError } from "../../provider/Errors.ts"; +import { ProjectionTurnRepositoryLive } from "../../persistence/Layers/ProjectionTurns.ts"; +import { ProjectionTurnRepository } from "../../persistence/Services/ProjectionTurns.ts"; import { TextGeneration } from "../../textGeneration/TextGeneration.ts"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import { ProviderRegistry } from "../../provider/Services/ProviderRegistry.ts"; @@ -302,6 +304,7 @@ const make = Effect.gen(function* () { const crypto = yield* Crypto.Crypto; const orchestrationEngine = yield* OrchestrationEngineService; const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; + const projectionTurns = yield* ProjectionTurnRepository; const providerService = yield* ProviderService; const providerRegistry = yield* ProviderRegistry; const gitWorkflow = yield* GitWorkflowService; @@ -1067,6 +1070,9 @@ const make = Effect.gen(function* () { const thread = yield* resolveThread(event.payload.threadId); if (!thread) { + yield* projectionTurns.deletePendingTurnStartByThreadId({ + threadId: event.payload.threadId, + }); return; } @@ -1440,4 +1446,6 @@ const make = Effect.gen(function* () { } satisfies ProviderCommandReactorShape; }); -export const ProviderCommandReactorLive = Layer.effect(ProviderCommandReactor, make); +export const ProviderCommandReactorLive = Layer.effect(ProviderCommandReactor, make).pipe( + Layer.provideMerge(ProjectionTurnRepositoryLive), +); diff --git a/apps/server/src/persistence/Layers/ProjectionTurns.ts b/apps/server/src/persistence/Layers/ProjectionTurns.ts index bd57a4eaa30a..8f1f1555ac82 100644 --- a/apps/server/src/persistence/Layers/ProjectionTurns.ts +++ b/apps/server/src/persistence/Layers/ProjectionTurns.ts @@ -169,6 +169,26 @@ const makeProjectionTurnRepository = Effect.gen(function* () { `, }); + const listPendingProjectionTurns = SqlSchema.findAll({ + Request: Schema.Void, + Result: ProjectionPendingTurnStart, + execute: () => + sql` + SELECT + thread_id AS "threadId", + pending_message_id AS "messageId", + source_proposed_plan_thread_id AS "sourceProposedPlanThreadId", + source_proposed_plan_id AS "sourceProposedPlanId", + requested_at AS "requestedAt" + FROM projection_turns + WHERE turn_id IS NULL + AND state = 'pending' + AND pending_message_id IS NOT NULL + AND checkpoint_turn_count IS NULL + ORDER BY thread_id ASC, requested_at DESC + `, + }); + const listProjectionTurnsByThread = SqlSchema.findAll({ Request: ListProjectionTurnsByThreadInput, Result: ProjectionTurnDbRowSchema, @@ -288,6 +308,16 @@ const makeProjectionTurnRepository = Effect.gen(function* () { ), ); + const listPendingTurnStarts: ProjectionTurnRepositoryShape["listPendingTurnStarts"] = () => + listPendingProjectionTurns(undefined).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionTurnRepository.listPendingTurnStarts:query", + "ProjectionTurnRepository.listPendingTurnStarts:decodeRows", + ), + ), + ); + const deletePendingTurnStartByThreadId: ProjectionTurnRepositoryShape["deletePendingTurnStartByThreadId"] = (input) => clearPendingProjectionTurnsByThread(input).pipe( @@ -341,6 +371,7 @@ const makeProjectionTurnRepository = Effect.gen(function* () { upsertByTurnId, replacePendingTurnStart, getPendingTurnStartByThreadId, + listPendingTurnStarts, deletePendingTurnStartByThreadId, listByThreadId, getByTurnId, diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index d93e57528da4..38bff8f5ca4c 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -56,6 +56,7 @@ import Migration0040 from "./Migrations/040_ProjectionProjectFaviconPath.ts"; import Migration0041 from "./Migrations/041_AuthSessionClientConnection.ts"; import Migration0042 from "./Migrations/042_ProjectionThreadAnnotation.ts"; import Migration0043 from "./Migrations/043_UpdateDrain.ts"; +import Migration0044 from "./Migrations/044_UpdateDrainClaim.ts"; /** * Migration loader with all migrations defined inline. @@ -111,6 +112,7 @@ export const migrationEntries = [ [41, "AuthSessionClientConnection", Migration0041], [42, "ProjectionThreadAnnotation", Migration0042], [43, "UpdateDrain", Migration0043], + [44, "UpdateDrainClaim", Migration0044], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/044_UpdateDrainClaim.test.ts b/apps/server/src/persistence/Migrations/044_UpdateDrainClaim.test.ts new file mode 100644 index 000000000000..2d99b2fb5d62 --- /dev/null +++ b/apps/server/src/persistence/Migrations/044_UpdateDrainClaim.test.ts @@ -0,0 +1,45 @@ +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import { runMigrations } from "../Migrations.ts"; +import * as NodeSqliteClient from "../NodeSqliteClient.ts"; + +const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); + +layer("044_UpdateDrainClaim", (it) => { + it.effect("preserves drain history and accepts one claimed transition", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* runMigrations({ toMigrationInclusive: 43 }); + yield* sql` + INSERT INTO update_drain_events ( + event_id, event_type, command_id, occurred_at, request_id, target_version, status + ) VALUES ( + 'event-start', 'update-drain.started', 'command-start', + '2026-08-21T00:00:00.000Z', 'request-1', '1.2.3', 'draining' + ) + `; + yield* runMigrations({ toMigrationInclusive: 44 }); + yield* sql` + INSERT INTO update_drain_events ( + event_id, event_type, command_id, occurred_at, request_id, target_version, status + ) VALUES ( + 'event-claim', 'update-drain.claimed', 'command-claim', + '2026-08-21T00:01:00.000Z', 'request-1', '1.2.3', 'claimed' + ) + `; + + const events = yield* sql<{ readonly eventType: string; readonly status: string }>` + SELECT event_type AS "eventType", status + FROM update_drain_events + ORDER BY sequence ASC + `; + assert.deepStrictEqual(events, [ + { eventType: "update-drain.started", status: "draining" }, + { eventType: "update-drain.claimed", status: "claimed" }, + ]); + }), + ); +}); diff --git a/apps/server/src/persistence/Migrations/044_UpdateDrainClaim.ts b/apps/server/src/persistence/Migrations/044_UpdateDrainClaim.ts new file mode 100644 index 000000000000..7b6e0ac7a54a --- /dev/null +++ b/apps/server/src/persistence/Migrations/044_UpdateDrainClaim.ts @@ -0,0 +1,55 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + yield* sql`ALTER TABLE update_drain_events RENAME TO update_drain_events_legacy_043`; + yield* sql` + CREATE TABLE update_drain_events ( + sequence INTEGER PRIMARY KEY AUTOINCREMENT, + event_id TEXT NOT NULL UNIQUE, + event_type TEXT NOT NULL CHECK (event_type IN ('update-drain.started', 'update-drain.cancelled', 'update-drain.claimed')), + command_id TEXT NOT NULL UNIQUE, + occurred_at TEXT NOT NULL, + request_id TEXT NOT NULL, + target_version TEXT NOT NULL, + status TEXT NOT NULL CHECK (status IN ('draining', 'cancelled', 'claimed')) + ) + `; + yield* sql` + INSERT INTO update_drain_events ( + sequence, event_id, event_type, command_id, occurred_at, request_id, target_version, status + ) + SELECT + sequence, event_id, event_type, command_id, occurred_at, request_id, target_version, status + FROM update_drain_events_legacy_043 + `; + yield* sql`DROP TABLE update_drain_events_legacy_043`; + + yield* sql`ALTER TABLE update_drain_command_receipts RENAME TO update_drain_command_receipts_legacy_043`; + yield* sql` + CREATE TABLE update_drain_command_receipts ( + command_id TEXT PRIMARY KEY, + command_type TEXT NOT NULL CHECK (command_type IN ('update-drain.start', 'update-drain.cancel', 'update-drain.claim')), + request_id TEXT NOT NULL, + target_version TEXT, + accepted_at TEXT NOT NULL, + result_sequence INTEGER NOT NULL, + status TEXT NOT NULL CHECK (status IN ('accepted', 'rejected')), + error_reason TEXT, + error TEXT + ) + `; + yield* sql` + INSERT INTO update_drain_command_receipts ( + command_id, command_type, request_id, target_version, accepted_at, + result_sequence, status, error_reason, error + ) + SELECT + command_id, command_type, request_id, target_version, accepted_at, + result_sequence, status, error_reason, error + FROM update_drain_command_receipts_legacy_043 + `; + yield* sql`DROP TABLE update_drain_command_receipts_legacy_043`; +}); diff --git a/apps/server/src/persistence/Services/ProjectionTurns.ts b/apps/server/src/persistence/Services/ProjectionTurns.ts index f3d5d5e47061..0732b54104da 100644 --- a/apps/server/src/persistence/Services/ProjectionTurns.ts +++ b/apps/server/src/persistence/Services/ProjectionTurns.ts @@ -128,6 +128,14 @@ export interface ProjectionTurnRepositoryShape { input: GetProjectionPendingTurnStartInput, ) => Effect.Effect, ProjectionRepositoryError>; + /** + * Lists threads with an accepted turn start that has not reached a concrete provider turn. + */ + readonly listPendingTurnStarts: () => Effect.Effect< + ReadonlyArray, + ProjectionRepositoryError + >; + /** * Deletes only pending-start placeholder rows (`turnId = null`) for a thread and leaves concrete turn rows untouched. */ diff --git a/apps/server/src/project/ProjectSetupScriptRunner.test.ts b/apps/server/src/project/ProjectSetupScriptRunner.test.ts index f950b2a9b8a2..174dee31c6b6 100644 --- a/apps/server/src/project/ProjectSetupScriptRunner.test.ts +++ b/apps/server/src/project/ProjectSetupScriptRunner.test.ts @@ -60,6 +60,8 @@ const makeTerminalManagerLayer = ( close: () => Effect.void, subscribe: () => Effect.succeed(() => undefined), subscribeMetadata: () => Effect.succeed(() => undefined), + metadata: Effect.succeed([]), + refreshMetadata: Effect.succeed([]), }); const testLayer = ( diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 8e12d343fa91..76d06bd1a6bd 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -8,6 +8,7 @@ import { AuthAccessTokenType, AuthEnvironmentBootstrapTokenType, AuthTokenExchangeGrantType, + ApprovalRequestId, CommandId, DEFAULT_SERVER_SETTINGS, EnvironmentId, @@ -30,6 +31,7 @@ import { ResolvedKeybindingRule, ThreadId, UpdateDrainRequestId, + UpdateDrainAdmissionError, UpdateDrainTargetVersion, WS_METHODS, WsRpcGroup, @@ -155,6 +157,7 @@ import * as ResourceTelemetry from "./resourceTelemetry/ResourceTelemetry.ts"; import * as UsageService from "./usage/UsageService.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; import * as UpdateDrain from "./updateDrain/UpdateDrain.ts"; +import * as UpdateDrainAdmission from "./updateDrain/UpdateDrainAdmission.ts"; import { UpdateDrainRepositoryLive } from "./persistence/Layers/UpdateDrainRepository.ts"; import * as Data from "effect/Data"; @@ -426,6 +429,7 @@ const buildAppUnderTest = (options?: { desktopTelemetryReceiver?: Partial< DesktopTelemetryReceiver.DesktopTelemetryReceiver["Service"] >; + updateDrainAdmission?: Partial; }; }) => Effect.gen(function* () { @@ -618,6 +622,35 @@ const buildAppUnderTest = (options?: { Layer.provide(UpdateDrainRepositoryLive), Layer.provide(SqlitePersistenceMemory), ); + const updateDrainAdmissionLayer = Layer.effect( + UpdateDrainAdmission.UpdateDrainAdmission, + Effect.gen(function* () { + const drain = yield* UpdateDrain.UpdateDrain; + return UpdateDrainAdmission.UpdateDrainAdmission.of({ + dispatch: drain.dispatch, + claimActivation: (input) => + drain.dispatch({ + type: "update-drain.claim", + commandId: CommandId.make(`update-drain:claim:${input.requestId}`), + requestId: input.requestId, + createdAt: DateTime.formatIso(TEST_EPOCH), + }), + admit: (_kind, effect) => effect, + admitOrElse: (_kind, effect) => effect, + status: drain.status.pipe( + Effect.map((state) => ({ + ...state, + admission: + state.intent !== null && state.intent.status !== "cancelled" + ? ("closed" as const) + : ("open" as const), + blockers: [], + })), + ), + ...options?.layers?.updateDrainAdmission, + }); + }), + ).pipe(Layer.provide(updateDrainLayer)); const servedRoutesLayer = HttpRouter.serve( makeRoutesLayer.pipe(Layer.provide(serviceLauncherClientLayer)), @@ -841,6 +874,7 @@ const buildAppUnderTest = (options?: { ); const appLayer = servedRoutesLayer.pipe( + Layer.provide(updateDrainAdmissionLayer), Layer.provide(updateDrainLayer), Layer.provide(resourceTelemetryLayer), Layer.provide(UsageService.layerTest), @@ -4504,6 +4538,170 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); + it.effect("closes work creation while drain controls remain operable", () => + Effect.gen(function* () { + const requestId = UpdateDrainRequestId.make("rpc-update-draining"); + const targetVersion = UpdateDrainTargetVersion.make("0.0.36-nightly.1"); + const maintenance = new UpdateDrainAdmissionError({ + reason: "update_draining", + requestId, + targetVersion, + message: "LastCode is draining for an update.", + }); + const terminalSnapshot = { + threadId: defaultThreadId, + terminalId: "existing", + cwd: "/tmp/project", + worktreePath: null, + status: "running" as const, + pid: 123, + history: "ready\n", + exitCode: null, + exitSignal: null, + label: "Existing", + updatedAt: "2026-08-21T00:00:00.000Z", + }; + yield* buildAppUnderTest({ + layers: { + updateDrainAdmission: { + admit: () => Effect.fail(maintenance), + admitOrElse: (_kind, _effect, whenClosed) => whenClosed, + }, + terminalManager: { + attachStream: (_input, listener) => + listener({ type: "snapshot", snapshot: terminalSnapshot }).pipe( + Effect.as(() => undefined), + ), + close: () => Effect.void, + }, + }, + }); + + const wsUrl = yield* getWsServerUrl("/ws"); + const results = yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + const blocked = yield* Effect.all([ + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.turn.start", + commandId: CommandId.make("blocked-turn"), + threadId: defaultThreadId, + message: { + messageId: MessageId.make("blocked-message"), + role: "user", + text: "wait for maintenance", + attachments: [], + }, + modelSelection: defaultModelSelection, + runtimeMode: "full-access", + interactionMode: "default", + createdAt: "2026-08-21T00:00:00.000Z", + }).pipe(Effect.result), + client[WS_METHODS.terminalOpen]({ + threadId: defaultThreadId, + terminalId: "new", + cwd: "/tmp/project", + }).pipe(Effect.result), + client[WS_METHODS.terminalWrite]({ + threadId: defaultThreadId, + terminalId: "existing", + data: "make test\n", + }).pipe(Effect.result), + client[WS_METHODS.terminalRestart]({ + threadId: defaultThreadId, + terminalId: "existing", + cwd: "/tmp/project", + cols: 120, + rows: 40, + }).pipe(Effect.result), + client[WS_METHODS.gitPreparePullRequestThread]({ + cwd: "/tmp/project", + reference: "1", + mode: "worktree", + threadId: defaultThreadId, + }).pipe(Effect.result), + ]); + + const attached = yield* client[WS_METHODS.terminalAttach]({ + threadId: defaultThreadId, + terminalId: "existing", + }).pipe(Stream.runHead); + yield* client[WS_METHODS.terminalClose]({ + threadId: defaultThreadId, + terminalId: "existing", + }); + const controls = yield* Effect.all([ + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.turn.interrupt", + commandId: CommandId.make("interrupt-during-drain"), + threadId: defaultThreadId, + createdAt: "2026-08-21T00:00:01.000Z", + }), + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.approval.respond", + commandId: CommandId.make("approval-during-drain"), + threadId: defaultThreadId, + requestId: ApprovalRequestId.make("approval-1"), + decision: "decline", + createdAt: "2026-08-21T00:00:02.000Z", + }), + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.user-input.respond", + commandId: CommandId.make("input-during-drain"), + threadId: defaultThreadId, + requestId: ApprovalRequestId.make("input-1"), + answers: { answer: "stop" }, + createdAt: "2026-08-21T00:00:03.000Z", + }), + ]); + return { blocked, attached, controls }; + }), + ), + ); + + assert.ok(results.blocked.every((result) => result._tag === "Failure")); + assert.ok( + results.blocked.every( + (result) => + result._tag === "Failure" && result.failure._tag === "UpdateDrainAdmissionError", + ), + ); + assert.isTrue(Option.isSome(results.attached)); + assert.equal(results.controls.length, 3); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + + it.effect("routes an activation claim using the drain request id", () => + Effect.gen(function* () { + yield* buildAppUnderTest(); + const wsUrl = yield* getWsServerUrl("/ws"); + const requestId = UpdateDrainRequestId.make("rpc-update-claim"); + const targetVersion = UpdateDrainTargetVersion.make("0.0.36-nightly.2"); + const result = yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + yield* client[WS_METHODS.serverStartUpdateDrain]({ + commandId: CommandId.make("rpc-update-claim-start"), + requestId, + targetVersion, + }); + const claim = yield* client[WS_METHODS.serverClaimUpdateActivation]({ requestId }); + const status = yield* client[WS_METHODS.serverGetUpdateDrainStatus]({}); + return { claim, status }; + }), + ), + ); + + assert.equal(result.claim.commandType, "update-drain.claim"); + assert.deepStrictEqual(result.status.intent, { + requestId, + targetVersion, + status: "claimed", + }); + assert.equal(result.status.admission, "closed"); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + it.effect("shares one preview automation broker across websocket sessions", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index f15a4c586665..81046e45bd68 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -110,6 +110,7 @@ import * as ResourceMonitorBinary from "./resourceTelemetry/ResourceMonitorBinar import * as ResourceTelemetry from "./resourceTelemetry/ResourceTelemetry.ts"; import * as UsageService from "./usage/UsageService.ts"; import * as UpdateDrain from "./updateDrain/UpdateDrain.ts"; +import * as UpdateDrainAdmission from "./updateDrain/UpdateDrainAdmission.ts"; import { UpdateDrainRepositoryLive } from "./persistence/Layers/UpdateDrainRepository.ts"; import { OrchestrationLayerLive } from "./orchestration/runtimeLayer.ts"; import { @@ -425,9 +426,12 @@ const RuntimeCoreDependenciesBaseLive = ReactorLayerLive.pipe( ), ); -const RuntimeCoreDependenciesLive = Layer.merge( - RuntimeCoreDependenciesBaseLive, - ActionResume.layer.pipe(Layer.provide(RuntimeCoreDependenciesBaseLive)), +const RuntimeCoreDependenciesWithDrainAdmissionLive = UpdateDrainAdmission.layer.pipe( + Layer.provideMerge(RuntimeCoreDependenciesBaseLive), +); + +const RuntimeCoreDependenciesLive = ActionResume.layer.pipe( + Layer.provideMerge(RuntimeCoreDependenciesWithDrainAdmissionLive), ); const RuntimeDependenciesLive = RuntimeCoreDependenciesLive.pipe( diff --git a/apps/server/src/terminal/Manager.test.ts b/apps/server/src/terminal/Manager.test.ts index 09de9bafe9ee..e76a9018c872 100644 --- a/apps/server/src/terminal/Manager.test.ts +++ b/apps/server/src/terminal/Manager.test.ts @@ -42,6 +42,7 @@ class FakePtyProcess implements PtyAdapter.PtyProcess { readonly pid: number; writeFailure: unknown | undefined; resizeFailure: unknown | undefined; + killFailure: unknown | undefined; private readonly dataListeners = new Set<(data: string) => void>(); private readonly exitListeners = new Set<(event: PtyAdapter.PtyExitEvent) => void>(); killed = false; @@ -67,6 +68,9 @@ class FakePtyProcess implements PtyAdapter.PtyProcess { kill(signal?: string): void { this.killed = true; this.killSignals.push(signal); + if (this.killFailure !== undefined) { + throw this.killFailure; + } } onData(callback: (data: string) => void): () => void { @@ -311,6 +315,7 @@ it.layer( rows: 40, }, (event) => Ref.update(attachEvents, (events) => [...events, event]), + false, ); yield* Effect.addFinalizer(() => Effect.sync(unsubscribe)); @@ -323,6 +328,22 @@ it.layer( }), ); + it.effect("refuses to create a missing terminal for read-only attach", () => + Effect.gen(function* () { + const { manager, ptyAdapter } = yield* createManager(); + const result = yield* manager + .attachStream(openInput(), () => Effect.void, false) + .pipe(Effect.result); + + assert.equal(result._tag, "Failure"); + assert.equal( + result._tag === "Failure" ? result.failure._tag : null, + "TerminalSessionLookupError", + ); + expect(ptyAdapter.spawnInputs).toHaveLength(0); + }), + ); + it.effect("keeps attach streams live when a terminal id is closed and reopened", () => Effect.gen(function* () { const { manager, ptyAdapter } = yield* createManager(); @@ -1323,6 +1344,55 @@ it.layer( }).pipe(Effect.provide(TestClock.layer())), ); + it.effect("keeps a closing terminal in blocker metadata until kill escalation finishes", () => + Effect.gen(function* () { + const { manager } = yield* createManager(5, { processKillGraceMs: 10 }); + yield* manager.open(openInput()); + + const closeFiber = yield* manager.close({ threadId: "thread-1" }).pipe(Effect.forkScoped); + yield* Effect.yieldNow; + + expect(yield* manager.refreshMetadata).toEqual([ + expect.objectContaining({ + threadId: "thread-1", + terminalId: DEFAULT_TERMINAL_ID, + status: "running", + hasRunningSubprocess: true, + }), + ]); + + yield* TestClock.adjust("10 millis"); + yield* Fiber.join(closeFiber); + yield* Effect.yieldNow; + + expect(yield* manager.metadata).toEqual([]); + expect(yield* manager.refreshMetadata).toEqual([]); + }).pipe(Effect.provide(TestClock.layer())), + ); + + it.effect("keeps a closing terminal blocked when process signaling fails", () => + Effect.gen(function* () { + const { manager, ptyAdapter } = yield* createManager(); + yield* manager.open(openInput()); + const process = ptyAdapter.processes[0]; + expect(process).toBeDefined(); + if (!process) return; + process.killFailure = new Error("simulated signal failure"); + + yield* manager.close({ threadId: "thread-1" }); + + expect(process.killSignals).toEqual(["SIGTERM"]); + expect(yield* manager.refreshMetadata).toEqual([ + expect.objectContaining({ + threadId: "thread-1", + terminalId: DEFAULT_TERMINAL_ID, + status: "running", + hasRunningSubprocess: true, + }), + ]); + }), + ); + it.effect("publishes closed events when terminals are explicitly closed", () => Effect.gen(function* () { const { manager, getEvents } = yield* createManager(); diff --git a/apps/server/src/terminal/Manager.ts b/apps/server/src/terminal/Manager.ts index 0f75c94be518..8fd3afb5424b 100644 --- a/apps/server/src/terminal/Manager.ts +++ b/apps/server/src/terminal/Manager.ts @@ -150,6 +150,7 @@ export class TerminalManager extends Context.Service< readonly attachStream: ( input: TerminalAttachInput, listener: (event: TerminalAttachStreamEvent) => Effect.Effect, + startIfNeeded?: boolean, ) => Effect.Effect<() => void, TerminalError>; /** @@ -203,6 +204,12 @@ export class TerminalManager extends Context.Service< readonly subscribeMetadata: ( listener: (event: TerminalMetadataStreamEvent) => Effect.Effect, ) => Effect.Effect<() => void>; + + /** Read current terminal metadata without subscribing to runtime events. */ + readonly metadata: Effect.Effect>; + + /** Refresh subprocess activity, then read current terminal metadata. */ + readonly refreshMetadata: Effect.Effect>; } >()("t3/terminal/Manager/TerminalManager") {} @@ -307,6 +314,7 @@ type DrainProcessEventAction = interface TerminalManagerState { sessions: Map; killFibers: Map>; + terminatingProcesses: Map; } function truncateTerminalWireLabel(value: string): string { @@ -1229,6 +1237,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const managerStateRef = yield* SynchronizedRef.make({ sessions: new Map(), killFibers: new Map(), + terminatingProcesses: new Map(), }); const threadLocksRef = yield* SynchronizedRef.make(new Map()); const terminalEventListeners = new Set<(event: TerminalEvent) => Effect.Effect>(); @@ -1341,12 +1350,12 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func ), ); if (!terminated) { - return; + return false; } yield* Effect.sleep(processKillGraceMs); - yield* Effect.try({ + return yield* Effect.try({ try: () => process.kill("SIGKILL"), catch: (cause) => new TerminalProcessSignalError({ @@ -1355,13 +1364,14 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func terminalPid: process.pid, }), }).pipe( + Effect.as(true), Effect.catch((error) => Effect.logWarning("failed to force-kill terminal process", { threadId, terminalId, signal: "SIGKILL", cause: error, - }), + }).pipe(Effect.as(false)), ), ); }); @@ -1372,6 +1382,19 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func terminalId: string, ) { const fiber = yield* runKillEscalation(process, threadId, terminalId).pipe( + Effect.tap((completed) => + completed + ? modifyManagerState((state) => { + if (!state.terminatingProcesses.has(process)) { + return [undefined, state] as const; + } + const terminatingProcesses = new Map(state.terminatingProcesses); + terminatingProcesses.delete(process); + return [undefined, { ...state, terminatingProcesses }] as const; + }) + : Effect.void, + ), + Effect.asVoid, Effect.ensuring( modifyManagerState((state) => { if (!state.killFibers.has(process)) { @@ -1787,6 +1810,12 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const updatedAt = yield* nowIso; yield* modifyManagerState((state) => { + const terminatingProcesses = new Map(state.terminatingProcesses); + terminatingProcesses.set(process, { + ...summary(session), + status: "running", + hasRunningSubprocess: true, + }); cleanupProcessHandles(session); session.process = null; session.pid = null; @@ -1798,7 +1827,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func session.pendingProcessEventIndex = 0; session.processEventDrainRunning = false; session.updatedAt = updatedAt; - return [undefined, state] as const; + return [undefined, { ...state, terminatingProcesses }] as const; }); yield* clearKillFiber(process); @@ -2162,13 +2191,17 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func yield* Effect.addFinalizer(() => Effect.gen(function* () { - const sessions = yield* modifyManagerState( + const { sessions, terminatingProcesses } = yield* modifyManagerState( (state) => [ - [...state.sessions.values()], + { + sessions: [...state.sessions.values()], + terminatingProcesses: [...state.terminatingProcesses.entries()], + }, { ...state, sessions: new Map(), + terminatingProcesses: new Map(), }, ] as const, ); @@ -2186,6 +2219,14 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func concurrency: "unbounded", discard: true, }); + yield* Effect.forEach( + terminatingProcesses, + ([process, terminal]) => + clearKillFiber(process).pipe( + Effect.andThen(runKillEscalation(process, terminal.threadId, terminal.terminalId)), + ), + { concurrency: "unbounded", discard: true }, + ); }).pipe(Effect.ignoreCause({ log: true })), ); @@ -2316,7 +2357,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const open: TerminalManager["Service"]["open"] = (input) => withThreadLock(input.threadId, openLocked(input)); - const openOrAttachForStream = (input: TerminalAttachInput) => + const openOrAttachForStream = (input: TerminalAttachInput, startIfNeeded = true) => withThreadLock( input.threadId, Effect.gen(function* () { @@ -2324,7 +2365,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const existing = yield* getSession(input.threadId, terminalId); if (Option.isNone(existing)) { - if (!input.cwd) { + if (!input.cwd || !startIfNeeded) { return yield* new TerminalSessionLookupError({ threadId: input.threadId, terminalId, @@ -2342,7 +2383,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func const targetCols = input.cols ?? session.cols; const targetRows = input.rows ?? session.rows; - if (!session.process && input.cwd && input.restartIfNotRunning === true) { + if (!session.process && input.cwd && input.restartIfNotRunning === true && startIfNeeded) { return yield* openLocked({ ...input, terminalId, @@ -2380,6 +2421,27 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func ), ); + const readDrainTerminalMetadata = () => + readManagerState.pipe( + Effect.map((state) => { + const terminals = new Map( + [...state.sessions.values()].map((session) => [ + toSessionKey(session.threadId, session.terminalId), + summary(session), + ]), + ); + for (const terminal of state.terminatingProcesses.values()) { + terminals.set(toSessionKey(terminal.threadId, terminal.terminalId), terminal); + } + return [...terminals.values()].sort( + (left, right) => + right.updatedAt.localeCompare(left.updatedAt) || + left.threadId.localeCompare(right.threadId) || + left.terminalId.localeCompare(right.terminalId), + ); + }), + ); + const readTerminalMetadata = (input: { readonly threadId: string; readonly terminalId: string; @@ -2396,7 +2458,11 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func }; }); - const attachStream: TerminalManager["Service"]["attachStream"] = (input, listener) => { + const attachStream: TerminalManager["Service"]["attachStream"] = ( + input, + listener, + startIfNeeded, + ) => { let unsubscribe: (() => void) | null = null; return Effect.gen(function* () { @@ -2417,7 +2483,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func return attachEvent ? listener(attachEvent) : Effect.void; }); - const initialSnapshot = yield* openOrAttachForStream(input); + const initialSnapshot = yield* openOrAttachForStream(input, startIfNeeded); yield* listener({ type: "snapshot", @@ -2716,6 +2782,8 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func close, subscribe, subscribeMetadata, + metadata: readAllTerminalMetadata(), + refreshMetadata: pollSubprocessActivity().pipe(Effect.andThen(readDrainTerminalMetadata())), }); }); diff --git a/apps/server/src/updateDrain/DrainState.ts b/apps/server/src/updateDrain/DrainState.ts index 686c2df5726d..57e0a4854034 100644 --- a/apps/server/src/updateDrain/DrainState.ts +++ b/apps/server/src/updateDrain/DrainState.ts @@ -51,6 +51,14 @@ export function decideUpdateDrainCommand( }), ); } + if (state.intent?.status === "claimed") { + return Effect.fail( + new UpdateDrainError({ + reason: "activation_claimed", + message: `Update drain '${state.intent.requestId}' has already been claimed for activation.`, + }), + ); + } if (usedRequestIds.has(command.requestId)) { return Effect.fail( new UpdateDrainError({ @@ -70,6 +78,42 @@ export function decideUpdateDrainCommand( }); } + if (command.type === "update-drain.claim") { + if (state.intent === null || state.intent.status === "cancelled") { + return Effect.fail( + new UpdateDrainError({ + reason: "no_active_drain", + message: "There is no active update drain to claim.", + }), + ); + } + if (state.intent.requestId !== command.requestId) { + return Effect.fail( + new UpdateDrainError({ + reason: "request_mismatch", + message: `Update drain '${command.requestId}' does not match active request '${state.intent.requestId}'.`, + }), + ); + } + if (state.intent.status === "claimed") { + return Effect.fail( + new UpdateDrainError({ + reason: "activation_claimed", + message: `Update drain '${command.requestId}' has already been claimed for activation.`, + }), + ); + } + return Effect.succeed({ + type: "update-drain.claimed", + eventId: eventIdFor(command), + commandId: command.commandId, + occurredAt: command.createdAt, + requestId: command.requestId, + targetVersion: state.intent.targetVersion, + status: "claimed", + }); + } + if (state.intent === null) { return Effect.fail( new UpdateDrainError({ @@ -94,6 +138,14 @@ export function decideUpdateDrainCommand( }), ); } + if (state.intent.status === "claimed") { + return Effect.fail( + new UpdateDrainError({ + reason: "activation_claimed", + message: `Update drain '${state.intent.requestId}' has already been claimed for activation.`, + }), + ); + } if (state.intent.requestId !== command.requestId) { return Effect.fail( new UpdateDrainError({ diff --git a/apps/server/src/updateDrain/UpdateDrain.test.ts b/apps/server/src/updateDrain/UpdateDrain.test.ts index 4205a2607d4e..21cababfafe0 100644 --- a/apps/server/src/updateDrain/UpdateDrain.test.ts +++ b/apps/server/src/updateDrain/UpdateDrain.test.ts @@ -113,6 +113,38 @@ tests("UpdateDrain", (it) => { ); }); +it.layer(testLayer)("UpdateDrain activation claim", (it) => { + it.effect("restores an activation claim and replays its receipt after restart", () => + Effect.gen(function* () { + const drain = yield* UpdateDrain; + const repository = yield* UpdateDrainRepository; + yield* drain.dispatch({ + type: "update-drain.start", + commandId: CommandId.make("claim-start"), + requestId, + targetVersion, + createdAt: startedAt, + }); + const claimCommand = { + type: "update-drain.claim" as const, + commandId: CommandId.make(`update-drain:claim:${requestId}`), + requestId, + createdAt: "2026-08-21T00:01:00.000Z", + }; + const claimed = yield* drain.dispatch(claimCommand); + + const restored = yield* makeUpdateDrain().pipe( + Effect.provideService(UpdateDrainRepository, repository), + ); + assert.deepStrictEqual(yield* restored.status, { + sequence: 2, + intent: { requestId, targetVersion, status: "claimed" }, + }); + assert.deepStrictEqual(yield* restored.dispatch(claimCommand), claimed); + }), + ); +}); + it.layer(testLayer)("UpdateDrain request history", (it) => { it.effect("never reuses a cancelled request id after later drains and restart", () => Effect.gen(function* () { diff --git a/apps/server/src/updateDrain/UpdateDrainAdmission.test.ts b/apps/server/src/updateDrain/UpdateDrainAdmission.test.ts new file mode 100644 index 000000000000..3ab6a12cbd88 --- /dev/null +++ b/apps/server/src/updateDrain/UpdateDrainAdmission.test.ts @@ -0,0 +1,285 @@ +import { + CommandId, + MessageId, + ProjectId, + ProviderInstanceId, + ThreadId, + TurnId, + UpdateDrainRequestId, + UpdateDrainTargetVersion, + type OrchestrationShellSnapshot, + type TerminalSummary, +} from "@t3tools/contracts"; +import { assert, it } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Ref from "effect/Ref"; + +import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts"; +import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; +import { ProjectionTurnRepositoryLive } from "../persistence/Layers/ProjectionTurns.ts"; +import { UpdateDrainRepositoryLive } from "../persistence/Layers/UpdateDrainRepository.ts"; +import { ProjectionTurnRepository } from "../persistence/Services/ProjectionTurns.ts"; +import { TerminalManager } from "../terminal/Manager.ts"; +import { layer as updateDrainLayer } from "./UpdateDrain.ts"; +import { makeUpdateDrainAdmission } from "./UpdateDrainAdmission.ts"; + +const requestId = UpdateDrainRequestId.make("update-1"); +const targetVersion = UpdateDrainTargetVersion.make("1.2.3"); +const threadId = ThreadId.make("thread-1"); +const turnId = TurnId.make("turn-1"); +const now = "2026-08-21T00:00:00.000Z"; + +const emptyShell = (): OrchestrationShellSnapshot => ({ + snapshotSequence: 0, + projects: [], + threads: [], + updatedAt: now, +}); + +const busyShell = (): OrchestrationShellSnapshot => ({ + snapshotSequence: 1, + projects: [], + updatedAt: now, + threads: [ + { + id: threadId, + projectId: ProjectId.make("project-1"), + title: "Busy thread", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + latestTurn: { + turnId, + state: "running", + requestedAt: now, + startedAt: now, + completedAt: null, + assistantMessageId: null, + }, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + session: { + threadId, + status: "running", + providerName: "codex", + runtimeMode: "full-access", + activeTurnId: turnId, + lastError: null, + updatedAt: now, + }, + latestUserMessageAt: now, + hasPendingApprovals: true, + hasPendingUserInput: true, + hasActionableProposedPlan: false, + backgroundLiveness: "working", + }, + ], +}); + +const busyTerminal = (): TerminalSummary => ({ + threadId, + terminalId: "terminal-1", + cwd: "/tmp/project", + worktreePath: null, + status: "running", + pid: 123, + exitCode: null, + exitSignal: null, + hasRunningSubprocess: true, + label: "tests", + updatedAt: now, +}); + +const durableLayer = updateDrainLayer.pipe( + Layer.provide(UpdateDrainRepositoryLive), + Layer.provide(SqlitePersistenceMemory), +); + +const makeHarness = Effect.fn("UpdateDrainAdmissionTest.makeHarness")(function* () { + const shell = yield* Ref.make(emptyShell()); + const terminals = yield* Ref.make>([]); + const dependencies = Layer.mergeAll( + durableLayer, + ProjectionTurnRepositoryLive.pipe(Layer.provide(SqlitePersistenceMemory)), + Layer.mock(ProjectionSnapshotQuery)({ getShellSnapshot: () => Ref.get(shell) }), + Layer.mock(TerminalManager)({ + metadata: Effect.succeed([]), + refreshMetadata: Ref.get(terminals), + }), + ); + return { shell, terminals, dependencies } as const; +}); + +it.effect("orders work admission before closing the drain", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + yield* Effect.gen(function* () { + const admission = yield* makeUpdateDrainAdmission(); + const entered = yield* Deferred.make(); + const release = yield* Deferred.make(); + const timeline = yield* Ref.make>([]); + const append = (value: string) => Ref.update(timeline, (values) => [...values, value]); + + const work = yield* admission + .admit( + "thread-turn", + append("work-start").pipe( + Effect.andThen(Deferred.succeed(entered, undefined)), + Effect.andThen(Deferred.await(release)), + Effect.andThen(append("work-admitted")), + ), + ) + .pipe(Effect.forkChild); + yield* Deferred.await(entered); + const close = yield* admission + .dispatch({ + type: "update-drain.start", + commandId: CommandId.make("start-1"), + requestId, + targetVersion, + createdAt: now, + }) + .pipe( + Effect.tap(() => append("drain-closed")), + Effect.forkChild, + ); + + yield* Effect.yieldNow; + yield* Deferred.succeed(release, undefined); + yield* Fiber.join(work); + yield* Fiber.join(close); + assert.deepStrictEqual(yield* Ref.get(timeline), [ + "work-start", + "work-admitted", + "drain-closed", + ]); + }).pipe(Effect.provide(harness.dependencies)); + }), +); + +it.effect("ignores pending starts left behind by a previous server lifetime", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + + yield* Effect.gen(function* () { + const projectionTurns = yield* ProjectionTurnRepository; + yield* projectionTurns.replacePendingTurnStart({ + threadId, + messageId: MessageId.make("message-stale-pending-turn"), + sourceProposedPlanThreadId: null, + sourceProposedPlanId: null, + requestedAt: now, + }); + const admission = yield* makeUpdateDrainAdmission(); + yield* admission.dispatch({ + type: "update-drain.start", + commandId: CommandId.make("start-after-restart"), + requestId, + targetVersion, + createdAt: now, + }); + + assert.deepStrictEqual((yield* admission.status).blockers, []); + assert.equal( + (yield* admission.claimActivation({ requestId })).commandType, + "update-drain.claim", + ); + }).pipe(Effect.provide(harness.dependencies)); + }), +); + +it.effect("reports only execution blockers and atomically claims when they clear", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + yield* Ref.set(harness.shell, busyShell()); + yield* Ref.set(harness.terminals, [busyTerminal()]); + + yield* Effect.gen(function* () { + const admission = yield* makeUpdateDrainAdmission(); + yield* admission.dispatch({ + type: "update-drain.start", + commandId: CommandId.make("start-2"), + requestId, + targetVersion, + createdAt: now, + }); + + const status = yield* admission.status; + assert.deepStrictEqual( + status.blockers.map((blocker) => blocker.type), + ["terminal-process", "thread-background", "thread-turn"], + ); + assert.ok(!status.blockers.some((blocker) => "approval" in blocker)); + assert.equal( + (yield* Effect.result(admission.claimActivation({ requestId })))._tag, + "Failure", + ); + + yield* Ref.set(harness.shell, emptyShell()); + yield* Ref.set(harness.terminals, []); + const claimed = yield* admission.claimActivation({ requestId }); + assert.equal(claimed.commandType, "update-drain.claim"); + assert.deepStrictEqual(yield* admission.status, { + sequence: 2, + intent: { requestId, targetVersion, status: "claimed" }, + admission: "closed", + blockers: [], + }); + + const blockedWrite = yield* Effect.result(admission.admit("terminal-write", Effect.void)); + assert.equal(blockedWrite._tag, "Failure"); + }).pipe(Effect.provide(harness.dependencies)); + }), +); + +it.effect("keeps accepted provider starts blocking until they reach the session projection", () => + Effect.gen(function* () { + const harness = yield* makeHarness(); + + yield* Effect.gen(function* () { + const admission = yield* makeUpdateDrainAdmission(); + const projectionTurns = yield* ProjectionTurnRepository; + yield* projectionTurns.replacePendingTurnStart({ + threadId, + messageId: MessageId.make("message-pending-turn"), + sourceProposedPlanThreadId: null, + sourceProposedPlanId: null, + requestedAt: now, + }); + yield* admission.dispatch({ + type: "update-drain.start", + commandId: CommandId.make("start-pending-turn"), + requestId, + targetVersion, + createdAt: now, + }); + + assert.deepStrictEqual((yield* admission.status).blockers, [ + { + type: "thread-turn", + threadId, + turnId: null, + status: "starting", + }, + ]); + assert.equal( + (yield* Effect.result(admission.claimActivation({ requestId })))._tag, + "Failure", + ); + + yield* projectionTurns.deletePendingTurnStartByThreadId({ threadId }); + assert.equal( + (yield* admission.claimActivation({ requestId })).commandType, + "update-drain.claim", + ); + }).pipe(Effect.provide(harness.dependencies)); + }), +); diff --git a/apps/server/src/updateDrain/UpdateDrainAdmission.ts b/apps/server/src/updateDrain/UpdateDrainAdmission.ts new file mode 100644 index 000000000000..028410650ae1 --- /dev/null +++ b/apps/server/src/updateDrain/UpdateDrainAdmission.ts @@ -0,0 +1,241 @@ +import { + CommandId, + ThreadId, + UpdateDrainAdmissionError, + type UpdateDrainBlocker, + type UpdateDrainCancelCommand, + type UpdateDrainClaimInput, + type UpdateDrainCommandReceipt, + UpdateDrainError, + type UpdateDrainStartCommand, + type UpdateDrainStatus, +} from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Semaphore from "effect/Semaphore"; + +import { ProjectionSnapshotQuery } from "../orchestration/Services/ProjectionSnapshotQuery.ts"; +import { ProjectionTurnRepository } from "../persistence/Services/ProjectionTurns.ts"; +import { TerminalManager } from "../terminal/Manager.ts"; +import { UpdateDrain } from "./UpdateDrain.ts"; + +export const UpdateDrainAdmissionKind = [ + "thread-turn", + "terminal-open", + "terminal-restart", + "terminal-write", + "action-resume", + "setup-script", +] as const; +export type UpdateDrainAdmissionKind = (typeof UpdateDrainAdmissionKind)[number]; + +type UpdateDrainLifecycleCommand = UpdateDrainStartCommand | UpdateDrainCancelCommand; + +export interface UpdateDrainAdmissionShape { + readonly dispatch: ( + command: UpdateDrainLifecycleCommand, + ) => Effect.Effect; + readonly claimActivation: ( + input: UpdateDrainClaimInput, + ) => Effect.Effect; + readonly status: Effect.Effect; + readonly admit: ( + kind: UpdateDrainAdmissionKind, + effect: Effect.Effect, + ) => Effect.Effect; + readonly admitOrElse: ( + kind: UpdateDrainAdmissionKind, + effect: Effect.Effect, + whenClosed: Effect.Effect, + ) => Effect.Effect; +} + +export class UpdateDrainAdmission extends Context.Service< + UpdateDrainAdmission, + UpdateDrainAdmissionShape +>()("t3/updateDrain/UpdateDrainAdmission") {} + +function internalError(_cause: unknown) { + return new UpdateDrainError({ + reason: "internal_error", + message: "Failed to derive current update drain blockers.", + }); +} + +function pendingTurnStartKey(threadId: string, messageId: string) { + return `${threadId}\u0000${messageId}`; +} + +export const makeUpdateDrainAdmission = Effect.fn("makeUpdateDrainAdmission")(function* () { + const drain = yield* UpdateDrain; + const projections = yield* ProjectionSnapshotQuery; + const projectionTurns = yield* ProjectionTurnRepository; + const terminals = yield* TerminalManager; + const mutex = yield* Semaphore.make(1); + // The provider event stream is hot, so accepted starts from a previous + // server lifetime cannot be resumed. Keep their exact identities out of the + // live blocker set; a new start replaces the row with a new message id. + const stalePendingTurnStartKeys = new Set( + (yield* projectionTurns.listPendingTurnStarts().pipe(Effect.mapError(internalError))).map( + (pending) => pendingTurnStartKey(pending.threadId, pending.messageId), + ), + ); + + const currentBlockers = Effect.fn("UpdateDrainAdmission.currentBlockers")(function* () { + // Read pending starts first. If one transitions while the shell snapshot is + // read, that newer snapshot contains the starting/running provider state. + const pendingTurnStarts = yield* projectionTurns + .listPendingTurnStarts() + .pipe(Effect.mapError(internalError)); + const [shell, terminalState] = yield* Effect.all([ + projections.getShellSnapshot().pipe(Effect.mapError(internalError)), + terminals.refreshMetadata, + ]); + const blockers: UpdateDrainBlocker[] = []; + const blockedTurnThreadIds = new Set(); + + for (const thread of shell.threads) { + if (thread.latestTurn?.state === "running") { + blockedTurnThreadIds.add(thread.id); + blockers.push({ + type: "thread-turn", + threadId: thread.id, + turnId: thread.latestTurn.turnId, + status: "running", + }); + } else if (thread.session?.status === "starting" || thread.session?.status === "running") { + blockedTurnThreadIds.add(thread.id); + blockers.push({ + type: "thread-turn", + threadId: thread.id, + turnId: thread.session.activeTurnId, + status: thread.session.status, + }); + } + + if (thread.backgroundLiveness) { + blockers.push({ + type: "thread-background", + threadId: thread.id, + status: thread.backgroundLiveness, + }); + } + } + + for (const pending of pendingTurnStarts) { + const pendingThreadId = pending.threadId; + if (stalePendingTurnStartKeys.has(pendingTurnStartKey(pendingThreadId, pending.messageId))) { + continue; + } + if (blockedTurnThreadIds.has(pendingThreadId)) continue; + blockers.push({ + type: "thread-turn", + threadId: pendingThreadId, + turnId: null, + status: "starting", + }); + } + + for (const terminal of terminalState) { + if (terminal.status !== "starting" && !terminal.hasRunningSubprocess) continue; + blockers.push({ + type: "terminal-process", + threadId: ThreadId.make(terminal.threadId), + terminalId: terminal.terminalId, + label: terminal.label, + status: terminal.status === "starting" ? "starting" : "running", + }); + } + + return blockers.sort((left, right) => { + const threadOrder = left.threadId.localeCompare(right.threadId); + if (threadOrder !== 0) return threadOrder; + const typeOrder = left.type.localeCompare(right.type); + if (typeOrder !== 0) return typeOrder; + if (left.type === "terminal-process" && right.type === "terminal-process") { + return left.terminalId.localeCompare(right.terminalId); + } + return 0; + }); + }); + + const statusUnlocked = Effect.fn("UpdateDrainAdmission.statusUnlocked")(function* () { + const durable = yield* drain.status; + if (durable.intent === null || durable.intent.status === "cancelled") { + return { + ...durable, + admission: "open" as const, + blockers: [], + } satisfies UpdateDrainStatus; + } + return { + ...durable, + admission: "closed" as const, + blockers: yield* currentBlockers(), + } satisfies UpdateDrainStatus; + }); + + const dispatch: UpdateDrainAdmissionShape["dispatch"] = (command) => + mutex.withPermits(1)(drain.dispatch(command)); + + const claimActivation: UpdateDrainAdmissionShape["claimActivation"] = (input) => + mutex.withPermits(1)( + Effect.gen(function* () { + const durable = yield* drain.status; + if (durable.intent?.status === "draining" && durable.intent.requestId === input.requestId) { + const blockers = yield* currentBlockers(); + if (blockers.length > 0) { + return yield* new UpdateDrainError({ + reason: "not_quiescent", + message: `Update drain '${input.requestId}' still has ${blockers.length} execution blocker${blockers.length === 1 ? "" : "s"}.`, + }); + } + } + + return yield* drain.dispatch({ + type: "update-drain.claim", + commandId: CommandId.make(`update-drain:claim:${input.requestId}`), + requestId: input.requestId, + createdAt: DateTime.formatIso(yield* DateTime.now), + }); + }), + ); + + const admit: UpdateDrainAdmissionShape["admit"] = (kind, effect) => + mutex.withPermits(1)( + Effect.gen(function* () { + const durable = yield* drain.status; + if (durable.intent !== null && durable.intent.status !== "cancelled") { + return yield* new UpdateDrainAdmissionError({ + reason: "update_draining", + requestId: durable.intent.requestId, + targetVersion: durable.intent.targetVersion, + message: `Cannot start ${kind.replaceAll("-", " ")} while LastCode is draining for update ${durable.intent.targetVersion}.`, + }); + } + return yield* effect; + }), + ); + + const admitOrElse: UpdateDrainAdmissionShape["admitOrElse"] = (_kind, effect, whenClosed) => + mutex.withPermits(1)( + Effect.gen(function* () { + const durable = yield* drain.status; + return yield* durable.intent !== null && durable.intent.status !== "cancelled" + ? whenClosed + : effect; + }), + ); + + return UpdateDrainAdmission.of({ + dispatch, + claimActivation, + status: mutex.withPermits(1)(statusUnlocked()), + admit, + admitOrElse, + }); +}); + +export const layer = Layer.effect(UpdateDrainAdmission, makeUpdateDrainAdmission()); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 0df33343e3b9..0d1a833920dc 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -58,6 +58,8 @@ import { type TerminalError, type TerminalEvent, type TerminalMetadataStreamEvent, + UpdateDrainAdmissionError, + UpdateDrainError, WS_METHODS, WsRpcGroup, } from "@t3tools/contracts"; @@ -113,7 +115,7 @@ import * as ResourceTelemetry from "./resourceTelemetry/ResourceTelemetry.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; import * as UsageService from "./usage/UsageService.ts"; import * as TraceDiagnostics from "./diagnostics/TraceDiagnostics.ts"; -import * as UpdateDrain from "./updateDrain/UpdateDrain.ts"; +import * as UpdateDrainAdmission from "./updateDrain/UpdateDrainAdmission.ts"; import * as PullRequestService from "./pullRequest/PullRequestService.ts"; import * as SourceControlDiscovery from "./sourceControl/SourceControlDiscovery.ts"; import * as SourceControlRepositoryService from "./sourceControl/SourceControlRepositoryService.ts"; @@ -131,6 +133,8 @@ import * as SessionStore from "./auth/SessionStore.ts"; import { failEnvironmentAuthInvalid, failEnvironmentInternal } from "./auth/http.ts"; import * as RelayClient from "@t3tools/shared/relayClient"; const isOrchestrationDispatchCommandError = Schema.is(OrchestrationDispatchCommandError); +const isUpdateDrainAdmissionError = Schema.is(UpdateDrainAdmissionError); +const isUpdateDrainError = Schema.is(UpdateDrainError); const nowIso = Effect.map(DateTime.now, DateTime.formatIso); const EDITOR_DISCOVERY_TIMEOUT = Duration.seconds(5); @@ -489,7 +493,7 @@ const makeWsRpcLayer = ( const processResourceMonitor = yield* ProcessResourceMonitor.ProcessResourceMonitor; const resourceTelemetry = yield* ResourceTelemetry.ResourceTelemetry; const usage = yield* UsageService.UsageService; - const updateDrain = yield* UpdateDrain.UpdateDrain; + const updateDrainAdmission = yield* UpdateDrainAdmission.UpdateDrainAdmission; const relayClient = yield* RelayClient.RelayClient; const authorizationError = (requiredScope: AuthEnvironmentScope) => new EnvironmentAuthorizationError({ @@ -1063,7 +1067,10 @@ const makeWsRpcLayer = ( const dispatchNormalizedCommand = ( normalizedCommand: OrchestrationCommand, - ): Effect.Effect<{ readonly sequence: number }, OrchestrationDispatchCommandError> => { + ): Effect.Effect< + { readonly sequence: number }, + OrchestrationDispatchCommandError | UpdateDrainAdmissionError | UpdateDrainError + > => { const dispatchEffect = normalizedCommand.type === "thread.turn.start" && normalizedCommand.bootstrap ? dispatchBootstrapTurnStart(normalizedCommand) @@ -1073,11 +1080,18 @@ const makeWsRpcLayer = ( ), ); + const admittedEffect = + normalizedCommand.type === "thread.turn.start" + ? updateDrainAdmission.admit("thread-turn", dispatchEffect) + : dispatchEffect; + return startup - .enqueueCommand(dispatchEffect) + .enqueueCommand(admittedEffect) .pipe( Effect.mapError((cause) => - toDispatchCommandError(cause, "Failed to dispatch orchestration command"), + isUpdateDrainAdmissionError(cause) || isUpdateDrainError(cause) + ? cause + : toDispatchCommandError(cause, "Failed to dispatch orchestration command"), ), ); }; @@ -1217,7 +1231,9 @@ const makeWsRpcLayer = ( return result; }).pipe( Effect.mapError((cause) => - isOrchestrationDispatchCommandError(cause) + isOrchestrationDispatchCommandError(cause) || + isUpdateDrainAdmissionError(cause) || + isUpdateDrainError(cause) ? cause : new OrchestrationDispatchCommandError({ message: "Failed to dispatch orchestration command", @@ -1705,7 +1721,7 @@ const makeWsRpcLayer = ( observeRpcEffect( WS_METHODS.serverStartUpdateDrain, Effect.flatMap(nowIso, (createdAt) => - updateDrain.dispatch({ + updateDrainAdmission.dispatch({ type: "update-drain.start", ...input, createdAt, @@ -1717,16 +1733,29 @@ const makeWsRpcLayer = ( observeRpcEffect( WS_METHODS.serverCancelUpdateDrain, Effect.flatMap(nowIso, (createdAt) => - updateDrain.dispatch({ + updateDrainAdmission.dispatch({ type: "update-drain.cancel", ...input, createdAt, }), + ).pipe( + Effect.tap(() => + Option.match(actionResume, { + onNone: () => Effect.void, + onSome: (service) => service.retryPendingFollowUps, + }), + ), ), { "rpc.aggregate": "update-drain" }, ), + [WS_METHODS.serverClaimUpdateActivation]: (input) => + observeRpcEffect( + WS_METHODS.serverClaimUpdateActivation, + updateDrainAdmission.claimActivation(input), + { "rpc.aggregate": "update-drain" }, + ), [WS_METHODS.serverGetUpdateDrainStatus]: (_input) => - observeRpcEffect(WS_METHODS.serverGetUpdateDrainStatus, updateDrain.status, { + observeRpcEffect(WS_METHODS.serverGetUpdateDrainStatus, updateDrainAdmission.status, { "rpc.aggregate": "update-drain", }), [WS_METHODS.cloudGetRelayClientStatus]: (_input) => @@ -2107,9 +2136,12 @@ const makeWsRpcLayer = ( [WS_METHODS.gitPreparePullRequestThread]: (input) => observeRpcEffect( WS_METHODS.gitPreparePullRequestThread, - gitWorkflow - .preparePullRequestThread(input) - .pipe(Effect.tap(() => refreshGitStatus(input.cwd))), + updateDrainAdmission.admit( + "setup-script", + gitWorkflow + .preparePullRequestThread(input) + .pipe(Effect.tap(() => refreshGitStatus(input.cwd))), + ), { "rpc.aggregate": "git" }, ), [WS_METHODS.vcsListRefs]: (input) => @@ -2159,24 +2191,37 @@ const makeWsRpcLayer = ( { "rpc.aggregate": "review" }, ), [WS_METHODS.terminalOpen]: (input) => - observeRpcEffect(WS_METHODS.terminalOpen, terminalManager.open(input), { - "rpc.aggregate": "terminal", - }), + observeRpcEffect( + WS_METHODS.terminalOpen, + updateDrainAdmission.admit("terminal-open", terminalManager.open(input)), + { "rpc.aggregate": "terminal" }, + ), [WS_METHODS.terminalAttach]: (input) => observeRpcStream( WS_METHODS.terminalAttach, - Stream.callback((queue) => - Effect.acquireRelease( - terminalManager.attachStream(input, (event) => Queue.offer(queue, event)), - (unsubscribe) => Effect.sync(unsubscribe), - ), - ), + Stream.callback< + TerminalAttachStreamEvent, + TerminalError | UpdateDrainAdmissionError | UpdateDrainError + >((queue) => { + const attach = (startIfNeeded: boolean) => + Effect.acquireRelease( + terminalManager.attachStream( + input, + (event) => Queue.offer(queue, event), + startIfNeeded, + ), + (unsubscribe) => Effect.sync(unsubscribe), + ); + return updateDrainAdmission.admitOrElse("terminal-open", attach(true), attach(false)); + }), { "rpc.aggregate": "terminal" }, ), [WS_METHODS.terminalWrite]: (input) => - observeRpcEffect(WS_METHODS.terminalWrite, terminalManager.write(input), { - "rpc.aggregate": "terminal", - }), + observeRpcEffect( + WS_METHODS.terminalWrite, + updateDrainAdmission.admit("terminal-write", terminalManager.write(input)), + { "rpc.aggregate": "terminal" }, + ), [WS_METHODS.terminalResize]: (input) => observeRpcEffect(WS_METHODS.terminalResize, terminalManager.resize(input), { "rpc.aggregate": "terminal", @@ -2186,9 +2231,11 @@ const makeWsRpcLayer = ( "rpc.aggregate": "terminal", }), [WS_METHODS.terminalRestart]: (input) => - observeRpcEffect(WS_METHODS.terminalRestart, terminalManager.restart(input), { - "rpc.aggregate": "terminal", - }), + observeRpcEffect( + WS_METHODS.terminalRestart, + updateDrainAdmission.admit("terminal-restart", terminalManager.restart(input)), + { "rpc.aggregate": "terminal" }, + ), [WS_METHODS.terminalClose]: (input) => observeRpcEffect(WS_METHODS.terminalClose, terminalManager.close(input), { "rpc.aggregate": "terminal", @@ -2196,16 +2243,19 @@ const makeWsRpcLayer = ( [WS_METHODS.actionResumeResume]: (input) => observeRpcEffect( WS_METHODS.actionResumeResume, - Option.match(actionResume, { - onNone: () => - Effect.fail( - new ActionResumeError({ - reason: "internal_error", - message: "Action resume is unavailable in this server runtime.", - }), - ), - onSome: (service) => service.resumeInterrupted(input.threadId), - }), + updateDrainAdmission.admit( + "action-resume", + Option.match(actionResume, { + onNone: () => + Effect.fail( + new ActionResumeError({ + reason: "internal_error", + message: "Action resume is unavailable in this server runtime.", + }), + ), + onSome: (service) => service.resumeInterrupted(input.threadId), + }), + ), { "rpc.aggregate": "action-resume" }, ), [WS_METHODS.actionResumeDiscard]: (input) => diff --git a/docs/lastcode/release.md b/docs/lastcode/release.md index fb2435845178..93af7201f07f 100644 --- a/docs/lastcode/release.md +++ b/docs/lastcode/release.md @@ -211,3 +211,17 @@ dedicated worktree, and retains the validated DMG for the managed local installer. It does not stage the updater ZIP through Electron/Squirrel. Checkpoint scheduling remains independent: the daemon never builds merely because it found a new nightly or revision. + +## Remote update drain admission + +An active remote update drain closes only entry points that can create new +execution: turn starts (including provider bootstrap or resume), terminal +creation and restart, terminal writes, and interrupted Action resume. Existing +read-only terminal attachment, terminal close, turn interruption, approvals, +and user-input responses remain available so current work can settle. + +Drain status reports only current execution blockers: starting or running +thread work, background agent work, and starting terminals or terminals with a +running subprocess. When that list is empty, the activation claim is committed +under the same server-lifetime admission lock. The claim survives a server +restart and keeps admission closed for the future activation helper. diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index 3b6bfd5e5316..459edc7ee55b 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -2,11 +2,13 @@ import * as Schema from "effect/Schema"; import * as Rpc from "effect/unstable/rpc/Rpc"; import * as RpcGroup from "effect/unstable/rpc/RpcGroup"; import { + UpdateDrainAdmissionError, UpdateDrainCancelInput, + UpdateDrainClaimInput, UpdateDrainCommandReceipt, UpdateDrainError, UpdateDrainStartInput, - UpdateDrainState, + UpdateDrainStatus, } from "./updateDrain.ts"; import { ExternalLauncherError, LaunchEditorInput } from "./editor.ts"; @@ -288,6 +290,7 @@ export const WS_METHODS = { serverGetUsageSummary: "server.getUsageSummary", serverStartUpdateDrain: "server.startUpdateDrain", serverCancelUpdateDrain: "server.cancelUpdateDrain", + serverClaimUpdateActivation: "server.claimUpdateActivation", serverGetUpdateDrainStatus: "server.getUpdateDrainStatus", // Cloud environment methods @@ -498,9 +501,15 @@ export const WsServerCancelUpdateDrainRpc = Rpc.make(WS_METHODS.serverCancelUpda error: Schema.Union([UpdateDrainError, EnvironmentAuthorizationError]), }); +export const WsServerClaimUpdateActivationRpc = Rpc.make(WS_METHODS.serverClaimUpdateActivation, { + payload: UpdateDrainClaimInput, + success: UpdateDrainCommandReceipt, + error: Schema.Union([UpdateDrainError, EnvironmentAuthorizationError]), +}); + export const WsServerGetUpdateDrainStatusRpc = Rpc.make(WS_METHODS.serverGetUpdateDrainStatus, { payload: Schema.Struct({}), - success: UpdateDrainState, + success: UpdateDrainStatus, error: Schema.Union([UpdateDrainError, EnvironmentAuthorizationError]), }); @@ -734,7 +743,12 @@ export const WsGitResolvePullRequestRpc = Rpc.make(WS_METHODS.gitResolvePullRequ export const WsGitPreparePullRequestThreadRpc = Rpc.make(WS_METHODS.gitPreparePullRequestThread, { payload: GitPreparePullRequestThreadInput, success: GitPreparePullRequestThreadResult, - error: Schema.Union([GitManagerServiceError, EnvironmentAuthorizationError]), + error: Schema.Union([ + GitManagerServiceError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), }); export const WsVcsListRefsRpc = Rpc.make(WS_METHODS.vcsListRefs, { @@ -791,19 +805,34 @@ export const WsReviewGetDiffFileContentsRpc = Rpc.make(WS_METHODS.reviewGetDiffF export const WsTerminalOpenRpc = Rpc.make(WS_METHODS.terminalOpen, { payload: TerminalOpenInput, success: TerminalSessionSnapshot, - error: Schema.Union([TerminalError, EnvironmentAuthorizationError]), + error: Schema.Union([ + TerminalError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), }); export const WsTerminalAttachRpc = Rpc.make(WS_METHODS.terminalAttach, { payload: TerminalAttachInput, success: TerminalAttachStreamEvent, - error: Schema.Union([TerminalError, EnvironmentAuthorizationError]), + error: Schema.Union([ + TerminalError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), stream: true, }); export const WsTerminalWriteRpc = Rpc.make(WS_METHODS.terminalWrite, { payload: TerminalWriteInput, - error: Schema.Union([TerminalError, EnvironmentAuthorizationError]), + error: Schema.Union([ + TerminalError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), }); export const WsTerminalResizeRpc = Rpc.make(WS_METHODS.terminalResize, { @@ -819,7 +848,12 @@ export const WsTerminalClearRpc = Rpc.make(WS_METHODS.terminalClear, { export const WsTerminalRestartRpc = Rpc.make(WS_METHODS.terminalRestart, { payload: TerminalRestartInput, success: TerminalSessionSnapshot, - error: Schema.Union([TerminalError, EnvironmentAuthorizationError]), + error: Schema.Union([ + TerminalError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), }); export const WsTerminalCloseRpc = Rpc.make(WS_METHODS.terminalClose, { @@ -829,7 +863,12 @@ export const WsTerminalCloseRpc = Rpc.make(WS_METHODS.terminalClose, { export const WsActionResumeResumeRpc = Rpc.make(WS_METHODS.actionResumeResume, { payload: Schema.Struct({ threadId: ThreadId }), - error: Schema.Union([ActionResumeError, EnvironmentAuthorizationError]), + error: Schema.Union([ + ActionResumeError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), }); export const WsActionResumeDiscardRpc = Rpc.make(WS_METHODS.actionResumeDiscard, { @@ -917,7 +956,12 @@ export const WsOrchestrationDispatchCommandRpc = Rpc.make( { payload: ClientOrchestrationCommand, success: OrchestrationRpcSchemas.dispatchCommand.output, - error: Schema.Union([OrchestrationDispatchCommandError, EnvironmentAuthorizationError]), + error: Schema.Union([ + OrchestrationDispatchCommandError, + UpdateDrainAdmissionError, + UpdateDrainError, + EnvironmentAuthorizationError, + ]), }, ); @@ -1050,6 +1094,7 @@ export const WsRpcGroup = RpcGroup.make( WsServerGetBackgroundPolicyRpc, WsServerStartUpdateDrainRpc, WsServerCancelUpdateDrainRpc, + WsServerClaimUpdateActivationRpc, WsServerGetUpdateDrainStatusRpc, WsCloudGetRelayClientStatusRpc, WsCloudInstallRelayClientRpc, diff --git a/packages/contracts/src/updateDrain.test.ts b/packages/contracts/src/updateDrain.test.ts index 60ce7354e0ac..42a1dda99c08 100644 --- a/packages/contracts/src/updateDrain.test.ts +++ b/packages/contracts/src/updateDrain.test.ts @@ -1,14 +1,20 @@ import { describe, expect, it } from "vite-plus/test"; import * as Schema from "effect/Schema"; -import { UpdateDrainCommand, UpdateDrainCommandReceipt, UpdateDrainState } from "./updateDrain.ts"; +import { + UpdateDrainCommand, + UpdateDrainCommandReceipt, + UpdateDrainState, + UpdateDrainStatus, +} from "./updateDrain.ts"; const decodeCommand = Schema.decodeUnknownSync(UpdateDrainCommand); const decodeState = Schema.decodeUnknownSync(UpdateDrainState); const decodeReceipt = Schema.decodeUnknownSync(UpdateDrainCommandReceipt); +const decodeStatus = Schema.decodeUnknownSync(UpdateDrainStatus); describe("update drain contracts", () => { - it("decodes start and cancel commands", () => { + it("decodes lifecycle commands", () => { expect( decodeCommand({ type: "update-drain.start", @@ -26,6 +32,14 @@ describe("update drain contracts", () => { createdAt: "2026-08-21T00:01:00.000Z", }), ).toMatchObject({ type: "update-drain.cancel" }); + expect( + decodeCommand({ + type: "update-drain.claim", + commandId: "update-drain:claim:request-1", + requestId: "request-1", + createdAt: "2026-08-21T00:02:00.000Z", + }), + ).toMatchObject({ type: "update-drain.claim" }); }); it("keeps blocker and grace data outside the durable status contract", () => { @@ -49,4 +63,23 @@ describe("update drain contracts", () => { expect(state).not.toHaveProperty("quietSince"); expect(receipt.status).toBe("accepted"); }); + + it("decodes the minimal live blocker status separately from durable state", () => { + const status = decodeStatus({ + sequence: 2, + intent: { requestId: "request-1", targetVersion: "1.2.3", status: "claimed" }, + admission: "closed", + blockers: [ + { + type: "terminal-process", + threadId: "thread-1", + terminalId: "terminal-1", + label: "tests", + status: "running", + }, + ], + }); + expect(status.blockers).toHaveLength(1); + expect(status).not.toHaveProperty("quietGrace"); + }); }); diff --git a/packages/contracts/src/updateDrain.ts b/packages/contracts/src/updateDrain.ts index df85d8392ca2..30a6903509df 100644 --- a/packages/contracts/src/updateDrain.ts +++ b/packages/contracts/src/updateDrain.ts @@ -5,6 +5,8 @@ import { EventId, IsoDateTime, NonNegativeInt, + ThreadId, + TurnId, TrimmedNonEmptyString, } from "./baseSchemas.ts"; @@ -18,7 +20,7 @@ export const UpdateDrainTargetVersion = TrimmedNonEmptyString.pipe( ); export type UpdateDrainTargetVersion = typeof UpdateDrainTargetVersion.Type; -export const UpdateDrainIntentStatus = Schema.Literals(["draining", "cancelled"]); +export const UpdateDrainIntentStatus = Schema.Literals(["draining", "cancelled", "claimed"]); export type UpdateDrainIntentStatus = typeof UpdateDrainIntentStatus.Type; export const UpdateDrainIntent = Schema.Struct({ @@ -34,6 +36,46 @@ export const UpdateDrainState = Schema.Struct({ }); export type UpdateDrainState = typeof UpdateDrainState.Type; +export const UpdateDrainBlocker = Schema.Union([ + Schema.Struct({ + type: Schema.Literal("thread-turn"), + threadId: ThreadId, + turnId: Schema.NullOr(TurnId), + status: Schema.Literals(["starting", "running"]), + }), + Schema.Struct({ + type: Schema.Literal("thread-background"), + threadId: ThreadId, + status: Schema.Literals(["working", "monitoring"]), + }), + Schema.Struct({ + type: Schema.Literal("terminal-process"), + threadId: ThreadId, + terminalId: TrimmedNonEmptyString, + label: Schema.String, + status: Schema.Literals(["starting", "running"]), + }), +]); +export type UpdateDrainBlocker = typeof UpdateDrainBlocker.Type; + +export const UpdateDrainStatus = Schema.Struct({ + sequence: NonNegativeInt, + intent: Schema.NullOr(UpdateDrainIntent), + admission: Schema.Literals(["open", "closed"]), + blockers: Schema.Array(UpdateDrainBlocker), +}); +export type UpdateDrainStatus = typeof UpdateDrainStatus.Type; + +export class UpdateDrainAdmissionError extends Schema.TaggedErrorClass()( + "UpdateDrainAdmissionError", + { + reason: Schema.Literal("update_draining"), + requestId: UpdateDrainRequestId, + targetVersion: UpdateDrainTargetVersion, + message: Schema.String, + }, +) {} + const UpdateDrainCommandBase = { commandId: CommandId, requestId: UpdateDrainRequestId, @@ -53,7 +95,17 @@ export const UpdateDrainCancelCommand = Schema.Struct({ }); export type UpdateDrainCancelCommand = typeof UpdateDrainCancelCommand.Type; -export const UpdateDrainCommand = Schema.Union([UpdateDrainStartCommand, UpdateDrainCancelCommand]); +export const UpdateDrainClaimCommand = Schema.Struct({ + ...UpdateDrainCommandBase, + type: Schema.Literal("update-drain.claim"), +}); +export type UpdateDrainClaimCommand = typeof UpdateDrainClaimCommand.Type; + +export const UpdateDrainCommand = Schema.Union([ + UpdateDrainStartCommand, + UpdateDrainCancelCommand, + UpdateDrainClaimCommand, +]); export type UpdateDrainCommand = typeof UpdateDrainCommand.Type; const UpdateDrainEventBase = { @@ -79,13 +131,26 @@ export const UpdateDrainCancelledEvent = Schema.Struct({ }); export type UpdateDrainCancelledEvent = typeof UpdateDrainCancelledEvent.Type; -export const UpdateDrainEvent = Schema.Union([UpdateDrainStartedEvent, UpdateDrainCancelledEvent]); +export const UpdateDrainClaimedEvent = Schema.Struct({ + ...UpdateDrainEventBase, + type: Schema.Literal("update-drain.claimed"), + status: Schema.Literal("claimed"), +}); +export type UpdateDrainClaimedEvent = typeof UpdateDrainClaimedEvent.Type; + +export const UpdateDrainEvent = Schema.Union([ + UpdateDrainStartedEvent, + UpdateDrainCancelledEvent, + UpdateDrainClaimedEvent, +]); export type UpdateDrainEvent = typeof UpdateDrainEvent.Type; export const UpdateDrainFailureReason = Schema.Literals([ "already_draining", + "activation_claimed", "command_id_conflict", "internal_error", + "not_quiescent", "no_active_drain", "request_already_cancelled", "request_mismatch", @@ -95,7 +160,7 @@ export type UpdateDrainFailureReason = typeof UpdateDrainFailureReason.Type; export const UpdateDrainCommandReceipt = Schema.Struct({ commandId: CommandId, requestId: UpdateDrainRequestId, - commandType: Schema.Literals(["update-drain.start", "update-drain.cancel"]), + commandType: Schema.Literals(["update-drain.start", "update-drain.cancel", "update-drain.claim"]), targetVersion: Schema.NullOr(UpdateDrainTargetVersion), acceptedAt: IsoDateTime, resultSequence: NonNegativeInt, @@ -118,6 +183,11 @@ export const UpdateDrainCancelInput = Schema.Struct({ }); export type UpdateDrainCancelInput = typeof UpdateDrainCancelInput.Type; +export const UpdateDrainClaimInput = Schema.Struct({ + requestId: UpdateDrainRequestId, +}); +export type UpdateDrainClaimInput = typeof UpdateDrainClaimInput.Type; + export class UpdateDrainError extends Schema.TaggedErrorClass()( "UpdateDrainError", {