Inherit the consciousness origin onto everything a wake starts; one admission door for its roots

Background Consciousness redesign, P3 (provenance and the door). Everything a
wake-up starts inherits its origin — label, ledger category
`consciousness_task`, autonomy level — through one helper
(`consciousness_origin_metadata`), and every producer path derives the level's
consequences before its first contract build:

- promote: the tool event carries the origin by value and the supervisor
  stamps the new root (`actor_id = "consciousness"`, metadata, derived
  contract) — the generalized presence rebind WITHOUT presence's
  project/workspace/source stripping (В9'); a refused promote now carries the
  queue's `_admission_detail` as `detail`.
- schedule_followup: the record template inherits the origin;
  `_task_from_schedule` derives the level before its contract build.
- schedule_subagent: the event's `origin_metadata` lands on the child's task
  metadata (its contract already inherits the parent's list by spread).
- evolution: `toggle_evolution` from a consciousness turn/tree names its
  origin on the event, the campaign record keeps it (`start_evolution_campaign
  origin=`), `enqueue_evolution_task_if_needed` copies it onto every cycle
  task, and the post-task request file carries it into `apply_pending_request`
  (PLAN 5.14 п.7). `post_task_evolution._eligible` refuses when
  `toggle_evolution` is withheld — one fact — so Act cannot reach a campaign
  through one indirection; the globalized promotion view keeps the contract.

ISSUER (PLAN 5.2a): a wake runs on the direct lane but nobody typed it, so
`_routing_issuer` no longer counts `is_direct_chat` as owner provenance for a
consciousness-origin task — its `steer_task` travels as
`[Message from independent task <id>]`; the two other triggers stay, so a
consciousness root relaying a REAL drained owner message keeps the owner's.

В12: `/evolve off` is sticky against the agent tool. `_handle_toggle_evolution`
clears `evolution_owner_stopped` only for an owner-sourced event
(`source == "owner_chat"`); an agent-tool enable while the flag stands is
refused with the same typed shape as the light block. The two GR4-6/GR5-1
pins now drive the owner-sourced event they describe.

The single admission door (PLAN 5.5 / 5.14 п.2): `supervisor.queue.enqueue_task`
derives a consciousness-origin task's level before attaching the contract and,
for a ROOT, under the queue lock refuses in the EXISTING shape
(`_admission_blocked` + `_admission_detail`) when live PENDING+RUNNING
consciousness roots reach `OUROBOROS_CONSCIOUSNESS_MAX_TASKS`
(`consciousness_task_limit`; 0 = never) or the rolling-24h window is spent
(`consciousness_allowance_exhausted`; 0 = may not spend) or unreadable
(`consciousness_allowance_unknown`). Subagents are their root's business; a
snapshot restore is not gated.

Tests: tests/test_consciousness_authority.py (levels, derivation, the
dispatch-only registry, the I3 serialized-request comparison, the level x
install-mode matrix, the lane, ISSUER, promote/followup/subagent inheritance,
eligibility, the sticky stop, campaign and cycle provenance),
tests/test_consciousness_admission.py (the door).

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-16 05:13:19 +03:00
parent ae9fac5ad3
commit d4d1699395
19 changed files with 924 additions and 29 deletions

View file

@ -42,7 +42,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
│ ├── task_admission.py ← Token-owned admission reservations fence duplicate user-ingress ids before Project/workspace/attachment side effects; schedule-dispatch refusals project through the same boundary; queue.py stays the state authority
│ ├── task_lifecycle.py ← Cancellation custody — the ONE settle owner of durable cancel intents: claim → capture → confirmed death → natural-completion re-check → owed delivery registration → settle → delivery/cleanup, plus the `sweep_cancel_intents` watchdog and the queue-owned root-budget admission fence; every custody rule it enforces is stated once in §10 (cancellation custody)
│ ├── cancel_publication.py ← Cancellation settlement publication, re-imported by `task_lifecycle.py`: typed CANCEL_* outcome vocabulary, artifact-honest cancelled result fields, physical-ledger cost reconstruction, salvage adapter, owed-before-settle outbox registration, publication of the STORED terminal truth, capture-miss terminalization/delivery adapter
│ ├── queue_transitions.py ← Queue-owned lifecycle transitions that are not cancellation custody: acceptance-fence open/inspect/seal, explicit budget resume, typed `stop_evolution_tasks` (per-task typed outcomes through the durable-intent ingress, never an in-place prune; an incomplete stop leaves the campaign OPEN via the durable `evolution_owner_stopped` flag, the settle-time backstop in events.py closes it when the last live evolution task settles, and both start ingresses clear the flag BEFORE minting a fresh campaign), and fenced Project deletion (cascade only lineage ROOTS — descendants fall with their trees, one cascade and one summary per tree; tombstone only after provable quiescence; a settled-but-LIVE root still mints the coordination intent, and wind-down defers and RE-CHECKS bounded instead of re-running the cancel pass over a settled-lingering set, which would deliver duplicate owner summaries); imports nothing from task_lifecycle; `supervisor.queue` re-exports these names
│ ├── queue_transitions.py ← Queue-owned lifecycle transitions that are not cancellation custody: acceptance-fence open/inspect/seal, explicit budget resume, typed `stop_evolution_tasks` (per-task typed outcomes through the durable-intent ingress, never an in-place prune; an incomplete stop leaves the campaign OPEN via the durable `evolution_owner_stopped` flag, the settle-time backstop in events.py closes it when the last live evolution task settles, and the OWNER start ingresses — `/evolve start`, an owner-sourced toggle event — clear the flag BEFORE minting a fresh campaign, while the agent's `toggle_evolution` against the set flag is refused: the owner's stop is sticky, В12), and fenced Project deletion (cascade only lineage ROOTS — descendants fall with their trees, one cascade and one summary per tree; tombstone only after provable quiescence; a settled-but-LIVE root still mints the coordination intent, and wind-down defers and RE-CHECKS bounded instead of re-running the cancel pass over a settled-lingering set, which would deliver duplicate owner summaries); imports nothing from task_lifecycle; `supervisor.queue` re-exports these names
│ ├── terminal_delivery.py ← Durable terminal-answer delivery seam: restart-surviving `delivery_id` dedupe + bounded PENDING outbox `state/terminal_deliveries.json` (owed before enqueue, cleared in the delivering write, replayed on boot and on the supervisor tick), shared by natural final answers (every root registers at durable-result persistence), cancel salvage, cascade digest, and non-retry reap; the cascade digest enumerates descendants by ANCESTRY (parent-chain walk, never `root_task_id` equality); eviction past outbox capacity is disclosed via typed `terminal_delivery_exhausted`, never a silent pop; salvage messages carry a bounded preview plus a full-copy receipt (path, size, full 64-hex sha256 or an explicit marker) and route by lineage chat — no resolvable chat records a typed `terminal_delivery_handoff` row; reads/mutations are row-strict per §10; the delivery id digests only the STABLE part (task id + status framing + core answer) so a replay whose rebuilt note shrank dedups instead of double-sending; the per-origin projection is `host_salvage` receipt / `host_notice` own text kept as a System row with its markdown / `model_final` assistant projection
│ ├── task_reaper.py ← Single-owner off-loop queue/pump for timeout teardown and health-prepared terminal-file/crash jobs, with same-job deferred replay bound to the captured worker/attempt/root, so old file recovery cannot replace a newer execution; keeps supervisor intake responsive. An unconfirmed death holds the slot reaping and the task RUNNING with task_reaper_wedged; confirmed death precedes delegated-custody reconciliation and retry. No cancel intents minted (§5 Supervisor Loop).
│ ├── owner_stop.py ← Owner graceful stop: `finalize_then_cancel` policy as an axis on the SAME durable cancel intent (monotonic — immediate HARDENS a pending graceful, never softens back; hardening revokes an unread control via mailbox revocation, and the loop revalidates durable policy at drain); one deterministic typed `finalize_now` control whose first line is the `owner_requested_finalization` literal, routed by the loop to its own rail (zero or one tool-less turn); descendants settle first with a bounded child projection fed to the root's final turn; the grace budget starts at the durable control DRAIN (first drain wins), bounded by `request + OWNER_STOP_OUTER_CAP_SEC`, and neither anchor is ever progress-extended; `running_owner_stop_tasks` bypasses only the generic idle/finalization-grace rails; a COMPLETED finalize root suppresses the redundant cascade summary
@ -170,6 +170,8 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
├── mcp_client.py ← MCP client: parses MCP_SERVERS, validates transports, masks tokens, prefixes tools `mcp_<server>__<tool>`, guarded SDK import; MCP descriptions/results stay untrusted data (§6)
├── safety.py ← Safety supervisor call with a bounded newest-first context budget (the omission marker is reserved INSIDE the budget); a 429 on the safety check is an infrastructure fact about the supervisor, not a verdict about the tool call — one deadline-capped retry, then the typed `⚠️ SAFETY_UNAVAILABLE` non-verdict telling the agent to retry the same call, not reword it; process-local storm latch answers in-window checks without provider calls; durable `safety_check_rate_limited` audit event; structured insufficient-quota keeps its PERMANENT classification and still blocks as a verdict; the serialized SUBJECT has its own 250k-char budget (`_SAFETY_SUBJECT_CHAR_BUDGET`, rendered `ensure_ascii=False`) and an over-budget subject is refused fail-closed with the typed `⚠️ SAFETY_SUBJECT_TOO_LARGE_BLOCKED` denial plus a durable `safety_subject_too_large` event — never truncated, because anything past a cut would run unreviewed; the fail-open cases are owned by prompts/SYSTEM.md
├── consciousness.py ← Background thinking loop with live progress emission (§6)
├── consciousness_authority.py ← The three autonomy levels of a consciousness wake-up (observe/act/full, carried as `metadata.consciousness_autonomy` beside `metadata.initiator`), the one helper that derives a level's two consequences at task build (`disabled_tools` as an exception list, `runtime_mode_cap=light` below Full), the origin keys everything a wake starts inherits, and the stricter-of-install-and-cap mode the tool dispatcher's light gates read; for a consciousness-origin task `disabled_tools` binds at dispatch only so its prompt prefix matches an owner turn's
├── consciousness_allowance.py ← Rolling-24h spend of consciousness (its wakes plus the roots they started) read off the usage ledger: roots by the category of final rows inside the 48 h fold horizon, every money row kind under them through the one `_usage_rows._summary` reducer, fingerprint-memoized selection with a per-call time filter, typed `allowance_unknown` on a read failure; the single admission door in `supervisor/queue.py` and the alarm read it
├── consolidator.py ← Dialogue consolidation with a generation-aware cursor over the ordered archive chain; an unfindable generation appends a loud durable `[MEMORY GAP]` block, never a silent offset reset; a run that advances without a new failure clears `last_consolidation_error`, and a knowledge-nomination batch that is not fully published leaves `last_unpublished_nominations` in `dialogue_meta.json` (a fully published batch clears it)
├── memory.py ← Scratchpad, identity, chat history
├── knowledge.py ← `ouroboros/knowledge.py`: shared linked-Markdown note addressing, exact source reads, revision-checked writes and generated shelf indexes for global and project knowledge; one writer preserves authored understanding, unknown metadata and provenance so concurrent cognition cannot silently overwrite a newer note or invent understanding from a truncated prefix

File diff suppressed because one or more lines are too long

View file

@ -390,6 +390,9 @@ def _run_global_backlog_promotion_only(
"type": str(task.get("type") or "task"),
"source": "project_scoped_global_improvement",
"metadata": {"globalized_from_project_task": True},
# The eligibility probe reads the contract (disabled_tools), so the
# globalized view keeps it: a level that may not evolve stays that way.
**({"task_contract": dict(task["task_contract"])} if isinstance(task.get("task_contract"), dict) else {}),
}
maybe_promote(env, global_task, sanitized_entry, llm)
except Exception as error:

View file

@ -26,6 +26,7 @@ import logging
import pathlib
from typing import Any, Dict, Optional
from ouroboros.consciousness_authority import consciousness_origin_metadata, task_disabled_tools
from ouroboros.evolution_fingerprint import _PLAN_REVIEW_SUFFIX
from ouroboros.config import runtime_setting
@ -70,6 +71,11 @@ def _eligible(task: Dict[str, Any]) -> bool:
return False
if str(task.get("delegation_role") or "") == "subagent":
return False
# ONE fact: a task whose contract withholds toggle_evolution (a consciousness
# wake-up below Full, and every root it started) may not propose evolution
# either — otherwise Act would reach a campaign through one indirection.
if "toggle_evolution" in task_disabled_tools(task):
return False
return True
@ -309,6 +315,9 @@ def _write_request(drive_root: pathlib.Path, decision: Dict[str, Any], task: Dic
"backlog_id": decision.get("backlog_id") or "",
"source": "post_task",
"origin_task_id": str(task.get("id") or ""),
# A Full-level consciousness tree proposing evolution keeps its origin, so
# the campaign and its cycle tasks stay inside the consciousness allowance.
**consciousness_origin_metadata(task.get("metadata")),
}
path = drive_root / _REQUEST_REL
# Atomic publish: the supervisor polls every tick, so a partial write must
@ -441,7 +450,8 @@ def apply_pending_request(drive_root: Any) -> bool:
return False
if bool(req.get("requires_plan_review", True)):
objective += _PLAN_REVIEW_SUFFIX
if not start_evolution_campaign(objective, source="post_task"):
origin = consciousness_origin_metadata(req)
if not start_evolution_campaign(objective, source="post_task", **({"origin": origin} if origin else {})):
return False
# Link the promoted backlog id to the campaign so close-on-commit (Phase 2 C)
# can mark it done when the cycle is absorbed. Validate it against the OPEN

View file

@ -166,12 +166,19 @@ def _routing_issuer(ctx: ToolContext) -> Dict[str, Any]:
contract a Swarm root never has, an empty client id read as "agent-issued",
a room veto keyed on the chat): the host now states it once, and the model
has no argument to claim otherwise.
A consciousness wake-up runs on the direct lane too, but nobody typed it: its
``is_direct_chat`` fact does NOT make it an owner turn (PLAN 5.2a) — it
speaks as a task. The two other triggers stay, so a consciousness turn that
relays a REAL owner message it drained keeps the owner's provenance.
"""
metadata = getattr(ctx, "task_metadata", None)
metadata = metadata if isinstance(metadata, dict) else {}
delivery = getattr(ctx, "last_owner_delivery", None)
from ouroboros.consciousness_authority import is_consciousness_origin
if (
bool(getattr(ctx, "is_direct_chat", False))
(bool(getattr(ctx, "is_direct_chat", False)) and not is_consciousness_origin(metadata))
or str(metadata.get("client_message_id") or "").strip()
or (isinstance(delivery, dict) and delivery)
):
@ -361,6 +368,13 @@ def _promote_chat_to_task(
"presence": dict(presence),
"task_contract": dict(getattr(ctx, "task_contract", {}) or {}),
})
# A promote from a consciousness turn/tree mints a consciousness root: the
# origin label, ledger category and level ride the event by value; the
# supervisor stamps them on the new root (worker_promotion) — no presence-style
# stripping, the wake chooses project/workspace like any Main turn (В9').
from ouroboros.consciousness_authority import consciousness_origin_metadata
evt.update(consciousness_origin_metadata(metadata))
_attach_origin_from_metadata(ctx, evt)
predecessor_error = _attach_predecessor_authority_from_metadata(
ctx, evt, predecessor_task_id,

View file

@ -325,11 +325,16 @@ def _toggle_evolution(ctx: ToolContext, enabled: bool, objective: str = "") -> s
block = ""
if block:
return block
from ouroboros.consciousness_authority import consciousness_origin_metadata
ctx.pending_events.append({
"type": "toggle_evolution",
"enabled": bool(enabled),
"objective": str(objective or "").strip(),
"ts": utc_now_iso(),
# A Full-level consciousness turn/tree names itself: the campaign and its
# cycle tasks then stay inside the consciousness allowance (PLAN 5.14 п.7).
**consciousness_origin_metadata(getattr(ctx, "task_metadata", None)),
})
state_str = "ON" if enabled else "OFF"
return f"OK: evolution mode toggled {state_str}."

View file

@ -21,6 +21,7 @@ from typing import Any, Dict, List, Optional
from ouroboros.artifacts import attachment_manifest_projection, resolve_attachment_manifest
from ouroboros.config import get_max_subagent_depth
from ouroboros.consciousness_authority import consciousness_origin_metadata
from ouroboros.depth_evidence import parse_task_depth
from ouroboros.contracts.task_contract import (
build_task_contract,
@ -813,6 +814,8 @@ def _schedule_task(ctx: ToolContext, internal: Dict[str, Any] | None = None, /,
"required_capabilities": required_caps,
**intent_fields,
"subagent_envelope": envelope,
# A child of a consciousness turn/tree carries the origin (label, category, level).
"origin_metadata": consciousness_origin_metadata(metadata),
}
_populate_subagent_event_extras(
evt, current_chat_id=current_chat_id, child_drive=child_drive,

View file

@ -23,6 +23,7 @@ from ouroboros.tools.tool_result import ToolResult, _publish_tool_result
import uuid
from typing import Any, Dict, List
from ouroboros.consciousness_authority import consciousness_origin_metadata
from ouroboros.deadline_utils import parse_deadline_ts
from ouroboros.tools.registry import ToolContext, ToolEntry
@ -224,6 +225,9 @@ def _handle_schedule_followup(ctx: ToolContext, **params) -> str:
**({"chat_id": source_chat_id} if source_chat_id not in (None, "") else {}),
},
}
# A follow-up from a consciousness turn/tree starts a consciousness root: the
# origin, category and level ride the template; admission derives the rest.
record["task"]["metadata"].update(consciousness_origin_metadata(metadata_src))
presence = metadata_src.get("presence") if isinstance(metadata_src, dict) else None
contract = getattr(ctx, "task_contract", None)
if isinstance(presence, dict) and presence and isinstance(contract, dict):

View file

@ -184,24 +184,39 @@ def _drive_cancel_task_event(evt: Dict[str, Any], ctx: Any) -> None:
def _handle_toggle_evolution(evt: Dict[str, Any], ctx: Any) -> None:
"""Toggle evolution mode from LLM tool call."""
"""Toggle evolution mode from an LLM tool call (or an owner-sourced event).
Owner decision В12: ``/evolve off`` is STICKY against the agent tool. Only an
event with owner provenance (``source == "owner_chat"``) may clear the
durable ``evolution_owner_stopped`` flag; an ``agent_tool`` enable while the
flag is set is refused with the same typed shape as the light-mode block.
"""
enabled = bool(evt.get("enabled"))
owner_sourced = str(evt.get("source") or "") == "owner_chat"
if enabled:
from supervisor.evolution_lifecycle import evolution_block_reason, start_evolution_campaign
from ouroboros.consciousness_authority import consciousness_origin_metadata
block = evolution_block_reason()
if not block and not owner_sourced and bool(ctx.load_state().get("evolution_owner_stopped")):
block = (
"🧬 Evolution stayed OFF: the owner stopped evolution (/evolve off), and that stop "
"is sticky against toggle_evolution. Only the owner's /evolve start re-arms it; "
"no campaign was started."
)
if block:
st = ctx.load_state()
if st.get("owner_chat_id"):
ctx.send_with_budget(int(st["owner_chat_id"]), block)
return
# GR4-6: clear the durable owner-stop flag BEFORE the campaign is
# minted. The old order (campaign first, flag cleared in a later state
# write) left a window where the owner-stop backstop — fired by an old
# evolution task settling — read flag=True + campaign=active and closed
# the FRESH campaign. This clear is owner-authorized (the owner is
# explicitly starting evolution). GR5-1: the prior value is captured in
# the same locked write so a failed start can restore it.
# GR4-6: an OWNER start clears the durable owner-stop flag BEFORE the
# campaign is minted. The old order (campaign first, flag cleared in a
# later state write) left a window where the owner-stop backstop — fired
# by an old evolution task settling — read flag=True + campaign=active and
# closed the FRESH campaign. GR5-1: the prior value is captured in the same
# locked write so a failed start can restore it. An agent-tool start
# reaches this point only with the flag already clear (checked above), so
# it clears nothing.
from supervisor.state import update_state as _update_state
_prior_owner_stop = {"value": False}
@ -210,9 +225,13 @@ def _handle_toggle_evolution(evt: Dict[str, Any], ctx: Any) -> None:
_prior_owner_stop["value"] = bool(live.get("evolution_owner_stopped"))
live["evolution_owner_stopped"] = False
_update_state(_clear_owner_stop)
if owner_sourced:
_update_state(_clear_owner_stop)
origin = consciousness_origin_metadata(evt)
source = "owner_chat" if owner_sourced else "agent_tool"
try:
if not start_evolution_campaign(str(evt.get("objective") or ""), source="agent_tool"):
if not start_evolution_campaign(str(evt.get("objective") or ""), source=source,
**({"origin": origin} if origin else {})):
raise RuntimeError("campaign write was refused")
except Exception:
log.warning("Failed to start evolution campaign from agent tool", exc_info=True)

View file

@ -524,6 +524,7 @@ def _handle_schedule_task(evt: Dict[str, Any], ctx: Any) -> None:
"configured_subagent": configured_subagent,
"parent_cognitive_route": parent_cognitive_route,
"parent_id": parent_id,
"origin_metadata": evt.get("origin_metadata"),
})
scheduled_failure_reason = ""
scheduled_failure_detail = ""

View file

@ -15,6 +15,7 @@ import pathlib
import uuid
from typing import Any, Dict, Optional
from ouroboros.consciousness_authority import consciousness_origin_metadata
from ouroboros.cost_projection import honest_cost_pair_amount
from ouroboros.evolution_fingerprint import canonical_objective_fingerprint
from ouroboros.outcomes import normalize_outcome_axes
@ -112,9 +113,7 @@ def enqueue_evolution_task_if_needed() -> None:
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.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.",
@ -200,10 +199,7 @@ def enqueue_evolution_task_if_needed() -> None:
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.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.",
@ -213,7 +209,7 @@ def enqueue_evolution_task_if_needed() -> None:
"id": tid, "type": "evolution",
"chat_id": int(owner_chat_id),
"text": q.build_evolution_task_text(cycle),
"metadata": {"evolution_transaction": transaction},
"metadata": {"evolution_transaction": transaction, **consciousness_origin_metadata(campaign)},
}
q.attach_task_contract(task)
q.enqueue_task(task)
@ -296,8 +292,8 @@ def evolution_block_reason() -> str:
return ""
def start_evolution_campaign(objective: str = "", *, source: str = "owner") -> Dict[str, Any]:
"""Start or resume the active evolution campaign."""
def start_evolution_campaign(objective: str = "", *, source: str = "owner", origin: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
"""Start or resume the active evolution campaign; ``origin`` = the consciousness origin keys of the starting turn, kept for its cycle tasks."""
from supervisor import state
state.assert_test_data_path(state.STATE_PATH)
@ -315,7 +311,7 @@ def start_evolution_campaign(objective: str = "", *, source: str = "owner") -> D
"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 ""),
"source": str(source or ""), **consciousness_origin_metadata(origin),
"started_at": now,
"updated_at": now,
"cycles_done": 0,
@ -333,6 +329,7 @@ def start_evolution_campaign(objective: str = "", *, source: str = "owner") -> D
campaign["objective"] = objective
if not str(campaign.get("source") or "").strip() and source:
campaign["source"] = str(source)
campaign.update({k: v for k, v in consciousness_origin_metadata(origin).items() if not campaign.get(k)})
campaign["status"] = "active"
campaign["updated_at"] = now
generation = current_evolution_boot_generation()

View file

@ -31,6 +31,7 @@ from ouroboros.config import (
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
@ -166,7 +167,7 @@ def enqueue_task(
"""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(t)
attach_task_contract(apply_consciousness_authority(t))
with _queue_lock:
require_unique_id = bool(t.pop("_require_unique_task_id", False))
require_worker_pool = bool(t.pop("_require_worker_pool", False))
@ -220,6 +221,12 @@ def enqueue_task(
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)
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
@ -273,6 +280,68 @@ def enqueue_task(
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_admission_block(task: Dict[str, Any]) -> 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.
"""
if (
not is_consciousness_origin(task.get("metadata"))
or str(task.get("delegation_role") or "root") != "root"
):
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"
))
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 "")

View file

@ -16,6 +16,7 @@ import pathlib
import time
import uuid
from typing import Any, Dict, List
from ouroboros.consciousness_authority import apply_consciousness_authority
from ouroboros.contracts.task_contract import build_task_contract, normalize_allowed_resources
from ouroboros.schedule_contract import RESERVED_TEMPLATE_FIELDS, schedule_slug
from ouroboros.skill_loader import skill_identity_collision_names
@ -264,7 +265,7 @@ def _task_from_schedule(record: Dict[str, Any]) -> Dict[str, Any]:
existing_contract = template.get("task_contract") if isinstance(template.get("task_contract"), dict) else {}
if existing_contract:
task["task_contract"] = existing_contract
task["task_contract"] = build_task_contract(task)
task["task_contract"] = build_task_contract(apply_consciousness_authority(task))
presence = metadata.get("presence")
workspace = task["task_contract"]["workspace"]
if isinstance(presence, dict) and presence and workspace["root"]:

View file

@ -48,6 +48,8 @@ def build_scheduled_task_payload(fields: Dict[str, Any]) -> Dict[str, Any]:
configured_subagent = fields.get("configured_subagent") if isinstance(fields.get("configured_subagent"), dict) else {}
parent_cognitive_route = fields.get("parent_cognitive_route") if isinstance(fields.get("parent_cognitive_route"), dict) else {}
directory_options = {key: fields[key] for key in ("directory_strategy", "scope_paths") if key in fields}
# A child of a consciousness turn/tree inherits its origin label, category and level.
origin_metadata = fields.get("origin_metadata") if isinstance(fields.get("origin_metadata"), dict) else {}
task: Dict[str, Any] = {
"id": tid,
"type": "task",
@ -115,6 +117,7 @@ def build_scheduled_task_payload(fields: Dict[str, Any]) -> Dict[str, Any]:
"parent_cognitive_route": parent_cognitive_route,
**directory_options,
"root_cost_ceiling_usd": root_cost_ceiling_usd,
**origin_metadata,
},
}
if not drive_root:

View file

@ -389,6 +389,7 @@ def promote_chat_to_task(evt: dict, ctx: Any) -> dict:
land in that thread) and the optional ``project_id`` scope; it competes for
the project writer lease like any other top-level project task.
"""
from ouroboros.consciousness_authority import apply_consciousness_authority, consciousness_origin_metadata
from ouroboros.contracts.task_contract import attach_task_contract
from ouroboros.project_naming import admission_names
@ -460,6 +461,13 @@ def promote_chat_to_task(evt: dict, ctx: Any) -> dict:
"promotion_admission_token": admission_token,
**_promoted_force_plan_metadata(evt),
}
origin = consciousness_origin_metadata(evt)
if origin:
# A root consciousness started keeps its origin/category/level by value (В9':
# ordinary Main flow, no presence-style project/workspace/source stripping).
task["actor_id"] = "consciousness"
task.setdefault("metadata", {}).update(origin)
apply_consciousness_authority(task)
inherited_attachment_manifest = _pool()._apply_presence_promotion_authority(
evt, task, objective=objective, expected_output=expected_output,
)
@ -561,6 +569,7 @@ def promote_chat_to_task(evt: dict, ctx: Any) -> dict:
return _pool()._reject_promoted_after_attachment_stage({
"status": "needs_manual_target",
"reason": str(admitted.get("_admission_blocked") or "admission_fence"),
"detail": str(admitted.get("_admission_detail") or ""),
"project_lifecycle": str(admitted.get("_project_lifecycle") or ""),
"task_id": tid,
}, attachment_manifest)

View file

@ -0,0 +1,151 @@
"""The single admission door for the roots consciousness starts (P3, PLAN 5.5 / 5.14 п.2).
``supervisor.queue.enqueue_task`` — the ONE place every pooled root passes —
refuses a consciousness-origin ROOT under the queue lock when the concurrency
cap (``OUROBOROS_CONSCIOUSNESS_MAX_TASKS``, 0 = never) or the rolling-24h
allowance (``OUROBOROS_CONSCIOUSNESS_DAILY_USD``, 0 = may not spend) says so,
in the queue's EXISTING refusal shape (``_admission_blocked`` + the
``_admission_detail`` the depth guard already uses). Subagents are their
root's business; owner work is never gated; a snapshot restore re-admits
already-admitted work.
"""
from __future__ import annotations
import types
import pytest
from ouroboros import consciousness_authority as ca
AVAILABLE = {"status": "available", "limit_usd": 20.0, "accounted_usd": 4.0, "remaining_usd": 16.0,
"unknown_unmetered": 0, "resets_at": ""}
@pytest.fixture
def door(tmp_path, monkeypatch):
from supervisor import queue
queue.init(tmp_path)
pending: list = []
running: dict = {}
queue.init_queue_refs(pending, running, {"value": 0})
monkeypatch.setattr(queue, "ACCEPTANCE_FENCES", {})
monkeypatch.setattr("ouroboros.config.get_consciousness_max_tasks", lambda: 2)
monkeypatch.setattr("ouroboros.consciousness_allowance.allowance_window", lambda root, now=None: dict(AVAILABLE))
return queue, pending, running
def _root(tid, level="act", **extra):
return {"id": tid, "type": "task", "text": tid, "chat_id": 1, "root_task_id": tid, "delegation_role": "root",
"metadata": {"initiator": "consciousness", "usage_category": "consciousness_task",
"consciousness_autonomy": level, **extra}}
def test_a_consciousness_root_is_admitted_with_its_level_derived(door):
queue, pending, _running = door
admitted = queue.enqueue_task(_root("c1"))
assert "_admission_blocked" not in admitted and [t["id"] for t in pending] == ["c1"]
assert admitted["task_contract"]["disabled_tools"] == list(ca.ACT_DISABLED)
assert admitted["metadata"]["runtime_mode_cap"] == "light"
assert queue.live_consciousness_root_count() == 1
def test_the_concurrency_cap_counts_live_pending_and_running_roots(door):
queue, pending, running = door
assert "_admission_blocked" not in queue.enqueue_task(_root("c1"))
running["c2"] = {"task": _root("c2")}
assert queue.live_consciousness_root_count() == 2
refused = queue.enqueue_task({**_root("c3"), "_admission_token": "tok"})
assert refused["_admission_blocked"] == "consciousness_task_limit"
assert refused["_admission_detail"] == "2 of 2 consciousness-started tasks already live"
assert [t["id"] for t in pending] == ["c1"]
assert "c3" not in queue.ADMISSION_RESERVATIONS
# A settled root frees its slot: the count is live, not historical.
running.clear()
assert "_admission_blocked" not in queue.enqueue_task(_root("c3"))
def test_max_tasks_zero_means_consciousness_never_starts_tasks(door, monkeypatch):
queue, pending, _running = door
monkeypatch.setattr("ouroboros.config.get_consciousness_max_tasks", lambda: 0)
refused = queue.enqueue_task(_root("c1"))
assert refused["_admission_blocked"] == "consciousness_task_limit"
assert "OUROBOROS_CONSCIOUSNESS_MAX_TASKS=0" in refused["_admission_detail"]
assert pending == []
def test_an_exhausted_allowance_refuses_with_the_window_facts(door, monkeypatch):
queue, pending, _running = door
monkeypatch.setattr("ouroboros.consciousness_allowance.allowance_window", lambda root, now=None: {
**AVAILABLE, "status": "exhausted", "accounted_usd": 21.5, "remaining_usd": 0.0,
"unknown_unmetered": 1, "resets_at": "2026-09-17T08:00:00+00:00"})
refused = queue.enqueue_task(_root("c1"))
assert refused["_admission_blocked"] == "consciousness_allowance_exhausted"
assert refused["_admission_detail"] == (
"$21.50 (at least) of $20.00 spent in the last 24 h; resets at 2026-09-17T08:00:00+00:00")
assert pending == []
monkeypatch.setattr("ouroboros.consciousness_allowance.allowance_window", lambda root, now=None: {
**AVAILABLE, "status": "exhausted", "limit_usd": 0.0, "accounted_usd": 0.0, "remaining_usd": 0.0})
zero = queue.enqueue_task(_root("c2"))
assert zero["_admission_blocked"] == "consciousness_allowance_exhausted"
assert "OUROBOROS_CONSCIOUSNESS_DAILY_USD=0" in zero["_admission_detail"]
def test_an_unreadable_ledger_refuses_honestly_instead_of_admitting(door, monkeypatch):
queue, pending, _running = door
monkeypatch.setattr("ouroboros.consciousness_allowance.allowance_window", lambda root, now=None: {
"status": "allowance_unknown", "error": "OSError: ledger locked", "limit_usd": 20.0,
"accounted_usd": None, "remaining_usd": None, "resets_at": ""})
refused = queue.enqueue_task(_root("c1"))
assert refused["_admission_blocked"] == "consciousness_allowance_unknown"
assert "ledger locked" in refused["_admission_detail"] and pending == []
def test_subagents_and_owner_work_pass_the_door_untouched(door, monkeypatch):
queue, pending, running = door
monkeypatch.setattr("ouroboros.config.get_consciousness_max_tasks", lambda: 0)
child = {**_root("kid"), "delegation_role": "subagent", "root_task_id": "wake-1", "parent_task_id": "wake-1"}
assert "_admission_blocked" not in queue.enqueue_task(child)
assert child["metadata"]["initiator"] == "consciousness" # the label rides; the door ignores children
owner = {"id": "o1", "type": "task", "text": "owner work", "chat_id": 1, "delegation_role": "root",
"metadata": {"client_message_id": "cm-1"}}
assert "_admission_blocked" not in queue.enqueue_task(owner)
assert [t["id"] for t in pending] == ["kid", "o1"]
assert queue.live_consciousness_root_count() == 0
def test_a_snapshot_restore_is_not_gated(door, monkeypatch):
queue, pending, _running = door
monkeypatch.setattr("ouroboros.config.get_consciousness_max_tasks", lambda: 0)
restored = queue.enqueue_task(_root("c1"), restoring_snapshot=True)
assert "_admission_blocked" not in restored and [t["id"] for t in pending] == ["c1"]
def test_the_refusal_uses_the_queue_vocabulary_the_promote_path_already_reads():
"""The promote handler turns ``_admission_blocked``/``_admission_detail`` into the
typed PROMOTE_REJECTED the tool renders (pinned in test_consciousness_authority);
the depth guard writes the same pair, so no new refusal form exists."""
import inspect
from supervisor import task_admission, worker_promotion
assert '_admission_detail' in inspect.getsource(task_admission.reject_invalid_task_depth)
assert '"detail": str(admitted.get("_admission_detail") or "")' in inspect.getsource(worker_promotion.promote_chat_to_task)
def test_promote_refusal_carries_the_door_detail(tmp_path, monkeypatch):
import supervisor.workers as workers
monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path)
ctx = types.SimpleNamespace(
enqueue_task=lambda task: {**task, "_admission_blocked": "consciousness_task_limit",
"_admission_detail": "2 of 2 consciousness-started tasks already live"},
persist_queue_snapshot=lambda **_k: True, load_state=lambda: {"owner_chat_id": 1},
)
evt = {"type": "promote_chat_to_task", "task_id": "c0000002", "objective": "x", "chat_id": 1,
"workspace": "none", "initiator": "consciousness"}
outcome = workers.promote_chat_to_task(evt, ctx)
assert outcome["status"] == "needs_manual_target"
assert outcome["reason"] == "consciousness_task_limit"
assert outcome["detail"] == "2 of 2 consciousness-started tasks already live"

