feat(runtime): enforce governed mutation boundaries

Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a06ba0-10aa-7ea0-b20a-4f3fac39efe9
This commit is contained in:
tegwick 2026-09-04 11:25:07 +02:00
parent e3c6124e22
commit 20e6f381f6
28 changed files with 1068 additions and 154 deletions

View file

@ -3,8 +3,8 @@
Consuming repos declare instances in `.kaizen/schedule.yml`; this package
runs them: resolve tool profile + budget, load persona orientation
(kaizen-agentic schedule prepare), run a bounded agentic session via an
llm-connect adapter, verify the local commit, write `.kaizen/metrics`,
and report to the Custodian State Hub.
llm-connect adapter, validate granted local commits, record metrics, and report
to the Custodian State Hub.
"""
__version__ = "0.1.0"

View file

@ -12,7 +12,9 @@ How to add a row:
from __future__ import annotations
import os
from dataclasses import dataclass, field
from datetime import date
from pathlib import Path
from typing import Any
@ -35,6 +37,21 @@ APPROACH_MAIL_TRIAGE = "mail-triage"
APPROACH_MAIL_PIPELINE = "mail-scan+triage"
APPROACH_AGENT_SESSION = "agent-session"
APPROACH_UNMATCHED = "unmatched"
LEGACY_APPROACHES_UNTIL_ENV = "AGENT_HARNESS_LEGACY_APPROACHES_UNTIL"
def legacy_approaches_enabled(
value: str | None = None,
*,
today: date | None = None,
) -> bool:
"""Require an explicit, non-expired ISO date for profile-absent routing."""
raw = os.environ.get(LEGACY_APPROACHES_UNTIL_ENV, "") if value is None else value
try:
expires = date.fromisoformat(raw.strip())
except (AttributeError, ValueError):
return False
return (today or date.today()) <= expires
@dataclass(frozen=True)
@ -185,6 +202,19 @@ def execute_approach(
reopen=False,
)
if not legacy_approaches_enabled():
configured = os.environ.get(LEGACY_APPROACHES_UNTIL_ENV, "").strip()
state = f"expired at {configured}" if configured else "not configured"
return ApproachResult(
ok=False,
approach=name,
reason=(
"refused: profile-absent compatibility routing is disabled "
f"({LEGACY_APPROACHES_UNTIL_ENV} {state})"
),
reopen=False,
)
try:
target = resolve_ops_target(run, cfg)
except TaskSpecError as exc:
@ -266,6 +296,7 @@ def execute_approach(
reopen=True,
)
def _run_fi(target: Path, *, report_to_hub: bool, commit: bool) -> ApproachResult:
from rein_aharness.fi_research_brief import run_fi_research_brief

View file

@ -305,7 +305,10 @@ def main(argv: list[str] | None = None) -> int:
run.add_argument(
"--no-metrics",
action="store_true",
help="Skip writing .kaizen/metrics in the target repo",
help=(
"Skip compatibility .kaizen/metrics writes "
"(not permitted with repository_grant)"
),
)
run.add_argument(
"--stream-tool-events",

View file

@ -420,20 +420,6 @@ def run_fi_research_brief(
)
committed = True
head_after = _git(repo, "rev-parse", "HEAD")
# Best-effort push so workstation / Forgejo see the brief.
# Diverged branches must not fail the run; log via reason only if push fails
# after a successful write.
if committed and os.environ.get("FI_RESEARCH_BRIEF_PUSH", "1").strip().lower() not in {
"0",
"false",
"no",
"off",
}:
try:
_git(repo, "push", "origin", "HEAD")
except FiResearchBriefError:
# Leave brief committed locally; operators reconcile git separately.
pass
except FiResearchBriefError as exc:
result = FiResearchBriefResult(
ok=False,

View file

@ -1,7 +1,8 @@
"""Write per-run records into the target repo's `.kaizen/metrics` tree.
"""Write per-run metrics to compatibility or durable external storage.
Follows kaizen-agentic ADR-004 conventions so the optimization loop can
observe harness-run agents:
Legacy grant-absent runs retain the kaizen-agentic ADR-004 repository layout.
Accepted grant runs use private external state and a projection descriptor so
metrics cannot dirty the validated checkout.
.kaizen/metrics/<agent>/
executions.jsonl # append-only
@ -10,7 +11,11 @@ observe harness-run agents:
from __future__ import annotations
import fcntl
import hashlib
import json
import os
import uuid
from dataclasses import asdict, dataclass, field
from datetime import datetime, timezone
from pathlib import Path
@ -46,6 +51,20 @@ def metrics_dir(project_root: Path, agent: str) -> Path:
return Path(project_root) / ".kaizen" / "metrics" / agent
def external_metrics_dir(
project_root: Path | str,
agent: str,
*,
state_dir: Path | None = None,
) -> Path:
"""Return the private durable metrics directory for a granted run."""
root = Path(project_root).expanduser().resolve()
base = state_dir.expanduser().resolve() if state_dir else _state_dir()
repo_id = hashlib.sha256(str(root).encode("utf-8")).hexdigest()[:32]
agent_id = hashlib.sha256(str(agent).encode("utf-8")).hexdigest()[:32]
return base / "execution-metrics" / repo_id / agent_id
def _utc_now() -> str:
return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace(
"+00:00", "Z"
@ -133,18 +152,17 @@ def record_execution(
directory = metrics_dir(root, agent)
directory.mkdir(parents=True, exist_ok=True)
record = ExecutionRecord(
timestamp=_utc_now(),
agent=agent,
record = _execution_record(
root,
agent,
success=success,
execution_time_s=float(execution_time_s),
session_id=session_id,
metadata=metadata or {},
repo=root.name,
execution_time_s=execution_time_s,
tokens=tokens,
committed=committed,
head_after=head_after,
reason=reason,
metadata=metadata,
session_id=session_id,
)
executions_path = directory / "executions.jsonl"
@ -158,3 +176,189 @@ def record_execution(
encoding="utf-8",
)
return executions_path
def record_external_execution(
project_root: Path | str,
agent: str,
*,
success: bool,
execution_time_s: float = 0.0,
tokens: int | None = None,
committed: bool | None = None,
head_after: str | None = None,
reason: str | None = None,
metadata: dict[str, Any] | None = None,
session_id: str | None = None,
state_dir: Path | None = None,
) -> Path:
"""Durably record granted-run metrics without dirtying the target checkout."""
root = Path(project_root).expanduser().resolve()
projection_target = _projection_target(agent)
directory = external_metrics_dir(root, agent, state_dir=state_dir)
_ensure_private_directory(directory)
lock_path = directory / ".lock"
lock_fd = os.open(lock_path, os.O_RDWR | os.O_CREAT, 0o600)
try:
os.fchmod(lock_fd, 0o600)
fcntl.flock(lock_fd, fcntl.LOCK_EX)
record = _execution_record(
root,
agent,
success=success,
execution_time_s=execution_time_s,
tokens=tokens,
committed=committed,
head_after=head_after,
reason=reason,
metadata=metadata,
session_id=session_id,
)
executions_path = directory / "executions.jsonl"
existing = (
executions_path.read_text(encoding="utf-8")
if executions_path.exists()
else ""
)
if existing and not existing.endswith("\n"):
raise OSError("external metrics ledger has an incomplete record")
records = _load_external_executions(executions_path) if existing else []
already_recorded = session_id is not None and any(
item.get("session_id") == session_id for item in records
)
if not already_recorded:
record_line = record.to_json_line()
_atomic_write_text(executions_path, existing + record_line + "\n")
records.append(json.loads(record_line))
_atomic_write_text(
directory / "summary.json",
json.dumps(regenerate_summary(agent, records), indent=2, sort_keys=True)
+ "\n",
)
repo_id = directory.parent.name
_atomic_write_text(
directory / "projection.json",
json.dumps(
{
"agent": agent,
"repository_id": repo_id,
"repository_name": root.name,
"target_relative_directory": projection_target,
},
indent=2,
sort_keys=True,
)
+ "\n",
)
_fsync_directory(directory)
return executions_path
finally:
try:
fcntl.flock(lock_fd, fcntl.LOCK_UN)
finally:
os.close(lock_fd)
def _projection_target(agent: str) -> str:
"""Return a checkout-relative projection target for a safe agent identity."""
if (
not agent
or agent in {".", ".."}
or "/" in agent
or "\\" in agent
or "\x00" in agent
or len(agent) > 200
):
raise OSError("agent identity cannot form a safe metrics projection")
return f".kaizen/metrics/{agent}"
def _execution_record(
root: Path,
agent: str,
*,
success: bool,
execution_time_s: float,
tokens: int | None,
committed: bool | None,
head_after: str | None,
reason: str | None,
metadata: dict[str, Any] | None,
session_id: str | None,
) -> ExecutionRecord:
return ExecutionRecord(
timestamp=_utc_now(),
agent=agent,
success=success,
execution_time_s=float(execution_time_s),
session_id=session_id,
metadata=metadata or {},
repo=root.name,
tokens=tokens,
committed=committed,
head_after=head_after,
reason=reason,
)
def _load_external_executions(path: Path) -> list[dict[str, Any]]:
records: list[dict[str, Any]] = []
for line in path.read_text(encoding="utf-8").splitlines():
try:
value = json.loads(line)
except json.JSONDecodeError as exc:
raise OSError("external metrics ledger contains invalid JSON") from exc
if not isinstance(value, dict):
raise OSError("external metrics ledger contains a non-object record")
records.append(value)
return records
def _state_dir() -> Path:
explicit = os.environ.get("REIN_AHARNESS_STATE_DIR", "").strip()
if explicit:
return Path(explicit).expanduser().resolve()
xdg = os.environ.get("XDG_STATE_HOME", "").strip()
if xdg:
return (Path(xdg).expanduser() / "rein-aharness").resolve()
return (Path.home() / ".local" / "state" / "rein-aharness").resolve()
def _ensure_private_directory(path: Path) -> None:
path.mkdir(parents=True, exist_ok=True)
current = path
while current.name and current != current.parent:
os.chmod(current, 0o700)
if current.name == "execution-metrics":
break
current = current.parent
def _atomic_write_text(path: Path, value: str) -> None:
temporary = path.parent / f".{path.name}.{uuid.uuid4().hex}.tmp"
fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
try:
with os.fdopen(fd, "w", encoding="utf-8") as handle:
fd = -1
handle.write(value)
handle.flush()
os.fsync(handle.fileno())
os.replace(temporary, path)
os.chmod(path, 0o600)
_fsync_directory(path.parent)
except BaseException:
if fd >= 0:
os.close(fd)
try:
temporary.unlink()
except FileNotFoundError:
pass
raise
def _fsync_directory(path: Path) -> None:
fd = os.open(path, os.O_RDONLY)
try:
os.fsync(fd)
finally:
os.close(fd)

View file

@ -1,10 +1,4 @@
"""Versioned authority contract for repository mutation.
The grant is parsed by TaskSpec but is not yet wired into run execution. A
supplied grant therefore causes run_task to refuse before adapter dispatch.
This keeps the contract reviewable without implying enforcement that the live
runner does not yet provide.
"""
"""Versioned authority contract for accepted local repository mutation."""
from __future__ import annotations

View file

@ -1,9 +1,9 @@
"""Process-safe transaction boundary for a local Git checkout.
Wired into `run_task`, profile-absent `execute_approach` mutators, and the
profiled claim path under REINAH-WP-0003-T02 / ADR-002. Repository
acceptance (T03) remains a separate read-only validator and is not applied
to live results yet.
profiled claim path under REINAH-WP-0003-T02 / ADR-002. Explicit local
TaskSpec grants activate repository acceptance; queued/profiled grant carriage
remains an upstream contract dependency.
"""
from __future__ import annotations

View file

@ -2,10 +2,10 @@
Flow: lock target repo → resolve tool profile / budget from instance
manifest → snapshot HEAD → persona bundle → prompt → agentic session →
verify a new commit exists → kaizen metrics + hub progress event
(+ task close). The run *fails* if the session pushed anywhere or left
the repo dirty in a way it should not — the worker never pushes;
publishing is a separate, explicitly-granted lane.
validate an explicit repository grant when present → metrics + hub progress
event (+ task close). Granted runs use durable external metrics so repository
acceptance remains clean. The worker never pushes; publishing is a separate,
explicitly-granted lane.
"""
from __future__ import annotations
@ -24,6 +24,7 @@ from rein_aharness.profiles import UnknownToolProfileError, get_profile
from rein_aharness.repository_transaction import (
DirtyRepositoryError,
GitRepositoryError,
RepositoryAcceptanceError,
RepositoryBusyError,
RepositoryTransaction,
RepositoryTransactionError,
@ -116,11 +117,11 @@ def run_task(
budget_tokens_override: int | None = None,
transaction: RepositoryTransaction | None = None,
) -> RunResult:
if spec.repository_grant is not None:
if spec.repository_grant is not None and not write_metrics:
return _refused_run(
reason=(
"refused: repository_grant enforcement is not enabled; "
"no adapter was dispatched"
"refused: repository_grant runs require durable external metrics; "
"--no-metrics is incompatible"
),
model=model,
)
@ -255,11 +256,62 @@ def _execute_locked_task(
head_after = _git(spec.target_repo, "rev-parse", "HEAD")
committed = head_after != head_before
if session_ok and spec.repository_grant is not None:
try:
tx.validate_acceptance(spec.repository_grant.acceptance_policy())
except RepositoryAcceptanceError as exc:
session_ok = False
reason = str(exc)
ok = session_ok and committed
if session_ok and not committed:
if session_ok and not committed and spec.repository_grant is None:
reason = "session completed without committing"
tokens_spent = budget_tracker.spent if budget_tracker is not None else None
transaction_evidence = tx.evidence()
if spec.repository_grant is not None:
transaction_evidence["repository_grant"] = spec.repository_grant.evidence()
metric_metadata = {
"task_title": spec.title,
"tool_profile": profile.name,
"labels": list(spec.labels),
"completion_event_type": spec.completion_event_type,
}
if spec.repository_grant is not None:
metric_metadata["repository_grant_id"] = spec.repository_grant.grant_id
if write_metrics:
try:
recorder = (
metrics.record_external_execution
if spec.repository_grant is not None
else metrics.record_execution
)
recorder(
spec.target_repo,
spec.agent,
success=ok,
execution_time_s=execution_time_s,
tokens=tokens_spent,
committed=committed,
head_after=head_after,
reason=reason or None,
metadata=metric_metadata,
session_id=tx.transaction_id,
)
if spec.repository_grant is not None:
transaction_evidence["metrics"] = {
"storage": "external",
"session_id": tx.transaction_id,
"projection_ready": True,
}
except OSError as exc:
if spec.repository_grant is not None:
ok = False
reason = (
"required external metrics persistence failed "
f"({type(exc).__name__})"
)
result = RunResult(
ok=ok,
@ -275,30 +327,9 @@ def _execute_locked_task(
tokens_spent=tokens_spent,
execution_time_s=execution_time_s,
tool_events=collected_events,
transaction=tx.evidence(),
transaction=transaction_evidence,
)
if write_metrics:
try:
metrics.record_execution(
spec.target_repo,
spec.agent,
success=ok,
execution_time_s=execution_time_s,
tokens=tokens_spent,
committed=committed,
head_after=head_after,
reason=reason or None,
metadata={
"task_title": spec.title,
"tool_profile": profile.name,
"labels": list(spec.labels),
"completion_event_type": spec.completion_event_type,
},
)
except OSError:
pass # metrics must not block run completion reporting
if report_to_hub:
detail = {
"repo": spec.target_repo.name,
@ -314,6 +345,7 @@ def _execute_locked_task(
"budget_tokens": budget_tokens,
"tokens_spent": tokens_spent,
"execution_time_s": round(execution_time_s, 3),
"repository_transaction": transaction_evidence,
}
hub.post_progress_event(
summary=f"executor run: {spec.title} ({'ok' if ok else 'failed'})",