Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions apps/server/src/orchestration/PairLifecycleReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ 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 * as TestClock from "effect/testing/TestClock";

import { ServerActivation } from "../serverActivation.ts";
import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts";
Expand Down Expand Up @@ -140,6 +141,8 @@ const makeHarness = Effect.fn(function* (input: {
Effect.succeed(Option.fromNullishOr(active.find((thread) => thread.id === id))),
getArchivedShellSnapshot: () =>
Effect.succeed({ snapshotSequence: 1, projects: [], threads: archived, updatedAt: NOW }),
getShellSnapshot: () =>
Effect.succeed({ snapshotSequence: 1, projects: [], threads: active, updatedAt: NOW }),
}),
Layer.mock(OrchestrationEngineService)({
subscribeDomainEvents: PubSub.subscribe(events).pipe(Effect.map(Stream.fromSubscription)),
Expand Down Expand Up @@ -248,6 +251,32 @@ describe("PairLifecycleReactor", () => {
),
);

it.effect("deletes executors left behind by abandoned drafts when the server starts", () =>
Effect.scoped(
Effect.gen(function* () {
// The test clock starts at the epoch, where nothing is a day old.
yield* TestClock.setTime(Date.parse(NOW));
const longAgo = "2026-01-01T00:00:00.000Z";
const abandoned = executorOf(ThreadId.make("abandoned-draft"));
const archivedLead = ThreadId.make("archived-lead");
const h = yield* makeHarness({
active: [
shell(LEAD),
shell(EXECUTOR, { createdAt: longAgo }),
shell(abandoned, { createdAt: longAgo }),
// Its lead is archived, not missing: a known thread.
shell(executorOf(archivedLead), { createdAt: longAgo }),
],
archived: [shell(archivedLead, { archivedAt: NOW })],
});
assert.deepStrictEqual(
h.commands.map((command) => [command.type, "threadId" in command && command.threadId]),
[["thread.delete", abandoned]],
);
}),
),
);

it.effect("does nothing for threads without an executor or for other threads' events", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
36 changes: 36 additions & 0 deletions apps/server/src/orchestration/PairLifecycleReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import * as OrchestrationEngine from "./Services/OrchestrationEngine.ts";
import * as ProjectionSnapshotQuery from "./Services/ProjectionSnapshotQuery.ts";
import {
type PairLifecycleIntent,
orphanedExecutorIds,
pairLifecycleApplies,
pairLifecycleIntent,
} from "./pairLifecycle.logic.ts";
Expand Down Expand Up @@ -133,7 +134,42 @@ export const make = Effect.gen(function* () {
return worker.enqueue({ event, intent });
};

const sweep = Effect.gen(function* () {
const activeSnapshot = yield* snapshots.getShellSnapshot();
const archivedSnapshot = yield* snapshots.getArchivedShellSnapshot();
const knownThreadIds = new Set<string>();
for (const thread of activeSnapshot.threads) {
knownThreadIds.add(thread.id);
}
for (const thread of archivedSnapshot.threads) {
knownThreadIds.add(thread.id);
}
const nowMs = DateTime.toEpochMillis(yield* DateTime.now);
const orphanIds = orphanedExecutorIds({
threads: activeSnapshot.threads,
knownThreadIds,
nowMs,
});
for (const executorId of orphanIds) {
const commandId = CommandId.make(`server:pair-lifecycle:orphan:${executorId}`);
yield* engine.dispatch({
type: "thread.delete",
commandId,
threadId: executorId,
});
}
}).pipe(
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause)
? Effect.interrupt
: Effect.logWarning("Pylon pair orphan sweep failed", {
cause: Cause.pretty(cause),
}),
),
);

const start = Effect.fn("PairLifecycleReactor.start")(function* () {
yield* sweep;
const events = yield* engine.subscribeDomainEvents;
yield* forkParked(Stream.runForEach(events, processEvent));
});
Expand Down
56 changes: 55 additions & 1 deletion apps/server/src/orchestration/pairLifecycle.logic.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,13 @@ import {
} from "@t3tools/contracts";
import { describe, expect, it } from "vite-plus/test";

import { pairLifecycleApplies, pairLifecycleIntent } from "./pairLifecycle.logic.ts";
import { pairExecutorThreadId } from "@t3tools/shared/delegatedThreads";
import {
ORPHAN_EXECUTOR_GRACE_MS,
orphanedExecutorIds,
pairLifecycleApplies,
pairLifecycleIntent,
} from "./pairLifecycle.logic.ts";

const NOW = "2026-09-18T00:00:00.000Z";
const LEAD = ThreadId.make("lead:with:colons");
Expand Down Expand Up @@ -154,3 +160,51 @@ describe("pairLifecycleApplies", () => {
expect(pairLifecycleApplies("interrupt", { ...running, archivedAt: NOW })).toBe(false);
});
});

describe("orphaned pair executors", () => {
const nowMs = Date.parse(NOW);
// One second past the grace period before NOW.
const old = "2026-09-16T23:59:59.000Z";
const gone = ThreadId.make("abandoned-draft");
const orphan = executor({ id: pairExecutorThreadId(gone), createdAt: old, latestTurn: null });

it("finds a never-briefed executor whose lead never came to exist", () => {
expect(
orphanedExecutorIds({ threads: [orphan], knownThreadIds: new Set([orphan.id]), nowMs }),
).toEqual([orphan.id]);
});

it("leaves alone anything a person could still be using", () => {
const known = new Set<string>([orphan.id]);
const cases: ReadonlyArray<readonly [string, OrchestrationThreadShell, ReadonlySet<string>]> = [
["a lead that exists", orphan, new Set([orphan.id, gone])],
["a draft paired a moment ago", { ...orphan, createdAt: NOW }, known],
["an executor that ran", { ...orphan, latestTurn: turn }, known],
[
"an executor with a session",
{
...orphan,
session: {
threadId: orphan.id,
status: "ready",
providerName: "antigravity",
runtimeMode: "full-access",
activeTurnId: null,
lastError: null,
updatedAt: NOW,
},
},
known,
],
[
"a fan-out child",
{ ...orphan, id: ThreadId.make(`delegated:${gone}:0123456789abcdef`) },
known,
],
["an ordinary thread", { ...orphan, id: ThreadId.make("plain") }, known],
];
for (const [label, thread, knownThreadIds] of cases) {
expect(orphanedExecutorIds({ threads: [thread], knownThreadIds, nowMs }), label).toEqual([]);
}
});
});
37 changes: 37 additions & 0 deletions apps/server/src/orchestration/pairLifecycle.logic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import type { OrchestrationEvent, OrchestrationThreadShell, ThreadId } from "@t3tools/contracts";
import { isDelegatedThreadId } from "../mcp/toolkits/delegation/logic.ts";
import { derivePairExecutorState } from "../mcp/toolkits/pair/logic.ts";
import { delegatedParentThreadId, isPairExecutorThreadId } from "@t3tools/shared/delegatedThreads";

/**
* - `archive`, `delete`, `settle`: the executor follows its lead.
Expand Down Expand Up @@ -64,3 +65,39 @@ export function pairLifecycleApplies(
return derivePairExecutorState(executor) === "running";
}
}

/** How long a never-briefed executor may wait for a lead that does not exist yet. */
export const ORPHAN_EXECUTOR_GRACE_MS = 24 * 60 * 60 * 1_000;

/**
* Pair executors nothing will ever brief. Turning Pair on in a draft creates the
* executor before the lead thread exists, so a missing lead is normal for a
* while; an abandoned draft, or a draft whose id changed before its first send,
* leaves that executor behind for good. Only an executor that never ran, whose
* lead is unknown, and that is older than the grace period counts.
*/
export function orphanedExecutorIds(input: {
readonly threads: ReadonlyArray<OrchestrationThreadShell>;
/** Every thread id the server knows, archived ones included. */
readonly knownThreadIds: ReadonlySet<string>;
readonly nowMs: number;
}): ReadonlyArray<ThreadId> {
const orphans: ThreadId[] = [];
for (const thread of input.threads) {
if (!isPairExecutorThreadId(thread.id)) {
continue;
}
const leadId = delegatedParentThreadId(thread.id);
if (leadId === null || input.knownThreadIds.has(leadId)) {
continue;
}
if (thread.latestTurn !== null || thread.session !== null) {
continue;
}
if (input.nowMs - Date.parse(thread.createdAt) <= ORPHAN_EXECUTOR_GRACE_MS) {
continue;
}
orphans.push(thread.id);
}
return orphans;
}
7 changes: 7 additions & 0 deletions docs/internals/delegation.md
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,13 @@ executor while its lead's session is starting or running, so it is not restarted
reactor reads no settings: following a lead is cleanup and keeps working after delegation is turned
off.

The same reactor sweeps once when the server starts. Turning Pair on in a draft creates the executor
before its lead thread exists, so an abandoned draft, or one whose id changed before its first send,
leaves an executor nothing will ever brief. `orphanedExecutorIds` finds them conservatively: a pair
executor that never ran, whose lead is not among the active or archived threads, and that is more
than a day old. A missing lead alone proves nothing, because that is what every freshly paired draft
looks like. The delete uses a deterministic command id, and a failed sweep only logs.

A brief can name the files its lead owns, normally its tests and contract, as `protectedPaths`.
The handlers hash each file when the brief is accepted and `pair_await` reports the ones whose
content changed or that disappeared, once the executor is no longer running. This is what makes
Expand Down
Loading