mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
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>
1498 lines
65 KiB
Python
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)
|