mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Allow an acknowledged owner-wait continuation to reach its existing budget controls while keeping root admission and physical reservations fenced. Remove required waiting copy once a quiz settles while preserving optional assumption history. Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
374 lines
19 KiB
Python
374 lines
19 KiB
Python
"""Handing a pending task to a free worker, and refusing the ones that must not run.
|
|
|
|
Admission order is the queue's; what this adds is the per-task gates: a cancelled
|
|
pending row is settled rather than dispatched, an evolution task without live
|
|
campaign authority is cancelled rather than started, and a repo-writing task waits
|
|
while the writer gate is closed.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import pathlib
|
|
import time
|
|
from typing import Any, Dict
|
|
|
|
from supervisor.queue import _queue_lock
|
|
|
|
|
|
def _pool():
|
|
"""The parent module, read at call time.
|
|
|
|
The pool owns the repo/drive roots, its size, the worker table, the shared
|
|
PENDING/RUNNING refs and the crash clock, and ``init`` REBINDS them.
|
|
Reading them through the module is what keeps one binding: a from-import
|
|
here would freeze the value this module saw at import time (the
|
|
owner-approved D18/D33 mechanical exception).
|
|
"""
|
|
from supervisor import workers
|
|
|
|
return workers
|
|
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
|
|
def _evolution_assignment_error(task: Dict[str, Any]) -> str:
|
|
"""Return the exact authority error for an evolution task about to run."""
|
|
if str(task.get("type") or "") != "evolution":
|
|
return ""
|
|
metadata = task.get("metadata") if isinstance(task.get("metadata"), dict) else {}
|
|
tx = metadata.get("evolution_transaction")
|
|
tx = tx if isinstance(tx, dict) else {}
|
|
task_id = str(task.get("id") or "")
|
|
if str(tx.get("task_id") or "") != task_id:
|
|
return "task_mismatch"
|
|
from supervisor.evolution_lifecycle import check_evolution_authority
|
|
|
|
try:
|
|
authority = check_evolution_authority(
|
|
campaign_id=str(tx.get("campaign_id") or ""),
|
|
transaction_id=str(tx.get("transaction_id") or ""),
|
|
task_id=task_id,
|
|
require_uncommitted=True,
|
|
)
|
|
except Exception:
|
|
log.warning("Evolution assignment authority check failed", exc_info=True)
|
|
return "authority_check_failed"
|
|
return "" if authority.get("ok") else str(authority.get("reason") or "unknown")
|
|
|
|
|
|
def _cancel_unauthorized_evolution(task: Dict[str, Any], reason: str) -> bool:
|
|
"""Terminally cancel a stale restored/retried evolution task."""
|
|
task_id = str(task.get("id") or "")
|
|
from ouroboros.task_results import STATUS_CANCELLED, write_task_result
|
|
|
|
try:
|
|
write_task_result(
|
|
_pool().DRIVE_ROOT,
|
|
task_id,
|
|
STATUS_CANCELLED,
|
|
reason_code="evolution_authority_missing",
|
|
authority_reason=str(reason or "unknown"),
|
|
metadata=task.get("metadata") if isinstance(task.get("metadata"), dict) else {},
|
|
result=f"Evolution authority is no longer active ({reason or 'unknown'}).",
|
|
)
|
|
except Exception:
|
|
log.debug("Failed to cancel unauthorized evolution task %s", task_id, exc_info=True)
|
|
return False
|
|
_pool()._emit_task_done_terminal(
|
|
task, task_id, "cancelled", reason_code="evolution_authority_missing",
|
|
)
|
|
_pool().append_jsonl(
|
|
_pool().DRIVE_ROOT / "logs" / "events.jsonl",
|
|
{
|
|
"ts": _pool().utc_now_iso(), "type": "evolution_assignment_rejected",
|
|
"task_id": task_id, "reason": str(reason or "unknown"),
|
|
},
|
|
)
|
|
return True
|
|
|
|
|
|
def assign_tasks() -> None:
|
|
from supervisor import queue
|
|
from supervisor.state import budget_remaining, EVOLUTION_BUDGET_RESERVE
|
|
from supervisor.worker_owner_wait import maintain_owner_wait_capacity
|
|
|
|
maintain_owner_wait_capacity()
|
|
with _queue_lock:
|
|
st = _pool().load_state()
|
|
# Cancellation/terminal custody wins before validating rows left in the
|
|
# queue. Then quarantine every malformed depth before budget, lease, or
|
|
# capacity filters can leave it waiting indefinitely.
|
|
if not _pool()._drop_cancelled_pending():
|
|
log.error(
|
|
"Task assignment blocked: cancellation authority or custody "
|
|
"state is indeterminate",
|
|
)
|
|
queue.persist_queue_snapshot(reason="cancellation_authority_indeterminate")
|
|
return
|
|
_pool()._retry_terminalization_pending_for_assignment(queue)
|
|
invalid_ids, unresolved_invalid_ids = _pool()._quarantine_invalid_pending_depths()
|
|
unresolved_invalid_id_set = set(unresolved_invalid_ids)
|
|
|
|
if invalid_ids:
|
|
queue.persist_queue_snapshot(reason="invalid_task_depth")
|
|
if unresolved_invalid_ids:
|
|
log.error(
|
|
"Invalid-depth rows deferred until terminal custody is available; continuing assignment for other tasks: %s",
|
|
", ".join(unresolved_invalid_ids),
|
|
)
|
|
try:
|
|
remaining = budget_remaining(st, strict=True)
|
|
except Exception:
|
|
log.error("Task assignment blocked: monetary authority unavailable")
|
|
return
|
|
if remaining <= 0:
|
|
planned = []
|
|
for task in _pool().PENDING:
|
|
if _pool()._invalid_depth_deferred(task, unresolved_invalid_id_set):
|
|
continue
|
|
if isinstance(task.get("_budget_pause"), dict):
|
|
continue
|
|
if task.get("_owner_wait_resume"):
|
|
continue # Restore the checkpoint; the loop still owns its budget stop.
|
|
task_id = str(task.get("id") or "")
|
|
cost_fields = _pool().reconstruct_task_cost(
|
|
task_id, fields=True,
|
|
drive_root=pathlib.Path(task.get("budget_drive_root") or _pool().DRIVE_ROOT),
|
|
)
|
|
if cost_fields.get("cost_accounting_status") != "available":
|
|
log.error("Budget pause blocked: task attempt history unavailable for %s", task_id)
|
|
return
|
|
retry_lineage = bool(
|
|
int(task.get("_attempt") or 1) > 1
|
|
or task.get("original_task_id") or task.get("timeout_retry_from")
|
|
)
|
|
replay_safe = (
|
|
int(cost_fields.get("total_rounds") or 0) == 0
|
|
and not bool(cost_fields.get("ledger_integrity_degraded"))
|
|
and not retry_lineage
|
|
)
|
|
pause = {
|
|
"status": "paused_before_dispatch" if replay_safe else "resource_limited",
|
|
"scope": "global",
|
|
"physical_calls": int(cost_fields.get("total_rounds") or 0),
|
|
"replay_safe": replay_safe,
|
|
"auto_resume": False,
|
|
"resume_policy": "manual_same_generation" if replay_safe else "cancel_or_new_run",
|
|
"paused_at": _pool().utc_now_iso(),
|
|
}
|
|
planned.append((task, pause, cost_fields))
|
|
newly_paused, terminal_ids = [], []
|
|
for task, pause, cost_fields in planned:
|
|
task_id = str(task.get("id") or "")
|
|
result_root = pathlib.Path(task.get("budget_drive_root") or _pool().DRIVE_ROOT)
|
|
try:
|
|
from ouroboros.task_results import STATUS_FAILED, STATUS_SCHEDULED, write_task_result
|
|
|
|
if pause["replay_safe"]:
|
|
task["_budget_pause"] = pause
|
|
newly_paused.append(task_id)
|
|
write_task_result(
|
|
result_root, task_id, STATUS_SCHEDULED,
|
|
reason_code="budget_exhausted", resource_limit=pause,
|
|
)
|
|
else:
|
|
write_task_result(
|
|
result_root, task_id, STATUS_FAILED,
|
|
reason_code="budget_exhausted", resource_limit=pause,
|
|
result="Budget exhausted after prior dispatch; cancel or start a new run.",
|
|
**cost_fields,
|
|
)
|
|
_pool()._emit_task_done_terminal(
|
|
task, task_id, "failed", reason_code="budget_exhausted",
|
|
cost_fields=cost_fields,
|
|
)
|
|
terminal_ids.append(task_id)
|
|
except Exception:
|
|
log.error("Failed to project budget stop for %s", task_id, exc_info=True)
|
|
if terminal_ids:
|
|
terminal = set(terminal_ids)
|
|
_pool().PENDING[:] = [task for task in _pool().PENDING if str(task.get("id") or "") not in terminal]
|
|
if newly_paused or terminal_ids:
|
|
_pool().append_jsonl(
|
|
_pool().DRIVE_ROOT / "logs" / "events.jsonl",
|
|
{
|
|
"ts": _pool().utc_now_iso(),
|
|
"type": "budget_tasks_paused",
|
|
"scope": "global",
|
|
"task_ids": newly_paused,
|
|
"resource_limited_task_ids": terminal_ids,
|
|
"auto_resume": False,
|
|
},
|
|
)
|
|
if st.get("owner_chat_id"):
|
|
_pool().send_with_budget(
|
|
int(st["owner_chat_id"]),
|
|
"🚫 Model budget reached. Queued tasks are paused before dispatch; "
|
|
"raising the limit does not resume them automatically.",
|
|
)
|
|
queue.persist_queue_snapshot(reason="budget_paused_before_dispatch")
|
|
if not any(task.get("_owner_wait_resume") for task in _pool().PENDING):
|
|
return
|
|
|
|
# Evolution is hard-blocked in light runtime mode at the assignment
|
|
# chokepoint too: a task restored from a snapshot or created before the
|
|
# mode switch must never actually run. Cancel them terminally.
|
|
from supervisor.evolution_lifecycle import evolution_block_reason
|
|
evo_block = evolution_block_reason()
|
|
blocked_ids = _pool()._drop_assignable_evolution_tasks(unresolved_invalid_id_set) if evo_block else []
|
|
if blocked_ids:
|
|
from ouroboros.task_results import STATUS_CANCELLED, write_task_result
|
|
for tid in blocked_ids:
|
|
try:
|
|
write_task_result(
|
|
_pool().DRIVE_ROOT, tid, STATUS_CANCELLED,
|
|
result="Evolution is disabled in light runtime mode.",
|
|
)
|
|
except Exception:
|
|
log.debug("Failed to cancel light-mode evolution task %s", tid, exc_info=True)
|
|
if st.get("owner_chat_id"):
|
|
_pool().send_with_budget(int(st["owner_chat_id"]), evo_block)
|
|
queue.persist_queue_snapshot(reason="evolution_blocked_light")
|
|
|
|
from ouroboros.project_lease import candidate_is_leasable, running_project_ids
|
|
from ouroboros.config import get_max_active_subagents_per_root
|
|
|
|
|
|
for w in _pool().WORKERS.values():
|
|
if (w.busy_task_id is None and not getattr(w, "reaping", False)
|
|
and getattr(w, "active_capacity", True) and _pool().PENDING):
|
|
# One-writer-per-project lease: recompute per assignment so a
|
|
# task assigned in THIS loop pass immediately occupies its lane.
|
|
leased = running_project_ids(_pool().RUNNING.values())
|
|
# Find first suitable task (skip over-budget evolution tasks
|
|
# and project-leased candidates)
|
|
chosen_idx = None
|
|
for i, candidate in enumerate(_pool().PENDING):
|
|
if remaining <= 0 and not candidate.get("_owner_wait_resume"):
|
|
continue
|
|
if _pool()._invalid_depth_deferred(candidate, unresolved_invalid_id_set):
|
|
continue
|
|
if not _pool().repo_writer_task_allowed(candidate):
|
|
continue
|
|
if isinstance(candidate.get("_budget_pause"), dict):
|
|
continue
|
|
root_task_id = str(candidate.get("root_task_id") or "").strip()
|
|
if root_task_id in queue.BUDGET_ROOT_FENCES and not candidate.get("_owner_wait_resume"):
|
|
continue
|
|
if str(candidate.get("type") or "") == "evolution" and remaining < EVOLUTION_BUDGET_RESERVE:
|
|
continue
|
|
if not candidate_is_leasable(candidate, leased):
|
|
continue
|
|
if str(candidate.get("delegation_role") or "") == "subagent":
|
|
root_task_id = str(candidate.get("root_task_id") or "")
|
|
if (
|
|
_pool()._running_subagent_count(root_task_id) >= get_max_active_subagents_per_root()
|
|
and not _pool()._assignment_depth_reservation_admits(candidate)
|
|
):
|
|
continue
|
|
chosen_idx = i
|
|
break
|
|
if chosen_idx is None:
|
|
# Project-leased rows wait; over-budget evolution rows are cleaned.
|
|
if remaining < EVOLUTION_BUDGET_RESERVE:
|
|
dropped_ids = _pool()._drop_assignable_evolution_tasks(unresolved_invalid_id_set)
|
|
if dropped_ids:
|
|
queue.persist_queue_snapshot(reason="evolution_dropped_budget")
|
|
continue
|
|
task = _pool().PENDING.pop(chosen_idx)
|
|
depth_error = _pool()._normalize_pending_task_depth(task)
|
|
if depth_error:
|
|
if _pool()._terminalize_invalid_pending_depth(task, depth_error):
|
|
queue.persist_queue_snapshot(reason="invalid_task_depth")
|
|
continue
|
|
# Keep failed terminalization in queue custody for retry.
|
|
_pool().PENDING.insert(chosen_idx, task)
|
|
log.error(
|
|
"Assignment blocked: invalid task depth could not be terminalized for %s",
|
|
task.get("id"),
|
|
)
|
|
break
|
|
evolution_error = _pool()._evolution_assignment_error(task)
|
|
if evolution_error:
|
|
if _pool()._cancel_unauthorized_evolution(task, evolution_error):
|
|
queue.persist_queue_snapshot(reason="evolution_authority_rejected")
|
|
else:
|
|
_pool().PENDING.insert(chosen_idx, task)
|
|
continue
|
|
if str(task.get("delegation_role") or "") == "subagent" and str(task.get("drive_root") or ""):
|
|
try:
|
|
from ouroboros.task_results import STATUS_RUNNING, write_task_result
|
|
from ouroboros.tools.control_delegation import stamp_task_assignment_depth
|
|
from ouroboros.config import get_max_subagent_depth
|
|
|
|
# Assignment is the first host-visible execution fact. Stamp
|
|
# the worker payload and canonical result from one projection.
|
|
_depth_fields = stamp_task_assignment_depth(
|
|
task, max_depth=get_max_subagent_depth(),
|
|
)
|
|
write_task_result(
|
|
_pool().DRIVE_ROOT,
|
|
str(task.get("id") or ""),
|
|
STATUS_RUNNING,
|
|
parent_task_id=task.get("parent_task_id"),
|
|
root_task_id=task.get("root_task_id"),
|
|
session_id=task.get("session_id"),
|
|
actor_id=task.get("actor_id"),
|
|
delegation_role=task.get("delegation_role"),
|
|
project_id=task.get("project_id"),
|
|
role=task.get("role"),
|
|
description=task.get("description"),
|
|
objective=task.get("objective") or task.get("description"),
|
|
expected_output=task.get("expected_output"),
|
|
constraints=task.get("constraints"),
|
|
context=task.get("context"),
|
|
memory_mode=task.get("memory_mode"),
|
|
drive_root=task.get("drive_root"),
|
|
child_drive_root=task.get("child_drive_root") or task.get("drive_root"),
|
|
budget_drive_root=task.get("budget_drive_root"),
|
|
task_constraint=task.get("task_constraint"),
|
|
**_depth_fields,
|
|
# INTENT ONLY. This mirror is written at ASSIGNMENT, one
|
|
# step before the worker dispatches and resolves the
|
|
# child; naming `effective_model_lane`/`model` here wrote
|
|
# whatever the record happened to hold, which on a retry
|
|
# is the PREVIOUS attempt's resolution and on a fresh
|
|
# child is nothing at all.
|
|
model_lane=task.get("model_lane"),
|
|
requested_model_lane=task.get("requested_model_lane"),
|
|
parent_model_lane=task.get("parent_model_lane"),
|
|
requested_executor=task.get("requested_executor"),
|
|
task_group_id=task.get("task_group_id"),
|
|
task_group=task.get("task_group"),
|
|
subagent_envelope=task.get("subagent_envelope"),
|
|
configured_subagent=task.get("configured_subagent"),
|
|
parent_cognitive_route=task.get("parent_cognitive_route"),
|
|
metadata=task.get("metadata") if isinstance(task.get("metadata"), dict) else {},
|
|
result="Subagent assigned to a worker.",
|
|
)
|
|
except Exception:
|
|
log.debug("Failed to mirror running subagent status", exc_info=True)
|
|
w.busy_task_id = task["id"]
|
|
w.in_q.put(task)
|
|
now_ts = time.time()
|
|
resume = task.get("_owner_wait_resume") or {}
|
|
_pool().RUNNING[task["id"]] = {
|
|
"task": dict(task), "worker_id": w.wid,
|
|
"started_at": float(resume.get("started_at") or now_ts), "last_heartbeat_at": now_ts,
|
|
"last_progress_at": now_ts,
|
|
**({"model_wait_quota_clock": dict(resume["model_wait_quota_clock"])}
|
|
if resume.get("model_wait_quota_clock") else {}),
|
|
"soft_sent": False, "attempt": int(task.get("_attempt") or 1),
|
|
}
|
|
task_type = str(task.get("type") or "")
|
|
if task_type in ("evolution", "review"):
|
|
st = _pool().load_state()
|
|
if st.get("owner_chat_id"):
|
|
emoji = '🧬' if task_type == 'evolution' else '🔎'
|
|
_pool().send_with_budget(
|
|
int(st["owner_chat_id"]),
|
|
f"{emoji} {task_type.capitalize()} task {task['id']} started.",
|
|
)
|
|
queue.persist_queue_snapshot(reason="assign_task")
|