From d791b7185b73b927453744f68633933035396ca7 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Fri, 18 Sep 2026 09:40:16 -0600 Subject: [PATCH 1/6] test(pair): contract and failing tests for the pair lifecycle --- .../PairLifecycleReactor.test.ts | 331 ++++++++++++++++++ .../src/orchestration/PairLifecycleReactor.ts | 30 ++ .../orchestration/pairLifecycle.logic.test.ts | 156 +++++++++ .../src/orchestration/pairLifecycle.logic.ts | 34 ++ .../Layers/ProviderSessionReaper.test.ts | 63 ++++ 5 files changed, 614 insertions(+) create mode 100644 apps/server/src/orchestration/PairLifecycleReactor.test.ts create mode 100644 apps/server/src/orchestration/PairLifecycleReactor.ts create mode 100644 apps/server/src/orchestration/pairLifecycle.logic.test.ts create mode 100644 apps/server/src/orchestration/pairLifecycle.logic.ts 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..375db70a14 --- /dev/null +++ b/apps/server/src/orchestration/PairLifecycleReactor.ts @@ -0,0 +1,30 @@ +/** + * 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. + * + * STUB: written by the lead so the tests compile. The executor replaces `make`; + * the service tag, its shape, and `layer` must not change. + * + * @module orchestration/PairLifecycleReactor + */ +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import type * as Scope from "effect/Scope"; + +export class PairLifecycleReactor extends Context.Service< + PairLifecycleReactor, + { + readonly start: () => Effect.Effect; + readonly drain: Effect.Effect; + } +>()("t3/orchestration/PairLifecycleReactor") {} + +export const make = Effect.sync(() => ({ + // Inert until implemented: the tests fail because nothing is dispatched. + start: (): Effect.Effect => Effect.void, + drain: Effect.void, +})); + +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..1036364ae6 --- /dev/null +++ b/apps/server/src/orchestration/pairLifecycle.logic.ts @@ -0,0 +1,34 @@ +/** + * Pure rules for what a lead's lifecycle means for its pair executor. + * + * STUB: written by the lead so the tests compile. The executor replaces the + * bodies; exported names and signatures must not change. + * + * @module orchestration/pairLifecycle.logic + */ +import type { OrchestrationEvent, OrchestrationThreadShell, ThreadId } from "@t3tools/contracts"; + +/** + * - `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 { + throw new Error("pairLifecycle.logic.pairLifecycleIntent is not implemented"); +} + +/** Whether the action still has something to do for this executor. */ +export function pairLifecycleApplies( + _action: PairLifecycleAction, + _executor: OrchestrationThreadShell, +): boolean { + throw new Error("pairLifecycle.logic.pairLifecycleApplies is not implemented"); +} 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) => { From 349243eef75bbfc578cdea0650d09d0d70ff140c Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Fri, 18 Sep 2026 09:41:44 -0600 Subject: [PATCH 2/6] feat(pair): pure rules for what a lead's lifecycle means for its executor --- .../src/orchestration/pairLifecycle.logic.ts | 48 +++++++++++++++---- 1 file changed, 40 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration/pairLifecycle.logic.ts b/apps/server/src/orchestration/pairLifecycle.logic.ts index 1036364ae6..b3e43b0ad9 100644 --- a/apps/server/src/orchestration/pairLifecycle.logic.ts +++ b/apps/server/src/orchestration/pairLifecycle.logic.ts @@ -1,12 +1,11 @@ /** * Pure rules for what a lead's lifecycle means for its pair executor. * - * STUB: written by the lead so the tests compile. The executor replaces the - * bodies; exported names and signatures must not change. - * * @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. @@ -21,14 +20,47 @@ export interface PairLifecycleIntent { } /** What a domain event asks of the pair, or null when it asks nothing. */ -export function pairLifecycleIntent(_event: OrchestrationEvent): PairLifecycleIntent | null { - throw new Error("pairLifecycle.logic.pairLifecycleIntent is not implemented"); +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, + action: PairLifecycleAction, + executor: OrchestrationThreadShell, ): boolean { - throw new Error("pairLifecycle.logic.pairLifecycleApplies is not implemented"); + 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"; + } } From 2df37c0042d12ecf35621af5071265cef85637c8 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Fri, 18 Sep 2026 09:43:16 -0600 Subject: [PATCH 3/6] feat(pair): keep the executor in step with its lead --- .../src/orchestration/PairLifecycleReactor.ts | 130 ++++++++++++++++-- 1 file changed, 122 insertions(+), 8 deletions(-) diff --git a/apps/server/src/orchestration/PairLifecycleReactor.ts b/apps/server/src/orchestration/PairLifecycleReactor.ts index 375db70a14..19ad93e8f0 100644 --- a/apps/server/src/orchestration/PairLifecycleReactor.ts +++ b/apps/server/src/orchestration/PairLifecycleReactor.ts @@ -3,15 +3,32 @@ * together, and stopped when the lead is rewound. A sidecar over existing * commands and domain events; it adds no command, event, or decider branch. * - * STUB: written by the lead so the tests compile. The executor replaces `make`; - * the service tag, its shape, and `layer` must not change. - * * @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, @@ -21,10 +38,107 @@ export class PairLifecycleReactor extends Context.Service< } >()("t3/orchestration/PairLifecycleReactor") {} -export const make = Effect.sync(() => ({ - // Inert until implemented: the tests fail because nothing is dispatched. - start: (): Effect.Effect => Effect.void, - drain: Effect.void, -})); +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); From 519b8b0db090ffe78484475af37dd9bb205f793a Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Fri, 18 Sep 2026 09:45:00 -0600 Subject: [PATCH 4/6] feat(pair): start the pair lifecycle reactor with the server --- .../OrchestrationEngineHarness.integration.ts | 7 +++++++ .../orchestration/Layers/OrchestrationReactor.test.ts | 11 +++++++++++ .../src/orchestration/Layers/OrchestrationReactor.ts | 3 +++ apps/server/src/server.ts | 2 ++ 4 files changed, 23 insertions(+) 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/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), From 4e22745238fabf6ceec4ae9db68b43ef0b64ea24 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Fri, 18 Sep 2026 09:45:52 -0600 Subject: [PATCH 5/6] feat(pair): do not reap an executor while its lead is mid-turn --- .../provider/Layers/ProviderSessionReaper.ts | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) 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", { From 9fa8cfc019276cf0d594c01755770eddd3cf1113 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Fri, 18 Sep 2026 09:48:46 -0600 Subject: [PATCH 6/6] docs(pair): make the lead protocol test-first and describe the pair lifecycle --- apps/server/src/provider/RuntimeInstructions.test.ts | 7 ++++++- apps/server/src/provider/RuntimeInstructions.ts | 5 +++-- docs/internals/delegation.md | 12 ++++++++++++ docs/user/agent-delegation.md | 3 ++- 4 files changed, 23 insertions(+), 4 deletions(-) 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/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