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
4 changes: 4 additions & 0 deletions apps/server/src/orchestration-v2/ProjectionSettlement.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,10 @@ it.effect.each([
[open, null, 2],
],
);
assert.deepEqual(
(yield* store.getThreadsWithPullRequests(open)).map((thread) => thread.id),
[open],
);
}).pipe(Effect.provide(testLayer)),
);

Expand Down
19 changes: 11 additions & 8 deletions apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -349,12 +349,12 @@ export interface ProjectionStoreV2Shape {
) => Effect.Effect<ReadonlyArray<ProjectionSettlementCandidate>, ProjectionStoreV2Error>;
/**
* Active (not deleted, not archived) threads with at least one pull request
* link, in shell snapshot order. Skips run, message and item reads.
* link, in shell snapshot order, or only `threadId` when given. Skips run,
* message and item reads.
*/
readonly getThreadsWithPullRequests: () => Effect.Effect<
ReadonlyArray<ProjectionThreadPullRequests>,
ProjectionStoreV2Error
>;
readonly getThreadsWithPullRequests: (
threadId?: ThreadId,
) => Effect.Effect<ReadonlyArray<ProjectionThreadPullRequests>, ProjectionStoreV2Error>;
readonly getTurnStartContext: (
threadId: ThreadId,
runId: RunId,
Expand Down Expand Up @@ -5133,12 +5133,14 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
)
.pipe(Effect.mapError((cause) => new ProjectionStoreSetupError({ cause })));