View file

@ -0,0 +1,597 @@
"""Authority levels of a consciousness wake-up and how they follow its work (P3).
Observe / Act / Full (owner decision В10', default Act) are carried as
``metadata.consciousness_autonomy`` and derived at task build into the
contract's ``disabled_tools`` and the per-task ``runtime_mode_cap`` (В21=A);
for a consciousness-origin task the disabled list binds at DISPATCH ONLY so the
wake's tool schemas and prompt prefix are byte-identical to an owner turn's
(В31=B, the I3 comparison below). The origin (label, ledger category, level)
is inherited by everything the wake starts; a wake speaks as a task through
``steer_task`` (ISSUER, PLAN 5.2a); ``/evolve off`` is sticky against the agent
tool (В12); a Full-level campaign stays inside the consciousness tree.
"""
from __future__ import annotations
import json
import pathlib
import types
import pytest
from ouroboros import consciousness_authority as ca
from ouroboros.tools.registry import ToolContext, ToolRegistry
WAKE_META = {
"initiator": "consciousness", "usage_category": "consciousness",
"wake_reason": "heartbeat", "consciousness_autonomy": "act", "model_role": "consciousness",
}
def _wake_task(level="act", **extra):
task = {"id": "wake-1", "type": "task", "text": "wake", "_is_direct_chat": True, "chat_id": 1,
"metadata": {**WAKE_META, "consciousness_autonomy": level, **extra}}
return ca.apply_consciousness_authority(task)
def _registry(tmp_path, metadata=None, *, task_id="turn-1"):
repo = tmp_path / "repo"
repo.mkdir(exist_ok=True)
(repo / "README.md").write_text("ok\n", encoding="utf-8")
drive = tmp_path / "drive"
drive.mkdir(exist_ok=True)
reg = ToolRegistry(repo_dir=repo, drive_root=drive)
reg.set_context(ToolContext(
repo_dir=repo, drive_root=drive, task_id=task_id, is_direct_chat=True,
task_metadata=dict(metadata or {}),
))
return reg
# --- the level tables -------------------------------------------------------------
def test_levels_and_their_two_consequences():
assert ca.LEVELS == ("observe", "act", "full")
assert ca.disabled_tools_for("full") == []
assert ca.disabled_tools_for("act") == list(ca.ACT_DISABLED)
observe = ca.disabled_tools_for("observe")
assert set(ca.ACT_DISABLED) <= set(observe)
assert {"promote_chat_to_task", "schedule_subagent", "write_file", "run_command",
"browser_action", "initiate_presence", "submit_skill_to_hub"} <= set(observe)
# The nanny of a running campaign is never withheld, at any level.
assert "steer_task" not in observe and "steer_task" not in ca.disabled_tools_for("act")
# Observe is an EXCEPTION list: reading and talking stay available by default.
for name in ("read_file", "web_search", "browse_page", "send_user_message", "escalate",
"knowledge_write", "update_scratchpad", "switch_model", "enable_tools"):
assert name not in observe
assert ca.runtime_mode_cap_for("act") == "light" == ca.runtime_mode_cap_for("observe")
assert ca.runtime_mode_cap_for("full") == ""
def test_unknown_level_falls_back_to_the_owner_setting(monkeypatch):
monkeypatch.setenv("OUROBOROS_CONSCIOUSNESS_AUTONOMY", "observe")
assert ca.normalize_level("bogus") == "observe"
assert ca.normalize_level("") == "observe"
assert ca.normalize_level("FULL") == "full"
def test_observe_table_covers_every_registry_entry_marked_mutates_worktree(tmp_path):
"""The registry marker is the second source of the same fact; the table cannot drift."""
reg = _registry(tmp_path)
marked = {e.name for e in reg._entries.values() if e.mutates_worktree and not e.alias_for}
assert marked, "the catalog carries mutates_worktree entries"
missing = marked - set(ca.OBSERVE_DISABLED)
assert not missing, f"mutates_worktree entries missing from OBSERVE_WORLD_MUTATION_TOOLS: {sorted(missing)}"
unknown = set(ca.OBSERVE_DISABLED) - {e.name for e in reg._entries.values()}
# Skill/project tools are registered lazily (skills, journal); the built-in names must exist.
assert unknown <= {"toggle_skill", "skill_owner_action", "journal_write", "workpad_write",
"configure_presence", "initiate_presence", "delegate_start"}, sorted(unknown)
# --- derivation at task build ---------------------------------------------------
def test_apply_consciousness_authority_derives_both_consequences_once():
task = _wake_task("act")
assert task["metadata"]["disabled_tools"] == list(ca.ACT_DISABLED)
assert task["metadata"]["runtime_mode_cap"] == "light"
full = _wake_task("full")
assert full["metadata"]["disabled_tools"] == [] and full["metadata"]["runtime_mode_cap"] == ""
# An explicit producer list stands; an owner turn is untouched.
explicit = _wake_task("act", disabled_tools=["web_search"])
assert explicit["metadata"]["disabled_tools"] == ["web_search"]
owner = ca.apply_consciousness_authority({"id": "o", "metadata": {"client_message_id": "cm"}})
assert "disabled_tools" not in owner["metadata"] and "runtime_mode_cap" not in owner["metadata"]
def test_contract_carries_the_derived_list_and_origin_helpers():
from ouroboros.contracts.task_contract import attach_task_contract
task = attach_task_contract(_wake_task("observe"))
assert task["task_contract"]["disabled_tools"] == ca.disabled_tools_for("observe")
assert "toggle_evolution" in ca.task_disabled_tools(task)
origin = ca.consciousness_origin_metadata(task["metadata"])
assert origin == {"initiator": "consciousness", "usage_category": "consciousness_task",
"consciousness_autonomy": "observe"}
assert ca.consciousness_origin_metadata({"client_message_id": "cm"}) == {}
assert ca.is_consciousness_origin(origin) and not ca.is_consciousness_origin(None)
@pytest.mark.parametrize(("install", "cap", "expected"), [
("cyber_pro", "light", "light"), ("pro", "light", "light"), ("advanced", "light", "light"),
("light", "light", "light"), ("light", "advanced", "light"), ("cyber_pro", "", "cyber_pro"),
("pro", "bogus", "pro"),
])
def test_effective_runtime_mode_is_the_stricter_of_install_and_cap(install, cap, expected):
assert ca.effective_runtime_mode(install, {"runtime_mode_cap": cap}) == expected
assert ca.effective_runtime_mode(install, None) == install
# --- dispatch-only enforcement (В31=B) -------------------------------------------
def test_consciousness_contract_keeps_the_full_schema_set_and_refuses_at_dispatch(tmp_path, monkeypatch):
monkeypatch.setenv("OUROBOROS_RUNTIME_MODE", "advanced")
main = _registry(tmp_path, {"client_message_id": "cm-1"})
wake = _registry(tmp_path, _wake_task("act")["metadata"])
# Same schemas, same advertised names, same initial envelope, same omission manifest.
assert wake.schemas() == main.schemas()
assert wake.available_tools() == main.available_tools()
assert wake.initial_tool_names() == main.initial_tool_names()
assert wake.capability_omissions() == main.capability_omissions()
assert not any(item.get("reason") == "disabled_by_contract" for item in wake.capability_omissions())
assert "toggle_evolution" in wake.available_tools()
assert wake.get_schema_by_name("toggle_evolution") is not None
assert wake.policy_hidden_reason("toggle_evolution") is None
# The dispatcher is the mechanism: the withheld name is refused with the typed block.
result = wake.execute("toggle_evolution", {"enabled": True, "objective": "x"})
assert "RESOURCE_CONSTRAINT_BLOCKED" in result and "toggle_evolution" in result
for name, args in (("request_restart", {}), ("set_tool_timeout", {"seconds": 30}),
("toggle_consciousness", {"action": "stop"})):
assert "RESOURCE_CONSTRAINT_BLOCKED" in wake.execute(name, args), name
def test_an_ordinary_contract_still_hides_its_disabled_tools(tmp_path):
"""The dispatch-only case is the consciousness special case, not a general change."""
reg = _registry(tmp_path, {"disabled_tools": ["toggle_evolution"]})
assert "toggle_evolution" not in reg.available_tools()
assert all(s["function"]["name"] != "toggle_evolution" for s in reg.schemas())
assert reg.get_schema_by_name("toggle_evolution") is None
assert reg.policy_hidden_reason("toggle_evolution") == "disabled by this task's contract (disabled_tools)"
assert any(item.get("reason") == "disabled_by_contract" for item in reg.capability_omissions())
def test_i3_serialized_request_prefix_matches_an_owner_turn(tmp_path, monkeypatch):
"""The provider request's cached prefix — the tool schema array and the two
cached system blocks up to the dynamic boundary (context_fit) — is byte-identical
for an owner turn and for a wake at Act and at Observe built from one snapshot."""
from ouroboros.context import build_llm_messages
from tests.test_cache_optimization import _make_env_and_memory
monkeypatch.setenv("OUROBOROS_RUNTIME_MODE", "advanced")
env, memory = _make_env_and_memory(tmp_path)
owner = {"id": "t-owner", "type": "task", "text": "hi", "_is_direct_chat": True, "chat_id": 1,
"metadata": {"client_message_id": "cm-1"}}
owner_msgs, _ = build_llm_messages(env=env, memory=memory, task=owner)
owner_prefix = [json.dumps(owner_msgs[0]["content"][i], sort_keys=True) for i in (0, 1)]
owner_reg = _registry(tmp_path, owner["metadata"], task_id="t-owner")
owner_tools = json.dumps(owner_reg.schemas(), sort_keys=True)
for level in ("act", "observe"):
task = _wake_task(level)
task.update(id="t-owner", text="hi") # the same turn: the wake differs only in its metadata
msgs, _ = build_llm_messages(env=env, memory=memory, task=task)
prefix = [json.dumps(msgs[0]["content"][i], sort_keys=True) for i in (0, 1)]
assert prefix == owner_prefix, level
assert "cache_control" not in msgs[0]["content"][2]
reg = _registry(tmp_path, task["metadata"], task_id="t-owner")
assert json.dumps(reg.schemas(), sort_keys=True) == owner_tools, level
assert reg.capability_omissions() == owner_reg.capability_omissions(), level
# --- the per-task mode cap (В21=A): level x install mode --------------------------
_BLOCKED_CALLS = (
("write_file", {"path": "README.md", "content": "changed\n"}),
("commit_reviewed", {"commit_message": "test"}),
("run_command", {"cmd": "touch x.py"}),
("start_service", {"cmd": ["sleep", "5"], "name": "svc"}),
)
@pytest.mark.parametrize("mode", ["light", "advanced", "pro", "cyber_pro"])
@pytest.mark.parametrize("level", ["act", "observe"])
@pytest.mark.parametrize(("tool_name", "args"), _BLOCKED_CALLS)
def test_act_and_observe_cannot_touch_the_repo_in_any_install_mode(tmp_path, monkeypatch, mode, level, tool_name, args):
monkeypatch.setenv("OUROBOROS_RUNTIME_MODE", mode)
reg = _registry(tmp_path, _wake_task(level)["metadata"])
result = reg.execute(tool_name, dict(args))
expected = "RESOURCE_CONSTRAINT_BLOCKED" if level == "observe" and tool_name != "commit_reviewed" else "LIGHT_MODE_BLOCKED"
assert expected in result, (mode, level, tool_name, result[:300])
assert not (tmp_path / "repo" / "x.py").exists()
assert (tmp_path / "repo" / "README.md").read_text(encoding="utf-8") == "ok\n"
@pytest.mark.parametrize("mode", ["advanced", "pro", "cyber_pro"])
def test_full_follows_the_install_mode(tmp_path, monkeypatch, mode):
monkeypatch.setenv("OUROBOROS_RUNTIME_MODE", mode)
reg = _registry(tmp_path, _wake_task("full")["metadata"])
assert "LIGHT_MODE_BLOCKED" not in reg.execute("write_file", {"path": "scratch.txt", "content": "changed\n"})
assert (tmp_path / "repo" / "scratch.txt").read_text(encoding="utf-8") == "changed\n"
assert "LIGHT_MODE_BLOCKED" not in reg.execute("run_command", {"cmd": "touch x.py"})
def test_full_in_a_light_install_is_still_light(tmp_path, monkeypatch):
monkeypatch.setenv("OUROBOROS_RUNTIME_MODE", "light")
reg = _registry(tmp_path, _wake_task("full")["metadata"])
assert "LIGHT_MODE_BLOCKED" in reg.execute("write_file", {"path": "README.md", "content": "x"})
def test_act_keeps_the_light_positive_paths(tmp_path, monkeypatch):
monkeypatch.setenv("OUROBOROS_RUNTIME_MODE", "pro")
reg = _registry(tmp_path, _wake_task("act")["metadata"])
for root in ("task_drive", "artifact_store"):
result = reg.execute("write_file", {"root": root, "path": "notes.txt", "content": "kept\n"})
assert "BLOCKED" not in result, (root, result[:300])
assert "BLOCKED" not in reg.execute("read_file", {"path": "README.md"})
# --- the wake through the real lane --------------------------------------------
def test_the_lane_attaches_the_level_to_the_wake_contract(monkeypatch, tmp_path):
import queue
import threading
from ouroboros import agent as agent_module
from supervisor import workers
from tests.test_consciousness_wake_lane import _lane, _wait_for
_lane(monkeypatch, tmp_path, event_q=queue.Queue())
seen: list = []
done = threading.Event()
class Actor:
def handle_task(self, task):
seen.append(task)
done.set()
return []
monkeypatch.setattr(agent_module, "make_agent", lambda **kw: Actor())
receipt = workers.handle_wake_direct(1, "wake", {**WAKE_META, "consciousness_autonomy": "act"})
assert receipt["admitted"] is True
assert done.wait(10) and _wait_for(lambda: bool(seen))
task = seen[0]
assert task["task_contract"]["disabled_tools"] == list(ca.ACT_DISABLED)
assert task["metadata"]["runtime_mode_cap"] == "light"
# --- ISSUER: a wake speaks as a task -------------------------------------------
def test_routing_issuer_treats_a_wake_as_a_task_but_keeps_real_owner_relays(tmp_path):
from ouroboros.tools.control_routing import ISSUER_OWNER_TURN, ISSUER_TASK, _routing_issuer
wake = types.SimpleNamespace(task_id="wake-1", is_direct_chat=True, last_owner_delivery=None,
task_metadata=dict(_wake_task("act")["metadata"]))
assert _routing_issuer(wake) == {"kind": ISSUER_TASK, "task_id": "wake-1", "root_task_id": "wake-1"}
owner = types.SimpleNamespace(task_id="turn-1", is_direct_chat=True, last_owner_delivery=None,
task_metadata={"client_message_id": "cm-1"})
assert _routing_issuer(owner) == {"kind": ISSUER_OWNER_TURN}
# The two other triggers stay: a consciousness root relaying a REAL owner message.
relaying = types.SimpleNamespace(task_id="c-root", is_direct_chat=False,
last_owner_delivery={"client_message_id": "cm-9", "text": "go"},
task_metadata={"initiator": "consciousness"})
assert _routing_issuer(relaying) == {"kind": ISSUER_OWNER_TURN}
stamped = types.SimpleNamespace(task_id="c-root", is_direct_chat=True, last_owner_delivery=None,
task_metadata={"initiator": "consciousness", "client_message_id": "cm-2"})
assert _routing_issuer(stamped) == {"kind": ISSUER_OWNER_TURN}
def test_steer_from_a_wake_is_written_as_an_independent_task_message(tmp_path, monkeypatch):
from ouroboros.tools import control_routing
sent: list = []
monkeypatch.setattr(control_routing, "_send_task_message",
lambda ctx, issuer, target, msg, chat_id: sent.append((issuer, target, msg)) or "WRITTEN")
ctx = types.SimpleNamespace(
pending_events=[], event_queue=None, current_chat_id=1, drive_root=tmp_path,
task_id="wake-1", is_direct_chat=True, last_owner_delivery=None,
task_metadata=dict(_wake_task("act")["metadata"]),
)
assert control_routing._steer_task(ctx, task_id="r-1", message="please also check X") == "WRITTEN"
assert sent == [({"kind": "task", "task_id": "wake-1", "root_task_id": "wake-1"}, "r-1", "please also check X")]
assert ctx.pending_events == []
# --- origin inheritance: promote / followup / subagent -------------------------
@pytest.fixture
def _promote_root(tmp_path, monkeypatch):
"""The real promote admission path (tool -> supervisor handler -> worker_promotion)."""
import ouroboros.config as cfg
import supervisor.message_bus as mb
import supervisor.queue as queue_mod
from supervisor import workers
monkeypatch.setattr(cfg, "DATA_DIR", tmp_path)
monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path)
monkeypatch.setattr(queue_mod, "DRIVE_ROOT", str(tmp_path))
monkeypatch.setattr(queue_mod, "ACCEPTANCE_FENCES", {})
monkeypatch.setattr(mb, "get_bridge", lambda: types.SimpleNamespace(broadcast=lambda payload: None))
monkeypatch.setattr(workers, "_announce_created_project", lambda *a, **kw: None)
(tmp_path / "logs").mkdir(parents=True, exist_ok=True)
return tmp_path
def test_promote_from_a_wake_mints_a_consciousness_root_through_the_real_admission(_promote_root):
from ouroboros.tools.control_routing import _promote_chat_to_task
from ouroboros.utils import append_jsonl
from supervisor.events_project_routing import _handle_promote_chat_to_task
tmp_path = _promote_root
enqueued: list = []
captured: dict = {}
supervisor = types.SimpleNamespace(
DRIVE_ROOT=tmp_path, RUNNING={}, PENDING=[], WORKERS={0: types.SimpleNamespace()},
bridge=types.SimpleNamespace(send_routing_ack=lambda *a, **k: None, broadcast=lambda *a, **k: None),
enqueue_task=lambda task: enqueued.append(task) or dict(task),
persist_queue_snapshot=lambda **_k: True, load_state=lambda: {"owner_chat_id": 1},
append_jsonl=append_jsonl,
)
ctx = types.SimpleNamespace(
pending_events=[], current_chat_id=1, drive_root=tmp_path, budget_drive_root=str(tmp_path),
task_id="wake-1", is_direct_chat=True, last_owner_delivery=None, project_id="",
task_metadata=dict(_wake_task("act")["metadata"]), task_contract={},
event_queue=types.SimpleNamespace(
put_nowait=lambda event: (captured.update(event), _handle_promote_chat_to_task(event, supervisor))),
)
out = _promote_chat_to_task(ctx, "audit the logs", workspace="none", predecessor_task_id="")
assert out.startswith("OK: task"), out
# The event carries the origin by value and nothing the wake asked for is stripped.
assert captured["initiator"] == "consciousness" and captured["consciousness_autonomy"] == "act"
assert captured["usage_category"] == "consciousness_task" and captured["workspace"] == "none"
assert "presence" not in captured
[root] = enqueued
assert root["actor_id"] == "consciousness" and root["delegation_role"] == "root"
assert root["metadata"]["initiator"] == "consciousness"
assert root["metadata"]["usage_category"] == "consciousness_task"
assert root["task_contract"]["disabled_tools"] == list(ca.ACT_DISABLED)
assert root["metadata"]["runtime_mode_cap"] == "light" and "_presence_origin" not in root
def test_promoted_root_is_stamped_and_its_contract_derives_the_level(tmp_path, monkeypatch):
import supervisor.workers as workers
monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path)
enqueued: list = []
def enqueue(task):
enqueued.append(task)
return dict(task)
ctx = types.SimpleNamespace(enqueue_task=enqueue, persist_queue_snapshot=lambda **_k: True,
load_state=lambda: {"owner_chat_id": 1})
evt = {"type": "promote_chat_to_task", "task_id": "c0000001", "objective": "audit the logs",
"chat_id": 1, "workspace": "none", "initiator": "consciousness",
"usage_category": "consciousness_task", "consciousness_autonomy": "act"}
assert workers.promote_chat_to_task(evt, ctx)["status"] == "scheduled"
task = enqueued[0]
assert task["actor_id"] == "consciousness" and task["delegation_role"] == "root"
assert task["metadata"]["initiator"] == "consciousness"
assert task["metadata"]["usage_category"] == "consciousness_task"
assert task["metadata"]["consciousness_autonomy"] == "act"
assert task["task_contract"]["disabled_tools"] == list(ca.ACT_DISABLED)
assert task["metadata"]["runtime_mode_cap"] == "light"
assert task["source"] == "promote_chat_to_task" and "_presence_origin" not in task
def test_followup_template_inherits_the_origin_and_admission_derives_the_level(tmp_path, monkeypatch):
from supervisor import queue
from tests.test_schedule_followup import _ctx, _followup
ctx = _ctx(tmp_path)
ctx.task_metadata.update(_wake_task("act")["metadata"])
assert _followup(ctx).startswith("FOLLOWUP_SCHEDULED")
record = queue.list_scheduled_tasks(tmp_path / "data")["tasks"][0]
meta = record["task"]["metadata"]
assert meta["initiator"] == "consciousness" and meta["usage_category"] == "consciousness_task"
assert meta["consciousness_autonomy"] == "act" and "task_contract" not in record["task"]
monkeypatch.setattr(queue, "load_state", lambda: {"owner_chat_id": 1})
task = queue._task_from_schedule(record)
assert task["delegation_role"] == "root" and task["metadata"]["initiator"] == "consciousness"
assert task["task_contract"]["disabled_tools"] == list(ca.ACT_DISABLED)
assert task["metadata"]["runtime_mode_cap"] == "light"
def test_owner_followup_template_carries_no_origin(tmp_path):
from supervisor import queue
from tests.test_schedule_followup import _ctx, _followup
assert _followup(_ctx(tmp_path)).startswith("FOLLOWUP_SCHEDULED")
meta = queue.list_scheduled_tasks(tmp_path / "data")["tasks"][0]["task"]["metadata"]
assert "initiator" not in meta and "consciousness_autonomy" not in meta
def test_subagent_payload_lands_the_origin_on_the_child_metadata():
from supervisor.task_dispatch import build_scheduled_task_payload
fields = {"tid": "kid1", "chat_id": 1, "text": "x", "desc": "x", "role": "researcher",
"root_task_id": "wake-1", "delegation_role": "subagent", "actor_id": "subagent:researcher",
"origin_metadata": ca.consciousness_origin_metadata(_wake_task("act")["metadata"])}
task = build_scheduled_task_payload(fields)
assert task["metadata"]["initiator"] == "consciousness"
assert task["metadata"]["usage_category"] == "consciousness_task"
assert task["metadata"]["consciousness_autonomy"] == "act"
plain = build_scheduled_task_payload({**fields, "origin_metadata": {}})
assert "initiator" not in plain["metadata"]
def test_schedule_subagent_event_names_the_origin():
"""The tool stamps ``origin_metadata`` on the schedule event beside the envelope."""
source = pathlib.Path("ouroboros/tools/control_scheduling.py").read_text(encoding="utf-8")
assert '"origin_metadata": consciousness_origin_metadata(metadata),' in source
handler = pathlib.Path("supervisor/events_schedule_task.py").read_text(encoding="utf-8")
assert '"origin_metadata": evt.get("origin_metadata"),' in handler
# --- evolution: eligibility, sticky owner stop, campaign provenance ------------
def test_post_task_promotion_is_refused_when_toggle_evolution_is_withheld():
from ouroboros.post_task_evolution import _eligible
assert _eligible({"type": "task"}) is True
assert _eligible({"type": "task", "task_contract": {"disabled_tools": ["toggle_evolution"]}}) is False
assert _eligible({"type": "task", "metadata": {"disabled_tools": ["toggle_evolution"]}}) is False
from ouroboros.contracts.task_contract import attach_task_contract
assert _eligible(attach_task_contract(_wake_task("act"))) is False
assert _eligible(attach_task_contract(_wake_task("full"))) is True
def test_globalized_promotion_view_keeps_the_contract(tmp_path, monkeypatch):
from ouroboros import agent_task_pipeline as pipeline
from ouroboros.contracts.task_contract import attach_task_contract
seen: list = []
monkeypatch.setattr(pipeline, "_update_improvement_backlog", lambda env, entry: None)
monkeypatch.setattr("ouroboros.post_task_evolution.maybe_promote",
lambda env, task, entry, llm: seen.append(task))
task = attach_task_contract({**_wake_task("act"), "project_id": "lab"})
pipeline._run_global_backlog_promotion_only(
types.SimpleNamespace(drive_root=tmp_path), task,
{"backlog_candidates": [{"summary": "tidy the logs"}]}, None,
)
assert seen and seen[0]["task_contract"]["disabled_tools"] == list(ca.ACT_DISABLED)
def test_request_file_and_pending_apply_carry_the_origin(tmp_path, monkeypatch):
from ouroboros import post_task_evolution as pte
task = _wake_task("full")
pte._write_request(tmp_path, {"objective": "improve X", "requires_plan_review": False}, task)
req = json.loads((tmp_path / pte._REQUEST_REL).read_text(encoding="utf-8"))
assert req["initiator"] == "consciousness" and req["consciousness_autonomy"] == "full"
assert req["usage_category"] == "consciousness_task"
calls: list = []
monkeypatch.setattr("ouroboros.config.get_post_task_evolution_enabled", lambda: True)
monkeypatch.setattr("supervisor.evolution_lifecycle.evolution_block_reason", lambda: "")
monkeypatch.setattr("supervisor.evolution_lifecycle.start_evolution_campaign",
lambda objective, source="", **kw: calls.append((objective, source, kw)) or {"id": "c1"})
monkeypatch.setattr("supervisor.state.load_state", lambda: {"owner_chat_id": 7})
def _update_state(mutator):
live: dict = {}
mutator(live)
return live
monkeypatch.setattr("supervisor.state.update_state", _update_state)
monkeypatch.setattr("ouroboros.config.get_post_task_evolution_budget_usd", lambda: 0.0)
assert pte.apply_pending_request(tmp_path) is True
assert calls == [("improve X", "post_task", {"origin": {
"initiator": "consciousness", "usage_category": "consciousness_task", "consciousness_autonomy": "full"}})]
def _toggle_ctx(state, sent):
return types.SimpleNamespace(load_state=state.load_state,
send_with_budget=lambda cid, text, **kw: sent.append(text))
def test_agent_tool_enable_is_refused_while_the_owner_stop_stands(tmp_path, monkeypatch):
"""В12: /evolve off is sticky against toggle_evolution — the typed refusal, no campaign."""
import supervisor.state as state
from supervisor import events as events_mod
from supervisor import evolution_lifecycle as el
state.init(tmp_path)
state.update_state(lambda live: live.update(owner_chat_id=7, evolution_owner_stopped=True))
started: list = []
monkeypatch.setattr(el, "evolution_block_reason", lambda: "")
monkeypatch.setattr(el, "start_evolution_campaign",
lambda objective, source="", **kw: started.append(source) or {"status": "active"})
sent: list = []
events_mod._handle_toggle_evolution({"enabled": True, "objective": "x"}, _toggle_ctx(state, sent))
assert started == []
assert bool(state.load_state().get("evolution_owner_stopped")) is True
assert not state.load_state().get("evolution_mode_enabled")
assert sent and "stayed OFF" in sent[0] and "sticky" in sent[0]
def test_agent_tool_enable_without_an_owner_stop_starts_a_campaign_with_the_origin(tmp_path, monkeypatch):
import supervisor.state as state
from supervisor import events as events_mod
from supervisor import evolution_lifecycle as el
state.init(tmp_path)
state.update_state(lambda live: live.update(owner_chat_id=7, evolution_owner_stopped=False))
started: list = []
monkeypatch.setattr(el, "evolution_block_reason", lambda: "")
monkeypatch.setattr(el, "start_evolution_campaign",
lambda objective, source="", **kw: started.append((source, kw)) or {"status": "active"})
sent: list = []
evt = {"enabled": True, "objective": "x", "initiator": "consciousness",
"usage_category": "consciousness_task", "consciousness_autonomy": "full"}
events_mod._handle_toggle_evolution(evt, _toggle_ctx(state, sent))
assert started == [("agent_tool", {"origin": {
"initiator": "consciousness", "usage_category": "consciousness_task", "consciousness_autonomy": "full"}})]
live = state.load_state()
assert live["evolution_mode_enabled"] is True and live["evolution_owner_stopped"] is False
def test_toggle_tool_stamps_the_turn_origin_on_its_event(monkeypatch):
from ouroboros.tools.control_runtime import _toggle_evolution
monkeypatch.setattr("supervisor.evolution_lifecycle.evolution_block_reason", lambda: "")
ctx = types.SimpleNamespace(pending_events=[], task_metadata=dict(_wake_task("full")["metadata"]))
assert _toggle_evolution(ctx, True, "improve X").startswith("OK")
evt = ctx.pending_events[0]
assert evt["type"] == "toggle_evolution" and evt["initiator"] == "consciousness"
assert evt["consciousness_autonomy"] == "full"
owner = types.SimpleNamespace(pending_events=[], task_metadata={"client_message_id": "cm"})
_toggle_evolution(owner, True, "improve X")
assert "initiator" not in owner.pending_events[0]
def test_campaign_keeps_the_origin_and_its_cycle_tasks_inherit_it(tmp_path, monkeypatch):
from supervisor import evolution_lifecycle, queue, state
state.init(tmp_path)
queue.init(tmp_path)
pending: list = []
queue.init_queue_refs(pending, {}, {"value": 0})
monkeypatch.setattr(state, "TOTAL_BUDGET_LIMIT", 0.0)
origin = {"initiator": "consciousness", "usage_category": "consciousness_task", "consciousness_autonomy": "full"}
campaign = evolution_lifecycle.start_evolution_campaign("Improve", source="agent_tool", origin=origin)
assert campaign["initiator"] == "consciousness" and campaign["consciousness_autonomy"] == "full"
# A resume keeps the recorded origin (the owner's later resume does not erase it).
campaign["status"] = "paused"
assert evolution_lifecycle._write_evolution_campaign(campaign) is True
resumed = evolution_lifecycle.start_evolution_campaign("", source="owner_chat")
assert resumed["initiator"] == "consciousness"
state.update_state(lambda live: live.update(owner_chat_id=1, evolution_mode_enabled=True,
evolution_owner_stopped=False))
monkeypatch.setattr(evolution_lifecycle, "evolution_block_reason", lambda: "")
monkeypatch.setattr(queue, "send_with_budget", lambda *a, **k: None)
monkeypatch.setattr(queue, "persist_queue_snapshot", lambda reason="": None)
monkeypatch.setattr("ouroboros.consciousness_allowance.allowance_window",
lambda root, now=None: {"status": "available", "limit_usd": 20.0, "accounted_usd": 0.0,
"remaining_usd": 20.0, "unknown_unmetered": 0, "resets_at": ""})
queue.enqueue_evolution_task_if_needed()
assert len(pending) == 1
task = pending[0]
assert task["type"] == "evolution" and task["metadata"]["initiator"] == "consciousness"
assert task["metadata"]["usage_category"] == "consciousness_task"
assert task["task_contract"]["disabled_tools"] == [] and task["metadata"]["runtime_mode_cap"] == ""
def test_owner_campaign_carries_no_origin(tmp_path):
from supervisor import evolution_lifecycle, queue, state
state.init(tmp_path)
queue.init(tmp_path)
campaign = evolution_lifecycle.start_evolution_campaign("Improve", source="owner_chat")
assert "initiator" not in campaign
assert ca.consciousness_origin_metadata(campaign) == {}

