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
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import * as ThreadPlanProgress from "../ThreadPlanProgress.ts";
import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts";
import { encodeThreadDetailPageCursor } from "../threadDetailCursor.ts";
import { projectThreadDetailSnapshot } from "../ActivityPayloadProjection.ts";
import { readSweepSnapshot } from "../ThreadPullRequestReactor.ts";
import { makeSqlStatementCounter } from "../../../integration/SqlStatementCounter.integration.ts";

const asProjectId = (value: string): ProjectId => ProjectId.make(value);
Expand Down Expand Up @@ -3574,6 +3575,74 @@ it.effect(
},
);

it.effect("reads one sweep thread and its projects like the shell snapshot", () => {
const layer = OrchestrationProjectionSnapshotQueryLive.pipe(
Layer.provide(ThreadBackgroundLiveness.layer),
Layer.provide(ThreadPlanProgress.layer),
Layer.provide(
Layer.succeed(RepositoryIdentityResolver.RepositoryIdentityResolver, {
resolve: () =>
Effect.succeed({
canonicalKey: "github.com/acme/web",
provider: "github",
displayName: "acme/web",
locator: {
source: "git-remote" as const,
remoteName: "origin",
remoteUrl: "https://github.com/acme/web.git",
},
}),
}),
),
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-02T00:00:00Z', '2026-09-02T00:00:00Z'),
('p3', 'Three', '/three', '[]', '2026-09-03T00:00:00Z', '2026-09-03T00:00:00Z')`;
yield* sql`INSERT INTO projection_threads (thread_id, project_id, title, model_selection_json, runtime_mode, interaction_mode, branch, worktree_path, branch_pull_request_json, latest_turn_id, latest_user_message_at, pending_approval_count, snoozed_until, snoozed_at, created_at, updated_at, settled_override, settled_at)
VALUES
('t-linked', 'p1', 'Linked', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', 'feature', '/one/wt', NULL, 'turn-1', '2026-09-02T00:00:00Z', 1, NULL, NULL, '2026-09-01T00:00:00Z', '2026-09-02T00:00:00Z', NULL, NULL),
('t-branch', 'p1', 'Branch', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', 'other', NULL,
'{"projectId":"p2","repository":"acme/web","number":8,"url":"https://github.com/acme/web/pull/8"}',
NULL, NULL, 0, '2026-09-10T00:00:00Z', '2026-09-02T00:00:00Z', '2026-09-01T00:00:00Z', '2026-09-02T00:00:00Z', 'settled', '2026-09-03T00:00:00Z'),
('t-other', 'p3', 'Other', '{"provider":"codex","model":"gpt-5"}', 'full-access', 'default', NULL, NULL, NULL, NULL, NULL, 0, NULL, NULL, '2026-09-01T00:00:00Z', '2026-09-01T00:00:00Z', NULL, NULL)`;
yield* sql`INSERT INTO projection_thread_pull_requests (thread_id, host, repository, number, url, source, linked_at)
VALUES ('t-linked', 'github.com', 'acme/web', 7, 'https://github.com/acme/web/pull/7', 'agent', '2026-09-02T00:00:00Z')`;
yield* sql`INSERT INTO projection_turns (thread_id, turn_id, state, requested_at, started_at, completed_at, checkpoint_files_json)
VALUES ('t-linked', 'turn-1', 'completed', '2026-09-02T00:00:00Z', '2026-09-02T00:00:01Z', '2026-09-02T00:00:02Z', '[]')`;
yield* sql`INSERT INTO projection_thread_sessions (thread_id, status, provider_name, active_turn_id, last_error, updated_at)
VALUES ('t-linked', 'ready', 'codex', NULL, NULL, '2026-09-02T00:00:03Z')`;
for (const projector of Object.values(ORCHESTRATION_PROJECTOR_NAMES)) {
yield* sql`INSERT INTO projection_state (projector, last_applied_sequence, updated_at)
VALUES (${projector}, 9, '2026-09-02T00:00:03Z')`;
}

const full = yield* query.getShellSnapshot();
// The seeded fields must reach the snapshot, or the parity check is empty.
const linked = full.threads.find((thread) => thread.id === ThreadId.make("t-linked"));
assert.strictEqual(full.snapshotSequence, 9);
assert.strictEqual(linked?.linkedPullRequest?.number, 7);
assert.strictEqual(linked?.latestTurn?.turnId, asTurnId("turn-1"));
assert.strictEqual(linked?.session?.status, "ready");

for (const [threadId, projectIds] of [
[ThreadId.make("t-linked"), [asProjectId("p1")]],
// Settlement also needs the project that the saved branch PR names.
[ThreadId.make("t-branch"), [asProjectId("p1"), asProjectId("p2")]],
] as const) {
assert.deepStrictEqual(yield* readSweepSnapshot(query, threadId), {
snapshotSequence: full.snapshotSequence,
projects: full.projects.filter((project) => projectIds.includes(project.id)),
threads: full.threads.filter((thread) => thread.id === threadId),
});
}
}).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
43 changes: 39 additions & 4 deletions apps/server/src/orchestration/ThreadPullRequestReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as PubSub from "effect/PubSub";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
Expand Down Expand Up @@ -142,7 +143,8 @@ const makeHarness = Effect.fn("makeThreadPullRequestHarness")(function* (options
threads: options.threads,
updatedAt: NOW,
});
const reads = yield* Queue.unbounded<void>();
// Each shell read: a thread id for a one-thread read, null for a full read.
const reads = yield* Queue.unbounded<ThreadId | null>();
const events = yield* PubSub.unbounded<OrchestrationEvent>();
const commands = yield* Ref.make<ReadonlyArray<SyncCommand>>([]);
const branchCalls = yield* Ref.make<
Expand All @@ -152,8 +154,24 @@ const makeHarness = Effect.fn("makeThreadPullRequestHarness")(function* (options
let uuid = 0;
const dependencies = Layer.mergeAll(
Layer.mock(ProjectionSnapshotQuery)({
getShellSnapshot: () =>
Ref.get(snapshots).pipe(Effect.tap(() => Queue.offer(reads, undefined))),
getShellSnapshot: () => Ref.get(snapshots).pipe(Effect.tap(() => Queue.offer(reads, null))),
getSnapshotSequence: () =>
Ref.get(snapshots).pipe(Effect.map(({ snapshotSequence }) => ({ snapshotSequence }))),
getThreadShellById: (threadId) =>
Ref.get(snapshots).pipe(
Effect.map(({ threads }) =>
Option.fromUndefinedOr(
threads.find((thread) => thread.id === threadId && thread.archivedAt === null),
),
),
Effect.tap(() => Queue.offer(reads, threadId)),
),
getProjectShells: (projectIds) =>
Ref.get(snapshots).pipe(
Effect.map(({ projects }) =>
projects.filter((project) => projectIds?.includes(project.id) ?? true),
),
),
}),
Layer.mock(GitManager)({
branchPullRequest: (input, readOptions) =>
Expand Down Expand Up @@ -376,7 +394,7 @@ describe("ThreadPullRequestReactor", () => {
: [checkpointEvent, sessionEvent];
for (const event of events) {
yield* fixture.publish(event);
yield* Queue.take(fixture.reads);
expect(yield* Queue.take(fixture.reads)).toBe(current.id);
yield* reactor.drain;
}
expect((yield* Ref.get(fixture.commands))[0]?.branchPullRequest).toEqual(reference(42));
Expand Down Expand Up @@ -516,6 +534,23 @@ describe("ThreadPullRequestReactor", () => {
yield* Effect.gen(function* () {
const reactor = yield* fixture.start();
expect(yield* Ref.get(fixture.commands)).toHaveLength(0);
// A one-thread read cannot show that other pending threads are gone.
const gone = ThreadId.make("gone");
yield* fixture.publish({
type: "thread.unarchived",
sequence: 2,
eventId: EventId.make("gone-unarchived"),
aggregateKind: "thread",
aggregateId: gone,
occurredAt: NOW,
commandId: null,
causationEventId: null,
correlationId: null,
metadata: {},
payload: { threadId: gone, updatedAt: NOW },
});
expect(yield* Queue.take(fixture.reads)).toBe(gone);
yield* reactor.drain;
yield* Ref.set(online, true);
yield* TestClock.adjust("1 minute");
yield* Queue.take(fixture.reads);
Expand Down
39 changes: 36 additions & 3 deletions apps/server/src/orchestration/ThreadPullRequestReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import {
CommandId,
type OrchestrationEvent,
type OrchestrationProjectShell,
type OrchestrationShellSnapshot,
type ThreadId,
type ThreadLinkedPullRequest,
} from "@t3tools/contracts";
Expand All @@ -16,11 +17,13 @@ import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Schedule from "effect/Schedule";
import type * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";

import * as GitManager from "../git/GitManager.ts";
import type { ProjectionRepositoryError } from "../persistence/Errors.ts";
import * as PullRequestService from "../pullRequest/PullRequestService.ts";
import * as RepositoryIdentityResolver from "../project/RepositoryIdentityResolver.ts";
import { forkParked } from "../serverActivation.ts";
Expand Down Expand Up @@ -69,6 +72,35 @@ export function pullRequestMatchesProject(
);
}

/**
* Read the shell state for a discovery or settlement sweep. A sweep for one
* thread reads that thread and the projects it names, not every thread.
*/
export const readSweepSnapshot = (
snapshots: ProjectionSnapshotQuery.ProjectionSnapshotQueryShape,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This production Effect helper accepts ProjectionSnapshotQueryShape as an explicit service dependency. Could it acquire ProjectionSnapshotQuery.ProjectionSnapshotQuery with yield* instead, and have both reactors and the parity test provide the query through their layers? That keeps service dependencies in the Effect environment rather than passing service instances between production modules.

Posted via Macroscope — Effect Service Conventions

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note

🤖 Claude Opus 5.5 responding on behalf of Theo

Keeping the argument. readSweepSnapshot is a read helper, not a service factory: both reactors already acquire ProjectionSnapshotQuery with yield* in make (ThreadPullRequestReactor.ts:107, ThreadSettlementReactor.ts:83) and pass that instance down. Reading it from the environment inside the helper would leak the requirement into ThreadSettlementReactor.start, which calls runSweep directly for merge events (ThreadSettlementReactor.ts:378), so it would need an extra Effect.provideService to keep its Effect<void, never, Scope> type.

threadId: ThreadId | null,
): Effect.Effect<
Pick<OrchestrationShellSnapshot, "snapshotSequence" | "projects" | "threads">,
ProjectionRepositoryError
> =>
threadId === null
? snapshots.getShellSnapshot()
: Effect.gen(function* () {
// Read the sequence first. The thread is then at least this new, so a
// command guarded by the sequence is rejected rather than missing a change.
const { snapshotSequence } = yield* snapshots.getSnapshotSequence();
const thread = yield* snapshots.getThreadShellById(threadId);
if (Option.isNone(thread)) return { snapshotSequence, projects: [], threads: [] };
// Settlement also checks the project a saved pull request names.
const reference = thread.value.linkedPullRequest ?? thread.value.branchPullRequest;
const projects = yield* snapshots.getProjectShells(
reference == null
? [thread.value.projectId]
: [thread.value.projectId, reference.projectId],
);
return { snapshotSequence, projects, threads: [thread.value] };
});

/** @public Service construction is part of the canonical Effect module API. */
export const make = Effect.gen(function* () {
const engine = yield* OrchestrationEngine.OrchestrationEngineService;
Expand Down Expand Up @@ -97,7 +129,7 @@ export const make = Effect.gen(function* () {
const synchronize = Effect.fn("ThreadPullRequestReactor.synchronize")(function* (
request: RefreshRequest,
) {
const snapshot = yield* snapshots.getShellSnapshot();
const snapshot = yield* readSweepSnapshot(snapshots, request.threadId);
const projects = new Map(snapshot.projects.map((project) => [project.id, project]));
if (request.backfill) {
for (const thread of snapshot.threads) {
Expand All @@ -109,14 +141,15 @@ export const make = Effect.gen(function* () {
}
}
}
// A single-thread read only shows whether its own thread is gone.
const threadIds = new Set(snapshot.threads.map((thread) => thread.id));
for (const threadId of pendingBackfill.keys()) {
const checkedIds = request.threadId === null ? pendingBackfill.keys() : [request.threadId];
for (const threadId of checkedIds) {
if (!threadIds.has(threadId)) pendingBackfill.delete(threadId);
}
const threads = snapshot.threads.filter(
(thread) =>
thread.archivedAt === null &&
(request.threadId === null || thread.id === request.threadId) &&
((thread.settledOverride !== "settled" && thread.settledAt === null) ||
request.threadId !== null ||
pendingBackfill.has(thread.id)) &&
Expand Down
31 changes: 25 additions & 6 deletions apps/server/src/orchestration/ThreadSettlementReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as PubSub from "effect/PubSub";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
Expand Down Expand Up @@ -181,7 +182,8 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
const activation = yield* Deferred.make<void>();
const snapshots = yield* Ref.make(options.snapshot);
const snapshotReadCount = yield* Ref.make(0);
const snapshotReads = yield* Queue.unbounded<number>();
// Each shell read: a thread id for a one-thread read, null for a full read.
const snapshotReads = yield* Queue.unbounded<ThreadId | null>();
const settings = yield* Ref.make(options.settings ?? DEFAULT_SERVER_SETTINGS);
const settingsReads = yield* Queue.unbounded<ServerSettings>();
const settingsChanges = yield* PubSub.unbounded<ServerSettings>();
Expand Down Expand Up @@ -255,10 +257,27 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
const dependencies = Layer.mergeAll(
Layer.mock(ProjectionSnapshotQuery)({
getShellSnapshot: () =>
Ref.updateAndGet(snapshotReadCount, (count) => count + 1).pipe(
Effect.tap((count) => Queue.offer(snapshotReads, count)),
Ref.update(snapshotReadCount, (count) => count + 1).pipe(
Effect.andThen(Queue.offer(snapshotReads, null)),
Effect.andThen(Ref.get(snapshots)),
),
getSnapshotSequence: () =>
Ref.get(snapshots).pipe(Effect.map(({ snapshotSequence }) => ({ snapshotSequence }))),
getThreadShellById: (threadId) =>
Ref.get(snapshots).pipe(
Effect.map(({ threads }) =>
Option.fromUndefinedOr(
threads.find((thread) => thread.id === threadId && thread.archivedAt === null),
),
),
Effect.tap(() => Queue.offer(snapshotReads, threadId)),
),
getProjectShells: (projectIds) =>
Ref.get(snapshots).pipe(
Effect.map(({ projects }) =>
projects.filter((project) => projectIds?.includes(project.id) ?? true),
),
),
}),
Layer.mock(GitManager)({
branchPullRequest,
Expand Down Expand Up @@ -313,7 +332,7 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options:
const startHarness = Effect.fn("startThreadSettlementHarness")(function* (
reactor: ThreadSettlementReactor.ThreadSettlementReactor["Service"],
activation: Deferred.Deferred<void>,
snapshotReads: Queue.Queue<number>,
snapshotReads: Queue.Queue<ThreadId | null>,
) {
yield* reactor.start();
yield* Deferred.succeed(activation, undefined);
Expand Down Expand Up @@ -443,7 +462,7 @@ describe("ThreadSettlementReactor", () => {
updatedAt: NOW,
},
});
yield* Queue.take(fixture.snapshotReads);
assert.strictEqual(yield* Queue.take(fixture.snapshotReads), thread.id);
yield* reactor.drain;
}
assert.deepStrictEqual(
Expand All @@ -463,7 +482,7 @@ describe("ThreadSettlementReactor", () => {
aggregateId: readySession.threadId,
payload: { threadId: readySession.threadId, session: readySession },
});
yield* Queue.take(fixture.snapshotReads);
assert.strictEqual(yield* Queue.take(fixture.snapshotReads), readySession.threadId);
yield* reactor.drain;
assert.deepStrictEqual(
(yield* Ref.get(fixture.commands)).map(({ threadId }) => threadId),
Expand Down
12 changes: 4 additions & 8 deletions apps/server/src/orchestration/ThreadSettlementReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ import * as ServerSettings from "../serverSettings.ts";
import { forkParked } from "../serverActivation.ts";
import * as OrchestrationEngine from "./Services/OrchestrationEngine.ts";
import * as ProjectionSnapshotQuery from "./Services/ProjectionSnapshotQuery.ts";
import { pullRequestMatchesProject } from "./ThreadPullRequestReactor.ts";
import { pullRequestMatchesProject, readSweepSnapshot } from "./ThreadPullRequestReactor.ts";
import {
isAutoSettlementCandidate,
resolveAutoSettlementAt,
Expand Down Expand Up @@ -95,20 +95,16 @@ export const make = Effect.gen(function* () {
if (!autoSettlementConfigured(settings)) {
return;
}
const snapshot = yield* snapshots.getShellSnapshot();
const snapshot = yield* readSweepSnapshot(snapshots, threadId ?? null);
const now = DateTime.formatIso(yield* DateTime.now);
const projects = new Map(snapshot.projects.map((project) => [project.id, project]));
// A merge rechecks all candidates, including branches that discovery has
// not linked yet. Those lookups can still have cached the PR as open.
const candidates = snapshot.threads.filter(
(thread) =>
(threadId === undefined || thread.id === threadId) &&
isAutoSettlementCandidate(thread, now),
);
const candidates = snapshot.threads.filter((thread) => isAutoSettlementCandidate(thread, now));

// Return the thread when it still needs a pull request decision. A rejected
// dispatch skips it for this snapshot instead of retrying through a lookup.
const settleThread = Effect.fn("ThreadSettlementReactor.settleThread")(
const settleThread = Effect.fnUntraced(
function* (thread: (typeof candidates)[number], pullRequest: SettlementPullRequest | null) {
const settings = resolveProjectSettings(
yield* settingsService.getSettings,
Expand Down
Loading