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
66 changes: 37 additions & 29 deletions .github/aw/work-queue.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,15 @@ description: Agent instructions for Git-backed work queue producers, dispatchers

# Work Queue

Choose the queue pattern first: use native Git storage for fair scheduling and
Claim authority, or [WorkQueueOps](../../docs/src/content/docs/patterns/workqueue-ops.md)
for issue-backed checklists/sub-issues, Discussions or cache-memory backlogs.

Use `tools.work-queue` for durable fair scheduling, immutable Work DAGs and
Claim-scoped effects. Treat the causal `work-queue.jsonl` log as the only
authority. Defaults are FIFO-like; configured weights share Claim opportunities,
not CPU time or successful completions. In this version-3 native protocol,
Issues/PRs are dependency nodes, not queue-storage backends. Git is always used;
do not configure a `storage` field. Do not confuse this with
issue-backed WorkQueueOps, which is a separate pattern.
Choose the queue pattern: Git storage for fair scheduling and Claim authority, or
[WorkQueueOps](../../docs/src/content/docs/patterns/workqueue-ops.md) for
issue-backed WorkQueueOps: lists, sub-issues, Discussions or cache backlogs.

Use `tools.work-queue` for fair scheduling, immutable Work DAGs and Claim-scoped
effects. The causal `work-queue.jsonl` log is authoritative. Weights distribute
Claim opportunities, not CPU or successful completions. Protocol v3 always uses
Git; Issues/PRs are dependencies, not storage. Never set `storage`. WorkQueueOps
is a separate pattern.

## Select the role

Expand All @@ -25,11 +23,10 @@ issue-backed WorkQueueOps, which is a separate pattern.
- **Worker:** require a compiler-supplied version-3 `claims` array. Never
construct/override reserved `work_queue_assignment` or caller context.

Trust the compiler/runtime role channel, not snapshot metadata, agent files or
assignment input alone. Never downgrade a declared worker when input is missing.
Trusted activation must authenticate the actual run/workflow/revision; reruns
cannot inherit attempt 1 authority. Report a genuinely absent queue as
uninitialized; fail explicitly on an existing empty, malformed or unsupported log.
Trust the compiler/runtime role, never snapshot metadata or assignment alone;
don't downgrade declared workers. Activation authenticates run/workflow/revision,
and reruns cannot inherit attempt 1 authority. Report absent queues as
uninitialized; reject empty, malformed or unsupported logs.

## Configure deployment

Expand All @@ -48,11 +45,24 @@ The allowlist is compiler approval, not authority: installed Policy binds each
profile's exact workflow path, immutable SHA, authenticated principal, trust
domain and effect scope. Queue dispatch ignores moving `target-ref`.

Require administrator-installed Policy and independently protected queue-branch
writers; frontmatter provisions neither. Administrator status does not grant
producer entitlement. Never give agent execution or snapshot MCP queue-write
credentials. Read the [deployment guide](../../docs/src/content/docs/guides/deploy-work-queue.md)
only for installation tasks; writer-restriction automation remains deferred.
Require administrator-installed Policy and protected queue-branch writers;
frontmatter provisions neither, and administrator status grants no producer
entitlement. Never give agents or snapshot MCP queue-write credentials. Read the
[deployment guide](../../docs/src/content/docs/guides/deploy-work-queue.md) only
for installation; writer-restriction automation remains deferred.

## Mirror admitted Work with Issues