View file

@ -415,7 +415,10 @@ def test_gr4_6_toggle_evolution_clears_owner_stop_before_the_campaign_mint(
load_state=state.load_state, send_with_budget=lambda *a, **kw: None,
)
events_mod._handle_toggle_evolution({"enabled": True, "objective": "x"}, ctx)
# An OWNER start (В12: only owner provenance clears the sticky stop).
events_mod._handle_toggle_evolution(
{"enabled": True, "objective": "x", "source": "owner_chat"}, ctx,
)
assert seen == [False], (
"GR4-6: the owner-stop flag is cleared BEFORE the campaign is minted — "

View file

@ -130,7 +130,11 @@ def test_gr5_1_toggle_start_failure_restores_a_prior_owner_stop(tmp_path, monkey
)
sent: list = []
events_mod._handle_toggle_evolution({"enabled": True, "objective": "x"}, _toggle_ctx(state, sent))
# An OWNER start reaches the mint (В12: the agent tool is refused outright
# while the flag stands — pinned in test_post_task_evolution).
events_mod._handle_toggle_evolution(
{"enabled": True, "objective": "x", "source": "owner_chat"}, _toggle_ctx(state, sent),
)
assert bool(state.load_state().get("evolution_owner_stopped")) is True, (
"GR5-1: a failed start must restore the captured owner-stop flag"