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 CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,11 @@ All notable changes to Agrippa are documented here. The format follows

## [Unreleased]

### Added

- **A follow-up can publish what it was steered to produce** (ADR-0019). When the base flow delivers to a branch, a follow-up that changed the workspace pauses at its own approval presenting the full cumulative patch, then pushes exactly one deterministic commit on top of what the chain last published (an expected-tip CAS — `runs.published_sha` is the chain's record) and re-targets the same PR. A steer that changed nothing publishes nothing; a byte-identical patch carries its earlier approval forward instead of re-asking; and a branch that no longer matches the chain's record — a human push, a post-merge deletion — fails typed as `publish_conflict` with the remote untouched.
- **Central runs route to the host that holds their workspace** (the per-host queue, ADR-0018 amendment). Hosts are identified by their storage (`WORKSPACE_ROOT/.agrippa-host-id`); repo checkouts stamp `runs.workspace_host` first-writer-wins; follow-ups inherit it; producers send host-pinned jobs to `run.host.<id>`, which only workers mounting that storage poll — follow-ups and resumes land where the directory lives instead of bouncing off claim-time declines, which remain as the deploy-skew fallback. A dead pin fails typed rather than circulating: runs parked on a host no live worker has advertised for five minutes fail `workspace_lost` with a notification (the `dead-host-runs` sweeper stage).

### Fixed

- A Codex resume whose session is gone (collected or migrated session home — the CLI dies before `thread.started` with "no rollout found") now reports the resume **rejected**, so the engine runs its context-loss disclosure and continues fresh, instead of failing the step `model_error` and re-offering the dead session on every retry.
Expand Down
4 changes: 2 additions & 2 deletions apps/api/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import path from "node:path";
import { createDb, migrateDb, seed } from "@agrippa/db";
import {
createRunQueue,
dbRunExecutorResolver,
dbRunQueueResolver,
RedisEventBus,
seedBuiltinTemplates,
} from "@agrippa/orchestration";
Expand All @@ -20,7 +20,7 @@ if (process.env.AGRIPPA_MIGRATE_ON_BOOT !== "0") {
}

const queue = await createRunQueue(process.env.DATABASE_URL as string, {
resolveRunExecutors: dbRunExecutorResolver(db),
resolveRunQueue: dbRunQueueResolver(db),
});
const bus = process.env.REDIS_URL ? new RedisEventBus(process.env.REDIS_URL) : null;
if (!bus) console.warn("[api] REDIS_URL not set — SSE falls back to DB polling");
Expand Down
59 changes: 56 additions & 3 deletions apps/worker/src/consumer.integration.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { afterAll, describe, expect, it } from "bun:test";
import { runExecuteQueueName, runExecuteSubsetQueues } from "@agrippa/core";
import { runExecuteQueueName, runExecuteSubsetQueues, runHostQueueName } from "@agrippa/core";
import { createDb } from "@agrippa/db";
import { type BossQueue, createRunQueue } from "@agrippa/orchestration";
import { sql } from "drizzle-orm";
Expand Down Expand Up @@ -60,9 +60,17 @@ describe.skipIf(!dbUp)("executor-set queue routing + fetch loop", () => {
await queue?.stop();
});

// name-level routes (host queues) win over executor-id derivation — the
// same precedence dbRunQueueResolver gives a stamped workspace_host
const routes = new Map<string, string>();
const setup = async (): Promise<BossQueue> => {
queue ??= await createRunQueue(TEST_DATABASE_URL, {
resolveRunExecutors: async (runId) => resolvers.get(runId) ?? [],
resolveRunQueue: async (runId) => {
const named = routes.get(runId);
if (named) return named;
const ids = resolvers.get(runId);
return ids && ids.length > 0 ? runExecuteQueueName(ids) : null;
},
});
return queue;
};
Expand All @@ -81,7 +89,7 @@ describe.skipIf(!dbUp)("executor-set queue routing + fetch loop", () => {
const runId = Bun.randomUUIDv7(); // no resolver entry → []
// post-M2-flush nothing consumes `run.execute` — parking the job there
// would strand it silently, so the enqueue must throw instead
await expect(q.enqueueRun(runId)).rejects.toThrow(/cannot derive an executor set/);
await expect(q.enqueueRun(runId)).rejects.toThrow(/cannot derive a queue/);
expect(await jobRow(runId)).toBeUndefined();
});

Expand Down Expand Up @@ -130,6 +138,51 @@ describe.skipIf(!dbUp)("executor-set queue routing + fetch loop", () => {
expect(deliveries).toBe(1);
});

it("a host-pinned run lands only on the worker mounting that storage", async () => {
// the two-root fleet, at the fetch level (ADR-0018 amendment): the holder
// polls its own run.host.<id> queue; a peer with the SAME executor set
// does not, so it can never fetch a chain whose directory it lacks
const q = await setup();
const execH = `exec-h-${suite}`;
const hostQueue = runHostQueueName(`host-fleet-${suite}`);
for (const name of [...runExecuteSubsetQueues([execH]), hostQueue]) {
await q.boss.createQueue(name);
}
const seenByHolder: string[] = [];
const seenByPeer: string[] = [];
loops.push(
startRunFetchLoop({
boss: q.boss,
queues: [...runExecuteSubsetQueues([execH]), hostQueue],
slots: 2,
handler: async (job) => void seenByHolder.push(job.data.runId),
logger: noopLogger,
tickMs: 100,
}),
startRunFetchLoop({
boss: q.boss,
queues: runExecuteSubsetQueues([execH]),
slots: 2,
handler: async (job) => void seenByPeer.push(job.data.runId),
logger: noopLogger,
tickMs: 100,
}),
);

const pinned = Bun.randomUUIDv7();
routes.set(pinned, hostQueue);
await q.enqueueRun(pinned);
await waitFor(() => seenByHolder.includes(pinned));
expect(seenByPeer).not.toContain(pinned);
await waitFor(async () => (await jobRow(pinned))?.state === "completed");

// an unpinned run of the same executor set stays fair game for either
const unpinned = Bun.randomUUIDv7();
resolvers.set(unpinned, [execH]);
await q.enqueueRun(unpinned);
await waitFor(() => seenByHolder.includes(unpinned) || seenByPeer.includes(unpinned));
});

it("updateQueues extends a live loop to a daemon-covered set queue", async () => {
const q = await setup();
// fresh executor ids so the loops from earlier tests cannot consume these
Expand Down
11 changes: 10 additions & 1 deletion apps/worker/src/consumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,16 @@ export function createRunConsumer(db: Db, deps: EngineDeps, queue: BossQueue): R
err instanceof WorkspaceBusyError
) {
const [run] = await db.select({ status: runs.status }).from(runs).where(eq(runs.id, runId));
if (run && (run.status === "queued" || run.status === "waiting_approval")) {
if (
run &&
(run.status === "queued" ||
run.status === "waiting_approval" ||
// a crashed run's pg-boss retry delivered to the wrong host: the
// pinned-elsewhere decline completes the job, and the LEASE
// sweeper re-enqueues the leaseless running run through the
// resolver — onto the host queue this delivery bypassed
(run.status === "running" && err instanceof WorkspaceElsewhereError))
) {
deps.logger.warn(`run ${runId}: declining — ${String((err as Error).message)}`);
// one timeline event so the SPA can show WHY the run is waiting
await appendRunEvent(db, {
Expand Down
50 changes: 39 additions & 11 deletions apps/worker/src/deps/scm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,13 @@ import {
type ScmService,
workspaceKeyOf,
} from "@agrippa/orchestration";
import { applyApprovedPatch, platformGitDirFor, workspaceIntact } from "@agrippa/workspace";
import { and, desc, eq } from "drizzle-orm";
import {
applyApprovedPatch,
platformGitDirFor,
TipConflictError,
workspaceIntact,
} from "@agrippa/workspace";
import { and, desc, eq, isNotNull } from "drizzle-orm";
import {
credentialedUrl,
git,
Expand Down Expand Up @@ -145,15 +150,38 @@ export class GitScmService implements ScmService {
baseSha = pinnedBase;
}

const { commitSha } = await applyApprovedPatch({
fetchSource,
fetchRef,
baseSha,
branch: spec.branch,
patch,
pushUrl,
});
return { status: "pushed", commitSha };
// The chain's expected tip E (ADR-0019): the latest snapshot any run of
// this workspace chain published. Absent → first publish (creation CAS);
// present → the new commit parents on it and the push advances it.
const expectedTip = await this.chainExpectedTip(key);
try {
const { commitSha } = await applyApprovedPatch({
fetchSource,
fetchRef,
baseSha,
branch: spec.branch,
patch,
pushUrl,
...(expectedTip === undefined ? {} : { expectedTip }),
});
return { status: "pushed", commitSha };
} catch (err) {
if (err instanceof TipConflictError) {
return { status: "tip_conflict", observedTip: err.observedTip };
}
throw err;
}
}

/** E = the latest non-null published_sha across the workspace chain. */
private async chainExpectedTip(workspaceKey: string): Promise<string | undefined> {
const [row] = await this.db
.select({ publishedSha: runs.publishedSha })
.from(runs)
.where(and(eq(runs.workspaceKey, workspaceKey), isNotNull(runs.publishedSha)))
.orderBy(desc(runs.number))
.limit(1);
return row?.publishedSha ?? undefined;
}

/** The base the SERVER pinned at checkout (runs.workspace_ref, remote runs). */
Expand Down
112 changes: 111 additions & 1 deletion apps/worker/src/deps/workspace.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import { mkdtempSync } from "node:fs";
import { appendFile, chmod, rm, symlink } from "node:fs/promises";
import { tmpdir } from "node:os";
import path from "node:path";
import { runExecuteQueueName, runHostQueueName } from "@agrippa/core";
import {
createDb,
encryptSecret,
Expand All @@ -20,7 +21,7 @@ import {
taskTypes,
users,
} from "@agrippa/db";
import { seedBuiltinTemplates } from "@agrippa/orchestration";
import { dbRunQueueResolver, seedBuiltinTemplates } from "@agrippa/orchestration";
import { eq, sql } from "drizzle-orm";

// WORKSPACE_ROOT is read at module load — point it at a scratch dir BEFORE
Expand Down Expand Up @@ -252,6 +253,39 @@ describe.skipIf(!dbUp)("GitWorkspaceManager + GitScmService (real git)", () => {
expect(await hostOf(pinnedRunId)).toBeNull();
});

it("dbRunQueueResolver: a host pin routes to the host queue; a runtime pin never does", async () => {
const resolve = dbRunQueueResolver(db);

const hostPinned = await newRunRow();
await db.update(runs).set({ workspaceHost: "host-q" }).where(eq(runs.id, hostPinned));
expect(await resolve(hostPinned)).toBe(runHostQueueName("host-q"));

const plain = await newRunRow();
expect(await resolve(plain)).toBe(runExecuteQueueName(["fake"]));

// a daemon-pinned run's affinity IS its runtime pin — it routes by
// executor set even if a workspace host somehow got stamped
const [runtime] = await db
.insert(runtimes)
.values({
orgId,
name: "route-guard",
tokenHash: "h",
tokenPrefix: `agrd_${Bun.randomUUIDv7().slice(-7)}`,
executors: [],
createdBy: userId,
})
.returning({ id: runtimes.id });
const daemonPinned = await newRunRow();
await db
.update(runs)
.set({ runtimeId: runtime?.id as string, workspaceHost: "host-q" })
.where(eq(runs.id, daemonPinned));
expect(await resolve(daemonPinned)).toBe(runExecuteQueueName(["fake"]));

expect(await resolve(Bun.randomUUIDv7())).toBeNull();
});

it("records the clone base and keeps sanitization out of the diff", async () => {
const dir = workspaceDirFor(workspaceKey);
const base = await platformBaseSha(workspaceKey);
Expand Down Expand Up @@ -347,6 +381,82 @@ describe.skipIf(!dbUp)("GitWorkspaceManager + GitScmService (real git)", () => {
expect(await workspace.diff(runId)).toContain("committed line");
});

it("a follow-up publish advances the chain tip; a human advance conflicts typed (ADR-0019)", async () => {
const branch = publishBranch;
const dir = workspaceDirFor(workspaceKey);
// the ancestor's publication record, as the engine's git.push handler writes it
const tip1 = (await git(["rev-parse", branch], sourceDir)).trim();
await db.update(runs).set({ publishedSha: tip1 }).where(eq(runs.id, runId));

const followupRow = async (): Promise<string> => {
runNumber += 1;
const [row] = await db
.insert(runs)
.values({
...newRunIdentity(),
workspaceKey, // the CHAIN's directory, inherited like the API does
kind: "followup",
parentRunId: runId,
taskId,
projectId,
number: runNumber,
templateVersionId,
faberId,
executorId: "fake",
paramsSnapshot: {},
modelResolution: {},
createdBy: userId,
})
.returning();
return row?.id as string;
};

// the steer adds one more thing; evidence stays cumulative against the base
const followupId = await followupRow();
await Bun.write(path.join(dir, "steer-addition.txt"), "asked for one more thing\n");
const approved = await workspace.diff(followupId);
const advanced = await scm.push(followupId, {
projectId,
repo: { repoConnectionId },
branch,
expectedPatch: approved,
});
if (advanced.status !== "pushed") throw new Error(`expected pushed, got ${advanced.status}`);
// exactly one new commit, parented on what the chain last published
expect((await git(["rev-parse", `${branch}^`], sourceDir)).trim()).toBe(tip1);
expect((await git(["rev-list", "--count", `main..${branch}`], sourceDir)).trim()).toBe("2");
// cumulative: the ancestor's work rides along with the steer's
const show = (spec: string) =>
Bun.spawnSync(["git", "show", spec], { cwd: sourceDir, stdout: "pipe", stderr: "pipe" });
expect(show(`${branch}:left-uncommitted.txt`).exitCode).toBe(0);
expect(show(`${branch}:steer-addition.txt`).exitCode).toBe(0);
await db.update(runs).set({ publishedSha: advanced.commitSha }).where(eq(runs.id, followupId));

// a human advances the branch; the next steer's publish must refuse typed
const human = (
await gitIn(sourceDir, [
"commit-tree",
`${advanced.commitSha}^{tree}`,
"-p",
advanced.commitSha,
"-m",
"human touch-up",
])
).trim();
await git(["update-ref", `refs/heads/${branch}`, human], sourceDir);

const conflictedId = await followupRow();
await Bun.write(path.join(dir, "another-steer.txt"), "and one more\n");
const conflicted = await scm.push(conflictedId, {
projectId,
repo: { repoConnectionId },
branch,
expectedPatch: await workspace.diff(conflictedId),
});
expect(conflicted).toEqual({ status: "tip_conflict", observedTip: human });
expect((await git(["rev-parse", branch], sourceDir)).trim()).toBe(human);
});

it("refuses to publish a run with no commits and no changes", async () => {
const emptyRunId = await newRunRow();
await workspace.checkout(emptyRunId, {
Expand Down
4 changes: 2 additions & 2 deletions apps/worker/src/heterogeneous-fleet.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ import {
type BossQueue,
buildParamsValidator,
createRunQueue,
dbRunExecutorResolver,
dbRunQueueResolver,
type EngineDeps,
FakeResourceMaterializer,
FakeWorkspaceManager,
Expand Down Expand Up @@ -179,7 +179,7 @@ describe.skipIf(!dbUp)("heterogeneous fleet routing (Phase A verify)", () => {
defaultFaberId = taskType?.defaultFaberId as string;

queue = await createRunQueue(TEST_DATABASE_URL, {
resolveRunExecutors: dbRunExecutorResolver(db),
resolveRunQueue: dbRunQueueResolver(db),
});

const workerA = createRunConsumer(db, depsFor({ fake: fakeOnA }), queue);
Expand Down
Loading
Loading