Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
379619d71f |
@@ -1,93 +0,0 @@
|
||||
# W-B — Measure Pi's real tool registry
|
||||
|
||||
- **Task / internal ref:** W-B from the lease-remediation orchestrator brief (no matching `docs/TASKS.md` row; workers do not modify that file)
|
||||
- **Objective:** identify the exact tool names emitted as `event.toolName` by the installed Pi runtime and compare them with the broker's Pi read-only carve-out.
|
||||
- **Scope:** measurement and report only; no broker or runtime source changes. W-C is out of scope.
|
||||
- **Budget:** no explicit token cap; constrained to this scratchpad and one local commit.
|
||||
- **Installed runtime:** `@earendil-works/pi-coding-agent` / `pi` `0.84.1`.
|
||||
|
||||
## Method
|
||||
|
||||
I created a throwaway extension at `/tmp/measure-pi-tool-registry.ts` (not in the worktree). On `session_start` it recorded `pi.getAllTools()` and `pi.getActiveTools()`; on every `tool_call` it appended the exact `event.toolName`. I then launched an isolated, ephemeral Pi session with all built-ins explicitly selected:
|
||||
|
||||
```text
|
||||
PI_OFFLINE=1 pi --mode print --no-session --no-approve \
|
||||
--no-context-files --no-skills --no-prompt-templates --no-extensions \
|
||||
-e /tmp/measure-pi-tool-registry.ts \
|
||||
--tools read,bash,edit,write,grep,find,ls <deterministic probe prompt>
|
||||
```
|
||||
|
||||
The prompt exercised file read, content search, file search, directory listing, shell execution, file write, and file edit. Pi exited `0`; every selected tool produced one `tool_call`. The write/edit control artifact ended with exact content `after`, proving the mutating calls executed in order.
|
||||
|
||||
This runtime observation was cross-checked against the installed distribution's canonical registry at `dist/core/tools/index.js:17`, which declares the same seven names. The gate consumes the measured field directly at `packages/mosaic/framework/runtime/pi/mosaic-extension.ts:368`.
|
||||
|
||||
## Exact distinct built-in set
|
||||
|
||||
The installed Pi built-in registry is exactly:
|
||||
|
||||
```text
|
||||
{bash, edit, find, grep, ls, read, write}
|
||||
```
|
||||
|
||||
| Tool | Runtime registry observation | `tool_call` observation | Installed definition |
|
||||
| --- | --- | --- | --- |
|
||||
| `read` | `<builtin:read>` | observed once | `dist/core/tools/read.js:138` |
|
||||
| `bash` | `<builtin:bash>` | observed once | `dist/core/tools/bash.js:231` |
|
||||
| `edit` | `<builtin:edit>` | observed once | `dist/core/tools/edit.js:170` |
|
||||
| `write` | `<builtin:write>` | observed once | `dist/core/tools/write.js:138` |
|
||||
| `grep` | `<builtin:grep>` | observed once | `dist/core/tools/grep.js:79` |
|
||||
| `find` | `<builtin:find>` | observed once | `dist/core/tools/find.js:79` |
|
||||
| `ls` | `<builtin:ls>` | observed once | `dist/core/tools/ls.js:61` |
|
||||
|
||||
The raw distinct `event.toolName` result was:
|
||||
|
||||
```json
|
||||
["bash", "edit", "find", "grep", "ls", "read", "write"]
|
||||
```
|
||||
|
||||
Pi registers all seven, but its default active set is only `read`, `bash`, `edit`, and `write` (`dist/core/sdk.js:132`). The probe explicitly activated all seven so the three search/list tools could be observed at the hook.
|
||||
|
||||
## Positive control
|
||||
|
||||
The known `read` tool was the control. The method surfaced it twice:
|
||||
|
||||
1. `pi.getAllTools()` returned `read` with source path `<builtin:read>`.
|
||||
2. Reading `/tmp/pi-registry-probe/seed.txt`, which contained `CONTROL_TOKEN`, produced one hook record with `event.toolName === "read"`.
|
||||
|
||||
The control was therefore positive; the seven-name result is measured, not an empty-probe inference.
|
||||
|
||||
## Carve-out comparison and collision result
|
||||
|
||||
The broker currently declares `{"read", "grep", "find", "ls"}` at `packages/mosaic/framework/tools/lease-broker/daemon.py:54`.
|
||||
|
||||
- `read`: real built-in.
|
||||
- `grep`: real built-in.
|
||||
- `find`: real built-in.
|
||||
- `ls`: real built-in.
|
||||
|
||||
All four carve-out names are exact, case-sensitive Pi tool names.
|
||||
|
||||
The general execution/writing tool names are `bash`, `edit`, and `write`. Their intersection with the carve-out is empty:
|
||||
|
||||
```text
|
||||
{bash, edit, write} ∩ {read, grep, find, ls} = ∅
|
||||
```
|
||||
|
||||
Therefore no general shell-exec or file-mutating Pi tool shares a name with a carve-out entry. `grep` and `find` may invoke constrained search helpers internally, but neither exposes an arbitrary command interface; the arbitrary command tool is distinctly named `bash`.
|
||||
|
||||
The Mosaic extension separately registers the non-built-in custom tool `mosaic_context_recover` at `packages/mosaic/framework/runtime/pi/mosaic-extension.ts:379`; the broker handles that identity through its dedicated recovery exemption rather than the read-only set (`daemon.py:722`). Unknown or third-party custom tools are not part of Pi's built-in seven-name registry and remain outside the carve-out.
|
||||
|
||||
## Verification evidence
|
||||
|
||||
- `pi --version` → `0.84.1`.
|
||||
- Isolated probe exit → `0`.
|
||||
- Runtime `getAllTools()` count → `7`, all with `sourceInfo.source === "builtin"`.
|
||||
- Distinct hook names → `bash`, `edit`, `find`, `grep`, `ls`, `read`, `write`.
|
||||
- Hook counts → exactly one call for each of the seven names.
|
||||
- Mutation artifact after `write` then `edit` → exact content `after`.
|
||||
- Installed registry source → `allToolNames = new Set(["read", "bash", "edit", "write", "grep", "find", "ls"])`.
|
||||
|
||||
## Risks / limitations
|
||||
|
||||
- The probe deliberately disabled all other extensions, so extension-defined third-party tools were excluded from the built-in registry measurement. The production gate still receives those names and treats names outside the broker carve-out as mutating/fail-closed.
|
||||
- Explicit `--tools` activation was required to exercise `grep`, `find`, and `ls`; this does not imply they are active in Pi's default four-tool configuration.
|
||||
@@ -1,99 +0,0 @@
|
||||
# PR merge squash message field
|
||||
|
||||
- **Charter:** `/home/hermes/agent-work/CHARTER-PRMERGE-MESSAGE-FIELD.md`
|
||||
- **Owner:** `be-coder-08`
|
||||
- **Branch:** `fix/pr-merge-message-field`
|
||||
- **Base:** remote `main` / local `origin/main` at `85d2108e4ed15c744ad3b87a5b629e7b2d39405a`
|
||||
- **Estate:** HOMELAB tooling shared by HOMELAB and USC
|
||||
|
||||
## Objective
|
||||
|
||||
Add an optional, identity-checked Gitea squash message to `pr-merge.sh` so genuine multi-author PRs retain non-poster branch authors without weakening hardcoded squash behavior.
|
||||
|
||||
## Binding requirements
|
||||
|
||||
1. `Do` remains hardcoded to `squash`; no provider/repository default may select merge style.
|
||||
2. A verified trailer uses a PR commit's linked `author.login` and that same commit's author email. No `/users/{login}` primary-email lookup occurs. Recorded rationale: this asks only what the provider can answer.
|
||||
3. A commit with `author.login` null blocks before merge, prints both the null provider fact and commit email fact, and names the escalation principal.
|
||||
4. The BLOCK arm must be observed firing; a normal canonical single-author API payload remains explicit squash plus its reviewed `head_commit_id`.
|
||||
5. Every provider mutation is read back from the provider; no real PR is merged during tests.
|
||||
|
||||
## Derived interface decisions
|
||||
|
||||
- Add `--co-author-trailers` rather than accepting arbitrary message text. The wrapper enumerates PR commits and constructs trailers, making an unchecked `Co-authored-by` line unexpressible.
|
||||
- Require `--escalate-to PRINCIPAL` with `--co-author-trailers`, so the BLOCK diagnostic always names a principal rather than a generic role.
|
||||
- Do not expose `MergeTitleField` separately. When trailers exist, set it from the provider PR title and set `MergeMessageField` only to construction-generated trailers. This preserves one provider source for the title and avoids an unrelated caller-controlled degree of freedom.
|
||||
- Preserve first-commit order and emit one trailer per distinct non-poster `author.login`, using that first linked commit's own email.
|
||||
|
||||
## Canonical delivery plan
|
||||
|
||||
1. Port the capability into the installed source of truth, `packages/mosaic/framework/tools/git/pr-merge.sh`; do not retain `infra/fleet/tools/git` as a second copy.
|
||||
2. Preserve canonical `--expect-head`, exact head branch/repository/SHA queue inspection, Gitea atomic head pinning, GitHub `--match-head-commit`, and delete-after-merge semantics.
|
||||
3. Do not port the deployed-only `--skip-queue-guard` bypass. Add the focused harness to the canonical framework-shell suite and re-establish RED/GREEN on the packaged baseline.
|
||||
4. Deliver through a reviewed package release followed by `mosaic update` with its default framework reseed. The installer snapshots, manifest-syncs framework-owned `tools/**`, and rolls back on failure.
|
||||
5. Before either estate relies on the change, require installed/package hash equality, `MergeMessageField` presence, and a green focused harness. Release/reseed ownership is currently unassigned and blocks activation after source merge.
|
||||
|
||||
## Evidence
|
||||
|
||||
- RED against the byte-identical deployed baseline (`sha256 08a65e8584c5…`): rc 1 with eight named failures. The wrapper rejected `--co-author-trailers`; the null-login path emitted none of the required BLOCK facts/principal; and both verified/ordinary API paths failed the stdin-config credential assertion (ordinary path exposed the fixture token through curl argv). Log: `/home/hermes/agent-work/be-coder-08/evidence/prmerge-message-field-red.log`.
|
||||
- GREEN on the deployed-baseline candidate: verified linked multi-author payload, null-login BLOCK, required named principal, explicit squash, stdin-config token transport, and absence of `/users` lookup all passed. Log: `/home/hermes/agent-work/be-coder-08/evidence/prmerge-message-field-green.log`.
|
||||
- RED against canonical packaged baseline `c581ef48…`: rc 1 with 32 assertions. It rejects the new option, and the first harness version did not satisfy canonical head branch/repository/SHA metadata. Log: `/home/hermes/agent-work/be-coder-08/evidence/prmerge-packaged-baseline-red.log`. The port adapts the fixture rather than weakening canonical head controls.
|
||||
- Provider capability probe against `git.mosaicstack.dev`: authenticated `be-coder-08` POST to deliberately nonexistent PR `2147483647` with both message fields returned JSON HTTP 404; the unauthenticated same request returned JSON HTTP 401 (not the charter's predicted 403). The authenticated-vs-unauthenticated differential proves write authorization resolved while no mergeable subject existed. `tl-mosaic` ruled the literal non-load-bearing: preserve the observed 404/401 pair and do not manufacture a 403 case. No cause was inferred and no real PR was targeted.
|
||||
- Provider-generated trailer behavior is not treated as exclusive or absent. The wrapper's VERIFIED/BLOCK decision binds each requested non-poster trailer to commit `author.login` plus that commit's email; it does not assume `MergeMessageField` is the squash's only trailer source. The poster is omitted from the constructed list because the resulting squash author already records the poster; any additional provider-generated trailer is outside this change's unmeasured mechanism.
|
||||
- An early candidate SHA-256 `5de32876990e4f26920448cb3220cc7f1146d558b4dd2bc1ee1a2abee2f2cbe6` passed the initial harness, then author-side review found credential-fallback and argv-exposure defects. The live deployed wrapper was atomically restored to baseline SHA-256 `08a65e8584c52c6d41ea1c686f8b95585c21e4b37320a2447eba09359a0e02c1`; the remediated candidate remains only in the worktree.
|
||||
|
||||
## Remediation and current review state
|
||||
|
||||
1. Token and Basic Auth now use stdin curl configuration, not argv. PR title, contributor email, and the JSON payload also remain out of child argv.
|
||||
2. Each credential attempt binds commit inspection and merge. A token failure during either inspection or mutation causes Basic fallback to repeat inspection before mutation; the payload pins the inspected `head_commit_id`.
|
||||
3. Focused tests cover token-resolution fail-closed behavior, both HTTP-401 fallback seams, metadata/credential argv absence, null-login BLOCK, explicit squash, canonical reviewed-head binding, unchanged ordinary payload, and retained log-safe provider diagnostics. Token-resolution RED: `/home/hermes/agent-work/be-coder-08/evidence/prmerge-token-resolution-red.log`.
|
||||
4. Codex review rounds 3–5 requested retained provider error text, log-safe provider diagnostics, fail-closed credential fallback, stable value-option parsing, and PR-title trailer-injection prevention. These are remediated with regression assertions. A post-remediation independent review is still required.
|
||||
5. **Accepted linkage limitation:** `author.login` resolution proves that the commit address maps to a registered provider account. It does not prove that the named principal authored the commit because Git author metadata is self-asserted. This gate checks attribution linkage, not authorship; commit signing is out of scope and currently unadopted. Coordinators explicitly ruled that this does not add a third state.
|
||||
6. Codex's sandbox could not execute the harness because its checkout was read-only; that environmental limitation is recorded separately from host-side test results.
|
||||
|
||||
## Disposable provider fixture acceptance
|
||||
|
||||
- Use a retained scratch repository only, with two branch authors and `author != committer` on at least one commit.
|
||||
- Arm A supplies a message-field trailer for one non-poster; record whether that value lands without forcing the partial-pair result into under-specified `APPENDS`/`REPLACES` labels. Demonstrate an absence control.
|
||||
- Arm B includes a registered trailer for a different non-poster on a branch commit; record whether it survives or drops. Verify identity through an existing commit whose `author.login` resolves and demonstrate an absence control.
|
||||
- Parse landed trailers key-agnostically with `^[A-Za-z-]+-[Bb]y:` and record generated poster pair presence/absence plus resulting poster attribution.
|
||||
- Record `/users/<login>` status and raw email only as non-gating estate telemetry. Never read `active`, `visibility`, or any profile field as an identity gate.
|
||||
- Use distinct principals: poster `be-coder-08`, merger `Mos`, Arm A `be-coder-07`, and Arm B `be-coder-06`. Capture every trailer-shaped line verbatim and in order. Zero trailer lines means the generator did not fire and the run is `VOID`, not evidence that either arm dropped.
|
||||
- Report the same read-back evidence to `mos-claude` on socket `default` and `tl-mosaic` on socket `mosaic-fleet`. Report values rather than mechanism inferences and stop on any poster-attribution regression.
|
||||
|
||||
## Fixture preflight
|
||||
|
||||
- Retained public repository: `mosaicstack/prmerge-trailer-fixture`; PR `#1`, posted by `be-coder-08` and reserved for merge by `Mos`.
|
||||
- Existing `mosaicstack/stack` commits resolve `be-coder-07` and `be-coder-06` through `author.login`; exact addresses are `[email protected]` and `[email protected]`.
|
||||
- Non-gating HOMELAB telemetry for authenticated reader `be-coder-08`: `/api/v1/users/be-coder-06` returned HTTP 200 with raw `email` value `[email protected]`.
|
||||
- Provider preflight showed PR commit enumeration is newest-first. A new RED test proved that deriving `head_commit_id` from the final array element selected the wrong commit. The candidate now reads `.head.sha` from the authenticated PR endpoint before enumeration, verifies it appears in the commit set, and atomically pins that SHA in the explicit squash payload. RED: `/home/hermes/agent-work/be-coder-08/evidence/prmerge-head-order-red.log`.
|
||||
- Fixture PR head `f6ba6e5105031fa21f5ff7bd8e4379d99c16e1de` has `author.login=be-coder-07`, `committer.login=be-coder-08`, and branch-message trailer `Co-authored-by: be-coder-06 <[email protected]>`.
|
||||
|
||||
## Fixture result
|
||||
|
||||
- `Mos` merged retained fixture PR `#1` through staged candidate SHA-256 `60e779a85fd13b729d859ea7c986d1e9b1641b97991611329226c1b3113ffb6e`; resulting squash commit: `3f550715d9bc716426fd355a65fe997b3a90fa7d` with one parent.
|
||||
- Provider read-back: poster/commit author `be-coder-08`, committer/merger `Mos`. The run is non-void.
|
||||
- Trailer-shaped lines, verbatim and in order:
|
||||
1. `Co-authored-by: be-coder-07 <[email protected]>`
|
||||
2. `Co-authored-by: be-coder-08 <[email protected]>`
|
||||
- Arm A supplied field value (`be-coder-07`) landed. Arm B branch trailer (`be-coder-06`) dropped. Both fabricated absence controls remained absent. No `Co-committed-by:` line landed.
|
||||
- The candidate payload construction explicitly excludes the poster and supplied only the Arm A `be-coder-07` line. Therefore the landed poster line was provider-generated, not candidate-composed. The raw result supports `FIELD LANDS`, `BRANCH DROPS`, and `POSTER GENERATED`; it does not support a claim that candidate code supplied the poster. Evidence: `/home/hermes/agent-work/be-coder-08/evidence/prmerge-fixture-readback.log` and the retained provider object.
|
||||
- Retained fixture PR `#2` measured the N=2 shape needed by `#1030`: supplied `be-coder-07` then `be-coder-06`; both landed in that order, followed by the provider-generated poster line. No truncation or dedup occurred at N=2. Resulting squash: `39db9d13aed0…`.
|
||||
|
||||
## Current hold point
|
||||
|
||||
PR `mosaicstack/stack#1066` is open. Its first frozen head `f4b162fa…` was terminal-green in Woodpecker `mosaic` pipeline `#2225`, but that evidence becomes stale when the canonical port moves the head. The deployed wrapper remains baseline `08a65e85…`; no manual copy will occur. Canonical port tests, commit amendment, rebase, one guarded force-with-lease, exact-head CI, and new independent review remain. Even after source merge, activation remains blocked on an assigned package-release/reseed owner and installed-byte read-back.
|
||||
|
||||
## Security review 96 remediation
|
||||
|
||||
Exact reviewed predecessor head: `1ceb11058f64dd7f4a817ceb2124f980a1c4dd23`.
|
||||
|
||||
RED-first focused harness produced 10 named failures: all curl calls lacked size/time/connect bounds; raw ESC email reached mutation; oversized and stalled curl failures were discarded and reached mutation; nonempty Basic output with resolver rc 91 authorized mutation.
|
||||
|
||||
Security remediation:
|
||||
|
||||
- Removed the cross-principal HTTP-401 Basic fallback. Both inspection-401 and merge-401 paths now refuse without Basic resolution or mutation; `get_gitea_basic_auth` references in the merge subject are 0.
|
||||
- Applied `--max-filesize`, `--max-time`, and `--connect-timeout` to all 3/3 provider curl sites and fail closed on curl transport rc at all 3/3 sites.
|
||||
- Required linked email bytes to be ASCII and printable before constructing `MergeMessageField`; guarded construction sites 1/1.
|
||||
|
||||
GREEN: message-field, exact-head, empty-UID/API, queue branch/repository/SHA, bash syntax, ShellCheck, and diff check pass. R7 total-removal mutants went RED: email guard 3 rows; bound switches 1 row; transport-rc guards 4 rows; HTTP-401 refusal 3 rows. R7 bound: mutants prove total removal only; explicit denominators above prove site coverage.
|
||||
@@ -1 +0,0 @@
|
||||
I invoked this registered command to authorize lease promotion; follow the local seat broker's injected receipt confirmation instruction exactly.
|
||||
@@ -32,18 +32,6 @@
|
||||
]
|
||||
}
|
||||
],
|
||||
"UserPromptSubmit": [
|
||||
{
|
||||
"matcher": "^/mosaic-promote$",
|
||||
"hooks": [
|
||||
{
|
||||
"type": "command",
|
||||
"command": "python3 ~/.config/mosaic/tools/lease-broker/promote-begin.py",
|
||||
"timeout": 15
|
||||
}
|
||||
]
|
||||
}
|
||||
],
|
||||
"PreToolUse": [
|
||||
{
|
||||
"matcher": ".*",
|
||||
@@ -93,8 +81,8 @@
|
||||
"hooks": [
|
||||
{
|
||||
"type": "command",
|
||||
"command": "python3 ~/.config/mosaic/tools/lease-broker/receipt-observer-client.py --runtime claude --latest-entry; observer_status=$?; python3 ~/.config/mosaic/tools/lease-broker/promote-complete.py; exit $observer_status",
|
||||
"timeout": 15
|
||||
"command": "python3 ~/.config/mosaic/tools/lease-broker/receipt-observer-client.py --runtime claude --latest-entry",
|
||||
"timeout": 3
|
||||
},
|
||||
{
|
||||
"type": "command",
|
||||
|
||||
@@ -153,24 +153,7 @@ if [[ $link_only -eq 1 ]]; then
|
||||
exit 0
|
||||
fi
|
||||
|
||||
# Skills are linked into the MOSAIC-OWNED harness homes, never a base install.
|
||||
# Paths mirror the config-dir env vars the launcher injects (HARNESS_HOME_ENV in
|
||||
# commands/launch.js):
|
||||
# claude CLAUDE_CONFIG_DIR -> <home>/skills
|
||||
# pi PI_CODING_AGENT_DIR -> <home>/skills (replaces ~/.pi/agent)
|
||||
# codex CODEX_HOME -> <home>/skills
|
||||
# opencode XDG_CONFIG_HOME -> <home>/opencode/skills (XDG adds a level)
|
||||
link_targets=(
|
||||
"$MOSAIC_HOME/.claude/skills"
|
||||
"$MOSAIC_HOME/.codex/skills"
|
||||
"$MOSAIC_HOME/.opencode/opencode/skills"
|
||||
"$MOSAIC_HOME/.pi/skills"
|
||||
)
|
||||
|
||||
# Pre-isolation installs planted the same symlink farm directly in the operator's
|
||||
# base installs. Those are now orphaned: the launcher no longer reads them, but
|
||||
# they persist and make a "clean" base install look mosaic-managed.
|
||||
legacy_link_targets=(
|
||||
"$HOME/.claude/skills"
|
||||
"$HOME/.codex/skills"
|
||||
"$HOME/.config/opencode/skills"
|
||||
@@ -262,72 +245,13 @@ prune_stale_links_in_target() {
|
||||
# -m resolves lexical dangling targets too. If resolution fails, ownership
|
||||
# is unproven and the link must be preserved.
|
||||
resolved="$(readlink -m "$link_path" 2>/dev/null || true)"
|
||||
# $canonical_real must be length-checked BEFORE use as a prefix: if it were
|
||||
# ever empty, "$resolved" == "$canonical_real/"* collapses to == "/"* and
|
||||
# matches every absolute path. Combined with the is_mosaic_skill_name skip
|
||||
# above, that inverts the function precisely — it would delete exactly the
|
||||
# FOREIGN symlinks and keep the mosaic ones. (#1087, reported by mos-claude.)
|
||||
if [[ -n "$resolved" && -n "$canonical_real" && "$resolved" == "$canonical_real/"* ]]; then
|
||||
if [[ -n "$resolved" && "$resolved" == "$canonical_real/"* ]]; then
|
||||
rm -f "$link_path"
|
||||
echo "[mosaic-skills] Removed stale retired skill link: $link_path"
|
||||
fi
|
||||
done < <(find "$target_dir" -mindepth 1 -maxdepth 1 -type l -print0)
|
||||
}
|
||||
|
||||
# Remove mosaic-owned symlinks left in a base install by a pre-isolation sync.
|
||||
#
|
||||
# Ownership is proven by RESOLUTION, not by name: only links resolving inside the
|
||||
# canonical or local skills dirs are removed. Anything else — a real directory, a
|
||||
# link elsewhere, an unresolvable link — is left untouched. This mirrors the
|
||||
# refusal in commands/skill.js ("only symlinks pointing inside the Mosaic skills
|
||||
# directory are managed") and preserves e.g. codex's own `.system` dir.
|
||||
#
|
||||
# The directory itself is kept: mosaic-doctor warns when ~/.pi/agent/skills is
|
||||
# missing, and an empty dir is the correct end state, not an absent one.
|
||||
cleanup_legacy_target() {
|
||||
local target_dir="$1"
|
||||
local removed=0 kept=0
|
||||
|
||||
[[ -d "$target_dir" ]] || return 0
|
||||
|
||||
while IFS= read -r -d '' link_path; do
|
||||
local resolved owned=0
|
||||
resolved="$(readlink -m "$link_path" 2>/dev/null || true)"
|
||||
|
||||
# Guard the empty-prefix trap: an unset *_real would make "$resolved" == "/"*
|
||||
# match every absolute path and delete foreign links.
|
||||
if [[ -n "$resolved" ]]; then
|
||||
if [[ -n "$canonical_real" && "$resolved" == "$canonical_real/"* ]]; then
|
||||
owned=1
|
||||
elif [[ -n "$local_real" && "$resolved" == "$local_real/"* ]]; then
|
||||
owned=1
|
||||
fi
|
||||
fi
|
||||
|
||||
if [[ $owned -eq 1 ]]; then
|
||||
rm -f "$link_path"
|
||||
removed=$((removed + 1))
|
||||
else
|
||||
kept=$((kept + 1))
|
||||
fi
|
||||
done < <(find "$target_dir" -mindepth 1 -maxdepth 1 -type l -print0)
|
||||
|
||||
if [[ $removed -gt 0 ]]; then
|
||||
echo "[mosaic-skills] Legacy cleanup: removed $removed mosaic symlink(s) from $target_dir (preserved $kept foreign)"
|
||||
fi
|
||||
}
|
||||
|
||||
for legacy in "${legacy_link_targets[@]}"; do
|
||||
# Skip anything that is also a current target, so isolation can never
|
||||
# self-destruct if the two lists ever overlap.
|
||||
skip=0
|
||||
for target in "${link_targets[@]}"; do
|
||||
[[ "$legacy" == "$target" ]] && skip=1
|
||||
done
|
||||
[[ $skip -eq 1 ]] && continue
|
||||
cleanup_legacy_target "$legacy"
|
||||
done
|
||||
|
||||
for target in "${link_targets[@]}"; do
|
||||
mkdir -p "$target"
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
#!/bin/bash
|
||||
# pr-merge.sh - Merge pull requests on Gitea or GitHub
|
||||
# Usage: pr-merge.sh -n PR_NUMBER [-m squash] [-d] [--expect-head SHA] [--co-author-trailers --escalate-to PRINCIPAL]
|
||||
# Usage: pr-merge.sh -n PR_NUMBER [-m squash] [-d]
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
@@ -14,8 +14,6 @@ MERGE_METHOD="squash"
|
||||
DELETE_BRANCH=false
|
||||
DRY_RUN=false
|
||||
EXPECT_HEAD=""
|
||||
CO_AUTHOR_TRAILERS=false
|
||||
ESCALATE_TO=""
|
||||
|
||||
usage() {
|
||||
cat <<EOF
|
||||
@@ -29,16 +27,12 @@ Options:
|
||||
-d, --delete-branch Delete the head branch after merge
|
||||
--dry-run Run metadata/login preflight without merging
|
||||
--expect-head SHA Refuse unless the PR head matches this full commit SHA
|
||||
--co-author-trailers Build verified trailers from linked PR commit authors
|
||||
--escalate-to NAME Named principal for an unresolved-author BLOCK
|
||||
-h, --help Show this help message
|
||||
|
||||
Examples:
|
||||
$(basename "$0") -n 42 # Merge PR #42
|
||||
$(basename "$0") -n 42 -m squash # Squash merge
|
||||
$(basename "$0") -n 42 -d # Squash merge and delete branch
|
||||
$(basename "$0") -n 42 --expect-head 0123456789abcdef0123456789abcdef01234567
|
||||
$(basename "$0") -n 42 --co-author-trailers --escalate-to tl-mosaic
|
||||
EOF
|
||||
exit "${1:-1}"
|
||||
}
|
||||
@@ -63,25 +57,9 @@ while [[ $# -gt 0 ]]; do
|
||||
shift
|
||||
;;
|
||||
--expect-head)
|
||||
if [[ $# -lt 2 ]]; then
|
||||
echo "Error: --expect-head requires one full commit SHA." >&2
|
||||
exit 1
|
||||
fi
|
||||
EXPECT_HEAD="$2"
|
||||
shift 2
|
||||
;;
|
||||
--co-author-trailers)
|
||||
CO_AUTHOR_TRAILERS=true
|
||||
shift
|
||||
;;
|
||||
--escalate-to)
|
||||
if [[ $# -lt 2 ]]; then
|
||||
echo "Error: --escalate-to requires one principal name." >&2
|
||||
exit 1
|
||||
fi
|
||||
ESCALATE_TO="$2"
|
||||
shift 2
|
||||
;;
|
||||
-h|--help)
|
||||
usage 0
|
||||
;;
|
||||
@@ -110,30 +88,17 @@ if [[ -n "$EXPECT_HEAD" && ! "$EXPECT_HEAD" =~ ^[0-9a-fA-F]{40}$ ]]; then
|
||||
echo "Error: --expect-head must be a full 40-character hexadecimal commit SHA." >&2
|
||||
exit 1
|
||||
fi
|
||||
if [[ "$CO_AUTHOR_TRAILERS" == true && -z "$ESCALATE_TO" ]]; then
|
||||
echo "Error: --co-author-trailers requires --escalate-to with a named principal." >&2
|
||||
exit 1
|
||||
fi
|
||||
if [[ -n "$ESCALATE_TO" && ! "$ESCALATE_TO" =~ ^[A-Za-z0-9_.-]+$ ]]; then
|
||||
echo "Error: --escalate-to must be one exact principal name." >&2
|
||||
exit 1
|
||||
fi
|
||||
if [[ "$CO_AUTHOR_TRAILERS" != true && -n "$ESCALATE_TO" ]]; then
|
||||
echo "Error: --escalate-to is valid only with --co-author-trailers." >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
PR_METADATA="$("$SCRIPT_DIR/pr-metadata.sh" -n "$PR_NUMBER")"
|
||||
BASE_BRANCH="$(printf '%s' "$PR_METADATA" | python3 -c 'import json, sys; print((json.load(sys.stdin).get("baseRefName") or "").strip())')"
|
||||
HEAD_BRANCH="$(printf '%s' "$PR_METADATA" | python3 -c 'import json, sys; print((json.load(sys.stdin).get("headRefName") or "").strip())')"
|
||||
HEAD_SHA="$(printf '%s' "$PR_METADATA" | python3 -c 'import json, sys; print((json.load(sys.stdin).get("headRefOid") or "").strip())')"
|
||||
HEAD_REPO="$(printf '%s' "$PR_METADATA" | python3 -c 'import json, sys; value=json.load(sys.stdin).get("headRepository") or ""; print((value.get("nameWithOwner") or value.get("full_name") or "") if isinstance(value, dict) else str(value).strip())')"
|
||||
PR_TITLE="$(printf '%s' "$PR_METADATA" | python3 -c 'import json, sys; print((json.load(sys.stdin).get("title") or "").strip())')"
|
||||
PR_AUTHOR="$(printf '%s' "$PR_METADATA" | python3 -c 'import json, sys; value=json.load(sys.stdin).get("author") or ""; print((value.get("login") or "").strip() if isinstance(value, dict) else str(value).strip())')"
|
||||
if [[ "$BASE_BRANCH" != "main" ]]; then
|
||||
echo "Error: Mosaic policy allows merges only for PRs targeting 'main' (found '$BASE_BRANCH')." >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
if [[ -z "$HEAD_BRANCH" || -z "$HEAD_REPO" || ! "$HEAD_SHA" =~ ^[0-9a-fA-F]{40}$ ]]; then
|
||||
echo "Error: Could not resolve the PR head branch, repository, and full commit SHA for queue inspection." >&2
|
||||
exit 1
|
||||
@@ -157,442 +122,70 @@ PLATFORM=$(detect_platform)
|
||||
OWNER=$(get_repo_owner)
|
||||
REPO=$(get_repo_name)
|
||||
|
||||
write_curl_auth_config() {
|
||||
local mode="$1" credential="$2"
|
||||
printf '%s' "$credential" | python3 -c '
|
||||
import sys
|
||||
mode = sys.argv[1]
|
||||
credential = sys.stdin.read()
|
||||
if not credential or any(char in credential for char in "\r\n"):
|
||||
raise SystemExit(1)
|
||||
escaped = credential.replace("\\", "\\\\").replace("\"", "\\\"")
|
||||
if mode == "token":
|
||||
print(f"header = \"Authorization: token {escaped}\"")
|
||||
elif mode == "basic":
|
||||
print(f"user = \"{escaped}\"")
|
||||
else:
|
||||
raise SystemExit(1)
|
||||
' "$mode"
|
||||
}
|
||||
|
||||
LAST_GITEA_HTTP_CODE="000"
|
||||
LAST_GITEA_ERROR=""
|
||||
MERGE_TEMP_DIRS=()
|
||||
GITEA_CURL_MAX_BYTES="${MOSAIC_GITEA_CURL_MAX_BYTES:-1048576}"
|
||||
GITEA_CURL_MAX_TIME="${MOSAIC_GITEA_CURL_MAX_TIME_SEC:-30}"
|
||||
GITEA_CURL_CONNECT_TIMEOUT="${MOSAIC_GITEA_CURL_CONNECT_TIMEOUT_SEC:-10}"
|
||||
for bound in "$GITEA_CURL_MAX_BYTES" "$GITEA_CURL_MAX_TIME" "$GITEA_CURL_CONNECT_TIMEOUT"; do
|
||||
if [[ ! "$bound" =~ ^[1-9][0-9]*$ ]]; then
|
||||
echo "Error: Gitea curl bounds must be positive integers; refusing request." >&2
|
||||
exit 1
|
||||
fi
|
||||
done
|
||||
GITEA_CURL_BOUNDS=(
|
||||
--max-filesize "$GITEA_CURL_MAX_BYTES"
|
||||
--max-time "$GITEA_CURL_MAX_TIME"
|
||||
--connect-timeout "$GITEA_CURL_CONNECT_TIMEOUT"
|
||||
)
|
||||
|
||||
format_gitea_error_response() {
|
||||
local response_file="$1"
|
||||
python3 - "$response_file" <<'PY'
|
||||
import json
|
||||
import sys
|
||||
|
||||
with open(sys.argv[1], "rb") as handle:
|
||||
raw = handle.read(65536)
|
||||
try:
|
||||
response = json.loads(raw.decode("utf-8", errors="replace"))
|
||||
except (UnicodeDecodeError, json.JSONDecodeError):
|
||||
message = "non-JSON response omitted"
|
||||
else:
|
||||
if isinstance(response, dict):
|
||||
message = response.get("message") or response.get("error")
|
||||
if not message and response.get("errors") is not None:
|
||||
message = json.dumps(response["errors"], separators=(",", ":"))
|
||||
else:
|
||||
message = None
|
||||
if not message:
|
||||
message = "JSON response contained no error message"
|
||||
message = str(message)
|
||||
if len(message) > 500:
|
||||
message = message[:500] + "..."
|
||||
print(ascii(message))
|
||||
PY
|
||||
}
|
||||
|
||||
cleanup_merge_temp_dirs() {
|
||||
local path
|
||||
for path in "${MERGE_TEMP_DIRS[@]}"; do
|
||||
[[ -n "$path" ]] && rm -rf -- "$path"
|
||||
done
|
||||
}
|
||||
trap cleanup_merge_temp_dirs EXIT
|
||||
trap 'exit 130' INT
|
||||
trap 'exit 143' TERM
|
||||
|
||||
fetch_gitea_pr_head() {
|
||||
local host="$1" auth_mode="$2" credential="$3" work_root="$4"
|
||||
local response_file raw_code api_url auth_config curl_rc
|
||||
response_file=$(mktemp "$work_root/pr-merge-pr.XXXXXX")
|
||||
api_url="https://${host}/api/v1/repos/${OWNER}/${REPO}/pulls/${PR_NUMBER}"
|
||||
if ! auth_config=$(write_curl_auth_config "$auth_mode" "$credential"); then
|
||||
echo "Error: Could not construct Gitea authentication config; refusing request." >&2
|
||||
rm -f "$response_file"
|
||||
return 1
|
||||
fi
|
||||
raw_code=$(curl -sS -K - "${GITEA_CURL_BOUNDS[@]}" -w '%{http_code}' -o "$response_file" \
|
||||
-H "User-Agent: curl/8" "$api_url" <<<"$auth_config")
|
||||
curl_rc=$?
|
||||
LAST_GITEA_HTTP_CODE="${raw_code:-000}"
|
||||
if [[ "$curl_rc" -ne 0 ]]; then
|
||||
LAST_GITEA_ERROR="curl transport failed (rc=$curl_rc)"
|
||||
rm -f "$response_file"
|
||||
return 1
|
||||
fi
|
||||
if [[ ! "$raw_code" =~ ^2 ]]; then
|
||||
LAST_GITEA_ERROR=$(format_gitea_error_response "$response_file")
|
||||
rm -f "$response_file"
|
||||
return 1
|
||||
fi
|
||||
if ! python3 - "$response_file" <<'PY'
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
|
||||
with open(sys.argv[1], encoding="utf-8") as handle:
|
||||
pull = json.load(handle)
|
||||
head = pull.get("head") if isinstance(pull, dict) else None
|
||||
sha = str(head.get("sha") or "") if isinstance(head, dict) else ""
|
||||
if not re.fullmatch(r"[0-9a-fA-F]{40}", sha):
|
||||
raise SystemExit(1)
|
||||
print(sha)
|
||||
PY
|
||||
then
|
||||
echo "Error: Gitea PR response has no valid head SHA; refusing merge." >&2
|
||||
rm -f "$response_file"
|
||||
return 1
|
||||
fi
|
||||
rm -f "$response_file"
|
||||
}
|
||||
|
||||
fetch_gitea_pr_commits() {
|
||||
local host="$1" auth_mode="$2" credential="$3" work_root="$4"
|
||||
local page page_file combined_file merged_file raw_code page_count api_url auth_config curl_rc
|
||||
mkdir -p "$work_root"
|
||||
if ! auth_config=$(write_curl_auth_config "$auth_mode" "$credential"); then
|
||||
echo "Error: Could not construct Gitea authentication config; refusing request." >&2
|
||||
return 1
|
||||
fi
|
||||
combined_file=$(mktemp "$work_root/pr-merge-commits.XXXXXX")
|
||||
printf '[]' > "$combined_file"
|
||||
|
||||
page=1
|
||||
while true; do
|
||||
page_file=$(mktemp "$work_root/pr-merge-commits-page.XXXXXX")
|
||||
api_url="https://${host}/api/v1/repos/${OWNER}/${REPO}/pulls/${PR_NUMBER}/commits?limit=50&page=${page}"
|
||||
raw_code=$(curl -sS -K - "${GITEA_CURL_BOUNDS[@]}" -w '%{http_code}' -o "$page_file" \
|
||||
-H "User-Agent: curl/8" "$api_url" <<<"$auth_config")
|
||||
curl_rc=$?
|
||||
LAST_GITEA_HTTP_CODE="${raw_code:-000}"
|
||||
if [[ "$curl_rc" -ne 0 ]]; then
|
||||
LAST_GITEA_ERROR="curl transport failed (rc=$curl_rc)"
|
||||
rm -f "$page_file" "$combined_file"
|
||||
return 1
|
||||
fi
|
||||
if [[ ! "$raw_code" =~ ^2 ]]; then
|
||||
LAST_GITEA_ERROR=$(format_gitea_error_response "$page_file")
|
||||
rm -f "$page_file" "$combined_file"
|
||||
return 1
|
||||
fi
|
||||
|
||||
if ! page_count=$(python3 - "$page_file" <<'PY'
|
||||
import json
|
||||
import sys
|
||||
|
||||
with open(sys.argv[1], encoding="utf-8") as handle:
|
||||
page = json.load(handle)
|
||||
if not isinstance(page, list):
|
||||
raise SystemExit(1)
|
||||
print(len(page))
|
||||
PY
|
||||
); then
|
||||
echo "Error: Gitea PR commits response is not a JSON array; refusing merge." >&2
|
||||
rm -f "$page_file" "$combined_file"
|
||||
return 1
|
||||
fi
|
||||
|
||||
merged_file=$(mktemp "$work_root/pr-merge-commits-merged.XXXXXX")
|
||||
if ! python3 - "$combined_file" "$page_file" > "$merged_file" <<'PY'
|
||||
import json
|
||||
import sys
|
||||
|
||||
with open(sys.argv[1], encoding="utf-8") as handle:
|
||||
combined = json.load(handle)
|
||||
with open(sys.argv[2], encoding="utf-8") as handle:
|
||||
page = json.load(handle)
|
||||
json.dump(combined + page, sys.stdout, separators=(",", ":"))
|
||||
PY
|
||||
then
|
||||
echo "Error: Could not combine paginated PR commit metadata; refusing merge." >&2
|
||||
rm -f "$page_file" "$combined_file" "$merged_file"
|
||||
return 1
|
||||
fi
|
||||
mv "$merged_file" "$combined_file"
|
||||
rm -f "$page_file"
|
||||
|
||||
if [[ "$page_count" -lt 50 ]]; then
|
||||
break
|
||||
fi
|
||||
page=$((page + 1))
|
||||
if [[ "$page" -gt 1000 ]]; then
|
||||
echo "Error: PR commit pagination exceeded 1000 pages; refusing merge." >&2
|
||||
rm -f "$combined_file"
|
||||
return 1
|
||||
fi
|
||||
done
|
||||
|
||||
cat "$combined_file"
|
||||
rm -f "$combined_file"
|
||||
}
|
||||
|
||||
# LIMITATION: author.login resolution proves the commit address maps to a registered account.
|
||||
# It does NOT prove the named principal authored the commit — git author metadata is self-asserted.
|
||||
# This gate checks ATTRIBUTION LINKAGE, not AUTHORSHIP. Commit signing is out of scope and unadopted.
|
||||
build_coauthor_message_fields() {
|
||||
local commits_file="$1" context_file="$2" head_file="$3"
|
||||
python3 - "$commits_file" "$context_file" "$head_file" <<'PY'
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
|
||||
commits_path, context_path, head_path = sys.argv[1:]
|
||||
with open(commits_path, encoding="utf-8") as handle:
|
||||
commits = json.load(handle)
|
||||
head_sha = open(head_path, encoding="utf-8").read().strip()
|
||||
context_parts = open(context_path, "rb").read().split(b"\0")
|
||||
if len(context_parts) != 4 or context_parts[-1] != b"":
|
||||
raise SystemExit(1)
|
||||
poster, title, principal = (part.decode("utf-8") for part in context_parts[:3])
|
||||
|
||||
if not isinstance(commits, list) or not commits:
|
||||
print(
|
||||
f"BLOCK: provider returned no PR commits; author identity is unmeasurable. "
|
||||
f"Refusing merge; escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
if not poster:
|
||||
print(
|
||||
f"BLOCK: PR poster login is empty; refusing merge; "
|
||||
f"escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
|
||||
if not re.fullmatch(r"[0-9a-fA-F]{40}", head_sha):
|
||||
print(
|
||||
f"BLOCK: inspected PR head SHA is invalid; refusing merge; "
|
||||
f"escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
|
||||
seen = set()
|
||||
trailers = []
|
||||
head_seen = False
|
||||
for item in commits:
|
||||
if not isinstance(item, dict):
|
||||
print(f"BLOCK: malformed PR commit metadata; escalate to named principal '{principal}'.", file=sys.stderr)
|
||||
raise SystemExit(75)
|
||||
sha = str(item.get("sha") or "<unknown>")
|
||||
if sha == head_sha:
|
||||
head_seen = True
|
||||
commit = item.get("commit") if isinstance(item.get("commit"), dict) else {}
|
||||
commit_author = commit.get("author") if isinstance(commit.get("author"), dict) else {}
|
||||
email = str(commit_author.get("email") or "").strip()
|
||||
provider_author = item.get("author") if isinstance(item.get("author"), dict) else {}
|
||||
login = str(provider_author.get("login") or "").strip()
|
||||
|
||||
if not login:
|
||||
diagnostic_email = email or "<missing>"
|
||||
print(
|
||||
f"BLOCK: commit {sha!r} has author.login=NULL while "
|
||||
f"commit.author.email={diagnostic_email!r}; refusing merge; "
|
||||
f"escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
if (
|
||||
not email.isascii()
|
||||
or not email.isprintable()
|
||||
or not re.fullmatch(r"[A-Za-z0-9_.-]+", login)
|
||||
or not re.fullmatch(r"[^<>\s]+@[^<>\s]+", email)
|
||||
):
|
||||
print(
|
||||
f"BLOCK: commit {sha!r} has unusable linked identity "
|
||||
f"author.login={login!r}, commit.author.email={email!r}; refusing merge; "
|
||||
f"escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
if login == poster or login in seen:
|
||||
continue
|
||||
seen.add(login)
|
||||
trailers.append(f"Co-authored-by: {login} <{email}>")
|
||||
|
||||
if not head_seen:
|
||||
print(
|
||||
f"BLOCK: inspected PR head is absent from commit enumeration; refusing merge; "
|
||||
f"escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
if not trailers:
|
||||
print("{}")
|
||||
raise SystemExit(0)
|
||||
if not title:
|
||||
print(
|
||||
f"BLOCK: PR title is empty; refusing merge; escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
if not title.isprintable() or re.match(r"^[A-Za-z-]+-[Bb]y:", title):
|
||||
print(
|
||||
f"BLOCK: PR title is not one printable, non-trailer line; refusing merge; "
|
||||
f"escalate to named principal '{principal}'.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
raise SystemExit(75)
|
||||
|
||||
print(json.dumps({
|
||||
"MergeTitleField": title,
|
||||
"MergeMessageField": "\n".join(trailers),
|
||||
}, separators=(",", ":")))
|
||||
PY
|
||||
}
|
||||
|
||||
merge_gitea_api_attempt() {
|
||||
local host="$1" auth_mode="$2" credential="$3"
|
||||
local api_url attempt_dir body_file raw_code commits_file fields_file context_file head_file payload_file work_root attempt_rc auth_config curl_rc
|
||||
LAST_GITEA_HTTP_CODE="000"
|
||||
LAST_GITEA_ERROR=""
|
||||
merge_gitea_with_api() {
|
||||
local host="$1" api_url token basic_auth body_file raw_code payload
|
||||
api_url="https://${host}/api/v1/repos/${OWNER}/${REPO}/pulls/${PR_NUMBER}/merge"
|
||||
work_root="${AGENT_WORK_ROOT:-${HOME:-/tmp}/mosaic/agent-work}"
|
||||
mkdir -p "$work_root"
|
||||
attempt_dir=$(mktemp -d "$work_root/pr-merge-attempt.XXXXXX")
|
||||
chmod 0700 "$attempt_dir"
|
||||
MERGE_TEMP_DIRS+=("$attempt_dir")
|
||||
body_file=$(mktemp "$attempt_dir/api-response.XXXXXX")
|
||||
fields_file=$(mktemp "$attempt_dir/message-fields.XXXXXX")
|
||||
payload_file=$(mktemp "$attempt_dir/payload.XXXXXX")
|
||||
printf '{}' > "$fields_file"
|
||||
|
||||
if [[ "$CO_AUTHOR_TRAILERS" == true ]]; then
|
||||
commits_file=$(mktemp "$attempt_dir/pr-merge-commits-input.XXXXXX")
|
||||
context_file=$(mktemp "$attempt_dir/pr-merge-message-context.XXXXXX")
|
||||
head_file=$(mktemp "$attempt_dir/pr-merge-head-input.XXXXXX")
|
||||
printf '%s\0%s\0%s\0' "$PR_AUTHOR" "$PR_TITLE" "$ESCALATE_TO" > "$context_file"
|
||||
if fetch_gitea_pr_head "$host" "$auth_mode" "$credential" "$attempt_dir" > "$head_file"; then
|
||||
:
|
||||
else
|
||||
attempt_rc=$?
|
||||
rm -f "$body_file" "$fields_file" "$payload_file" "$commits_file" "$context_file" "$head_file"
|
||||
return "$attempt_rc"
|
||||
fi
|
||||
if [[ "$(<"$head_file")" != "$HEAD_SHA" ]]; then
|
||||
echo "BLOCK: authenticated PR head moved from reviewed $HEAD_SHA to $(<"$head_file"); refusing merge; escalate to named principal '$ESCALATE_TO'." >&2
|
||||
rm -f "$body_file" "$fields_file" "$payload_file" "$commits_file" "$context_file" "$head_file"
|
||||
return 75
|
||||
fi
|
||||
if fetch_gitea_pr_commits "$host" "$auth_mode" "$credential" "$attempt_dir" > "$commits_file"; then
|
||||
:
|
||||
else
|
||||
attempt_rc=$?
|
||||
rm -f "$body_file" "$fields_file" "$payload_file" "$commits_file" "$context_file" "$head_file"
|
||||
return "$attempt_rc"
|
||||
fi
|
||||
if build_coauthor_message_fields "$commits_file" "$context_file" "$head_file" > "$fields_file"; then
|
||||
:
|
||||
else
|
||||
attempt_rc=$?
|
||||
rm -f "$body_file" "$fields_file" "$payload_file" "$commits_file" "$context_file" "$head_file"
|
||||
return "$attempt_rc"
|
||||
fi
|
||||
rm -f "$commits_file" "$context_file" "$head_file"
|
||||
fi
|
||||
|
||||
if ! python3 - "$fields_file" "$HEAD_SHA" "$DELETE_BRANCH" > "$payload_file" <<'PY'
|
||||
mkdir -p "${AGENT_WORK_ROOT:-${HOME:-/tmp}/mosaic/agent-work}"
|
||||
body_file=$(mktemp "${AGENT_WORK_ROOT:-${HOME:-/tmp}/mosaic/agent-work}/pr-merge-api-response.XXXXXX")
|
||||
payload=$(python3 - "$HEAD_SHA" "$DELETE_BRANCH" <<'PY'
|
||||
import json
|
||||
import sys
|
||||
|
||||
with open(sys.argv[1], encoding="utf-8") as handle:
|
||||
fields = json.load(handle)
|
||||
head_sha, delete_branch = sys.argv[2:]
|
||||
head_sha, delete_branch = sys.argv[1:]
|
||||
payload = {"Do": "squash", "head_commit_id": head_sha}
|
||||
if delete_branch == "true":
|
||||
payload["delete_branch_after_merge"] = True
|
||||
payload.update(fields)
|
||||
allowed = {"Do", "head_commit_id", "delete_branch_after_merge", "MergeTitleField", "MergeMessageField"}
|
||||
if payload.get("Do") != "squash" or set(payload) - allowed:
|
||||
raise SystemExit(1)
|
||||
print(json.dumps(payload, separators=(",", ":")))
|
||||
PY
|
||||
then
|
||||
rm -f "$body_file" "$fields_file" "$payload_file"
|
||||
return 1
|
||||
fi
|
||||
rm -f "$fields_file"
|
||||
)
|
||||
|
||||
if ! auth_config=$(write_curl_auth_config "$auth_mode" "$credential"); then
|
||||
echo "Error: Could not construct Gitea authentication config; refusing request." >&2
|
||||
rm -f "$body_file" "$payload_file"
|
||||
return 1
|
||||
token=$(get_gitea_token "$host" || true)
|
||||
if [[ -n "$token" ]]; then
|
||||
raw_code=$(curl -sS -w '%{http_code}' -o "$body_file" \
|
||||
-X POST \
|
||||
-H "User-Agent: curl/8" \
|
||||
-H "Authorization: token $token" \
|
||||
-H 'Content-Type: application/json' \
|
||||
-d "$payload" \
|
||||
"$api_url" || true)
|
||||
if [[ "$raw_code" =~ ^2 ]]; then
|
||||
rm -f "$body_file"
|
||||
return 0
|
||||
fi
|
||||
fi
|
||||
raw_code=$(curl -sS -K - "${GITEA_CURL_BOUNDS[@]}" -w '%{http_code}' -o "$body_file" \
|
||||
-X POST -H "User-Agent: curl/8" \
|
||||
-H 'Content-Type: application/json' \
|
||||
--data-binary "@$payload_file" "$api_url" <<<"$auth_config")
|
||||
curl_rc=$?
|
||||
LAST_GITEA_HTTP_CODE="${raw_code:-000}"
|
||||
if [[ "$curl_rc" -ne 0 ]]; then
|
||||
LAST_GITEA_ERROR="curl transport failed (rc=$curl_rc)"
|
||||
rm -f "$body_file" "$payload_file"
|
||||
rm -rf -- "$attempt_dir"
|
||||
return 1
|
||||
fi
|
||||
if [[ ! "$raw_code" =~ ^2 ]]; then
|
||||
LAST_GITEA_ERROR=$(format_gitea_error_response "$body_file")
|
||||
fi
|
||||
rm -f "$body_file" "$payload_file"
|
||||
rm -rf -- "$attempt_dir"
|
||||
[[ "$raw_code" =~ ^2 ]]
|
||||
}
|
||||
|
||||
merge_gitea_with_api() {
|
||||
local host="$1" token attempt_rc
|
||||
basic_auth=$(get_gitea_basic_auth "$host" || true)
|
||||
if [[ -n "$basic_auth" ]]; then
|
||||
raw_code=$(curl -sS -w '%{http_code}' -o "$body_file" \
|
||||
-X POST \
|
||||
-u "$basic_auth" \
|
||||
-H "User-Agent: curl/8" \
|
||||
-H 'Content-Type: application/json' \
|
||||
-d "$payload" \
|
||||
"$api_url" || true)
|
||||
if [[ "$raw_code" =~ ^2 ]]; then
|
||||
rm -f "$body_file"
|
||||
return 0
|
||||
fi
|
||||
fi
|
||||
|
||||
if ! token=$(get_gitea_token "$host"); then
|
||||
echo "Error: Could not resolve the required Gitea token; refusing merge without changing principals." >&2
|
||||
return 1
|
||||
fi
|
||||
if [[ -z "$token" ]]; then
|
||||
echo "Error: Required Gitea token resolved empty; refusing merge without changing principals." >&2
|
||||
return 1
|
||||
fi
|
||||
if merge_gitea_api_attempt "$host" token "$token"; then
|
||||
return 0
|
||||
else
|
||||
attempt_rc=$?
|
||||
fi
|
||||
if [[ "$attempt_rc" -eq 75 ]]; then
|
||||
return 75
|
||||
fi
|
||||
if [[ "$LAST_GITEA_HTTP_CODE" != "401" ]]; then
|
||||
echo "Error: Gitea API merge failed with the identity-bound token (HTTP ${LAST_GITEA_HTTP_CODE:-000}).${LAST_GITEA_ERROR:+ Provider response: $LAST_GITEA_ERROR}" >&2
|
||||
return 1
|
||||
fi
|
||||
echo "Error: Gitea API rejected the identity-bound token with HTTP 401; refusing cross-principal credential fallback." >&2
|
||||
python3 - "${raw_code:-000}" "$body_file" <<'PY' >&2
|
||||
import json
|
||||
import sys
|
||||
code, path = sys.argv[1], sys.argv[2]
|
||||
try:
|
||||
with open(path, encoding="utf-8", errors="replace") as handle:
|
||||
raw = handle.read(500)
|
||||
data = json.loads(raw) if raw else {}
|
||||
message = data.get("message") or data.get("error") or raw or "empty response"
|
||||
except Exception:
|
||||
try:
|
||||
message = open(path, encoding="utf-8", errors="replace").read(500) or "empty response"
|
||||
except Exception:
|
||||
message = "unreadable response"
|
||||
print(f"Error: Gitea API merge failed with HTTP {code}: {message}")
|
||||
PY
|
||||
rm -f "$body_file"
|
||||
return 1
|
||||
}
|
||||
|
||||
@@ -602,10 +195,11 @@ if [[ "$DRY_RUN" == true ]]; then
|
||||
echo "Error: Cannot determine host from origin remote URL" >&2
|
||||
exit 1
|
||||
}
|
||||
if [[ "$CO_AUTHOR_TRAILERS" == true ]]; then
|
||||
echo "Dry run: would verify PR commit authors and merge PR #$PR_NUMBER on $HOST with authenticated Gitea API message fields (base=$BASE_BRANCH, method=squash)."
|
||||
TEA_LOGIN="$(get_gitea_login_for_host "$HOST" || true)"
|
||||
if [[ -n "$TEA_LOGIN" ]]; then
|
||||
echo "Dry run: would merge PR #$PR_NUMBER on $HOST with tea login '$TEA_LOGIN' (base=$BASE_BRANCH, method=squash)."
|
||||
else
|
||||
echo "Dry run: would merge PR #$PR_NUMBER on $HOST with the authenticated exact-head Gitea API path (base=$BASE_BRANCH, method=squash)."
|
||||
echo "Dry run: would merge PR #$PR_NUMBER on $HOST with authenticated Gitea API fallback (base=$BASE_BRANCH, method=squash)."
|
||||
fi
|
||||
else
|
||||
echo "Dry run: would merge PR #$PR_NUMBER on $PLATFORM (base=$BASE_BRANCH, method=squash)."
|
||||
@@ -615,10 +209,6 @@ fi
|
||||
|
||||
case "$PLATFORM" in
|
||||
github)
|
||||
if [[ "$CO_AUTHOR_TRAILERS" == true ]]; then
|
||||
echo "Error: --co-author-trailers currently requires the Gitea REST message-field contract." >&2
|
||||
exit 1
|
||||
fi
|
||||
cmd=(gh pr merge "$PR_NUMBER" --squash --match-head-commit "$HEAD_SHA")
|
||||
[[ "$DELETE_BRANCH" == true ]] && cmd+=(--delete-branch)
|
||||
"${cmd[@]}"
|
||||
@@ -629,7 +219,7 @@ case "$PLATFORM" in
|
||||
exit 1
|
||||
}
|
||||
# Gitea's API head_commit_id is an atomic compare-and-merge precondition.
|
||||
# tea cannot express it, so every Gitea merge uses the authenticated API path.
|
||||
# tea cannot express it, so exact-head merges use the authenticated API path.
|
||||
merge_gitea_with_api "$HOST"
|
||||
;;
|
||||
*)
|
||||
|
||||
@@ -51,23 +51,22 @@ for arg in "$@"; do
|
||||
prev=""
|
||||
continue
|
||||
fi
|
||||
if [[ "$prev" == "data" ]]; then
|
||||
if [[ "$prev" == "-d" ]]; then
|
||||
post_data="$arg"
|
||||
[[ "$post_data" == @* ]] && post_data=$(<"${post_data#@}")
|
||||
prev=""
|
||||
continue
|
||||
fi
|
||||
if [[ "$prev" == "config" ]]; then
|
||||
[[ "$arg" == "-" ]] && cat >/dev/null
|
||||
prev=""
|
||||
if [[ "$arg" == "-o" ]]; then
|
||||
prev="-o"
|
||||
continue
|
||||
fi
|
||||
case "$arg" in
|
||||
-o) prev="-o" ;;
|
||||
-d|--data|--data-binary) prev="data" ;;
|
||||
-K|--config) prev="config" ;;
|
||||
-w) write_code=true ;;
|
||||
esac
|
||||
if [[ "$arg" == "-d" ]]; then
|
||||
prev="-d"
|
||||
continue
|
||||
fi
|
||||
if [[ "$arg" == "-w" ]]; then
|
||||
write_code=true
|
||||
fi
|
||||
done
|
||||
emit_response() {
|
||||
local body="$1"
|
||||
|
||||
@@ -36,30 +36,13 @@ cat > "$WORK_DIR/gitea/curl" <<'SH'
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
payload=""
|
||||
out_file=""
|
||||
while [[ $# -gt 0 ]]; do
|
||||
case "$1" in
|
||||
-d|--data|--data-binary)
|
||||
payload="$2"
|
||||
[[ "$payload" == @* ]] && payload=$(<"${payload#@}")
|
||||
shift 2
|
||||
;;
|
||||
-o)
|
||||
out_file="$2"
|
||||
shift 2
|
||||
;;
|
||||
-K|--config)
|
||||
[[ "$2" == "-" ]] && cat >/dev/null
|
||||
shift 2
|
||||
;;
|
||||
-w|-X|-H)
|
||||
shift 2
|
||||
;;
|
||||
*) shift ;;
|
||||
esac
|
||||
for ((i=1; i<=$#; i++)); do
|
||||
if [[ "${!i}" == "-d" ]]; then
|
||||
j=$((i + 1))
|
||||
payload="${!j}"
|
||||
fi
|
||||
done
|
||||
printf '%s' "$payload" > "${MOSAIC_MERGE_PAYLOAD_LOG:?}"
|
||||
[[ -n "$out_file" ]] && printf '{}' > "$out_file"
|
||||
printf '200'
|
||||
SH
|
||||
chmod +x "$WORK_DIR/gitea/curl"
|
||||
|
||||
@@ -1,541 +0,0 @@
|
||||
#!/usr/bin/env bash
|
||||
# Regression harness for the optional, identity-checked Gitea squash message.
|
||||
|
||||
set -u
|
||||
|
||||
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
||||
SUBJECT="${MOSAIC_TEST_SUBJECT:-$SCRIPT_DIR/pr-merge.sh}"
|
||||
WORK_DIR="${MOSAIC_TEST_WORK_DIR:-$PWD/.mosaic-test-work/pr-merge-message-field}"
|
||||
ORIG_PATH="$PATH"
|
||||
failures=0
|
||||
|
||||
rm -rf "$WORK_DIR"
|
||||
mkdir -p "$WORK_DIR"
|
||||
|
||||
fail() {
|
||||
echo "FAIL $1" >&2
|
||||
failures=$((failures + 1))
|
||||
}
|
||||
|
||||
make_case() {
|
||||
local name="$1" case_dir
|
||||
case_dir="$WORK_DIR/$name"
|
||||
mkdir -p "$case_dir/bin" "$case_dir/agent"
|
||||
cp "$SUBJECT" "$case_dir/pr-merge.sh"
|
||||
chmod +x "$case_dir/pr-merge.sh"
|
||||
|
||||
cat > "$case_dir/detect-platform.sh" <<'SH'
|
||||
#!/usr/bin/env bash
|
||||
detect_platform() { PLATFORM=gitea; printf 'gitea\n'; }
|
||||
get_repo_owner() { printf 'acme\n'; }
|
||||
get_repo_name() { printf 'widgets\n'; }
|
||||
get_remote_host() { printf 'git.example.test\n'; }
|
||||
get_gitea_token() {
|
||||
printf 'resolved\n' >> "${MOSAIC_TEST_TOKEN_RESOLUTION_LOG:?}"
|
||||
if [[ "${MOSAIC_TEST_TOKEN_AVAILABLE:-true}" != "true" ]]; then
|
||||
return 1
|
||||
fi
|
||||
printf 'fixture-token\n'
|
||||
}
|
||||
get_gitea_basic_auth() {
|
||||
printf 'resolved\n' >> "${MOSAIC_TEST_BASIC_RESOLUTION_LOG:?}"
|
||||
if [[ "${MOSAIC_TEST_BASIC_AVAILABLE:-false}" == "true" ]]; then
|
||||
printf 'fixture-user:fixture-password\n'
|
||||
return "${MOSAIC_TEST_BASIC_RC:-0}"
|
||||
fi
|
||||
return 1
|
||||
}
|
||||
get_gitea_login_for_host() { return 1; }
|
||||
SH
|
||||
|
||||
cat > "$case_dir/pr-metadata.sh" <<'SH'
|
||||
#!/usr/bin/env bash
|
||||
if [[ "${MOSAIC_TEST_TITLE_MODE:-safe}" == "injection" ]]; then
|
||||
title='Preserve authors\n\nCo-authored-by: victim <[email protected]>'
|
||||
else
|
||||
title='Preserve both branch authors'
|
||||
fi
|
||||
case "${MOSAIC_TEST_COMMITS_MODE:?}" in
|
||||
verified) head_sha=2222222222222222222222222222222222222222 ;;
|
||||
null-login|unsafe-identity) head_sha=3333333333333333333333333333333333333333 ;;
|
||||
single) head_sha=1111111111111111111111111111111111111111 ;;
|
||||
*) echo "unknown commits mode" >&2; exit 2 ;;
|
||||
esac
|
||||
printf '{"number":42,"title":"%s","author":"poster","baseRefName":"main","headRefName":"feature/fixture","headRefOid":"%s","headRepository":"acme/widgets"}\n' "$title" "$head_sha"
|
||||
SH
|
||||
|
||||
cat > "$case_dir/ci-queue-wait.sh" <<'SH'
|
||||
#!/usr/bin/env bash
|
||||
exit 0
|
||||
SH
|
||||
|
||||
cat > "$case_dir/bin/python3" <<'SH'
|
||||
#!/usr/bin/env bash
|
||||
for arg in "$@"; do
|
||||
case "$arg" in
|
||||
*"Preserve both branch authors"*|*"[email protected]"*)
|
||||
: > "${MOSAIC_TEST_METADATA_ARGV_MARKER:?}"
|
||||
;;
|
||||
esac
|
||||
done
|
||||
exec "${MOSAIC_TEST_REAL_PYTHON:?}" "$@"
|
||||
SH
|
||||
|
||||
cat > "$case_dir/bin/curl" <<'SH'
|
||||
#!/usr/bin/env bash
|
||||
set -eu
|
||||
|
||||
for arg in "$@"; do
|
||||
case "$arg" in
|
||||
*"Preserve both branch authors"*|*"[email protected]"*)
|
||||
: > "${MOSAIC_TEST_METADATA_ARGV_MARKER:?}"
|
||||
;;
|
||||
esac
|
||||
done
|
||||
|
||||
url=""
|
||||
method="GET"
|
||||
out_file=""
|
||||
data=""
|
||||
config=""
|
||||
auth_mode="none"
|
||||
has_max_filesize=0
|
||||
has_max_time=0
|
||||
has_connect_timeout=0
|
||||
while [[ $# -gt 0 ]]; do
|
||||
case "$1" in
|
||||
-o)
|
||||
out_file="$2"
|
||||
shift 2
|
||||
;;
|
||||
-w)
|
||||
shift 2
|
||||
;;
|
||||
-X)
|
||||
method="$2"
|
||||
shift 2
|
||||
;;
|
||||
-d|--data|--data-binary)
|
||||
data="$2"
|
||||
if [[ "$data" == @* ]]; then
|
||||
data=$(<"${data#@}")
|
||||
fi
|
||||
shift 2
|
||||
;;
|
||||
-K|--config)
|
||||
if [[ "$2" == "-" ]]; then
|
||||
config=$(cat)
|
||||
fi
|
||||
shift 2
|
||||
;;
|
||||
--max-filesize)
|
||||
has_max_filesize=1
|
||||
shift 2
|
||||
;;
|
||||
--max-time)
|
||||
has_max_time=1
|
||||
shift 2
|
||||
;;
|
||||
--connect-timeout)
|
||||
has_connect_timeout=1
|
||||
shift 2
|
||||
;;
|
||||
-H|--header|-u|--user)
|
||||
if [[ "$2" == *"fixture-token"* ]]; then
|
||||
: > "${MOSAIC_TEST_TOKEN_ARGV_MARKER:?}"
|
||||
fi
|
||||
if [[ "$2" == *"fixture-password"* ]]; then
|
||||
: > "${MOSAIC_TEST_BASIC_ARGV_MARKER:?}"
|
||||
fi
|
||||
shift 2
|
||||
;;
|
||||
http://*|https://*)
|
||||
url="$1"
|
||||
shift
|
||||
;;
|
||||
*)
|
||||
shift
|
||||
;;
|
||||
esac
|
||||
done
|
||||
|
||||
if [[ "$config" == *"Authorization: token fixture-token"* ]]; then
|
||||
auth_mode="token"
|
||||
: > "${MOSAIC_TEST_AUTH_CONFIG_MARKER:?}"
|
||||
elif [[ "$config" == *"user = \"fixture-user:fixture-password\""* ]]; then
|
||||
auth_mode="basic"
|
||||
: > "${MOSAIC_TEST_BASIC_CONFIG_MARKER:?}"
|
||||
fi
|
||||
printf '%s %s %s\n' "$method" "$auth_mode" "$url" >> "${MOSAIC_TEST_CURL_LOG:?}"
|
||||
printf '%s:%s:%s\n' "$has_max_filesize" "$has_max_time" "$has_connect_timeout" >> "${MOSAIC_TEST_CURL_BOUNDS_LOG:?}"
|
||||
|
||||
case "$url" in
|
||||
*/pulls/42)
|
||||
case "${MOSAIC_TEST_COMMITS_MODE:?}" in
|
||||
verified) head_sha=2222222222222222222222222222222222222222 ;;
|
||||
null-login|unsafe-identity) head_sha=3333333333333333333333333333333333333333 ;;
|
||||
single) head_sha=1111111111111111111111111111111111111111 ;;
|
||||
*) echo "unknown commits mode" >&2; exit 2 ;;
|
||||
esac
|
||||
if [[ "${MOSAIC_TEST_HEAD_MODE:-stable}" == "moved" ]]; then
|
||||
head_sha=4444444444444444444444444444444444444444
|
||||
fi
|
||||
body="{\"head\":{\"sha\":\"$head_sha\"}}"
|
||||
code=200
|
||||
if [[ "${MOSAIC_TEST_FALLBACK_MODE:-none}" == "inspection" && "$auth_mode" == "token" ]]; then
|
||||
body='{"message":"token rejected"}'
|
||||
code=401
|
||||
fi
|
||||
;;
|
||||
*/pulls/42/commits*)
|
||||
case "${MOSAIC_TEST_COMMITS_MODE:?}" in
|
||||
verified)
|
||||
if [[ "${MOSAIC_TEST_EMAIL_MODE:-safe}" == "escape" ]]; then
|
||||
body='[{"sha":"2222222222222222222222222222222222222222","commit":{"author":{"name":"Alice","email":"alice+\u001b[[email protected]"}},"author":{"login":"alice"}},{"sha":"1111111111111111111111111111111111111111","commit":{"author":{"name":"Poster","email":"[email protected]"}},"author":{"login":"poster"}}]'
|
||||
else
|
||||
body='[{"sha":"2222222222222222222222222222222222222222","commit":{"author":{"name":"Alice","email":"[email protected]"}},"author":{"login":"alice"}},{"sha":"1111111111111111111111111111111111111111","commit":{"author":{"name":"Poster","email":"[email protected]"}},"author":{"login":"poster"}}]'
|
||||
fi
|
||||
;;
|
||||
null-login)
|
||||
body='[{"sha":"1111111111111111111111111111111111111111","commit":{"author":{"name":"Poster","email":"[email protected]"}},"author":{"login":"poster"}},{"sha":"3333333333333333333333333333333333333333","commit":{"author":{"name":"Unresolved Author","email":"[email protected]\n\u001b[31m"}},"author":null}]'
|
||||
;;
|
||||
unsafe-identity)
|
||||
body='[{"sha":"unsafe\n\u001b[31m","commit":{"author":{"name":"Unsafe","email":"not-an-email"}},"author":{"login":"unsafe"}},{"sha":"3333333333333333333333333333333333333333","commit":{"author":{"name":"Poster","email":"[email protected]"}},"author":{"login":"poster"}}]'
|
||||
;;
|
||||
single)
|
||||
body='[{"sha":"1111111111111111111111111111111111111111","commit":{"author":{"name":"Poster","email":"[email protected]"}},"author":{"login":"poster"}}]'
|
||||
;;
|
||||
*)
|
||||
echo "unknown commits mode" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
code=200
|
||||
if [[ "${MOSAIC_TEST_FALLBACK_MODE:-none}" == "inspection" && "$auth_mode" == "token" ]]; then
|
||||
body='{"message":"token rejected"}'
|
||||
code=401
|
||||
fi
|
||||
;;
|
||||
*/pulls/42/merge)
|
||||
body='{}'
|
||||
code=200
|
||||
if [[ "${MOSAIC_TEST_FALLBACK_MODE:-none}" == "merge" && "$auth_mode" == "token" ]]; then
|
||||
body='{"message":"token rejected"}'
|
||||
code=401
|
||||
elif [[ "${MOSAIC_TEST_FALLBACK_MODE:-none}" == "provider-error" ]]; then
|
||||
body='{"message":"branch policy rejected\n\u001b[31m"}'
|
||||
code=409
|
||||
elif [[ "${MOSAIC_TEST_FALLBACK_MODE:-none}" == "forbidden" ]]; then
|
||||
body='{"message":"permission denied"}'
|
||||
code=403
|
||||
else
|
||||
printf '%s' "$data" > "${MOSAIC_TEST_MERGE_PAYLOAD:?}"
|
||||
fi
|
||||
;;
|
||||
*/users/*)
|
||||
body='{"message":"not found"}'
|
||||
code=404
|
||||
;;
|
||||
*)
|
||||
body='{"message":"unexpected URL"}'
|
||||
code=500
|
||||
;;
|
||||
esac
|
||||
|
||||
if [[ -n "$out_file" ]]; then
|
||||
printf '%s' "$body" > "$out_file"
|
||||
else
|
||||
printf '%s' "$body"
|
||||
fi
|
||||
printf '%s' "$code"
|
||||
case "${MOSAIC_TEST_CURL_FAILURE:-none}" in
|
||||
oversize) exit 63 ;;
|
||||
stalled) exit 28 ;;
|
||||
esac
|
||||
SH
|
||||
|
||||
chmod +x "$case_dir/detect-platform.sh" "$case_dir/pr-metadata.sh" \
|
||||
"$case_dir/ci-queue-wait.sh" "$case_dir/bin/curl" "$case_dir/bin/python3"
|
||||
printf '%s\n' "$case_dir"
|
||||
}
|
||||
|
||||
run_case() {
|
||||
local case_dir="$1" mode="$2"
|
||||
shift 2
|
||||
MOSAIC_TEST_COMMITS_MODE="$mode" \
|
||||
MOSAIC_TEST_CURL_LOG="$case_dir/curl.log" \
|
||||
MOSAIC_TEST_CURL_BOUNDS_LOG="$case_dir/curl-bounds.log" \
|
||||
MOSAIC_TEST_MERGE_PAYLOAD="$case_dir/merge-payload.json" \
|
||||
MOSAIC_TEST_TOKEN_ARGV_MARKER="$case_dir/token-in-argv" \
|
||||
MOSAIC_TEST_BASIC_ARGV_MARKER="$case_dir/basic-in-argv" \
|
||||
MOSAIC_TEST_AUTH_CONFIG_MARKER="$case_dir/auth-via-config" \
|
||||
MOSAIC_TEST_BASIC_CONFIG_MARKER="$case_dir/basic-via-config" \
|
||||
MOSAIC_TEST_TOKEN_RESOLUTION_LOG="$case_dir/token-resolution.log" \
|
||||
MOSAIC_TEST_BASIC_RESOLUTION_LOG="$case_dir/basic-resolution.log" \
|
||||
MOSAIC_TEST_METADATA_ARGV_MARKER="$case_dir/metadata-in-argv" \
|
||||
MOSAIC_TEST_REAL_PYTHON="$(command -v python3)" \
|
||||
AGENT_WORK_ROOT="$case_dir/agent" \
|
||||
PATH="$case_dir/bin:$ORIG_PATH" \
|
||||
"$case_dir/pr-merge.sh" -n 42 "$@"
|
||||
}
|
||||
|
||||
# Verified multi-author path: the non-poster trailer is built from one commit's
|
||||
# linked author.login and that same commit's author email. No /users lookup.
|
||||
verified_dir=$(make_case verified)
|
||||
set +e
|
||||
verified_output=$(run_case "$verified_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
verified_rc=$?
|
||||
set -e
|
||||
if [[ "$verified_rc" -ne 0 ]]; then
|
||||
fail "verified multi-author merge expected rc=0, got rc=$verified_rc: $verified_output"
|
||||
elif [[ ! -s "$verified_dir/merge-payload.json" ]]; then
|
||||
fail "verified multi-author merge did not reach the API payload"
|
||||
else
|
||||
python3 - "$verified_dir/merge-payload.json" <<'PY' || fail "verified payload did not preserve squash and exact message fields"
|
||||
import json
|
||||
import sys
|
||||
payload = json.load(open(sys.argv[1], encoding="utf-8"))
|
||||
assert payload == {
|
||||
"Do": "squash",
|
||||
"head_commit_id": "2222222222222222222222222222222222222222",
|
||||
"MergeTitleField": "Preserve both branch authors",
|
||||
"MergeMessageField": "Co-authored-by: alice <[email protected]>",
|
||||
}, payload
|
||||
PY
|
||||
fi
|
||||
[[ -e "$verified_dir/auth-via-config" ]] || fail "verified path did not authenticate curl through stdin config"
|
||||
[[ ! -e "$verified_dir/token-in-argv" ]] || fail "verified path placed the Gitea token in curl argv"
|
||||
[[ ! -e "$verified_dir/metadata-in-argv" ]] || fail "verified path placed PR title or contributor email in child argv"
|
||||
[[ "$(wc -l < "$verified_dir/token-resolution.log")" -eq 1 ]] || fail "verified path did not bind inspection and merge to one credential resolution"
|
||||
if grep -q '/users/' "$verified_dir/curl.log" 2>/dev/null; then
|
||||
fail "verified path performed a forbidden second /users lookup"
|
||||
fi
|
||||
if grep -qv '^1:1:1$' "$verified_dir/curl-bounds.log"; then
|
||||
fail "verified path did not apply size/max-time/connect-time bounds to every provider download"
|
||||
fi
|
||||
|
||||
# A linked email containing a terminal escape must block before mutation.
|
||||
escape_email_dir=$(make_case escape-email)
|
||||
set +e
|
||||
escape_email_output=$(MOSAIC_TEST_EMAIL_MODE=escape run_case "$escape_email_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
escape_email_rc=$?
|
||||
set -e
|
||||
[[ "$escape_email_rc" -ne 0 ]] || fail "control-byte email unexpectedly passed"
|
||||
[[ "$escape_email_output" == *"unusable linked identity"* ]] || fail "control-byte email refusal lost its diagnostic"
|
||||
[[ ! -e "$escape_email_dir/merge-payload.json" ]] || fail "control-byte email reached the merge API"
|
||||
|
||||
# Curl transfer and duration failures must remain failures even with HTTP 200.
|
||||
for failure_mode in oversize stalled; do
|
||||
failure_dir=$(make_case "curl-$failure_mode")
|
||||
set +e
|
||||
failure_output=$(MOSAIC_TEST_CURL_FAILURE="$failure_mode" run_case "$failure_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
failure_rc=$?
|
||||
set -e
|
||||
[[ "$failure_rc" -ne 0 ]] || fail "curl $failure_mode failure was discarded: $failure_output"
|
||||
[[ ! -e "$failure_dir/merge-payload.json" ]] || fail "curl $failure_mode failure reached the merge API"
|
||||
done
|
||||
|
||||
# The authenticated head is re-read under the mutation credential but cannot
|
||||
# replace the canonical preflight/review head. A move blocks before enumeration
|
||||
# or mutation even though the provider returned a valid new SHA.
|
||||
moved_dir=$(make_case moved-head)
|
||||
set +e
|
||||
moved_output=$(MOSAIC_TEST_HEAD_MODE=moved \
|
||||
run_case "$moved_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
moved_rc=$?
|
||||
set -e
|
||||
[[ "$moved_rc" -ne 0 ]] || fail "moved authenticated head unexpectedly passed"
|
||||
[[ "$moved_output" == *"authenticated PR head moved from reviewed"* ]] || fail "moved head refusal lost its diagnostic"
|
||||
[[ "$moved_output" == *"tl-mosaic"* ]] || fail "moved head refusal omitted the named escalation principal"
|
||||
[[ ! -e "$moved_dir/merge-payload.json" ]] || fail "moved head refusal reached the merge API"
|
||||
moved_sequence=$(awk '{print $1 ":" $2}' "$moved_dir/curl.log" | paste -sd, -)
|
||||
[[ "$moved_sequence" == "GET:token" ]] || fail "moved head refusal performed post-move inspection/mutation (calls=$moved_sequence)"
|
||||
|
||||
# Token resolution failure is not an authentication response. It must fail
|
||||
# closed instead of borrowing a Basic credential under a different principal.
|
||||
token_missing_dir=$(make_case token-missing)
|
||||
set +e
|
||||
token_missing_output=$(MOSAIC_TEST_TOKEN_AVAILABLE=false MOSAIC_TEST_BASIC_AVAILABLE=true \
|
||||
run_case "$token_missing_dir" single 2>&1)
|
||||
token_missing_rc=$?
|
||||
set -e
|
||||
[[ "$token_missing_rc" -ne 0 ]] || fail "missing token unexpectedly borrowed Basic Auth"
|
||||
[[ "$token_missing_output" == *"required Gitea token"* ]] || fail "missing token refusal lost its diagnostic"
|
||||
[[ ! -e "$token_missing_dir/basic-resolution.log" ]] || fail "missing token resolved Basic Auth after identity failure"
|
||||
[[ ! -e "$token_missing_dir/curl.log" ]] || fail "missing token reached a provider request"
|
||||
|
||||
# A failed Basic resolver must never use its nonempty output or reach mutation.
|
||||
basic_rc_dir=$(make_case basic-resolver-rc)
|
||||
set +e
|
||||
basic_rc_output=$(MOSAIC_TEST_BASIC_AVAILABLE=true MOSAIC_TEST_BASIC_RC=91 MOSAIC_TEST_FALLBACK_MODE=inspection \
|
||||
run_case "$basic_rc_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
basic_rc_rc=$?
|
||||
set -e
|
||||
[[ "$basic_rc_rc" -ne 0 ]] || fail "failed Basic resolver output unexpectedly authorized a merge: $basic_rc_output"
|
||||
[[ ! -e "$basic_rc_dir/merge-payload.json" ]] || fail "failed Basic resolver reached the merge API"
|
||||
|
||||
# HTTP 401 never changes principals: inspection rejection fails closed without
|
||||
# resolving or attempting Basic Auth.
|
||||
fallback_inspect_dir=$(make_case fallback-inspection)
|
||||
set +e
|
||||
fallback_inspect_output=$(MOSAIC_TEST_BASIC_AVAILABLE=true MOSAIC_TEST_FALLBACK_MODE=inspection \
|
||||
run_case "$fallback_inspect_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
fallback_inspect_rc=$?
|
||||
set -e
|
||||
[[ "$fallback_inspect_rc" -ne 0 ]] || fail "inspection token rejection unexpectedly changed principals"
|
||||
[[ "$fallback_inspect_output" == *"refusing cross-principal credential fallback"* ]] || fail "inspection token rejection lost its refusal diagnostic"
|
||||
[[ ! -e "$fallback_inspect_dir/basic-resolution.log" ]] || fail "inspection token rejection resolved Basic Auth"
|
||||
[[ ! -e "$fallback_inspect_dir/merge-payload.json" ]] || fail "inspection token rejection reached merge mutation"
|
||||
inspect_sequence=$(awk '{print $1 ":" $2}' "$fallback_inspect_dir/curl.log" | paste -sd, -)
|
||||
[[ "$inspect_sequence" == "GET:token" ]] || fail "inspection rejection made unexpected provider calls (calls=$inspect_sequence)"
|
||||
|
||||
# Token rejection at merge likewise fails closed without cross-principal retry.
|
||||
fallback_merge_dir=$(make_case fallback-merge)
|
||||
set +e
|
||||
fallback_merge_output=$(MOSAIC_TEST_BASIC_AVAILABLE=true MOSAIC_TEST_FALLBACK_MODE=merge \
|
||||
run_case "$fallback_merge_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
fallback_merge_rc=$?
|
||||
set -e
|
||||
[[ "$fallback_merge_rc" -ne 0 ]] || fail "merge token rejection unexpectedly changed principals"
|
||||
[[ "$fallback_merge_output" == *"refusing cross-principal credential fallback"* ]] || fail "merge token rejection lost its refusal diagnostic"
|
||||
[[ ! -e "$fallback_merge_dir/basic-resolution.log" ]] || fail "merge token rejection resolved Basic Auth"
|
||||
[[ ! -e "$fallback_merge_dir/merge-payload.json" ]] || fail "merge token rejection recorded a successful payload"
|
||||
merge_sequence=$(awk '{print $1 ":" $2}' "$fallback_merge_dir/curl.log" | paste -sd, -)
|
||||
[[ "$merge_sequence" == "GET:token,GET:token,POST:token" ]] || fail "merge rejection made unexpected provider calls (calls=$merge_sequence)"
|
||||
|
||||
# BLOCK path: a commit email exists but author.login is null. It must name both
|
||||
# facts, name the escalation principal, and never reach the merge endpoint.
|
||||
null_dir=$(make_case null-login)
|
||||
set +e
|
||||
null_output=$(run_case "$null_dir" null-login --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
null_rc=$?
|
||||
set -e
|
||||
[[ "$null_rc" -ne 0 ]] || fail "null-login author expected a non-zero BLOCK"
|
||||
[[ "$null_output" == *"BLOCK"* ]] || fail "null-login author omitted BLOCK diagnostic"
|
||||
[[ "$null_output" == *"author.login=NULL"* ]] || fail "null-login author omitted the null provider fact"
|
||||
[[ "$null_output" == *"[email protected]"* ]] || fail "null-login author omitted the commit email fact"
|
||||
[[ "$null_output" == *'\n\x1b[31m'* ]] || fail "null-login author diagnostic did not escape control characters"
|
||||
[[ "$null_output" != *$'\033'* ]] || fail "null-login author diagnostic emitted a raw terminal escape"
|
||||
[[ "$(printf '%s\n' "$null_output" | wc -l)" -eq 1 ]] || fail "null-login author diagnostic permitted newline injection"
|
||||
[[ "$null_output" == *"tl-mosaic"* ]] || fail "null-login author omitted the named escalation principal"
|
||||
[[ ! -e "$null_dir/merge-payload.json" ]] || fail "null-login BLOCK still reached the merge API"
|
||||
|
||||
# Every provider-derived field in alternate BLOCK diagnostics is log-safe too,
|
||||
# including an invalid non-head SHA that contains control characters.
|
||||
unsafe_dir=$(make_case unsafe-identity)
|
||||
set +e
|
||||
unsafe_output=$(run_case "$unsafe_dir" unsafe-identity --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
unsafe_rc=$?
|
||||
set -e
|
||||
[[ "$unsafe_rc" -ne 0 ]] || fail "unsafe identity expected a non-zero BLOCK"
|
||||
[[ "$unsafe_output" == *"unusable linked identity"* ]] || fail "unsafe identity omitted its BLOCK reason"
|
||||
[[ "$unsafe_output" == *'\n\x1b[31m'* ]] || fail "unsafe identity SHA did not escape control characters"
|
||||
[[ "$unsafe_output" != *$'\033'* ]] || fail "unsafe identity diagnostic emitted a raw terminal escape"
|
||||
[[ "$(printf '%s\n' "$unsafe_output" | wc -l)" -eq 1 ]] || fail "unsafe identity diagnostic permitted newline injection"
|
||||
[[ ! -e "$unsafe_dir/merge-payload.json" ]] || fail "unsafe identity BLOCK still reached the merge API"
|
||||
|
||||
# The provider PR title cannot add an unchecked trailer outside the constructed
|
||||
# message field: multi-line and trailer-shaped titles block before mutation.
|
||||
title_dir=$(make_case title-injection)
|
||||
set +e
|
||||
title_output=$(MOSAIC_TEST_TITLE_MODE=injection \
|
||||
run_case "$title_dir" verified --co-author-trailers --escalate-to tl-mosaic 2>&1)
|
||||
title_rc=$?
|
||||
set -e
|
||||
[[ "$title_rc" -ne 0 ]] || fail "title trailer injection unexpectedly passed"
|
||||
[[ "$title_output" == *"not one printable, non-trailer line"* ]] || fail "title injection refusal lost its diagnostic"
|
||||
[[ ! -e "$title_dir/merge-payload.json" ]] || fail "title injection reached the merge API"
|
||||
|
||||
# Provider failures remain diagnosable after their temporary response file is
|
||||
# removed, but provider-controlled control characters stay log-safe.
|
||||
error_dir=$(make_case provider-error)
|
||||
set +e
|
||||
error_output=$(MOSAIC_TEST_BASIC_AVAILABLE=true MOSAIC_TEST_FALLBACK_MODE=provider-error \
|
||||
run_case "$error_dir" single 2>&1)
|
||||
error_rc=$?
|
||||
set -e
|
||||
[[ "$error_rc" -ne 0 ]] || fail "provider error unexpectedly passed"
|
||||
[[ "$error_output" == *"HTTP 409"* ]] || fail "provider error omitted the HTTP status"
|
||||
[[ "$error_output" == *"branch policy rejected"* ]] || fail "provider error response was discarded"
|
||||
[[ "$error_output" == *'\n\x1b[31m'* ]] || fail "provider error response did not escape control characters"
|
||||
[[ "$error_output" != *$'\033'* ]] || fail "provider error response emitted a raw terminal escape"
|
||||
[[ "$error_output" != *"Basic Auth fallback"* ]] || fail "provider error advertised removed Basic Auth fallback"
|
||||
[[ ! -e "$error_dir/basic-resolution.log" ]] || fail "HTTP 409 policy denial incorrectly triggered Basic Auth fallback"
|
||||
|
||||
# Authorization denials likewise fail closed instead of changing principals.
|
||||
forbidden_dir=$(make_case forbidden)
|
||||
set +e
|
||||
forbidden_output=$(MOSAIC_TEST_BASIC_AVAILABLE=true MOSAIC_TEST_FALLBACK_MODE=forbidden \
|
||||
run_case "$forbidden_dir" single 2>&1)
|
||||
forbidden_rc=$?
|
||||
set -e
|
||||
[[ "$forbidden_rc" -ne 0 ]] || fail "HTTP 403 authorization denial unexpectedly passed"
|
||||
[[ "$forbidden_output" == *"HTTP 403"* ]] || fail "authorization denial omitted the HTTP status"
|
||||
[[ "$forbidden_output" != *"Basic Auth fallback"* ]] || fail "authorization denial advertised removed Basic Auth fallback"
|
||||
[[ ! -e "$forbidden_dir/basic-resolution.log" ]] || fail "HTTP 403 authorization denial incorrectly triggered Basic Auth fallback"
|
||||
|
||||
# The BLOCK destination cannot be generic or inferred after failure: opting in
|
||||
# without a named principal is refused before any provider operation.
|
||||
principal_dir=$(make_case missing-principal)
|
||||
set +e
|
||||
principal_output=$(run_case "$principal_dir" verified --co-author-trailers 2>&1)
|
||||
principal_rc=$?
|
||||
set -e
|
||||
[[ "$principal_rc" -ne 0 ]] || fail "co-author mode without a named principal unexpectedly passed"
|
||||
[[ "$principal_output" == *"requires --escalate-to with a named principal"* ]] || fail "missing-principal refusal lost its diagnostic"
|
||||
[[ ! -e "$principal_dir/merge-payload.json" ]] || fail "missing-principal refusal reached the merge API"
|
||||
|
||||
# A trailing value-taking option receives a stable CLI diagnostic instead of a
|
||||
# set -u unbound-variable crash.
|
||||
value_dir=$(make_case missing-principal-value)
|
||||
set +e
|
||||
value_output=$(run_case "$value_dir" verified --co-author-trailers --escalate-to 2>&1)
|
||||
value_rc=$?
|
||||
set -e
|
||||
[[ "$value_rc" -ne 0 ]] || fail "missing --escalate-to value unexpectedly passed"
|
||||
[[ "$value_output" == *"--escalate-to requires one principal name"* ]] || fail "missing --escalate-to value lost its diagnostic"
|
||||
[[ "$value_output" != *"unbound variable"* ]] || fail "missing --escalate-to value crashed under set -u"
|
||||
[[ ! -e "$value_dir/merge-payload.json" ]] || fail "missing --escalate-to value reached the merge API"
|
||||
|
||||
# Negative control: ordinary single-author merge remains byte-for-byte payload
|
||||
# compatible and hardcoded to squash, with no optional message fields.
|
||||
single_dir=$(make_case single)
|
||||
set +e
|
||||
single_output=$(run_case "$single_dir" single 2>&1)
|
||||
single_rc=$?
|
||||
set -e
|
||||
if [[ "$single_rc" -ne 0 ]]; then
|
||||
fail "ordinary single-author merge expected rc=0, got rc=$single_rc: $single_output"
|
||||
elif [[ ! -s "$single_dir/merge-payload.json" ]]; then
|
||||
fail "ordinary single-author merge did not reach the API payload"
|
||||
else
|
||||
python3 - "$single_dir/merge-payload.json" <<'PY' || fail "ordinary single-author payload changed"
|
||||
import json
|
||||
import sys
|
||||
payload = json.load(open(sys.argv[1], encoding="utf-8"))
|
||||
assert payload == {
|
||||
"Do": "squash",
|
||||
"head_commit_id": "1111111111111111111111111111111111111111",
|
||||
}, payload
|
||||
PY
|
||||
fi
|
||||
[[ -e "$single_dir/auth-via-config" ]] || fail "ordinary path did not authenticate curl through stdin config"
|
||||
[[ ! -e "$single_dir/token-in-argv" ]] || fail "ordinary path placed the Gitea token in curl argv"
|
||||
[[ "$(wc -l < "$single_dir/token-resolution.log")" -eq 1 ]] || fail "ordinary path did not use exactly one credential resolution"
|
||||
|
||||
# Squash is not defaultable: an explicit non-squash method must remain refused.
|
||||
method_dir=$(make_case method-refusal)
|
||||
set +e
|
||||
method_output=$(run_case "$method_dir" single -m merge 2>&1)
|
||||
method_rc=$?
|
||||
set -e
|
||||
[[ "$method_rc" -ne 0 ]] || fail "non-squash method unexpectedly passed"
|
||||
[[ "$method_output" == *"enforces squash merge only"* ]] || fail "non-squash refusal lost its policy diagnostic"
|
||||
[[ ! -e "$method_dir/merge-payload.json" ]] || fail "non-squash refusal reached the merge API"
|
||||
|
||||
if [[ "$failures" -ne 0 ]]; then
|
||||
echo "pr-merge message-field regression failed ($failures assertions)" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
echo "pr-merge message-field regression passed (verified, BLOCK, and unchanged squash control)"
|
||||
@@ -39,7 +39,7 @@ MAX_FRAME: Final = 64 * 1024
|
||||
MAX_STATE: Final = 4 * 1024 * 1024
|
||||
MAX_PENDING_TOKENS: Final = 256
|
||||
MAX_IN_FLIGHT_CONNECTIONS: Final = 16
|
||||
MAX_LEASE_TTL_SECONDS: Final = 3600
|
||||
MAX_LEASE_TTL_SECONDS: Final = 300
|
||||
STATE_VERSION: Final = 1
|
||||
READ_DEADLINE_SECONDS: Final = 1.0
|
||||
HANDLE_QUEUE_TIMEOUT_SECONDS: Final = 1.0
|
||||
@@ -50,8 +50,8 @@ LEASE_PENDING: Final = "PENDING_VERIFICATION"
|
||||
LEASE_PENDING_PROMOTION: Final = "PENDING_PROMOTION"
|
||||
LEASE_VERIFIED: Final = "VERIFIED"
|
||||
READ_ONLY_TOOLS: Final = {
|
||||
"claude": frozenset({"Read", "Grep", "Glob"}),
|
||||
"pi": frozenset({"read", "ls"}),
|
||||
"claude": frozenset({"Read", "Grep", "Glob", "Ls", "Find"}),
|
||||
"pi": frozenset({"read", "grep", "find", "ls"}),
|
||||
}
|
||||
RECOVERY_TOOL: Final = "mosaic_context_recover"
|
||||
|
||||
|
||||
@@ -8,9 +8,7 @@ import json
|
||||
import os
|
||||
import socket
|
||||
import sys
|
||||
import time
|
||||
from collections.abc import Callable, Mapping, Sequence
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Final
|
||||
|
||||
@@ -55,48 +53,6 @@ def broker_request(socket_path: Path, request: dict[str, object]) -> dict[str, o
|
||||
return value
|
||||
|
||||
|
||||
def _self_starttime() -> str | None:
|
||||
"""Field 22 of our own /proc stat — the anchor starttime the broker records.
|
||||
|
||||
Read past the comm field's parens, since a process name may contain them.
|
||||
"""
|
||||
try:
|
||||
raw = Path(f"/proc/{os.getpid()}/stat").read_text()
|
||||
return raw.rsplit(")", 1)[1].split()[19]
|
||||
except (OSError, IndexError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _append_launch_record(environ: Mapping[str, str], record: dict[str, object]) -> None:
|
||||
"""Append one NDJSON event to the #797 Runtime Session Ledger.
|
||||
|
||||
`fleet/run/sessions/` is operator-classified in framework-manifest.txt and is
|
||||
already covered by test-upgrade-manifest-guard.sh, so an upgrade can neither
|
||||
overwrite nor prune it. Files 0600 under a 0700 dir, matching what that guard
|
||||
asserts.
|
||||
|
||||
Never raises: a launch must not be denied over bookkeeping. But it also never
|
||||
fails silently — a missing record is exactly the kind of gap that made the
|
||||
2026-08-06 MUTATOR_UNVERIFIED investigation cost a day.
|
||||
"""
|
||||
try:
|
||||
mosaic_home = environ.get("MOSAIC_HOME") or str(Path.home() / ".config" / "mosaic")
|
||||
directory = Path(mosaic_home) / "fleet" / "run" / "sessions"
|
||||
directory.mkdir(parents=True, exist_ok=True)
|
||||
os.chmod(directory, 0o700)
|
||||
framed = {
|
||||
"seq": time.time_ns() // 1_000_000,
|
||||
"ts": datetime.now(timezone.utc).isoformat(),
|
||||
**record,
|
||||
}
|
||||
path = directory / "events.ndjson"
|
||||
descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o600)
|
||||
with os.fdopen(descriptor, "w") as handle:
|
||||
handle.write(json.dumps(framed, separators=(",", ":")) + "\n")
|
||||
except (OSError, ValueError, TypeError) as error:
|
||||
print(f"[mosaic] WARNING: launch record not written: {error}", file=sys.stderr)
|
||||
|
||||
|
||||
def main(
|
||||
argv: Sequence[str] | None = None,
|
||||
*,
|
||||
@@ -138,9 +94,8 @@ def main(
|
||||
# silent pass and never folded into the generic registration-failure
|
||||
# branch.
|
||||
try:
|
||||
activation_capability = probe_activation_capability(source_environment)
|
||||
assert_activation_capability_matches(
|
||||
activation_capability,
|
||||
probe_activation_capability(source_environment),
|
||||
expected_activation_capability,
|
||||
)
|
||||
except VersionCouplingError as version_error:
|
||||
@@ -173,32 +128,6 @@ def main(
|
||||
print("Mosaic lease broker registration failed; runtime launch denied.", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
# Immutable launch record, half two. `mosaic` wrote `session.launch` with the
|
||||
# config/provenance it knows; only this process knows the broker session id
|
||||
# and the activation capability it just asserted. os.execvpe preserves the
|
||||
# PID, so this PID is BOTH the anchor pid and the join key back to that
|
||||
# record. Never fatal — bookkeeping must not deny a launch — but never
|
||||
# silent either.
|
||||
_append_launch_record(
|
||||
source_environment,
|
||||
{
|
||||
"kind": "lease.register",
|
||||
# Joins back to `mosaic`'s session.launch record. NOT pid: execRuntime()
|
||||
# spawns rather than execs, so this process is a CHILD of mosaic with a
|
||||
# different pid. This pid IS the broker anchor pid (os.execvpe below
|
||||
# preserves it), which is a separate and still-useful fact.
|
||||
"launch_id": source_environment.get("MOSAIC_LAUNCH_ID"),
|
||||
"pid": os.getpid(),
|
||||
"runtime": arguments.runtime,
|
||||
"session_id": session_id,
|
||||
"runtime_generation": generation,
|
||||
"generation_file": str(generation_file),
|
||||
"anchor_starttime": _self_starttime(),
|
||||
"activation_capability": activation_capability,
|
||||
"command": Path(command[0]).name,
|
||||
},
|
||||
)
|
||||
|
||||
environment = dict(source_environment)
|
||||
environment["MOSAIC_LEASE_SESSION_ID"] = session_id
|
||||
environment["MOSAIC_RUNTIME_GENERATION"] = str(generation)
|
||||
|
||||
@@ -1,337 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Lease promotion client — the half the enforcement toolkit never shipped.
|
||||
|
||||
The enforcement half (``daemon.py`` + ``mutator-gate.py``) ships and denies. The
|
||||
promotion half has no production caller anywhere in the package: as of 0.0.48,
|
||||
0.0.49 and 0.0.50-next.2207, ``begin_verification`` / ``observe_receipt`` /
|
||||
``promote_lease`` are invoked only by ``broker-test-client.ts``, the acceptance
|
||||
spec, unit tests, and two probes under ``docs/``. Consequence: **no lease on any
|
||||
host can reach VERIFIED**, so every mutator is denied ``MUTATOR_UNVERIFIED`` by a
|
||||
gate nothing can satisfy.
|
||||
|
||||
THE PROTOCOL (``daemon.py:578-754``)
|
||||
------------------------------------
|
||||
1. ``begin_verification`` — broker revokes, mints a challenge, and returns the
|
||||
exact ``receipt`` text the MODEL must emit
|
||||
2. *the model emits that text verbatim as its ENTIRE latest message*
|
||||
3. the runtime adapter ships that message to the daemon-owned observer socket
|
||||
4. ``observe_receipt`` -> ``PENDING_PROMOTION``
|
||||
5. ``promote_lease`` -> ``VERIFIED``
|
||||
|
||||
THIS MODULE IMPLEMENTS 1, 4 AND 5 — NEVER 2
|
||||
-------------------------------------------
|
||||
Step 2 is the security property, not a formality. ``is_verbatim_receipt`` uses
|
||||
``hmac.compare_digest`` against the exact minted string — explicitly "not a
|
||||
transcript substring" (``receipt_challenge.py``). Promotion therefore requires a
|
||||
live model that received the challenge in its context and echoed it exactly.
|
||||
|
||||
``receipt-observer-client.py`` will post ANY string as the latest assistant
|
||||
message. A promotion client that posted its own receipt would satisfy the broker
|
||||
while proving nothing — a gate-disabler indistinguishable from a working fix
|
||||
unless someone looks for it. **This module never posts a receipt.** Emitting it
|
||||
belongs to the runtime adapter, where a real model turn happens.
|
||||
|
||||
The construction binds the exact normative source bytes. ``h_source`` /
|
||||
``h_payload`` are derived by the framework's own
|
||||
``normative_fragments.build_payload`` rather than reimplemented: the broker
|
||||
derives them the same way and any divergence yields ``PAYLOAD_BINDING_MISMATCH``.
|
||||
There must be exactly one implementation.
|
||||
|
||||
WHAT THE BINDING DOES *NOT* PROVE
|
||||
---------------------------------
|
||||
It is tempting to read a VERIFIED lease as "this agent is running THIS law".
|
||||
**It does not mean that**, and writing it down that way is how the belief spread.
|
||||
The broker holds no reference copy of any normative source and never opens one;
|
||||
it recomputes ``h_source`` / ``h_payload`` from the fragment bytes THIS CLIENT
|
||||
sent and compares them to the binding THIS CLIENT sent (``daemon.py:602-616``).
|
||||
Both sides of that comparison originate here, so it detects corruption in
|
||||
transit and nothing else. What the binding actually asserts is "the client
|
||||
claims these bytes, self-consistently".
|
||||
|
||||
Making it mean the stronger thing requires the broker to re-read the on-disk
|
||||
sources itself, against a manifest the agent cannot rewrite — i.e. broker code
|
||||
attestation under its own uid. Until then, do not cite a VERIFIED lease as
|
||||
evidence of law integrity.
|
||||
|
||||
Usage
|
||||
-----
|
||||
lease_promote.py --begin # prints the receipt the MODEL must emit
|
||||
lease_promote.py --complete <challenge> # after the adapter observed it
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import base64
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import socket
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Final
|
||||
|
||||
# Isolated (`python -I`) adapter invocations must still import co-located
|
||||
# framework modules; never depend on the caller's PYTHONPATH.
|
||||
_MODULE_DIRECTORY = str(Path(__file__).resolve().parent)
|
||||
if _MODULE_DIRECTORY not in sys.path:
|
||||
sys.path.insert(0, _MODULE_DIRECTORY)
|
||||
|
||||
from normative_fragments import NormativeFragment, build_payload # noqa: E402
|
||||
|
||||
MAX_FRAME: Final = 64 * 1024
|
||||
BROKER_TIMEOUT_SECONDS: Final = 3.0
|
||||
SCHEMA_VERSION: Final = 1
|
||||
MANIFEST_VERSION: Final = 1
|
||||
GENERATOR_VERSION: Final = "mosaic/lease_promote@1"
|
||||
DEFAULT_TTL_SECONDS: Final = 3600
|
||||
|
||||
# Normative sources whose exact bytes bind the lease, in binding order. Order is
|
||||
# load-bearing: ``h_source`` frames the resolved sequence, so reordering changes
|
||||
# the derivation. Never fabricate a source that is not on disk.
|
||||
FRAGMENT_SOURCES: Final = (
|
||||
"CONSTITUTION.md",
|
||||
"AGENTS.md",
|
||||
"SOUL.md",
|
||||
"USER.md",
|
||||
"STANDARDS.md",
|
||||
"TOOLS.md",
|
||||
)
|
||||
|
||||
# Framework-owned sources, reconciled on every upgrade — `install.sh:76`
|
||||
# FRAMEWORK_OWNED and `config/file-adapter.ts` FRAMEWORK_OWNED_FILES — plus the
|
||||
# per-runtime contract shipped under `framework/runtime/<runtime>/`. A deployment
|
||||
# missing one of these is broken, not minimal, so their absence is refused rather
|
||||
# than silently dropped from the binding.
|
||||
#
|
||||
# SOUL.md and USER.md are deliberately excluded: install.sh does not seed them
|
||||
# ("intentionally NOT seeded here — they are generated by `mosaic init`"), so a
|
||||
# fresh install legitimately lacks both. TOOLS.md is user-seeded on first install
|
||||
# only. Absence of those three is reported, not fatal.
|
||||
REQUIRED_SOURCES: Final = frozenset({"CONSTITUTION.md", "AGENTS.md", "STANDARDS.md"})
|
||||
|
||||
|
||||
class IncompleteBinding(RuntimeError):
|
||||
"""A source that must bind this lease could not be read.
|
||||
|
||||
**Never downgrade this to a skip.** The broker recomputes the hashes from the
|
||||
fragments it is sent, so an omitted fragment is internally consistent and
|
||||
``PAYLOAD_BINDING_MISMATCH`` cannot fire — a partial law promotes exactly like
|
||||
a complete one, and nothing downstream can tell the difference. Dropping an
|
||||
unreadable source therefore does not degrade the binding, it forges a smaller
|
||||
one. Fail here, where the omission is still visible.
|
||||
"""
|
||||
|
||||
|
||||
def mosaic_home() -> Path:
|
||||
return Path(os.environ.get("MOSAIC_HOME") or Path.home() / ".config" / "mosaic")
|
||||
|
||||
|
||||
def broker_socket() -> Path:
|
||||
value = os.environ.get("MOSAIC_LEASE_BROKER_SOCKET")
|
||||
if value:
|
||||
return Path(value)
|
||||
runtime_dir = os.environ.get("XDG_RUNTIME_DIR")
|
||||
if runtime_dir:
|
||||
return Path(runtime_dir) / "mosaic-lease" / "broker.sock"
|
||||
return Path(f"/run/user/{os.getuid()}/mosaic-lease/broker.sock")
|
||||
|
||||
|
||||
def session_identity() -> tuple[str, int, str]:
|
||||
"""Session id, CURRENT generation, runtime.
|
||||
|
||||
The generation file wins over the env var, matching ``lease_generation.py``.
|
||||
Sending a generation HIGHER than the broker's would revoke this session's own
|
||||
authority (``daemon.py:342-344``), so this never guesses.
|
||||
"""
|
||||
session_id = os.environ["MOSAIC_LEASE_SESSION_ID"]
|
||||
runtime = os.environ["MOSAIC_LEASE_RUNTIME"]
|
||||
state_file = os.environ.get("MOSAIC_LEASE_GENERATION_FILE")
|
||||
if state_file:
|
||||
try:
|
||||
return session_id, int(Path(state_file).read_text().strip()), runtime
|
||||
except (OSError, ValueError):
|
||||
pass
|
||||
return session_id, int(os.environ["MOSAIC_RUNTIME_GENERATION"]), runtime
|
||||
|
||||
|
||||
def build_construction(runtime: str) -> tuple[dict[str, object], object]:
|
||||
"""Assemble the wire construction and derive its hashes with the sole builder."""
|
||||
runtime_contract = f"runtime/{runtime}/RUNTIME.md"
|
||||
sources = list(FRAGMENT_SOURCES) + [runtime_contract]
|
||||
required = REQUIRED_SOURCES | {runtime_contract}
|
||||
wire_fragments: list[dict[str, str]] = []
|
||||
objects: list[NormativeFragment] = []
|
||||
absent: list[str] = []
|
||||
|
||||
for source_id in sources:
|
||||
try:
|
||||
content = (mosaic_home() / source_id).read_bytes()
|
||||
except FileNotFoundError:
|
||||
# Genuinely not on disk. Legitimate only for operator-owned sources.
|
||||
if source_id in required:
|
||||
raise IncompleteBinding(
|
||||
f"required normative source is absent: {source_id}"
|
||||
) from None
|
||||
absent.append(source_id)
|
||||
continue
|
||||
except OSError as exc:
|
||||
# The path resolves but will not read — EACCES, EIO, EISDIR, ELOOP.
|
||||
# That is an anomaly for EVERY source, optional ones included: an
|
||||
# unreadable file is not an un-configured one, and treating it as
|
||||
# absent is what lets a permission change quietly shrink the law.
|
||||
raise IncompleteBinding(
|
||||
f"normative source is present but unreadable: {source_id} "
|
||||
f"({type(exc).__name__})"
|
||||
) from exc
|
||||
|
||||
digest = hashlib.sha256(content).hexdigest()
|
||||
wire_fragments.append(
|
||||
{
|
||||
"source_id": source_id,
|
||||
"content_base64": base64.b64encode(content).decode("ascii"),
|
||||
"expected_sha256": digest,
|
||||
}
|
||||
)
|
||||
objects.append(NormativeFragment(source_id, content, digest))
|
||||
|
||||
if not wire_fragments:
|
||||
raise IncompleteBinding("no normative sources found — refusing an empty binding")
|
||||
|
||||
# Absence is legitimate here but never invisible. The omission is already
|
||||
# baked into h_source (the framed source sequence differs), but nothing
|
||||
# compares h_source to an expected value, so this line is the only place a
|
||||
# human learns the binding was narrower than the full set.
|
||||
if absent:
|
||||
print(
|
||||
f"lease_promote: binding omits absent operator sources: {', '.join(absent)}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
|
||||
result = build_payload(
|
||||
manifest_version=MANIFEST_VERSION,
|
||||
generator_version=GENERATOR_VERSION,
|
||||
fragments=objects,
|
||||
)
|
||||
if result.injectionDecision != "ACCEPTED" or not result.promotion:
|
||||
raise RuntimeError(f"construction refused locally: {result.source_reason}")
|
||||
|
||||
return (
|
||||
{
|
||||
"manifest_version": MANIFEST_VERSION,
|
||||
"generator_version": GENERATOR_VERSION,
|
||||
"fragments": wire_fragments,
|
||||
},
|
||||
result,
|
||||
)
|
||||
|
||||
|
||||
def broker_request(payload: dict[str, object]) -> dict[str, object]:
|
||||
raw = (json.dumps(payload, separators=(",", ":")) + "\n").encode()
|
||||
if len(raw) > MAX_FRAME:
|
||||
raise ValueError(
|
||||
f"request too large ({len(raw)} bytes); broker frame cap is {MAX_FRAME}"
|
||||
)
|
||||
response = bytearray()
|
||||
with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as connection:
|
||||
connection.settimeout(BROKER_TIMEOUT_SECONDS)
|
||||
connection.connect(str(broker_socket()))
|
||||
connection.sendall(raw)
|
||||
connection.shutdown(socket.SHUT_WR)
|
||||
while len(response) <= MAX_FRAME:
|
||||
chunk = connection.recv(4096)
|
||||
if not chunk:
|
||||
break
|
||||
response.extend(chunk)
|
||||
if len(response) > MAX_FRAME or not response.endswith(b"\n"):
|
||||
raise ValueError("invalid broker reply")
|
||||
value = json.loads(response)
|
||||
if not isinstance(value, dict):
|
||||
raise ValueError("invalid broker reply")
|
||||
return value
|
||||
|
||||
|
||||
def begin(
|
||||
ttl_seconds: int = DEFAULT_TTL_SECONDS,
|
||||
compaction_epoch: int = 0,
|
||||
request_epoch: int = 0,
|
||||
) -> dict[str, object]:
|
||||
"""Step 1. Returns the broker reply, including the exact ``receipt`` text."""
|
||||
session_id, generation, runtime = session_identity()
|
||||
construction, derived = build_construction(runtime)
|
||||
return broker_request(
|
||||
{
|
||||
"action": "begin_verification",
|
||||
"session_id": session_id,
|
||||
"runtime_generation": generation,
|
||||
"runtime": runtime,
|
||||
"ttl_seconds": ttl_seconds,
|
||||
"binding": {
|
||||
"compaction_epoch": compaction_epoch,
|
||||
"request_epoch": request_epoch,
|
||||
"h_source": derived.h_source,
|
||||
"h_payload": derived.h_payload,
|
||||
"schema_version": SCHEMA_VERSION,
|
||||
},
|
||||
"construction": construction,
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def complete(challenge: str) -> dict[str, object]:
|
||||
"""Steps 4-5. Assumes the model already emitted the receipt and the adapter
|
||||
shipped it to the observer socket."""
|
||||
session_id, generation, _ = session_identity()
|
||||
observed = broker_request(
|
||||
{
|
||||
"action": "observe_receipt",
|
||||
"session_id": session_id,
|
||||
"runtime_generation": generation,
|
||||
"receipt_challenge": challenge,
|
||||
}
|
||||
)
|
||||
if observed.get("ok") is not True or observed.get("state") != "PENDING_PROMOTION":
|
||||
return {"stage": "observe_receipt", **observed}
|
||||
promoted = broker_request(
|
||||
{
|
||||
"action": "promote_lease",
|
||||
"session_id": session_id,
|
||||
"runtime_generation": generation,
|
||||
"receipt_challenge": challenge,
|
||||
}
|
||||
)
|
||||
return {"stage": "promote_lease", **promoted}
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Mosaic lease promotion client.")
|
||||
group = parser.add_mutually_exclusive_group(required=True)
|
||||
group.add_argument(
|
||||
"--begin",
|
||||
action="store_true",
|
||||
help="mint a challenge; prints the receipt the MODEL must emit verbatim",
|
||||
)
|
||||
group.add_argument(
|
||||
"--complete",
|
||||
metavar="CHALLENGE",
|
||||
help="observe the emitted receipt and promote the lease",
|
||||
)
|
||||
parser.add_argument("--ttl-seconds", type=int, default=DEFAULT_TTL_SECONDS)
|
||||
arguments = parser.parse_args(argv)
|
||||
|
||||
try:
|
||||
if arguments.begin:
|
||||
print(json.dumps(begin(ttl_seconds=arguments.ttl_seconds), indent=2))
|
||||
else:
|
||||
print(json.dumps(complete(arguments.complete), indent=2))
|
||||
except KeyError as exc:
|
||||
print(f"missing lease environment: {exc}; not a lease-gated session", file=sys.stderr)
|
||||
return 2
|
||||
except (OSError, ValueError, RuntimeError, json.JSONDecodeError) as exc:
|
||||
print(f"{type(exc).__name__}: {exc}", file=sys.stderr)
|
||||
return 2
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -1,333 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Claude UserPromptSubmit hook for operator-triggered lease promotion."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import fcntl
|
||||
import json
|
||||
import os
|
||||
import secrets
|
||||
import stat
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
from collections.abc import Callable, Mapping
|
||||
from pathlib import Path
|
||||
from typing import Final, TextIO
|
||||
|
||||
_MODULE_DIRECTORY = str(Path(__file__).resolve().parent)
|
||||
if _MODULE_DIRECTORY not in sys.path:
|
||||
sys.path.insert(0, _MODULE_DIRECTORY)
|
||||
|
||||
from receipt_challenge import receipt_for # noqa: E402
|
||||
|
||||
MAX_FRAME: Final = 64 * 1024
|
||||
PENDING_MAX_AGE_SECONDS: Final = 60 * 60
|
||||
PROMOTER_TIMEOUT_SECONDS: Final = 10.0
|
||||
PROMOTION_PROMPT: Final = "/mosaic-promote"
|
||||
PROMOTER: Final = Path(__file__).resolve().with_name("lease_promote.py")
|
||||
PENDING_DIRECTORY: Final = "mosaic-lease"
|
||||
LOCK_FILE: Final = "promotion.lock"
|
||||
EXPECTED_BEGIN_KEYS: Final = frozenset(
|
||||
{"ok", "state", "receipt_challenge", "receipt", "binding"}
|
||||
)
|
||||
EXPECTED_BINDING_KEYS: Final = frozenset(
|
||||
{
|
||||
"compaction_epoch",
|
||||
"request_epoch",
|
||||
"h_source",
|
||||
"h_payload",
|
||||
"runtime_generation",
|
||||
"schema_version",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
class PromotionAlreadyInProgress(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
def reject_duplicate_json_keys(pairs: list[tuple[str, object]]) -> dict[str, object]:
|
||||
value: dict[str, object] = {}
|
||||
for key, item in pairs:
|
||||
if key in value:
|
||||
raise ValueError("duplicate promoter JSON key")
|
||||
value[key] = item
|
||||
return value
|
||||
|
||||
|
||||
def read_hook_input(stream: object) -> dict[str, object]:
|
||||
raw = getattr(stream, "buffer", stream).read(MAX_FRAME + 1)
|
||||
if not isinstance(raw, bytes) or len(raw) > MAX_FRAME:
|
||||
raise ValueError("invalid UserPromptSubmit input")
|
||||
value = json.loads(raw, object_pairs_hook=reject_duplicate_json_keys)
|
||||
if not isinstance(value, dict):
|
||||
raise ValueError("invalid UserPromptSubmit input")
|
||||
return value
|
||||
|
||||
|
||||
def emit_context(stream: TextIO, message: str) -> None:
|
||||
json.dump(
|
||||
{
|
||||
"hookSpecificOutput": {
|
||||
"hookEventName": "UserPromptSubmit",
|
||||
"additionalContext": message,
|
||||
}
|
||||
},
|
||||
stream,
|
||||
separators=(",", ":"),
|
||||
)
|
||||
stream.write("\n")
|
||||
|
||||
|
||||
def session_pending_name(environ: Mapping[str, str]) -> tuple[Path, str]:
|
||||
runtime_dir = Path(environ["XDG_RUNTIME_DIR"])
|
||||
session_id = environ["MOSAIC_LEASE_SESSION_ID"]
|
||||
if not runtime_dir.is_absolute():
|
||||
raise ValueError("XDG_RUNTIME_DIR must be absolute")
|
||||
if len(session_id) != 64 or any(character not in "0123456789abcdef" for character in session_id):
|
||||
raise ValueError("invalid lease session id")
|
||||
return runtime_dir, f"pending-{session_id}"
|
||||
|
||||
|
||||
def open_pending_directory(runtime_dir: Path) -> int:
|
||||
directory_flags = (
|
||||
os.O_RDONLY
|
||||
| getattr(os, "O_CLOEXEC", 0)
|
||||
| getattr(os, "O_DIRECTORY", 0)
|
||||
| getattr(os, "O_NOFOLLOW", 0)
|
||||
)
|
||||
runtime_descriptor = os.open(runtime_dir, directory_flags)
|
||||
try:
|
||||
runtime_metadata = os.fstat(runtime_descriptor)
|
||||
if (
|
||||
not stat.S_ISDIR(runtime_metadata.st_mode)
|
||||
or runtime_metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(runtime_metadata.st_mode) != 0o700
|
||||
):
|
||||
raise ValueError("unsafe XDG runtime directory")
|
||||
try:
|
||||
os.mkdir(PENDING_DIRECTORY, mode=0o700, dir_fd=runtime_descriptor)
|
||||
except FileExistsError:
|
||||
pass
|
||||
descriptor = os.open(PENDING_DIRECTORY, directory_flags, dir_fd=runtime_descriptor)
|
||||
finally:
|
||||
os.close(runtime_descriptor)
|
||||
|
||||
metadata = os.fstat(descriptor)
|
||||
if (
|
||||
not stat.S_ISDIR(metadata.st_mode)
|
||||
or metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(metadata.st_mode) != 0o700
|
||||
):
|
||||
os.close(descriptor)
|
||||
raise ValueError("unsafe promotion pending directory")
|
||||
return descriptor
|
||||
|
||||
|
||||
def acquire_lock(directory_descriptor: int) -> int:
|
||||
flags = (
|
||||
os.O_RDWR
|
||||
| os.O_CREAT
|
||||
| getattr(os, "O_CLOEXEC", 0)
|
||||
| getattr(os, "O_NOFOLLOW", 0)
|
||||
)
|
||||
descriptor = os.open(LOCK_FILE, flags, 0o600, dir_fd=directory_descriptor)
|
||||
metadata = os.fstat(descriptor)
|
||||
if (
|
||||
not stat.S_ISREG(metadata.st_mode)
|
||||
or metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(metadata.st_mode) != 0o600
|
||||
):
|
||||
os.close(descriptor)
|
||||
raise ValueError("unsafe promotion lock file")
|
||||
try:
|
||||
fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
except BlockingIOError as error:
|
||||
os.close(descriptor)
|
||||
raise PromotionAlreadyInProgress() from error
|
||||
return descriptor
|
||||
|
||||
|
||||
def sweep_stale_pending(directory_descriptor: int, current_time: float) -> None:
|
||||
cutoff = current_time - PENDING_MAX_AGE_SECONDS
|
||||
removed = False
|
||||
with os.scandir(directory_descriptor) as entries:
|
||||
for candidate in entries:
|
||||
if not (
|
||||
candidate.name.startswith("pending-")
|
||||
or candidate.name.startswith(".pending-")
|
||||
):
|
||||
continue
|
||||
try:
|
||||
metadata = candidate.stat(follow_symlinks=False)
|
||||
if metadata.st_mtime < cutoff and not stat.S_ISDIR(metadata.st_mode):
|
||||
os.unlink(candidate.name, dir_fd=directory_descriptor)
|
||||
removed = True
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
if removed:
|
||||
os.fsync(directory_descriptor)
|
||||
|
||||
|
||||
def write_pending(directory_descriptor: int, name: str, challenge: str) -> None:
|
||||
temporary = f".{name}.tmp-{secrets.token_hex(8)}"
|
||||
flags = (
|
||||
os.O_WRONLY
|
||||
| os.O_CREAT
|
||||
| os.O_EXCL
|
||||
| getattr(os, "O_CLOEXEC", 0)
|
||||
| getattr(os, "O_NOFOLLOW", 0)
|
||||
)
|
||||
descriptor = os.open(temporary, flags, 0o600, dir_fd=directory_descriptor)
|
||||
try:
|
||||
os.fchmod(descriptor, 0o600)
|
||||
with os.fdopen(descriptor, "w", encoding="utf-8", closefd=False) as stream:
|
||||
stream.write(challenge)
|
||||
stream.flush()
|
||||
os.fsync(stream.fileno())
|
||||
os.replace(
|
||||
temporary,
|
||||
name,
|
||||
src_dir_fd=directory_descriptor,
|
||||
dst_dir_fd=directory_descriptor,
|
||||
)
|
||||
os.fsync(directory_descriptor)
|
||||
except Exception:
|
||||
try:
|
||||
os.unlink(temporary, dir_fd=directory_descriptor)
|
||||
except FileNotFoundError:
|
||||
pass
|
||||
raise
|
||||
finally:
|
||||
os.close(descriptor)
|
||||
|
||||
|
||||
def parse_begin_reply(
|
||||
completed: subprocess.CompletedProcess[str],
|
||||
) -> tuple[str, dict[str, object] | None]:
|
||||
if completed.returncode != 0:
|
||||
return f"PROMOTER_EXIT_{completed.returncode}", None
|
||||
try:
|
||||
value = json.loads(
|
||||
completed.stdout,
|
||||
object_pairs_hook=reject_duplicate_json_keys,
|
||||
)
|
||||
except (json.JSONDecodeError, RecursionError, TypeError, ValueError):
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
if not isinstance(value, dict):
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
if value.get("ok") is False and set(value) == {"ok", "code"}:
|
||||
code = value.get("code")
|
||||
return code if isinstance(code, str) and code else "PROMOTION_BEGIN_REFUSED", value
|
||||
if set(value) != EXPECTED_BEGIN_KEYS or value.get("ok") is not True:
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
if value.get("state") != "PENDING_VERIFICATION":
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
challenge = value.get("receipt_challenge")
|
||||
receipt = value.get("receipt")
|
||||
binding = value.get("binding")
|
||||
if (
|
||||
not isinstance(challenge, str)
|
||||
or len(challenge) != 64
|
||||
or any(character not in "0123456789abcdef" for character in challenge)
|
||||
or not isinstance(receipt, str)
|
||||
or not isinstance(binding, dict)
|
||||
or set(binding) != EXPECTED_BINDING_KEYS
|
||||
):
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
integer_fields = (
|
||||
"compaction_epoch",
|
||||
"request_epoch",
|
||||
"runtime_generation",
|
||||
"schema_version",
|
||||
)
|
||||
if any(type(binding.get(field)) is not int or binding[field] < 0 for field in integer_fields):
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
if not all(
|
||||
isinstance(binding.get(field), str)
|
||||
and len(binding[field]) == 64
|
||||
and all(character in "0123456789abcdef" for character in binding[field])
|
||||
for field in ("h_source", "h_payload")
|
||||
):
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
if not secrets.compare_digest(
|
||||
receipt.encode("utf-8"),
|
||||
receipt_for(challenge, binding).encode("utf-8"),
|
||||
):
|
||||
return "INVALID_PROMOTER_REPLY", None
|
||||
return "", value
|
||||
|
||||
|
||||
def main(
|
||||
*,
|
||||
environ: Mapping[str, str] | None = None,
|
||||
stdin: object | None = None,
|
||||
stdout: TextIO | None = None,
|
||||
stderr: TextIO | None = None,
|
||||
run: Callable[..., subprocess.CompletedProcess[str]] = subprocess.run,
|
||||
now: Callable[[], float] = time.time,
|
||||
) -> int:
|
||||
source_environment = os.environ if environ is None else environ
|
||||
input_stream = sys.stdin if stdin is None else stdin
|
||||
output_stream = sys.stdout if stdout is None else stdout
|
||||
error_stream = sys.stderr if stderr is None else stderr
|
||||
|
||||
try:
|
||||
hook_input = read_hook_input(input_stream)
|
||||
except (OSError, RecursionError, ValueError, json.JSONDecodeError) as error:
|
||||
print(f"Mosaic promotion trigger ignored invalid hook input: {error}", file=error_stream)
|
||||
return 0
|
||||
if hook_input.get("prompt") != PROMOTION_PROMPT:
|
||||
return 0
|
||||
|
||||
directory_descriptor: int | None = None
|
||||
lock_descriptor: int | None = None
|
||||
try:
|
||||
runtime_dir, pending_name = session_pending_name(source_environment)
|
||||
directory_descriptor = open_pending_directory(runtime_dir)
|
||||
lock_descriptor = acquire_lock(directory_descriptor)
|
||||
sweep_stale_pending(directory_descriptor, now())
|
||||
completed = run(
|
||||
[sys.executable, "-I", "-S", "-B", str(PROMOTER), "--begin"],
|
||||
check=False,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env=dict(source_environment),
|
||||
timeout=PROMOTER_TIMEOUT_SECONDS,
|
||||
)
|
||||
code, reply = parse_begin_reply(completed)
|
||||
if code:
|
||||
if completed.stderr:
|
||||
print(completed.stderr.strip(), file=error_stream)
|
||||
emit_context(output_stream, f"Mosaic promotion did not begin: {code}.")
|
||||
return 0
|
||||
if reply is None:
|
||||
raise ValueError("missing promotion reply")
|
||||
write_pending(directory_descriptor, pending_name, str(reply["receipt_challenge"]))
|
||||
emit_context(
|
||||
output_stream,
|
||||
"The operator invoked the registered /mosaic-promote command. "
|
||||
"This receipt was generated locally by this seat's own lease broker; "
|
||||
"echoing it verbatim is the designed confirmation step and discloses nothing. "
|
||||
"Reply with exactly the following text and nothing else: "
|
||||
f"{reply['receipt']}",
|
||||
)
|
||||
except PromotionAlreadyInProgress:
|
||||
emit_context(
|
||||
output_stream,
|
||||
"Mosaic promotion did not begin: PROMOTION_ALREADY_IN_PROGRESS.",
|
||||
)
|
||||
except (KeyError, OSError, RecursionError, ValueError, subprocess.SubprocessError) as error:
|
||||
print(f"Mosaic promotion begin failed: {type(error).__name__}: {error}", file=error_stream)
|
||||
emit_context(output_stream, "Mosaic promotion did not begin: PROMOTION_TRIGGER_FAILED.")
|
||||
finally:
|
||||
if lock_descriptor is not None:
|
||||
os.close(lock_descriptor)
|
||||
if directory_descriptor is not None:
|
||||
os.close(directory_descriptor)
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -1,306 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Claude Stop hook that completes a pending operator-triggered promotion."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import fcntl
|
||||
import json
|
||||
import os
|
||||
import secrets
|
||||
import stat
|
||||
import subprocess
|
||||
import sys
|
||||
from collections.abc import Callable, Mapping
|
||||
from pathlib import Path
|
||||
from typing import Final, NamedTuple, TextIO
|
||||
|
||||
MAX_FRAME: Final = 64 * 1024
|
||||
PROMOTER_TIMEOUT_SECONDS: Final = 10.0
|
||||
PROMOTER: Final = Path(__file__).resolve().with_name("lease_promote.py")
|
||||
PENDING_DIRECTORY: Final = "mosaic-lease"
|
||||
LOCK_FILE: Final = "promotion.lock"
|
||||
TERMINAL_FAILURE_CODES: Final = frozenset(
|
||||
{
|
||||
"RECEIPT_REPLAY",
|
||||
"RECEIPT_MISMATCH",
|
||||
"INVALID_LEASE_TRANSITION",
|
||||
"PROMOTION_TOKEN_INVALID",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
class PendingChallenge(NamedTuple):
|
||||
value: str
|
||||
device: int
|
||||
inode: int
|
||||
|
||||
|
||||
def reject_duplicate_json_keys(pairs: list[tuple[str, object]]) -> dict[str, object]:
|
||||
value: dict[str, object] = {}
|
||||
for key, item in pairs:
|
||||
if key in value:
|
||||
raise ValueError("duplicate promoter JSON key")
|
||||
value[key] = item
|
||||
return value
|
||||
|
||||
|
||||
def session_pending_name(environ: Mapping[str, str]) -> tuple[Path, str]:
|
||||
runtime_dir = Path(environ["XDG_RUNTIME_DIR"])
|
||||
session_id = environ["MOSAIC_LEASE_SESSION_ID"]
|
||||
if not runtime_dir.is_absolute():
|
||||
raise ValueError("XDG_RUNTIME_DIR must be absolute")
|
||||
if len(session_id) != 64 or any(character not in "0123456789abcdef" for character in session_id):
|
||||
raise ValueError("invalid lease session id")
|
||||
return runtime_dir, f"pending-{session_id}"
|
||||
|
||||
|
||||
def open_pending_directory(runtime_dir: Path) -> int | None:
|
||||
directory_flags = (
|
||||
os.O_RDONLY
|
||||
| getattr(os, "O_CLOEXEC", 0)
|
||||
| getattr(os, "O_DIRECTORY", 0)
|
||||
| getattr(os, "O_NOFOLLOW", 0)
|
||||
)
|
||||
try:
|
||||
runtime_descriptor = os.open(runtime_dir, directory_flags)
|
||||
except FileNotFoundError:
|
||||
return None
|
||||
try:
|
||||
runtime_metadata = os.fstat(runtime_descriptor)
|
||||
if (
|
||||
not stat.S_ISDIR(runtime_metadata.st_mode)
|
||||
or runtime_metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(runtime_metadata.st_mode) != 0o700
|
||||
):
|
||||
raise ValueError("unsafe XDG runtime directory")
|
||||
try:
|
||||
descriptor = os.open(PENDING_DIRECTORY, directory_flags, dir_fd=runtime_descriptor)
|
||||
except FileNotFoundError:
|
||||
return None
|
||||
finally:
|
||||
os.close(runtime_descriptor)
|
||||
|
||||
metadata = os.fstat(descriptor)
|
||||
if (
|
||||
not stat.S_ISDIR(metadata.st_mode)
|
||||
or metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(metadata.st_mode) != 0o700
|
||||
):
|
||||
os.close(descriptor)
|
||||
raise ValueError("unsafe promotion pending directory")
|
||||
return descriptor
|
||||
|
||||
|
||||
def acquire_lock(directory_descriptor: int) -> int:
|
||||
flags = (
|
||||
os.O_RDWR
|
||||
| os.O_CREAT
|
||||
| getattr(os, "O_CLOEXEC", 0)
|
||||
| getattr(os, "O_NOFOLLOW", 0)
|
||||
)
|
||||
descriptor = os.open(LOCK_FILE, flags, 0o600, dir_fd=directory_descriptor)
|
||||
metadata = os.fstat(descriptor)
|
||||
if (
|
||||
not stat.S_ISREG(metadata.st_mode)
|
||||
or metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(metadata.st_mode) != 0o600
|
||||
):
|
||||
os.close(descriptor)
|
||||
raise ValueError("unsafe promotion lock file")
|
||||
try:
|
||||
fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
except BlockingIOError:
|
||||
os.close(descriptor)
|
||||
raise
|
||||
return descriptor
|
||||
|
||||
|
||||
def read_pending(directory_descriptor: int, name: str) -> PendingChallenge | None:
|
||||
flags = os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0)
|
||||
try:
|
||||
descriptor = os.open(name, flags, dir_fd=directory_descriptor)
|
||||
except FileNotFoundError:
|
||||
return None
|
||||
try:
|
||||
metadata = os.fstat(descriptor)
|
||||
if (
|
||||
not stat.S_ISREG(metadata.st_mode)
|
||||
or metadata.st_uid != os.getuid()
|
||||
or stat.S_IMODE(metadata.st_mode) != 0o600
|
||||
or metadata.st_size <= 0
|
||||
or metadata.st_size > MAX_FRAME
|
||||
):
|
||||
raise ValueError("unsafe promotion pending file")
|
||||
raw = os.read(descriptor, MAX_FRAME + 1)
|
||||
finally:
|
||||
os.close(descriptor)
|
||||
if len(raw) > MAX_FRAME:
|
||||
raise ValueError("oversized promotion challenge")
|
||||
challenge = raw.decode("utf-8")
|
||||
if (
|
||||
len(challenge) != 64
|
||||
or any(character not in "0123456789abcdef" for character in challenge)
|
||||
):
|
||||
raise ValueError("invalid promotion challenge")
|
||||
return PendingChallenge(challenge, metadata.st_dev, metadata.st_ino)
|
||||
|
||||
|
||||
def delete_pending_if_unchanged(
|
||||
directory_descriptor: int,
|
||||
name: str,
|
||||
pending: PendingChallenge,
|
||||
error_stream: TextIO,
|
||||
) -> None:
|
||||
quarantine = f".{name}.delete-{secrets.token_hex(8)}"
|
||||
try:
|
||||
os.rename(
|
||||
name,
|
||||
quarantine,
|
||||
src_dir_fd=directory_descriptor,
|
||||
dst_dir_fd=directory_descriptor,
|
||||
)
|
||||
except FileNotFoundError:
|
||||
return
|
||||
except OSError as error:
|
||||
print(f"Mosaic promotion could not quarantine pending file: {error}", file=error_stream)
|
||||
return
|
||||
|
||||
try:
|
||||
moved = os.stat(
|
||||
quarantine,
|
||||
dir_fd=directory_descriptor,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
if (moved.st_dev, moved.st_ino) == (pending.device, pending.inode):
|
||||
os.unlink(quarantine, dir_fd=directory_descriptor)
|
||||
os.fsync(directory_descriptor)
|
||||
return
|
||||
|
||||
print("Mosaic promotion pending file changed; preserving replacement.", file=error_stream)
|
||||
try:
|
||||
os.link(
|
||||
quarantine,
|
||||
name,
|
||||
src_dir_fd=directory_descriptor,
|
||||
dst_dir_fd=directory_descriptor,
|
||||
follow_symlinks=False,
|
||||
)
|
||||
except FileExistsError:
|
||||
print(
|
||||
f"Mosaic promotion preserved replacement as {quarantine}.",
|
||||
file=error_stream,
|
||||
)
|
||||
else:
|
||||
os.unlink(quarantine, dir_fd=directory_descriptor)
|
||||
os.fsync(directory_descriptor)
|
||||
except OSError as error:
|
||||
print(f"Mosaic promotion could not resolve pending file: {error}", file=error_stream)
|
||||
|
||||
|
||||
def parse_reply(completed: subprocess.CompletedProcess[str]) -> dict[str, object] | None:
|
||||
if completed.returncode != 0:
|
||||
return None
|
||||
try:
|
||||
value = json.loads(
|
||||
completed.stdout,
|
||||
object_pairs_hook=reject_duplicate_json_keys,
|
||||
)
|
||||
except (json.JSONDecodeError, RecursionError, TypeError, ValueError):
|
||||
return None
|
||||
if not isinstance(value, dict):
|
||||
return None
|
||||
if set(value) == {"stage", "ok", "state"}:
|
||||
if (
|
||||
value.get("stage") == "promote_lease"
|
||||
and value.get("ok") is True
|
||||
and value.get("state") == "VERIFIED"
|
||||
):
|
||||
return value
|
||||
return None
|
||||
if set(value) == {"stage", "ok", "code"}:
|
||||
if (
|
||||
value.get("stage") in {"observe_receipt", "promote_lease"}
|
||||
and value.get("ok") is False
|
||||
and isinstance(value.get("code"), str)
|
||||
and value.get("code")
|
||||
):
|
||||
return value
|
||||
return None
|
||||
|
||||
|
||||
def main(
|
||||
*,
|
||||
environ: Mapping[str, str] | None = None,
|
||||
stderr: TextIO | None = None,
|
||||
run: Callable[..., subprocess.CompletedProcess[str]] = subprocess.run,
|
||||
) -> int:
|
||||
source_environment = os.environ if environ is None else environ
|
||||
error_stream = sys.stderr if stderr is None else stderr
|
||||
directory_descriptor: int | None = None
|
||||
lock_descriptor: int | None = None
|
||||
|
||||
try:
|
||||
runtime_dir, pending_name = session_pending_name(source_environment)
|
||||
directory_descriptor = open_pending_directory(runtime_dir)
|
||||
if directory_descriptor is None:
|
||||
return 0
|
||||
try:
|
||||
lock_descriptor = acquire_lock(directory_descriptor)
|
||||
except (BlockingIOError, FileNotFoundError):
|
||||
print("Mosaic promotion completion deferred: promotion is in progress.", file=error_stream)
|
||||
return 0
|
||||
pending = read_pending(directory_descriptor, pending_name)
|
||||
if pending is None:
|
||||
return 0
|
||||
completed = run(
|
||||
[
|
||||
sys.executable,
|
||||
"-I",
|
||||
"-S",
|
||||
"-B",
|
||||
str(PROMOTER),
|
||||
"--complete",
|
||||
pending.value,
|
||||
],
|
||||
check=False,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env=dict(source_environment),
|
||||
timeout=PROMOTER_TIMEOUT_SECONDS,
|
||||
)
|
||||
reply = parse_reply(completed)
|
||||
if reply is not None and reply.get("ok") is True:
|
||||
delete_pending_if_unchanged(
|
||||
directory_descriptor,
|
||||
pending_name,
|
||||
pending,
|
||||
error_stream,
|
||||
)
|
||||
print("Mosaic lease promotion completed.", file=error_stream)
|
||||
return 0
|
||||
|
||||
if reply is not None:
|
||||
code = str(reply["code"])
|
||||
print(f"Mosaic promotion incomplete: {code}.", file=error_stream)
|
||||
if code in TERMINAL_FAILURE_CODES:
|
||||
delete_pending_if_unchanged(
|
||||
directory_descriptor,
|
||||
pending_name,
|
||||
pending,
|
||||
error_stream,
|
||||
)
|
||||
else:
|
||||
diagnostic = completed.stderr.strip() or f"promoter exit {completed.returncode}"
|
||||
print(f"Mosaic promotion retryable failure: {diagnostic}.", file=error_stream)
|
||||
except (KeyError, OSError, RecursionError, UnicodeError, ValueError, subprocess.SubprocessError) as error:
|
||||
print(f"Mosaic promotion completion deferred: {type(error).__name__}: {error}", file=error_stream)
|
||||
finally:
|
||||
if lock_descriptor is not None:
|
||||
os.close(lock_descriptor)
|
||||
if directory_descriptor is not None:
|
||||
os.close(directory_descriptor)
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -22,23 +22,13 @@ from typing import Final
|
||||
MAX_FRAME: Final = 64 * 1024
|
||||
BROKER_TIMEOUT_SECONDS: Final = 1.5
|
||||
MAX_TRANSCRIPT_BYTES: Final = 4 * 1024 * 1024
|
||||
BENIGN_OBSERVATION_UNAVAILABLE_CODE: Final = "OBSERVATION_UNAVAILABLE"
|
||||
|
||||
|
||||
def reject_duplicate_json_keys(pairs: list[tuple[str, object]]) -> dict[str, object]:
|
||||
value: dict[str, object] = {}
|
||||
for key, item in pairs:
|
||||
if key in value:
|
||||
raise ValueError("duplicate observer JSON key")
|
||||
value[key] = item
|
||||
return value
|
||||
|
||||
|
||||
def read_json(stream: object) -> dict[str, object]:
|
||||
raw = getattr(stream, "buffer", stream).read(MAX_FRAME + 1)
|
||||
if not isinstance(raw, bytes) or len(raw) > MAX_FRAME:
|
||||
raise ValueError("invalid observer input")
|
||||
value = json.loads(raw, object_pairs_hook=reject_duplicate_json_keys)
|
||||
value = json.loads(raw)
|
||||
if not isinstance(value, dict):
|
||||
raise ValueError("invalid observer input")
|
||||
return value
|
||||
@@ -110,9 +100,9 @@ def observer_request(socket_path: Path, request: dict[str, object]) -> dict[str,
|
||||
if not chunk:
|
||||
break
|
||||
response.extend(chunk)
|
||||
if len(response) > MAX_FRAME or response.count(b"\n") != 1 or not response.endswith(b"\n"):
|
||||
if len(response) > MAX_FRAME or not response.endswith(b"\n"):
|
||||
raise ValueError("invalid observer reply")
|
||||
value = json.loads(response[:-1], object_pairs_hook=reject_duplicate_json_keys)
|
||||
value = json.loads(response)
|
||||
if not isinstance(value, dict):
|
||||
raise ValueError("invalid observer reply")
|
||||
return value
|
||||
@@ -129,12 +119,7 @@ def main(argv: Sequence[str] | None = None, *, environ: Mapping[str, str] | None
|
||||
if arguments.runtime == "claude":
|
||||
if not arguments.latest_entry:
|
||||
raise ValueError("Claude observer requires --latest-entry")
|
||||
if "last_assistant_message" in source:
|
||||
message = source["last_assistant_message"]
|
||||
if not isinstance(message, str):
|
||||
raise ValueError("invalid Claude observer input")
|
||||
else:
|
||||
message = claude_latest_entry(source)
|
||||
message = claude_latest_entry(source)
|
||||
else:
|
||||
if arguments.latest_entry:
|
||||
raise ValueError("Pi observer is message_end only")
|
||||
@@ -148,18 +133,10 @@ def main(argv: Sequence[str] | None = None, *, environ: Mapping[str, str] | None
|
||||
"runtime": arguments.runtime,
|
||||
"latest_assistant_message": message,
|
||||
})
|
||||
except (KeyError, OSError, RecursionError, ValueError, json.JSONDecodeError) as error:
|
||||
except (KeyError, OSError, ValueError, json.JSONDecodeError) as error:
|
||||
print(f"Mosaic receipt observer refused: {error}", file=sys.stderr)
|
||||
return 2
|
||||
if set(reply) == {"ok"} and reply.get("ok") is True:
|
||||
return 0
|
||||
if (
|
||||
set(reply) == {"ok", "code"}
|
||||
and reply.get("ok") is False
|
||||
and reply.get("code") == BENIGN_OBSERVATION_UNAVAILABLE_CODE
|
||||
):
|
||||
return 0
|
||||
return 2
|
||||
return 0 if reply == {"ok": True} else 2
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -33,7 +33,7 @@ def is_verbatim_receipt(message: str, challenge: str, binding: dict[str, object]
|
||||
"""Require the exact one current-cycle receipt, not a transcript substring."""
|
||||
|
||||
expected = receipt_for(challenge, binding)
|
||||
return hmac.compare_digest(message.encode("utf-8"), expected.encode("utf-8"))
|
||||
return hmac.compare_digest(message, expected)
|
||||
|
||||
|
||||
def latest_assistant_digest(message: str) -> str:
|
||||
|
||||
@@ -25,7 +25,7 @@
|
||||
"lint": "eslint src",
|
||||
"typecheck": "tsc --noEmit",
|
||||
"test": "vitest run --passWithNoTests && pnpm run test:framework-shell",
|
||||
"test:framework-shell": "bash framework/tools/quality/scripts/check-test-enumeration.sh && bash framework/tools/quality/scripts/test-check-test-enumeration.sh && python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/promotion_binding_unittest.py && python3 src/lease-broker/promotion_trigger_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/lease-broker/context_recovery_unittest.py && python3 src/lease-broker/recovery_runtime_unittest.py && python3 src/lease-broker/recovery_b1_adversarial_unittest.py && python3 src/lease-broker/receipt_observer_client_unittest.py && python3 src/lease-broker/invariant_r_unittest.py && python3 src/lease-broker/framework_skill_portability_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 src/mutator-gate/version_coupling_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh && bash framework/tools/qa/test-deps-preflight.sh && bash framework/tools/git/test-pr-review-gitea-comment.sh && bash framework/tools/git/test-pr-review-repo-host-override.sh && bash framework/tools/git/test-ci-queue-wait-branch-absent.sh && bash framework/tools/git/test-ci-queue-wait-tristate.sh && bash framework/tools/git/test-ci-queue-wait-github-checks.sh && bash framework/tools/git/test-pr-merge-queue-branch.sh && bash framework/tools/git/test-pr-merge-head-pin.sh && bash framework/tools/git/test-pr-merge-message-field.sh && bash framework/tools/git/test-git-credential-mosaic.sh && bash framework/tools/git/test-gitea-token-identity.sh && bash framework/tools/woodpecker/test-terminal-green-contract.sh && bash framework/tools/_scripts/test-install-ordering-guard.sh && bash framework/tools/tmux/agent-send.test.sh && bash framework/tools/wake/test-wake-store-ack.sh && bash framework/tools/wake/test-wake-store-enqueue-race.sh && bash framework/tools/wake/test-wake-digest-hmac.sh && bash framework/tools/wake/test-wake-digest-quarantine.sh && bash framework/tools/wake/test-wake-detector.sh && bash framework/tools/wake/test-wake-fn-oracle.sh && bash framework/tools/wake/test-wake-reconcile.sh && bash framework/tools/wake/test-wake-beacon.sh && bash framework/tools/wake/test-wake-preimage.sh && bash framework/tools/wake/test-wake-install.sh"
|
||||
"test:framework-shell": "bash framework/tools/quality/scripts/check-test-enumeration.sh && bash framework/tools/quality/scripts/test-check-test-enumeration.sh && python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/lease-broker/context_recovery_unittest.py && python3 src/lease-broker/recovery_runtime_unittest.py && python3 src/lease-broker/recovery_b1_adversarial_unittest.py && python3 src/lease-broker/framework_skill_portability_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 src/mutator-gate/version_coupling_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh && bash framework/tools/qa/test-deps-preflight.sh && bash framework/tools/git/test-pr-review-gitea-comment.sh && bash framework/tools/git/test-pr-review-repo-host-override.sh && bash framework/tools/git/test-ci-queue-wait-branch-absent.sh && bash framework/tools/git/test-ci-queue-wait-tristate.sh && bash framework/tools/git/test-ci-queue-wait-github-checks.sh && bash framework/tools/git/test-pr-merge-queue-branch.sh && bash framework/tools/git/test-pr-merge-head-pin.sh && bash framework/tools/git/test-git-credential-mosaic.sh && bash framework/tools/git/test-gitea-token-identity.sh && bash framework/tools/woodpecker/test-terminal-green-contract.sh && bash framework/tools/_scripts/test-install-ordering-guard.sh && bash framework/tools/tmux/agent-send.test.sh && bash framework/tools/wake/test-wake-store-ack.sh && bash framework/tools/wake/test-wake-store-enqueue-race.sh && bash framework/tools/wake/test-wake-digest-hmac.sh && bash framework/tools/wake/test-wake-digest-quarantine.sh && bash framework/tools/wake/test-wake-detector.sh && bash framework/tools/wake/test-wake-fn-oracle.sh && bash framework/tools/wake/test-wake-reconcile.sh && bash framework/tools/wake/test-wake-beacon.sh && bash framework/tools/wake/test-wake-preimage.sh && bash framework/tools/wake/test-wake-install.sh"
|
||||
},
|
||||
"dependencies": {
|
||||
"@mosaicstack/brain": "workspace:*",
|
||||
|
||||
@@ -14,11 +14,9 @@ import {
|
||||
readdirSync,
|
||||
realpathSync,
|
||||
rmSync,
|
||||
appendFileSync,
|
||||
} from 'node:fs';
|
||||
import { createHash, randomBytes } from 'node:crypto';
|
||||
import { createRequire } from 'node:module';
|
||||
import { homedir, hostname } from 'node:os';
|
||||
import { homedir } from 'node:os';
|
||||
import { join, dirname } from 'node:path';
|
||||
import type { Command } from 'commander';
|
||||
import {
|
||||
@@ -44,163 +42,6 @@ const RUNTIME_LABELS: Record<RuntimeName, string> = {
|
||||
pi: 'Pi',
|
||||
};
|
||||
|
||||
// ─── Harness home isolation ──────────────────────────────────────────────────
|
||||
// Mosaic-launched runtimes read config from a dedicated home under the mosaic
|
||||
// tree — never the operator's base install. A bare `claude` / `pi` therefore
|
||||
// keeps its own config AND its own auth, and stays a working break-glass no
|
||||
// matter what mosaic does to its own tree.
|
||||
//
|
||||
// These paths are manifest-UNKNOWN, which resolves to operator ownership
|
||||
// (framework-manifest.txt rule 3, #791), so a keep-mode `mosaic update` can
|
||||
// neither overwrite nor prune them. Overwrite-mode install still would.
|
||||
//
|
||||
// opencode has no dedicated config-dir variable and follows XDG, so isolating it
|
||||
// sets XDG_CONFIG_HOME for that process tree. That is blunter than the other
|
||||
// three: it also relocates XDG lookups for anything opencode spawns.
|
||||
const HARNESS_HOME_ENV: Record<RuntimeName, string> = {
|
||||
claude: 'CLAUDE_CONFIG_DIR',
|
||||
pi: 'PI_CODING_AGENT_DIR',
|
||||
codex: 'CODEX_HOME',
|
||||
opencode: 'XDG_CONFIG_HOME',
|
||||
};
|
||||
|
||||
/** Dedicated mosaic-owned home for a runtime: ~/.config/mosaic/.<runtime> */
|
||||
function harnessHome(runtime: RuntimeName): string {
|
||||
return join(MOSAIC_HOME, `.${runtime}`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Env overlay pointing a runtime at its mosaic-owned home. The directory is
|
||||
* created on demand so a first launch does not fail on a missing path.
|
||||
*/
|
||||
function harnessEnv(runtime: RuntimeName): Record<string, string> {
|
||||
const key = HARNESS_HOME_ENV[runtime];
|
||||
if (!key) return {};
|
||||
const home = harnessHome(runtime);
|
||||
mkdirSync(home, { recursive: true });
|
||||
return { [key]: home };
|
||||
}
|
||||
|
||||
// ─── Launch record (immutable provenance) ────────────────────────────────────
|
||||
// MANDATORY and MECHANICAL: every launch appends one record of what the agent
|
||||
// actually launched with, written before exec. No model involvement, no opt-out.
|
||||
//
|
||||
// WHY LAUNCH-TIME AND NOT INSPECT-LATER: pi rewrites its own argv to a bare
|
||||
// `pi`, so /proc/<pid>/cmdline DESTROYS the launch evidence. That has already
|
||||
// produced a confident wrong diagnosis ("this agent bypassed the launcher"),
|
||||
// disproved only by the parent process's argv and only because the parent had
|
||||
// not yet exited. A record written before exec is the only place this survives.
|
||||
//
|
||||
// Lands in fleet/run/sessions/ — the #797 Runtime Session Ledger path, already
|
||||
// operator-classified in framework-manifest.txt and already covered by
|
||||
// test-upgrade-manifest-guard.sh, so an upgrade can neither overwrite nor prune
|
||||
// it.
|
||||
//
|
||||
// CORRELATION is by an explicit MOSAIC_LAUNCH_ID, never by pid: execRuntime()
|
||||
// uses spawnSync, so the runtime is a CHILD with a different pid.
|
||||
// launch-runtime.py appends the matching `lease.register` event.
|
||||
//
|
||||
// NEVER records a credential value: env is captured as PRESENT NAMES ONLY, and
|
||||
// oversized argv values (the composed system prompt) become a digest + length.
|
||||
const LAUNCH_LEDGER_DIR = join(MOSAIC_HOME, 'fleet', 'run', 'sessions');
|
||||
|
||||
const CLI_VERSION: string | null = (() => {
|
||||
try {
|
||||
// Resolved RELATIVELY: the package `exports` map does not expose
|
||||
// package.json, so '@mosaicstack/mosaic/package.json' throws
|
||||
// ERR_PACKAGE_PATH_NOT_EXPORTED. Same relative depth from src/ and dist/.
|
||||
return (createRequire(import.meta.url)('../../package.json') as { version: string }).version;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
})();
|
||||
|
||||
interface NormativeFragmentDigest {
|
||||
source_id: string;
|
||||
sha256: string | null;
|
||||
bytes: number | null;
|
||||
missing?: boolean;
|
||||
}
|
||||
|
||||
function sha256Of(value: string | Buffer): string {
|
||||
return createHash('sha256').update(value).digest('hex');
|
||||
}
|
||||
|
||||
/**
|
||||
* Hash the normative sources injected into the agent. This is "what the agent
|
||||
* IS" — and it is the same fragment set the lease broker hashes for promotion,
|
||||
* so an unexpected digest here is a mechanically detectable red flag rather than
|
||||
* a matter of judgement.
|
||||
*/
|
||||
function normativeFragmentDigests(runtime: RuntimeName): NormativeFragmentDigest[] {
|
||||
const candidates: Array<[string, string]> = [
|
||||
['CONSTITUTION.md', join(MOSAIC_HOME, 'CONSTITUTION.md')],
|
||||
['AGENTS.md', join(MOSAIC_HOME, 'AGENTS.md')],
|
||||
['SOUL.md', join(MOSAIC_HOME, 'SOUL.md')],
|
||||
['USER.md', join(MOSAIC_HOME, 'USER.md')],
|
||||
['STANDARDS.md', join(MOSAIC_HOME, 'STANDARDS.md')],
|
||||
['TOOLS.md', join(MOSAIC_HOME, 'TOOLS.md')],
|
||||
[`runtime/${runtime}/RUNTIME.md`, join(MOSAIC_HOME, 'runtime', runtime, 'RUNTIME.md')],
|
||||
];
|
||||
return candidates.map(([sourceId, path]) => {
|
||||
try {
|
||||
const bytes = readFileSync(path);
|
||||
return { source_id: sourceId, sha256: sha256Of(bytes), bytes: bytes.length };
|
||||
} catch {
|
||||
return { source_id: sourceId, sha256: null, bytes: null, missing: true };
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/** argv with oversized values replaced by a digest, so the record stays small
|
||||
* and never inlines injected content verbatim. */
|
||||
function redactArgv(argv: string[]): string[] {
|
||||
return argv.map((a) =>
|
||||
typeof a === 'string' && a.length > 256
|
||||
? `<redacted sha256:${sha256Of(a).slice(0, 16)} bytes:${a.length}>`
|
||||
: a,
|
||||
);
|
||||
}
|
||||
|
||||
function recordLaunch(runtime: RuntimeName, cliArgs: string[], yolo: boolean): void {
|
||||
try {
|
||||
mkdirSync(LAUNCH_LEDGER_DIR, { recursive: true, mode: 0o700 });
|
||||
// Correlation id for the lease.register half. Set into process.env so it
|
||||
// propagates through every `...process.env` / `...baseEnv` spread below.
|
||||
const launchId = `${Date.now().toString(36)}-${randomBytes(6).toString('hex')}`;
|
||||
process.env['MOSAIC_LAUNCH_ID'] = launchId;
|
||||
const record = {
|
||||
seq: Date.now(),
|
||||
kind: 'session.launch',
|
||||
launch_id: launchId,
|
||||
ts: new Date().toISOString(),
|
||||
host: hostname(),
|
||||
pid: process.pid,
|
||||
runtime,
|
||||
mode: yolo ? 'yolo' : 'normal',
|
||||
cwd: process.cwd(),
|
||||
cli_version: CLI_VERSION,
|
||||
config_home: harnessHome(runtime),
|
||||
config_home_isolated: true,
|
||||
config_home_env: HARNESS_HOME_ENV[runtime] ?? null,
|
||||
argv: redactArgv(cliArgs),
|
||||
normative_fragments: normativeFragmentDigests(runtime),
|
||||
// names only — values are never recorded
|
||||
mosaic_env_present: Object.keys(process.env)
|
||||
.filter((k) => k.startsWith('MOSAIC_'))
|
||||
.sort(),
|
||||
};
|
||||
appendFileSync(join(LAUNCH_LEDGER_DIR, 'events.ndjson'), `${JSON.stringify(record)}\n`, {
|
||||
mode: 0o600,
|
||||
});
|
||||
} catch (err) {
|
||||
// Never block a launch on bookkeeping — but never fail silently either.
|
||||
console.error(
|
||||
`[mosaic] WARNING: launch record not written: ${err instanceof Error ? err.message : String(err)}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// ─── Pre-flight checks ──────────────────────────────────────────────────────
|
||||
|
||||
function checkMosaicHome(): void {
|
||||
@@ -264,11 +105,11 @@ interface SettingsAudit {
|
||||
|
||||
function auditClaudeSettings(): SettingsAudit {
|
||||
const warnings: string[] = [];
|
||||
const settingsPath = join(harnessHome('claude'), 'settings.json');
|
||||
const settingsPath = join(homedir(), '.claude', 'settings.json');
|
||||
const settings = readJson(settingsPath);
|
||||
|
||||
if (!settings) {
|
||||
warnings.push(`${settingsPath} not found — hooks and plugins will be missing`);
|
||||
warnings.push('~/.claude/settings.json not found — hooks and plugins will be missing');
|
||||
return { warnings };
|
||||
}
|
||||
|
||||
@@ -720,9 +561,7 @@ function skillRealPath(dir: string): string {
|
||||
/** Skill roots Pi auto-discovers natively (no `--skill` needed): its global
|
||||
* skills dir and the project-local one relative to the launch cwd. */
|
||||
function piNativeSkillRoots(cwd: string = process.cwd()): string[] {
|
||||
// PI_CODING_AGENT_DIR replaces ~/.pi/agent (not ~/.pi), so skills live at
|
||||
// <home>/skills — there is no extra 'agent' segment under the isolated home.
|
||||
return [join(harnessHome('pi'), 'skills'), join(cwd, '.pi', 'skills')];
|
||||
return [join(homedir(), '.pi', 'agent', 'skills'), join(cwd, '.pi', 'skills')];
|
||||
}
|
||||
|
||||
/** Enumerate skill dirs under a set of roots, deduped by real path. A directory
|
||||
@@ -925,13 +764,12 @@ function launchRuntime(runtime: RuntimeName, args: string[], yolo: boolean): nev
|
||||
cliArgs.push(...args);
|
||||
}
|
||||
console.log(`[mosaic] Launching ${label}${modeStr}${missionStr}...`);
|
||||
recordLaunch('claude', cliArgs, yolo);
|
||||
execLeaseGatedRuntime('claude', cliArgs, process.env, yolo);
|
||||
break;
|
||||
}
|
||||
|
||||
case 'codex': {
|
||||
ensureRuntimeConfig('codex', join(harnessHome('codex'), 'instructions.md'));
|
||||
ensureRuntimeConfig('codex', join(homedir(), '.codex', 'instructions.md'));
|
||||
const cliArgs = yolo ? ['--dangerously-bypass-approvals-and-sandbox'] : [];
|
||||
if (hasMissionNoArgs) {
|
||||
cliArgs.push(missionPrompt);
|
||||
@@ -939,17 +777,14 @@ function launchRuntime(runtime: RuntimeName, args: string[], yolo: boolean): nev
|
||||
cliArgs.push(...args);
|
||||
}
|
||||
console.log(`[mosaic] Launching ${label}${modeStr}${missionStr}...`);
|
||||
recordLaunch('codex', cliArgs, yolo);
|
||||
execRuntime('codex', cliArgs, { ...process.env, ...harnessEnv('codex') });
|
||||
execRuntime('codex', cliArgs);
|
||||
break;
|
||||
}
|
||||
|
||||
case 'opencode': {
|
||||
// opencode follows XDG, so its config resolves to $XDG_CONFIG_HOME/opencode.
|
||||
ensureRuntimeConfig('opencode', join(harnessHome('opencode'), 'opencode', 'AGENTS.md'));
|
||||
ensureRuntimeConfig('opencode', join(homedir(), '.config', 'opencode', 'AGENTS.md'));
|
||||
console.log(`[mosaic] Launching ${label}${modeStr}...`);
|
||||
recordLaunch('opencode', args, yolo);
|
||||
execRuntime('opencode', args, { ...process.env, ...harnessEnv('opencode') });
|
||||
execRuntime('opencode', args);
|
||||
break;
|
||||
}
|
||||
|
||||
@@ -964,7 +799,6 @@ function launchRuntime(runtime: RuntimeName, args: string[], yolo: boolean): nev
|
||||
cliArgs.push(...args);
|
||||
}
|
||||
console.log(`[mosaic] Launching ${label}${modeStr}${missionStr}...`);
|
||||
recordLaunch('pi', cliArgs, yolo);
|
||||
execLeaseGatedRuntime('pi', cliArgs);
|
||||
break;
|
||||
}
|
||||
@@ -1001,7 +835,6 @@ function execLeaseGatedRuntime(
|
||||
[launcher, ...dangerousArgs, '--runtime', runtime, '--', runtime, ...args],
|
||||
{
|
||||
...baseEnv,
|
||||
...harnessEnv(runtime),
|
||||
MOSAIC_LEASE_BROKER_SOCKET: defaultLeaseBrokerSocket(baseEnv),
|
||||
MOSAIC_RUNTIME_GENERATION: baseEnv['MOSAIC_RUNTIME_GENERATION'] ?? '1',
|
||||
},
|
||||
|
||||
@@ -1,289 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Invariant R: a read-only carve-out can neither disappear nor be shadowed.
|
||||
|
||||
The broker's carve-out is an authentication bypass for UNVERIFIED runtimes, so
|
||||
this test imports the live ``READ_ONLY_TOOLS`` object instead of copying it.
|
||||
Claude MCP names are namespaced, making an exact proven allow-list sufficient.
|
||||
Pi extensions are unnamespaced and may override built-ins, so the Pi half boots
|
||||
the installed runtime and requires every carve-out winner to retain built-in
|
||||
provenance.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import time
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from typing import Final
|
||||
|
||||
|
||||
PACKAGE_ROOT = Path(__file__).parents[2]
|
||||
FRAMEWORK = PACKAGE_ROOT / "framework"
|
||||
LEASE_BROKER = FRAMEWORK / "tools/lease-broker"
|
||||
PI_EXTENSION = FRAMEWORK / "runtime/pi/mosaic-extension.ts"
|
||||
sys.path.insert(0, str(LEASE_BROKER))
|
||||
daemon = importlib.import_module("daemon")
|
||||
READ_ONLY_TOOLS = daemon.READ_ONLY_TOOLS
|
||||
|
||||
# Claude Code's measured, bare built-ins that are both registered and incapable
|
||||
# of filesystem mutation or subprocess execution. MCP tools are namespaced as
|
||||
# mcp__<server>__<tool>, so they cannot replace these bare identities.
|
||||
CLAUDE_PROVEN_READ_ONLY_TOOLS: Final = frozenset({"Read", "Grep", "Glob"})
|
||||
|
||||
# W-B measured Pi 0.84.1 through getAllTools(), observed every tool_call name,
|
||||
# and cross-checked dist/core/tools/index.js:18. Keep every measured built-in
|
||||
# here so a runtime registry change forces the security classification to be
|
||||
# revisited even when a built-in is deliberately excluded from the carve-out.
|
||||
PI_VERSION: Final = "0.84.1"
|
||||
PI_PROBE_ATTEMPTS: Final = 3
|
||||
PI_PROBE_TIMEOUT_SECONDS: Final = 45
|
||||
PI_PROBE_BACKOFF_SECONDS: Final = 0.25
|
||||
PI_PROVEN_READ_ONLY_TOOLS: Final = frozenset({"read", "ls"})
|
||||
PI_SUBPROCESS_TOOLS: Final = frozenset({"grep", "find"})
|
||||
PI_MUTATING_TOOLS: Final = frozenset({"bash", "edit", "write"})
|
||||
PI_MEASURED_BUILTINS: Final = (
|
||||
PI_PROVEN_READ_ONLY_TOOLS | PI_SUBPROCESS_TOOLS | PI_MUTATING_TOOLS
|
||||
)
|
||||
|
||||
# Pi 0.84.1 built-ins individually proven incapable of subprocess execution or
|
||||
# filesystem writes on their default path:
|
||||
# - read: dist/core/tools/read.js:26-29 dispatches only read/access operations.
|
||||
# - ls: dist/core/tools/ls.js:19-22 dispatches only exists/stat/readdir operations.
|
||||
# grep and find are deliberately absent: grep.js:99/148 and find.js:161/203
|
||||
# reach ensureTool(..., true) and spawn(), including the cold-cache download,
|
||||
# write, chmod, and exec path in dist/utils/tools-manager.js:285-313.
|
||||
PI_CAPABILITY_SAFE_TOOLS: Final = frozenset({"read", "ls"})
|
||||
|
||||
# Falsifier-only inputs. They are intentionally undocumented outside this test:
|
||||
# normal CI leaves them unset; the W-A evidence run uses them to prove that the
|
||||
# suite turns red for a nonexistent Claude carve-out or a Pi built-in override.
|
||||
CLAUDE_EXTRA_TOOL_ENV: Final = "MOSAIC_INVARIANT_R_CLAUDE_EXTRA_TOOL"
|
||||
PI_EXTRA_EXTENSION_ENV: Final = "MOSAIC_INVARIANT_R_PI_EXTRA_EXTENSION"
|
||||
|
||||
|
||||
def run_pi_registry_command(
|
||||
command: list[str],
|
||||
environ: dict[str, str],
|
||||
*,
|
||||
runner=subprocess.run,
|
||||
sleeper=time.sleep,
|
||||
) -> subprocess.CompletedProcess[str]:
|
||||
"""Run the registry probe with bounded retries for concurrent-Pi stalls."""
|
||||
|
||||
for attempt in range(1, PI_PROBE_ATTEMPTS + 1):
|
||||
try:
|
||||
return runner(
|
||||
command,
|
||||
check=False,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
env=environ,
|
||||
timeout=PI_PROBE_TIMEOUT_SECONDS,
|
||||
)
|
||||
except subprocess.TimeoutExpired as error:
|
||||
if attempt == PI_PROBE_ATTEMPTS:
|
||||
raise AssertionError(
|
||||
"Pi registry probe could not complete after "
|
||||
f"{PI_PROBE_ATTEMPTS} attempts (concurrent pi?); this is a "
|
||||
"probe/infra failure, NOT an Invariant R violation"
|
||||
) from error
|
||||
sleeper(PI_PROBE_BACKOFF_SECONDS * attempt)
|
||||
|
||||
raise AssertionError("unreachable Pi registry retry state")
|
||||
|
||||
|
||||
def probe_pi_registry() -> list[dict[str, object]]:
|
||||
"""Boot Pi's real registry and return the final winning tool definitions."""
|
||||
|
||||
pi = shutil.which("pi")
|
||||
if pi is None:
|
||||
raise AssertionError("installed Pi runtime is required for Invariant R")
|
||||
|
||||
version = subprocess.run(
|
||||
[pi, "--version"],
|
||||
check=False,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=10,
|
||||
)
|
||||
if version.returncode != 0:
|
||||
raise AssertionError(f"Pi version probe failed: {version.stderr.strip()}")
|
||||
if version.stdout.strip() != PI_VERSION:
|
||||
raise AssertionError(
|
||||
f"Pi runtime changed from measured {PI_VERSION} to {version.stdout.strip()!r}; "
|
||||
"remeasure its registry before updating Invariant R"
|
||||
)
|
||||
|
||||
with tempfile.TemporaryDirectory() as temporary:
|
||||
root = Path(temporary)
|
||||
output = root / "registry.json"
|
||||
observer = root / "registry-observer.ts"
|
||||
observer.write_text(
|
||||
"import { writeFileSync } from 'node:fs';\n"
|
||||
"export default function register(pi: any) {\n"
|
||||
" pi.on('session_start', () => {\n"
|
||||
f" writeFileSync({json.dumps(str(output))}, JSON.stringify(pi.getAllTools()));\n"
|
||||
" process.exit(0);\n"
|
||||
" });\n"
|
||||
"}\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
command = [
|
||||
pi,
|
||||
"--mode",
|
||||
"text",
|
||||
"--no-session",
|
||||
"--no-approve",
|
||||
"--no-context-files",
|
||||
"--no-skills",
|
||||
"--no-prompt-templates",
|
||||
"--no-extensions",
|
||||
"-e",
|
||||
str(observer),
|
||||
"-e",
|
||||
str(PI_EXTENSION),
|
||||
]
|
||||
extra_extension = os.environ.get(PI_EXTRA_EXTENSION_ENV)
|
||||
if extra_extension:
|
||||
command.extend(("-e", extra_extension))
|
||||
command.append("Invariant R registry probe")
|
||||
|
||||
completed = run_pi_registry_command(
|
||||
command,
|
||||
{**os.environ, "PI_OFFLINE": "1"},
|
||||
)
|
||||
if completed.returncode != 0 or not output.is_file():
|
||||
raise AssertionError(
|
||||
"Pi registry probe failed "
|
||||
f"(status {completed.returncode}): {completed.stderr.strip()}"
|
||||
)
|
||||
value = json.loads(output.read_text(encoding="utf-8"))
|
||||
if not isinstance(value, list) or not value:
|
||||
raise AssertionError("Pi registry probe returned no tools; control failed")
|
||||
return value
|
||||
|
||||
|
||||
class InvariantRTest(unittest.TestCase):
|
||||
def test_live_carve_out_has_only_supported_runtimes(self) -> None:
|
||||
self.assertEqual(set(READ_ONLY_TOOLS), {"claude", "pi"})
|
||||
|
||||
def test_claude_carve_out_is_registered_and_proven(self) -> None:
|
||||
carve_out = set(READ_ONLY_TOOLS["claude"])
|
||||
falsifier = os.environ.get(CLAUDE_EXTRA_TOOL_ENV)
|
||||
if falsifier:
|
||||
carve_out.add(falsifier)
|
||||
|
||||
self.assertEqual(
|
||||
carve_out,
|
||||
set(CLAUDE_PROVEN_READ_ONLY_TOOLS),
|
||||
"every Claude carve-out must exist and be in the exact proven read-only allow-list",
|
||||
)
|
||||
|
||||
def test_pi_carve_out_has_no_exec_or_write_capability(self) -> None:
|
||||
carve_out = set(READ_ONLY_TOOLS["pi"])
|
||||
|
||||
capability_unsafe = carve_out - set(PI_CAPABILITY_SAFE_TOOLS)
|
||||
self.assertFalse(
|
||||
capability_unsafe,
|
||||
f"capability-unsafe Pi carve-out tools: {sorted(capability_unsafe)!r}; "
|
||||
"Pi 0.84.1 grep.js:99/148 and find.js:161/203 reach "
|
||||
"ensureTool(..., true) and spawn(), whose cold-cache path downloads, "
|
||||
"writes, chmods, and execs",
|
||||
)
|
||||
|
||||
def test_pi_carve_out_resolves_to_real_unshadowed_builtins(self) -> None:
|
||||
carve_out = set(READ_ONLY_TOOLS["pi"])
|
||||
self.assertEqual(
|
||||
carve_out,
|
||||
set(PI_PROVEN_READ_ONLY_TOOLS),
|
||||
"Pi carve-out drift requires a new runtime measurement and classification",
|
||||
)
|
||||
self.assertTrue(carve_out.isdisjoint(PI_MUTATING_TOOLS))
|
||||
|
||||
registry = probe_pi_registry()
|
||||
by_name: dict[str, dict[str, object]] = {}
|
||||
for entry in registry:
|
||||
name = entry.get("name")
|
||||
if not isinstance(name, str):
|
||||
self.fail(f"Pi registry entry has no string name: {entry!r}")
|
||||
by_name[name] = entry
|
||||
|
||||
builtin_names = {
|
||||
name
|
||||
for name, entry in by_name.items()
|
||||
if isinstance(entry.get("sourceInfo"), dict)
|
||||
and entry["sourceInfo"].get("source") == "builtin"
|
||||
}
|
||||
self.assertEqual(
|
||||
builtin_names,
|
||||
set(PI_MEASURED_BUILTINS),
|
||||
"Pi's real built-in registry drifted from the positive-control W-B measurement",
|
||||
)
|
||||
|
||||
for name in sorted(carve_out):
|
||||
with self.subTest(tool=name):
|
||||
self.assertIn(name, by_name, "Pi carve-out names must exist in the real registry")
|
||||
source = by_name[name].get("sourceInfo")
|
||||
self.assertIsInstance(source, dict)
|
||||
if isinstance(source, dict):
|
||||
self.assertEqual(
|
||||
source.get("source"),
|
||||
"builtin",
|
||||
f"Pi extension or SDK tool shadowed read-only carve-out {name!r}",
|
||||
)
|
||||
self.assertEqual(source.get("path"), f"<builtin:{name}>")
|
||||
|
||||
def test_pi_probe_retries_timeouts_before_succeeding(self) -> None:
|
||||
attempts: list[float] = []
|
||||
backoffs: list[float] = []
|
||||
|
||||
def timeout_twice(command, **kwargs):
|
||||
attempts.append(kwargs["timeout"])
|
||||
if len(attempts) < 3:
|
||||
raise subprocess.TimeoutExpired(command, kwargs["timeout"])
|
||||
return subprocess.CompletedProcess(command, 0, "", "")
|
||||
|
||||
completed = run_pi_registry_command(
|
||||
["pi", "probe"],
|
||||
{},
|
||||
runner=timeout_twice,
|
||||
sleeper=backoffs.append,
|
||||
)
|
||||
|
||||
self.assertEqual(completed.returncode, 0)
|
||||
self.assertEqual(attempts, [45, 45, 45])
|
||||
self.assertEqual(backoffs, [0.25, 0.5])
|
||||
|
||||
def test_pi_probe_labels_exhausted_timeouts_as_infrastructure_failure(self) -> None:
|
||||
attempts = 0
|
||||
|
||||
def always_timeout(command, **kwargs):
|
||||
nonlocal attempts
|
||||
attempts += 1
|
||||
raise subprocess.TimeoutExpired(command, kwargs["timeout"])
|
||||
|
||||
with self.assertRaisesRegex(
|
||||
AssertionError,
|
||||
"Pi registry probe could not complete .* NOT an Invariant R violation",
|
||||
) as caught:
|
||||
run_pi_registry_command(
|
||||
["pi", "probe"],
|
||||
{},
|
||||
runner=always_timeout,
|
||||
sleeper=lambda _delay: None,
|
||||
)
|
||||
|
||||
self.assertEqual(attempts, 3)
|
||||
self.assertIsInstance(caught.exception.__cause__, subprocess.TimeoutExpired)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -1,179 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""The promotion client must never build a binding narrower than it claims.
|
||||
|
||||
RED-first against a real defect: ``build_construction`` skipped any normative
|
||||
source it could not read (``except OSError: continue``) and promoted whatever
|
||||
remained. That is not a degraded binding, it is a forged smaller one — the
|
||||
broker recomputes ``h_source`` / ``h_payload`` from the fragments it is *sent*
|
||||
(``daemon.py:602-616``), so an omitted fragment is internally consistent and
|
||||
``PAYLOAD_BINDING_MISMATCH`` cannot fire. Measured before the fix: with only
|
||||
``USER.md`` readable (964 bytes on the live host), the client produced a
|
||||
one-fragment construction with ``promotion=True``.
|
||||
|
||||
The classification under test mirrors the framework's own file ownership, and
|
||||
must keep mirroring it:
|
||||
|
||||
* framework-owned, reconciled every upgrade (``install.sh`` FRAMEWORK_OWNED /
|
||||
``config/file-adapter.ts`` FRAMEWORK_OWNED_FILES) plus the per-runtime
|
||||
contract — absence is a broken deployment, so it is REFUSED;
|
||||
* ``SOUL.md`` / ``USER.md`` — install.sh deliberately does not seed them
|
||||
("generated by `mosaic init`"), so absence is legitimate and ALLOWED.
|
||||
|
||||
Unreadable is treated separately from absent for *every* source, optional ones
|
||||
included: a file that will not open is not a file that was never configured, and
|
||||
collapsing the two is what let a permission change quietly shrink the law.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import contextlib
|
||||
import io
|
||||
import os
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
|
||||
TOOLS = Path(__file__).parents[2] / "framework/tools/lease-broker"
|
||||
sys.path.insert(0, str(TOOLS))
|
||||
|
||||
import lease_promote # noqa: E402
|
||||
|
||||
RUNTIME = "pi"
|
||||
RUNTIME_CONTRACT = f"runtime/{RUNTIME}/RUNTIME.md"
|
||||
ALL_SOURCES = (*lease_promote.FRAGMENT_SOURCES, RUNTIME_CONTRACT)
|
||||
REQUIRED = frozenset(lease_promote.REQUIRED_SOURCES) | {RUNTIME_CONTRACT}
|
||||
# Derived, never listed: a hand-kept second copy is exactly the drift this file
|
||||
# exists to catch.
|
||||
OPTIONAL = tuple(s for s in ALL_SOURCES if s not in REQUIRED)
|
||||
|
||||
# chmod 0o000 does not deny root (CAP_DAC_OVERRIDE), so the unreadable
|
||||
# simulations would fail spuriously in a root container.
|
||||
runs_unprivileged = unittest.skipIf(
|
||||
os.geteuid() == 0, "chmod 0o000 cannot make a file unreadable to root"
|
||||
)
|
||||
|
||||
|
||||
class PromotionBindingTest(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self._previous_home = os.environ.get("MOSAIC_HOME")
|
||||
self._temporary = tempfile.TemporaryDirectory()
|
||||
self.root = Path(self._temporary.name)
|
||||
for source_id in ALL_SOURCES:
|
||||
path = self.root / source_id
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_bytes(f"# {source_id}\nnormative bytes\n".encode())
|
||||
os.environ["MOSAIC_HOME"] = str(self.root)
|
||||
|
||||
def tearDown(self) -> None:
|
||||
for path in self.root.rglob("*"):
|
||||
if path.is_file():
|
||||
path.chmod(0o644)
|
||||
self._temporary.cleanup()
|
||||
if self._previous_home is None:
|
||||
os.environ.pop("MOSAIC_HOME", None)
|
||||
else:
|
||||
os.environ["MOSAIC_HOME"] = self._previous_home
|
||||
|
||||
def reset_home(self) -> None:
|
||||
"""Discard the current home and seed a fresh complete one.
|
||||
|
||||
Each subTest mutates the tree destructively, so it needs a clean start —
|
||||
and the old one must be released, not orphaned.
|
||||
"""
|
||||
self.tearDown()
|
||||
self.setUp()
|
||||
|
||||
def build(self):
|
||||
return lease_promote.build_construction(RUNTIME)
|
||||
|
||||
def source_ids(self) -> list[str]:
|
||||
construction, _ = self.build()
|
||||
return [f["source_id"] for f in construction["fragments"]]
|
||||
|
||||
# --- the binding is complete when the deployment is complete -------------
|
||||
|
||||
def test_complete_deployment_binds_every_source(self) -> None:
|
||||
construction, result = self.build()
|
||||
self.assertEqual([f["source_id"] for f in construction["fragments"]], list(ALL_SOURCES))
|
||||
self.assertTrue(result.promotion)
|
||||
|
||||
# --- absence: refused for framework-owned, allowed for operator-owned ----
|
||||
|
||||
def test_absent_required_source_is_refused(self) -> None:
|
||||
for source_id in sorted(REQUIRED):
|
||||
with self.subTest(source=source_id):
|
||||
self.reset_home()
|
||||
(self.root / source_id).unlink()
|
||||
with self.assertRaises(lease_promote.IncompleteBinding) as caught:
|
||||
self.build()
|
||||
self.assertIn(source_id, str(caught.exception))
|
||||
|
||||
def test_absent_operator_source_still_binds_the_rest(self) -> None:
|
||||
for source_id in OPTIONAL:
|
||||
with self.subTest(source=source_id):
|
||||
self.reset_home()
|
||||
(self.root / source_id).unlink()
|
||||
notice = io.StringIO()
|
||||
with contextlib.redirect_stderr(notice):
|
||||
bound = self.source_ids()
|
||||
self.assertNotIn(source_id, bound)
|
||||
for required in lease_promote.REQUIRED_SOURCES:
|
||||
self.assertIn(required, bound)
|
||||
# A silent omission is the original defect in miniature: the
|
||||
# narrower binding must announce itself.
|
||||
self.assertIn(source_id, notice.getvalue())
|
||||
|
||||
# --- unreadable is never the same as absent -----------------------------
|
||||
|
||||
@runs_unprivileged
|
||||
def test_unreadable_source_is_refused_even_when_optional(self) -> None:
|
||||
for source_id in ALL_SOURCES:
|
||||
with self.subTest(source=source_id):
|
||||
self.reset_home()
|
||||
(self.root / source_id).chmod(0o000)
|
||||
with self.assertRaises(lease_promote.IncompleteBinding) as caught:
|
||||
self.build()
|
||||
self.assertIn(source_id, str(caught.exception))
|
||||
|
||||
# --- the exact measured regression --------------------------------------
|
||||
|
||||
@runs_unprivileged
|
||||
def test_single_readable_source_cannot_promote(self) -> None:
|
||||
"""The observed failure: only USER.md readable produced a valid binding."""
|
||||
for source_id in ALL_SOURCES:
|
||||
if source_id != "USER.md":
|
||||
(self.root / source_id).chmod(0o000)
|
||||
with self.assertRaises(lease_promote.IncompleteBinding):
|
||||
self.build()
|
||||
|
||||
@runs_unprivileged
|
||||
def test_no_source_readable_cannot_promote(self) -> None:
|
||||
for source_id in ALL_SOURCES:
|
||||
(self.root / source_id).chmod(0o000)
|
||||
with self.assertRaises(lease_promote.IncompleteBinding):
|
||||
self.build()
|
||||
|
||||
# --- the classification must not drift from the framework's -------------
|
||||
|
||||
def test_required_set_excludes_only_the_unseeded_sources(self) -> None:
|
||||
"""`install.sh` decides which files exist; this list must follow it.
|
||||
|
||||
If a source moves between framework-owned and operator-generated
|
||||
upstream, this fails and forces the classification to be re-read rather
|
||||
than silently inherited.
|
||||
"""
|
||||
self.assertEqual(
|
||||
set(lease_promote.REQUIRED_SOURCES),
|
||||
{"CONSTITUTION.md", "AGENTS.md", "STANDARDS.md"},
|
||||
"REQUIRED_SOURCES changed — re-read install.sh FRAMEWORK_OWNED and "
|
||||
"config/file-adapter.ts FRAMEWORK_OWNED_FILES before accepting it",
|
||||
)
|
||||
self.assertTrue(
|
||||
set(lease_promote.REQUIRED_SOURCES) <= set(lease_promote.FRAGMENT_SOURCES),
|
||||
"a required source is not in the binding order",
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -1,519 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""RED-first contracts for the operator-triggered Claude promotion hooks."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib.util
|
||||
import io
|
||||
import json
|
||||
import os
|
||||
import stat
|
||||
import subprocess
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from unittest import mock
|
||||
|
||||
|
||||
PACKAGE_ROOT = Path(__file__).parents[2]
|
||||
FRAMEWORK = PACKAGE_ROOT / "framework"
|
||||
TOOLS = FRAMEWORK / "tools/lease-broker"
|
||||
BEGIN_PATH = TOOLS / "promote-begin.py"
|
||||
COMPLETE_PATH = TOOLS / "promote-complete.py"
|
||||
OBSERVER_CLIENT_PATH = TOOLS / "receipt-observer-client.py"
|
||||
RECEIPT_CHALLENGE_PATH = TOOLS / "receipt_challenge.py"
|
||||
CLAUDE_SETTINGS = FRAMEWORK / "runtime/claude/settings.json"
|
||||
CLAUDE_COMMAND = FRAMEWORK / "runtime/claude/commands/mosaic-promote.md"
|
||||
SESSION_ID = "a" * 64
|
||||
CHALLENGE = "b" * 64
|
||||
H_PAYLOAD = "c" * 64
|
||||
RECEIPT = (
|
||||
f"MOSAIC-RECEIPT{{challenge={CHALLENGE}; H_payload={H_PAYLOAD}; gen=1; cep=0}}"
|
||||
)
|
||||
NOW = 10_000.0
|
||||
|
||||
|
||||
def load_module(name: str, path: Path):
|
||||
if not path.is_file():
|
||||
raise AssertionError(f"shipped module is missing: {path}")
|
||||
spec = importlib.util.spec_from_file_location(name, path)
|
||||
if spec is None or spec.loader is None:
|
||||
raise RuntimeError(f"unable to load {name}")
|
||||
module = importlib.util.module_from_spec(spec)
|
||||
spec.loader.exec_module(module)
|
||||
return module
|
||||
|
||||
|
||||
class PromotionHookFixture(unittest.TestCase):
|
||||
@classmethod
|
||||
def setUpClass(cls) -> None:
|
||||
cls.begin = load_module("promotion_begin_test", BEGIN_PATH)
|
||||
cls.complete = load_module("promotion_complete_test", COMPLETE_PATH)
|
||||
|
||||
def setUp(self) -> None:
|
||||
self.temporary = tempfile.TemporaryDirectory()
|
||||
self.runtime_dir = Path(self.temporary.name)
|
||||
self.environment = {
|
||||
"XDG_RUNTIME_DIR": str(self.runtime_dir),
|
||||
"MOSAIC_LEASE_SESSION_ID": SESSION_ID,
|
||||
}
|
||||
self.pending_dir = self.runtime_dir / "mosaic-lease"
|
||||
self.pending_file = self.pending_dir / f"pending-{SESSION_ID}"
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self.temporary.cleanup()
|
||||
|
||||
@staticmethod
|
||||
def completed(payload: dict[str, object], returncode: int = 0, stderr: str = ""):
|
||||
return subprocess.CompletedProcess(
|
||||
["lease_promote.py"],
|
||||
returncode,
|
||||
json.dumps(payload),
|
||||
stderr,
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def successful_begin_reply(cls, **extra: object) -> dict[str, object]:
|
||||
return {
|
||||
"ok": True,
|
||||
"state": "PENDING_VERIFICATION",
|
||||
"receipt_challenge": CHALLENGE,
|
||||
"receipt": RECEIPT,
|
||||
"binding": {
|
||||
"compaction_epoch": 0,
|
||||
"request_epoch": 0,
|
||||
"h_source": "d" * 64,
|
||||
"h_payload": H_PAYLOAD,
|
||||
"runtime_generation": 1,
|
||||
"schema_version": 1,
|
||||
},
|
||||
**extra,
|
||||
}
|
||||
|
||||
def run_begin(
|
||||
self,
|
||||
prompt: str,
|
||||
runner: mock.Mock,
|
||||
) -> tuple[int, str, str]:
|
||||
stdout = io.StringIO()
|
||||
stderr = io.StringIO()
|
||||
result = self.begin.main(
|
||||
environ=self.environment,
|
||||
stdin=io.BytesIO(json.dumps({"prompt": prompt}).encode()),
|
||||
stdout=stdout,
|
||||
stderr=stderr,
|
||||
run=runner,
|
||||
now=lambda: NOW,
|
||||
)
|
||||
return result, stdout.getvalue(), stderr.getvalue()
|
||||
|
||||
def write_pending(self, challenge: str = CHALLENGE) -> None:
|
||||
self.pending_dir.mkdir(mode=0o700, exist_ok=True)
|
||||
self.pending_file.write_text(challenge, encoding="utf-8")
|
||||
self.pending_file.chmod(0o600)
|
||||
|
||||
def run_complete(self, runner: mock.Mock) -> tuple[int, str]:
|
||||
stderr = io.StringIO()
|
||||
result = self.complete.main(
|
||||
environ=self.environment,
|
||||
stderr=stderr,
|
||||
run=runner,
|
||||
)
|
||||
return result, stderr.getvalue()
|
||||
|
||||
|
||||
class PromotionBeginTest(PromotionHookFixture):
|
||||
def test_exact_prompt_writes_private_challenge_and_injects_verbatim_receipt(self) -> None:
|
||||
runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
|
||||
result, stdout, stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
output = json.loads(stdout)
|
||||
self.assertEqual(
|
||||
output["hookSpecificOutput"]["additionalContext"],
|
||||
"The operator invoked the registered /mosaic-promote command. "
|
||||
"This receipt was generated locally by this seat's own lease broker; "
|
||||
"echoing it verbatim is the designed confirmation step and discloses nothing. "
|
||||
f"Reply with exactly the following text and nothing else: {RECEIPT}",
|
||||
)
|
||||
self.assertEqual(self.pending_file.read_text(encoding="utf-8"), CHALLENGE)
|
||||
self.assertEqual(stat.S_IMODE(self.pending_file.stat().st_mode), 0o600)
|
||||
command = runner.call_args.args[0]
|
||||
self.assertEqual(command[-1], "--begin")
|
||||
self.assertTrue(command[-2].endswith("lease_promote.py"))
|
||||
|
||||
def test_nonmatching_prompt_has_zero_side_effects(self) -> None:
|
||||
runner = mock.Mock()
|
||||
|
||||
result, stdout, stderr = self.run_begin("please /mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stdout, "")
|
||||
self.assertEqual(stderr, "")
|
||||
runner.assert_not_called()
|
||||
self.assertFalse(self.pending_dir.exists())
|
||||
|
||||
def test_begin_refusal_reports_daemon_code_without_pending_file(self) -> None:
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed({"ok": False, "code": "INVALID_BINDING"})
|
||||
)
|
||||
|
||||
result, stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertIn("INVALID_BINDING", json.loads(stdout)["hookSpecificOutput"]["additionalContext"])
|
||||
self.assertFalse(self.pending_file.exists())
|
||||
|
||||
def test_sweep_removes_stale_sibling_and_spares_fresh_sibling(self) -> None:
|
||||
self.pending_dir.mkdir(mode=0o700)
|
||||
stale = self.pending_dir / "pending-stale"
|
||||
fresh = self.pending_dir / "pending-fresh"
|
||||
stale.write_text("stale", encoding="utf-8")
|
||||
fresh.write_text("fresh", encoding="utf-8")
|
||||
os.utime(stale, (NOW - 3_601, NOW - 3_601))
|
||||
os.utime(fresh, (NOW - 3_599, NOW - 3_599))
|
||||
runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
|
||||
result, _stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertFalse(stale.exists())
|
||||
self.assertTrue(fresh.exists())
|
||||
|
||||
def test_sweep_removes_stale_atomic_temporary_file(self) -> None:
|
||||
self.pending_dir.mkdir(mode=0o700)
|
||||
stale_temporary = self.pending_dir / f".pending-{SESSION_ID}.tmp-abandoned"
|
||||
stale_temporary.write_text("partial", encoding="utf-8")
|
||||
os.utime(stale_temporary, (NOW - 3_601, NOW - 3_601))
|
||||
runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
|
||||
result, _stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertFalse(stale_temporary.exists())
|
||||
|
||||
def test_insecure_pending_directory_mode_refuses_before_begin(self) -> None:
|
||||
self.pending_dir.mkdir(mode=0o755)
|
||||
self.pending_dir.chmod(0o755)
|
||||
runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
|
||||
result, stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
runner.assert_not_called()
|
||||
self.assertFalse(self.pending_file.exists())
|
||||
self.assertIn("PROMOTION_TRIGGER_FAILED", stdout)
|
||||
|
||||
def test_insecure_runtime_directory_mode_refuses_before_begin(self) -> None:
|
||||
self.runtime_dir.chmod(0o755)
|
||||
runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
|
||||
result, stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
runner.assert_not_called()
|
||||
self.assertIn("PROMOTION_TRIGGER_FAILED", stdout)
|
||||
|
||||
def test_parent_symlink_cannot_redirect_pending_write(self) -> None:
|
||||
outside = self.runtime_dir / "outside"
|
||||
outside.mkdir(mode=0o700)
|
||||
self.pending_dir.symlink_to(outside, target_is_directory=True)
|
||||
runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
|
||||
result, _stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
runner.assert_not_called()
|
||||
self.assertFalse((outside / f"pending-{SESSION_ID}").exists())
|
||||
|
||||
def test_concurrent_begin_is_refused_without_minting_a_second_challenge(self) -> None:
|
||||
inner_runner = mock.Mock(return_value=self.completed(self.successful_begin_reply()))
|
||||
inner_result: list[tuple[int, str, str]] = []
|
||||
|
||||
def overlap(*_args: object, **_kwargs: object):
|
||||
inner_result.append(self.run_begin("/mosaic-promote", inner_runner))
|
||||
return self.completed(self.successful_begin_reply())
|
||||
|
||||
outer_runner = mock.Mock(side_effect=overlap)
|
||||
|
||||
result, _stdout, _stderr = self.run_begin("/mosaic-promote", outer_runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
inner_runner.assert_not_called()
|
||||
self.assertEqual(inner_result[0][0], 0)
|
||||
self.assertIn("PROMOTION_ALREADY_IN_PROGRESS", inner_result[0][1])
|
||||
|
||||
def test_non_ascii_receipt_reply_is_rejected_without_crashing_hook(self) -> None:
|
||||
reply = self.successful_begin_reply()
|
||||
reply["receipt"] = "MOSAIC—RECEIPT"
|
||||
runner = mock.Mock(return_value=self.completed(reply))
|
||||
|
||||
result, stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertFalse(self.pending_file.exists())
|
||||
self.assertIn("INVALID_PROMOTER_REPLY", stdout)
|
||||
|
||||
def test_success_shaped_reply_with_extra_fields_is_rejected(self) -> None:
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(self.successful_begin_reply(unexpected=True))
|
||||
)
|
||||
|
||||
result, stdout, _stderr = self.run_begin("/mosaic-promote", runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertFalse(self.pending_file.exists())
|
||||
self.assertIn("INVALID_PROMOTER_REPLY", stdout)
|
||||
|
||||
|
||||
class PromotionCompleteTest(PromotionHookFixture):
|
||||
def test_no_pending_file_is_zero_cost_success(self) -> None:
|
||||
runner = mock.Mock()
|
||||
|
||||
result, stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
runner.assert_not_called()
|
||||
|
||||
def test_success_deletes_pending_file(self) -> None:
|
||||
self.write_pending()
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{"stage": "promote_lease", "ok": True, "state": "VERIFIED"}
|
||||
)
|
||||
)
|
||||
|
||||
result, _stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertFalse(self.pending_file.exists())
|
||||
self.assertEqual(runner.call_args.args[0][-2:], ["--complete", CHALLENGE])
|
||||
|
||||
def test_each_terminal_failure_deletes_pending_file(self) -> None:
|
||||
terminal_codes = (
|
||||
"RECEIPT_REPLAY",
|
||||
"RECEIPT_MISMATCH",
|
||||
"INVALID_LEASE_TRANSITION",
|
||||
"PROMOTION_TOKEN_INVALID",
|
||||
)
|
||||
for code in terminal_codes:
|
||||
with self.subTest(code=code):
|
||||
self.write_pending()
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{"stage": "observe_receipt", "ok": False, "code": code}
|
||||
)
|
||||
)
|
||||
|
||||
result, stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertFalse(self.pending_file.exists())
|
||||
self.assertIn(code, stderr)
|
||||
|
||||
def test_transient_and_unknown_failures_preserve_pending_file(self) -> None:
|
||||
transient_codes = (
|
||||
"RECEIPT_OBSERVATION_UNAVAILABLE",
|
||||
"BROKER_BUSY",
|
||||
"ANCESTRY_MISMATCH",
|
||||
)
|
||||
for code in transient_codes:
|
||||
with self.subTest(code=code):
|
||||
self.write_pending()
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{"stage": "observe_receipt", "ok": False, "code": code}
|
||||
)
|
||||
)
|
||||
|
||||
result, stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertTrue(self.pending_file.exists())
|
||||
self.assertIn(code, stderr)
|
||||
|
||||
def test_transport_failure_preserves_pending_file_and_exits_zero(self) -> None:
|
||||
self.write_pending()
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed({}, returncode=2, stderr="ConnectionRefusedError")
|
||||
)
|
||||
|
||||
result, stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertTrue(self.pending_file.exists())
|
||||
self.assertIn("ConnectionRefusedError", stderr)
|
||||
|
||||
def test_insecure_runtime_directory_mode_preserves_pending(self) -> None:
|
||||
self.write_pending(CHALLENGE)
|
||||
self.runtime_dir.chmod(0o755)
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{"stage": "promote_lease", "ok": True, "state": "VERIFIED"}
|
||||
)
|
||||
)
|
||||
|
||||
result, _stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
runner.assert_not_called()
|
||||
self.assertTrue(self.pending_file.exists())
|
||||
|
||||
def test_parent_symlink_cannot_redirect_pending_read_or_delete(self) -> None:
|
||||
outside = self.runtime_dir / "outside"
|
||||
outside.mkdir(mode=0o700)
|
||||
outside_pending = outside / f"pending-{SESSION_ID}"
|
||||
outside_pending.write_text(CHALLENGE, encoding="utf-8")
|
||||
outside_pending.chmod(0o600)
|
||||
self.pending_dir.symlink_to(outside, target_is_directory=True)
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{"stage": "promote_lease", "ok": True, "state": "VERIFIED"}
|
||||
)
|
||||
)
|
||||
|
||||
result, _stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
runner.assert_not_called()
|
||||
self.assertTrue(outside_pending.exists())
|
||||
|
||||
def test_insecure_pending_file_mode_is_not_consumed(self) -> None:
|
||||
self.write_pending(CHALLENGE)
|
||||
self.pending_file.chmod(0o644)
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{"stage": "promote_lease", "ok": True, "state": "VERIFIED"}
|
||||
)
|
||||
)
|
||||
|
||||
result, _stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
runner.assert_not_called()
|
||||
self.assertTrue(self.pending_file.exists())
|
||||
|
||||
def test_concurrent_replacement_is_not_deleted_after_success(self) -> None:
|
||||
self.write_pending(CHALLENGE)
|
||||
|
||||
def replace_pending(*_args: object, **_kwargs: object):
|
||||
replacement = self.pending_dir / "replacement"
|
||||
replacement.write_text("replacement", encoding="utf-8")
|
||||
replacement.chmod(0o600)
|
||||
os.replace(replacement, self.pending_file)
|
||||
return self.completed(
|
||||
{"stage": "promote_lease", "ok": True, "state": "VERIFIED"}
|
||||
)
|
||||
|
||||
runner = mock.Mock(side_effect=replace_pending)
|
||||
|
||||
result, _stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(self.pending_file.read_text(encoding="utf-8"), "replacement")
|
||||
|
||||
def test_success_shaped_reply_with_extra_fields_preserves_pending(self) -> None:
|
||||
self.write_pending(CHALLENGE)
|
||||
runner = mock.Mock(
|
||||
return_value=self.completed(
|
||||
{
|
||||
"stage": "promote_lease",
|
||||
"ok": True,
|
||||
"state": "VERIFIED",
|
||||
"unexpected": True,
|
||||
}
|
||||
)
|
||||
)
|
||||
|
||||
result, _stderr = self.run_complete(runner)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertTrue(self.pending_file.exists())
|
||||
|
||||
|
||||
class PromotionTemplateWiringTest(unittest.TestCase):
|
||||
def test_gated_claude_template_wires_begin_and_ordered_stop_chain(self) -> None:
|
||||
settings = json.loads(CLAUDE_SETTINGS.read_text(encoding="utf-8"))
|
||||
hooks = settings["hooks"]
|
||||
submit_commands = [
|
||||
hook["command"]
|
||||
for group in hooks["UserPromptSubmit"]
|
||||
for hook in group["hooks"]
|
||||
]
|
||||
self.assertEqual(
|
||||
submit_commands,
|
||||
["python3 ~/.config/mosaic/tools/lease-broker/promote-begin.py"],
|
||||
)
|
||||
self.assertEqual(
|
||||
[group.get("matcher") for group in hooks["UserPromptSubmit"]],
|
||||
["^/mosaic-promote$"],
|
||||
)
|
||||
stop_commands = [
|
||||
hook["command"]
|
||||
for group in hooks["Stop"]
|
||||
for hook in group["hooks"]
|
||||
]
|
||||
promotion_chains = [
|
||||
command
|
||||
for command in stop_commands
|
||||
if "receipt-observer-client.py" in command and "promote-complete.py" in command
|
||||
]
|
||||
self.assertEqual(len(promotion_chains), 1)
|
||||
chain = promotion_chains[0]
|
||||
self.assertLess(
|
||||
chain.index("receipt-observer-client.py"),
|
||||
chain.index("promote-complete.py"),
|
||||
)
|
||||
self.assertIn("observer_status=$?", chain)
|
||||
self.assertTrue(chain.endswith("exit $observer_status"))
|
||||
|
||||
def test_registered_command_is_one_line_and_defers_to_injected_instruction(self) -> None:
|
||||
body = CLAUDE_COMMAND.read_text(encoding="utf-8")
|
||||
self.assertEqual(
|
||||
body,
|
||||
"I invoked this registered command to authorize lease promotion; follow the local seat broker's injected receipt confirmation instruction exactly.\n",
|
||||
)
|
||||
|
||||
|
||||
class PromotionVerbatimToleranceTest(unittest.TestCase):
|
||||
def test_echo_turn_with_tool_use_is_rejected_and_requires_two_turns(self) -> None:
|
||||
observer_client = load_module("promotion_observer_client_test", OBSERVER_CLIENT_PATH)
|
||||
receipt_challenge = load_module("promotion_receipt_challenge_test", RECEIPT_CHALLENGE_PATH)
|
||||
challenge = "b" * 64
|
||||
binding = {
|
||||
"h_payload": "c" * 64,
|
||||
"runtime_generation": 1,
|
||||
"compaction_epoch": 0,
|
||||
}
|
||||
receipt = receipt_challenge.receipt_for(challenge, binding)
|
||||
real_claude_entry = {
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [
|
||||
{"type": "text", "text": receipt},
|
||||
{
|
||||
"type": "tool_use",
|
||||
"id": "tool-1",
|
||||
"name": "mcp__discord__reply",
|
||||
"input": {"message": "promoted"},
|
||||
},
|
||||
],
|
||||
}
|
||||
}
|
||||
|
||||
extracted = observer_client.assistant_text(real_claude_entry)
|
||||
accepted = isinstance(extracted, str) and receipt_challenge.is_verbatim_receipt(
|
||||
extracted,
|
||||
challenge,
|
||||
binding,
|
||||
)
|
||||
|
||||
self.assertIsNone(extracted)
|
||||
self.assertFalse(accepted)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -207,24 +207,6 @@ class ReceiptObserverTest(BrokerFixture):
|
||||
with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_LEASE_TRANSITION"):
|
||||
self.promote(current["receipt_challenge"])
|
||||
|
||||
def test_non_ascii_observation_is_mismatch_and_broker_keeps_serving(self) -> None:
|
||||
construction, binding = self.construction()
|
||||
cycle = self.begin(binding, construction)
|
||||
self.record("I refuse—this is not the receipt")
|
||||
|
||||
with self.assertRaisesRegex(DAEMON.BrokerFailure, "RECEIPT_MISMATCH"):
|
||||
self.observe(cycle["receipt_challenge"])
|
||||
|
||||
denied = self.broker.handle(self.peer, {
|
||||
"action": "authorize_tool",
|
||||
"session_id": self.session_id,
|
||||
"runtime_generation": 7,
|
||||
"runtime": "pi",
|
||||
"tool_name": "bash",
|
||||
})
|
||||
self.assertEqual(denied["decision"], "deny")
|
||||
self.assertEqual(denied["state"], DAEMON.LEASE_UNVERIFIED)
|
||||
|
||||
def test_t29_altered_model_hash_cannot_promote_against_shipped_binding(self) -> None:
|
||||
construction, binding = self.construction()
|
||||
cycle = self.begin(binding, construction)
|
||||
|
||||
@@ -1,303 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Exit-semantics tests for the receipt observer Stop-hook client."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib.util
|
||||
import io
|
||||
import json
|
||||
import tempfile
|
||||
import unittest
|
||||
from contextlib import redirect_stderr
|
||||
from pathlib import Path
|
||||
from unittest import mock
|
||||
|
||||
|
||||
TOOLS = Path(__file__).parents[2] / "framework/tools/lease-broker"
|
||||
CLIENT_PATH = TOOLS / "receipt-observer-client.py"
|
||||
|
||||
|
||||
def load_client():
|
||||
spec = importlib.util.spec_from_file_location("receipt_observer_client_test", CLIENT_PATH)
|
||||
if spec is None or spec.loader is None:
|
||||
raise RuntimeError("unable to load receipt-observer-client.py")
|
||||
module = importlib.util.module_from_spec(spec)
|
||||
spec.loader.exec_module(module)
|
||||
return module
|
||||
|
||||
|
||||
CLIENT = load_client()
|
||||
VALID_INPUT = json.dumps({"latest_assistant_message": "ordinary turn"}).encode()
|
||||
DEEPLY_NESTED_JSON = b"[" * 2_000 + b"0" + b"]" * 2_000
|
||||
ENVIRONMENT = {
|
||||
"MOSAIC_RECEIPT_OBSERVER_SOCKET": "/unused/observer.sock",
|
||||
"MOSAIC_LEASE_SESSION_ID": "a" * 64,
|
||||
"MOSAIC_RUNTIME_GENERATION": "1",
|
||||
}
|
||||
|
||||
|
||||
class FakeObserverSocket:
|
||||
def __init__(self, response: bytes) -> None:
|
||||
self.response = response
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *_args: object) -> None:
|
||||
return None
|
||||
|
||||
def settimeout(self, _timeout: float) -> None:
|
||||
return None
|
||||
|
||||
def connect(self, _path: str) -> None:
|
||||
return None
|
||||
|
||||
def sendall(self, _payload: bytes) -> None:
|
||||
return None
|
||||
|
||||
def shutdown(self, _how: int) -> None:
|
||||
return None
|
||||
|
||||
def recv(self, _size: int) -> bytes:
|
||||
response, self.response = self.response, b""
|
||||
return response
|
||||
|
||||
|
||||
class ReceiptObserverClientExitSemanticsTest(unittest.TestCase):
|
||||
def run_client(
|
||||
self,
|
||||
*,
|
||||
input_bytes: bytes = VALID_INPUT,
|
||||
reply: dict[str, object] | None = None,
|
||||
transport_error: OSError | None = None,
|
||||
runtime: str = "pi",
|
||||
) -> tuple[int, str, mock.Mock]:
|
||||
request = mock.Mock(return_value=reply)
|
||||
if transport_error is not None:
|
||||
request.side_effect = transport_error
|
||||
stderr = io.StringIO()
|
||||
with (
|
||||
mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(input_bytes)),
|
||||
mock.patch.object(CLIENT, "observer_request", request),
|
||||
redirect_stderr(stderr),
|
||||
):
|
||||
arguments = ["--runtime", runtime]
|
||||
if runtime == "claude":
|
||||
arguments.append("--latest-entry")
|
||||
result = CLIENT.main(arguments, environ=ENVIRONMENT)
|
||||
return result, stderr.getvalue(), request
|
||||
|
||||
def test_claude_prefers_inline_last_assistant_message(self) -> None:
|
||||
inline = "the just-finished assistant message"
|
||||
input_bytes = json.dumps({
|
||||
"last_assistant_message": inline,
|
||||
"transcript_path": "/must/not/be/opened.jsonl",
|
||||
}).encode()
|
||||
|
||||
result, stderr, request = self.run_client(
|
||||
input_bytes=input_bytes,
|
||||
reply={"ok": True},
|
||||
runtime="claude",
|
||||
)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
self.assertEqual(
|
||||
request.call_args.args[1]["latest_assistant_message"],
|
||||
inline,
|
||||
)
|
||||
|
||||
def test_claude_falls_back_to_transcript_when_inline_field_is_absent(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as directory:
|
||||
transcript = Path(directory) / "transcript.jsonl"
|
||||
transcript.write_text(
|
||||
json.dumps({
|
||||
"message": {
|
||||
"role": "assistant",
|
||||
"content": [{"type": "text", "text": "fallback message"}],
|
||||
}
|
||||
})
|
||||
+ "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
input_bytes = json.dumps({"transcript_path": str(transcript)}).encode()
|
||||
|
||||
result, stderr, request = self.run_client(
|
||||
input_bytes=input_bytes,
|
||||
reply={"ok": True},
|
||||
runtime="claude",
|
||||
)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
self.assertEqual(
|
||||
request.call_args.args[1]["latest_assistant_message"],
|
||||
"fallback message",
|
||||
)
|
||||
|
||||
def test_claude_present_invalid_inline_field_fails_without_fallback(self) -> None:
|
||||
fallback = mock.Mock(return_value="must not be used")
|
||||
input_bytes = json.dumps({
|
||||
"last_assistant_message": None,
|
||||
"transcript_path": "/unused/transcript.jsonl",
|
||||
}).encode()
|
||||
|
||||
with mock.patch.object(CLIENT, "claude_latest_entry", fallback):
|
||||
result, stderr, request = self.run_client(
|
||||
input_bytes=input_bytes,
|
||||
reply={"ok": True},
|
||||
runtime="claude",
|
||||
)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr)
|
||||
request.assert_not_called()
|
||||
fallback.assert_not_called()
|
||||
|
||||
def test_claude_inline_message_size_guard_stays_enforced(self) -> None:
|
||||
request = mock.Mock(return_value={"ok": True})
|
||||
stderr = io.StringIO()
|
||||
with (
|
||||
mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(b"{}")),
|
||||
mock.patch.object(
|
||||
CLIENT,
|
||||
"read_json",
|
||||
return_value={"last_assistant_message": "x" * (CLIENT.MAX_FRAME + 1)},
|
||||
),
|
||||
mock.patch.object(CLIENT, "observer_request", request),
|
||||
redirect_stderr(stderr),
|
||||
):
|
||||
result = CLIENT.main(
|
||||
["--runtime", "claude", "--latest-entry"],
|
||||
environ=ENVIRONMENT,
|
||||
)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr.getvalue())
|
||||
request.assert_not_called()
|
||||
|
||||
def test_pi_still_posts_only_its_runtime_message(self) -> None:
|
||||
result, stderr, request = self.run_client(reply={"ok": True})
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
self.assertEqual(
|
||||
request.call_args.args[1]["latest_assistant_message"],
|
||||
"ordinary turn",
|
||||
)
|
||||
|
||||
def test_nothing_pending_observation_refusal_is_benign(self) -> None:
|
||||
result, stderr, request = self.run_client(
|
||||
reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE"}
|
||||
)
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
request.assert_called_once()
|
||||
|
||||
def test_success_reply_remains_successful(self) -> None:
|
||||
result, stderr, _request = self.run_client(reply={"ok": True})
|
||||
|
||||
self.assertEqual(result, 0)
|
||||
self.assertEqual(stderr, "")
|
||||
|
||||
def test_pending_cycle_auth_failure_stays_fail_closed(self) -> None:
|
||||
result, _stderr, _request = self.run_client(
|
||||
reply={"ok": False, "code": "ANCESTRY_MISMATCH"}
|
||||
)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
|
||||
def test_transport_failure_stays_fail_closed(self) -> None:
|
||||
result, stderr, _request = self.run_client(
|
||||
transport_error=ConnectionRefusedError("observer unavailable")
|
||||
)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr)
|
||||
|
||||
def test_parse_failure_stays_fail_closed(self) -> None:
|
||||
result, stderr, request = self.run_client(input_bytes=b"{")
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr)
|
||||
request.assert_not_called()
|
||||
|
||||
def test_malformed_wire_replies_stay_fail_closed(self) -> None:
|
||||
for response in (
|
||||
b"not-json\n",
|
||||
b'{"ok":true}',
|
||||
b"{}\n{}\n",
|
||||
b'{"ok":false,"code":"OBSERVATION_UNAVAILABLE"}\n\n',
|
||||
b'{"ok":true,"ok":false,"code":"OBSERVATION_UNAVAILABLE"}\n',
|
||||
DEEPLY_NESTED_JSON + b"\n",
|
||||
b"x" * (CLIENT.MAX_FRAME + 1),
|
||||
):
|
||||
with self.subTest(response=response):
|
||||
stderr = io.StringIO()
|
||||
with (
|
||||
mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(VALID_INPUT)),
|
||||
mock.patch.object(
|
||||
CLIENT.socket,
|
||||
"socket",
|
||||
return_value=FakeObserverSocket(response),
|
||||
),
|
||||
redirect_stderr(stderr),
|
||||
):
|
||||
result = CLIENT.main(["--runtime", "pi"], environ=ENVIRONMENT)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr.getvalue())
|
||||
|
||||
def test_oversized_input_stays_fail_closed(self) -> None:
|
||||
result, stderr, request = self.run_client(input_bytes=b"x" * (CLIENT.MAX_FRAME + 1))
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr)
|
||||
request.assert_not_called()
|
||||
|
||||
def test_deeply_nested_input_stays_fail_closed(self) -> None:
|
||||
result, stderr, request = self.run_client(input_bytes=DEEPLY_NESTED_JSON)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr)
|
||||
request.assert_not_called()
|
||||
|
||||
def test_json_recursion_failure_stays_fail_closed(self) -> None:
|
||||
request = mock.Mock()
|
||||
stderr = io.StringIO()
|
||||
with (
|
||||
mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(VALID_INPUT)),
|
||||
mock.patch.object(CLIENT, "observer_request", request),
|
||||
mock.patch.object(
|
||||
CLIENT.json,
|
||||
"loads",
|
||||
side_effect=RecursionError("maximum JSON nesting exceeded"),
|
||||
),
|
||||
redirect_stderr(stderr),
|
||||
):
|
||||
result = CLIENT.main(["--runtime", "pi"], environ=ENVIRONMENT)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
self.assertIn("Mosaic receipt observer refused", stderr.getvalue())
|
||||
request.assert_not_called()
|
||||
|
||||
def test_observation_unavailable_with_unexpected_fields_stays_fail_closed(self) -> None:
|
||||
result, _stderr, _request = self.run_client(
|
||||
reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE", "unexpected": True}
|
||||
)
|
||||
|
||||
self.assertEqual(result, 2)
|
||||
|
||||
def test_non_boolean_ok_values_stay_fail_closed(self) -> None:
|
||||
for reply in (
|
||||
{"ok": 1},
|
||||
{"ok": 0, "code": "OBSERVATION_UNAVAILABLE"},
|
||||
):
|
||||
with self.subTest(reply=reply):
|
||||
result, _stderr, _request = self.run_client(reply=reply)
|
||||
self.assertEqual(result, 2)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user