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
139 changes: 139 additions & 0 deletions .github/scripts/work-queue-stress.test.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,145 @@ const fs = require("node:fs/promises");
const path = require("node:path");
const { test } = require("node:test");
const { LocalGitHub, parseArgs, runSimulator } = require("./work-queue-stress.cjs");
const { canonical } = require("../../actions/setup/js/work_queue_codec.cjs");
const { newWork } = require("../../actions/setup/js/work_queue_graph.cjs");
const { actorFromContext, defaultPolicy } = require("../../actions/setup/js/work_queue_policy.cjs");
const { newRequest, prepareCheckpoint, replayTransactions, serializeProjection, serializeTransactionLog } = require("../../actions/setup/js/work_queue_replay.cjs");
const { authenticatePublisher } = require("../../actions/setup/js/work_queue_native.cjs");
const { compactWorkQueue, initializeWorkQueue, publishWorkQueueRequest, readWorkQueueLog } = require("../../actions/setup/js/work_queue_store.cjs");
const { PRINCIPAL, PUBLISHER_WORKFLOW, REPOSITORY } = require("./work-queue-stress-git.cjs");

test("checkpoint loses real-Git publication race, refreshes, and preserves cold replay", { timeout: 60000 }, async () => {
const directory = await fs.mkdtemp(path.join(process.cwd(), ".work-queue-stress-checkpoint-"));
const host = await new LocalGitHub(path.join(directory, "queue.git")).initialize();
const branch = "checkpoint-race";
const args = { githubClient: host.client, owner: "local", repo: "queue", branch };
const timestamp = 1800000000000;
try {
const context = async role =>
authenticatePublisher({
githubClient: host.client,
context: host.nativeContext(),
role,
workflowRef: `${REPOSITORY}/${PUBLISHER_WORKFLOW}@${host.ref}`,
});
const admin = await context("administrator");
const producer = await context("producer");
const dispatcher = await context("dispatcher");
const policy = defaultPolicy({ repository: REPOSITORY, principal: PRINCIPAL, ref: host.ref });
await initializeWorkQueue({ ...args, context: admin, policyProposal: policy, now: () => timestamp });
const submit = async (name, at) =>
publishWorkQueueRequest({
...args,
context: producer,
now: () => at,
request: newRequest(name, "submit", actorFromContext(producer), {
nodes: [newWork({ item: name, effect_contract: { kind: "none" } }, "checkpoint-graph", name, "default", policy, at)],
}),
});
await submit("first", timestamp + 1);
const stale = await readWorkQueueLog(args);
const staleCheckpoint = prepareCheckpoint(stale.state, stale.sha, actorFromContext(admin), timestamp + 2);
await submit("racing", timestamp + 2);
const grant = await publishWorkQueueRequest({
...args,
context: dispatcher,
now: () => timestamp + 2,
request: newRequest("checkpoint-grant", "dispatch_next", actorFromContext(dispatcher), {
pool: "default",
max_claims: 1,
max_dispatches: 1,
max_bytes: 48 * 1024,
}),
});
assert.equal(grant.assignments.length, 1);
const publish = async (base, transactions) => {
const content = serializeTransactionLog(transactions);
const blob = await host.client.rest.git.createBlob({ owner: "local", repo: "queue", content, encoding: "utf-8" });
const tree = await host.client.rest.git.createTree({
owner: "local",
repo: "queue",
base_tree: base.treeSha,
tree: [{ path: "work-queue.jsonl", mode: "100644", type: "blob", sha: blob.data.sha }],
});
const candidate = await host.client.rest.git.createCommit({
owner: "local",
repo: "queue",
tree: tree.data.sha,
parents: [base.sha],
message: "checkpoint candidate",
});
return host.client.rest.git.updateRef({ owner: "local", repo: "queue", ref: `heads/${branch}`, sha: candidate.data.sha, force: false });
};
await assert.rejects(publish(stale, staleCheckpoint), error => error.status === 422 && /not a fast-forward/.test(error.message));
const fresh = await readWorkQueueLog(args);
assert.equal(fresh.state.works.size, 2, "failed checkpoint must not erase competing submission");
assert.equal(fresh.state.claims.size, 1, "checkpoint must preserve an open Claim");
const published = await compactWorkQueue({ ...args, context: admin, now: () => timestamp + 3 });
assert.equal(published.publishedNow, true);
const checkpoint = published.transactions;
const restored = await readWorkQueueLog(args);
assert.equal(restored.transactions.length, 1);
assert.equal(restored.transactions[0].operations[0].kind, "Checkpoint");
const snapshot = restored.transactions[0].operations[0];
assert.equal(Object.hasOwn(snapshot, "history"), false, "checkpoint must store minimal state, not a compressed historical log");
assert.equal(Object.hasOwn(snapshot.state, "transactions"), false, "checkpoint snapshot must not embed the transaction log");
assert.equal(restored.state.works.size, 2);
assert.equal(restored.state.claims.size, 1);
assert.equal(restored.state.dispatches.size, 1);
assert.equal(canonical(snapshot.state.clocks), canonical(serializeProjection(fresh.state).clocks), "checkpoint must retain nonzero scheduler debt");
assert.ok(
Object.values(snapshot.state.clocks.default.classes.pass).some(value => BigInt(value) > 0n),
"stress fixture must actually charge the scheduling clock"
);
for (const [id, work] of fresh.state.works) assert.equal(canonical(restored.state.works.get(id)), canonical(work));
assert.equal(restored.state.requests.size, fresh.state.requests.size + 1);
assert.equal(host.metrics.cas_conflicts, 1);
assert.throws(() => replayTransactions([{ ...checkpoint[0], operations: [{ ...snapshot, state_sha256: "0".repeat(64) }] }]), /checkpoint_invalid|checkpoint state/);
await submit("after-checkpoint", timestamp + 4);
const continued = await readWorkQueueLog(args);
assert.equal(continued.transactions.length, 2);
assert.equal(continued.state.works.size, 3, "post-checkpoint append must preserve snapshot Work");
assert.equal(continued.transactions[1].previous, checkpoint[0].id);
host.armFaults(branch);
const outcomes = await Promise.allSettled([0, 1].map(() => compactWorkQueue({ ...args, context: admin, now: () => timestamp + 5 })));
assert.equal(outcomes.filter(result => result.status === "fulfilled").length, 2, "both checkpoint publishers must recover after one wins a ref race");
assert.ok(
outcomes.every(result => result.status === "fulfilled" && result.value.publishedNow === false),
"neither ambiguous response is a fresh publication"
);
assert.equal((await readWorkQueueLog(args)).state.works.size, 3, "competing checkpoint must not lose snapshot Work");
assert.equal(host.metrics.forced_races, 1);
assert.equal(host.metrics.lost_success_responses, 1);
assert.ok(host.metrics.cas_conflicts >= 2);
await submit("intervening-baseline", timestamp + 6);
const updateRef = host.client.rest.git.updateRef;
let injected = false;
host.client.rest.git.updateRef = async request => {
if (!injected && request.ref === `heads/${branch}`) {
injected = true;
await submit("intervening", timestamp + 7);
await compactWorkQueue({ ...args, context: admin, now: () => timestamp + 8 });
}
return updateRef(request);
};
try {
const recovered = await compactWorkQueue({ ...args, context: admin, now: () => timestamp + 9 });
assert.equal(injected, true, "the checkpoint writer must encounter a competing mutation and checkpoint");
assert.equal(recovered.publishedNow, false);
const current = await readWorkQueueLog(args);
assert.equal(current.state.works.size, 5, "checkpoint retry must preserve Work published between source read and competing checkpoint");
assert.equal(current.transactions.length, 1);
assert.equal(current.transactions[0].operations[0].kind, "Checkpoint");
} finally {
host.client.rest.git.updateRef = updateRef;
}
await host.git(["fsck", "--strict", "--no-dangling"]);
} finally {
await host.close();
await fs.rm(directory, { recursive: true, force: true });
}
});