const getThreadsWithPullRequests: ProjectionStoreV2Shape["getThreadsWithPullRequests"] = () =>
const getThreadsWithPullRequests: ProjectionStoreV2Shape["getThreadsWithPullRequests"] = (
threadId,
) =>
Effect.gen(function* () {
const rows = yield* sql<PayloadRow>`
SELECT payload_json
FROM orchestration_v2_projection_threads
WHERE deleted_at IS NULL
WHERE deleted_at IS NULL${threadId === undefined ? sql`` : sql` AND thread_id = ${threadId}`}
AND json_extract(payload_json, '$.archivedAt') IS NULL
AND json_array_length(payload_json, '$.pullRequests') > 0
ORDER BY updated_at ASC, thread_id ASC
Expand Down Expand Up @@ -5583,13 +5585,14 @@ export const layerMemory: Layer.Layer<ProjectionStoreV2> = Layer.effect(
left.id.localeCompare(right.id),
);
}),
getThreadsWithPullRequests: () =>
getThreadsWithPullRequests: (threadId) =>
Ref.get(replayState).pipe(
Effect.map((state) =>
[...state.projections.values()]
.map(({ thread }) => thread)
.filter(
(thread) =>
(threadId === undefined || thread.id === threadId) &&
thread.deletedAt === null &&
thread.archivedAt === null &&
(thread.pullRequests ?? []).length > 0,
Expand Down
137 changes: 134 additions & 3 deletions apps/server/src/orchestration-v2/PullRequestSyncReactor.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,16 @@
import * as Stream from "effect/Stream";
import {
EventId,
MessageId,
ProjectId,
ProviderInstanceId,
PullRequestOperationError,
RunId,
ThreadId,
TurnItemId,
type OrchestrationV2Command as OrchestrationCommand,
type OrchestrationV2DomainEvent,
type OrchestrationV2Run,
type OrchestrationProjectShell,
type PullRequestRef,
type PullRequestStack,
Expand Down Expand Up @@ -171,6 +177,7 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:
const linkCommands = yield* Ref.make<ReadonlyArray<LinkCommand>>([]);
const summaryCalls = yield* Ref.make<ReadonlyArray<PullRequestRef>>([]);
const stackCalls = yield* Ref.make<ReadonlyArray<PullRequestRef>>([]);
const domainEvents = yield* Queue.unbounded<OrchestrationV2DomainEvent>();

const summary: PullRequestService.PullRequestService["Service"]["summary"] = (
input,
Expand Down Expand Up @@ -211,12 +218,17 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:
}),
Layer.mock(ProjectionStore.ProjectionStoreV2)({
// Mirrors the store's filter: active threads that have at least one link.
getThreadsWithPullRequests: () =>
getThreadsWithPullRequests: (threadId) =>
Queue.offer(snapshotReads, undefined).pipe(
Effect.andThen(Ref.get(snapshots)),
Effect.map((snapshot) =>
snapshot.threads
.filter((thread) => thread.archivedAt === null && thread.pullRequests.length > 0)
.filter(
(thread) =>
(threadId === undefined || thread.id === threadId) &&
thread.archivedAt === null &&
thread.pullRequests.length > 0,
)
.map((thread) => ({
id: thread.id,
projectId: thread.projectId,
Expand All @@ -233,7 +245,7 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:
Effect.andThen(Effect.die(new Error("pull request sync must not read the shell"))),
),
dispatch,
streamDomainEvents: Stream.empty,
streamDomainEvents: Stream.fromQueue(domainEvents),
}),
Layer.succeed(ServerActivation.ServerActivation, Deferred.await(activation)),
Layer.succeed(Crypto.Crypto, testCrypto),
Expand All @@ -248,6 +260,7 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:
linkCommands,
summaryCalls,
stackCalls,
domainEvents,
layer: PullRequestSyncReactor.layer.pipe(Layer.provide(dependencies)),
};
});
Expand All @@ -272,6 +285,68 @@ const sweepAgain = Effect.fn("sweepPullRequestSyncHarness")(function* (
yield* reactor.drain;
});

function runUpdated(
threadId: ThreadId,
status: OrchestrationV2Run["status"],
): OrchestrationV2DomainEvent {
const runId = RunId.make(`run:${threadId}:1`);
const at = DateTime.makeUnsafe(NOW);
return {
type: "run.updated",
id: EventId.make(`event:${threadId}:${status}`),
threadId,
runId,
occurredAt: at,
payload: {
id: runId,
threadId,
ordinal: 1,
providerInstanceId: ProviderInstanceId.make("codex"),
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5" },
providerThreadId: null,
userMessageId: MessageId.make(`message:${threadId}:1`),
rootNodeId: null,
activeAttemptId: null,
status,
requestedAt: at,
startedAt: at,
completedAt: at,
checkpointId: null,
contextHandoffId: null,
},
};
}

function commandRan(threadId: ThreadId, input: string): OrchestrationV2DomainEvent {
const runId = RunId.make(`run:${threadId}:1`);
const at = DateTime.makeUnsafe(NOW);
return {
type: "turn-item.updated",
id: EventId.make(`event:${threadId}:command`),
threadId,
runId,
occurredAt: at,
payload: {
id: TurnItemId.make(`item:${threadId}:command`),
threadId,
runId,
nodeId: null,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: 1,
status: "completed",
title: null,
startedAt: at,
completedAt: at,
updatedAt: at,
type: "command_execution",
input,
},
};
}

/** What the reactor would have persisted, so the next sweep sees its own writes. */
function applySync(
snapshot: TestShellSnapshot,
Expand Down Expand Up @@ -832,4 +907,60 @@ describe("PullRequestSyncReactor", () => {
}),
),
);

it.effect("reads open links fresh only when a run that ran a merge command ends", () =>
Effect.scoped(
Effect.gen(function* () {
yield* TestClock.setTime(Date.parse(NOW));
const state = yield* Ref.make<PullRequestSummary["state"]>("open");
const invalidated = yield* Ref.make<ReadonlyArray<number>>([]);
const fixture = yield* makeHarness({
snapshot: makeSnapshot([
makeThread("agent", { pullRequests: [makeLink(7, { state: "open" })] }),
makeThread("other", { pullRequests: [makeLink(9, { state: "open" })] }),
]),
summary: (input) =>
Ref.get(state).pipe(
Effect.map((current) =>
makeSummary(input, current === "merged" ? { state: current, mergedAt: NOW } : {}),
),
),
invalidate: ({ reference }) =>
Ref.update(invalidated, (numbers) => [...numbers, reference?.number ?? -1]),
});

yield* Effect.gen(function* () {
const reactor = yield* startAndSweep(fixture);
// The agent merges from a shell during its turn, inside the summary cache window.
yield* Ref.set(state, "merged");
yield* Ref.set(fixture.summaryCalls, []);

yield* Queue.offerAll(fixture.domainEvents, [
// A run that only reads its pull request costs no host read when it ends.
commandRan(ThreadId.make("other"), "gh pr view 9"),
runUpdated(ThreadId.make("other"), "completed"),
commandRan(ThreadId.make("agent"), "gh pr merge 7 --squash 2>&1 | tail -3"),
runUpdated(ThreadId.make("agent"), "completed"),
]);
// The agent thread's lookup, then the requested sweep.
yield* Queue.take(fixture.snapshotReads);
yield* Queue.take(fixture.snapshotReads);
yield* reactor.drain;

assert.deepStrictEqual(yield* Ref.get(invalidated), [7]);
assert.deepStrictEqual(
(yield* Ref.get(fixture.summaryCalls)).map((call) => call.number),
[7],
);
assert.deepStrictEqual(
(yield* Ref.get(fixture.syncCommands)).map((command) => [
command.number,
command.snapshot.state,
]),
[[7, "merged"]],
);
}).pipe(Effect.provide(fixture.layer));
}),
),
);
});
54 changes: 48 additions & 6 deletions apps/server/src/orchestration-v2/PullRequestSyncReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { siblingPullRequestUrl } from "@t3tools/shared/changeRequestUrl";
import {
CommandId,
type PullRequestSummary,
type ThreadId,
type ThreadPullRequestKey,
type ThreadPullRequestLink,
type ThreadPullRequestSnapshot,
Expand Down Expand Up @@ -29,8 +30,11 @@ import * as PullRequestService from "../pullRequest/PullRequestService.ts";
import { forkParked } from "../serverActivation.ts";
import * as Orchestrator from "./Orchestrator.ts";
import * as ProjectionStore from "./ProjectionStore.ts";
import { isTerminalRunStatus } from "./ThreadManagementService.ts";

const SLOW_SYNC_INTERVAL_MS = 15 * 60 * 1_000;
/** Shell commands that can merge or close a pull request without a merge notification. */
const PULL_REQUEST_CLOSE_COMMAND = /\b(?:gh\s+pr|glab\s+mr)\s+(?:merge|close)\b/u;

type SnapshotFields = Omit<ThreadPullRequestSnapshot, "syncedAt">;

Expand Down Expand Up @@ -330,22 +334,60 @@ export const make = Effect.gen(function* () {
}).pipe(Effect.catchCause(logSkipped("pull request sync sweep failed", {}))),
);

// Threads whose current run ran a merge or close command, until that run ends.
const closeCommandThreads = new Set<ThreadId>();
const refreshOpenLinks = (threadId: ThreadId) =>
projections.getThreadsWithPullRequests(threadId).pipe(
Effect.flatMap((threads) =>
Effect.forEach(
threads.flatMap((thread) =>
visibleThreadPullRequests(thread.pullRequests ?? []).filter(
(link) => link.snapshot?.state === "open",
),
),
requestSync,
{ discard: true },
),
),
Effect.catchCause(logSkipped("pull request refresh after run skipped", { threadId })),
);

const start: PullRequestSyncReactor["Service"]["start"] = Effect.fn(
"PullRequestSyncReactor.start",
)(function* () {
const events = engine.streamDomainEvents;
yield* forkParked(
Stream.runForEach(events, (event) =>
event.type === "thread.pull-request-synced"
? Effect.forEach(
Stream.runForEach(events, (event) => {
switch (event.type) {
case "thread.pull-request-synced":
return Effect.forEach(
visibleThreadPullRequests(event.payload.pullRequests ?? []).filter(
(link) => link.snapshot === null,
),
requestSync,
{ discard: true },
)
: Effect.void,
).pipe(Effect.catchCause(logSkipped("pull request sync event stream failed", {}))),
);
// An agent can merge or close its pull request from a shell (`gh pr merge`), which
// sends no merge notification. When a run that ran such a command ends, read the
// thread's open links fresh, so settlement does not wait for the next sweep and the
// cached summary. Other runs add no host reads.
case "turn-item.updated":
if (
event.payload.type === "command_execution" &&
PULL_REQUEST_CLOSE_COMMAND.test(event.payload.input)
) {
closeCommandThreads.add(event.threadId);
}
return Effect.void;
case "run.updated":
return isTerminalRunStatus(event.payload.status) &&
closeCommandThreads.delete(event.threadId)
? refreshOpenLinks(event.threadId)
: Effect.void;
default:
return Effect.void;
}
}).pipe(Effect.catchCause(logSkipped("pull request sync event stream failed", {}))),
);
yield* forkParked(
Effect.gen(function* () {
Expand Down
7 changes: 5 additions & 2 deletions docs/internals/overview.md
Original file line number Diff line number Diff line change
Expand Up @@ -87,8 +87,11 @@ must reject that operation before changing the filesystem.
Thread settlement is server-owned. The
[settlement service](../../apps/server/src/orchestration-v2/ThreadSettlementService.ts) evaluates PR
and inactivity settings without a connected client. Merge notifications invalidate cached PR state
and trigger a check. The guarded `thread.auto-settle` command rejects newer activity, explicit
settlement overrides, and live or blocked work. It records the activity timestamp for stable
and trigger a check. A merge outside T3, such as an agent running `gh pr merge`, sends no
notification, so the [PR sync reactor](../../apps/server/src/orchestration-v2/PullRequestSyncReactor.ts)
re-reads a thread's open links when a run that ran a merge or close command ends. The guarded
`thread.auto-settle` command rejects newer activity, explicit settlement overrides, and live or
blocked work. It records the activity timestamp for stable
sorting and detaches idle provider sessions. Clients render the persisted result; they do not
derive settlement from their own clocks or PR caches.

Expand Down
Loading