`tools.work-queue.issues: true` mirrors admitted Work with the `work` label.
An object may set `label` and a pre-provisioned native organization
`status-field`. Require installed projector authority: only protected hooks
project their own admissions/original Claims. Immutable `backing_issue` binds
one Work per Issue. Agents never write mirrors; human edits never establish
Result. Pre-existing Issues need exact installed `backing_issues` grants.
Only trusted admission-time `completion_policy` permits closure, never payload.
Uncertain writes stay pending, never recreated or globally repaired.
Upgrade all closed-schema readers first. See the
[backing Issue reference](../../docs/src/content/docs/reference/work-queue.md#backing-issues).

## Plan and dispatch

Expand Down Expand Up @@ -85,11 +95,9 @@ only for installation tasks; writer-restriction automation remains deferred.
When `<mcp-clis>` advertises the wrapper, use
`work-queue work_queue_read '{}'` or
`work-queue work_queue_claim_finish '{"claim_handle":"h1","outcome":"completed"}'`.
These are MCP subcommands, not `gh aw work-queue` operator commands.
Pass one JSON argument to each wrapper subcommand; do not invoke a wrapper as
a structured tool with `command`/`description`. Check `queue_state` before
interpreting counts: `uninitialized` is a deployment failure, not an empty
backlog. A dispatch response with `status: "staged"` is not a grant or launch.
These are MCP subcommands, not operator commands. Pass one JSON argument, not
structured-tool `command`/`description`. Check `queue_state`: `uninitialized`
means deployment failure, not empty; `status: "staged"` is not a grant or launch.

## Dependencies and recovery

Expand All @@ -108,10 +116,10 @@ automatic upgrades; preserve old evidence before explicit redeployment.

- Operator commands, TUI and diagnostic artifacts:
[queue reference](../../docs/src/content/docs/reference/work-queue.md).
Cancellation is terminal for Work, not a native-worker stop; reconcile separately.
In this checkout, build with `go build -o ./gh-aw ./cmd/gh-aw` and use
Cancellation is terminal for Work, not a worker stop; reconcile separately.
Build with `go build -o ./gh-aw ./cmd/gh-aw`; use
`./gh-aw work-queue --repo OWNER/REPO stats --json`. `dispatch-next` grants
reservations but never launches workers; use the authorized dispatcher workflow.
reservations, never launches; use the authorized dispatcher.
- Daily report rotation, dedicated Policy and Claim examples:
[portfolio walkthrough](../../docs/src/content/docs/patterns/daily-report-portfolio.md)
and [shared worker instructions](../workflows/shared/daily-report-worker.md).
Expand Down
7 changes: 7 additions & 0 deletions actions/setup/js/setup_sh_file_lists.test.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,13 @@ describe("setup.sh SAFE_OUTPUTS_FILES", () => {
expect(safeOutputsFiles).toContain("work_queue_logging.cjs");
});

it("deploys checked transport and Issue binding dependencies with the queue runtime", () => {
expect(getDirectLocalRequires("work_queue_store.cjs")).toContain("work_queue_checked_transport.cjs");
expect(getDirectLocalRequires("work_queue_replay.cjs")).toContain("work_queue_issue_contract.cjs");
expect(safeOutputsFiles).toContain("work_queue_checked_transport.cjs");
expect(safeOutputsFiles).toContain("work_queue_issue_contract.cjs");
});

it("deploys Claim authority and worker route provisioning dependencies", () => {
expect(getDirectLocalRequires("safe_outputs_handlers.cjs")).toContain("work_queue_claim_scope.cjs");
expect(getDirectLocalRequires("work_queue_provisioning.cjs")).toContain("work_queue_yaml.cjs");
Expand Down
121 changes: 121 additions & 0 deletions actions/setup/js/work_queue_checked_transport.cjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
// @ts-check
"use strict";

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 MAX_BYTES = 80 * 1024 * 1024;
const JOURNAL_PATH = /^\.gh-aw\/issue-projection\/[a-f0-9]{64}\.json$/;

async function checkedBlobText(githubClient, owner, repo, blob, oid, maxBytes) {
if (!blob || blob.oid !== oid || !Number.isSafeInteger(blob.byteSize) || blob.byteSize > maxBytes || blob.byteSize < 0) throw queueError("ledger_invalid", "blob does not belong to the immutable checked head or is oversized");
if (blob.isTruncated === false && typeof blob.text === "string" && blob.byteSize === Buffer.byteLength(blob.text, "utf8")) return blob.text;
if (blob.isTruncated !== true) throw queueError("ledger_invalid", "checked blob is unreadable");
const response = await githubClient.rest.git.getBlob({ owner, repo, file_sha: oid, request: { retries: 0 } });
const value = response?.data;
if (value?.sha !== oid || value.encoding !== "base64" || typeof value.content !== "string" || value.size !== blob.byteSize) throw queueError("ledger_invalid", "immutable blob fallback identity or encoding is invalid");
const encoded = value.content.replace(/\s/g, "");
if (encoded.length > Math.ceil(maxBytes / 3) * 4 || encoded.length % 4 !== 0 || /[^A-Za-z0-9+/=]/.test(encoded)) throw queueError("ledger_invalid", "immutable blob fallback is malformed or oversized");
const bytes = Buffer.from(encoded, "base64");
if (bytes.length !== blob.byteSize || bytes.toString("base64") !== encoded) throw queueError("ledger_invalid", "immutable blob fallback is truncated or malformed");
try {
return new TextDecoder("utf-8", { fatal: true, ignoreBOM: true }).decode(bytes);
} catch {
throw queueError("ledger_invalid", "immutable blob fallback contains malformed UTF-8");
}
}

/** @param {import("./work_queue_store.cjs").QueueReadOptions & {paths?: string[], issueRead?: ReturnType<typeof import("./work_queue_issue_api.cjs").issueReadQuery>}} options */
async function readCheckedQueue({ githubClient, owner, repo, branch = "work-queue", paths = [], issueRead = undefined }) {
validateBranch(branch);
if (paths.length > 25 || paths.some(path => !JOURNAL_PATH.test(path))) throw queueError("projection_invalid", "invalid bounded projection journal paths");
const journals = paths.map((path, index) => `j${index}: object(expression:$j${index}) { oid ... on Blob { text byteSize isTruncated } }`).join("\n");
const variables = { owner, repo, ref: `refs/heads/${branch}`, log: `refs/heads/${branch}:work-queue.jsonl` };
paths.forEach((path, index) => {
variables[`j${index}`] = `refs/heads/${branch}:${path}`;
});
const parameters = paths.map((_, index) => `$j${index}:String!`).join(",");
const nativeParameters = issueRead?.declarations?.join(",") || "";
Object.assign(variables, issueRead?.variables || {});
const response = await githubClient.graphql(
`query CheckedWorkQueue($owner:String!,$repo:String!,$ref:String!,$log:String!${parameters ? "," + parameters : ""}${nativeParameters ? "," + nativeParameters : ""}) {
repository(owner:$owner,name:$repo) {
id nameWithOwner isEmpty defaultBranchRef { target { oid } }
legacy0: ref(qualifiedName:"refs/heads/dispatch-coordinator") { target { oid } }
legacy1: ref(qualifiedName:"refs/heads/gh-aw-work-queue") { target { oid } }
ref(qualifiedName:$ref) { target { ... on Commit { oid tree { oid entries { name type mode oid object {
... on Tree { entries { name type mode oid object { ... on Tree { entries { name type mode oid } } } } }
} } } } } }
log: object(expression:$log) { oid ... on Blob { text byteSize isTruncated } }
${journals}
}
${issueRead?.selections?.join("\n") || ""}
}`,
variables
);
const repository = response?.repository;
if (!repository || repository.nameWithOwner.toLowerCase() !== `${owner}/${repo}`.toLowerCase()) throw queueError("repository_unavailable", "checked queue repository identity is unavailable");
if (repository.ref === null) {
if (repository.isEmpty !== true && !repository.defaultBranchRef?.target?.oid) throw queueError("repository_unavailable", "cannot establish contents access before treating a queue as absent");
if (branch === "work-queue" && (repository.legacy0 || repository.legacy1)) throw queueError("unsupported_protocol", "a legacy queue cannot be implicitly adopted");
return { sha: null, treeSha: null, transactions: [], state: replayTransactions([]), branch, logPath: "work-queue.jsonl", repositoryId: repository.id, journal: new Map() };
}
const commit = repository.ref?.target;
if (!commit?.oid || !commit.tree?.oid || !Array.isArray(commit.tree.entries)) throw queueError("ledger_invalid", "checked transport requires an initialized queue");
const entries = commit.tree.entries;
if (entries.length > 16384) throw queueError("resource_limit", "queue tree entry limit exceeded");
if (entries.some(entry => entry.name === "dispatch-work-coordinator.jsonl")) throw queueError("unsupported_protocol", "legacy queue storage is unsupported");
const logs = entries.filter(entry => entry.name === "work-queue.jsonl");
const entry = logs[0];
if (logs.length !== 1 || entry.type !== "blob" || entry.mode !== 33188) throw queueError("ledger_invalid", "queue requires one regular canonical log");
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");
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()) {
const value = repository[`j${index}`];
const expected = directory.find(entry => entry.name === path.split("/").at(-1));
if (!expected && value === null) continue;
if (!expected || expected.type !== "blob" || expected.mode !== 33188 || !value || value.oid !== expected.oid) throw queueError("projection_journal_conflict", "journal does not belong to the immutable checked ledger head");
journal.set(path, parseStrictJSON(await checkedBlobText(githubClient, owner, repo, value, expected.oid, 1024 * 1024), { maxBytes: 1024 * 1024 }));
}
return { sha: commit.oid, treeSha: commit.tree.oid, transactions: state.transactions, state, branch, logPath: "work-queue.jsonl", repositoryId: repository.id, journal, nativeResponse: issueRead ? response : undefined };
}

async function writeCheckedFiles({ githubClient, owner, repo, branch, expectedHeadOid, files }) {
validateBranch(branch);
if (
!/^(?:[a-f0-9]{40}|[a-f0-9]{64})$/.test(expectedHeadOid) ||
files.length < 1 ||
files.length > 26 ||
new Set(files.map(file => file.path)).size !== files.length ||
files.some(
file =>
typeof file.path !== "string" ||
(file.path !== "work-queue.jsonl" && !JOURNAL_PATH.test(file.path)) ||
typeof file.content !== "string" ||
Buffer.byteLength(file.content, "utf8") > (file.path === "work-queue.jsonl" ? MAX_BYTES : 1024 * 1024)
)
)
throw queueError("publication_invalid", "checked publication requires explicit head and files");
const response = await githubClient.graphql("mutation CheckedWorkQueuePublication($input:CreateCommitOnBranchInput!) { createCommitOnBranch(input:$input) { commit { oid } } }", {
input: {
branch: { repositoryNameWithOwner: `${owner}/${repo}`, branchName: branch },
expectedHeadOid,
message: { headline: "Publish checked work queue projection" },
fileChanges: { additions: files.map(file => ({ path: file.path, contents: Buffer.from(file.content, "utf8").toString("base64") })) },
},
request: { retries: 0, timeout: 30000 },
});
const oid = response?.createCommitOnBranch?.commit?.oid;
if (typeof oid !== "string" || !/^[a-f0-9]{40}$|^[a-f0-9]{64}$/.test(oid)) throw queueError("publication_unresolved", "checked publication returned no commit identity");
return oid;
}

function writeCheckedCandidate({ githubClient, owner, repo, current, transactions }) {
return writeCheckedFiles({ githubClient, owner, repo, branch: current.branch, expectedHeadOid: current.sha, files: [{ path: current.logPath, content: transactions.map(commit => canonical(commit)).join("\n") + "\n" }] });
}

module.exports = { readCheckedQueue, writeCheckedFiles, writeCheckedCandidate };
2 changes: 2 additions & 0 deletions actions/setup/js/work_queue_dispatch.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,7 @@ async function resolveAdmissionResources(options, state, parameters, trustedCont
references.push({ target: edge, field: "resource", resource: { ...edge.resource, condition: edge.condition } });
}
if (node.subject) references.push({ target: node, field: "subject", resource: { ...node.subject, condition: node.subject.kind === "issue" ? "completed" : "merged" } });
if (node.backing_issue) references.push({ target: node, field: "backing_issue", resource: { ...node.backing_issue, condition: "completed" } });
}
if (!references.length) return parameters;
log.debug("admission.resolve.references", { references: references.length });
Expand Down Expand Up @@ -308,6 +309,7 @@ function acceptedSubmissionParameters(state, trustedContext, parameters, prior)
return edge.kind !== "work" && original?.kind === edge.kind && original.condition === edge.condition ? { ...edge, resource: restoreIdentity(edge.resource, original.resource) } : edge;
});
if (node.subject) node.subject = restoreIdentity(node.subject, accepted?.subject);
if (node.backing_issue) node.backing_issue = restoreIdentity(node.backing_issue, accepted?.backing_issue);
}
return normalized;
}
Expand Down
13 changes: 12 additions & 1 deletion actions/setup/js/work_queue_documentation.test.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,24 @@ describe("work-queue deployment documentation", () => {
it("installs the documented policy after replacing identity placeholders", () => {
const source = readRepositoryFile(deploymentPath);
const examples = [...source.matchAll(/```json(?: [^\n]*)?\n([\s\S]*?)\n```/g)];
expect(examples).toHaveLength(1);
expect(examples).toHaveLength(2);
const policy = JSON.parse(examples[0][1].replaceAll("REPLACE_WITH_PRODUCER_ACTOR_ID", "11").replaceAll("REPLACE_WITH_WORKER_CREDENTIAL_ACTOR_ID", "12").replaceAll("REPLACE_WITH_40_OR_64_HEX_COMMIT_SHA", "a".repeat(40)));

expect(validatePolicy(policy)).toBe(policy);
expect(policy.accounting_weights).toEqual({ "": 1 });
expect(policy.producers["11"].fairness_keys).toEqual([""]);
expect(policy.limits).toEqual(DEFAULT_LIMITS);
const projection = JSON.parse(
examples[1][1]
.replaceAll("REPLACE_WITH_NATIVE_PRINCIPAL_ID", "12")
.replaceAll("REPLACE_WITH_40_OR_64_HEX_COMMIT_SHA", "a".repeat(40))
.replaceAll("REPLACE_WITH_NUMERIC_REPOSITORY_ID", "9876")
.replaceAll("REPLACE_WITH_NUMERIC_ISSUE_ID", "9007199254740993")
.replaceAll("REPLACE_WITH_ISSUE_NUMBER", "42")
);
expect(validatePolicy({ ...policy, ...projection }).projectors).toEqual(projection.projectors);
expect(projection.projectors[0].completion_policy).toBe("keep-open");
expect(projection.projectors[0].backing_issues[0]).toEqual({ kind: "issue", host: "github.com", repository: "github/gh-aw", repository_id: "9876", resource_id: "9007199254740993", number: "42" });
});

it("separates published docs and specifications from bounded agent instructions", () => {
Expand Down
Loading
Loading