test("real-Git adapter refuses stale non-fast-forward candidates and foreign repositories", { timeout: 30000 }, async () => {
const directory = await fs.mkdtemp(path.join(process.cwd(), ".work-queue-stress-adapter-"));
Expand Down
37 changes: 37 additions & 0 deletions .github/workflows/agentics-maintenance.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1984,3 +1984,40 @@ jobs:
- name: Report ledger compaction application failures
if: ${{ steps.apply_0.outcome == 'failure' || steps.apply_1.outcome == 'failure' || steps.apply_2.outcome == 'failure' || steps.apply_3.outcome == 'failure' || steps.apply_4.outcome == 'failure' || steps.apply_5.outcome == 'failure' }}
run: exit 1

work_queue_compaction:
name: Compact work queue
if: ${{ (!(github.event.repository.fork)) && github.event_name != 'push' && (github.event_name != 'workflow_dispatch' && github.event_name != 'workflow_call' || inputs.operation == '' || inputs.operation == 'none') }}
runs-on: ubuntu-slim
permissions:
contents: write
steps:
- name: Checkout actions folder
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
with:
sparse-checkout: |
actions
clean: false
persist-credentials: false

- name: Setup Scripts
uses: ./actions/setup
with:
destination: ${{ runner.temp }}/gh-aw/actions

- name: Plan and apply work queue compaction
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
env:
GH_AW_WORK_QUEUE_COMPACTION_PLAN_FILE: ${{ runner.temp }}/gh-aw/work-queue-compaction/plan.json
with:
script: |
const path = require('path');
const actionsDir = path.join(process.env.RUNNER_TEMP, 'gh-aw', 'actions');
const { setupGlobals } = require(path.join(actionsDir, 'setup_globals.cjs'));
setupGlobals(core, github, context, exec, io, getOctokit);
const { main: plan } = require(path.join(actionsDir, 'work_queue_compaction_plan.cjs'));
const result = await plan();
if (result.status === 'planned') {
const { main: apply } = require(path.join(actionsDir, 'work_queue_compaction_apply.cjs'));
await apply();
}
9 changes: 8 additions & 1 deletion .github/workflows/work-queue-stress.yml
Original file line number Diff line number Diff line change
Expand Up @@ -51,9 +51,16 @@ jobs:
id: invariants
run: |
mkdir -p stress-results
node --test .github/scripts/work-queue-stress.test.cjs .github/scripts/work-queue-stress-summary.test.cjs \
node --test --test-skip-pattern='checkpoint loses real-Git' \
.github/scripts/work-queue-stress.test.cjs .github/scripts/work-queue-stress-summary.test.cjs \
| tee stress-results/invariants.tap

- name: Exercise checkpoint publication races
id: checkpoint
run: |
node --test --test-name-pattern='checkpoint loses real-Git' .github/scripts/work-queue-stress.test.cjs \
| tee stress-results/checkpoint.tap

- name: Replay retained Work and Claim items
id: history
run: |
Expand Down
5 changes: 4 additions & 1 deletion actions/setup/js/work_queue_checked_transport.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@

const { canonical, parseStrictJSON, queueError } = require("./work_queue_codec.cjs");
const { replayTransactionLog, replayTransactions } = require("./work_queue_replay.cjs");
const { validateBranch } = require("./work_queue_store.cjs");
const { validateBranch, verifyCheckpointGitHistory } = require("./work_queue_store.cjs");

const MAX_BYTES = 80 * 1024 * 1024;
const JOURNAL_PATH = /^\.gh-aw\/issue-projection\/[a-f0-9]{64}\.json$/;
Expand Down Expand Up @@ -72,6 +72,9 @@ async function readCheckedQueue({ githubClient, owner, repo, branch = "work-queu
const text = await checkedBlobText(githubClient, owner, repo, repository.log, entry.oid, MAX_BYTES);
const state = replayTransactionLog(text);
if (state.transactions.some(transaction => transaction.actor.repository.toLowerCase() !== repository.nameWithOwner.toLowerCase())) throw queueError("actor_unauthorized", "ledger contains a foreign repository");
if (state.transactions[0]?.operations[0]?.kind === "Checkpoint") {
await verifyCheckpointGitHistory({ githubClient, owner, repo, branch, repository: { full_name: repository.nameWithOwner } }, commit.oid, state.transactions[0]);
}
const journal = new Map();
const directory = entries.find(entry => entry.name === ".gh-aw")?.object?.entries?.find(entry => entry.name === "issue-projection")?.object?.entries || [];
for (const [index, path] of paths.entries()) {
Expand Down
Loading
Loading