diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index 3923812218..9376b74989 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -67,6 +67,7 @@ import { import { ThreadDeletionReactor } from "../src/orchestration/Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../src/orchestration/ThreadSettlementReactor.ts"; import * as DelegationFollowThroughReactor from "../src/orchestration/DelegationFollowThroughReactor.ts"; +import * as PairLifecycleReactor from "../src/orchestration/PairLifecycleReactor.ts"; import * as PullRequestSyncReactor from "../src/orchestration/PullRequestSyncReactor.ts"; import * as ThreadPullRequestReactor from "../src/orchestration/ThreadPullRequestReactor.ts"; import * as ProjectSettingsReactor from "../src/orchestration/ProjectSettingsReactor.ts"; @@ -401,6 +402,12 @@ export const makeOrchestrationIntegrationHarness = ( drain: Effect.void, }), ), + Layer.provide( + Layer.succeed(PairLifecycleReactor.PairLifecycleReactor, { + start: () => Effect.void, + drain: Effect.void, + }), + ), Layer.provideMerge( Layer.succeed(ThreadSettlementReactor.ThreadSettlementReactor, { start: () => Effect.void, diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts index 037abf85a7..69cf65f858 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.test.ts @@ -11,6 +11,7 @@ import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeInge import { ThreadDeletionReactor } from "../Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../ThreadSettlementReactor.ts"; import * as DelegationFollowThroughReactor from "../DelegationFollowThroughReactor.ts"; +import * as PairLifecycleReactor from "../PairLifecycleReactor.ts"; import * as PullRequestSyncReactor from "../PullRequestSyncReactor.ts"; import * as ProjectSettingsReactor from "../ProjectSettingsReactor.ts"; import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts"; @@ -96,6 +97,15 @@ describe("OrchestrationReactor", () => { drain: Effect.void, }), ), + Layer.provideMerge( + Layer.succeed(PairLifecycleReactor.PairLifecycleReactor, { + start: () => + Effect.sync(() => { + started.push("pair-lifecycle"); + }), + drain: Effect.void, + }), + ), Layer.provideMerge( Layer.succeed(ThreadSettlementReactor.ThreadSettlementReactor, { start: () => { @@ -142,6 +152,7 @@ describe("OrchestrationReactor", () => { "pull-request-sync-reactor", "agent-awareness-relay", "delegation-follow-through", + "pair-lifecycle", ]); await Effect.runPromise(Scope.close(scope, Exit.void)); diff --git a/apps/server/src/orchestration/Layers/OrchestrationReactor.ts b/apps/server/src/orchestration/Layers/OrchestrationReactor.ts index 7706d27d50..4ab186fa35 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationReactor.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationReactor.ts @@ -11,6 +11,7 @@ import { ProviderRuntimeIngestionService } from "../Services/ProviderRuntimeInge import { ThreadDeletionReactor } from "../Services/ThreadDeletionReactor.ts"; import * as ThreadSettlementReactor from "../ThreadSettlementReactor.ts"; import * as DelegationFollowThroughReactor from "../DelegationFollowThroughReactor.ts"; +import * as PairLifecycleReactor from "../PairLifecycleReactor.ts"; import * as PullRequestSyncReactor from "../PullRequestSyncReactor.ts"; import * as ProjectSettingsReactor from "../ProjectSettingsReactor.ts"; import * as ThreadPullRequestReactor from "../ThreadPullRequestReactor.ts"; @@ -23,6 +24,7 @@ export const makeOrchestrationReactor = Effect.gen(function* () { const threadDeletionReactor = yield* ThreadDeletionReactor; const delegationFollowThrough = yield* DelegationFollowThroughReactor.DelegationFollowThroughReactor; + const pairLifecycle = yield* PairLifecycleReactor.PairLifecycleReactor; const threadSettlementReactor = yield* ThreadSettlementReactor.ThreadSettlementReactor; const pullRequestSyncReactor = yield* PullRequestSyncReactor.PullRequestSyncReactor; const threadPullRequestReactor = yield* ThreadPullRequestReactor.ThreadPullRequestReactor; @@ -40,6 +42,7 @@ export const makeOrchestrationReactor = Effect.gen(function* () { yield* pullRequestSyncReactor.start(); yield* agentAwarenessRelay.start(); yield* delegationFollowThrough.start(); + yield* pairLifecycle.start(); }); return { diff --git a/apps/server/src/orchestration/PairLifecycleReactor.test.ts b/apps/server/src/orchestration/PairLifecycleReactor.test.ts new file mode 100644 index 0000000000..3ce9e45059 --- /dev/null +++ b/apps/server/src/orchestration/PairLifecycleReactor.test.ts @@ -0,0 +1,331 @@ +import * as NodeCrypto from "node:crypto"; + +import { + CommandId, + EventId, + ProjectId, + ProviderInstanceId, + ThreadId, + TurnId, + type OrchestrationCommand, + type OrchestrationEvent, + type OrchestrationThreadShell, +} from "@t3tools/contracts"; +import { assert, describe, it } from "@effect/vitest"; +import * as Crypto from "effect/Crypto"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as PubSub from "effect/PubSub"; +import * as Stream from "effect/Stream"; + +import { ServerActivation } from "../serverActivation.ts"; +import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; +import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; +import * as Reactor from "./PairLifecycleReactor.ts"; + +const NOW = "2026-09-18T00:00:00.000Z"; +const LEAD = ThreadId.make("lead"); +const OTHER = ThreadId.make("other"); +const executorOf = (lead: ThreadId) => + ThreadId.make( + `delegated:${lead}:${NodeCrypto.createHash("sha256").update(`${lead}\npair`).digest("hex").slice(0, 16)}`, + ); +const EXECUTOR = executorOf(LEAD); +const FAN_OUT_CHILD = ThreadId.make(`delegated:${LEAD}:0123456789abcdef`); + +function shell( + id: ThreadId, + overrides: Partial = {}, +): OrchestrationThreadShell { + return { + id, + projectId: ProjectId.make("project"), + title: id, + modelSelection: { instanceId: ProviderInstanceId.make("antigravity"), model: "test" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + pullRequests: [], + latestTurn: null, + createdAt: NOW, + updatedAt: NOW, + archivedAt: null, + settledOverride: null, + settledAt: null, + session: null, + latestUserMessageAt: NOW, + hasPendingApprovals: false, + hasPendingUserInput: false, + hasActionableProposedPlan: false, + ...overrides, + }; +} + +const runningTurnId = TurnId.make("turn-running"); +const running = (id: ThreadId) => + shell(id, { + latestTurn: { + turnId: runningTurnId, + state: "running", + requestedAt: NOW, + startedAt: NOW, + completedAt: null, + assistantMessageId: null, + }, + session: { + threadId: id, + status: "running", + providerName: "antigravity", + runtimeMode: "full-access", + activeTurnId: runningTurnId, + lastError: null, + updatedAt: NOW, + }, + }); + +const targetOf = (command: OrchestrationCommand | undefined) => + command !== undefined && "threadId" in command ? command.threadId : null; + +type LifecycleEventType = + | "thread.archived" + | "thread.deleted" + | "thread.settled" + | "thread.checkpoint-revert-requested" + | "thread.unarchived"; + +let nextEvent = 0; +function lifecycleEvent( + type: LifecycleEventType, + threadId: ThreadId, + historyImport = false, +): OrchestrationEvent { + const base = { + eventId: EventId.make(`event-${++nextEvent}`), + sequence: nextEvent, + aggregateKind: "thread" as const, + aggregateId: threadId, + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: { historyImport }, + }; + switch (type) { + case "thread.archived": + return { ...base, type, payload: { threadId, archivedAt: NOW, updatedAt: NOW } }; + case "thread.deleted": + return { ...base, type, payload: { threadId, deletedAt: NOW } }; + case "thread.settled": + return { ...base, type, payload: { threadId, settledAt: NOW, updatedAt: NOW } }; + case "thread.checkpoint-revert-requested": + return { ...base, type, payload: { threadId, turnCount: 1, createdAt: NOW } }; + case "thread.unarchived": + return { ...base, type, payload: { threadId, updatedAt: NOW } }; + } +} + +const makeHarness = Effect.fn(function* (input: { + readonly active?: readonly OrchestrationThreadShell[]; + readonly archived?: readonly OrchestrationThreadShell[]; +}) { + const active = [...(input.active ?? [])]; + const archived = [...(input.archived ?? [])]; + const commands: OrchestrationCommand[] = []; + const events = yield* PubSub.unbounded(); + const dependencies = Layer.mergeAll( + Layer.mock(ProjectionSnapshotQuery)({ + getThreadShellById: (id) => + Effect.succeed(Option.fromNullishOr(active.find((thread) => thread.id === id))), + getArchivedShellSnapshot: () => + Effect.succeed({ snapshotSequence: 1, projects: [], threads: archived, updatedAt: NOW }), + }), + Layer.mock(OrchestrationEngineService)({ + subscribeDomainEvents: PubSub.subscribe(events).pipe(Effect.map(Stream.fromSubscription)), + dispatch: (command) => + Effect.sync(() => { + commands.push(command); + return { sequence: commands.length }; + }), + }), + Layer.succeed(ServerActivation, Effect.void), + Layer.succeed( + Crypto.Crypto, + Crypto.make({ + randomBytes: (size) => new Uint8Array(size), + digest: (_algorithm, data) => + Effect.succeed(new Uint8Array(NodeCrypto.createHash("sha256").update(data).digest())), + }), + ), + ); + const service = yield* Reactor.make.pipe(Effect.provide(dependencies)); + yield* service.start(); + const emit = (event: OrchestrationEvent) => + Effect.gen(function* () { + yield* PubSub.publish(events, event); + // Let the subscriber fiber take the event, then wait for the worker. + yield* Effect.yieldNow; + yield* Effect.yieldNow; + yield* service.drain; + }); + return { commands, emit }; +}); + +describe("PairLifecycleReactor", () => { + it.effect("archives, settles, and deletes the executor with its lead", () => + Effect.scoped( + Effect.gen(function* () { + const h = yield* makeHarness({ active: [shell(LEAD), shell(EXECUTOR)] }); + const archive = lifecycleEvent("thread.archived", LEAD); + const settle = lifecycleEvent("thread.settled", LEAD); + const remove = lifecycleEvent("thread.deleted", LEAD); + yield* h.emit(archive); + yield* h.emit(settle); + yield* h.emit(remove); + assert.deepStrictEqual(h.commands, [ + { + type: "thread.archive", + commandId: CommandId.make( + `server:pair-lifecycle:archive:${EXECUTOR}:${archive.eventId}`, + ), + threadId: EXECUTOR, + }, + { + type: "thread.settle", + commandId: CommandId.make(`server:pair-lifecycle:settle:${EXECUTOR}:${settle.eventId}`), + threadId: EXECUTOR, + }, + { + type: "thread.delete", + commandId: CommandId.make(`server:pair-lifecycle:delete:${EXECUTOR}:${remove.eventId}`), + threadId: EXECUTOR, + }, + ]); + }), + ), + ); + + it.effect("stops a running executor when its lead is rewound, and only then", () => + Effect.scoped( + Effect.gen(function* () { + const busy = yield* makeHarness({ active: [shell(LEAD), running(EXECUTOR)] }); + const rewind = lifecycleEvent("thread.checkpoint-revert-requested", LEAD); + yield* busy.emit(rewind); + assert.strictEqual(busy.commands.length, 1); + const [interrupt] = busy.commands; + assert.strictEqual(interrupt?.type, "thread.turn.interrupt"); + assert.strictEqual(targetOf(interrupt), EXECUTOR); + assert.strictEqual( + interrupt?.commandId, + `server:pair-lifecycle:interrupt:${EXECUTOR}:${rewind.eventId}`, + ); + + const idle = yield* makeHarness({ active: [shell(LEAD), shell(EXECUTOR)] }); + yield* idle.emit(lifecycleEvent("thread.checkpoint-revert-requested", LEAD)); + assert.deepStrictEqual(idle.commands, []); + }), + ), + ); + + it.effect("deletes an executor that was archived when the pair was turned off", () => + Effect.scoped( + Effect.gen(function* () { + const h = yield* makeHarness({ + active: [shell(LEAD)], + archived: [shell(EXECUTOR, { archivedAt: NOW })], + }); + // Archiving and settling have nothing left to do for an archived executor. + yield* h.emit(lifecycleEvent("thread.archived", LEAD)); + yield* h.emit(lifecycleEvent("thread.settled", LEAD)); + assert.deepStrictEqual(h.commands, []); + yield* h.emit(lifecycleEvent("thread.deleted", LEAD)); + assert.deepStrictEqual( + h.commands.map((command) => ({ type: command.type, threadId: targetOf(command) })), + [{ type: "thread.delete", threadId: EXECUTOR }], + ); + }), + ), + ); + + it.effect("does nothing for threads without an executor or for other threads' events", () => + Effect.scoped( + Effect.gen(function* () { + const h = yield* makeHarness({ + active: [shell(LEAD), shell(OTHER), shell(EXECUTOR), shell(FAN_OUT_CHILD)], + }); + // OTHER has no executor. A fan-out child of LEAD is not the pair. + yield* h.emit(lifecycleEvent("thread.archived", OTHER)); + yield* h.emit(lifecycleEvent("thread.deleted", OTHER)); + // The executor's own lifecycle is how a user turns the pair off; it must not echo. + yield* h.emit(lifecycleEvent("thread.archived", EXECUTOR)); + yield* h.emit(lifecycleEvent("thread.archived", FAN_OUT_CHILD)); + // Unarchiving a lead does not turn the pair back on, and imported history is inert. + yield* h.emit(lifecycleEvent("thread.unarchived", LEAD)); + yield* h.emit(lifecycleEvent("thread.archived", LEAD, true)); + assert.deepStrictEqual(h.commands, []); + }), + ), + ); + + it.effect("keeps working after a dispatch is rejected", () => + Effect.scoped( + Effect.gen(function* () { + const commands: OrchestrationCommand[] = []; + const events = yield* PubSub.unbounded(); + let first = true; + const dependencies = Layer.mergeAll( + Layer.mock(ProjectionSnapshotQuery)({ + getThreadShellById: (id) => + Effect.succeed( + Option.fromNullishOr( + [shell(LEAD), shell(EXECUTOR)].find((thread) => thread.id === id), + ), + ), + getArchivedShellSnapshot: () => + Effect.succeed({ snapshotSequence: 1, projects: [], threads: [], updatedAt: NOW }), + }), + Layer.mock(OrchestrationEngineService)({ + subscribeDomainEvents: PubSub.subscribe(events).pipe( + Effect.map(Stream.fromSubscription), + ), + dispatch: (command) => { + if (first) { + first = false; + return Effect.die(new Error("rejected")); + } + return Effect.sync(() => { + commands.push(command); + return { sequence: commands.length }; + }); + }, + }), + Layer.succeed(ServerActivation, Effect.void), + Layer.succeed( + Crypto.Crypto, + Crypto.make({ + randomBytes: (size) => new Uint8Array(size), + digest: (_algorithm, data) => + Effect.succeed( + new Uint8Array(NodeCrypto.createHash("sha256").update(data).digest()), + ), + }), + ), + ); + const service = yield* Reactor.make.pipe(Effect.provide(dependencies)); + yield* service.start(); + for (const type of ["thread.archived", "thread.settled"] as const) { + yield* PubSub.publish(events, lifecycleEvent(type, LEAD)); + yield* Effect.yieldNow; + yield* Effect.yieldNow; + yield* service.drain; + } + // The first dispatch failed and was logged; the reactor did not die. + assert.deepStrictEqual( + commands.map(({ type }) => type), + ["thread.settle"], + ); + }), + ), + ); +}); diff --git a/apps/server/src/orchestration/PairLifecycleReactor.ts b/apps/server/src/orchestration/PairLifecycleReactor.ts new file mode 100644 index 0000000000..19ad93e8f0 --- /dev/null +++ b/apps/server/src/orchestration/PairLifecycleReactor.ts @@ -0,0 +1,144 @@ +/** + * Keeps a pair executor in step with its lead: archived, deleted, and settled + * together, and stopped when the lead is rewound. A sidecar over existing + * commands and domain events; it adds no command, event, or decider branch. + * + * @module orchestration/PairLifecycleReactor + */ +import { + CommandId, + type OrchestrationEvent, + type OrchestrationThreadShell, +} from "@t3tools/contracts"; +import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; +import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import type * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; +import { forkParked } from "../serverActivation.ts"; +import { pairExecutorThreadId } from "../mcp/toolkits/pair/logic.ts"; +import * as OrchestrationEngine from "./Services/OrchestrationEngine.ts"; +import * as ProjectionSnapshotQuery from "./Services/ProjectionSnapshotQuery.ts"; +import { + type PairLifecycleIntent, + pairLifecycleApplies, + pairLifecycleIntent, +} from "./pairLifecycle.logic.ts"; + +export class PairLifecycleReactor extends Context.Service< + PairLifecycleReactor, + { + readonly start: () => Effect.Effect; + readonly drain: Effect.Effect; + } +>()("t3/orchestration/PairLifecycleReactor") {} + +export const make = Effect.gen(function* () { + const engine = yield* OrchestrationEngine.OrchestrationEngineService; + const snapshots = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery; + const crypto = yield* Crypto.Crypto; + + type Work = { + readonly event: OrchestrationEvent; + readonly intent: PairLifecycleIntent; + }; + + const process = Effect.fn("PairLifecycleReactor.process")(function* (work: Work) { + const { event, intent } = work; + const digest = yield* crypto.digest( + "SHA-256", + new TextEncoder().encode(`${intent.leadThreadId}\npair`), + ); + const hex = Array.from(digest, (byte) => byte.toString(16).padStart(2, "0")).join(""); + const executorId = pairExecutorThreadId(intent.leadThreadId, () => hex); + + let executor: OrchestrationThreadShell | null = null; + const shellOption = yield* snapshots.getThreadShellById(executorId); + if (Option.isSome(shellOption)) { + executor = shellOption.value; + } else if (intent.action === "delete") { + const archivedSnapshot = yield* snapshots.getArchivedShellSnapshot(); + const found = archivedSnapshot.threads.find((thread) => thread.id === executorId); + if (found !== undefined) { + executor = found; + } + } + if (executor === null) { + return; + } + if (!pairLifecycleApplies(intent.action, executor)) { + return; + } + + const commandId = CommandId.make( + `server:pair-lifecycle:${intent.action}:${executorId}:${event.eventId}`, + ); + + switch (intent.action) { + case "archive": + yield* engine.dispatch({ + type: "thread.archive", + commandId, + threadId: executorId, + }); + break; + case "settle": + yield* engine.dispatch({ + type: "thread.settle", + commandId, + threadId: executorId, + }); + break; + case "delete": + yield* engine.dispatch({ + type: "thread.delete", + commandId, + threadId: executorId, + }); + break; + case "interrupt": { + const createdAt = DateTime.formatIso(yield* DateTime.now); + yield* engine.dispatch({ + type: "thread.turn.interrupt", + commandId, + threadId: executorId, + createdAt, + }); + break; + } + } + }); + + const worker = yield* makeDrainableWorker((work: Work) => + process(work).pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning("Pylon pair lifecycle dispatch failed", { + leadThreadId: work.intent.leadThreadId, + cause: Cause.pretty(cause), + }), + ), + ), + ); + + const processEvent = (event: OrchestrationEvent) => { + const intent = pairLifecycleIntent(event); + if (intent === null) return Effect.void; + return worker.enqueue({ event, intent }); + }; + + const start = Effect.fn("PairLifecycleReactor.start")(function* () { + const events = yield* engine.subscribeDomainEvents; + yield* forkParked(Stream.runForEach(events, processEvent)); + }); + + return { start, drain: worker.drain }; +}); + +export const layer = Layer.effect(PairLifecycleReactor, make); diff --git a/apps/server/src/orchestration/pairLifecycle.logic.test.ts b/apps/server/src/orchestration/pairLifecycle.logic.test.ts new file mode 100644 index 0000000000..4167ecd22b --- /dev/null +++ b/apps/server/src/orchestration/pairLifecycle.logic.test.ts @@ -0,0 +1,156 @@ +import { + EventId, + ProjectId, + ProviderInstanceId, + ThreadId, + TurnId, + type OrchestrationEvent, + type OrchestrationThreadShell, +} from "@t3tools/contracts"; +import { describe, expect, it } from "vite-plus/test"; + +import { pairLifecycleApplies, pairLifecycleIntent } from "./pairLifecycle.logic.ts"; + +const NOW = "2026-09-18T00:00:00.000Z"; +const LEAD = ThreadId.make("lead:with:colons"); +const EXECUTOR = ThreadId.make("delegated:lead:with:colons:0123456789abcdef"); + +const base = (threadId: ThreadId, historyImport = false) => ({ + eventId: EventId.make(`event-${threadId}`), + sequence: 1, + aggregateKind: "thread" as const, + aggregateId: threadId, + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: { historyImport }, +}); + +const archived = (threadId: ThreadId, historyImport = false): OrchestrationEvent => ({ + ...base(threadId, historyImport), + type: "thread.archived", + payload: { threadId, archivedAt: NOW, updatedAt: NOW }, +}); +const deleted = (threadId: ThreadId): OrchestrationEvent => ({ + ...base(threadId), + type: "thread.deleted", + payload: { threadId, deletedAt: NOW }, +}); +const settled = (threadId: ThreadId): OrchestrationEvent => ({ + ...base(threadId), + type: "thread.settled", + payload: { threadId, settledAt: NOW, updatedAt: NOW }, +}); +const rewound = (threadId: ThreadId): OrchestrationEvent => ({ + ...base(threadId), + type: "thread.checkpoint-revert-requested", + payload: { threadId, turnCount: 1, createdAt: NOW }, +}); +const unarchived = (threadId: ThreadId): OrchestrationEvent => ({ + ...base(threadId), + type: "thread.unarchived", + payload: { threadId, updatedAt: NOW }, +}); + +const turn = { + turnId: TurnId.make("turn-1"), + state: "completed" as const, + requestedAt: NOW, + startedAt: NOW, + completedAt: NOW, + assistantMessageId: null, +}; +function executor(overrides: Partial = {}): OrchestrationThreadShell { + return { + id: EXECUTOR, + projectId: ProjectId.make("project"), + title: "Executor", + modelSelection: { instanceId: ProviderInstanceId.make("antigravity"), model: "test" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + pullRequests: [], + latestTurn: turn, + createdAt: NOW, + updatedAt: NOW, + archivedAt: null, + settledOverride: null, + settledAt: null, + session: null, + latestUserMessageAt: NOW, + hasPendingApprovals: false, + hasPendingUserInput: false, + hasActionableProposedPlan: false, + ...overrides, + }; +} + +describe("pairLifecycleIntent", () => { + it("maps a lead's lifecycle to what its executor should do", () => { + expect(pairLifecycleIntent(archived(LEAD))).toEqual({ leadThreadId: LEAD, action: "archive" }); + expect(pairLifecycleIntent(deleted(LEAD))).toEqual({ leadThreadId: LEAD, action: "delete" }); + expect(pairLifecycleIntent(settled(LEAD))).toEqual({ leadThreadId: LEAD, action: "settle" }); + expect(pairLifecycleIntent(rewound(LEAD))).toEqual({ + leadThreadId: LEAD, + action: "interrupt", + }); + }); + + it("never reacts to an executor's or a fan-out child's own lifecycle", () => { + // Archiving the executor is how a user turns the pair off; it must not loop. + for (const event of [ + archived(EXECUTOR), + deleted(EXECUTOR), + settled(EXECUTOR), + rewound(EXECUTOR), + ]) + expect(pairLifecycleIntent(event)).toBeNull(); + }); + + it("ignores imported history and events that ask nothing of the pair", () => { + expect(pairLifecycleIntent(archived(LEAD, true))).toBeNull(); + // Unarchiving a lead does not turn the pair back on: the user may have turned it off. + expect(pairLifecycleIntent(unarchived(LEAD))).toBeNull(); + }); +}); + +describe("pairLifecycleApplies", () => { + it("archives only an executor that is not archived yet", () => { + expect(pairLifecycleApplies("archive", executor())).toBe(true); + expect(pairLifecycleApplies("archive", executor({ archivedAt: NOW }))).toBe(false); + }); + + it("always deletes, even an archived executor", () => { + expect(pairLifecycleApplies("delete", executor())).toBe(true); + expect(pairLifecycleApplies("delete", executor({ archivedAt: NOW }))).toBe(true); + }); + + it("settles only an executor that is active and not settled", () => { + expect(pairLifecycleApplies("settle", executor())).toBe(true); + expect(pairLifecycleApplies("settle", executor({ settledOverride: "settled" }))).toBe(false); + expect(pairLifecycleApplies("settle", executor({ archivedAt: NOW }))).toBe(false); + }); + + it("interrupts only an executor that is running", () => { + const running = executor({ + latestTurn: { ...turn, state: "running", completedAt: null }, + session: { + threadId: EXECUTOR, + status: "running", + providerName: "antigravity", + runtimeMode: "full-access", + activeTurnId: turn.turnId, + lastError: null, + updatedAt: NOW, + }, + }); + expect(pairLifecycleApplies("interrupt", running)).toBe(true); + expect(pairLifecycleApplies("interrupt", executor())).toBe(false); + expect(pairLifecycleApplies("interrupt", executor({ session: null, latestTurn: null }))).toBe( + false, + ); + expect(pairLifecycleApplies("interrupt", { ...running, archivedAt: NOW })).toBe(false); + }); +}); diff --git a/apps/server/src/orchestration/pairLifecycle.logic.ts b/apps/server/src/orchestration/pairLifecycle.logic.ts new file mode 100644 index 0000000000..b3e43b0ad9 --- /dev/null +++ b/apps/server/src/orchestration/pairLifecycle.logic.ts @@ -0,0 +1,66 @@ +/** + * Pure rules for what a lead's lifecycle means for its pair executor. + * + * @module orchestration/pairLifecycle.logic + */ +import type { OrchestrationEvent, OrchestrationThreadShell, ThreadId } from "@t3tools/contracts"; +import { isDelegatedThreadId } from "../mcp/toolkits/delegation/logic.ts"; +import { derivePairExecutorState } from "../mcp/toolkits/pair/logic.ts"; + +/** + * - `archive`, `delete`, `settle`: the executor follows its lead. + * - `interrupt`: the lead is being rewound. Both threads share one worktree, so + * an executor that keeps editing would write over the restored files. + */ +export type PairLifecycleAction = "archive" | "delete" | "settle" | "interrupt"; + +export interface PairLifecycleIntent { + readonly leadThreadId: ThreadId; + readonly action: PairLifecycleAction; +} + +/** What a domain event asks of the pair, or null when it asks nothing. */ +export function pairLifecycleIntent(event: OrchestrationEvent): PairLifecycleIntent | null { + if (event.metadata.historyImport === true) { + return null; + } + let action: PairLifecycleAction; + switch (event.type) { + case "thread.archived": + action = "archive"; + break; + case "thread.deleted": + action = "delete"; + break; + case "thread.settled": + action = "settle"; + break; + case "thread.checkpoint-revert-requested": + action = "interrupt"; + break; + default: + return null; + } + const leadThreadId = event.payload.threadId; + if (isDelegatedThreadId(leadThreadId)) { + return null; + } + return { leadThreadId, action }; +} + +/** Whether the action still has something to do for this executor. */ +export function pairLifecycleApplies( + action: PairLifecycleAction, + executor: OrchestrationThreadShell, +): boolean { + switch (action) { + case "archive": + return executor.archivedAt === null; + case "delete": + return true; + case "settle": + return executor.archivedAt === null && executor.settledOverride !== "settled"; + case "interrupt": + return derivePairExecutorState(executor) === "running"; + } +} diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts index 465f0c5ddf..794f2cb2b6 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.test.ts @@ -1,3 +1,5 @@ +import * as NodeCrypto from "node:crypto"; + import * as NodeServices from "@effect/platform-node/NodeServices"; import { ProjectId, @@ -434,6 +436,67 @@ describe("ProviderSessionReaper", () => { expect(Option.isSome(remaining)).toBe(true); }); + it.each([ + ["running", false], + ["starting", false], + ["ready", true], + ] as const)( + "while a pair's lead session is %s, reaping its idle executor is %s", + async (leadStatus, reaped) => { + const lead = ThreadId.make(`thread-reaper-pair-lead-${leadStatus}`); + const executor = ThreadId.make( + `delegated:${lead}:${NodeCrypto.createHash("sha256") + .update(`${lead}\npair`) + .digest("hex") + .slice(0, 16)}`, + ); + const now = "2026-01-01T00:00:00.000Z"; + const session = (threadId: ThreadId, status: "running" | "starting" | "ready") => ({ + threadId, + status, + providerName: "claudeAgent" as const, + runtimeMode: "full-access" as const, + activeTurnId: null, + lastError: null, + updatedAt: now, + }); + const harness = await createHarness({ + readModel: makeReadModel([ + { id: lead, session: session(lead, leadStatus) }, + { id: executor, session: session(executor, "ready") }, + ]), + }); + const repository = await runtime!.runPromise( + Effect.service(ProviderSessionRuntime.ProviderSessionRuntimeRepository), + ); + // Only the executor has a stale binding. A lead that is mid-turn, often + // blocked in pair_await, is about to brief it again: reaping it now would + // restart its provider process between every brief. + await runtime!.runPromise( + repository.upsert({ + threadId: executor, + providerName: "antigravity", + providerInstanceId: null, + adapterKey: "antigravity", + runtimeMode: "full-access", + status: "running", + lastSeenAt: "2026-04-14T00:00:00.000Z", + resumeCursor: { opaque: "resume-pair-executor" }, + runtimePayload: null, + }), + ); + + await startReaper(); + if (reaped) { + await waitFor(() => harness.stopSession.mock.calls.length === 1); + expect(harness.stopSession.mock.calls[0]?.[0]).toEqual({ threadId: executor }); + } else { + await Effect.runPromise(drainFibers); + expect(harness.stopSession).not.toHaveBeenCalled(); + } + }, + ); + it.each(["ready", "interrupted", "error"] as const)( "gives a long turn a full idle window after becoming %s", async (status) => { diff --git a/apps/server/src/provider/Layers/ProviderSessionReaper.ts b/apps/server/src/provider/Layers/ProviderSessionReaper.ts index 8f04bd1821..cb982e8f7b 100644 --- a/apps/server/src/provider/Layers/ProviderSessionReaper.ts +++ b/apps/server/src/provider/Layers/ProviderSessionReaper.ts @@ -4,6 +4,9 @@ import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schedule from "effect/Schedule"; +import * as NodeCrypto from "node:crypto"; +import { delegatedParentThreadId } from "@t3tools/shared/delegatedThreads"; +import { isPairExecutorThreadId } from "../../mcp/toolkits/pair/logic.ts"; import { ProjectionSnapshotQuery } from "../../orchestration/Services/ProjectionSnapshotQuery.ts"; import { ProviderSessionDirectory } from "../Services/ProviderSessionDirectory.ts"; @@ -94,6 +97,30 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) = continue; } + // A lead that is mid-turn, often blocked in pair_await, is about to + // brief the executor again, and reaping it would restart its provider + // process between every brief. + if ( + isPairExecutorThreadId(binding.threadId, (input) => + NodeCrypto.createHash("sha256").update(input).digest("hex"), + ) + ) { + const leadId = delegatedParentThreadId(binding.threadId); + if (leadId !== null) { + const lead = yield* projectionSnapshotQuery + .getThreadShellById(leadId) + .pipe(Effect.map(Option.getOrUndefined)); + if (lead?.session?.status === "running" || lead?.session?.status === "starting") { + yield* Effect.logDebug("provider.session.reaper.skipped-paired-lead-active", { + threadId: binding.threadId, + leadThreadId: leadId, + idleDurationMs, + }); + continue; + } + } + } + const reaped = yield* providerService.stopSession({ threadId: binding.threadId }).pipe( Effect.tap(() => Effect.logInfo("provider.session.reaped", { diff --git a/apps/server/src/provider/RuntimeInstructions.test.ts b/apps/server/src/provider/RuntimeInstructions.test.ts index 914254d8ac..cd301b86b5 100644 --- a/apps/server/src/provider/RuntimeInstructions.test.ts +++ b/apps/server/src/provider/RuntimeInstructions.test.ts @@ -76,12 +76,17 @@ describe("buildRuntimeInstructions", () => { it("states the rules a lead must not get wrong", () => { for (const rule of [ + "Work test-first", + "fail for the right reason", + "without editing them", + "code only", "pair_handoff", "pair_await", "Never poll in a loop", + "confirm your tests are unchanged", "re-run the checks yourself", "not verification", - "Only you commit", + "Only you push and open pull requests", "never approve on their behalf", ]) { expect(PAIR_LEAD_PROTOCOL).toContain(rule); diff --git a/apps/server/src/provider/RuntimeInstructions.ts b/apps/server/src/provider/RuntimeInstructions.ts index 6e40f9e535..fc6bad024e 100644 --- a/apps/server/src/provider/RuntimeInstructions.ts +++ b/apps/server/src/provider/RuntimeInstructions.ts @@ -33,8 +33,9 @@ Keep small or tightly coupled work local. Before choosing a delegation method fo * pair started mid-session reaches the lead before its next session start does. */ export const PAIR_LEAD_PROTOCOL = `You are the lead of a Pylon pair. One executor thread on a faster, cheaper model is linked to this thread and works in your worktree. Hand implementation to it instead of using your own subagents. -Brief it with pair_handoff: exact files, the behavior wanted, the acceptance checks to run, exclusions, and the report format. Hand off a whole plan step, not small nudges; write tests or exact acceptance checks first when you can. Wait with pair_await, or end your turn: Pylon wakes you when the executor finishes or needs the user. Never poll in a loop. -When it reports, read the diff in your worktree and re-run the checks yourself. The executor's report is not verification, and its prose is less reliable than its code. Send one consolidated correction, at most three rounds, then ask the user. Only you commit, push, and open pull requests. Make small or tightly coupled changes yourself. If the executor needs an approval or an answer, tell the user; never approve on their behalf.`; +Work test-first. Write the contract (types, signatures, stubs) and the failing tests yourself, run them, and confirm they fail for the right reason. Commit them. Then brief the executor with pair_handoff: make these tests pass without editing them, with the exact files, the behavior wanted, the commands to run, the exclusions, and the report format. Hand off a whole plan step, not small nudges. Where tests cannot express the work, give exact acceptance checks instead. Ask the executor for code only, and write documentation and pull request text yourself. +Wait with pair_await, or end your turn: Pylon wakes you when the executor finishes or needs the user. Never poll in a loop. +When it reports, confirm your tests are unchanged (git diff --stat on those paths is empty), read the diff in your worktree, and re-run the checks yourself. The executor's report is not verification, and its prose is less reliable than its code. Send one consolidated correction, at most three rounds, then ask the user. Only you push and open pull requests. Make small or tightly coupled changes yourself. If the executor needs an approval or an answer, tell the user; never approve on their behalf.`; /** * Included only while this thread has a pair executor. It replaces the diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index a044d6dec4..dbbf3ad0f6 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -77,6 +77,7 @@ import { ThreadDeletionReactorLive } from "./orchestration/Layers/ThreadDeletion import * as RollbackSagaRunner from "./rollback/RollbackSagaRunner.ts"; import * as ThreadSettlementReactor from "./orchestration/ThreadSettlementReactor.ts"; import * as DelegationFollowThroughReactor from "./orchestration/DelegationFollowThroughReactor.ts"; +import * as PairLifecycleReactor from "./orchestration/PairLifecycleReactor.ts"; import * as PullRequestSyncReactor from "./orchestration/PullRequestSyncReactor.ts"; import * as ProjectSettingsReactor from "./orchestration/ProjectSettingsReactor.ts"; import * as ThreadPullRequestReactor from "./orchestration/ThreadPullRequestReactor.ts"; @@ -304,6 +305,7 @@ const ReactorLayerLive = Layer.empty.pipe( Layer.provideMerge(ThreadDeletionReactorLive), Layer.provideMerge(ThreadSettlementReactor.layer), Layer.provideMerge(DelegationFollowThroughReactor.layer), + Layer.provideMerge(PairLifecycleReactor.layer), Layer.provideMerge(PullRequestSyncReactor.layer), Layer.provideMerge(ThreadPullRequestReactor.layer), Layer.provideMerge(ProjectSettingsReactor.layer), diff --git a/docs/internals/delegation.md b/docs/internals/delegation.md index 0827b7a66d..e77a927432 100644 --- a/docs/internals/delegation.md +++ b/docs/internals/delegation.md @@ -109,6 +109,18 @@ same protocol text; the tool denial applies from the next session start. Antigra because it offers no per-session control over its own subagents, but it can be the executor. Prime Agent's own subagent depth is not yet held while paired. +The executor follows its lead. `PairLifecycleReactor` watches domain events and dispatches existing +commands: archiving, settling, or deleting a lead does the same to its executor, including an +executor that was archived when the pair was turned off. An executor's own lifecycle never echoes +back, because archiving it is how a user turns the pair off, and unarchiving a lead does not turn a +pair back on. A rewind of the lead is not blocked, since that would mean changing rewind admission in +the decider. Instead the reactor interrupts a running executor the moment the rewind is requested: the +two threads share one worktree, and an executor that kept editing would write over the restored +files. It may still write for a moment before the interrupt lands. The idle-session reaper skips an +executor while its lead's session is starting or running, so it is not restarted between briefs. The +reactor reads no settings: following a lead is cleanup and keeps working after delegation is turned +off. + ## Accepted limits - A follow-up is refused while the child is running, but a user message can arrive between the diff --git a/docs/user/agent-delegation.md b/docs/user/agent-delegation.md index fd56d8ebab..c085a7f272 100644 --- a/docs/user/agent-delegation.md +++ b/docs/user/agent-delegation.md @@ -137,7 +137,8 @@ unaffected; this takes effect the next time the lead's session starts. Antigravi executor but cannot lead a pair. The executor appears under its parent in the sidebar like any delegated thread. Archive it to turn -the pair off for that thread. Pairing needs **Pylon delegation** turned on in **Settings → +the pair off for that thread. Archiving, settling, or deleting the lead does the same to its +executor. Rewinding the lead stops the executor first, because both work on the same files. Pairing needs **Pylon delegation** turned on in **Settings → Integrations**. ## Things to know