mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
The pre-push panel's third reviewer ran 56 mutations and drove the real
supervisor handler. One defect and several unpinned branches:
- A refusal that writes no result of its own (rejected attachments, and the
duplicate-id exit) left this promote's emitted stub on disk, so the
reconciliation read answered 'admission pending, do not promote the same
work a second time' for ever and the owner's work never started. Both exits
now replace the stub of their OWN routing token; with no stub on disk they
still write nothing. This corrects the claim in 8d57aabe3 that every refused
promote already overwrote the stub.
- The pending sentence states only what the row proves (emitted, no receipt,
the time now) and hands the decision back: an event can be lost with a
restarted supervisor, so a standing ban was a lie waiting to happen.
- The admission gates use a strict form, is_own_admission_stub: no token is
no claim, so a tokenless admission meets any stub as a row that owns the id,
and _promote_duplicate_reason requires its token. This closes the disclosed
tokenless residual structurally.
- An emitted stub is neither the project's last result nor a steerable
predecessor: the lookup skips it like a child, and the door says what it is.
- prompts/SYSTEM.md, the same one sentence: a read-only subagent reads
deliverables, every subagent its parent's and root's task files (the earlier
'never a sibling's' is false on a single-drive geometry).
- The light-mode refusal says 'runtime control data', since it lists the task
drive under the writable roots in the same text.
New pins, each re-run against the reviewer's surviving mutation and now red
without its line: the stub's token and status conditions, the token passed
by the promote handler, the retry window of the compaction floor, the
read-only gate of the lineage redirect, the project filter and omission count
of a room's live roots, the secret-name filter of a Deliverables listing, and
the host-issued authority source of a listed row.
650 lines
30 KiB
Python
650 lines
30 KiB
Python
"""Supervisor task queue, persistence, timeouts, and evolution scheduling."""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import math
|
||
import pathlib
|
||
import queue as _stdqueue # noqa: F401 — re-exported for the test suite's reap-queue isolation
|
||
import threading
|
||
import time # noqa: F401 -- facade name tests read (queue.time)
|
||
import uuid
|
||
from typing import Any, Dict, List, Optional, Tuple
|
||
|
||
from supervisor.state import (
|
||
load_state,
|
||
append_jsonl, # noqa: F401 -- queue_snapshot leaf reads it via the _queue() handle
|
||
atomic_write_text, # noqa: F401 -- queue_snapshot leaf reads it via the _queue() handle
|
||
budget_remaining, EVOLUTION_BUDGET_RESERVE,
|
||
reconstruct_task_cost as reconstruct_task_cost,
|
||
)
|
||
from supervisor.message_bus import (
|
||
coerce_chat_identity, # noqa: F401 -- queue_timeouts leaf reads it via the _queue() handle
|
||
notification_chat_route,
|
||
send_with_budget,
|
||
)
|
||
from ouroboros.config import (
|
||
DATA_DIR,
|
||
FINALIZATION_GRACE_DEFAULT_SEC,
|
||
get_finalization_grace_sec,
|
||
get_per_call_timeout_ceiling_sec, # noqa: F401 -- queue_timeouts leaf reads it via the _queue() handle
|
||
get_task_abs_ceiling_sec, # noqa: F401 -- queue_timeouts leaf reads it via the _queue() handle
|
||
get_task_idle_timeout_sec, # noqa: F401 -- queue_timeouts leaf reads it via the _queue() handle
|
||
)
|
||
from ouroboros.consciousness_authority import apply_consciousness_authority, is_consciousness_origin
|
||
from ouroboros.contracts.task_contract import attach_task_contract, build_task_contract, normalize_allowed_resources # noqa: F401
|
||
from ouroboros.schedule_contract import RESERVED_TEMPLATE_FIELDS, schedule_slug # noqa: F401
|
||
from ouroboros.skill_loader import skill_identity_collision_names # noqa: F401
|
||
from ouroboros.outcomes import terminal_outcome_axes
|
||
from ouroboros.utils import atomic_write_json, read_json_dict, utc_now_iso # noqa: F401
|
||
from supervisor.evolution_lifecycle import ( # noqa: F401 -- public queue API and lazy scheduler dependencies
|
||
_deliver_pending_owner_report,
|
||
_read_evolution_campaign,
|
||
begin_evolution_transaction,
|
||
build_evolution_task_text,
|
||
disable_evolution_authority,
|
||
disable_evolution_projection,
|
||
deliver_pending_owner_report,
|
||
enqueue_evolution_task_if_needed,
|
||
evolution_block_reason,
|
||
notify_owner_cycle_outcome,
|
||
pause_evolution_campaign,
|
||
start_evolution_campaign, # noqa: F401 -- historical queue API re-export
|
||
)
|
||
from supervisor.task_lifecycle import ( # noqa: F401 -- public queue API re-exports
|
||
BUDGET_ROOT_FENCES, apply_budget_root_admission_fence, cancel_task_by_id,
|
||
clear_acceptance_fence_for_root,
|
||
resume_budget_paused_task, restore_queue_fences, transition_acceptance_fence,
|
||
)
|
||
log = logging.getLogger(__name__)
|
||
|
||
|
||
DRIVE_ROOT: pathlib.Path = pathlib.Path(DATA_DIR)
|
||
# The queue snapshot path has ONE authority (MIGRATION row 1030, D18): this
|
||
# module. init() rebinds it per drive root; the queue_snapshot leaf reads it
|
||
# through the _queue() handle.
|
||
QUEUE_SNAPSHOT_PATH: pathlib.Path = DRIVE_ROOT / "state" / "queue_snapshot.json"
|
||
HEARTBEAT_STALE_SEC: int = 120
|
||
QUEUE_MAX_RETRIES: int = 1
|
||
FINALIZATION_GRACE_SEC: int = FINALIZATION_GRACE_DEFAULT_SEC
|
||
SCHEDULED_TASKS_FILE = pathlib.Path("state") / "scheduled_tasks.json"
|
||
# BUG3: pause a campaign whose objective fails to absorb after this many reviewed cycles.
|
||
# Mirrors the consecutive-failures threshold; keyed on the objective fingerprint, not failures.
|
||
OBJECTIVE_REPEAT_CAP: int = 3
|
||
|
||
|
||
def init(drive_root: pathlib.Path) -> None:
|
||
global DRIVE_ROOT, FINALIZATION_GRACE_SEC, QUEUE_SNAPSHOT_PATH
|
||
DRIVE_ROOT = drive_root
|
||
QUEUE_SNAPSHOT_PATH = drive_root / "state" / "queue_snapshot.json"
|
||
FINALIZATION_GRACE_SEC = get_finalization_grace_sec()
|
||
BUDGET_ROOT_FENCES.clear()
|
||
# A previous process's direct-chat turns must not outlive it in the roster,
|
||
# and this clear is the last moment their ids exist: the roster is taken over
|
||
# here and handed to snapshot restore below, which fences them like any other
|
||
# row the stop caught.
|
||
from supervisor.direct_roots import take_direct_roots
|
||
|
||
PRIOR_DIRECT_ROOTS.clear()
|
||
PRIOR_DIRECT_ROOTS.update(take_direct_roots(drive_root))
|
||
|
||
|
||
def refresh_timeouts_from_settings(settings: dict) -> None:
|
||
"""Hot-reload the active liveness settings.
|
||
|
||
The flat wall-clock pair this once also had to absorb (soft/hard) is retired
|
||
in 7.0: load_settings strips the keys, so there is no stored value left to
|
||
accept, warn about, or lie about honoring.
|
||
"""
|
||
global FINALIZATION_GRACE_SEC
|
||
FINALIZATION_GRACE_SEC = get_finalization_grace_sec(settings)
|
||
|
||
|
||
# The previous process's direct-chat roots, taken from `state/direct_roots.json`
|
||
# by init above and consumed once by snapshot restore. A process-local handover
|
||
# of the SAME fragment, never a second store.
|
||
PRIOR_DIRECT_ROOTS: Dict[str, Any] = {}
|
||
|
||
# Set by workers.init_queue_refs().
|
||
PENDING: List[Dict[str, Any]] = []
|
||
RUNNING: Dict[str, Dict[str, Any]] = {}
|
||
QUEUE_SEQ_COUNTER_REF: Dict[str, int] = {"value": 0}
|
||
ACCEPTANCE_FENCES: Dict[str, Dict[str, Any]] = {}
|
||
ADMISSION_RESERVATIONS: Dict[str, str] = {}
|
||
|
||
# Guards PENDING/RUNNING mutations across main loop, direct chat, watchdog.
|
||
_queue_lock = threading.RLock()
|
||
from supervisor.task_admission import ( # noqa: E402,F401 - public queue API
|
||
coerce_queue_order, prefer_terminalization_retry_rows, record_scheduled_admission,
|
||
reject_invalid_task_depth, release_task_admission, restore_invalid_depth_admission,
|
||
restore_terminalization_retry, restore_terminalization_retry_rows,
|
||
reserve_task_admission,
|
||
)
|
||
# Variant A off-loop worker reaper lives in supervisor/task_reaper.py (module size); re-export
|
||
# the thin names the enforce path and tests use — monkeypatching these queue names still works.
|
||
from supervisor.task_reaper import ( # noqa: E402,F401 — re-exported for enforce path + tests
|
||
ensure_reaper_started as _ensure_reaper_started,
|
||
reap_queue as _reap_queue,
|
||
reap_timed_out_task as _reap_timed_out_task,
|
||
request_finalization_grace as _request_finalization_grace,
|
||
resolve_grace_episode_for_spared_task as _resolve_grace_episode_for_spared_task,
|
||
)
|
||
|
||
|
||
def init_queue_refs(pending: List[Dict[str, Any]], running: Dict[str, Dict[str, Any]],
|
||
seq_counter_ref: Dict[str, int]) -> None:
|
||
"""Bind queue structures owned by workers.py."""
|
||
global PENDING, RUNNING, QUEUE_SEQ_COUNTER_REF
|
||
PENDING = pending
|
||
RUNNING = running
|
||
QUEUE_SEQ_COUNTER_REF = seq_counter_ref
|
||
ADMISSION_RESERVATIONS.clear()
|
||
|
||
|
||
def _task_priority(task_type: str) -> int:
|
||
t = str(task_type or "").strip().lower()
|
||
if t in ("skill_publish", "task", "review", "deep_self_review"):
|
||
return 0
|
||
if t == "evolution":
|
||
return 1
|
||
return 2
|
||
|
||
|
||
def _queue_sort_key(task: Dict[str, Any]) -> Tuple[int, int]:
|
||
pr = coerce_queue_order(task.get("priority"), _task_priority(str(task.get("type") or "")))
|
||
seq = coerce_queue_order(task.get("_queue_seq"))
|
||
return pr, seq
|
||
|
||
|
||
def sort_pending() -> None:
|
||
"""Sort pending queue by priority and insertion sequence."""
|
||
PENDING.sort(key=_queue_sort_key)
|
||
|
||
|
||
def drain_all_pending(*, persist: bool = True) -> list:
|
||
"""Drain pending tasks; optionally defer snapshot persistence until custody settles."""
|
||
drained = list(PENDING)
|
||
PENDING.clear()
|
||
if persist:
|
||
persist_queue_snapshot(reason="drain_all_pending")
|
||
return drained
|
||
|
||
|
||
def enqueue_task(
|
||
task: Dict[str, Any], front: bool = False, *, restoring_snapshot: bool = False,
|
||
) -> Dict[str, Any]:
|
||
"""Add task to PENDING (thread-safe: HTTP handlers enqueue concurrently
|
||
with the supervisor main loop, so the mutation must hold the queue lock)."""
|
||
t = dict(task)
|
||
attach_task_contract(apply_consciousness_authority(t))
|
||
# The allowance read takes the cross-process ledger lock: read it BEFORE the queue
|
||
# lock so a contended ledger never stalls every queue reader; only the live-root
|
||
# count and the append must be one transaction with the lock (the window gates
|
||
# starts — a window stale by milliseconds changes nothing).
|
||
consciousness_window = None
|
||
if not restoring_snapshot and _consciousness_root(t):
|
||
from ouroboros.consciousness_allowance import allowance_window
|
||
|
||
consciousness_window = allowance_window(DRIVE_ROOT)
|
||
with _queue_lock:
|
||
require_unique_id = bool(t.pop("_require_unique_task_id", False))
|
||
require_worker_pool = bool(t.pop("_require_worker_pool", False))
|
||
admission_token = str(t.pop("_admission_token", "") or "")
|
||
task_id = str(t.get("id") or "").strip()
|
||
reserved_token = str(ADMISSION_RESERVATIONS.get(task_id) or "")
|
||
if reserved_token and admission_token != reserved_token:
|
||
# A reservation owns this id until its request either enqueues or releases it.
|
||
# Tokenless internal callers and competing ingress must not consume/collide with it.
|
||
t["_admission_blocked"] = "admission_reservation_owned"
|
||
return t
|
||
if require_unique_id and task_id:
|
||
# Exact-id ownership wins over malformed-depth replay.
|
||
live_duplicate = task_id in RUNNING or any(
|
||
isinstance(row, dict) and str(row.get("id") or "") == task_id
|
||
for row in PENDING
|
||
)
|
||
if live_duplicate:
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
t["_admission_blocked"] = "duplicate_task_id"
|
||
return t
|
||
try:
|
||
from ouroboros.routing_wait import is_own_admission_stub
|
||
from ouroboros.task_results import load_task_result
|
||
stored = load_task_result(DRIVE_ROOT, task_id, strict=True)
|
||
# The emitted promote stub (#1160) belongs to THIS admission token:
|
||
# its own enqueue reads around it, any other row still owns the id.
|
||
if stored and not is_own_admission_stub(stored, admission_token):
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
t["_admission_blocked"] = "duplicate_task_id"
|
||
return t
|
||
except Exception:
|
||
log.warning("Fresh task-id lookup failed for %s", task_id, exc_info=True)
|
||
t["_admission_blocked"] = "task_id_lookup_failed"
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
retry = restore_terminalization_retry(t, pending=PENDING, running=RUNNING, queue_seq_counter_ref=QUEUE_SEQ_COUNTER_REF, sort_pending=sort_pending) if restoring_snapshot else None
|
||
if retry:
|
||
return retry
|
||
if reject_invalid_task_depth(t, reservations=ADMISSION_RESERVATIONS, admission_token=admission_token):
|
||
return t
|
||
if require_worker_pool:
|
||
try:
|
||
from supervisor import workers
|
||
|
||
pool_state = workers._worker_pool_execution_state()
|
||
except Exception:
|
||
pool_state = {"available": False, "disabled_reason": "state_unavailable"}
|
||
if not pool_state["available"]:
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
t["_admission_blocked"] = "worker_pool_unavailable"
|
||
t["_worker_pool_disabled_reason"] = pool_state["disabled_reason"]
|
||
return t
|
||
consciousness_block = None if restoring_snapshot else _consciousness_admission_block(t, consciousness_window)
|
||
if consciousness_block is not None:
|
||
t["_admission_blocked"], t["_admission_detail"] = consciousness_block
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
if admission_token and reserved_token != admission_token:
|
||
t["_admission_blocked"] = "admission_reservation_lost"
|
||
return t
|
||
project_id = str(t.get("project_id") or "").strip()
|
||
if project_id:
|
||
try:
|
||
from ouroboros.projects_registry import get_reserved_project
|
||
|
||
project = get_reserved_project(DRIVE_ROOT, project_id)
|
||
lifecycle = str((project or {}).get("lifecycle") or "active")
|
||
if project is not None and lifecycle != "active":
|
||
t["_admission_blocked"] = "project_routing_fence"
|
||
t["_project_lifecycle"] = lifecycle
|
||
t["_project_id"] = project_id
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
except Exception:
|
||
log.warning("Project admission check failed for %s", project_id, exc_info=True)
|
||
t["_admission_blocked"] = "project_routing_fence_lookup_failed"
|
||
t["_project_id"] = project_id
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
root_id = str(t.get("root_task_id") or "").strip()
|
||
if root_id and not restoring_snapshot and apply_budget_root_admission_fence(t, root_id):
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
fence = ACCEPTANCE_FENCES.get(root_id) if root_id else None
|
||
if isinstance(fence, dict) and str(fence.get("status") or "") in {"active", "sealed"}:
|
||
t["_admission_blocked"] = "task_acceptance_fence"
|
||
t["_acceptance_fence_token"] = str(fence.get("token") or "")
|
||
t["_acceptance_fence_status"] = str(fence.get("status") or "active")
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
QUEUE_SEQ_COUNTER_REF["value"] += 1
|
||
seq = QUEUE_SEQ_COUNTER_REF["value"]
|
||
t["priority"] = coerce_queue_order(t.get("priority"), _task_priority(str(t.get("type") or "")))
|
||
_att = t.get("_attempt")
|
||
t.setdefault("_attempt", int(_att) if _att is not None else 1)
|
||
t["_queue_seq"] = -seq if front else seq
|
||
t["queued_at"] = utc_now_iso()
|
||
if admission_token:
|
||
t["_admission_owner_token"] = admission_token
|
||
PENDING.append(t)
|
||
sort_pending()
|
||
if ADMISSION_RESERVATIONS.get(task_id) == admission_token:
|
||
ADMISSION_RESERVATIONS.pop(task_id, None)
|
||
return t
|
||
|
||
|
||
def live_consciousness_root_count() -> int:
|
||
"""Live PENDING+RUNNING roots that consciousness started (its origin marker on the
|
||
task metadata; subagents are their root's business). Sibling of ``queue_has_task_type``."""
|
||
def _counts(task: Any) -> bool:
|
||
return (
|
||
isinstance(task, dict)
|
||
and str(task.get("delegation_role") or "root") == "root"
|
||
and is_consciousness_origin(task.get("metadata"))
|
||
)
|
||
|
||
live = sum(1 for task in PENDING if _counts(task))
|
||
return live + sum(
|
||
1 for meta in RUNNING.values() if isinstance(meta, dict) and _counts(meta.get("task"))
|
||
)
|
||
|
||
|
||
def _consciousness_root(task: Dict[str, Any]) -> bool:
|
||
return (is_consciousness_origin(task.get("metadata"))
|
||
and str(task.get("delegation_role") or "root") == "root")
|
||
|
||
|
||
def _consciousness_admission_block(task: Dict[str, Any], window: Optional[Dict[str, Any]] = None) -> Optional[Tuple[str, str]]:
|
||
"""The ONE admission door for the roots consciousness starts (owner decisions В11/В18).
|
||
|
||
Called under the queue lock, so the live count and the admission are one
|
||
transaction. Returns ``(reason, detail)`` in the queue's existing refusal
|
||
vocabulary when a consciousness-origin ROOT may not start — the concurrency
|
||
cap over live roots with the same marker (``OUROBOROS_CONSCIOUSNESS_MAX_TASKS``,
|
||
0 = never) or the rolling-24h allowance (``OUROBOROS_CONSCIOUSNESS_DAILY_USD``,
|
||
0 = consciousness may not spend; an unreadable ledger refuses honestly as
|
||
``allowance_unknown``) — and ``None`` when it may. Subagents are bounded by
|
||
their root's own cap and the per-root child cap, never counted twice; a
|
||
snapshot restore re-admits already-admitted work and is not gated here. ``window``
|
||
is the allowance the caller read before taking the lock; read here only when absent.
|
||
"""
|
||
if not _consciousness_root(task):
|
||
return None
|
||
from ouroboros.config import get_consciousness_max_tasks
|
||
from ouroboros.consciousness_allowance import (
|
||
STATUS_AVAILABLE, STATUS_UNKNOWN, allowance_window,
|
||
)
|
||
|
||
max_tasks = get_consciousness_max_tasks()
|
||
live = live_consciousness_root_count()
|
||
if live >= max_tasks:
|
||
return ("consciousness_task_limit", (
|
||
f"{live} of {max_tasks} consciousness-started tasks already live"
|
||
if max_tasks else "OUROBOROS_CONSCIOUSNESS_MAX_TASKS=0: consciousness never starts tasks"
|
||
))
|
||
if window is None:
|
||
window = allowance_window(DRIVE_ROOT)
|
||
if window["status"] == STATUS_UNKNOWN:
|
||
return ("consciousness_allowance_unknown",
|
||
f"the usage ledger could not be read: {window.get('error') or 'unknown error'}")
|
||
if window["status"] != STATUS_AVAILABLE:
|
||
if not window["limit_usd"]:
|
||
return ("consciousness_allowance_exhausted",
|
||
"OUROBOROS_CONSCIOUSNESS_DAILY_USD=0: consciousness may not spend")
|
||
at_least = " (at least)" if window["unknown_unmetered"] else ""
|
||
return ("consciousness_allowance_exhausted", (
|
||
f"${window['accounted_usd']:.2f}{at_least} of ${window['limit_usd']:.2f} "
|
||
f"spent in the last 24 h; resets at {window['resets_at'] or 'unknown'}"
|
||
))
|
||
return None
|
||
|
||
|
||
def queue_has_task_type(task_type: str) -> bool:
|
||
"""Return whether this task type is pending or running."""
|
||
tt = str(task_type or "")
|
||
if any(str(t.get("type") or "") == tt for t in PENDING):
|
||
return True
|
||
for meta in RUNNING.values():
|
||
task = meta.get("task") if isinstance(meta, dict) else None
|
||
if isinstance(task, dict) and str(task.get("type") or "") == tt:
|
||
return True
|
||
return False
|
||
|
||
|
||
# Cron/timezone schedule helpers live in supervisor/schedule_time.py (P7
|
||
# module-size relief); imported under their historical private names.
|
||
from supervisor.schedule_time import ( # noqa: E402
|
||
next_cron_time as _next_cron_time, # noqa: F401
|
||
once_due as _once_due, # noqa: F401
|
||
parse_schedule_time as _parse_schedule_time, # noqa: F401
|
||
prune_consumed_once_records as _prune_consumed_once, record_last_error as _record_last_error, # noqa: F401
|
||
schedule_next_run as _schedule_next_run, # noqa: F401
|
||
timezone_for_schedule as _timezone_for_schedule, # noqa: F401
|
||
)
|
||
|
||
|
||
def _emit_cancel_task_done(
|
||
task: Optional[Dict[str, Any]],
|
||
task_id: str,
|
||
*,
|
||
cost_fields: Optional[Dict[str, Any]] = None,
|
||
status: str = "cancelled",
|
||
) -> None:
|
||
"""Emit a task_done event after a cancel so the UI live card resolves.
|
||
Covers both the agent-tool path (_handle_cancel_task) and the HTTP path.
|
||
``status`` carries the STORED terminal truth: when a worker wrote its own
|
||
natural result just before the kill, the card must resolve to THAT outcome
|
||
rather than be left unresolved until a reload.
|
||
``cost_fields`` is the caller's accounting authority — a reconstructed
|
||
ledger projection or a CONFIRMED pre-start zero. An absent projection emits
|
||
an honest nullable unknown; the old default fabricated a final $0 for every
|
||
cancel (Poltergeist A1.10, owner 10=B)."""
|
||
try:
|
||
from supervisor import workers
|
||
chat_id = int((task or {}).get("chat_id") or 0) if isinstance(task, dict) else 0
|
||
workers.get_event_q().put({
|
||
"type": "task_done",
|
||
"task_id": str(task_id),
|
||
# The tree identity survives even though the row already left
|
||
# PENDING/RUNNING: the fence-release seam resolves the root
|
||
# from the event when the queue no longer holds the task.
|
||
"root_task_id": str((task or {}).get("root_task_id") or "") if isinstance(task, dict) else "",
|
||
"task_type": str((task or {}).get("type") or ""),
|
||
"chat_id": chat_id,
|
||
"status": status,
|
||
"outcome_axes": terminal_outcome_axes(
|
||
lifecycle=status, execution=status, reason_code=status,
|
||
review_trigger="supervisor_terminal",
|
||
),
|
||
**(cost_fields or {
|
||
"cost_accounting_status": "unavailable", "cost_final": False,
|
||
# ABI-3: honest name only — the retired alias is read-only.
|
||
"accounted_upper_bound_usd": None,
|
||
}),
|
||
"metadata": (task or {}).get("metadata") if isinstance((task or {}).get("metadata"), dict) else {},
|
||
})
|
||
except Exception:
|
||
log.debug("Failed to emit task_done for cancelled task %s", task_id, exc_info=True)
|
||
|
||
|
||
# Cancellation custody and the terminal-cancel result-field builder live in
|
||
# supervisor.task_lifecycle (module-size boundary); re-exported so
|
||
# `supervisor.queue` stays the single import surface for callers.
|
||
from supervisor.task_lifecycle import ( # noqa: E402, F401 -- intentional public re-exports
|
||
CANCEL_ALREADY_SETTLED,
|
||
CANCEL_CANCELLED,
|
||
CANCEL_FAILED,
|
||
CANCEL_NOT_FOUND,
|
||
_CANCEL_TERMINALIZED,
|
||
_cancel_result_fields,
|
||
cancel_task_custody,
|
||
drive_cancel_intent_scope,
|
||
task_has_live_ownership,
|
||
task_subtree_is_live,
|
||
)
|
||
from supervisor.queue_transitions import ( # noqa: E402, F401 -- intentional public re-exports
|
||
evolution_stop_report,
|
||
stop_evolution_tasks,
|
||
sweep_orphaned_budget_fences,
|
||
)
|
||
|
||
|
||
def _cancel_task_by_id_single(task_id: str) -> bool:
|
||
"""Boolean facade for the pre-v6.82 single-task callers."""
|
||
return cancel_task_custody(task_id) in {CANCEL_CANCELLED, CANCEL_ALREADY_SETTLED}
|
||
|
||
|
||
# Evolution-stop transitions (GR2-13) live in supervisor.queue_transitions
|
||
# (module-size boundary); re-exported below with the other transition helpers
|
||
# so `supervisor.queue` stays the single import surface for callers.
|
||
|
||
|
||
def queue_deep_self_review_task(reason: str, model: str = "", force: bool = False, chat_id: Optional[int] = None,
|
||
origin: Optional[Dict[str, Any]] = None) -> Optional[str]:
|
||
"""Queue a deep self-review task.
|
||
|
||
``chat_id`` targets a specific chat (e.g. the external transport chat that ran
|
||
``/review``) so the queued ack and the task results return to the requester
|
||
instead of always defaulting to the web owner's ``owner_chat_id``. ``origin`` is
|
||
the requester's consciousness origin (a wake-up or its tree), stamped on the
|
||
root so the admission door and the ledger see it; empty for the owner.
|
||
"""
|
||
# Membership, not truthiness: a review asked for from the hidden partition
|
||
# is answered there, not silently re-routed to the owner's main chat.
|
||
target_chat_id = notification_chat_route(chat_id, load_state().get("owner_chat_id"))
|
||
if target_chat_id is None:
|
||
return None
|
||
if (not force) and queue_has_task_type("deep_self_review"):
|
||
return None
|
||
tid = uuid.uuid4().hex[:8]
|
||
admitted = enqueue_task({
|
||
"id": tid,
|
||
"type": "deep_self_review",
|
||
"chat_id": int(target_chat_id),
|
||
"text": reason or "Deep self-review",
|
||
"model": model,
|
||
"_require_worker_pool": True,
|
||
**({"metadata": dict(origin)} if origin else {}),
|
||
})
|
||
if admitted.get("_admission_blocked"):
|
||
reason = admitted.get("_worker_pool_disabled_reason") or admitted["_admission_blocked"]
|
||
detail = str(admitted.get("_admission_detail") or "")
|
||
hint = f" {detail}." if detail else " Use /restart to restore the worker pool."
|
||
send_with_budget(
|
||
int(target_chat_id),
|
||
f"Deep self-review could not be queued: {reason}.{hint}",
|
||
role="system", system_type="deep_self_review_unavailable",
|
||
)
|
||
return None
|
||
persist_queue_snapshot(reason="deep_self_review_enqueued")
|
||
# Typed SYSTEM row: an acknowledgement is never a task's answer, and the bench
|
||
# trajectory reader takes the last UNTYPED outbound row as one.
|
||
send_with_budget(int(target_chat_id), f"🔎 Deep self-review queued: {tid} ({reason})", role="system", system_type="deep_self_review_queued")
|
||
return tid
|
||
|
||
|
||
def get_evolution_status_snapshot(*, budget_projection: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||
"""Return a non-mutating evolution scheduling snapshot.
|
||
|
||
``budget_projection``: optional pre-computed global usage projection from a
|
||
caller that already replayed the ledger this request (``/api/state``), so the
|
||
snapshot does not replay it again. Default ``None`` keeps the self-computing,
|
||
strict fail-closed behavior — a caller whose own computation FAILED must pass
|
||
nothing, so the paused-evolution disclosure still comes from this snapshot.
|
||
"""
|
||
st = load_state()
|
||
enabled = bool(st.get("evolution_mode_enabled"))
|
||
owner_chat_id = int(st.get("owner_chat_id") or 0)
|
||
consecutive_failures = int(st.get("evolution_consecutive_failures") or 0)
|
||
try:
|
||
remaining: Optional[float] = round(float(budget_remaining(st, strict=True, projection=budget_projection)), 2)
|
||
accounting_available = True
|
||
except Exception:
|
||
remaining = None
|
||
accounting_available = False
|
||
queued_task = next((t for t in PENDING if str(t.get("type") or "") == "evolution"), None)
|
||
running_task = next(
|
||
(
|
||
(meta.get("task") if isinstance(meta, dict) else None)
|
||
for meta in RUNNING.values()
|
||
if isinstance(meta, dict)
|
||
and isinstance(meta.get("task"), dict)
|
||
and str(meta["task"].get("type") or "") == "evolution"
|
||
),
|
||
None,
|
||
)
|
||
status = "disabled"
|
||
detail = "Evolution mode is off."
|
||
|
||
campaign = _read_evolution_campaign()
|
||
active_tx = campaign.get("active_transaction") if isinstance(campaign.get("active_transaction"), dict) else {}
|
||
restart_blocked = bool(
|
||
active_tx
|
||
and str(active_tx.get("commit_sha") or "").strip()
|
||
and (bool(active_tx.get("restart_required")) or not bool(active_tx.get("restart_verified")))
|
||
)
|
||
|
||
if restart_blocked:
|
||
status = "waiting_for_restart_verify"
|
||
detail = "Waiting for restart verification before the next absorbed evolution cycle."
|
||
elif isinstance(running_task, dict):
|
||
status = "running"
|
||
detail = "Evolution task is running now."
|
||
elif isinstance(queued_task, dict):
|
||
status = "queued"
|
||
detail = "Evolution task is queued and waiting for a worker."
|
||
elif not accounting_available:
|
||
status = "accounting_unavailable"
|
||
detail = "Cost accounting is unavailable; evolution dispatch is paused without changing the campaign."
|
||
elif consecutive_failures >= 3:
|
||
status = "paused_failures"
|
||
detail = (
|
||
f"Paused after {consecutive_failures} consecutive failures. "
|
||
"Use Evolve again after investigating the failure."
|
||
)
|
||
elif enabled and not owner_chat_id:
|
||
status = "waiting_for_owner_chat"
|
||
detail = "Waiting for the first owner chat binding before scheduling evolution."
|
||
elif enabled and remaining is not None and remaining < EVOLUTION_BUDGET_RESERVE:
|
||
status = "budget_blocked"
|
||
detail = (
|
||
f"Budget reserve active: ${remaining:.2f} remaining, "
|
||
f"${EVOLUTION_BUDGET_RESERVE:.0f} reserved for conversations."
|
||
)
|
||
elif enabled and (PENDING or RUNNING):
|
||
status = "waiting_for_idle"
|
||
detail = "Waiting for active tasks to finish before the next evolution cycle."
|
||
elif enabled:
|
||
status = "idle_ready"
|
||
detail = "Idle and ready to queue the next evolution cycle."
|
||
elif remaining is not None and remaining < EVOLUTION_BUDGET_RESERVE and str(st.get("last_evolution_task_at") or "").strip():
|
||
status = "budget_stopped"
|
||
detail = (
|
||
f"Evolution auto-stopped because only ${remaining:.2f} remains, "
|
||
f"below the ${EVOLUTION_BUDGET_RESERVE:.0f} conversation reserve."
|
||
)
|
||
|
||
return {
|
||
"enabled": enabled,
|
||
"status": status,
|
||
"detail": detail,
|
||
"campaign": campaign,
|
||
"cycle": int(st.get("evolution_cycle") or 0),
|
||
"owner_chat_bound": bool(owner_chat_id),
|
||
"last_task_at": str(st.get("last_evolution_task_at") or ""),
|
||
"consecutive_failures": consecutive_failures,
|
||
"cost_accounting_status": "available" if accounting_available else "unavailable",
|
||
# Unbounded budget (supervisor not initialized / TOTAL_BUDGET<=0)
|
||
# is float('inf'), which strict JSON cannot carry — surface None so
|
||
# /api/state stays serializable on onboarding installs.
|
||
"budget_remaining_usd": remaining if remaining is not None and math.isfinite(remaining) else None,
|
||
"budget_reserve_usd": float(EVOLUTION_BUDGET_RESERVE),
|
||
"pending_count": len(PENDING),
|
||
"running_count": len(RUNNING),
|
||
"queued_task_id": str((queued_task or {}).get("id") or ""),
|
||
"running_task_id": str((running_task or {}).get("id") or ""),
|
||
}
|
||
|
||
|
||
# v7next F1 (D08): moved spans live in their owner leaves; re-exported here
|
||
# so this facade stays the single import surface for callers and tests.
|
||
from supervisor.queue_schedules import ( # noqa: E402, F401 -- intentional public re-exports
|
||
_SKILL_SCHEDULE_SYNC_INTERVAL_SEC,
|
||
_last_skill_schedule_sync,
|
||
_schedule_running_or_queued,
|
||
_scheduled_tasks_path,
|
||
_task_from_schedule,
|
||
_write_scheduled_tasks,
|
||
check_scheduled_tasks,
|
||
list_scheduled_tasks,
|
||
remove_scheduled_task,
|
||
resync_skill_schedules,
|
||
sync_skill_schedules,
|
||
upsert_scheduled_task,
|
||
)
|
||
from supervisor.queue_snapshot import ( # noqa: E402, F401 -- intentional public re-exports
|
||
_kept_service_pids,
|
||
_retained_daemon_pids,
|
||
parse_iso_to_ts,
|
||
persist_queue_snapshot,
|
||
restore_pending_from_snapshot,
|
||
)
|
||
from supervisor.queue_timeouts import ( # noqa: E402, F401 -- intentional public re-exports
|
||
_enforce_task_timeouts_locked,
|
||
_has_live_descendant,
|
||
_has_pending_descendant,
|
||
_is_descendant_of,
|
||
_subtree_progressing,
|
||
_task_deadline_ts,
|
||
_task_drive_for_task,
|
||
enforce_task_timeouts,
|
||
)
|