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
5 changes: 5 additions & 0 deletions apps/server/src/checkpointing/CheckpointDiffQuery.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ describe("CheckpointDiffQuery.layer", () => {
getShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request the orchestration shell snapshot"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request archived shell snapshots"),
getSnapshotSequence: () => Effect.succeed({ snapshotSequence: 0 }),
Expand Down Expand Up @@ -201,6 +202,7 @@ describe("CheckpointDiffQuery.layer", () => {
getShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request the orchestration shell snapshot"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request archived shell snapshots"),
getSnapshotSequence: () => Effect.succeed({ snapshotSequence: 0 }),
Expand Down Expand Up @@ -293,6 +295,7 @@ describe("CheckpointDiffQuery.layer", () => {
getShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request the orchestration shell snapshot"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request archived shell snapshots"),
getSnapshotSequence: () => Effect.succeed({ snapshotSequence: 0 }),
Expand Down Expand Up @@ -370,6 +373,7 @@ describe("CheckpointDiffQuery.layer", () => {
getShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request the orchestration shell snapshot"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request archived shell snapshots"),
getSnapshotSequence: () => Effect.succeed({ snapshotSequence: 0 }),
Expand Down Expand Up @@ -432,6 +436,7 @@ describe("CheckpointDiffQuery.layer", () => {
getShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request the orchestration shell snapshot"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () =>
Effect.die("CheckpointDiffQuery should not request archived shell snapshots"),
getSnapshotSequence: () => Effect.succeed({ snapshotSequence: 0 }),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -435,6 +435,7 @@ describe("OrchestrationEngine", () => {
updatedAt: projectionSnapshot.updatedAt,
}),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () =>
Effect.succeed({
snapshotSequence: projectionSnapshot.snapshotSequence,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3503,6 +3503,77 @@ it.effect("omits foreign-host PRs from legacy snapshots while preserving native
}).pipe(Effect.provide(layer));
});

it.effect(
"lists linked threads like the shell snapshot, in one query and without identities",
() => {
const resolved: string[] = [];
const layer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(
Layer.succeed(RepositoryIdentityResolver.RepositoryIdentityResolver, {
resolve: (root) =>
Effect.sync(() => {
resolved.push(root);
return null;
}),
}),
),
Layer.provideMerge(SqlitePersistenceMemory),
);
return Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const query = yield* ProjectionSnapshotQuery;
yield* sql`INSERT INTO projection_projects (project_id, title, workspace_root, scripts_json, created_at, updated_at)
VALUES ('p1', 'One', '/one', '[]', '2026-09-01T00:00:00Z', '2026-09-01T00:00:00Z'),
('p2', 'Two', '/two', '[]', '2026-09-01T00:00:00Z', '2026-09-01T00:00:00Z')`;
yield* sql`INSERT INTO projection_threads (thread_id, project_id, title, model_selection_json, runtime_mode, interaction_mode, created_at, updated_at, archived_at, deleted_at, settled_override, settled_at)
VALUES
('t-late', 'p1', 'Late', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', '2026-09-03T00:00:00Z', '2026-09-03T00:00:00Z', NULL, NULL, 'settled', '2026-09-04T00:00:00Z'),
('t-early', 'p2', 'Early', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', '2026-09-01T00:00:00Z', '2026-09-01T00:00:00Z', NULL, NULL, NULL, NULL),
('t-first', 'p1', 'First', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', '2026-09-02T00:00:00Z', '2026-09-02T00:00:00Z', NULL, NULL, NULL, NULL),
('t-plain', 'p1', 'Plain', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', '2026-09-02T00:00:00Z', '2026-09-02T00:00:00Z', NULL, NULL, NULL, NULL),
('t-archived', 'p1', 'Archived', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', '2026-09-02T00:00:00Z', '2026-09-02T00:00:00Z', '2026-09-05T00:00:00Z', NULL, NULL, NULL),
('t-deleted', 'p1', 'Deleted', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', '2026-09-02T00:00:00Z', '2026-09-02T00:00:00Z', NULL, '2026-09-05T00:00:00Z', NULL, NULL)`;
yield* sql`INSERT INTO projection_thread_pull_requests (thread_id, host, repository, number, url, source, linked_at, snapshot_json)
VALUES
('t-late', 'github.com', 'acme/web', 3, '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/acme/web/pull/3', 'manual', '2026-09-03T00:00:00Z', NULL),
('t-early', 'github.com', 'acme/api', 4, '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/acme/api/pull/4', 'agent', '2026-09-01T00:00:00Z', NULL),
('t-first', 'github.com', 'acme/web', 2, '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/acme/web/pull/2', 'stack-dismissed', '2026-09-02T00:00:00Z', NULL),
('t-first', 'github.com', 'acme/web', 1, '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/acme/web/pull/1', 'created', '2026-09-02T00:00:00Z',
'{"state":"open","title":"One","headBranch":"one","baseBranch":"main","isDraft":false,"updatedAt":null,"syncedAt":"2026-09-02T00:00:00Z"}'),
('t-archived', 'github.com', 'acme/web', 5, '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/acme/web/pull/5', 'manual', '2026-09-02T00:00:00Z', NULL),
('t-deleted', 'github.com', 'acme/web', 6, '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/acme/web/pull/6', 'manual', '2026-09-02T00:00:00Z', NULL)`;
const expected = (yield* query.getShellSnapshot()).threads
.filter((thread) => thread.pullRequests.length > 0)
.map(({ id, projectId, settledOverride, settledAt, pullRequests }) => ({
id,
projectId,
settledOverride,
settledAt,
pullRequests,
}));
resolved.length = 0;

const counter = makeSqlStatementCounter();
const threads = yield* query
.listThreadsWithPullRequests()
.pipe(Effect.withTracer(counter.tracer));
assert.deepStrictEqual(
threads.map((thread) => [thread.id, thread.pullRequests.map((link) => link.number)]),
[
["t-first", [1, 2]],
["t-late", [3]],
["t-early", [4]],
],
);
assert.deepStrictEqual(threads, expected);
assert.strictEqual(counter.count(), 1);
assert.deepStrictEqual(resolved, []);
}).pipe(Effect.provide(layer));
},
);

projectionSnapshotLayer("ProjectionSnapshotQuery activities by kind", (it) => {
it.effect("lists one kind across active threads only, without hydrating the threads", () =>
Effect.gen(function* () {
Expand Down
66 changes: 66 additions & 0 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ import {
type ProjectionSnapshotCounts,
type ProjectionThreadCheckpointContext,
type ProjectionThreadDetailQuery,
type ProjectionThreadPullRequests,
type ProjectionSnapshotQueryShape,
} from "../Services/ProjectionSnapshotQuery.ts";

Expand Down Expand Up @@ -802,6 +803,41 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () {
`,
});

// One row per link, in the shell snapshot's thread order and link order.
const listActiveThreadPullRequestSyncRows = SqlSchema.findAll({
Request: Schema.Void,
Result: ProjectionThreadPullRequestDbRowSchema.mapFields(
Struct.assign({
projectId: ProjectionThread.fields.projectId,
settledOverride: ProjectionThread.fields.settledOverride,
settledAt: ProjectionThread.fields.settledAt,
}),
),
execute: () =>
sql`
SELECT
links.thread_id AS "threadId",
threads.project_id AS "projectId",
threads.settled_override AS "settledOverride",
threads.settled_at AS "settledAt",
links.host,
links.repository,
links.number,
links.url,
links.source,
links.linked_at AS "linkedAt",
links.snapshot_json AS "snapshot",
links.stack_json AS "stack"
FROM projection_thread_pull_requests links
INNER JOIN projection_threads threads
ON threads.thread_id = links.thread_id
WHERE threads.deleted_at IS NULL
AND threads.archived_at IS NULL
ORDER BY threads.project_id ASC, threads.created_at ASC, threads.thread_id ASC,
links.linked_at ASC, links.number ASC
`,
});

const listArchivedThreadPullRequestRows = SqlSchema.findAll({
Request: Schema.Void,
Result: ProjectionThreadPullRequestDbRowSchema,
Expand Down Expand Up @@ -2783,6 +2819,35 @@ pending_approval_requests AS (
}),
);

const listThreadsWithPullRequests: ProjectionSnapshotQueryShape["listThreadsWithPullRequests"] =
() =>
listActiveThreadPullRequestSyncRows(undefined).pipe(
Effect.map((rows) => {
const threads = new Map<
ThreadId,
ProjectionThreadPullRequests & { readonly pullRequests: Array<ThreadPullRequestLink> }
>();
for (const row of rows) {
const thread = threads.get(row.threadId) ?? {
id: row.threadId,
projectId: row.projectId,
settledOverride: row.settledOverride,
settledAt: row.settledAt,
pullRequests: [],
};
thread.pullRequests.push(mapPullRequestRow(row));
threads.set(row.threadId, thread);
}
return [...threads.values()];
}),
Effect.mapError(
toPersistenceSqlOrDecodeError(
"ProjectionSnapshotQuery.listThreadsWithPullRequests:query",
"ProjectionSnapshotQuery.listThreadsWithPullRequests:decodeRows",
),
),
);

const getArchivedShellSnapshot: ProjectionSnapshotQueryShape["getArchivedShellSnapshot"] = () =>
sql
.withTransaction(
Expand Down Expand Up @@ -3778,6 +3843,7 @@ pending_approval_requests AS (
listActivitiesByKind,
getSnapshot,
getShellSnapshot,
listThreadsWithPullRequests,
getArchivedShellSnapshot,
getDeletedWorktreeThreads,
searchThreads,
Expand Down
14 changes: 13 additions & 1 deletion apps/server/src/orchestration/PullRequestSyncReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,7 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:
const snapshots = yield* Ref.make(options.snapshot);
const events = yield* PubSub.unbounded<OrchestrationEvent>();
const snapshotReads = yield* Queue.unbounded<void>();
const shellSnapshotReads = yield* Ref.make(0);
const syncCommands = yield* Ref.make<ReadonlyArray<SyncCommand>>([]);
const linkCommands = yield* Ref.make<ReadonlyArray<LinkCommand>>([]);
const summaryCalls = yield* Ref.make<ReadonlyArray<PullRequestRef>>([]);
Expand Down Expand Up @@ -209,8 +210,16 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:

const dependencies = Layer.mergeAll(
Layer.mock(ProjectionSnapshotQuery)({
listThreadsWithPullRequests: () =>
Queue.offer(snapshotReads, undefined).pipe(
Effect.andThen(Ref.get(snapshots)),
Effect.map((snapshot) => snapshot.threads),
),
getShellSnapshot: () =>
Queue.offer(snapshotReads, undefined).pipe(Effect.andThen(Ref.get(snapshots))),
Ref.update(shellSnapshotReads, (count) => count + 1).pipe(
Effect.andThen(Queue.offer(snapshotReads, undefined)),
Effect.andThen(Ref.get(snapshots)),
),
}),
Layer.mock(PullRequestService)({
summary,
Expand All @@ -235,6 +244,7 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options:
activation,
snapshots,
snapshotReads,
shellSnapshotReads,
syncCommands,
linkCommands,
summaryCalls,
Expand Down Expand Up @@ -511,6 +521,8 @@ describe("PullRequestSyncReactor", () => {
],
);
assert.strictEqual((yield* Ref.get(fixture.stackCalls)).length, 1);
// Reads only linked threads, never the full shell snapshot of every thread.
assert.strictEqual(yield* Ref.get(fixture.shellSnapshotReads), 0);
}).pipe(Effect.provide(fixture.layer));
}),
),
Expand Down
16 changes: 7 additions & 9 deletions apps/server/src/orchestration/PullRequestSyncReactor.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import { siblingPullRequestUrl } from "@t3tools/shared/changeRequestUrl";
import {
CommandId,
type OrchestrationThreadShell,
type PullRequestSummary,
type ThreadPullRequestKey,
type ThreadPullRequestLink,
Expand Down Expand Up @@ -36,7 +35,7 @@ const SLOW_SYNC_INTERVAL_MS = 15 * 60 * 1_000;
type SnapshotFields = Omit<ThreadPullRequestSnapshot, "syncedAt">;

interface LinkEntry {
readonly thread: OrchestrationThreadShell;
readonly thread: ProjectionSnapshotQuery.ProjectionThreadPullRequests;
readonly link: ThreadPullRequestLink;
}

Expand Down Expand Up @@ -104,15 +103,15 @@ function stacksEqual(
);
}

function isUnsettled(thread: OrchestrationThreadShell): boolean {
function isUnsettled(thread: ProjectionSnapshotQuery.ProjectionThreadPullRequests): boolean {
return thread.settledOverride !== "settled" && thread.settledAt === null;
}

/**
* Keeps every thread ↔ pull request link's host snapshot current. One sweep a minute reads
* the shell snapshot, groups visible links by pull request so the host is asked once per PR
* no matter how many threads share it, and writes back only what changed. Native stacks the
* host reports are auto-linked to the thread as `source: "stack"`.
* only the active threads that have links, groups visible links by pull request so the host
* is asked once per PR no matter how many threads share it, and writes back only what
* changed. Native stacks the host reports are auto-linked to the thread as `source: "stack"`.
*/
export class PullRequestSyncReactor extends Context.Service<
PullRequestSyncReactor,
Expand Down Expand Up @@ -153,14 +152,13 @@ export const make = Effect.gen(function* () {
Cause.hasInterruptsOnly(cause) ? Effect.failCause(cause) : Effect.logWarning(message, fields);

const sweep = Effect.fn("PullRequestSyncReactor.sweep")(function* (requestedKey?: string) {
const snapshot = yield* snapshots.getShellSnapshot();
const threads = yield* snapshots.listThreadsWithPullRequests();
const now = yield* DateTime.now;
const nowMs = DateTime.toEpochMillis(now);
const nowIso = DateTime.formatIso(now);

const groups = new Map<string, Array<LinkEntry>>();
for (const thread of snapshot.threads) {
if (thread.archivedAt !== null) continue;
for (const thread of threads) {
for (const link of visibleThreadPullRequests(thread.pullRequests)) {
const key = threadPullRequestKeyOf(link);
const entries = groups.get(key) ?? [];
Expand Down
16 changes: 16 additions & 0 deletions apps/server/src/orchestration/Services/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,12 @@ export interface ProjectionFullThreadDiffContext {
readonly toCheckpointRef: CheckpointRef | null;
}

/** The thread fields pull request sync reads, for a thread with at least one link. */
export type ProjectionThreadPullRequests = Pick<
OrchestrationThreadShell,
"id" | "projectId" | "settledOverride" | "settledAt" | "pullRequests"
>;

export interface ProjectionThreadDetailQuery {
/**
* Limit activities before SQLite returns and decodes their payloads.
Expand Down Expand Up @@ -131,6 +137,16 @@ export interface ProjectionSnapshotQueryShape {
ProjectionRepositoryError
>;

/**
* Read active (not deleted, not archived) threads that have at least one pull
* request link, in shell snapshot order. Skips repository identity, so no
* legacy `linkedPullRequest` is derived.
*/
readonly listThreadsWithPullRequests: () => Effect.Effect<
ReadonlyArray<ProjectionThreadPullRequests>,
ProjectionRepositoryError
>;

/** Durable worktree ownership retained after thread deletion, including across restarts. */
readonly getDeletedWorktreeThreads: () => Effect.Effect<
ReadonlyArray<{
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/project/AgentSessionScanner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ const makeProjectionSnapshotQueryLayer = (importedWorkspaceRoots: ReadonlyArray<
updatedAt: "2026-01-01T00:00:00.000Z",
}),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () => Effect.die("unused"),
getSnapshotSequence: () => Effect.die("unused"),
getCounts: () => Effect.die("unused"),
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/project/ProjectSetupScriptRunner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ const makeProjectionSnapshotQueryLayer = (project: OrchestrationProject) =>
getSnapshot: () => Effect.die("unused"),
getShellSnapshot: () => Effect.die("unused"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () => Effect.die("unused"),
getSnapshotSequence: () => Effect.succeed({ snapshotSequence: 1 }),
getCounts: () => Effect.die("unused"),
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4969,6 +4969,7 @@ describe("agent browser access", () => {
getSnapshot: () => Effect.die("unused"),
getShellSnapshot: () => Effect.die("unused"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () => Effect.die("unused"),
getSnapshotSequence: () => Effect.die("unused"),
getCounts: () => Effect.die("unused"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,7 @@ describe("ProviderSessionReaper", () => {
getSnapshot: () => Effect.die("unused"),
getShellSnapshot: () => Effect.die("unused"),
getDeletedWorktreeThreads: () => Effect.die("unused"),
listThreadsWithPullRequests: () => Effect.die("unused"),
getArchivedShellSnapshot: () => Effect.die("unused"),
getSnapshotSequence: () =>
Effect.succeed({ snapshotSequence: input.readModel.snapshotSequence }),
Expand Down
Loading
Loading