ouroboros/supervisor/evolution_lifecycle.py
Ouroboros 7726fca3ff Fix the rc.11 tag CI: Windows and macOS full-test
Four classes, none Linux-visible (the exact-SHA battery on this host was
green):

- Windows has no os.fchmod: the live stand's write_settings raised on
  five tests. The 0600-before-content write keeps fchmod where it exists
  and falls back to chmod after the write elsewhere (the shape the
  system_e2e harness already uses).
- The orphan scan reads /proc environ (Linux only): run_lane's finally
  raised FileNotFoundError on macOS and Windows. Without procfs the lane
  records a typed fact (orphan_scan=unavailable:no_procfs) and no check —
  never a passed check that did not run.
- supervisor/evolution_lifecycle._write_evolution_campaign treated a
  campaign file that exists but cannot be read as "no campaign" and let
  a stale write win (windows-latest: test_stale_campaign_cannot_overwrite_
  a_new_campaign — a transient read failure is enough). Present-but-
  unreadable now refuses the write with a warning; an absent file stays
  writable.
- macos-latest: the password resolution pin answered '' with settings
  patched to a password — the module wrapper had been replaced on that
  xdist worker by a started-and-never-stopped patch from an earlier
  module. The resolution order is now a pure function
  (resolve_network_password(env_value, settings_loader)) that the
  wrapper calls with os.environ and load_settings; the pin exercises the
  pure function and cannot be reached by either polluter class (a leaked
  environment writer or a leaked patch), and a second test pins the
  wrapper through the resolver seam.

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
2026-09-05 02:54:15 +00:00

1498 lines
65 KiB
Python

