mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 20:27:56 +00:00
Keep private-snapshot provisioning off the machine-wide worktree lock
A mutating delegate_start provisioned the child's execution snapshot inside subagent_worktrees' ops lock, and the untracked-file classification inside it spawned one `git diff --no-index --numstat` per file. One project with 67,692 untracked files held the lock for 40 minutes; every other mutating start on the machine timed out after 120 s (issue #1241). subagent_worktrees: the ops lock now guards shared metadata only. Provisioning takes it twice for milliseconds - registry row FIRST, then the baseline pin, then `worktree add --no-checkout` - so a crash after the row is reclaimable by the startup GC; the tree walk, hashing, populate (`reset --hard --quiet --no-recurse-submodules`, what `worktree add` runs internally, minus the target's post-checkout hook, which no longer executes project-authored code at provision) and the raw-bytes copy run outside it. Removal deletes files outside the lock and forgets admin dir, pin and row inside. Payload snapshots copy and commit outside the lock. A malformed registry refuses before the tree is hashed. The lock file names its holder (pid/task/op/since/target) so the owner-aware stale check evicts a SIGKILLed holder at once, and a timeout is a typed WorktreeOpsLockBusy. workspace_patch_capture: untracked_binary_verdicts stages every regular file as the empty blob into a scratch index and asks ONE index-versus-worktree `git diff --numstat -z`, so the verdict stays git's own (attributes, diff drivers, clean filters, working-tree encodings) with one process per inventory; a failed batch falls back to the per-file verdict with a warning. Both the snapshot and the finalization patch capture use it. delegate: every pre-POST provisioning refusal is definitely_unrun (a configured leaf ends at $0), carries the lock holder when the cause is a busy lock, and settles its invocation with a durable START_FAILED row; the started receipt discloses the snapshot's entries, untracked files, file-input bytes and provisioning seconds. Docs: ARCHITECTURE section 6 snapshot paragraph, section 1 module row, DEVELOPMENT delegated-lane bullet; two chapter byte budgets raised for the added rationale text. Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
parent
567fa5330d
commit
226285a1f8
11 changed files with 986 additions and 295 deletions
|
|
@ -236,7 +236,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
|
|||
├── subagent_messages.py ← Bounded durable child-message identity shared by the final frame, recovery, persistence and replay; `executor_observation_meta` validates task-bound progress actor facts
|
||||
├── subagents.py ← Subagent envelopes + bounded legacy compatibility; `configured_subagent` snapshots dispatch through subagent_runtime
|
||||
├── subagent_history.py ← Existing compact helper receipt: dated API attempts, session settlement/recovery and typed unrun starts; dynamic context and owner UI only, never admission (§6 Route health)
|
||||
├── subagent_worktrees.py ← Worktree lifecycle + durable registry `state/subagent_worktrees.json`; `provision_genesis_project` (never registry/GC); execution snapshots pinned by `refs/ouroboros/delegated/`; standalone payload snapshots; removal only explicit or custody-cross-checked startup GC, fail-closed on an unreadable custody log (§6 Delegated subagents)
|
||||
├── subagent_worktrees.py ← Worktree lifecycle + durable registry `state/subagent_worktrees.json`; `provision_genesis_project` (never registry/GC); execution snapshots pinned by `refs/ouroboros/delegated/`; standalone payload snapshots; removal only explicit or custody-cross-checked startup GC, fail-closed on an unreadable custody log; the ops lock covers only registry/admin-dir/ref writes (issue #1241) (§6 Delegated subagents)
|
||||
├── artifacts.py ← Attachment staging into `artifact_store/attachments/`; artifact records; scratch fingerprints (`.scratch_manifest.json`) that gate patch exclusion only while content matches; the undeclared-output guard; `delegated_capture_read_target`; partial tool evidence first reads its verified actor source, then a matching redacted observability projection; recovered text uses the existing exact-text source writer best-effort, and a later budget cut without a verified handle remains source-unavailable
|
||||
├── retention.py ← Unified GC retention SSOT: clamp/age-cutoff + legacy-key seed picker
|
||||
├── workspace_preflight.py ← Read-only external-workspace git/manifest/toolchain snapshot used by gateway task creation
|
||||
|
|
|
|||
|
|
@ -323,7 +323,7 @@ WHERE a mutating run's changes are destined is the second, separate record — t
|
|||
|
||||
A payload target gets a standalone private Git snapshot (`subagent_worktrees.provision_payload_snapshot`): the live payload is never initialized as Git, and capture trusts nothing under the child-writable snapshot's `.git`. Disposition (`integrate_payload_patch`) applies a live, index-free `git apply` in no-repository mode under a whole-payload content-hash CAS (drift = typed conflict; identical content = idempotent applied) — `GIT_CEILING_DIRECTORIES` is pinned at the payload's resolved PARENT, because git still searches the ceiling entry itself and an ancestor Git worktree above the runtime data root could otherwise make git skip every hunk at rc=0 — and reserved paths refuse the WHOLE apply as `blocked_reserved_paths` with the candidate preserved. The post-apply outcome set is complete: a live loader hash equal to the recorded RESULT hash is the success; a hash equal to the recorded BASELINE hash with a non-empty touched set is a provable non-mutation that RESOLVES the apply intent (typed `INTEGRATE_APPLY_NO_OP`, the `apply_no_op` arm: no success, no disposition, no reconcile queued, retry lane open); anything else is the ambiguous mismatch, whose intent stays PENDING and whose reconcile marker IS queued because the payload did mutate. A successful apply queues the extension reconcile (`request_extension_reconcile`) and the skill's review goes STALE pending fresh `skill_preflight`/`skill_review`; a run whose ONLY change is a mode flip is already refused at CAPTURE time as `unreviewable_metadata_change`, so a live hash still equal to the baseline means nothing was written.
|
||||
|
||||
**A mutating run normally executes in a PRIVATE EXECUTION SNAPSHOT** (metered children keep sharing the tree; their patches integrate through `integrate_subagent_patch` — sha256-bound, 3-way `--index`, protected-path gated, genesis refused, `coop_already_in_tree` a no-op). At `delegate_start` the host snapshots tracked, staged and eligible untracked state, deciding sensitive/credential vetoes BEFORE hashing because `git add -A` would put secrets such as `.env` in the shared object database. The baseline is pinned by `refs/ouroboros/delegated/` and checked out as a detached worktree; `scope.root` remains the authority target and `execution.workspaceRoot` names the snapshot. The host appends a separate typed binding after the immutable work order: the snapshot is writable and the authority is read-only until integration. Directory-copy runs use the engine-created copy and never the source folder. Full native access has no filesystem sandbox; the binding names the writable root but does not enforce it. Terminal capture records authority drift as evidence: ready-no-changes stays no-change with unknown authorship, while ready-with-changes keeps its private artifact and the locked baseline proof decides integration. Nested Git directories are excluded and disclosed; skill payloads use content-hash CAS. The binding is durable before POST and retry replays it; a GC-collected snapshot is a typed `execution_snapshot_missing` refusal. Worktrees live in `state/subagent_worktrees.json`; removal is explicit or custody-cross-checked startup GC, fail-closed on unreadable custody. The run still uses `live` from the engine view, so the scoped-HOME/`delegated` marker below applies.
|
||||
**A mutating run normally executes in a PRIVATE EXECUTION SNAPSHOT** (metered children keep sharing the tree; their patches integrate through `integrate_subagent_patch` — sha256-bound, 3-way `--index`, protected-path gated, genesis refused, `coop_already_in_tree` a no-op). At `delegate_start` the host snapshots tracked, staged and eligible untracked state, deciding sensitive/credential vetoes BEFORE hashing because `git add -A` would put secrets such as `.env` in the shared object database; git's binary verdict for the whole untracked inventory comes from ONE index-versus-worktree `git diff --numstat` over a scratch index (`workspace_patch_capture.untracked_binary_verdicts`, shared with patch capture), never one process per file. The machine-wide worktree ops lock (`subagent_worktrees._ops_lock`) guards SHARED metadata only — the registry file, a target's `.git/worktrees` and its baseline pin — held twice for milliseconds (row FIRST, then the ref, then `worktree add --no-checkout`, so a crash after the row is GC-nameable); listing, classifying, hashing, populating (`reset --hard --no-recurse-submodules`: what `worktree add` runs internally, minus the target's `post-checkout` hook) and copying run outside it (issue #1241: one 67k-file provision held the lock 40 minutes and every other mutating start timed out). A held lock refuses typed (`cause: lock_busy` + the holder's pid/task/op), a SIGKILLed holder is evicted by the owner-aware stale check, every pre-POST provisioning refusal is `definitely_unrun` with a durable `START_FAILED` row, and the start receipt discloses `snapshot` size/time. The baseline is pinned by `refs/ouroboros/delegated/` and checked out as a detached worktree; `scope.root` remains the authority target and `execution.workspaceRoot` names the snapshot. The host appends a separate typed binding after the immutable work order: the snapshot is writable and the authority is read-only until integration. Directory-copy runs use the engine-created copy and never the source folder. Full native access has no filesystem sandbox; the binding names the writable root but does not enforce it. Terminal capture records authority drift as evidence: ready-no-changes stays no-change with unknown authorship, while ready-with-changes keeps its private artifact and the locked baseline proof decides integration. Nested Git directories are excluded and disclosed; skill payloads use content-hash CAS. The binding is durable before POST and retry replays it; a GC-collected snapshot is a typed `execution_snapshot_missing` refusal. Worktrees live in `state/subagent_worktrees.json`; removal is explicit or custody-cross-checked startup GC, fail-closed on unreadable custody. The run still uses `live` from the engine view, so the scoped-HOME/`delegated` marker below applies.
|
||||
|
||||
At terminal, `delegate_wait` captures the run's diff against the baseline durably into the task's artifact store; NOTHING reaches the target automatically — the nanny explicitly applies or rejects through `integrate_delegated_patch`. Git and skill captures are whole-result operations: omitted `paths` and an exact empty list select the same captured result, disclosed when explicit; nonempty selectors remain directory-only. Engine-directory empty/subset semantics are unchanged. The staging substrate differs: a GIT workspace target applies under the repo git lock after PROVING no touched path drifted from `baseline_sha` — a plain `git apply` relocates hunks by offset, and the touched-path set is read from `git apply --numstat` in BOTH directions because each direction names only the paths it writes — then applies and STAGES, never commits. A SKILL-PAYLOAD target captures through the payload adapter over a parent-owned trusted index and applies LIVE into the non-Git payload — nothing is staged into any active root and no `.git` or index is created in the payload. The protected-path gate applies only when the target IS the Ouroboros body; a conflict (proven drift) is owned by the still-running nanny, with snapshot and patch persisting until explicit resolution or discard. Mutation rides an apply-intent protocol: a durable `delegate_run_patch_apply_started` row lands before any tree mutation, so on replay a pending intent without a disposition answers typed `INTEGRATE_DELEGATED_APPLY_AMBIGUOUS`, resolved only by explicit `acknowledge_ambiguous=true`, while the provably non-mutating outcomes (a lock error, proven baseline drift, a failed apply, a verified revert, a baseline-equal payload hash) RESOLVE the intent. `patch_verdict.py` is the ONE verdict writer for both pipelines: subjects are minted `run_<rid>` by the writer, never prefix-matched by readers, and each decision lands twice — artifact plus typed `delegate_run_patch_verdict` custody row — with a failed artifact write disclosed on the row. `artifacts.delegated_capture_read_target` narrowly rebinds `artifact_store` READS for the owning task's own `delegated_runs/` prefix, and `delegate_shared.orphan_capture_read_target` extends the same read-only, one-directory READ to the terminal-owner ORPHAN the disposition rule authorizes (confirmed by `orphan_disposition_status`), so the actor that may dispose a patch can inspect it without widened write authority. A read-only child stays in Claudexor's default envelope — one transport with one derived difference, not a second pipeline.
|
||||
|
||||
|
|
|
|||
|
|
@ -326,7 +326,8 @@ and 23 (`delegated_transport`), both critical. The imperatives:
|
|||
staged-never-committed) is unchanged
|
||||
(`tests/test_delegated_run_isolation_orphans.py`). A copy failure or a
|
||||
source change against the baseline leaves no registered snapshot or pinned
|
||||
ref (`tests/test_snapshot_file_inputs.py`).
|
||||
ref; no tree walk or per-file git process runs under the worktree ops lock
|
||||
(`tests/test_snapshot_file_inputs.py`, `tests/test_subagent_worktrees_lock_scope.py`).
|
||||
- Outcome honesty: a delegating parent must not produce a clean no-tool final
|
||||
answer while direct children run undecided — one bounded absorption
|
||||
reminder, then best-effort (`children_unabsorbed`); the delivery candidate
|
||||
|
|
|
|||
|
|
@ -135,6 +135,19 @@ def _fail(tool: str, code: str, detail: str, **extra: Any) -> ToolResult:
|
|||
return delegate_result(payload)
|
||||
|
||||
|
||||
def lock_busy_facts(exc: BaseException) -> Dict[str, Any]:
|
||||
"""Typed facts when a HELD worktree ops lock refused a snapshot provision (#1241):
|
||||
who holds it and for what (``subagent_worktrees.WorktreeOpsLockBusy``), so the
|
||||
nanny can wait for that provision instead of guessing. ``{}`` for any other cause."""
|
||||
holder = getattr(exc, "holder", None)
|
||||
if not isinstance(exc, TimeoutError) or holder is None:
|
||||
return {}
|
||||
return {"cause": "lock_busy", "holder": dict(holder), "retryable": True,
|
||||
"waited_sec": round(float(getattr(exc, "waited_sec", 0.0) or 0.0), 1),
|
||||
"retry_hint": "Another snapshot is being provisioned under the shared worktree "
|
||||
"lock; wait for it (see holder) and retry delegate_start."}
|
||||
|
||||
|
||||
def _emit(ctx: ToolContext, kind: str, payload: Dict[str, Any]) -> None:
|
||||
custody.emit(custody.custody_root(ctx), kind, {
|
||||
"task_id": str(getattr(ctx, "task_id", "") or ""), **payload,
|
||||
|
|
|
|||
|
|
@ -7,9 +7,12 @@ returns a ``workspace.patch``; the parent integrates and is the sole committer.
|
|||
|
||||
git has no automatic worktree garbage collection, so we keep a durable JSON
|
||||
registry (``data/state/subagent_worktrees.json``) and prune orphans on startup.
|
||||
All worktree mutations are serialized by a portable cross-process lock because
|
||||
``git worktree add/remove/prune`` mutate shared ``.git/worktrees`` metadata and
|
||||
the existing repo git lock is drive-root scoped, not ``.git`` scoped.
|
||||
Mutations of SHARED metadata — a target's ``.git/worktrees`` (add/prune), its
|
||||
``refs/ouroboros/delegated/*`` pins and the registry file — are serialized by a
|
||||
portable cross-process lock (the existing repo git lock is drive-root scoped,
|
||||
not ``.git`` scoped). Tree-proportional work never runs under it (#1241):
|
||||
listing, classifying, hashing, populating, copying and deleting a snapshot's
|
||||
files happen outside the lock, so one huge inventory delays only its own task.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -29,7 +32,7 @@ from pathlib import Path
|
|||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from ouroboros.platform_layer import acquire_exclusive_file_lock, release_exclusive_file_lock
|
||||
from ouroboros.utils import atomic_write_json
|
||||
from ouroboros.utils import atomic_write_json, utc_now_iso
|
||||
from ouroboros.config import DATA_DIR, get_subagent_projects_root, get_subagent_worktree_root
|
||||
from ouroboros.retention import age_cutoff, get_gc_retention_days
|
||||
|
||||
|
|
@ -172,21 +175,56 @@ def _save_registry(entries: List[Dict[str, Any]], data_dir: Optional[Any] = None
|
|||
# --------------------------------------------------------------------------- #
|
||||
# Locking
|
||||
# --------------------------------------------------------------------------- #
|
||||
class WorktreeOpsLockBusy(TimeoutError):
|
||||
"""The worktree ops lock stayed held for the whole wait.
|
||||
|
||||
``holder`` is what the holder wrote into the lock file (``pid``, ``task``,
|
||||
``op``, ``since``, ``target``; ``{}`` when unreadable), so a refusal names who
|
||||
is doing what instead of a bare timeout."""
|
||||
|
||||
def __init__(self, lock_path: Path, waited_sec: float, holder: Dict[str, str]):
|
||||
self.lock_path, self.waited_sec, self.holder = str(lock_path), waited_sec, holder
|
||||
who = (f" (held by pid {holder.get('pid')} task {holder.get('task') or '?'} op "
|
||||
f"{holder.get('op') or '?'} since {holder.get('since') or '?'})") if holder else ""
|
||||
super().__init__(f"subagent worktree ops lock timeout after {waited_sec:.0f}s: {lock_path}{who}")
|
||||
|
||||
|
||||
def _lock_holder(lock_path: Path) -> Dict[str, str]:
|
||||
"""Parse ``pid=… task=… op=… since=… target=…``; ``target`` is the tail of the
|
||||
line so a path with spaces survives. Readable while held: flock never blocks reads."""
|
||||
try:
|
||||
text = lock_path.read_text(encoding="utf-8", errors="replace").strip()
|
||||
except OSError:
|
||||
return {}
|
||||
head, sep, target = text.partition(" target=")
|
||||
holder = {key: value for key, eq, value in (field.partition("=") for field in head.split())
|
||||
if eq and key in ("pid", "task", "op", "since")}
|
||||
if sep:
|
||||
holder["target"] = target
|
||||
return holder
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _ops_lock(root: Path):
|
||||
"""Serialize worktree mutations in-process (threading.Lock) and across
|
||||
processes via the shared portable file-lock SSOT (platform_layer)."""
|
||||
def _ops_lock(root: Path, *, op: str, task_id: str = "", target: str = ""):
|
||||
"""Serialize SHARED-METADATA mutations in-process (threading.Lock) and across
|
||||
processes via the portable file-lock SSOT (platform_layer).
|
||||
|
||||
Held for milliseconds plus one registry read-modify-write — never for a tree
|
||||
walk (#1241). The holder names itself in the lock file, ``pid=`` first: the
|
||||
owner-aware stale check reads it, so a SIGKILLed holder is evicted at once
|
||||
instead of blacking out every waiter for ``_LOCK_STALE_SEC``; a live holder
|
||||
is never evicted (its kernel flock is probed). A timeout is typed with the
|
||||
holder's facts."""
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
lock_path = root / _LOCK_NAME
|
||||
metadata = f"pid={os.getpid()} task={task_id or '-'} op={op} since={utc_now_iso()} target={target}"
|
||||
with _inproc_lock:
|
||||
started = time.monotonic()
|
||||
fd = acquire_exclusive_file_lock(
|
||||
lock_path,
|
||||
timeout_sec=_LOCK_TIMEOUT_SEC,
|
||||
stale_sec=_LOCK_STALE_SEC,
|
||||
metadata=str(os.getpid()),
|
||||
)
|
||||
lock_path, timeout_sec=_LOCK_TIMEOUT_SEC, stale_sec=_LOCK_STALE_SEC,
|
||||
metadata=metadata, owner_aware_stale=True)
|
||||
if fd is None:
|
||||
raise TimeoutError(f"subagent worktree ops lock timeout: {lock_path}")
|
||||
raise WorktreeOpsLockBusy(lock_path, time.monotonic() - started, _lock_holder(lock_path))
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
|
|
@ -228,6 +266,14 @@ def _git(repo_dir: Path, *args: str, check: bool = True,
|
|||
)
|
||||
|
||||
|
||||
def _git_quiet(repo_dir: Path, *args: str) -> None:
|
||||
"""Best-effort git: a failing command or a vanished repo is not an error here."""
|
||||
try:
|
||||
_git(repo_dir, *args, check=False)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _remove_paths(repo_dir: Path, wt_path: Path, branch: str, *, allowed_root: Optional[Any] = None) -> None:
|
||||
"""Best-effort teardown: drop the worktree checkout, dir, and branch.
|
||||
|
||||
|
|
@ -241,21 +287,51 @@ def _remove_paths(repo_dir: Path, wt_path: Path, branch: str, *, allowed_root: O
|
|||
not wt_text or wt_text in (".", "/", "//") or not _is_within(wt_path, Path(allowed_root))
|
||||
):
|
||||
return
|
||||
try:
|
||||
_git(repo_dir, "worktree", "remove", "--force", str(wt_path), check=False)
|
||||
except Exception:
|
||||
pass
|
||||
_git_quiet(repo_dir, "worktree", "remove", "--force", str(wt_path))
|
||||
if wt_path.exists():
|
||||
_force_rmtree(wt_path)
|
||||
try:
|
||||
_git(repo_dir, "worktree", "prune", check=False)
|
||||
except Exception:
|
||||
pass
|
||||
_git_quiet(repo_dir, "worktree", "prune")
|
||||
if branch:
|
||||
try:
|
||||
_git(repo_dir, "branch", "-D", branch, check=False)
|
||||
except Exception:
|
||||
pass
|
||||
_git_quiet(repo_dir, "branch", "-D", branch)
|
||||
|
||||
|
||||
def _register_snapshot(handle: "ExecutionSnapshotHandle", excluded: List[Dict[str, Any]],
|
||||
data_dir: Optional[Any]) -> None:
|
||||
"""Registry read-modify-write (under the ops lock): upsert the snapshot's row by path."""
|
||||
record = asdict(handle)
|
||||
record["kind"] = _KIND_DELEGATED_EXEC
|
||||
record["excluded_untracked"] = list(excluded)
|
||||
entries = [e for e in _load_registry(data_dir, strict=True, op="register_execution_snapshot")
|
||||
if e.get("path") != handle.path]
|
||||
entries.append(record)
|
||||
_save_registry(entries, data_dir)
|
||||
|
||||
|
||||
def _unregister_snapshot(snapshot_id: str, data_dir: Optional[Any], op: str) -> None:
|
||||
"""Registry read-modify-write (under the ops lock): drop the snapshot's row."""
|
||||
survivors = [e for e in _load_registry(data_dir, strict=True, op=op) if not (
|
||||
e.get("kind") == _KIND_DELEGATED_EXEC and e.get("snapshot_id") == snapshot_id)]
|
||||
_save_registry(survivors, data_dir)
|
||||
|
||||
|
||||
def _discard_snapshot_checkout(target: Path, wt_path: Path, ref: str, snapshot_id: str, *,
|
||||
root: Path, data_dir: Optional[Any], task_id: str) -> None:
|
||||
"""Undo a Git execution snapshot once its row and pin exist: the checkout's
|
||||
files go first, OUTSIDE the lock (a 1.8 GB tree takes minutes), then one short
|
||||
section forgets the admin dir, unpins the baseline and drops the row — the
|
||||
postcondition ``tests/test_snapshot_file_inputs.py`` asserts (no row, no ref,
|
||||
no ``dlg_*`` directory). Best-effort: the startup GC reconciles what a crash
|
||||
leaves."""
|
||||
if _is_within(wt_path, root) and wt_path.exists():
|
||||
_force_rmtree(wt_path)
|
||||
try:
|
||||
with _ops_lock(root, op="discard", task_id=task_id, target=str(target)):
|
||||
_git_quiet(target, "worktree", "prune")
|
||||
_git_quiet(target, "update-ref", "-d", ref)
|
||||
_unregister_snapshot(snapshot_id, data_dir, "discard_execution_snapshot")
|
||||
except Exception:
|
||||
log.warning("Failed to discard execution snapshot %s after a provisioning failure",
|
||||
snapshot_id, exc_info=True)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
|
@ -290,7 +366,7 @@ def provision_worktree(
|
|||
root = _resolve_root(worktree_root)
|
||||
_assert_root_isolated(root, repo_dir, _data_dir(data_dir))
|
||||
safe_task = _safe_name(task_id)
|
||||
with _ops_lock(root):
|
||||
with _ops_lock(root, op="worktree", task_id=str(task_id or ""), target=str(repo_dir)):
|
||||
if base_sha:
|
||||
_git(repo_dir, "rev-parse", "--verify", f"{base_sha}^{{commit}}")
|
||||
base_sha = _git(repo_dir, "rev-parse", base_sha).stdout.strip()
|
||||
|
|
@ -355,7 +431,7 @@ def provision_genesis_project(
|
|||
root = root.expanduser().resolve()
|
||||
_assert_root_isolated(root, repo_dir, _data_dir(data_dir))
|
||||
safe_task = _safe_name(dir_name or task_id)
|
||||
with _ops_lock(root):
|
||||
with _ops_lock(root, op="genesis", task_id=str(task_id or "")):
|
||||
proj = (root / safe_task).resolve()
|
||||
# Genesis projects are durable: never clobber an existing one -> unique name. Since
|
||||
# dir_name can repeat across projects (a shared display name), count up under the
|
||||
|
|
@ -446,6 +522,9 @@ class ExecutionSnapshotHandle:
|
|||
file_baseline: Dict[str, Dict[str, Any]] = field(default_factory=dict)
|
||||
untracked_baseline: Dict[str, Dict[str, Any]] = field(default_factory=dict)
|
||||
capture_warnings: tuple = ()
|
||||
# Wall-clock seconds the provision took; a disclosure on the start receipt and
|
||||
# the baseline manifest (a heavy tree is visible the moment it is snapshotted).
|
||||
provisioning_sec: float = 0.0
|
||||
|
||||
|
||||
def _git_env_index(index_path: Path) -> Dict[str, str]:
|
||||
|
|
@ -488,8 +567,15 @@ def provision_execution_snapshot(
|
|||
that commit, registered durably BEFORE the caller records any start intent,
|
||||
and removed only by an explicit disposition (``remove_execution_snapshot``)
|
||||
or by the startup GC after custody says the run is closed.
|
||||
|
||||
The ops lock is held twice, briefly (#1241): once to register the row, pin
|
||||
the baseline and create the worktree's admin dir — row FIRST, so everything
|
||||
after it is nameable by the startup GC — and once to finalize the row.
|
||||
Listing, classifying (one git process for every binary verdict), hashing,
|
||||
populating and copying run OUTSIDE it, so a huge untracked inventory delays
|
||||
only its own task instead of refusing every other mutating start.
|
||||
"""
|
||||
from ouroboros.headless import untracked_capture_veto_reason
|
||||
from ouroboros.workspace_patch_capture import untracked_binary_verdicts, untracked_capture_veto_reason
|
||||
|
||||
target = Path(target_root).resolve()
|
||||
if not (target / ".git").exists():
|
||||
|
|
@ -504,165 +590,153 @@ def provision_execution_snapshot(
|
|||
if not snap:
|
||||
raise ValueError("snapshot_id is required for a delegated execution snapshot")
|
||||
safe_snap = _safe_name(snap)
|
||||
safe_task = _safe_name(task_id)
|
||||
with _ops_lock(root):
|
||||
wt_path = (root / f"dlg_{safe_task}_{safe_snap[:16]}").resolve()
|
||||
baseline_ref = f"{_BASELINE_REF_PREFIX}{safe_snap}"
|
||||
# Clear any stale checkout/ref left by a crashed earlier attempt of the
|
||||
# SAME snapshot id (idempotent re-provision).
|
||||
_remove_paths(target, wt_path, "", allowed_root=root)
|
||||
_git(target, "update-ref", "-d", baseline_ref, check=False)
|
||||
head_proc = _git(target, "rev-parse", "--verify", "HEAD", check=False)
|
||||
target_head = head_proc.stdout.strip() if head_proc.returncode == 0 else ""
|
||||
index_path = root / f".baseline_index_{safe_snap[:16]}_{os.getpid()}"
|
||||
env = _git_env_index(index_path)
|
||||
try:
|
||||
if target_head:
|
||||
_git_env(target, "read-tree", target_head, env=env)
|
||||
task = str(task_id or "")
|
||||
started = time.monotonic()
|
||||
wt_path = (root / f"dlg_{_safe_name(task_id)}_{safe_snap[:16]}").resolve()
|
||||
baseline_ref = f"{_BASELINE_REF_PREFIX}{safe_snap}"
|
||||
# A malformed registry refuses HERE, before the tree is hashed (strict read).
|
||||
_load_registry(data_dir, strict=True, op="provision_execution_snapshot")
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
# A stale checkout of the SAME snapshot id (a crashed earlier attempt) is plain
|
||||
# files; its admin dir and its pin are replaced under the lock below.
|
||||
if wt_path.exists():
|
||||
_force_rmtree(wt_path)
|
||||
head_proc = _git(target, "rev-parse", "--verify", "HEAD", check=False)
|
||||
target_head = head_proc.stdout.strip() if head_proc.returncode == 0 else ""
|
||||
index_path = root / f".baseline_index_{safe_snap[:16]}_{os.getpid()}"
|
||||
env = _git_env_index(index_path)
|
||||
try:
|
||||
if target_head:
|
||||
_git_env(target, "read-tree", target_head, env=env)
|
||||
else:
|
||||
_git_env(target, "read-tree", "--empty", env=env)
|
||||
# ELIGIBILITY IS DECIDED BEFORE ANYTHING IS HASHED. `git add -A` writes a
|
||||
# blob for EVERY untracked file into the target's object database —
|
||||
# including `.env` / `credentials.json` — and removing the index entry
|
||||
# afterwards does not unwrite the object: the execution worktree shares
|
||||
# that ODB, so a vetoed secret stayed readable there by hash. The staged
|
||||
# set is therefore computed first (tracked/staged paths from the REAL
|
||||
# index, plus only the ELIGIBLE untracked ones) and fed to plumbing that
|
||||
# touches nothing else. NUL-delimited stdin: byte-safe and immune to argv
|
||||
# limits.
|
||||
real_env = dict(os.environ)
|
||||
indexed = [
|
||||
p for p in _git_env(target, "ls-files", "-z", env=real_env)
|
||||
.stdout.decode("utf-8", errors="surrogateescape").split("\0") if p
|
||||
]
|
||||
untracked_raw = _git_env(
|
||||
target, "ls-files", "-z", "--others", "--exclude-standard",
|
||||
env=real_env,
|
||||
).stdout.decode("utf-8", errors="surrogateescape")
|
||||
untracked = [p for p in untracked_raw.split("\0") if p]
|
||||
excluded: List[Dict[str, Any]] = []
|
||||
eligible: List[str] = []
|
||||
untracked_baseline: Dict[str, Dict[str, Any]] = {}
|
||||
file_inputs: List[str] = []
|
||||
capture_warnings: List[Dict[str, Any]] = []
|
||||
# git's binary verdict for the whole inventory in ONE process (#1241).
|
||||
binary_verdicts = untracked_binary_verdicts(target, untracked, warnings=capture_warnings)
|
||||
from ouroboros.workspace_file_outputs import _side
|
||||
for rel in untracked:
|
||||
candidate = target / rel
|
||||
if candidate.is_dir() and not candidate.is_symlink():
|
||||
excluded.append({"path": rel, "reason": "nested_repository"})
|
||||
continue
|
||||
reference: List[str] = []
|
||||
reason = untracked_capture_veto_reason(
|
||||
target, rel, file_outputs=reference, warnings=capture_warnings, binary_verdicts=binary_verdicts)
|
||||
if reference:
|
||||
file_inputs.append(rel)
|
||||
elif reason:
|
||||
excluded.append({"path": rel, "reason": reason, "baseline": _side(target / rel)})
|
||||
else:
|
||||
_git_env(target, "read-tree", "--empty", env=env)
|
||||
# ELIGIBILITY IS DECIDED BEFORE ANYTHING IS HASHED. `git add -A` writes a
|
||||
# blob for EVERY untracked file into the target's object database —
|
||||
# including `.env` / `credentials.json` — and removing the index entry
|
||||
# afterwards does not unwrite the object: the execution worktree shares
|
||||
# that ODB, so a vetoed secret stayed readable there by hash. The staged
|
||||
# set is therefore computed first (tracked/staged paths from the REAL
|
||||
# index, plus only the ELIGIBLE untracked ones) and fed to plumbing that
|
||||
# touches nothing else. NUL-delimited stdin: byte-safe and immune to argv
|
||||
# limits.
|
||||
real_env = dict(os.environ)
|
||||
indexed = [
|
||||
p for p in _git_env(target, "ls-files", "-z", env=real_env)
|
||||
.stdout.decode("utf-8", errors="surrogateescape").split("\0") if p
|
||||
]
|
||||
untracked_raw = _git_env(
|
||||
target, "ls-files", "-z", "--others", "--exclude-standard",
|
||||
env=real_env,
|
||||
).stdout.decode("utf-8", errors="surrogateescape")
|
||||
excluded: List[Dict[str, Any]] = []
|
||||
eligible: List[str] = []
|
||||
untracked_baseline: Dict[str, Dict[str, Any]] = {}
|
||||
file_inputs: List[str] = []
|
||||
capture_warnings: List[Dict[str, Any]] = []
|
||||
for rel in (p for p in untracked_raw.split("\0") if p):
|
||||
candidate = target / rel
|
||||
if candidate.is_dir() and not candidate.is_symlink():
|
||||
excluded.append({"path": rel, "reason": "nested_repository"})
|
||||
continue
|
||||
reference: List[str] = []
|
||||
reason = untracked_capture_veto_reason(target, rel, file_outputs=reference, warnings=capture_warnings)
|
||||
if reference:
|
||||
file_inputs.append(rel)
|
||||
elif reason:
|
||||
from ouroboros.workspace_file_outputs import _side
|
||||
excluded.append({"path": rel, "reason": reason,
|
||||
"baseline": _side(target / rel)})
|
||||
else:
|
||||
eligible.append(rel)
|
||||
from ouroboros.workspace_file_outputs import _side
|
||||
baseline = _side(target / rel)
|
||||
if baseline is not None:
|
||||
untracked_baseline[rel] = baseline
|
||||
staged_paths = indexed + eligible
|
||||
if staged_paths:
|
||||
# `--add --remove` stages each named path's CURRENT worktree content
|
||||
# (and drops the entry when the file is gone), which is exactly what
|
||||
# `git add -A` did for these paths — without visiting any other file.
|
||||
_git_env(
|
||||
target, "update-index", "-z", "--add", "--remove", "--stdin", env=env,
|
||||
input_bytes=b"\0".join(
|
||||
p.encode("utf-8", errors="surrogateescape") for p in staged_paths) + b"\0")
|
||||
tree_sha = _git_env(target, "write-tree", env=env).stdout.decode("utf-8").strip()
|
||||
manifest_raw = _git_env(target, "ls-tree", "-r", "-z", tree_sha,
|
||||
env=dict(os.environ)).stdout
|
||||
import hashlib
|
||||
manifest_digest = hashlib.sha256(manifest_raw).hexdigest()
|
||||
entry_count = sum(1 for chunk in manifest_raw.split(b"\0") if chunk)
|
||||
commit_args = ["commit-tree", tree_sha, "-m",
|
||||
f"ouroboros: delegated-run baseline {snap}"]
|
||||
if target_head:
|
||||
commit_args[2:2] = ["-p", target_head]
|
||||
baseline_sha = _git_env(
|
||||
target, *commit_args,
|
||||
env={**env,
|
||||
"GIT_AUTHOR_NAME": "Ouroboros", "GIT_AUTHOR_EMAIL": "ouroboros@localhost",
|
||||
"GIT_COMMITTER_NAME": "Ouroboros", "GIT_COMMITTER_EMAIL": "ouroboros@localhost"},
|
||||
).stdout.decode("utf-8").strip()
|
||||
finally:
|
||||
try:
|
||||
index_path.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
eligible.append(rel)
|
||||
baseline = _side(target / rel)
|
||||
if baseline is not None:
|
||||
untracked_baseline[rel] = baseline
|
||||
staged_paths = indexed + eligible
|
||||
if staged_paths:
|
||||
# `--add --remove` stages each named path's CURRENT worktree content
|
||||
# (and drops the entry when the file is gone), which is exactly what
|
||||
# `git add -A` did for these paths — without visiting any other file.
|
||||
_git_env(
|
||||
target, "update-index", "-z", "--add", "--remove", "--stdin", env=env,
|
||||
input_bytes=b"\0".join(
|
||||
p.encode("utf-8", errors="surrogateescape") for p in staged_paths) + b"\0")
|
||||
tree_sha = _git_env(target, "write-tree", env=env).stdout.decode("utf-8").strip()
|
||||
manifest_raw = _git_env(target, "ls-tree", "-r", "-z", tree_sha,
|
||||
env=dict(os.environ)).stdout
|
||||
import hashlib
|
||||
manifest_digest = hashlib.sha256(manifest_raw).hexdigest()
|
||||
entry_count = sum(1 for chunk in manifest_raw.split(b"\0") if chunk)
|
||||
commit_args = ["commit-tree", tree_sha, "-m",
|
||||
f"ouroboros: delegated-run baseline {snap}"]
|
||||
if target_head:
|
||||
commit_args[2:2] = ["-p", target_head]
|
||||
baseline_sha = _git_env(
|
||||
target, *commit_args,
|
||||
env={**env,
|
||||
"GIT_AUTHOR_NAME": "Ouroboros", "GIT_AUTHOR_EMAIL": "ouroboros@localhost",
|
||||
"GIT_COMMITTER_NAME": "Ouroboros", "GIT_COMMITTER_EMAIL": "ouroboros@localhost"},
|
||||
).stdout.decode("utf-8").strip()
|
||||
finally:
|
||||
try:
|
||||
index_path.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
fields: Dict[str, Any] = dict(
|
||||
snapshot_id=snap, task_id=task, path=str(wt_path), target_root=str(target),
|
||||
baseline_ref=baseline_ref, baseline_sha=baseline_sha, baseline_tree=tree_sha,
|
||||
manifest_digest=manifest_digest, target_head=target_head, created_at=time.time(),
|
||||
entry_count=entry_count, excluded_untracked=tuple(excluded),
|
||||
untracked_baseline=untracked_baseline, capture_warnings=tuple(capture_warnings))
|
||||
with _ops_lock(root, op="provision", task_id=task, target=str(target)):
|
||||
# Row FIRST, then the pin, then the admin dir: a crash after any of these
|
||||
# leaves a REGISTERED snapshot custody never opened, which the startup GC
|
||||
# removes (checkout, ref, row). A pin without a row would be invisible.
|
||||
_git_quiet(target, "worktree", "prune")
|
||||
_register_snapshot(ExecutionSnapshotHandle(**fields), excluded, data_dir)
|
||||
_git(target, "update-ref", baseline_ref, baseline_sha)
|
||||
try:
|
||||
wt_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
_git(target, "worktree", "add", "--detach", str(wt_path), baseline_sha)
|
||||
from ouroboros.artifacts import copy_artifact_file
|
||||
wt_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
_git(target, "worktree", "add", "--detach", "--no-checkout", str(wt_path), baseline_sha)
|
||||
try:
|
||||
# Populate outside the lock: the same reset git's own `worktree add` runs
|
||||
# (no submodule recursion) — minus the target's post-checkout hook.
|
||||
_git(wt_path, "reset", "--hard", "--quiet", "--no-recurse-submodules")
|
||||
from ouroboros.artifacts import copy_artifact_file
|
||||
|
||||
# Git owns the baseline representation, but the child must see the
|
||||
# source's actual working bytes, not checkout's CRLF/smudge rewrite.
|
||||
# Read the existing tree inventory so deletions, links and gitlinks
|
||||
# keep Git's semantics and excluded paths can never enter the copy.
|
||||
for item in manifest_raw.split(b"\0"):
|
||||
metadata, separator, raw_path = item.partition(b"\t")
|
||||
if separator and metadata.split()[0] in (b"100644", b"100755"):
|
||||
relative = raw_path.decode("utf-8", errors="surrogateescape")
|
||||
original = target / relative
|
||||
if original.is_symlink():
|
||||
raise OSError(f"snapshot input changed from a regular file: {relative}")
|
||||
copy_artifact_file(original, wt_path / relative)
|
||||
# A concurrent source edit must not appear as the child's work.
|
||||
# Use the same Git representation as ordinary patch capture, once
|
||||
# for the whole tree, before the separately tracked file inputs.
|
||||
_git(wt_path, "diff", "--quiet", "--no-ext-diff", baseline_sha, "--")
|
||||
from ouroboros.workspace_file_outputs import copy_snapshot_file_inputs
|
||||
file_baseline = copy_snapshot_file_inputs(target, wt_path, file_inputs)
|
||||
if file_baseline:
|
||||
manifest_digest = hashlib.sha256(
|
||||
manifest_raw + json.dumps(file_baseline, sort_keys=True).encode("utf-8")
|
||||
).hexdigest()
|
||||
entry_count += len(file_baseline)
|
||||
except Exception:
|
||||
# Do not leak the pinned baseline ref when the checkout failed.
|
||||
_remove_paths(target, wt_path, "", allowed_root=root)
|
||||
_git(target, "update-ref", "-d", baseline_ref, check=False)
|
||||
raise
|
||||
handle = ExecutionSnapshotHandle(
|
||||
snapshot_id=snap,
|
||||
task_id=str(task_id or ""),
|
||||
path=str(wt_path),
|
||||
target_root=str(target),
|
||||
baseline_ref=baseline_ref,
|
||||
baseline_sha=baseline_sha,
|
||||
baseline_tree=tree_sha,
|
||||
manifest_digest=manifest_digest,
|
||||
target_head=target_head,
|
||||
created_at=time.time(),
|
||||
entry_count=entry_count,
|
||||
excluded_untracked=tuple(excluded),
|
||||
file_baseline=file_baseline,
|
||||
untracked_baseline=untracked_baseline,
|
||||
capture_warnings=tuple(capture_warnings),
|
||||
)
|
||||
try:
|
||||
entries = [
|
||||
e for e in _load_registry(data_dir, strict=True, op="provision_execution_snapshot")
|
||||
if e.get("path") != str(wt_path)
|
||||
]
|
||||
record = asdict(handle)
|
||||
record["kind"] = _KIND_DELEGATED_EXEC
|
||||
record["excluded_untracked"] = excluded
|
||||
entries.append(record)
|
||||
_save_registry(entries, data_dir)
|
||||
except Exception:
|
||||
# Registration is INSIDE the cleanup scope, exactly like the payload
|
||||
# branch: a snapshot nothing registered is invisible to disposal and
|
||||
# retention, and on this branch it also strands the baseline ref,
|
||||
# which pins its commit against git's own GC for good.
|
||||
_remove_paths(target, wt_path, "", allowed_root=root)
|
||||
_git(target, "update-ref", "-d", baseline_ref, check=False)
|
||||
raise
|
||||
return handle
|
||||
# Git owns the baseline representation, but the child must see the
|
||||
# source's actual working bytes, not checkout's CRLF/smudge rewrite.
|
||||
# Read the existing tree inventory so deletions, links and gitlinks
|
||||
# keep Git's semantics and excluded paths can never enter the copy.
|
||||
for item in manifest_raw.split(b"\0"):
|
||||
metadata, separator, raw_path = item.partition(b"\t")
|
||||
if separator and metadata.split()[0] in (b"100644", b"100755"):
|
||||
relative = raw_path.decode("utf-8", errors="surrogateescape")
|
||||
original = target / relative
|
||||
if original.is_symlink():
|
||||
raise OSError(f"snapshot input changed from a regular file: {relative}")
|
||||
copy_artifact_file(original, wt_path / relative)
|
||||
# A concurrent source edit must not appear as the child's work.
|
||||
# Use the same Git representation as ordinary patch capture, once
|
||||
# for the whole tree, before the separately tracked file inputs.
|
||||
_git(wt_path, "diff", "--quiet", "--no-ext-diff", baseline_sha, "--")
|
||||
from ouroboros.workspace_file_outputs import copy_snapshot_file_inputs
|
||||
file_baseline = copy_snapshot_file_inputs(target, wt_path, file_inputs)
|
||||
if file_baseline:
|
||||
manifest_digest = hashlib.sha256(
|
||||
manifest_raw + json.dumps(file_baseline, sort_keys=True).encode("utf-8")
|
||||
).hexdigest()
|
||||
entry_count += len(file_baseline)
|
||||
handle = ExecutionSnapshotHandle(**{
|
||||
**fields, "manifest_digest": manifest_digest, "entry_count": entry_count,
|
||||
"file_baseline": file_baseline, "provisioning_sec": round(time.monotonic() - started, 3)})
|
||||
with _ops_lock(root, op="provision", task_id=task, target=str(target)):
|
||||
_register_snapshot(handle, excluded, data_dir)
|
||||
except Exception:
|
||||
_discard_snapshot_checkout(target, wt_path, baseline_ref, snap, root=root, data_dir=data_dir, task_id=task)
|
||||
raise
|
||||
return handle
|
||||
|
||||
|
||||
def isolated_git_env() -> Dict[str, str]:
|
||||
|
|
@ -892,74 +966,73 @@ def provision_payload_snapshot(
|
|||
if not snap:
|
||||
raise ValueError("snapshot_id is required for a delegated payload snapshot")
|
||||
safe_snap = _safe_name(snap)
|
||||
safe_task = _safe_name(task_id)
|
||||
with _ops_lock(root):
|
||||
wt_path = (root / f"dlgp_{safe_task}_{safe_snap[:16]}").resolve()
|
||||
if wt_path.exists():
|
||||
_force_rmtree(wt_path) # idempotent re-provision of the SAME snapshot id
|
||||
source_hash = payload_content_hash(target)
|
||||
env = isolated_git_env()
|
||||
try:
|
||||
entry_count = _copy_payload_inventory(target, wt_path)
|
||||
# Fully config-isolated baseline: empty template dir (no user hooks),
|
||||
# no global/system config (Fable F1). RAW staging instead of `git
|
||||
# add` (Sol P1 modes/filters): a payload .gitattributes (eol/clean)
|
||||
# must not normalize the recorded baseline away from the live raw
|
||||
# bytes, and the baseline records the REAL file modes.
|
||||
_git(wt_path, "init", "--template=", env=env)
|
||||
from ouroboros.skill_loader import _iter_payload_files
|
||||
task = str(task_id or "")
|
||||
started = time.monotonic()
|
||||
wt_path = (root / f"dlgp_{_safe_name(task_id)}_{safe_snap[:16]}").resolve()
|
||||
_load_registry(data_dir, strict=True, op="provision_payload_snapshot") # refuse before the copy
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
if wt_path.exists():
|
||||
_force_rmtree(wt_path) # idempotent re-provision of the SAME snapshot id
|
||||
source_hash = payload_content_hash(target)
|
||||
env = isolated_git_env()
|
||||
try:
|
||||
# Copy, init, stage and commit run OUTSIDE the ops lock (#1241): the
|
||||
# standalone snapshot's .git is private, so its only shared metadata is
|
||||
# the registry row written under the lock at the end.
|
||||
entry_count = _copy_payload_inventory(target, wt_path)
|
||||
# Fully config-isolated baseline: empty template dir (no user hooks),
|
||||
# no global/system config (Fable F1). RAW staging instead of `git
|
||||
# add` (Sol P1 modes/filters): a payload .gitattributes (eol/clean)
|
||||
# must not normalize the recorded baseline away from the live raw
|
||||
# bytes, and the baseline records the REAL file modes.
|
||||
_git(wt_path, "init", "--template=", env=env)
|
||||
from ouroboros.skill_loader import _iter_payload_files
|
||||
|
||||
stage_raw_payload_inventory(
|
||||
wt_path,
|
||||
(p.relative_to(wt_path).as_posix() for p in _iter_payload_files(wt_path)),
|
||||
env)
|
||||
_git(
|
||||
wt_path,
|
||||
"-c", "user.email=ouroboros@localhost", "-c", "user.name=Ouroboros",
|
||||
"commit", "--allow-empty", "-m",
|
||||
f"ouroboros: delegated payload baseline {snap}", env=env,
|
||||
)
|
||||
baseline_sha = _git(wt_path, "rev-parse", "HEAD", env=env).stdout.strip()
|
||||
baseline_tree = _git(wt_path, "rev-parse", "HEAD^{tree}", env=env).stdout.strip()
|
||||
manifest_raw = _git(wt_path, "ls-tree", "-r", "-z", baseline_tree, env=env).stdout
|
||||
import hashlib
|
||||
stage_raw_payload_inventory(
|
||||
wt_path,
|
||||
(p.relative_to(wt_path).as_posix() for p in _iter_payload_files(wt_path)),
|
||||
env)
|
||||
_git(
|
||||
wt_path,
|
||||
"-c", "user.email=ouroboros@localhost", "-c", "user.name=Ouroboros",
|
||||
"commit", "--allow-empty", "-m",
|
||||
f"ouroboros: delegated payload baseline {snap}", env=env,
|
||||
)
|
||||
baseline_sha = _git(wt_path, "rev-parse", "HEAD", env=env).stdout.strip()
|
||||
baseline_tree = _git(wt_path, "rev-parse", "HEAD^{tree}", env=env).stdout.strip()
|
||||
manifest_raw = _git(wt_path, "ls-tree", "-r", "-z", baseline_tree, env=env).stdout
|
||||
import hashlib
|
||||
|
||||
manifest_digest = hashlib.sha256(
|
||||
manifest_raw.encode("utf-8", errors="surrogateescape")).hexdigest()
|
||||
if payload_content_hash(target) != source_hash:
|
||||
raise RuntimeError(
|
||||
"the live payload changed while it was being snapshotted "
|
||||
"(another writer raced the copy); retry the delegation")
|
||||
handle = ExecutionSnapshotHandle(
|
||||
snapshot_id=snap,
|
||||
task_id=str(task_id or ""),
|
||||
path=str(wt_path),
|
||||
target_root=str(target),
|
||||
baseline_ref="",
|
||||
baseline_sha=baseline_sha,
|
||||
baseline_tree=baseline_tree,
|
||||
manifest_digest=manifest_digest,
|
||||
target_head="",
|
||||
created_at=time.time(),
|
||||
entry_count=entry_count,
|
||||
standalone=True,
|
||||
payload_hash=source_hash,
|
||||
)
|
||||
# Registry write INSIDE the cleanup scope: an unregistered snapshot
|
||||
# directory would be invisible to disposal/retention (orphan leak).
|
||||
entries = [
|
||||
e for e in _load_registry(data_dir, strict=True, op="provision_payload_snapshot")
|
||||
if e.get("path") != str(wt_path)
|
||||
]
|
||||
record = asdict(handle)
|
||||
record["kind"] = _KIND_DELEGATED_EXEC
|
||||
record["excluded_untracked"] = []
|
||||
entries.append(record)
|
||||
_save_registry(entries, data_dir)
|
||||
except Exception:
|
||||
_force_rmtree(wt_path)
|
||||
raise
|
||||
return handle
|
||||
manifest_digest = hashlib.sha256(
|
||||
manifest_raw.encode("utf-8", errors="surrogateescape")).hexdigest()
|
||||
if payload_content_hash(target) != source_hash:
|
||||
raise RuntimeError(
|
||||
"the live payload changed while it was being snapshotted "
|
||||
"(another writer raced the copy); retry the delegation")
|
||||
handle = ExecutionSnapshotHandle(
|
||||
snapshot_id=snap,
|
||||
task_id=task,
|
||||
path=str(wt_path),
|
||||
target_root=str(target),
|
||||
baseline_ref="",
|
||||
baseline_sha=baseline_sha,
|
||||
baseline_tree=baseline_tree,
|
||||
manifest_digest=manifest_digest,
|
||||
target_head="",
|
||||
created_at=time.time(),
|
||||
entry_count=entry_count,
|
||||
standalone=True,
|
||||
payload_hash=source_hash,
|
||||
provisioning_sec=round(time.monotonic() - started, 3),
|
||||
)
|
||||
# Registry write INSIDE the cleanup scope: an unregistered snapshot
|
||||
# directory would be invisible to disposal/retention (orphan leak).
|
||||
with _ops_lock(root, op="provision_payload", task_id=task, target=str(target)):
|
||||
_register_snapshot(handle, [], data_dir)
|
||||
except Exception:
|
||||
_force_rmtree(wt_path)
|
||||
raise
|
||||
return handle
|
||||
|
||||
|
||||
def find_execution_snapshot(snapshot_id: str, data_dir: Optional[Any] = None) -> Optional[Dict[str, Any]]:
|
||||
|
|
@ -989,28 +1062,25 @@ def remove_execution_snapshot(
|
|||
if entry is None:
|
||||
return False
|
||||
root = _resolve_root(worktree_root)
|
||||
with _ops_lock(root):
|
||||
if entry.get("standalone"):
|
||||
# Standalone payload snapshot (R1 §10.4): its .git lives INSIDE the
|
||||
# snapshot directory and the target is a non-Git payload — remove
|
||||
# only the private directory and the registry row; no target-repo
|
||||
# worktree/ref command exists to run.
|
||||
wt_path = Path(str(entry.get("path") or ""))
|
||||
if str(wt_path).strip() and _is_within(wt_path, root) and wt_path.exists():
|
||||
_force_rmtree(wt_path)
|
||||
else:
|
||||
wt_path = Path(str(entry.get("path") or ""))
|
||||
# The checkout's files go first, OUTSIDE the lock (#1241: a 1.8 GB snapshot
|
||||
# took minutes to delete and timed every other mutating start out). Only a
|
||||
# path strictly inside the snapshot root is ever deleted: the registry is
|
||||
# durable state and a malformed row must never name an arbitrary path.
|
||||
if str(wt_path).strip() and _is_within(wt_path, root) and wt_path.exists():
|
||||
_force_rmtree(wt_path)
|
||||
with _ops_lock(root, op="remove", task_id=str(entry.get("task_id") or ""),
|
||||
target=str(entry.get("target_root") or "")):
|
||||
if not entry.get("standalone"):
|
||||
# A standalone payload snapshot (R1 §10.4) keeps its .git INSIDE the
|
||||
# directory just deleted; a Git snapshot also owns an admin dir and a
|
||||
# baseline pin in the TARGET repository.
|
||||
target = Path(str(entry.get("target_root") or "."))
|
||||
_remove_paths(target, Path(str(entry.get("path") or "")), "", allowed_root=root)
|
||||
_git_quiet(target, "worktree", "prune")
|
||||
ref = str(entry.get("baseline_ref") or "")
|
||||
if ref.startswith(_BASELINE_REF_PREFIX):
|
||||
try:
|
||||
_git(target, "update-ref", "-d", ref, check=False)
|
||||
except Exception:
|
||||
pass
|
||||
survivors = [e for e in _load_registry(data_dir, strict=True, op="remove_execution_snapshot") if not (
|
||||
e.get("kind") == _KIND_DELEGATED_EXEC and e.get("snapshot_id") == entry.get("snapshot_id")
|
||||
)]
|
||||
_save_registry(survivors, data_dir)
|
||||
_git_quiet(target, "update-ref", "-d", ref)
|
||||
_unregister_snapshot(str(entry.get("snapshot_id") or ""), data_dir, "remove_execution_snapshot")
|
||||
return True
|
||||
|
||||
|
||||
|
|
@ -1064,7 +1134,7 @@ def remove_worktree(
|
|||
match = entry
|
||||
break
|
||||
root = _resolve_root(worktree_root)
|
||||
with _ops_lock(root):
|
||||
with _ops_lock(root, op="remove_worktree", task_id=str(task_id or "")):
|
||||
if match is not None:
|
||||
_remove_paths(Path(match.get("repo_dir") or "."), Path(match.get("path") or ""), match.get("branch") or "", allowed_root=root)
|
||||
survivors = [
|
||||
|
|
@ -1097,7 +1167,7 @@ def prune_orphans(
|
|||
removed: List[Dict[str, Any]] = []
|
||||
kept: List[Dict[str, Any]] = []
|
||||
repos: set[str] = set()
|
||||
with _ops_lock(root):
|
||||
with _ops_lock(root, op="prune"):
|
||||
for entry in _load_registry(data_dir, strict=True, op="prune_orphans"):
|
||||
if entry.get("kind") == _KIND_DELEGATED_EXEC:
|
||||
# Delegated execution snapshots have their OWN lifecycle: they persist
|
||||
|
|
|
|||
|
|
@ -474,6 +474,7 @@ def _delegate_start(ctx: ToolContext, prompt: str, max_seconds: Optional[int] =
|
|||
definitely_unrun=True)
|
||||
snapshot, snap_error = _provision_snapshot(ctx, drive, target_root, invocation_id)
|
||||
if snap_error:
|
||||
_settle_refused_provision(ctx, gateway, snap_error, invocation_id, history_facts)
|
||||
return snap_error
|
||||
if snapshot is not None:
|
||||
snapshot_id, baseline_sha, root = snapshot.snapshot_id, snapshot.baseline_sha, snapshot.path
|
||||
|
|
@ -618,13 +619,15 @@ def _delegate_start(ctx: ToolContext, prompt: str, max_seconds: Optional[int] =
|
|||
baseline_sha=baseline_sha,
|
||||
resource_ref=resource_ref,
|
||||
processing=processing_info,
|
||||
engine_version=str(getattr(gateway, "engine_version", "") or ""))
|
||||
engine_version=str(getattr(gateway, "engine_version", "") or ""),
|
||||
snapshot_facts=_snapshot_facts(snapshot))
|
||||
|
||||
|
||||
def _started_payload(handle: Dict[str, Any], run_id: str, route: Any, access: str,
|
||||
authority: "DelegatedRunShape", root: str, *, durable: bool,
|
||||
recovering: bool, invocation_id: str, snapshot_id: str, target_root: str,
|
||||
baseline_sha: str, engine_version: str = "", resource_ref=None, processing=None) -> ToolResult:
|
||||
baseline_sha: str, engine_version: str = "", resource_ref=None, processing=None,
|
||||
snapshot_facts=None) -> ToolResult:
|
||||
"""The one author of delegate_start's started result (note + payload).
|
||||
|
||||
The AUTHORITY guidance and the CUSTODY warning are independent facts about the same
|
||||
|
|
@ -677,6 +680,8 @@ def _started_payload(handle: Dict[str, Any], run_id: str, route: Any, access: st
|
|||
payload["pending_invocation_id"] = str(invocation_id or "")
|
||||
if processing:
|
||||
payload["processing"] = processing
|
||||
if snapshot_facts:
|
||||
payload["snapshot"] = snapshot_facts
|
||||
if snapshot_id:
|
||||
# The C1 binding, stated where the nanny can read it: the run edits the
|
||||
# EXECUTION snapshot; the authority target receives nothing until apply.
|
||||
|
|
@ -698,6 +703,31 @@ def _started_payload(handle: Dict[str, Any], run_id: str, route: Any, access: st
|
|||
return delegate_result(payload)
|
||||
|
||||
|
||||
def _settle_refused_provision(ctx: ToolContext, gateway: Any, refusal: ToolResult,
|
||||
invocation_id: str, history_facts: Optional[dict]) -> None:
|
||||
"""A pre-POST provisioning refusal (no start row, no run) still settles its
|
||||
invocation durably (START_FAILED), so the refusal exists outside this process —
|
||||
the incident's bootstrap refusals left no row at all (#1241)."""
|
||||
from ouroboros.delegate_shared import delegate_payload
|
||||
|
||||
_retire_orphaned_registration(
|
||||
ctx, gateway, "", definite_refusal=True, invocation_id=invocation_id,
|
||||
reason=str(delegate_payload(refusal).get("reason") or "execution_snapshot_failed"),
|
||||
history_facts=history_facts)
|
||||
|
||||
|
||||
def _snapshot_facts(handle: Any) -> Dict[str, Any]:
|
||||
"""Disclosure on the start receipt: how big the private snapshot is and how long
|
||||
it took to provision (#1241). Facts only — nothing refuses or truncates on them."""
|
||||
if handle is None:
|
||||
return {}
|
||||
file_baseline = getattr(handle, "file_baseline", {}) or {}
|
||||
return {"entries": int(getattr(handle, "entry_count", 0) or 0),
|
||||
"untracked_files": len(getattr(handle, "untracked_baseline", {}) or {}) + len(file_baseline),
|
||||
"file_input_bytes": sum(int(item.get("size", 0) or 0) for item in file_baseline.values()),
|
||||
"provisioning_sec": float(getattr(handle, "provisioning_sec", 0.0) or 0.0)}
|
||||
|
||||
|
||||
def _retire_orphaned_registration(ctx: ToolContext, gateway: Any, project_id: str, *,
|
||||
definite_refusal: bool, reason: str,
|
||||
project_persistent: bool = False,
|
||||
|
|
|
|||
|
|
@ -22,7 +22,7 @@ from ouroboros.delegate_custody import RunCustody as _RunCustody
|
|||
# ONE refusal author for the whole delegate surface: the neutral leaf
|
||||
# `delegate_shared` (phase B's facade split), never a local twin that could drift.
|
||||
from ouroboros.delegate_registration_policy import record_persistent as _record_persistent
|
||||
from ouroboros.delegate_shared import _fail
|
||||
from ouroboros.delegate_shared import _fail, lock_busy_facts
|
||||
from ouroboros.configured_subagents import SESSION_ACCESS_PROFILES
|
||||
from ouroboros.tools.tool_result import ToolResult
|
||||
from ouroboros.tools.registry import ToolContext, active_repo_dir_for
|
||||
|
|
@ -379,7 +379,7 @@ def _provision_snapshot(ctx: ToolContext, drive: pathlib.Path, target_root: str,
|
|||
"A private execution snapshot of the write root could not be provisioned "
|
||||
f"({type(exc).__name__}: {exc}). The run was NOT started: a mutating "
|
||||
"delegated run executes only in its own snapshot, never in the shared tree.",
|
||||
target_root=target_root)
|
||||
target_root=target_root, definitely_unrun=True, **lock_busy_facts(exc))
|
||||
_record_baseline_manifest(drive, task_id, invocation_id, handle)
|
||||
return handle, None
|
||||
|
||||
|
|
@ -405,6 +405,7 @@ def _record_baseline_manifest(drive: pathlib.Path, task_id: str, invocation_id:
|
|||
"entry_count": handle.entry_count,
|
||||
"file_input_count": len(getattr(handle, "file_baseline", {})),
|
||||
"file_input_bytes": sum(item.get("size", 0) for item in getattr(handle, "file_baseline", {}).values()),
|
||||
"provisioning_sec": float(getattr(handle, "provisioning_sec", 0.0) or 0.0),
|
||||
"target_root": handle.target_root,
|
||||
"target_head": handle.target_head,
|
||||
"execution_root": handle.path,
|
||||
|
|
@ -904,7 +905,7 @@ def _provision_payload_snapshot(
|
|||
f"provisioned ({type(exc).__name__}: {exc}). The run was NOT started: "
|
||||
"a mutating delegated run executes only in its own snapshot, never "
|
||||
"in the live payload.",
|
||||
target_root=record["target_root"])
|
||||
target_root=record["target_root"], definitely_unrun=True, **lock_busy_facts(exc))
|
||||
record["resource_ref"]["payload_hash"] = handle.payload_hash
|
||||
_record_baseline_manifest(drive, task_id, invocation_id, handle,
|
||||
payload_hash=handle.payload_hash,
|
||||
|
|
|
|||
|
|
@ -3,7 +3,8 @@
|
|||
Owns the streamed `workspace.patch` and `workspace_patch.json` pair — patch
|
||||
baseline resolution (including the unborn-HEAD empty-tree case and the acting
|
||||
subagent `base_sha` binding), the bounded git process helpers the capture runs
|
||||
on, the declared-scratch and untracked eligibility filtering, the moved-HEAD
|
||||
on, the declared-scratch and untracked eligibility filtering (git's binary
|
||||
verdict for a whole inventory in one process), the moved-HEAD
|
||||
tripwire for a private self worktree, and the empty manifest a failed
|
||||
finalization falls back to. The static eligibility rules live in
|
||||
``workspace_patch_rules``; the task-drive, child-result and artifact
|
||||
|
|
@ -16,6 +17,7 @@ import json
|
|||
import os
|
||||
import pathlib
|
||||
import re
|
||||
import stat
|
||||
import subprocess
|
||||
import tempfile
|
||||
import threading
|
||||
|
|
@ -136,6 +138,10 @@ def write_workspace_patch_artifacts(
|
|||
except Exception:
|
||||
scratch_sha_by_rel = {}
|
||||
scratch_sha_by_abs = {}
|
||||
# git's binary verdict for the whole inventory in ONE process (#1241); the loop
|
||||
# below keeps its per-file order and reads the verdict instead of spawning
|
||||
# ``git diff --numstat`` per file.
|
||||
binary_verdicts = untracked_binary_verdicts(root, untracked, warnings=diagnostics)
|
||||
for rel in untracked:
|
||||
_want_sha = scratch_sha_by_rel.get(rel) or scratch_sha_by_abs.get(os.path.normcase(str((root / rel).resolve(strict=False))))
|
||||
if _want_sha:
|
||||
|
|
@ -154,7 +160,8 @@ def write_workspace_patch_artifacts(
|
|||
if reason:
|
||||
excluded.append({"path": rel, "reason": reason})
|
||||
continue
|
||||
blob_reason = _untracked_blob_exclude_reason(root, rel, file_outputs=file_output_paths, warnings=diagnostics)
|
||||
blob_reason = _untracked_blob_exclude_reason(
|
||||
root, rel, file_outputs=file_output_paths, warnings=diagnostics, binary_verdicts=binary_verdicts)
|
||||
if blob_reason:
|
||||
excluded.append({"path": rel, "reason": blob_reason})
|
||||
continue
|
||||
|
|
@ -657,13 +664,18 @@ def pem_capture_refusal(root: pathlib.Path, rel: str, *, warnings=None) -> str:
|
|||
return reason
|
||||
|
||||
|
||||
def _untracked_blob_exclude_reason(root: pathlib.Path, rel: str, *, file_outputs: Optional[List[str]] = None, warnings=None) -> str:
|
||||
def _untracked_blob_exclude_reason(root: pathlib.Path, rel: str, *, file_outputs: Optional[List[str]] = None,
|
||||
warnings=None, binary_verdicts: Optional[Dict[str, bool]] = None) -> str:
|
||||
"""Reason to drop an untracked file from the workspace patch when it is a
|
||||
build/runtime BINARY, exceeds the per-file size cap, or carries a PEM
|
||||
private-key header in its head bytes. Keeps real-usage patches
|
||||
source-shaped without losing data (the file stays in the workspace
|
||||
and is recorded under ``untracked_excluded``). On any git/stat failure the
|
||||
file is INCLUDED (conservative — the main binary diff still applies)."""
|
||||
file is INCLUDED (conservative — the main binary diff still applies).
|
||||
|
||||
``binary_verdicts`` is one :func:`untracked_binary_verdicts` batch over the whole
|
||||
inventory; ``None`` keeps git's per-file ``--numstat`` verdict (one subprocess per
|
||||
file) for a caller that did not batch."""
|
||||
|
||||
try:
|
||||
size = (root / rel).lstat().st_size
|
||||
|
|
@ -675,21 +687,85 @@ def _untracked_blob_exclude_reason(root: pathlib.Path, rel: str, *, file_outputs
|
|||
if file_outputs is not None:
|
||||
file_outputs.append(rel)
|
||||
return f"untracked file exceeds size cap ({size}B > {_PATCH_MAX_UNTRACKED_FILE_BYTES}B)"
|
||||
numstat = _git_stdout(
|
||||
["git", "diff", "--no-index", "--numstat", "--no-ext-diff", "--no-color", "--", os.devnull, rel],
|
||||
root,
|
||||
allow_rc={0, 1},
|
||||
errors=None,
|
||||
)
|
||||
first = numstat.strip().splitlines()[0] if numstat.strip() else ""
|
||||
if first.startswith("-\t-"):
|
||||
if binary_verdicts is None:
|
||||
numstat = _git_stdout(
|
||||
["git", "diff", "--no-index", "--numstat", "--no-ext-diff", "--no-color", "--", os.devnull, rel],
|
||||
root,
|
||||
allow_rc={0, 1},
|
||||
errors=None,
|
||||
)
|
||||
first = numstat.strip().splitlines()[0] if numstat.strip() else ""
|
||||
binary = first.startswith("-\t-")
|
||||
else:
|
||||
binary = bool(binary_verdicts.get(rel, False))
|
||||
if binary:
|
||||
if file_outputs is not None:
|
||||
file_outputs.append(rel)
|
||||
return "binary file"
|
||||
return ""
|
||||
|
||||
|
||||
def untracked_capture_veto_reason(root: pathlib.Path, rel: str, *, file_outputs: Optional[List[str]] = None, warnings=None) -> str:
|
||||
def untracked_binary_verdicts(root: pathlib.Path, rels: Sequence[str], *, warnings=None) -> Optional[Dict[str, bool]]:
|
||||
"""git's own binary verdict for every path in ``rels`` from ONE ``git diff``.
|
||||
|
||||
Every REGULAR file is staged as the empty blob into a scratch index
|
||||
(``update-index --index-info`` touches no file and writes no content object), and
|
||||
one index-versus-worktree ``git diff --numstat -z`` then runs exactly the machinery
|
||||
a per-file ``git diff --no-index --numstat`` ran — attributes, diff drivers, clean
|
||||
filters and working-tree encodings included — so the answer is git's, not a
|
||||
re-implementation (parity probed on git 2.53 for 16 path classes). ``-\\t-`` is
|
||||
binary; a text file, an empty file (absent from the diff) and a file that vanished
|
||||
meanwhile (``0\\t0``) are text, exactly as the per-file verdict reads them.
|
||||
Symlinks and other non-regular paths are text (git diffs the link text). One
|
||||
inventory of tens of thousands of files therefore costs one process, not one per
|
||||
file (#1241). ``None`` — the "did not batch" signal — when git cannot answer, after
|
||||
an advisory warning: the callers keep the per-file verdict (the same answer, one
|
||||
process per file), never a new refusal."""
|
||||
regular: List[str] = []
|
||||
for rel in rels:
|
||||
try:
|
||||
if rel and stat.S_ISREG(os.lstat(root / rel).st_mode):
|
||||
regular.append(rel)
|
||||
except OSError:
|
||||
continue
|
||||
if not regular:
|
||||
return {}
|
||||
env = dict(os.environ)
|
||||
fd, scratch = tempfile.mkstemp(prefix="ouroboros-binary-verdict-", suffix=".index")
|
||||
os.close(fd)
|
||||
env["GIT_INDEX_FILE"] = scratch
|
||||
|
||||
def _git(*args: str, data: bytes = b"") -> bytes:
|
||||
return subprocess.run(["git", *args], cwd=str(root), capture_output=True, input=data,
|
||||
env=env, timeout=300, check=True).stdout
|
||||
|
||||
try:
|
||||
empty_blob = _git("hash-object", "-w", "--stdin").strip()
|
||||
_git("read-tree", "--empty")
|
||||
_git("update-index", "-z", "--index-info",
|
||||
data=b"".join(b"100644 " + empty_blob + b"\t" + os.fsencode(rel) + b"\0" for rel in regular))
|
||||
rows = _git("diff", "--numstat", "-z", "--no-renames", "--no-ext-diff", "--no-color")
|
||||
except Exception as exc:
|
||||
if warnings is not None:
|
||||
warnings.append({"reason": "binary_verdict_batch_unavailable", "advisory": True,
|
||||
"detail": f"{type(exc).__name__}: {exc}"[:300]})
|
||||
return None
|
||||
finally:
|
||||
try:
|
||||
os.unlink(scratch)
|
||||
except OSError:
|
||||
pass
|
||||
verdicts = {rel: False for rel in regular}
|
||||
for row in rows.split(b"\0"):
|
||||
added, sep, rest = row.partition(b"\t")
|
||||
deleted, sep2, path = rest.partition(b"\t")
|
||||
if sep and sep2:
|
||||
verdicts[os.fsdecode(path)] = added == b"-" and deleted == b"-"
|
||||
return verdicts
|
||||
|
||||
|
||||
def untracked_capture_veto_reason(root: pathlib.Path, rel: str, *, file_outputs: Optional[List[str]] = None,
|
||||
warnings=None, binary_verdicts: Optional[Dict[str, bool]] = None) -> str:
|
||||
"""Classify an untracked file for Git, file-reference transfer, or exclusion.
|
||||
|
||||
The delegated-run baseline snapshot
|
||||
|
|
@ -703,6 +779,8 @@ def untracked_capture_veto_reason(root: pathlib.Path, rel: str, *, file_outputs:
|
|||
Returns the human-readable patch exclusion reason, or "" for a Git input.
|
||||
``file_outputs`` receives eligible binary/large files which use file
|
||||
artifacts rather than Git blobs; credential/junk exclusions never enter it.
|
||||
``binary_verdicts`` comes from one :func:`untracked_binary_verdicts` batch over
|
||||
the whole inventory (``None`` = per-file git verdict, see the blob check).
|
||||
"""
|
||||
reason = _sensitive_untracked_reason(rel)
|
||||
if reason:
|
||||
|
|
@ -710,7 +788,8 @@ def untracked_capture_veto_reason(root: pathlib.Path, rel: str, *, file_outputs:
|
|||
reason = _patch_exclude_reason(rel)
|
||||
if reason:
|
||||
return reason
|
||||
return _untracked_blob_exclude_reason(root, rel, file_outputs=file_outputs, warnings=warnings)
|
||||
return _untracked_blob_exclude_reason(root, rel, file_outputs=file_outputs, warnings=warnings,
|
||||
binary_verdicts=binary_verdicts)
|
||||
|
||||
|
||||
def _preflight_head_from_task(task: Dict[str, Any]) -> str:
|
||||
|
|
|
|||
|
|
@ -96,7 +96,10 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
|
|||
# the digest-selected historical read and the reader admission rule.
|
||||
# 295650 -> 297250: document the new diagnostic-only source/coverage contract,
|
||||
# unavailable evidence and no-effects ordering without removing review/custody rules.
|
||||
"docs/architecture/06-agent-core.md": 297250,
|
||||
# 297250 -> 298400: the private-snapshot paragraph now states the worktree ops
|
||||
# lock's scope (issue #1241: shared metadata only, row-then-ref order, batched
|
||||
# binary verdict, typed busy refusal) — rationale-layer text BIBLE P6 requires.
|
||||
"docs/architecture/06-agent-core.md": 298400,
|
||||
"docs/architecture/07-configuration.md": 36991,
|
||||
# 18947 -> 19287: CI failure collection now documents diagnostic desktop builds while release remains gated.
|
||||
"docs/architecture/08-git-branching-ci-and-build.md": 19287,
|
||||
|
|
@ -139,7 +142,9 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
|
|||
# never does; a pre-check's refusal takes the exact read). The one sentence it touches
|
||||
# (the lock's caller wait) is replaced; the rest is a rule the chapter lacked, and the
|
||||
# chapter had 5 bytes left. Sized to the text: 5 bytes of margin.
|
||||
"docs/development/06-rules-by-change-class.md": 94520,
|
||||
# 94520 -> 94650: the delegated-lane bullet names the worktree ops lock rule
|
||||
# (issue #1241: no tree walk or per-file git process under the lock).
|
||||
"docs/development/06-rules-by-change-class.md": 94650,
|
||||
"docs/development/07-managed-update-rule.md": 4166,
|
||||
"docs/development/08-mutation-attribution-rule.md": 2899,
|
||||
"docs/development/09-process-custody-rule.md": 10028,
|
||||
|
|
|
|||
361
tests/test_subagent_worktrees_lock_scope.py
Normal file
361
tests/test_subagent_worktrees_lock_scope.py
Normal file
|
|
@ -0,0 +1,361 @@
|
|||
"""The worktree ops lock guards shared metadata only (#1241).
|
||||
|
||||
One delegated snapshot of a project with 67,692 untracked files held the
|
||||
machine-wide ``.worktree_ops.lock`` for 40 minutes, and every other mutating
|
||||
``delegate_start`` on the machine waited 120 s and was refused. Listing,
|
||||
classifying, hashing, populating, copying and deleting a snapshot's files now run
|
||||
OUTSIDE the lock; the registry row, the baseline pin and the worktree admin dir
|
||||
are the only things written under it — row first, so everything after it is
|
||||
nameable by the startup GC. A held lock is a typed refusal that names its holder,
|
||||
a dead holder is evicted at once, and a pre-POST refusal is definitely unrun.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import pathlib
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from ouroboros import artifacts, subagent_worktrees as wt, workspace_patch_capture as capture
|
||||
from ouroboros.platform_layer import _lock_identity
|
||||
from tests._delegated_transport_shared import _owned_gateway_uses_each_test_transport # noqa: F401
|
||||
from tests.test_delegated_full_access import full_run # noqa: F401
|
||||
from tests.test_delegated_run_isolation import _git, _nanny_ctx, _seed_target
|
||||
|
||||
REPO = pathlib.Path(__file__).resolve().parents[1]
|
||||
|
||||
|
||||
def _lock_held(snaps: pathlib.Path) -> bool:
|
||||
return bool(_lock_identity(pathlib.Path(snaps) / wt._LOCK_NAME))
|
||||
|
||||
|
||||
def _provision(target, snaps, data, snapshot_id="snap1", task_id="t1"):
|
||||
return wt.provision_execution_snapshot(
|
||||
target_root=target, task_id=task_id, snapshot_id=snapshot_id, worktree_root=snaps, data_dir=data)
|
||||
|
||||
|
||||
def _phase_spies(monkeypatch, snaps):
|
||||
"""Record whether the lock was held when each phase ran."""
|
||||
seen: dict = {}
|
||||
real_verdicts, real_copy, real_git, real_save = (
|
||||
capture.untracked_binary_verdicts, artifacts.copy_artifact_file, wt._git, wt._save_registry)
|
||||
|
||||
def spy_verdicts(root, rels, **kw):
|
||||
seen["classify"] = _lock_held(snaps)
|
||||
return real_verdicts(root, rels, **kw)
|
||||
|
||||
def spy_copy(source, destination, **kw):
|
||||
seen["copy"] = _lock_held(snaps)
|
||||
return real_copy(source, destination, **kw)
|
||||
|
||||
def spy_git(repo_dir, *args, **kw):
|
||||
for marker in ("update-index", "commit-tree", "reset", "update-ref", "worktree"):
|
||||
if marker in args:
|
||||
seen[marker + ("_add" if marker == "worktree" and "add" in args else "")] = _lock_held(snaps)
|
||||
return real_git(repo_dir, *args, **kw)
|
||||
|
||||
def spy_save(entries, data_dir=None):
|
||||
seen.setdefault("save_registry", []).append(_lock_held(snaps))
|
||||
return real_save(entries, data_dir)
|
||||
|
||||
monkeypatch.setattr(capture, "untracked_binary_verdicts", spy_verdicts)
|
||||
monkeypatch.setattr(artifacts, "copy_artifact_file", spy_copy)
|
||||
monkeypatch.setattr(wt, "_git", spy_git)
|
||||
monkeypatch.setattr(wt, "_save_registry", spy_save)
|
||||
return seen
|
||||
|
||||
|
||||
def test_tree_walk_runs_outside_the_lock_and_shared_metadata_inside(tmp_path, monkeypatch):
|
||||
target = _seed_target(tmp_path)
|
||||
snaps, data = tmp_path / "snaps", tmp_path / "data"
|
||||
seen = _phase_spies(monkeypatch, snaps)
|
||||
|
||||
handle = _provision(target, snaps, data)
|
||||
|
||||
assert seen["classify"] is False and seen["copy"] is False and seen["reset"] is False
|
||||
assert seen["update-ref"] is True and seen["worktree_add"] is True
|
||||
assert seen["save_registry"] == [True, True] # provisional row, then the final row
|
||||
assert not _lock_held(snaps)
|
||||
assert handle.provisioning_sec >= 0 and wt.find_execution_snapshot("snap1", data_dir=data)["file_baseline"] == {}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("same_target", [False, True])
|
||||
def test_a_parked_provision_does_not_block_another(tmp_path, monkeypatch, same_target):
|
||||
target_a = _seed_target(tmp_path)
|
||||
(tmp_path / "other").mkdir()
|
||||
target_b = target_a if same_target else _seed_target(tmp_path / "other")
|
||||
snaps, data = tmp_path / "snaps", tmp_path / "data"
|
||||
parked, release = threading.Event(), threading.Event()
|
||||
real_verdicts = capture.untracked_binary_verdicts
|
||||
observed: dict = {}
|
||||
|
||||
def spy_verdicts(root, rels, **kw):
|
||||
if threading.current_thread().name == "provision-a":
|
||||
parked.set()
|
||||
assert release.wait(timeout=60), "the parked provision was never released"
|
||||
else:
|
||||
observed["b_saw_lock_held"] = _lock_held(snaps)
|
||||
return real_verdicts(root, rels, **kw)
|
||||
|
||||
monkeypatch.setattr(capture, "untracked_binary_verdicts", spy_verdicts)
|
||||
errors: list = []
|
||||
|
||||
def run_a():
|
||||
try:
|
||||
_provision(target_a, snaps, data, snapshot_id="snapA", task_id="ta")
|
||||
except BaseException as exc: # pragma: no cover - surfaced by the assertion below
|
||||
errors.append(exc)
|
||||
|
||||
thread_a = threading.Thread(target=run_a, name="provision-a")
|
||||
thread_a.start()
|
||||
assert parked.wait(timeout=60)
|
||||
index_before = (target_b / ".git" / "index").read_bytes()
|
||||
head_before = _git(target_b, "rev-parse", "HEAD").stdout
|
||||
|
||||
b_done = threading.Event()
|
||||
|
||||
def run_b():
|
||||
_provision(target_b, snaps, data, snapshot_id="snapB", task_id="tb")
|
||||
b_done.set()
|
||||
|
||||
thread_b = threading.Thread(target=run_b, name="provision-b")
|
||||
thread_b.start()
|
||||
# With the old whole-function lock B would sit behind A forever; A is still parked.
|
||||
assert b_done.wait(timeout=120), "provision B did not finish while A was parked"
|
||||
assert observed["b_saw_lock_held"] is False
|
||||
release.set()
|
||||
thread_a.join(timeout=120)
|
||||
assert not errors and not thread_a.is_alive()
|
||||
assert (target_b / ".git" / "index").read_bytes() == index_before
|
||||
assert _git(target_b, "rev-parse", "HEAD").stdout == head_before
|
||||
rows = {row["snapshot_id"]: row for row in wt.list_worktrees(data_dir=data)}
|
||||
assert set(rows) == {"snapA", "snapB"}
|
||||
assert all(pathlib.Path(rows[s]["path"]).is_dir() for s in rows)
|
||||
|
||||
|
||||
_HOLDER = """
|
||||
import os, pathlib, sys, time
|
||||
sys.path.insert(0, sys.argv[3])
|
||||
from ouroboros.platform_layer import acquire_exclusive_file_lock
|
||||
fd = acquire_exclusive_file_lock(
|
||||
pathlib.Path(sys.argv[1]), timeout_sec=5, stale_sec=600,
|
||||
metadata=f"pid={os.getpid()} task=t-holder op=provision since=2026-09-24T00:00:00Z target=/tmp/with space")
|
||||
assert fd is not None
|
||||
print("HELD", flush=True)
|
||||
time.sleep(float(sys.argv[2]))
|
||||
"""
|
||||
|
||||
|
||||
def _hold_lock(snaps: pathlib.Path, seconds: float) -> subprocess.Popen:
|
||||
snaps.mkdir(parents=True, exist_ok=True)
|
||||
proc = subprocess.Popen(
|
||||
[sys.executable, "-c", _HOLDER, str(snaps / wt._LOCK_NAME), str(seconds), str(REPO)],
|
||||
stdout=subprocess.PIPE, text=True)
|
||||
assert proc.stdout.readline().strip() == "HELD"
|
||||
return proc
|
||||
|
||||
|
||||
def test_a_live_holder_is_a_typed_refusal_and_a_dead_holder_is_evicted(tmp_path, monkeypatch):
|
||||
target = _seed_target(tmp_path)
|
||||
snaps, data = tmp_path / "snaps", tmp_path / "data"
|
||||
monkeypatch.setattr(wt, "_LOCK_TIMEOUT_SEC", 0.5)
|
||||
holder = _hold_lock(snaps, 30)
|
||||
try:
|
||||
with pytest.raises(wt.WorktreeOpsLockBusy) as info:
|
||||
_provision(target, snaps, data)
|
||||
finally:
|
||||
holder.kill()
|
||||
holder.wait()
|
||||
busy = info.value
|
||||
assert busy.holder == {"pid": str(holder.pid), "task": "t-holder", "op": "provision",
|
||||
"since": "2026-09-24T00:00:00Z", "target": "/tmp/with space"}
|
||||
assert busy.waited_sec >= 0.5 and "t-holder" in str(busy)
|
||||
# Nothing was registered, pinned or checked out for the refused attempt.
|
||||
assert wt.find_execution_snapshot("snap1", data_dir=data) is None
|
||||
assert _git(target, "for-each-ref", "refs/ouroboros/").stdout == ""
|
||||
assert not list(snaps.glob("dlg_*"))
|
||||
|
||||
# A holder that died (SIGKILL, panic) left its lock file behind: the owner-aware
|
||||
# stale check evicts it at once instead of refusing for _LOCK_STALE_SEC.
|
||||
lock_path = snaps / wt._LOCK_NAME
|
||||
lock_path.write_text(f"pid={holder.pid} task=t-dead op=provision since=x target=y", encoding="utf-8")
|
||||
started = time.monotonic()
|
||||
handle = _provision(target, snaps, data)
|
||||
assert time.monotonic() - started < 10 and pathlib.Path(handle.path).is_dir()
|
||||
assert not _lock_held(snaps)
|
||||
|
||||
|
||||
def test_a_provisioning_refusal_is_definitely_unrun_and_names_the_holder(tmp_path, monkeypatch):
|
||||
from ouroboros import delegate_custody as custody
|
||||
from ouroboros.delegate_shared import delegate_payload
|
||||
from ouroboros.subagent_bootstrap import _startup_refusal_definite
|
||||
from ouroboros.tools.delegate_integration import _provision_snapshot
|
||||
|
||||
target = _seed_target(tmp_path)
|
||||
ctx = _nanny_ctx(tmp_path, target, monkeypatch)
|
||||
drive = custody.custody_root(ctx)
|
||||
|
||||
def busy(**kw):
|
||||
raise wt.WorktreeOpsLockBusy(tmp_path / "lock", 120.0, {"pid": "4242", "task": "t-x", "op": "provision"})
|
||||
|
||||
monkeypatch.setattr(wt, "provision_execution_snapshot", busy)
|
||||
handle, refusal = _provision_snapshot(ctx, drive, str(target), "inv-busy")
|
||||
payload = delegate_payload(refusal)
|
||||
assert handle is None and payload["reason"] == "execution_snapshot_failed"
|
||||
assert payload["definitely_unrun"] is True and payload["cause"] == "lock_busy"
|
||||
assert payload["holder"]["pid"] == "4242" and payload["waited_sec"] == 120.0 and payload["retryable"] is True
|
||||
assert "t-x" in payload["detail"] and _startup_refusal_definite(payload)
|
||||
|
||||
def broken(**kw):
|
||||
raise RuntimeError("git exploded")
|
||||
|
||||
monkeypatch.setattr(wt, "provision_execution_snapshot", broken)
|
||||
_handle, refusal = _provision_snapshot(ctx, drive, str(target), "inv-broken")
|
||||
payload = delegate_payload(refusal)
|
||||
assert payload["definitely_unrun"] is True and "cause" not in payload
|
||||
assert _startup_refusal_definite(payload)
|
||||
|
||||
|
||||
def test_a_provisioning_refusal_leaves_a_durable_start_failed_row(full_run, monkeypatch): # noqa: F811
|
||||
"""The incident's two bootstrap refusals existed only inside the child's first
|
||||
prompt: no custody row, no event. A refused provision now settles its
|
||||
invocation durably, and never reaches the daemon."""
|
||||
from ouroboros import delegate_custody as custody
|
||||
from ouroboros.delegate_shared import delegate_payload
|
||||
from ouroboros.tools import delegate
|
||||
|
||||
ctx, _target, facts = full_run
|
||||
|
||||
def refused(**_kw):
|
||||
raise wt.WorktreeOpsLockBusy(pathlib.Path("/lock"), 120.0, {"pid": "7", "task": "t-other", "op": "provision"})
|
||||
|
||||
monkeypatch.setattr(wt, "provision_execution_snapshot", refused)
|
||||
payload = delegate_payload(delegate._delegate_start(ctx, "Implement the fixture."))
|
||||
|
||||
assert payload["reason"] == "execution_snapshot_failed" and payload["definitely_unrun"] is True
|
||||
assert payload["cause"] == "lock_busy" and payload["holder"]["task"] == "t-other"
|
||||
assert facts["requests"] == [], "a refused provision must never POST a run"
|
||||
rows = [json.loads(line) for line in custody.event_log_path(custody.custody_root(ctx))
|
||||
.read_text(encoding="utf-8").splitlines() if line.strip()]
|
||||
failed = [row for row in rows if row.get("type") == custody.START_FAILED]
|
||||
assert len(failed) == 1 and failed[0]["definite"] is True and failed[0]["invocation_id"]
|
||||
assert failed[0]["reason"] == "execution_snapshot_failed" and failed[0]["run_id"] == ""
|
||||
|
||||
|
||||
def test_removal_deletes_files_outside_the_lock_and_forgets_metadata_inside(tmp_path, monkeypatch):
|
||||
target = _seed_target(tmp_path)
|
||||
snaps, data = tmp_path / "snaps", tmp_path / "data"
|
||||
handle = _provision(target, snaps, data)
|
||||
seen: dict = {}
|
||||
real_rmtree, real_git = wt._force_rmtree, wt._git
|
||||
|
||||
def spy_rmtree(path):
|
||||
seen["rmtree"] = _lock_held(snaps)
|
||||
return real_rmtree(path)
|
||||
|
||||
def spy_git(repo_dir, *args, **kw):
|
||||
if "update-ref" in args and "-d" in args:
|
||||
seen["unpin"] = _lock_held(snaps)
|
||||
return real_git(repo_dir, *args, **kw)
|
||||
|
||||
monkeypatch.setattr(wt, "_force_rmtree", spy_rmtree)
|
||||
monkeypatch.setattr(wt, "_git", spy_git)
|
||||
assert wt.remove_execution_snapshot("snap1", worktree_root=snaps, data_dir=data)
|
||||
assert seen == {"rmtree": False, "unpin": True}
|
||||
assert not pathlib.Path(handle.path).exists()
|
||||
assert wt.find_execution_snapshot("snap1", data_dir=data) is None
|
||||
assert _git(target, "rev-parse", handle.baseline_ref, check=False).returncode != 0
|
||||
assert not _lock_held(snaps)
|
||||
|
||||
|
||||
def test_a_failure_after_the_provisional_row_leaves_nothing_and_a_crash_is_gc_reclaimable(tmp_path, monkeypatch):
|
||||
target = _seed_target(tmp_path)
|
||||
snaps, data = tmp_path / "snaps", tmp_path / "data"
|
||||
|
||||
def boom(*_a, **_k):
|
||||
raise OSError("copy failed")
|
||||
|
||||
monkeypatch.setattr(artifacts, "copy_artifact_file", boom)
|
||||
with pytest.raises(OSError, match="copy failed"):
|
||||
_provision(target, snaps, data, snapshot_id="snapFail")
|
||||
assert wt.find_execution_snapshot("snapFail", data_dir=data) is None
|
||||
assert _git(target, "for-each-ref", "refs/ouroboros/").stdout == ""
|
||||
assert not list(snaps.glob("dlg_*"))
|
||||
|
||||
# A worker killed between the provisional row and the final row: the row is
|
||||
# registered, custody never opened -> the startup GC removes checkout, ref, row.
|
||||
monkeypatch.setattr(wt, "_discard_snapshot_checkout", lambda *a, **k: None)
|
||||
with pytest.raises(OSError, match="copy failed"):
|
||||
_provision(target, snaps, data, snapshot_id="snapCrash")
|
||||
assert wt.find_execution_snapshot("snapCrash", data_dir=data) is not None
|
||||
assert "refs/ouroboros/delegated/snapCrash" in _git(target, "for-each-ref", "refs/ouroboros/").stdout
|
||||
report = wt.prune_execution_snapshots(set(), worktree_root=snaps, data_dir=data)
|
||||
assert report["removed"] == ["snapCrash"]
|
||||
assert wt.find_execution_snapshot("snapCrash", data_dir=data) is None
|
||||
assert _git(target, "for-each-ref", "refs/ouroboros/").stdout == ""
|
||||
assert not list(snaps.glob("dlg_*"))
|
||||
|
||||
|
||||
def test_populate_matches_worktree_add_and_runs_no_target_hook(tmp_path, monkeypatch):
|
||||
"""``worktree add --no-checkout`` + ``reset --hard --no-recurse-submodules`` is what
|
||||
git's own ``worktree add`` runs — a target with ``submodule.recurse=true`` must
|
||||
still snapshot — minus the target's post-checkout hook, which no longer executes
|
||||
project-authored code at provision."""
|
||||
sub = tmp_path / "sub"
|
||||
sub.mkdir()
|
||||
_git(sub, "init", "-q")
|
||||
(sub / "s.txt").write_text("s\n", encoding="utf-8")
|
||||
_git(sub, "add", "s.txt")
|
||||
_git(sub, "-c", "user.email=t@t", "-c", "user.name=t", "commit", "-qm", "sub")
|
||||
target = _seed_target(tmp_path)
|
||||
_git(target, "-c", "protocol.file.allow=always", "submodule", "add", "-q", str(sub), "vendored")
|
||||
_git(target, "-c", "user.email=t@t", "-c", "user.name=t", "commit", "-qm", "vendored")
|
||||
_git(target, "config", "submodule.recurse", "true")
|
||||
hooks = tmp_path / "hooks"
|
||||
hooks.mkdir()
|
||||
(hooks / "post-checkout").write_text("#!/bin/sh\ntouch .hook_ran\n", encoding="utf-8")
|
||||
(hooks / "post-checkout").chmod(0o755)
|
||||
_git(target, "config", "core.hooksPath", str(hooks))
|
||||
snaps, data = tmp_path / "snaps", tmp_path / "data"
|
||||
|
||||
handle = _provision(target, snaps, data)
|
||||
|
||||
exec_root = pathlib.Path(handle.path)
|
||||
assert (exec_root / "vendored").is_dir() and (exec_root / "tracked.txt").read_text(encoding="utf-8") == "one\ntwo\n"
|
||||
assert not (exec_root / ".hook_ran").exists() and not (target / ".hook_ran").exists()
|
||||
assert _git(exec_root, "status", "--porcelain").stdout == ""
|
||||
assert "160000" in _git(exec_root, "ls-files", "-s", "vendored").stdout
|
||||
|
||||
# The guard: without --no-recurse-submodules the populate fails on this target.
|
||||
real_git = wt._git
|
||||
|
||||
def recursing_git(repo_dir, *args, **kw):
|
||||
return real_git(repo_dir, *tuple(a for a in args if a != "--no-recurse-submodules"), **kw)
|
||||
|
||||
monkeypatch.setattr(wt, "_git", recursing_git)
|
||||
with pytest.raises(subprocess.CalledProcessError):
|
||||
_provision(target, snaps, data, snapshot_id="snapRecurse")
|
||||
|
||||
|
||||
def test_snapshot_facts_ride_the_start_receipt_as_disclosure_only(tmp_path):
|
||||
from ouroboros.tools.delegate import _snapshot_facts
|
||||
|
||||
target = _seed_target(tmp_path)
|
||||
(target / "blob.bin").write_bytes(b"\x7fELF\x00" * 10)
|
||||
handle = _provision(target, tmp_path / "snaps", tmp_path / "data")
|
||||
facts = _snapshot_facts(handle)
|
||||
assert set(facts) == {"entries", "untracked_files", "file_input_bytes", "provisioning_sec"}
|
||||
assert facts["entries"] == handle.entry_count and facts["untracked_files"] == 2 # untracked.txt + blob.bin
|
||||
assert facts["file_input_bytes"] == 50 and facts["provisioning_sec"] == handle.provisioning_sec
|
||||
assert _snapshot_facts(None) == {}
|
||||
# Facts, not a threshold: no runtime module consumes them to refuse or truncate.
|
||||
consumers = sorted(
|
||||
str(path.relative_to(REPO)) for path in (REPO / "ouroboros").rglob("*.py")
|
||||
if "provisioning_sec" in path.read_text(encoding="utf-8"))
|
||||
assert consumers == ["ouroboros/subagent_worktrees.py", "ouroboros/tools/delegate.py",
|
||||
"ouroboros/tools/delegate_integration.py"]
|
||||
131
tests/test_untracked_binary_verdict.py
Normal file
131
tests/test_untracked_binary_verdict.py
Normal file
|
|
@ -0,0 +1,131 @@
|
|||
"""git's binary verdict for an untracked inventory comes from ONE process (#1241).
|
||||
|
||||
One heavy delegated snapshot spawned ``git diff --no-index --numstat`` once per
|
||||
untracked file — 66,327 processes, 36 of the 40 minutes it held the machine-wide
|
||||
worktree lock. ``untracked_binary_verdicts`` stages every regular file as the
|
||||
empty blob into a scratch index and asks one index-versus-worktree ``git diff``,
|
||||
so the answer stays git's own (attributes, diff drivers, clean filters and
|
||||
working-tree encodings included) and parity with the per-file verdict is exact.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import pathlib
|
||||
import subprocess
|
||||
|
||||
from ouroboros import workspace_patch_capture as capture
|
||||
|
||||
|
||||
def _git(cwd, *args, **kw):
|
||||
return subprocess.run(["git", *args], cwd=str(cwd), capture_output=True, check=True, **kw)
|
||||
|
||||
|
||||
def _fixture_repo(root: pathlib.Path) -> tuple[pathlib.Path, list[str]]:
|
||||
"""Every class the verdict can be asked about, plus the config git consults."""
|
||||
repo = root / "repo"
|
||||
repo.mkdir()
|
||||
_git(repo, "init", "-q")
|
||||
_git(repo, "config", "user.email", "t@example.com")
|
||||
_git(repo, "config", "user.name", "T")
|
||||
_git(repo, "config", "diff.drvbin.binary", "true")
|
||||
_git(repo, "config", "diff.drvtext.binary", "false")
|
||||
_git(repo, "config", "filter.stripnul.clean", "tr -d '\\000'")
|
||||
(repo / ".gitattributes").write_text(
|
||||
"*.dat -diff\n*.bin binary\n*.txt diff\n*.drv diff=drvbin\n*.drt diff=drvtext\n"
|
||||
"*.lfs filter=stripnul\n*.u16 working-tree-encoding=UTF-16\n", encoding="utf-8")
|
||||
_git(repo, "add", ".gitattributes")
|
||||
_git(repo, "commit", "-qm", "attrs")
|
||||
files = {
|
||||
"plain.dat": b"hello\n", # -diff attribute: binary whatever the bytes
|
||||
"plain.bin": b"hello\n", # binary macro
|
||||
"nul.txt": b"x\0y", # diff set: text despite the NUL
|
||||
"nul.md": b"x\0y", # unspecified: NUL in the first 8000 bytes
|
||||
"late.md": b"a" * 9000 + b"\0", # NUL past git's probe: text
|
||||
"edge7999.md": b"a" * 7999 + b"\0", # NUL at offset 7999: binary
|
||||
"edge8000.md": b"a" * 8000 + b"\0", # NUL at offset 8000: text
|
||||
"nonul.drv": b"hello\n", # driver with binary=true and no NUL: binary
|
||||
"nul.drt": b"x\0y", # driver with binary=false and a NUL: text
|
||||
"nul.lfs": b"x\0y", # clean filter strips the NUL: text
|
||||
"text.u16": "hi\n".encode("utf-16"), # working-tree encoding: text
|
||||
"empty.md": b"",
|
||||
"target.bin2": b"x\0y",
|
||||
"new\nline.md": b"nl\n", # a newline in the name survives -z
|
||||
"vanish.md": b"gone\n",
|
||||
}
|
||||
for name, data in files.items():
|
||||
(repo / name).write_bytes(data)
|
||||
os.symlink("target.bin2", repo / "link_to_bin")
|
||||
os.symlink("nowhere", repo / "dangling")
|
||||
rels = sorted(files) + ["link_to_bin", "dangling"]
|
||||
return repo, rels
|
||||
|
||||
|
||||
def _oracle(repo: pathlib.Path, rel: str) -> bool:
|
||||
"""Today's per-file verdict: ``-\\t-`` from ``git diff --no-index --numstat``."""
|
||||
proc = subprocess.run(
|
||||
["git", "diff", "--no-index", "--numstat", "--no-ext-diff", "--no-color", "--", os.devnull, rel],
|
||||
cwd=str(repo), capture_output=True)
|
||||
first = proc.stdout.decode("utf-8", errors="replace").strip().splitlines()
|
||||
return bool(first) and first[0].startswith("-\t-")
|
||||
|
||||
|
||||
def test_batch_verdict_matches_git_for_every_path_class(tmp_path):
|
||||
repo, rels = _fixture_repo(tmp_path)
|
||||
oracle = {rel: _oracle(repo, rel) for rel in rels}
|
||||
(repo / "vanish.md").unlink() # vanished between the listing and the verdict: text, like git says
|
||||
oracle["vanish.md"] = False
|
||||
warnings: list = []
|
||||
|
||||
verdicts = capture.untracked_binary_verdicts(repo, rels, warnings=warnings)
|
||||
|
||||
assert verdicts is not None and warnings == []
|
||||
assert {rel: verdicts.get(rel, False) for rel in rels} == oracle
|
||||
# The classes that a NUL sniff or an attribute lookup alone would get wrong.
|
||||
assert verdicts["nonul.drv"] and not verdicts["nul.drt"] and not verdicts["nul.lfs"] and not verdicts["text.u16"]
|
||||
assert verdicts["edge7999.md"] and not verdicts["edge8000.md"]
|
||||
assert "link_to_bin" not in verdicts and "dangling" not in verdicts # symlinks: text, never followed
|
||||
# Only the empty blob entered the target's object database: no content was hashed.
|
||||
objects = _git(repo, "count-objects").stdout.decode()
|
||||
assert objects.startswith("4 objects"), objects
|
||||
|
||||
|
||||
def test_capture_asks_one_process_and_never_one_per_file(tmp_path, monkeypatch):
|
||||
repo, rels = _fixture_repo(tmp_path)
|
||||
calls: list = []
|
||||
real_run = subprocess.run
|
||||
|
||||
def spy(cmd, *args, **kwargs):
|
||||
calls.append(list(cmd))
|
||||
return real_run(cmd, *args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(capture.subprocess, "run", spy)
|
||||
artifacts, manifest = capture.write_workspace_patch_artifacts(repo, tmp_path / "artifacts", task={})
|
||||
|
||||
assert manifest["status"] == "ready_with_changes", manifest["errors"]
|
||||
numstat_calls = [c for c in calls if "--no-index" in c and "--numstat" in c]
|
||||
assert numstat_calls == [], "the per-file --no-index --numstat spawn is back"
|
||||
assert sum(1 for c in calls if c[:2] == ["git", "diff"] and "--numstat" in c) == 1
|
||||
excluded = {row["path"]: row["reason"] for row in manifest["untracked_excluded"]}
|
||||
expected_binary = {rel for rel in rels if _oracle(repo, rel)}
|
||||
assert {rel for rel, reason in excluded.items() if reason == "binary file"} == expected_binary
|
||||
|
||||
|
||||
def test_a_failed_batch_falls_back_to_the_per_file_verdict_with_a_warning(tmp_path, monkeypatch):
|
||||
repo, rels = _fixture_repo(tmp_path)
|
||||
real_run = subprocess.run
|
||||
|
||||
def broken_read_tree(cmd, *args, **kwargs):
|
||||
if list(cmd[:2]) == ["git", "read-tree"]:
|
||||
raise OSError("scratch index unavailable")
|
||||
return real_run(cmd, *args, **kwargs)
|
||||
|
||||
monkeypatch.setattr(capture.subprocess, "run", broken_read_tree)
|
||||
warnings: list = []
|
||||
assert capture.untracked_binary_verdicts(repo, rels, warnings=warnings) is None
|
||||
assert warnings and warnings[0]["reason"] == "binary_verdict_batch_unavailable"
|
||||
# ``None`` keeps git's per-file verdict: the same answer, one process per file.
|
||||
for rel in ("plain.dat", "nul.drt", "late.md"):
|
||||
reason = capture.untracked_capture_veto_reason(repo, rel, binary_verdicts=None)
|
||||
assert (reason == "binary file") == _oracle(repo, rel), rel
|
||||
assert capture.untracked_binary_verdicts(repo, [], warnings=warnings) == {}
|
||||
Loading…
Add table
Add a link
Reference in a new issue