"""Evolution campaign state and transaction lifecycle.
Owns the campaign file (``state/evolution_campaign.json``), the campaign/
transaction lifecycle transitions, and the owner-facing cycle reporting.
``supervisor.queue`` imports this module top-level; anything here that needs
queue state (``DRIVE_ROOT``, locks) must import the queue lazily at call time
to avoid a module-load cycle.
"""
from __future__ import annotations
import logging
import os
import pathlib
import uuid
from typing import Any, Dict, Optional
from ouroboros.cost_projection import honest_cost_pair_amount
from ouroboros.evolution_fingerprint import canonical_objective_fingerprint
from ouroboros.outcomes import normalize_outcome_axes
from ouroboros.utils import atomic_write_json, read_json_dict, utc_now_iso
log = logging.getLogger(__name__)
EVOLUTION_CAMPAIGN_FILE = pathlib.Path("state") / "evolution_campaign.json"
EVOLUTION_CAMPAIGN_CAS_TIMEOUT_SEC = 0.05
def current_evolution_boot_generation() -> str:
"""Return the existing custody generation used by boot reconciliation."""
try:
from ouroboros.process_custody import current_custody_session_id
return str(current_custody_session_id() or "")
except Exception:
return ""
def _evolution_campaign_path() -> pathlib.Path:
from supervisor import queue
return pathlib.Path(queue.DRIVE_ROOT) / EVOLUTION_CAMPAIGN_FILE
def _read_evolution_campaign() -> Dict[str, Any]:
data = read_json_dict(_evolution_campaign_path()) or {}
return data if isinstance(data, dict) else {}
def disable_evolution_projection() -> None:
"""Clear the scheduler projection without changing campaign history."""
from supervisor import state
state.update_state(lambda live: live.update({
"evolution_mode_enabled": False,
"post_task_autostop": False,
}))
def disable_evolution_authority(
action: str, *, campaign_id: str = "", task_id: str = "",
) -> None:
"""Clear an invalid scheduler projection and record the typed refusal."""
from supervisor import queue, state
disable_evolution_projection()
event = {
"ts": utc_now_iso(),
"type": "evolution_authority_missing",
"action": str(action or ""),
"campaign_id": str(campaign_id or ""),
}
if task_id:
event["task_id"] = str(task_id)
state.append_jsonl(pathlib.Path(queue.DRIVE_ROOT) / "logs" / "events.jsonl", event)
def deliver_pending_owner_report(notify: Any = None) -> None:
"""Deliver and clear a worker-staged owner report from the server process."""
try:
campaign = _read_evolution_campaign()
report = campaign.get("pending_owner_report")
if not isinstance(report, dict):
return
(notify or notify_owner_cycle_outcome)(campaign, report)
clear_pending_owner_report(report)
except Exception:
log.debug("failed to deliver pending owner report", exc_info=True)
def _deliver_pending_owner_report() -> None:
"""Queue-compatibility wrapper for the server-process owner report."""
from supervisor import queue as q
q.deliver_pending_owner_report(q.notify_owner_cycle_outcome)
def enqueue_evolution_task_if_needed() -> None:
"""Queue evolution only when idle, enabled, within budget, and not failure-paused."""
from supervisor import queue as q
q._deliver_pending_owner_report()
if q.PENDING or q.RUNNING:
return
st = q.load_state()
if not bool(st.get("evolution_mode_enabled")):
return
owner_chat_id = st.get("owner_chat_id")
if not owner_chat_id:
return
campaign = q._read_evolution_campaign()
from supervisor.state import update_state
has_authority = all(str(campaign.get(key) or "").strip() for key in ("id", "source"))
if campaign.get("status") != "active" or not has_authority:
q.disable_evolution_authority(
"bare_flag_disabled", campaign_id=str(campaign.get("id") or ""),
)
q.send_with_budget(
int(owner_chat_id),
"🧬 Evolution stayed off: the enable flag had no active campaign authority. Use /evolve start to begin a fresh campaign.",
)
return
active_tx = campaign.get("active_transaction") if isinstance(campaign.get("active_transaction"), dict) else {}
if active_tx and (
str(active_tx.get("commit_sha") or "").strip()
or str(active_tx.get("dispatch_status") or "") == "reaping"
):
return
# Defensive net: light mode must never run evolution even if the flag was
# left enabled (e.g. carried across a restart into light mode). Disable and
# pause once; entry points already refuse new starts up front.
block = q.evolution_block_reason()
if block:
q.pause_evolution_campaign("blocked in light runtime mode")
q.disable_evolution_projection()
q.send_with_budget(int(owner_chat_id), block)
return
consecutive_failures = int(st.get("evolution_consecutive_failures") or 0)
if consecutive_failures >= 3:
q.pause_evolution_campaign("paused after consecutive failures")
q.disable_evolution_projection()
q.send_with_budget(
int(owner_chat_id),
f"🧬⚠️ Evolution paused: {consecutive_failures} consecutive failures. "
f"Use /evolve start to resume after investigating the issue."
)
return
# BUG3: pause if the SAME objective has been re-proposed and no-op'd
# OBJECTIVE_REPEAT_CAP times without ever absorbing. This is a SEPARATE
# breaker from consecutive_failures above: that counter is reset to 0 by
# ANY non-failing cycle (events.py), so it cannot catch a self-maintenance
# loop where a blocked objective is re-proposed NON-consecutively
# (interleaved with other no_op work). The per-objective count is keyed on
# the same canonical fingerprint the transaction stamps, accumulates across
# non-consecutive recurrence, and is cleared only on a genuine absorb.
objective_repeat_counts = campaign.get("objective_repeat_counts") or {}
active_objective_fp = canonical_objective_fingerprint(
str(campaign.get("objective") or ""),
)
objective_repeats = (
int(objective_repeat_counts.get(active_objective_fp, 0))
if active_objective_fp else 0
)
if objective_repeats >= q.OBJECTIVE_REPEAT_CAP:
q.pause_evolution_campaign(
"paused: objective re-proposed without ever absorbing",
)
q.disable_evolution_projection()
q.send_with_budget(
int(owner_chat_id),
f"🧬⚠️ Evolution paused: the current objective ran {objective_repeats} reviewed "
f"cycles WITHOUT ever being absorbed — it keeps getting re-proposed and never lands "
f"(a self-maintenance loop, not progress). A plain resume won't help; use "
f"/evolve start with a DIFFERENT objective."
)
return
try:
remaining = q.budget_remaining(st, strict=True)
except Exception:
log.error("Evolution scheduling deferred: cost accounting unavailable", exc_info=True)
q.append_jsonl(q.DRIVE_ROOT / "logs" / "events.jsonl", {
"ts": utc_now_iso(), "type": "evolution_accounting_unavailable",
"action": "dispatch_deferred", "owner_visible": True,
})
return
if remaining < q.EVOLUTION_BUDGET_RESERVE:
q.pause_evolution_campaign("budget reserve reached")
q.disable_evolution_projection()
q.send_with_budget(
int(owner_chat_id),
f"💸 Evolution stopped: ${remaining:.2f} remaining "
f"(reserve ${q.EVOLUTION_BUDGET_RESERVE:.0f} for conversations).",
)
return
cycle = int(st.get("evolution_cycle") or 0) + 1
tid = uuid.uuid4().hex[:8]
transaction = q.begin_evolution_transaction(tid, cycle=cycle, campaign=campaign)
if not transaction:
q.disable_evolution_authority(
"transaction_attach_failed",
campaign_id=str(campaign.get("id") or ""), task_id=tid,
)
q.send_with_budget(
int(owner_chat_id),
"🧬 Evolution stayed off: the campaign changed before its next task could be attached. Start it again when ready.",
)
return
task = {
"id": tid, "type": "evolution",
"chat_id": int(owner_chat_id),
"text": q.build_evolution_task_text(cycle),
"metadata": {"evolution_transaction": transaction},
}
q.attach_task_contract(task)
q.enqueue_task(task)
def _record_cycle(live: Dict[str, Any]) -> None:
live["evolution_cycle"] = cycle
live["last_evolution_task_at"] = utc_now_iso()
update_state(_record_cycle)
def _write_evolution_campaign(
data: Dict[str, Any], *, expected_campaign_id: Optional[str] = None,
expected_active_transaction: Optional[Dict[str, Any]] = None,
_state_lock_held: bool = False, _lock_timeout_sec: float = 4.0,
) -> bool:
path = _evolution_campaign_path()
from supervisor import state
state.assert_test_data_path(path)
lock_fd = None
if not _state_lock_held:
lock_fd = state.acquire_file_lock(
state.STATE_LOCK_PATH, timeout_sec=float(_lock_timeout_sec),
)
if lock_fd is None:
return False
try:
current = read_json_dict(path)
if current is None and path.is_file(): # present but unreadable is NOT "no campaign"
log.warning("Refusing evolution campaign write: %s exists but is unreadable", path)
return False
current = current or {}
current_id = str(current.get("id") or "")
data_id = str(data.get("id") or "")
expected_id = data_id if expected_campaign_id is None else str(expected_campaign_id or "")
if current_id and current_id != expected_id:
log.warning(
"Refusing stale evolution campaign write: current=%s expected=%s",
current_id, expected_id,
)
return False
if (
current_id == data_id
and current.get("status") not in {"active", "paused"}
and data.get("status") in {"active", "paused"}
):
log.warning("Refusing stale resurrection of terminal evolution campaign %s", data_id)
return False
if expected_active_transaction is not None:
current_tx = current.get("active_transaction", {})
if current_tx != expected_active_transaction:
log.warning("Refusing stale evolution transaction write for campaign %s", data_id)
return False
path.parent.mkdir(parents=True, exist_ok=True)
atomic_write_json(path, data, trailing_newline=True)
return True
finally:
if lock_fd is not None:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def evolution_block_reason() -> str:
"""Refusal message when evolution may not run in the current runtime mode.
Evolution campaigns are self-modification work, so they require runtime
mode ``advanced`` or ``pro``. In ``light`` (conversation-only) mode they are
hard-blocked before any campaign state, queue entry, or expensive round.
Returns ``""`` when evolution is allowed.
"""
from ouroboros.config import get_runtime_mode
if get_runtime_mode() == "light":
return (
"🧬 Evolution campaigns are self-modification work and require runtime "
"mode 'advanced' or 'pro'. The runtime is in 'light' mode "
"(self-modification is disabled), so no campaign was started. Switch "
"the runtime mode in Settings to evolve."
)
return ""
def start_evolution_campaign(objective: str = "", *, source: str = "owner") -> Dict[str, Any]:
"""Start or resume the active evolution campaign."""
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return {}
try:
campaign = _read_evolution_campaign()
prior_campaign_id = str(campaign.get("id") or "")
now = utc_now_iso()
objective = str(objective or "").strip()
if campaign.get("status") not in {"active", "paused"}:
campaign = {
"schema_version": 1,
"id": uuid.uuid4().hex[:8],
"status": "active",
"objective": objective or "Autonomously improve Ouroboros by acting on the highest-value backlog or process-memory signal.",
"source": str(source or ""),
"started_at": now,
"updated_at": now,
"cycles_done": 0,
"absorbed_cycles_done": 0,
"objective_repeat_counts": {}, # BUG3: fp -> non-absorbing-cycle count
"dropped_objective_fps": [], # BUG3 Layer B: attempted-and-dropped objective fps
"budget_spent_usd": 0.0,
"last_task_id": "",
"progress_notes": "",
"completed_at": "",
"completion_reason": "",
}
else:
if objective:
campaign["objective"] = objective
if not str(campaign.get("source") or "").strip() and source:
campaign["source"] = str(source)
campaign["status"] = "active"
campaign["updated_at"] = now
generation = current_evolution_boot_generation()
if generation:
campaign["last_boot_reconcile_gen"] = generation
return campaign if _write_evolution_campaign(
campaign,
expected_campaign_id=prior_campaign_id,
_state_lock_held=True,
) else {}
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def pause_evolution_campaign(reason: str = "") -> Dict[str, Any]:
"""Pause the active evolution campaign without deleting its state.
RESUMABLE: a later ``/evolve start`` resumes the SAME campaign in place
(start_evolution_campaign treats ``paused`` as resumable). Used by the system
breakers (light mode, consecutive failures, objective-repeat cap, budget reserve,
restart-blocked) — NOT by an owner stop. For an owner stop use
``complete_evolution_campaign`` (terminal).
"""
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return {}
try:
campaign = _read_evolution_campaign()
if campaign:
campaign["status"] = "paused"
campaign["updated_at"] = utc_now_iso()
campaign["pause_reason"] = str(reason or "")
if not _write_evolution_campaign(campaign, _state_lock_held=True):
return {}
return campaign
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def complete_evolution_campaign(
reason: str = "", *, status: str = "stopped", cleanup_worktree: bool = True
) -> Dict[str, Any]:
"""Terminally CLOSE the active campaign — the OWNER-stop counterpart of the
resumable pause. ``status`` is non-{active,paused}, so a later ``/evolve start``
mints a FRESH campaign instead of resurrecting this one. Archives + pops any
in-flight ``active_transaction`` (and ``post_task_backlog_id``) so a terminally
stopped campaign carries no dangling commit for a boot reconcile to absorb. The
durable gate against autonomous re-arm is the ``evolution_owner_stopped`` state
flag set at the owner-stop sites (read by ``apply_pending_request``); this terminal
status is the observability/audit marker plus a clean campaign. Never raises.
``cleanup_worktree`` (default True) runs the deterministic per-cycle worktree reset
for an in-flight transaction. PANIC passes ``False``: the Emergency Stop Invariant
(BIBLE) forbids delaying panic, so panic must NOT run git stash/reset work before its
hard exit — the panic flag + boot reconcile own that recovery instead."""
try:
from supervisor import state
snapshot = _read_evolution_campaign()
snapshot_tx = snapshot.get("active_transaction")
cleanup_updates: Dict[str, Any] = {}
if cleanup_worktree and isinstance(snapshot_tx, dict):
try:
_cleanup_worktree_after_cycle(snapshot_tx, str(snapshot_tx.get("task_id") or ""))
cleanup_updates = {
key: value for key, value in snapshot_tx.items()
if key.startswith("cleanup_") or key == "recovery_hint"
}
except Exception:
pass
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(
state.STATE_LOCK_PATH, timeout_sec=0.001 if not cleanup_worktree else 4.0,
)
if lock_fd is None:
return {}
try:
campaign = _read_evolution_campaign()
if not campaign:
return campaign or {}
now = utc_now_iso()
tx = campaign.get("active_transaction")
if isinstance(tx, dict):
if (
cleanup_updates
and isinstance(snapshot_tx, dict)
and str(tx.get("transaction_id") or "")
== str(snapshot_tx.get("transaction_id") or "")
):
tx.update(cleanup_updates)
tx = {**tx, "cycle_outcome": tx.get("cycle_outcome") or "owner_stopped"}
append_unique_transaction(campaign, tx)
campaign.pop("active_transaction", None)
campaign.pop("post_task_backlog_id", None)
campaign.pop("pause_reason", None)
campaign["status"] = str(status or "stopped")
campaign["updated_at"] = now
campaign["completed_at"] = now
campaign["completion_reason"] = str(reason or "")
return campaign if _write_evolution_campaign(
campaign, _state_lock_held=True,
) else {}
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
except Exception:
log.debug("complete_evolution_campaign failed", exc_info=True)
return {}
def begin_evolution_transaction(task_id: str, *, cycle: int, campaign: Dict[str, Any]) -> Dict[str, Any]:
"""Attach a compact self-modification transaction to the active campaign."""
try:
from supervisor import git_ops
rc_head, head, _ = git_ops.git_capture(["git", "rev-parse", "HEAD"])
rc_branch, branch, _ = git_ops.git_capture(["git", "rev-parse", "--abbrev-ref", "HEAD"])
base_head = head.strip() if rc_head == 0 else ""
base_branch = branch.strip() if rc_branch == 0 else ""
except Exception:
base_head = ""
base_branch = ""
transaction = {
"schema_version": 2,
"transaction_id": uuid.uuid4().hex[:12],
"campaign_id": str((campaign or {}).get("id") or ""),
"task_id": str(task_id or ""),
"cycle": int(cycle or 0),
# BUG3: capture the objective this cycle will run, at cycle START, as the SSOT
# per-cycle fingerprint. Read here (not at outcome time) because campaign["objective"]
# can be overwritten by a later promotion before the outcome is recorded.
"objective_fp": canonical_objective_fingerprint(str((campaign or {}).get("objective") or "")),
"created_at": utc_now_iso(),
"updated_at": utc_now_iso(),
"base_head": base_head,
"base_branch": base_branch,
"preflight_status": "pending",
"advisory_status": "pending",
"triad_scope_status": "pending",
"commit_sha": "",
"push_status": "pending",
"restart_decision": "",
"restart_required": False,
"restart_verified": False,
"restart_verified_at": "",
"rescue_ref": "",
"rescue_path": "",
"recovery_hint": "",
}
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return {}
try:
current = _read_evolution_campaign()
live_state = state.json_load_file(state.STATE_PATH) or {}
existing_tx = current.get("active_transaction")
existing_tx = existing_tx if isinstance(existing_tx, dict) else {}
if (
bool(live_state.get("evolution_owner_stopped"))
or not bool(live_state.get("evolution_mode_enabled"))
or current.get("status") != "active"
or str(current.get("id") or "") != str(campaign.get("id") or "")
or bool(str(existing_tx.get("commit_sha") or "").strip())
):
return {}
if existing_tx:
existing_tx.update({
"cycle_outcome": "abandoned",
"abandoned_reason": "dispatch_not_persisted",
"updated_at": utc_now_iso(),
})
append_unique_transaction(current, existing_tx)
current["active_transaction"] = transaction
current["updated_at"] = utc_now_iso()
if not _write_evolution_campaign(current, _state_lock_held=True):
return {}
stored = _read_evolution_campaign()
stored_tx = stored.get("active_transaction")
if not isinstance(stored_tx, dict) or any(
str(stored_tx.get(key) or "") != str(transaction.get(key) or "")
for key in ("campaign_id", "transaction_id", "task_id")
):
return {}
return transaction
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def evolution_commit_receipt_error(
tx: Dict[str, Any], *, campaign_id: str, transaction_id: str,
task_id: str, commit_sha: str,
) -> str:
if str(tx.get("commit_sha") or "") != commit_sha:
return "commit_receipt_mismatch"
receipt = tx.get("commit_receipt")
if not isinstance(receipt, dict) or not bool(receipt.get("ok")):
return "commit_receipt_missing"
expected = {
"campaign_id": campaign_id,
"transaction_id": transaction_id,
"task_id": task_id,
"commit_sha": commit_sha,
}
if any(str(receipt.get(key) or "") != value for key, value in expected.items()):
return "commit_receipt_mismatch"
return ""
def _evolution_claim_error(
campaign: Dict[str, Any], live_state: Dict[str, Any], *,
campaign_id: str, transaction_id: str, task_id: str, commit_sha: str = "",
require_uncommitted: bool = False,
) -> str:
if not campaign_id or not transaction_id or not task_id:
return "claim_identity_missing"
if bool(live_state.get("evolution_owner_stopped")):
return "owner_stopped"
if not bool(live_state.get("evolution_mode_enabled")) and not commit_sha:
return "evolution_disabled"
if campaign.get("status") != "active":
return "campaign_not_active"
if str(campaign.get("id") or "") != campaign_id:
return "campaign_mismatch"
tx = campaign.get("active_transaction")
if not isinstance(tx, dict):
return "transaction_missing"
if str(tx.get("transaction_id") or "") != transaction_id:
return "transaction_mismatch"
if str(tx.get("task_id") or "") != task_id:
return "task_mismatch"
if require_uncommitted and str(tx.get("commit_sha") or "").strip():
return "transaction_already_committed"
if commit_sha:
return evolution_commit_receipt_error(
tx,
campaign_id=campaign_id,
transaction_id=transaction_id,
task_id=task_id,
commit_sha=commit_sha,
)
return ""
def check_evolution_authority(
campaign_id: str, transaction_id: str, task_id: str, *, commit_sha: str = "",
require_uncommitted: bool = False,
) -> Dict[str, Any]:
"""Check the exact campaign claim under the existing state lock."""
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return {"ok": False, "reason": "state_lock_unavailable"}
try:
campaign = _read_evolution_campaign()
live_state = state.json_load_file(state.STATE_PATH) or {}
reason = _evolution_claim_error(
campaign, live_state,
campaign_id=str(campaign_id or ""),
transaction_id=str(transaction_id or ""),
task_id=str(task_id or ""),
commit_sha=str(commit_sha or ""),
require_uncommitted=bool(require_uncommitted),
)
return {
"ok": not reason,
"reason": reason,
"campaign_id": str(campaign_id or ""),
"transaction_id": str(transaction_id or ""),
"task_id": str(task_id or ""),
"commit_sha": str(commit_sha or ""),
}
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def _build_commit_receipt(
campaign_id: str, transaction_id: str, task_id: str, commit_sha: str,
*, reason: str, recorded_at: str,
) -> Dict[str, Any]:
"""The exact receipt shape ``evolution_commit_receipt_error`` validates.
One constructor so the receipt written by the tool path and the one
re-derived by boot recovery are the same durable fact, differing only in
``reason`` (which records HOW the SHA was attributed).
"""
return {
"ok": True,
"reason": reason,
"campaign_id": str(campaign_id),
"transaction_id": str(transaction_id),
"task_id": str(task_id),
"commit_sha": str(commit_sha),
"recorded_at": recorded_at,
}
def adopt_evolution_commit_intent(
campaign: Dict[str, Any], tx: Dict[str, Any], head_sha: str = "",
) -> str:
"""Finish an interrupted commit receipt from the pre-commit intent.
``git commit`` and its SHA receipt are two writes; a crash between them used
to leave a reviewed commit on HEAD that no boot path could attribute (the
markerless reconcile short-circuits on an empty ``commit_sha``). The intent
written before the commit carries the exact reviewed tree and parents, so
attribution here is structural rather than a guess: the commit at HEAD is
adopted only when its tree AND its full parent list are identical to that
reviewed material. A failed commit, a contained orphan (the branch is rewound
to the parent) or any later HEAD movement fails the match and recovery stays
fail-closed. The caller owns persistence, so the recovered SHA and its receipt
land in the SAME write as the decision that consumed them.
"""
intent = tx.get("commit_intent") if isinstance(tx, dict) else None
if not isinstance(intent, dict):
return ""
tree_sha = str(intent.get("tree_sha") or "")
parents = [str(value) for value in (intent.get("parents") or [])]
if not tree_sha or not parents:
return ""
head = str(head_sha or "").strip() or "HEAD"
try:
from supervisor import git_ops
rc_tree, actual_tree, _ = git_ops.git_capture(["git", "rev-parse", f"{head}^{{tree}}"])
rc_parents, parent_line, _ = git_ops.git_capture(
["git", "rev-list", "--parents", "-n", "1", head]
)
except Exception:
return ""
fields = parent_line.strip().split() if rc_parents == 0 else []
if rc_tree != 0 or not fields or actual_tree.strip() != tree_sha or fields[1:] != parents:
return ""
commit_sha, now = fields[0], utc_now_iso()
tx.update({
"commit_sha": commit_sha,
"commit_receipt": _build_commit_receipt(
str(campaign.get("id") or ""), str(tx.get("transaction_id") or ""),
str(tx.get("task_id") or ""), commit_sha,
reason="recovered_from_commit_intent", recorded_at=now,
),
"commit_receipt_recovered": True,
"restart_required": True,
"restart_verified": False,
"updated_at": now,
})
campaign["active_transaction"] = tx
return commit_sha
def record_evolution_commit(
campaign_id: str, transaction_id: str, task_id: str, commit_sha: str,
) -> Dict[str, Any]:
"""CAS the exact reviewed local commit into its still-authorized transaction."""
from supervisor import state
commit_sha = str(commit_sha or "").strip()
if not commit_sha:
return {"ok": False, "reason": "commit_sha_missing"}
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return {"ok": False, "reason": "state_lock_unavailable", "commit_sha": commit_sha}
try:
live_state = state.json_load_file(state.STATE_PATH) or {}
from ouroboros.utils import update_json_locked
receipt: Dict[str, Any] = {}
refused = {"reason": "campaign_write_refused"}
def _mutate(campaign: Dict[str, Any]) -> Optional[Dict[str, Any]]:
reason = _evolution_claim_error(
campaign, live_state,
campaign_id=str(campaign_id or ""),
transaction_id=str(transaction_id or ""),
task_id=str(task_id or ""),
)
tx = campaign.get("active_transaction")
prior_sha = str((tx or {}).get("commit_sha") or "") if isinstance(tx, dict) else ""
if not reason and prior_sha and prior_sha != commit_sha:
reason = "commit_receipt_conflict"
if reason:
refused["reason"] = reason
return None
now = utc_now_iso()
receipt.update(_build_commit_receipt(
campaign_id, transaction_id, task_id, commit_sha,
reason="recorded", recorded_at=now,
))
tx.update({
"preflight_status": "passed",
"advisory_status": "fresh_or_bypassed",
"triad_scope_status": "passed",
"commit_sha": commit_sha,
"commit_receipt": dict(receipt),
"restart_required": True,
"restart_verified": False,
"updated_at": now,
})
campaign["active_transaction"] = tx
campaign["updated_at"] = now
return campaign
update_json_locked(
_evolution_campaign_path(), _mutate,
timeout_sec=EVOLUTION_CAMPAIGN_CAS_TIMEOUT_SEC,
strict_existing_dict=True,
)
if not receipt:
return {"ok": False, "reason": refused["reason"], "commit_sha": commit_sha}
return receipt
except Exception as exc:
log.warning("Failed to record exact evolution commit receipt", exc_info=True)
return {"ok": False, "reason": f"campaign_write_failed:{type(exc).__name__}", "commit_sha": commit_sha}
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def link_evolution_rescue(drive_root: pathlib.Path, rescue_info: Dict[str, Any]) -> Dict[str, Any]:
"""Attach rescue pointers without racing an exact commit receipt."""
from supervisor import state
root = pathlib.Path(drive_root)
path = root / EVOLUTION_CAMPAIGN_FILE
lock_path = root / "locks" / "state.lock"
state.assert_test_data_path(path)
lock_fd = state.acquire_file_lock(lock_path)
if lock_fd is None:
return {}
try:
from ouroboros.utils import update_json_locked
linked: Dict[str, Any] = {}
def _mutate(campaign: Dict[str, Any]) -> Optional[Dict[str, Any]]:
if campaign.get("status") not in {"active", "paused"}:
return None
tx = campaign.get("active_transaction")
if not isinstance(tx, dict):
return None
tx["rescue_ref"] = str(rescue_info.get("rescue_ref") or "")
tx["dirty_snapshot_ref"] = str(rescue_info.get("rescue_ref") or "")
tx["rescue_path"] = str(rescue_info.get("path") or "")
tx["restart_decision"] = "rescue_snapshot_created"
tx["recovery_hint"] = (
f"Recover with {tx['rescue_ref']} or inspect {tx['rescue_path']}"
if tx.get("rescue_ref") or tx.get("rescue_path")
else "Rescue attempted; inspect supervisor logs."
)
tx["updated_at"] = utc_now_iso()
campaign["active_transaction"] = tx
campaign["updated_at"] = utc_now_iso()
linked.update(tx)
return campaign
update_json_locked(
path, _mutate,
timeout_sec=EVOLUTION_CAMPAIGN_CAS_TIMEOUT_SEC,
strict_existing_dict=True,
)
return linked
finally:
state.release_file_lock(lock_path, lock_fd)
def _bump_objective_repeat_count(campaign: Dict[str, Any], tx: Dict[str, Any]) -> None:
"""BUG3: count one non-absorbing cycle against its objective fingerprint.
Cumulative PER-FINGERPRINT (not a consecutive streak), so a blocked objective that is
re-proposed NON-consecutively (interleaved with other no_op work) still accumulates toward
the pause gate. ``setdefault`` tolerates campaigns persisted before this field existed; a
transaction without an ``objective_fp`` (e.g. a tx-less idle cycle) is skipped, never
bucketed under the empty key.
"""
fp = str((tx or {}).get("objective_fp") or "")
if not fp:
return
counts = campaign.setdefault("objective_repeat_counts", {})
counts[fp] = int(counts.get(fp, 0) or 0) + 1
# Layer B: also mark this objective attempted-and-dropped so the chooser (Layer A) can be
# told not to re-propose it. This is a campaign-local signal, NOT a backlog status flip:
# the backlog item stays "open" (the work is genuinely unsolved), we only stop FEEDING it
# back to the evolution objective chooser.
dropped = campaign.setdefault("dropped_objective_fps", [])
if fp not in dropped:
dropped.append(fp)
def _clear_objective_repeat_count(campaign: Dict[str, Any], tx: Dict[str, Any]) -> None:
"""BUG3: a genuine absorb clears ONLY this objective's repeat tally and dropped flag.
Called at every site that sets ``cycle_outcome == "absorbed"`` (task-done here, plus the
two durable boot/restart-verify absorb sites in agent_startup_checks). Keyed on the same
SSOT fingerprint as the bump so success on the looping objective resets exactly its bucket
and un-drops it (it landed, so it is no longer do-not-re-propose).
"""
fp = str((tx or {}).get("objective_fp") or "")
if not fp:
return
counts = campaign.setdefault("objective_repeat_counts", {})
counts.pop(fp, None)
dropped = campaign.setdefault("dropped_objective_fps", [])
if fp in dropped:
dropped.remove(fp)
def update_evolution_transaction(task_id: str, **updates: Any) -> bool:
"""Best-effort update of the active/lightweight evolution transaction."""
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return False
try:
campaign = _read_evolution_campaign()
tx = campaign.get("active_transaction")
if not isinstance(tx, dict) or str(tx.get("task_id") or "") != str(task_id or ""):
return False
for key, value in updates.items():
if value is not None:
tx[key] = value
tx["updated_at"] = utc_now_iso()
campaign["active_transaction"] = tx
campaign["updated_at"] = utc_now_iso()
return _write_evolution_campaign(campaign, _state_lock_held=True)
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def _cleanup_worktree_after_cycle(tx: Dict[str, Any], task_id: str) -> None:
"""Deterministic worktree cleanup when a cycle closes WITHOUT absorption.
A no_op/abandoned evolution cycle must leave the repo at its recorded
``base_head``: abandoned edits or unreviewed local commits otherwise leak
into the next cycle (and into unrelated tasks) as mystery state. Recovery
is never silent — dirty files go into a git stash and an ahead HEAD is
preserved as a local branch before the hard reset; both refs are recorded
on the transaction. Skipped (with a recorded reason) when other tasks are
running in the shared worktree or the base is unknown. Never raises.
Kill-switch: OUROBOROS_EVOLUTION_CYCLE_CLEANUP=false.
"""
if str(os.environ.get("OUROBOROS_EVOLUTION_CYCLE_CLEANUP", "true") or "true").lower() in {"0", "false", "no", "off"}:
tx["cleanup_status"] = "disabled"
return
base_head = str(tx.get("base_head") or "").strip()
if not base_head:
tx["cleanup_status"] = "skipped_no_base"
return
update_lock_fh = None
release_update_lock = None
try:
from supervisor import git_ops, queue
from supervisor.update_merge import acquire_update_lock, release_update_lock
from supervisor.workers import repo_writer_admission_closed
try:
update_lock_fh = acquire_update_lock()
except RuntimeError:
tx["cleanup_status"] = "skipped_managed_update_lock_busy"
return
admission_reason = repo_writer_admission_closed()
if admission_reason:
tx["cleanup_status"] = "skipped_repo_writer_admission_closed"
tx["cleanup_deferred_reason"] = admission_reason
return
# Same protection class as git_ops._guard_live_repo_destructive_git, but
# covering the stash too: a unit test that never re-pointed
# git_ops.REPO_DIR must not stash/reset the LIVE repo's working tree.
if os.environ.get("OUROBOROS_ALLOW_LIVE_REPO_TESTS") != "1":
import sys as _sys
try:
live_repo = git_ops.REPO_DIR.resolve(strict=False) == (
pathlib.Path.home() / "Ouroboros" / "repo"
).resolve(strict=False)
except OSError:
live_repo = False
if live_repo and ("PYTEST_CURRENT_TEST" in os.environ or "pytest" in _sys.modules):
tx["cleanup_status"] = "skipped_live_repo_test_guard"
return
# Lock-free RUNNING snapshot is safe because this runs on the SAME
# single supervisor thread that assigns tasks (dispatch_event ->
# assign_tasks are sequential); only cancel paths mutate RUNNING from
# HTTP threads, which can only shrink the set.
running_others = [tid for tid in list(queue.RUNNING.keys()) if str(tid) != str(task_id)]
if running_others:
# The live worktree is shared: a reset would destroy concurrent
# tasks' work. Leave state for the boot reconcile / next cycle.
tx["cleanup_status"] = "skipped_other_tasks_running"
return
rc_status, status_out, _ = git_ops.git_capture(["git", "status", "--porcelain"])
rc_head, head_out, _ = git_ops.git_capture(["git", "rev-parse", "HEAD"])
if rc_status != 0 or rc_head != 0:
tx["cleanup_status"] = "skipped_git_unavailable"
return
dirty = bool(status_out.strip())
head = head_out.strip()
if not dirty and head == base_head:
tx["cleanup_status"] = "already_clean"
return
if dirty:
stash_label = f"evolution-cycle-cleanup-{tx.get('transaction_id') or task_id}"
rc_stash, _, stash_err = git_ops.git_capture(
["git", "stash", "push", "--include-untracked", "-m", stash_label]
)
if rc_stash != 0:
from ouroboros.utils import truncate_review_artifact
# Refuse to reset over unsaved changes (P1: no silent loss).
tx["cleanup_status"] = "skipped_stash_failed"
tx["recovery_hint"] = (
"worktree dirty and stash failed: "
+ truncate_review_artifact(str(stash_err or "").strip(), 400)
)
return
tx["cleanup_stash"] = stash_label
if head != base_head:
preserved, ref_name = git_ops.preserve_local_ref_branch("HEAD", prefix="evolution-leftover")
if not preserved:
tx["cleanup_status"] = "skipped_preserve_failed"
tx["recovery_hint"] = (
"HEAD ahead of base_head and preserve-branch failed; left as-is"
+ (f" (dirty files already saved in stash {tx['cleanup_stash']})." if tx.get("cleanup_stash") else ".")
)
return
tx["cleanup_preserved_ref"] = ref_name
rc_reset, _, reset_err = git_ops.git_capture(["git", "reset", "--hard", base_head])
if rc_reset != 0:
from ouroboros.utils import truncate_review_artifact
tx["cleanup_status"] = "reset_failed"
tx["recovery_hint"] = (
f"git reset --hard {base_head[:12]} failed: "
+ truncate_review_artifact(str(reset_err or "").strip(), 400)
)
return
tx["cleanup_status"] = "reset_to_base"
else:
tx["cleanup_status"] = "stashed_dirty" # HEAD already at base; only the stash happened
log.info(
"Evolution cycle %s cleanup: worktree restored to base %s (stash=%s, preserved=%s)",
tx.get("transaction_id") or task_id, base_head[:12],
tx.get("cleanup_stash") or "-", tx.get("cleanup_preserved_ref") or "-",
)
except Exception:
tx["cleanup_status"] = "error"
log.debug("Evolution cycle worktree cleanup failed", exc_info=True)
finally:
if update_lock_fh is not None and release_update_lock is not None:
release_update_lock(update_lock_fh)
def _persist_evolution_cleanup(campaign_id: str, tx: Dict[str, Any]) -> bool:
"""Patch cleanup evidence into the exact terminal transaction under the state lock."""
from supervisor import state
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return False
try:
campaign = _read_evolution_campaign()
if str(campaign.get("id") or "") != str(campaign_id or ""):
return False
tx_id = str(tx.get("transaction_id") or "")
cleanup = {
key: tx[key]
for key in (
"cleanup_status", "cleanup_stash", "cleanup_preserved_ref",
"recovery_hint",
)
if key in tx
}
updated = False
for item in list(campaign.get("transaction_history") or []):
if isinstance(item, dict) and str(item.get("transaction_id") or "") == tx_id:
item.update(cleanup)
updated = True
for row in list(campaign.get("history") or []):
item = row.get("transaction") if isinstance(row, dict) else None
if isinstance(item, dict) and str(item.get("transaction_id") or "") == tx_id:
item.update(cleanup)
updated = True
if not updated:
return False
campaign["updated_at"] = utc_now_iso()
return _write_evolution_campaign(
campaign, expected_campaign_id=campaign_id, _state_lock_held=True,
)
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def _resume_evolution_terminal_effects(
campaign_id: str, task_id: str, tx: Dict[str, Any],
) -> Dict[str, Any]:
"""Finish idempotent effects that follow the durable terminal transition."""
from supervisor import queue
resumed = dict(tx or {})
if resumed.get("cleanup_status") == "pending":
_cleanup_worktree_after_cycle(resumed, str(task_id or ""))
try:
if not _persist_evolution_cleanup(campaign_id, resumed):
log.warning("Failed to persist evolution cleanup evidence for %s", task_id)
except Exception:
log.warning("Failed to persist evolution cleanup evidence for %s", task_id, exc_info=True)
if (
resumed.get("cycle_outcome") == "waiting_for_restart"
and resumed.get("restart_required")
and not resumed.get("restart_verified")
):
try:
request_evolution_restart(queue.DRIVE_ROOT, resumed, log=log)
except Exception:
log.warning("Failed to resume evolution restart request for %s", task_id, exc_info=True)
# Abandon/absorb reports are staged in the same campaign write as the
# transition. Delivery clears only the exact report after a successful send.
deliver_pending_owner_report()
return resumed
def update_evolution_campaign_after_task(
task_id: str,
*,
cost_usd: Optional[float],
cost_accounting_status: str = "available",
outcome_axes: Dict[str, Any],
rounds: int,
transaction: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
"""Record an evolution cycle outcome in the active campaign file."""
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return {
"accepted": True, "persisted": False, "replay": False,
"reason": "state_lock_unavailable", "transaction": {},
}
replay_found = False
replay_tx: Dict[str, Any] = {}
campaign_id = ""
tx: Dict[str, Any] = {}
try:
# The complete read/check/mutate/write is protected by the same state lock
# used by pause, owner stop, boot reconciliation, and transaction updates.
campaign = _read_evolution_campaign()
if campaign.get("status") not in {"active", "paused"}:
return {
"accepted": False, "persisted": False, "replay": False,
"reason": "campaign_not_active", "transaction": {},
}
metadata_tx = transaction if isinstance(transaction, dict) else {}
active_tx = (
campaign.get("active_transaction")
if isinstance(campaign.get("active_transaction"), dict) else {}
)
campaign_id = str(campaign.get("id") or "")
def _matches(candidate: Dict[str, Any]) -> bool:
return bool(
campaign_id
and str(candidate.get("campaign_id") or "") == campaign_id
and str(candidate.get("transaction_id") or "")
and str(candidate.get("task_id") or "") == str(task_id or "")
)
metadata_tx_id = (
str(metadata_tx.get("transaction_id") or "") if _matches(metadata_tx) else ""
)
if metadata_tx_id:
for existing in list(campaign.get("history") or []):
if (
not isinstance(existing, dict)
or str(existing.get("task_id") or "") != str(task_id or "")
):
continue
existing_tx = existing.get("transaction")
if (
isinstance(existing_tx, dict)
and str(existing_tx.get("campaign_id") or "") == campaign_id
and str(existing_tx.get("transaction_id") or "") == metadata_tx_id
):
replay_found = True
replay_tx = dict(existing_tx)
break
if not replay_found and not active_tx and not metadata_tx:
# Preserve idempotency for old history-only records, but never let a
# new metadata-less terminal mutate whichever campaign happens to be active.
for existing in list(campaign.get("history") or []):
if (
isinstance(existing, dict)
and str(existing.get("task_id") or "") == str(task_id or "")
):
replay_found = True
replay_tx = dict(existing.get("transaction") or {})
break
if not replay_found:
return {
"accepted": False, "persisted": False, "replay": False,
"reason": "transaction_missing", "transaction": {},
}
if not replay_found:
if (
not _matches(active_tx)
or not _matches(metadata_tx)
or str(active_tx.get("transaction_id") or "")
!= str(metadata_tx.get("transaction_id") or "")
):
return {
"accepted": False, "persisted": False, "replay": False,
"reason": "transaction_mismatch", "transaction": {},
}
tx = {**metadata_tx, **active_tx}
axes = normalize_outcome_axes({"outcome_axes": outcome_axes or {}})
tx["outcome_axes"] = axes
tx["updated_at"] = utc_now_iso()
history = list(campaign.get("history") or [])
cost_available = cost_accounting_status == "available" and cost_usd is not None
row = {
"task_id": str(task_id or ""),
"ts": utc_now_iso(),
# ABI-3 (fix-round-3): the honest cost name — this row reaches
# /api/state through the evolution snapshot. Stored legacy rows
# (cost_usd) keep resolving deprecated-wins at the readers and
# at the /api/state projection boundary.
"accounted_upper_bound_usd": float(cost_usd) if cost_available else None,
"cost_accounting_status": "available" if cost_available else "unavailable",
"outcome_axes": axes,
"rounds": int(rounds or 0),
"transaction": tx,
}
history.append(row)
campaign["history"] = history[-50:]
has_commit = bool(str(tx.get("commit_sha") or "").strip())
if not has_commit:
# A crash between the reviewed commit and its SHA receipt can reach
# task-done before boot recovery: adopt the commit the intent proves,
# so the cycle is classified as commit-bearing instead of ``no_op``.
has_commit = bool(adopt_evolution_commit_intent(campaign, tx))
restart_verified = bool(tx.get("restart_verified"))
has_rescue = bool(str(tx.get("rescue_ref") or "").strip())
if has_commit and restart_verified:
tx["cycle_outcome"] = "absorbed"
campaign["absorbed_cycles_done"] = int(
campaign.get("absorbed_cycles_done") or 0
) + 1
append_unique_transaction(campaign, tx)
campaign.pop("active_transaction", None)
_clear_objective_repeat_count(campaign, tx)
elif has_rescue:
tx["cycle_outcome"] = "abandoned"
tx["abandoned_reason"] = "rescue_ref_present"
tx["cleanup_status"] = "pending"
append_unique_transaction(campaign, tx)
campaign.pop("active_transaction", None)
campaign.pop("post_task_backlog_id", None)
_bump_objective_repeat_count(campaign, tx)
elif not has_commit:
tx["cycle_outcome"] = "no_op"
tx["restart_required"] = False
tx["recovery_hint"] = ""
tx["cleanup_status"] = "pending"
append_unique_transaction(campaign, tx)
campaign.pop("active_transaction", None)
campaign.pop("post_task_backlog_id", None)
_bump_objective_repeat_count(campaign, tx)
else:
tx["cycle_outcome"] = "waiting_for_restart"
tx["recovery_hint"] = tx.get("recovery_hint") or (
"Task ended without a reviewed commit plus restart verification; active "
"transaction retained until restart verifies, repo state is recovered, or it is superseded."
)
tx["restart_required"] = True
if not tx.get("restart_decision"):
tx["restart_decision"] = "supervisor_auto_requested"
campaign["active_transaction"] = tx
if tx.get("cycle_outcome") in {"absorbed", "abandoned"}:
campaign["pending_owner_report"] = dict(tx)
campaign["last_task_id"] = str(task_id or "")
campaign["cycles_done"] = int(campaign.get("cycles_done") or 0) + 1
execution_status = str((axes.get("execution") or {}).get("status") or "unknown")
objective_status = str((axes.get("objective") or {}).get("status") or "not_evaluated")
cost_note = f"${float(cost_usd):.4f}" if cost_available else "unavailable"
campaign["progress_notes"] = (
f"Last cycle {task_id}: execution={execution_status}, objective={objective_status}, "
f"rounds={int(rounds or 0)}, cost={cost_note}."
)
if cost_available:
campaign["budget_spent_usd"] = round(
float(campaign.get("budget_spent_usd") or 0.0) + float(cost_usd), 6,
)
else:
campaign["cost_accounting_status"] = "unavailable"
campaign["updated_at"] = utc_now_iso()
try:
persisted = _write_evolution_campaign(
campaign,
expected_campaign_id=campaign_id,
_state_lock_held=True,
)
except Exception:
log.warning("Failed to persist evolution terminal for %s", task_id, exc_info=True)
return {
"accepted": True, "persisted": False, "replay": False,
"reason": "campaign_write_failed", "transaction": {},
}
if not persisted:
return {
"accepted": True, "persisted": False, "replay": False,
"reason": "campaign_write_refused", "transaction": {},
}
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
if replay_found:
replay_tx = _resume_evolution_terminal_effects(campaign_id, str(task_id or ""), replay_tx)
return {
"accepted": True, "persisted": True, "replay": True,
"reason": "duplicate_terminal", "transaction": replay_tx,
}
tx = _resume_evolution_terminal_effects(campaign_id, str(task_id or ""), tx)
return {
"accepted": True, "persisted": True, "replay": False,
"reason": "", "transaction": tx,
}
def build_evolution_task_text(cycle: int) -> str:
"""Build the next evolution-campaign task prompt."""
from ouroboros.config import get_evolution_persistent_objective
campaign = _read_evolution_campaign()
if campaign.get("status") != "active":
return f"EVOLUTION #{cycle}"
parts = [
f"EVOLUTION CAMPAIGN {campaign.get('id') or 'active'} — CYCLE #{cycle}",
"",
"## Objective",
str(campaign.get("objective") or "Autonomously improve Ouroboros."),
]
steer = get_evolution_persistent_objective()
if steer:
parts.extend([
"",
"## Owner Standing Steer (optional bias — does NOT override the Objective above)",
steer,
])
progress = str(campaign.get("progress_notes") or "").strip()
if progress:
parts.extend(["", "## Progress So Far", progress])
history = list(campaign.get("history") or [])[-3:]
if history:
parts.extend(["", "## Recent Campaign Cycles"])
for row in history:
axes = normalize_outcome_axes(row)
execution_status = str((axes.get("execution") or {}).get("status") or "unknown")
objective_status = str((axes.get("objective") or {}).get("status") or "not_evaluated")
# ABI-3 read tolerance: a stored legacy row's spelling resolves
# deprecated-wins; new rows carry the honest name only.
_, row_amount = honest_cost_pair_amount(
row, "accounted_upper_bound_usd", "cost_usd",
)
row_cost = (
f"${row_amount:.4f}"
if row.get("cost_accounting_status") != "unavailable"
and row_amount is not None else "unavailable"
)
parts.append(
f"- {row.get('task_id')}: execution={execution_status}, objective={objective_status}; "
f"rounds={row.get('rounds', 0)}; cost={row_cost}"
)
# Fix B (C10.2): surface the durable improvement backlog and recent solve-capability
# as optional CONTEXT, never a directive. Ouroboros decides what (if anything) to act
# on — an evolution cycle is NOT obligated to draw from the backlog or repeat past
# patterns. Injecting them is LLM-first steering, not a hardcoded work order.
try:
from ouroboros.evolution_checkpoints import build_solve_capability_digest
from ouroboros.improvement_backlog import format_backlog_digest
from ouroboros.utils import truncate_review_artifact
from supervisor import queue
_digest_root = pathlib.Path(queue.DRIVE_ROOT)
_backlog_digest = format_backlog_digest(_digest_root, limit=8, max_chars=3000)
if _backlog_digest:
parts.extend([
"",
"## Improvement Backlog (context only — NOT a work order)",
"Standing nominations from past cycles. Weigh them if useful, but you are "
"free to pursue the Objective however you judge best; you need not pick "
"from this list.",
"",
_backlog_digest,
])
_capability_digest = truncate_review_artifact(build_solve_capability_digest(_digest_root), 2000)
if _capability_digest:
parts.extend([
"",
"## Recent Solve-Capability (context only)",
_capability_digest,
])
except Exception:
log.debug("evolution task digest injection failed", exc_info=True)
parts.extend([
"",
"## Execution Contract",
"- Work as a normal Ouroboros self-improvement task.",
"- Use standard tests and the normal advisory + triad + scope review flow before committing code.",
"- Land at most ONE reviewed self-modification commit in this cycle. Fold reviewer fixes into that commit before committing; do not churn follow-up commits.",
"- After a reviewed commit lands, call request_restart once and stop. Restart verification is the absorption boundary for the cycle.",
"- An honest no-op is a legitimate outcome when the objective is unsafe, already solved, too broad, or needs owner input; do not commit just to make a cycle non-empty.",
"- If the best next step is memory/identity/backlog rather than code, update those durable artifacts with provenance, but do not treat that as an absorbed self-evolution cycle.",
"- A true absorbed self-evolution cycle requires one reviewed self-modification commit followed by successful restart verification before the next campaign cycle.",
"- The review enforcement mode (advisory vs blocking) is the owner's setting. Do NOT hardcode review findings to always block (or always pass) regardless of that mode: forcing per-finding blocks under an owner-chosen advisory mode is forbidden self-modification (BIBLE P3), not a hardening. If advisory pass-through of a critical finding feels wrong, surface it to the owner — never patch the enforcement gate to override their choice.",
"- If the objective is complete or needs owner input, say so clearly in the final result.",
])
return "\n".join(parts)
def notify_owner_cycle_outcome(campaign: Dict[str, Any], tx: Dict[str, Any]) -> None:
"""WS-13.5 (e5=ux_absorb_report): owner-facing chat note for a finished
self-evolution cycle. Absorbed -> short what/why; abandoned -> honest
warning; no_op / waiting -> quiet (the lifecycle event already records it).
No web/UI edits; chat only, budget-gated. Lazy imports avoid an import
cycle with supervisor.state / supervisor.message_bus.
"""
outcome = str(tx.get("cycle_outcome") or "")
if outcome not in ("absorbed", "abandoned"):
return # no_op / waiting_for_restart: event-only, stay quiet
from supervisor.state import load_state
from supervisor.message_bus import send_with_budget
owner_chat_id = int(load_state().get("owner_chat_id") or 0)
if not owner_chat_id:
return
objective = str(campaign.get("objective") or "").strip()
obj_short = (objective[:160] + "…") if len(objective) > 160 else objective
if outcome == "absorbed":
commit_sha = str(tx.get("commit_sha") or "").strip()[:12]
msg = (
f"🧬 Evolution cycle absorbed (commit {commit_sha}).\n"
f"Objective: {obj_short or 'autonomous self-improvement'}\n"
"The reviewed self-modification is now live (restart verified). Reply if you want it reverted."
)
else:
reason = str(tx.get("abandoned_reason") or "unspecified")
msg = (
f"⚠️ Evolution cycle abandoned (reason: {reason}).\n"
f"Objective: {obj_short or 'autonomous self-improvement'}\n"
"No change was absorbed; the transaction was rolled back/closed to unblock the next cycle."
)
send_with_budget(owner_chat_id, msg)
def clear_pending_owner_report(expected: Dict[str, Any]) -> bool:
"""Clear only the report that was sent, without overwriting newer campaign state."""
from ouroboros.utils import update_json_locked
from supervisor import state
path = _evolution_campaign_path()
state.assert_test_data_path(path)
lock_fd = state.acquire_file_lock(state.STATE_LOCK_PATH)
if lock_fd is None:
return False
cleared = {"ok": False}
def _mutate(campaign: Dict[str, Any]) -> Optional[Dict[str, Any]]:
if campaign.get("pending_owner_report") != expected:
return None
campaign.pop("pending_owner_report", None)
campaign["updated_at"] = utc_now_iso()
cleared["ok"] = True
return campaign
try:
update_json_locked(path, _mutate, strict_existing_dict=True)
return bool(cleared["ok"])
finally:
state.release_file_lock(state.STATE_LOCK_PATH, lock_fd)
def append_unique_transaction(campaign: Dict[str, Any], tx: Dict[str, Any]) -> None:
tx_history = list(campaign.get("transaction_history") or [])
tx_id = str(tx.get("transaction_id") or "")
if tx_id and any(isinstance(item, dict) and str(item.get("transaction_id") or "") == tx_id for item in tx_history):
campaign["transaction_history"] = tx_history[-50:]
return
tx_history.append(dict(tx))
campaign["transaction_history"] = tx_history[-50:]
def write_pending_restart_marker(
drive_root: pathlib.Path, *, expected_sha: str, expected_branch: str, reason: str,
evolution_claim: Optional[Dict[str, Any]] = None,
) -> pathlib.Path:
"""One schema of ``state/pending_restart_verify.json`` for both writers (the
supervisor's evolution restart, the agent's ``restart`` tool); ``verify_restart``
reads it at boot. The claim key exists only for an exact evolution claim — an
empty one would read as a claim mismatch."""
path = pathlib.Path(drive_root) / "state" / "pending_restart_verify.json"
atomic_write_json(path, {
"ts": utc_now_iso(),
"expected_sha": str(expected_sha or "").strip(),
"expected_branch": str(expected_branch or "").strip(),
"reason": str(reason or "").strip(),
**({"evolution_claim": dict(evolution_claim)} if evolution_claim else {}),
}, trailing_newline=True)
return path
def request_evolution_restart(drive_root: pathlib.Path, tx: Dict[str, Any], log: Any = None) -> None:
"""Write the exact restart-verify claim, then ask the server to restart.
``OUROBOROS_EVOLUTION_AUTO_RESTART`` off skips ONLY the restart: the marker is
what lets the next boot — the owner's manual one included — attribute the cycle
by exact claim rather than the weaker markerless reconcile (W4-F3)."""
commit_sha = str(tx.get("commit_sha") or "").strip()
if not commit_sha:
return
claim = {key: str(tx.get(key) or "") for key in ("campaign_id", "transaction_id", "task_id")}
claim["commit_sha"] = commit_sha
authority = check_evolution_authority(**claim)
if not authority.get("ok"):
if log is not None:
log.warning("Automatic evolution restart cancelled: exact authority changed (%s)",
authority.get("reason") or "unknown")
return
try:
existing = read_json_dict(pathlib.Path(drive_root) / "state" / "pending_restart_verify.json") or {}
restart_reason = (
str(existing.get("reason") or "").strip() if existing.get("evolution_claim") == claim else ""
) or "supervisor_auto_evolution_restart"
write_pending_restart_marker(
drive_root, expected_sha=commit_sha, expected_branch=str(tx.get("base_branch") or ""),
reason=restart_reason, evolution_claim=claim,
)
auto_restart = str(os.environ.get("OUROBOROS_EVOLUTION_AUTO_RESTART", "true") or "true").lower()
if auto_restart in {"0", "false", "no", "off"}:
if log is not None:
log.info("Automatic evolution restart is off; the restart-verify marker for %s awaits "
"a manual restart", commit_sha[:12])
return
from supervisor import workers
workers.get_event_q().put({
"type": "restart_request", "reason": restart_reason,
"evolution_restart": True, "ts": utc_now_iso(),
})
except Exception:
if log is not None:
log.debug("Failed to request automatic evolution restart", exc_info=True)