ouroboros: checkpoint after task 4b8cb25784ee4f17 — ТЗ2 — завершить самокоординацию Presence

This commit is contained in:
Ouroboros 2026-09-25 04:24:24 +03:00
parent 2f7a3c9e45
commit 4d149ff677
34 changed files with 1787 additions and 95 deletions

View file

@ -734,8 +734,9 @@ ceiling and reply context rather than widening authority.
Transport custody preserves provider arrival order before Host admission; the
host serializes one conversation and enforces the installation-wide active-turn
limit across processes. A current Presence turn may cancel only its own
binding-and-conversation-correlated `work_ref`. Owner chat or Background
limit across processes. A Presence turn may find, read, message and cancel only
independent work started from its own binding, in any of that binding's
conversations; other bindings and the owner's tasks stay out of reach. Owner chat or Background
Consciousness may initiate an existing binding, but the resulting cycle must use
an explicitly selected transport tool and finish `tool_delivered` to claim that
an external message was sent.
@ -765,6 +766,8 @@ status notices stay in the owner task; an empty deferred body sends nothing but
still requires polling. Cached and late results preserve that empty body rather
than substituting the task diagnostic. This does not turn failure into success.
Ordinary implicit replies and genuine authored best-effort answers remain valid.
A forced final separates the task record from the reply: only a `presence_finish`
declared in that answer is spoken, so an undeclared record sends nothing new.
`GET /identity` advertises `presence_delivery_version: 1` on supporting hosts.
Only then request `delivery_reporting_version: 1` alongside `binding_id` and

View file

@ -16,7 +16,7 @@ Between the sends of one loop execution the transcript is append-only — each s
Host plan/orphan disclosures ride beside the model answer as `terminal_host_notice`; delivery sends the model answer alone and the disclosure stays a field of the result (an owed pre-upgrade notice row still replays as the untyped System row it was), and the custody audit is its own typed row (`terminal_custody_notice` on the send event → `system_type="custody_notice"`, `card_row="timeline"`, its own owed delivery id), so the open-delegation fact is a row of the task's card and the unreconciled-runs note never enters the assistant text. CLI gets one host-labelled status (`terminal_host_notice_text`); external Presence speech excludes host notices (chapter 12). Stored answer bytes do not change: the answer keeps its hash and its real PASS or FAIL when only the notice changes, and equal text cannot revive a superseded verdict. Every reader that hands a result on — synthesis, parent handoff, `get_task_result`/`wait_task`/`wait_tasks`, the `task:<id>` plan-evidence reader — carries the notice as a separate field, and the child-result and plan-evidence hashes include it, so a changed child limitation invalidates an old parent disposition or an old review. A host-salvaged terminal is labelled, not hidden: the durable row carries `Preserved intermediate output (not a final answer):` plus the bounded excerpt every terminal row uses, the untruncated copy stays with `get_task_result` and the stop receipt, and only a row written to the chat the stop receipt itself reached (`cancel_receipt.delivered_chat_id`, recorded after a successful send) reduces to its label beside the pointer; the two durable writes are not atomic, and no writer trades its only pointer for the label. One event is disclosed once, at the layer that owns it: the forced orphan note leaves a child to its own terminal row only where that row reached the reader the note addresses, and a rail that ended a routing turn is named by that turn's own row, never inferred from another layer's stamp.
The delivery-control protocol is resolved here and only here (DEVELOPMENT keeps the rule and points here). The candidate carries sticky loop-local provenance that its lineage has seen a host-issued delivery-control episode; without one, exact JSON is ordinary text. In a marked lineage both the ordinary and the forced resolver intercept recognizable whole-body envelopes and balanced trailing protocol attempts — valid `keep` resolves to the retained candidate, valid `replace` to `full_answer`, anything malformed preserves the retained candidate. Both strip one whole-body fence and treat a balanced protocol object at the very END of prose as a protocol attempt (`utils.extract_trailing_json_object` + `loop_delivery._parse_delivery_control_body`); the trailing-object rule deliberately refuses substring scanning, because quoted protocol literals mid-prose are legitimate text, so a control object quoted mid-prose stays prose and a truncated trailing fragment remains prose. During an ordinary acceptance continuation, `finalization_control="acceptance_feedback"` makes a complete revised answer ordinary prose and keeps `keep`/`replace`/`pending_review` optional. Prose resets the pending-review choice to `wait`; it never means `finish`. An outstanding effect, owner-revision or child-action control retains its stricter rule, and unread owner source still requires acknowledgement. Empty or recognizable malformed control bodies retain the candidate rather than becoming its replacement. The ordinary resolver takes one repair round then degraded-preserve; the forced resolver resolves purely and never re-loops — malformation preserves the retained candidate with the typed `delivery_control_degraded` reason, which is this forced rail's own code, while the ordinary repair path records `invalid_delivery_control_after_repair`. Every degradation carries the cause it computed: `outcomes.derive_loop_outcome` falls back to `delivery_control_degraded` only for a degradation that reports no cause and publishes the `loop_outcome.degraded`/`degraded_reason` pair the benchmark ledgers read. A malformed attempt in a marked lineage never leaks JSON, even after the transient latch clears.
The delivery-control protocol is resolved here and only here (DEVELOPMENT keeps the rule and points here). The candidate carries sticky loop-local provenance that its lineage has seen a host-issued delivery-control episode; without one, exact JSON is ordinary text. In a marked lineage both the ordinary and the forced resolver intercept recognizable whole-body envelopes and balanced trailing protocol attempts — valid `keep` resolves to the retained candidate, valid `replace` to `full_answer`, anything malformed preserves the retained candidate. Both strip one whole-body fence and treat a balanced protocol object at the very END of prose as a protocol attempt (`utils.extract_trailing_json_object` + `loop_delivery._parse_delivery_control_body`); the trailing-object rule deliberately refuses substring scanning, because quoted protocol literals mid-prose are legitimate text, so a control object quoted mid-prose stays prose and a truncated trailing fragment remains prose. During an ordinary acceptance continuation, `finalization_control="acceptance_feedback"` makes a complete revised answer ordinary prose and keeps `keep`/`replace`/`pending_review` optional. Prose resets the pending-review choice to `wait`; it never means `finish`. An outstanding effect, owner-revision or child-action control retains its stricter rule, and unread owner source still requires acknowledgement. Empty or recognizable malformed control bodies retain the candidate rather than becoming its replacement. The ordinary resolver takes one repair round then degraded-preserve; the forced resolver resolves purely and never re-loops — malformation preserves the retained candidate with the typed `delivery_control_degraded` reason, which is this forced rail's own code, while the ordinary repair path records `invalid_delivery_control_after_repair`. Every degradation carries the cause it computed: `outcomes.derive_loop_outcome` falls back to `delivery_control_degraded` only for a degradation that reports no cause and publishes the `loop_outcome.degraded`/`degraded_reason` pair the benchmark ledgers read. A malformed attempt in a marked lineage never leaks JSON, even after the transient latch clears. A Presence turn or root's forced call (not a ceiling-only child) is always armed; its nested `presence_finish` is an admitted envelope key read from the original bytes, so duplicates still refuse (chapter 12).
The child-absorption gate is an action gate: while undispositioned direct children remain, the loop HOLDS the candidate (`child_absorption_or_revision_required`) instead of arming the JSON-only control instruction — the hold-vs-arm split exists so the model never receives two contradictory instructions in one round. A typed keep cannot close the gate; after the one bounded reminder it forces the best-effort `children_unabsorbed` rail with a current `id [status] sha256` listing. The absorption digest's `## child` header carries the child's typed custody debt (`delegated_runs_unreconciled`, bounded, with a `get_task_result` pointer) as visibility only — the parent's authority over that patch is exactly the orphan rule, and a child's debt never relabels the root card (DESIGN §4). Finalizing over an UNDISPOSED OWN delegated patch is deliberately NOT gated: the consequence is disclosed where the decision is made (the `integrate_delegated_patch` schema, the apply receipts) and lands as the additive Done-with-warnings custody overlay rather than a hold; a pre-finalization reminder and propagation of child custody debt into the acceptance-subtree snapshot remain disclosed deferred gaps.
@ -68,7 +68,7 @@ A workspace task's completion compares against the captured preflight base — t
`promote_chat_to_task`, `route_to_project`, `steer_task` and `ensure_project_scope` ride one receipt rail: an act succeeds only once its token-matched supervisor facts are durable in the task result, queue snapshot, annotation or mailbox authority; among several possible tasks the LLM chooses, code auto-delivers only the unambiguous one-target case, and an unconfirmed or stale receipt fails visibly instead of launching a second root. Receipts are retained per `(owner message, routing token)`: an earlier act's receipt stays readable by its token (`chat_annotation_receipt`) while the message's latest row is the UI projection and the picker's liveness test. A KNOWN rejection returns `rejected` with its reason, never a timeout, so `UNCONFIRMED` keeps its one meaning: no matching receipt exists. A receipt proves admission, not completion; an unread indicator proves a visible revision, not memory isolation.
WHO is speaking is ONE host-minted fact on the event (`control_routing._routing_issuer`; the model has no argument): an OWNER TURN (a direct turn the owner door stamped: `is_direct_chat` and `run_origin.owner_ingress`) or a TASK speaking for itself (a promoted root inherits the stamp as ancestry, not as issuer; a root relaying an owner message it just drained; a consciousness wake-up, a Presence event or the auto-resume template on the direct lane, which nobody typed — a wake's promotes mint consciousness roots inheriting its origin, ledger category and autonomy level). An owner turn's steer travels as owner text: `[Message from my human]`, the owner corpus, the generation bump that supersedes a reviewed answer, the room veto from the registry lane of the issuing chat (a Project room reaches its own roots, Main every host-listed root) and the owner acknowledgement. A task's own words NEVER travel as owner text: they go through the one task-message writer `forward_to_worker` also uses, as `independent_task` provenance, to any host-listed active independent root (hidden roots included, no room veto), render as `[Message from independent task <id>]`, enter no owner corpus (`owner_source_sha256` and the acceptance premises stay the owner's), carry no attachments, and are confirmed WRITTEN or refused with the host's reason. A relay keys and publishes its acknowledgement on the owner message it drained; other task acts use their own synthetic receipt id without a chat acknowledgement. Neither receipt grants owner authority; author and target ride the receipt and one `task_message_routed` Logs row, and the receiver's `task_message_injected` row names the sender. Independent roots learn the roster from a `[INDEPENDENT_ROOTS]` TAIL note (`peer_roster.py`, `ROSTER_NOTE_CAP` = 40 rows shown, the cut disclosed), appended only when the roster changed and never merged into a sent row. Rows identify live direct conversations without claiming an owner initiator; messaging one is the model's call. A root may publish one bounded `update_focus(text, source_ref)` record, exposed in grouped notes and paginated `live_roots` and retained by `recent_tasks` on dormant results. A `source_ref` names a reader; `update_focus` answers it through that reader under the caller's registry admission (`disabled_tools` included) and retains the exact answer (≤256 KiB) write-once on the canonical root as `focus.source_handle`, a native `task_source` ref peers read from any drive via `get_task_result(include_focus_source=True, focus_source_sha256=…)` against the PHYSICAL author's record: the digest selects the immutable historical file, so no later focus or retry can substitute its evidence (the roster's `retained_source`); a refusing reader or changed snapshot refuses the focus (`FOCUS_SOURCE_UNRESOLVED`), and the roster drops a settled root's focus. Focus has authored time separate from host observation and carries no TTL or owner authority. Explicit authorized project journal/workpad reads return exactly the requested source; foreign scoped writes refuse; children/Presence keep their capability ceiling.
WHO is speaking is ONE host-minted fact on the event (`control_routing._routing_issuer`; the model has no argument): an OWNER TURN (a direct turn the owner door stamped: `is_direct_chat` and `run_origin.owner_ingress`) or a TASK speaking for itself (a promoted root inherits the stamp as ancestry, not as issuer; a root relaying an owner message it just drained; a consciousness wake-up, a Presence event or the auto-resume template on the direct lane, which nobody typed — a wake's promotes mint consciousness roots inheriting its origin, ledger category and autonomy level). An owner turn's steer travels as owner text: `[Message from my human]`, the owner corpus, the generation bump that supersedes a reviewed answer, the room veto from the registry lane of the issuing chat (a Project room reaches its own roots, Main every host-listed root) and the owner acknowledgement. A task's own words NEVER travel as owner text: they go through the one task-message writer `forward_to_worker` also uses, as `independent_task` provenance, to any host-listed active independent root (hidden roots included, no room veto; a Presence sender only to its own binding's work, chapter 12), render as `[Message from independent task <id>]`, enter no owner corpus (`owner_source_sha256` and the acceptance premises stay the owner's), carry no attachments, and are confirmed WRITTEN or refused with the host's reason. A relay keys and publishes its acknowledgement on the owner message it drained; other task acts use their own synthetic receipt id without a chat acknowledgement. Neither receipt grants owner authority; author and target ride the receipt and one `task_message_routed` Logs row, and the receiver's `task_message_injected` row names the sender. Independent roots learn the roster from a `[INDEPENDENT_ROOTS]` TAIL note (`peer_roster.py`, `ROSTER_NOTE_CAP` = 40 rows shown, the cut disclosed), appended only when the roster changed and never merged into a sent row. Rows identify live direct conversations without claiming an owner initiator; messaging one is the model's call. A root may publish one bounded `update_focus(text, source_ref)` record, exposed in grouped notes and paginated `live_roots` and retained by `recent_tasks` on dormant results. A `source_ref` names a reader; `update_focus` answers it through that reader under the caller's registry admission (`disabled_tools` included) and retains the exact answer (≤256 KiB) write-once on the canonical root as `focus.source_handle`, a native `task_source` ref peers read from any drive via `get_task_result(include_focus_source=True, focus_source_sha256=…)` against the PHYSICAL author's record: the digest selects the immutable historical file, so no later focus or retry can substitute its evidence (the roster's `retained_source`); a refusing reader or changed snapshot refuses the focus (`FOCUS_SOURCE_UNRESOLVED`), and the roster drops a settled root's focus. Focus has authored time separate from host observation and carries no TTL or owner authority. Explicit authorized project journal/workpad reads return exactly the requested source; foreign scoped writes refuse; children/Presence keep their capability ceiling.
`steer_task` relays the owner's exact ingress bytes only on the turn's FIRST routing act, while it still acts on the message that started it; the window ends with the latest owner message the turn DRAINED (`ToolContext.last_owner_delivery`, stamped at the loop's mailbox drain; a message only written to the mailbox ends nothing) or with a landed promote/route/steer receipt already on the origin message (a refused or unconfirmed act carried nothing, so the next act still relays). Past either, the turn RELAYS its own words and its receipt is keyed on the message relayed (the drained delivery's own client id, else the synthetic `agent-steer:<routing token>` id), never again on an origin this turn already routed. An agent-authored steer belonging to no owner message earns its receipt under that synthetic id, confirmable through the same `routing_wait` poll; no chat row carries the id and the owner message's own receipt (what a later decision turn reads) stands. Each steer's mailbox entry is keyed by its routing token, so several instructions under one origin are several deliveries while a retried emit of one steer stays one.

View file

@ -12,9 +12,9 @@ Admission is one limiter with two policies (`host_service._RateLimiter`): the WS
Operation correlation: a named injected message has `operation_ref=<chat_id>:<client_message_id>` on 202, 200, 504 and disconnect responses. `supervisor.message_bus.accept_local_message` serializes check, canonical inbound-row acceptance and enqueue, and the `log_chat` write must succeed before work is queued. Repeated same-id, same-text, same-skill delivery rejoins even before supervisor dequeue; changed content or source is refused with 409. The queue stays in-memory: a crash after acceptance can lose delivery and reads honestly as `lost`, never authorizing a second enqueue. Routing annotations and outbound task ids are discovery hints only — task reads and cancellation require the actual queue/task record's complete `origin_message_ref` to match the authenticated skill's canonical source (`DirectActivityRegistry` carries the same origin), and named response waits poll that exact operation and its retry-aware effective task result, because chat ordering alone never proves a reply. Cancel enters the durable intent and cascade-custody owner (§5) only for work with that origin and the same installation root; a different or unavailable owner root and unaddressable or foreign work are `cancel_unsupported` before any intent is written, and unresolved custody never becomes a false `cancelled`.
Presence (`presence_runner.py`): `POST /presence/turn` admits an exact provider/account/conversation/thread event and skill-state-confined files under the hash-bound `presence` permission and an owner binding (`state/presence_bindings.json`). Agents share one autobiography. Event-derived IDs deduplicate retries; cross-process locks cap concurrency and serialize conversations. Inbound, initiated and receipt paths derive one conversation key and chat id (`presence_bindings.conversation_key`, empty thread = `0`). A host-reconciled placeholder, or a still-running row of a dead attempt, is not a result: the same event id re-runs told what that attempt confirmed sending (unknown once its rows left the live chat generation), and the row shows `superseded_placeholder`. A presence-local liveness set shields an executing turn from orphan reconciliation without making it an owner-addressable direct actor, so update drains do not wait for it and after a restart the transport's retry re-runs it. Each conversation's last executed turn is a rebuildable pointer (`state/presence_turn_gate/last-<sha256>.json`; canonical sources: task results and receipts; a replay of a settled, authored turn that lost it rebuilds it) shown to the next turn. Turns, receipts and inject have separate per-skill in-flight budgets. Outcomes are message/silent/tool_delivered/deferred; deferred needs a correlated `work_ref`, read via `GET /presence/work/{work_ref}`, not the general task API; an orphan-reconciled child replays silent (nobody re-runs it). Promotion discards requested Project/workspace/source widening and copies the ceiling, cost and return context; unusable folders return `workspace_unusable` with repair detail. `presence_cancel_work` requires the current binding and conversation. Owner chat and consciousness at Act or above may `initiate_presence` on an enabled binding.
Presence (`presence_runner.py`): `POST /presence/turn` admits an exact provider/account/conversation/thread event and skill-state-confined files under the hash-bound `presence` permission and an owner binding (`state/presence_bindings.json`). Agents share one autobiography. Event-derived IDs deduplicate retries; cross-process locks cap concurrency and serialize conversations. Inbound, initiated and receipt paths derive one conversation key and chat id (`presence_bindings.conversation_key`, empty thread = `0`). A host-reconciled placeholder, or a still-running row of a dead attempt, is not a result: the same event id re-runs told what that attempt confirmed sending (unknown once its rows left the live chat generation), and the row shows `superseded_placeholder`. A presence-local liveness set shields an executing turn from orphan reconciliation without making it an owner-addressable direct actor, so update drains do not wait for it and after a restart the transport's retry re-runs it. Each conversation's last executed turn is a rebuildable pointer (`state/presence_turn_gate/last-<sha256>.json`; canonical sources: task results and receipts; a replay of a settled, authored turn that lost it rebuilds it) shown to the next turn. Turns, receipts and inject have separate per-skill in-flight budgets. Outcomes are message/silent/tool_delivered/deferred; deferred needs a correlated `work_ref`, read via `GET /presence/work/{work_ref}`, not the general task API; an orphan-reconciled child replays silent (nobody re-runs it). Promotion discards requested Project/workspace/source widening and copies the ceiling, cost and return context, and its scheduled row already carries the Presence provenance; unusable folders return `workspace_unusable` with repair detail. A binding's own work is an independent root carrying its nonempty binding id, from any of its conversations (`presence_related_work`); other bindings, owner roots, inline turns, children and rows without that provenance are never attributed. The context lists its first page (canonical root, read gaps stated); new ceilings add `recent_tasks`/`get_task_result` host-bound to `presence_scope=own_binding` plus `steer_task`, and steer, `presence_cancel_work` and a selected cancel or forward reach only that work (cancel/forward also the caller's own tree); a selected global reader stays global. Owner chat and consciousness at Act or above may `initiate_presence` on an enabled binding.
The ceiling includes knowledge, scratchpad, identity and chat history, not correspondent tool or owner-command authority. Unselected baseline `chat_history` is unbound; older ceilings keep their digest without gaining baseline tools. `presence_context.py` supplies route facts and frames incoming speech as observation, distinct from initiation or inherited work. Source text stays separate from host attachment declarations; optional provider identity/mention/root facts do not determine the addressee. The model chooses useful, concise participation or silence; early speech follows that choice through selected send tools, never Working forwarding. `presence_finish` enters common completion checks after the tool/control/budget tail; an omitted message/deferred body keeps an answer round. Revised/failed work discards stale completion text. Authored best-effort replies remain speech; host diagnostics stay owner-side. Host-only terminals retain scheduled child custody as deferred with an empty body. Live/cache/work readers share authorship and empty-body rules; unknown-origin legacy text speaks only for completed rows. Synthesis preserves the recorded adapter body, outcome/origin/work reference and internal result, not today's replay policy; preparation is no delivery receipt.
The ceiling includes knowledge, scratchpad, identity and chat history, not correspondent tool or owner-command authority. Unselected baseline `chat_history` is unbound; older ceilings keep their digest without gaining baseline tools. `presence_context.py` supplies route facts and frames incoming speech as observation, distinct from initiation or inherited work. Source text stays separate from host attachment declarations; optional provider identity/mention/root facts do not determine the addressee. The model chooses useful, concise participation or silence; early speech follows that choice through selected send tools, never Working forwarding. `presence_finish` enters common completion checks after the tool/control/budget tail; an omitted message/deferred body keeps an answer round. Revised/failed work discards stale completion text, and the next round is told so with the task's confirmed sends. Ordinary authored best-effort replies remain speech. A forced final is the internal record: only a valid `presence_finish` declared inside that one call speaks (a declared partial even on failure, never a `tool_delivered` note); anything else says nothing new and records `presence_declaration`. Host diagnostics stay owner-side. Host-only terminals retain scheduled child custody as deferred with an empty body. Live/cache/work readers share authorship and empty-body rules; unknown-origin legacy text speaks only for completed rows. Synthesis preserves the recorded adapter body, outcome/origin/work reference and internal result, not today's replay policy; preparation is no delivery receipt.
Receipt reporting is negotiated through `/identity` (`presence_delivery_version: 1`) and optional `delivery_reporting_version: 1` on a turn; mode survives cached/deferred results. Mode1 writes outgoing history on provider receipts; mode0 retains an authored, delivery-unconfirmed row. `presence_delivery.py` accepts authenticated `POST /presence/delivery` observations through the chat writer. Parts retain target, text/format and provider facts; SMTP `accepted` means provider acceptance only. Speech is typed `presence_delivery`, never a task-finalizing untyped row; failed/uncertain attempts are System facts, queued remains in tool/outbox receipts. Memory, history and consolidation retain destination/state/details. A Host-context projection rebuilds once from retained chat generations and updates after required writes; identical report retries deduplicate, changed facts conflict. Transport outboxes own provider receipts and separate report ACK/backoff: a slow or failed report neither resends nor blocks provider delivery. No second store or scheduler; old queues are not imported. Wire fields and procedure: CREATING_SKILLS, “Reporting actual Presence delivery”.

View file

@ -51,10 +51,10 @@ Enforcement: `tests/test_protected_artifacts_policy.py` and `tests/test_acceptan
### Skill-defined Presence
- Keep behavior portable and authority installation-local: a reviewed `presence:` profile declares instructions, context topics, bounded runtime defaults and conceptual tool/script/resource requests — never provider credentials, room ids or one installed tool spelling; `presence_capabilities.py` stores the owner's exact selections outside the payload, fingerprinted by the request semantics that authorize them. Preserve its optional `workspace_root` (an owner-local external folder, validated through the existing workspace admission and copied into each task contract) when editing runtime/capability selections; unset profiles retain their prior serialized state and fingerprint. Presence keeps canonical shared memory without deriving a Project or creating a forked drive from that folder (ARCHITECTURE §6 "Skills and extensions").
- Presence authority is a positive immutable ceiling, not a denylist or a prompt promise: admission requires the owner-created binding plus an installed, enabled, freshly executable behavior skill and every required selection, then freezes skill/profile/state/selection fingerprints, exact grants (the profile's selections plus the constant cognitive-memory baseline `tool_capabilities.COGNITIVE_MEMORY_TOOL_NAMES`; a selected grant keeps its bindings), argument bindings, runtime slot and round limit into `task_contract.capability_ceiling`. Schema discovery and execution enforce that same ceiling for built-ins, extensions, MCP tools, scripts and resource roots.
- Presence authority is a positive immutable ceiling, not a denylist or a prompt promise: admission requires the owner-created binding plus an installed, enabled, freshly executable behavior skill and every required selection, then freezes skill/profile/state/selection fingerprints, exact grants (the profile's selections plus the constant cognitive-memory baseline `tool_capabilities.COGNITIVE_MEMORY_TOOL_NAMES` and the own-work baseline — both readers host-bound to `presence_scope=own_binding`, `steer_task`; a selected grant keeps its bindings), argument bindings, runtime slot and round limit into `task_contract.capability_ceiling`. Schema discovery and execution enforce that same ceiling for built-ins, extensions, MCP tools, scripts and resource roots.
- `state/presence_bindings.json` is host-owned authority: a transport token resolves only bindings naming that exact transport skill, and the submitted provider/account/conversation/thread must match the binding origin — never recover those identities from message text. Staged files stay inside the calling skill's state root before entering the ordinary attachment store (the turn flow: ARCHITECTURE §12).
- Run each admitted event with a fresh agent, a deterministic binding-plus-source-event task id, the cross-process installation-wide concurrency gate and per-conversation serialization; the transport's durable provider custody owns arrival FIFO before Host admission. Do not add a transport-specific task scheduler, memory silo, core terminal outbox or resident cross-room agent.
- Completion is exactly `message`, `silent`, `tool_delivered` or `deferred` (deferred requires a successfully promoted `work_ref`; correlated lookup stays behind the same transport token and binding, and `presence_cancel_work` additionally requires the current binding and conversation to match). Promotion and `schedule_followup` copy the Presence metadata, admitted workspace and capability ceiling by value; any new descendant producer preserves this ceiling or refuses the transition — reconstructing authority from mutable current state is forbidden.
- Completion is exactly `message`, `silent`, `tool_delivered` or `deferred` (deferred requires a successfully promoted `work_ref`; correlated lookup stays behind the same transport token and binding). A Presence caller reads, steers and cancels only independent roots of its own nonempty binding (`presence_authority.presence_work_refusal`), and a forced final speaks only its nested `presence_finish` declaration. Promotion and `schedule_followup` copy the Presence metadata, admitted workspace and capability ceiling by value; any new descendant producer preserves this ceiling or refuses the transition — reconstructing authority from mutable current state is forbidden.
- Knowledge-topic and scratchpad mutation each use one stable lock, so concurrent owner and Presence turns cannot overwrite a newer projection with an older render. Test the boundary at both layers — strict profile/state/ceiling parsing, stale/missing review admission, schema and direct-execution filtering, argument binding, binding/token/origin checks, event idempotency and conversation ordering, typed outcomes, late-work correlation, promotion/follow-up inheritance; provider adapter E2E is separate evidence. Enforcement: `tests/test_presence_admission.py` plus the both-layer boundary tests this list requires.
### Devtools isolation

View file

@ -1359,6 +1359,8 @@ def _capture_context_core(
presence_section = build_presence_context_section(
pathlib.Path(env.drive_root),
task_metadata.get("presence"),
str(task.get("id") or ""),
status_root=canonical_root, # a forked promoted root finds its binding's work canonically
)
if presence_section:
dynamic_parts.append(presence_section)

View file

@ -3,6 +3,7 @@
from __future__ import annotations
import json
from pathlib import Path
from typing import Any, Mapping
from ouroboros.contracts.chat_id_policy import HIDDEN_CHAT_ID, WEB_UI_CHAT_ID
@ -25,6 +26,112 @@ def is_presence_task(task: Mapping[str, Any]) -> bool:
)
# A Presence binding's own work is reached through this scope (owner Q1/Q2).
PRESENCE_OWN_WORK_SCOPE = "own_binding"
def presence_record_binding(record: Any) -> str:
"""The nonempty host binding id one task/queue record carries, else ``""``."""
metadata = record.get("metadata") if isinstance(record, Mapping) else None
presence = metadata.get("presence") if isinstance(metadata, Mapping) else None
value = presence.get("binding_id") if isinstance(presence, Mapping) else None
return value.strip() if isinstance(value, str) else ""
def presence_related_work(binding_id: str, record: Any) -> bool:
"""Independent work started from this same nonempty binding (owner Q1).
Related work is a promoted or follow-up ROOT carrying the host's Presence
provenance for exactly this binding id, whichever of its conversations it
came from. An inline Presence turn, a delegated child, an owner root and a
record without that provenance are never attributed to the binding.
"""
binding = str(binding_id or "").strip()
return bool(
binding
and isinstance(record, Mapping)
and presence_record_binding(record) == binding
and str(record.get("delegation_role") or "") == "root"
and not str(record.get("parent_task_id") or "").strip()
)
def presence_caller_binding(ctx: Any) -> str | None:
"""``None`` for a non-Presence caller; otherwise its binding id (may be empty)."""
metadata = getattr(ctx, "task_metadata", None)
presence = metadata.get("presence") if isinstance(metadata, Mapping) else None
if not isinstance(presence, Mapping):
return None
value = presence.get("binding_id")
return value.strip() if isinstance(value, str) else ""
def presence_sender_origin(ctx: Any) -> dict[str, str]:
"""Where a Presence caller's run started (its ``run_origin`` room/event facts).
It names the sending run's origin only: later arrivals in that turn may have
been written by other people, so it never claims authorship of quoted words.
"""
metadata = getattr(ctx, "task_metadata", None)
return dict(run_origin({"metadata": metadata if isinstance(metadata, Mapping) else {}}).get("presence") or {})
def presence_queue_task(drive_root: Any, task_id: str) -> dict[str, Any] | None:
"""The persisted queue row of one pending/running task, if the snapshot lists it."""
from ouroboros.utils import read_json_dict
snapshot = read_json_dict(Path(drive_root) / "state" / "queue_snapshot.json") or {}
for key in ("pending", "running"):
for item in snapshot.get(key) or []:
task = item.get("task") if isinstance(item, Mapping) else None
if isinstance(task, Mapping) and str(item.get("id") or task.get("id") or "") == task_id:
return {**dict(task), "id": task_id}
return None
def presence_target_record(drive_root: Any, task_id: str) -> Mapping[str, Any] | None:
"""The record that decides whose work ``task_id`` is.
The canonical task record decides; a legacy row without Presence provenance
may be established only by the queue's own task metadata.
"""
from ouroboros.task_results import load_task_result
target = str(task_id or "").strip()
try:
stored = load_task_result(Path(drive_root), target) if target else None
except (OSError, ValueError):
stored = None # an unreadable or invalid id is no evidence of relation
record = stored if isinstance(stored, Mapping) and stored else None
if record is None or not presence_record_binding(record):
record = presence_queue_task(drive_root, target) or record
return record
def presence_effective_hops(task_id: str, effective: Any) -> list[str]:
"""The OTHER tasks an effective projection of ``task_id`` carries: retry lineage and successor."""
if not isinstance(effective, Mapping):
return []
ids = [value for hop in effective.get("retry_lineage") or [] if isinstance(hop, Mapping)
for value in (hop.get("task_id"), hop.get("retry_task_id"))]
ids.append(effective.get("task_id") or effective.get("id"))
requested = str(task_id or "").strip()
return [hop for hop in dict.fromkeys(str(value or "").strip() for value in ids) if hop and hop != requested]
def presence_effective_related(binding: str, task_id: str, effective: Any, *, drive_root: Any) -> bool:
"""Whether every task an effective projection of ``task_id`` reaches is this binding's own work."""
return all(presence_related_work(binding, presence_target_record(drive_root, hop))
for hop in presence_effective_hops(task_id, effective))
def presence_provenance_from_task(task: Mapping[str, Any]) -> dict[str, str]:
"""Return the stable, non-secret presence facts carried by one task.

View file

@ -12,6 +12,7 @@ import logging
import pathlib
from typing import Dict, List
from ouroboros.presence_authority import presence_record_binding
from ouroboros.task_result_schema import (
quarantine_task_result,
task_result_schema_refusal,
@ -78,6 +79,9 @@ def raw_result_facts(results_dir: pathlib.Path, *, reader=None) -> tuple[Dict[st
malformed.append(name)
continue
facts = {field: str(data.get(field) or "") for field in _RESULT_FACT_KEYS}
# One derived scalar selects a Presence binding's own work without a
# second read; the full row still decides once the selection loads it.
facts["presence_binding_id"] = presence_record_binding(data)
facts["schema_refusal"] = task_result_schema_refusal(data)
rows[name] = facts
_RAW_TS_MEMO[key] = (signature, tuple(facts.items()))

View file

@ -118,6 +118,10 @@ def _finalize_loop_candidate(content, limit_ctx, tools, emit_progress, *, after_
if transcript_growth_signature(limit_ctx.messages) == spoken_before:
wait_for_acceptance_feedback(tools, limit_ctx, limit_ctx.llm_trace,
limit_ctx.tool_schemas, limit_ctx.owner_msg_seen)
elif isinstance(completion, dict): # that owed round also learns its finish is void
from ouroboros.presence_context import presence_finish_not_accepted_note
_append_or_merge_user_message(limit_ctx.messages, presence_finish_not_accepted_note(ctx, completion), slot=ctx)
return result
@ -420,6 +424,7 @@ def _record_transcript_prefix(ctx, messages, round_idx, accumulated_usage,
def _reset_turn_state(ctx: Any) -> None:
"""Clear the per-turn state this turn owns; nothing durable is touched."""
ctx._presence_completion, ctx._presence_completion_accepted = None, False
ctx._presence_forced_declaration = ctx._presence_forced_pending = None
ctx._delivery_candidate, ctx._delivery_candidate_revision, ctx._delivery_control_required = None, 0, False
ctx._delivery_evidence_revision, ctx._delivery_evidence_fingerprint = 0, ""
ctx.model_turn_state, ctx._authoring_handover, ctx._pending_model_wait_handover = ModelTurnState(), None, None

View file

@ -867,9 +867,9 @@ def _parse_delivery_control_object(
def _classify_parsed_delivery_control(
parsed: Optional[Dict[str, Any]],
duplicate_protocol_key: bool,
embedded: bool,
embedded: bool, *, envelope_keys: Tuple[str, ...] = (),
) -> Tuple[str, str, str]:
"""Return ``(kind, replacement, error)`` for a parsed control body."""
"""Return ``(kind, replacement, error)``; ``envelope_keys`` are members an armed caller reads itself."""
exact_error = "control must be one exact JSON object"
if embedded:
@ -890,7 +890,7 @@ def _classify_parsed_delivery_control(
selected = str(parsed.get("delivery_control") or "")
if "pending_review" in parsed and str(parsed.get("pending_review") or "").strip().lower() not in {"wait", "finish"}:
return "invalid", "", 'pending_review must be "wait" or "finish"'
keys = set(parsed) - {"acceptance_subject", "pending_review"}
keys = set(parsed) - {"acceptance_subject", "pending_review", *envelope_keys}
if selected == "keep" and keys == {"delivery_control"}:
return "keep", "", ""
if selected == "replace" and keys == {"delivery_control", "full_answer"}:
@ -905,7 +905,7 @@ def _resolve_forced_delivery_control_body(
raw: str,
candidate: Optional[DeliveryCandidate],
*,
armed: bool,
armed: bool, envelope_keys: Tuple[str, ...] = (), # members the armed caller reads itself
) -> Tuple[str, bool, bool, bool, bool]:
"""Return text plus retained/degraded/consumed/replaced facts."""
@ -913,7 +913,7 @@ def _resolve_forced_delivery_control_body(
candidate = None
parsed, duplicate_protocol_key, embedded_protocol = _parse_delivery_control_body(raw)
control_kind, replacement, _error = _classify_parsed_delivery_control(
parsed, duplicate_protocol_key, embedded_protocol,
parsed, duplicate_protocol_key, embedded_protocol, envelope_keys=envelope_keys,
)
historical = bool(
not armed

View file

@ -388,8 +388,8 @@ def _maybe_enforce_child_absorption_gate(
prompt=(
"[FINALIZE_WITH_UNABSORBED_CHILDREN]\n"
"You still have child results without exact dispositions and already received one "
"child-absorption reminder. Produce an honest best-effort final answer now; name the "
"unabsorbed or unfinished children explicitly. Current child state: "
"child-absorption reminder. Produce an honest best-effort final answer now that says "
"what remains unabsorbed or unfinished; the exact child state is: "
f"{_undecided_children_listing(undecided)}."
),
fallback_text="⚠️ Finalized best-effort with undispositioned child results.",
@ -637,10 +637,101 @@ def _prepare_forced_prompt(
prompt
+ _loop()._forced_delegation_note(tools_ctx, llm_trace)
+ _forced_state_facts(ctx, llm_trace)
+ _presence_forced_contract(ctx, tools_ctx)
+ _forced_subject_prompt(ctx, llm_trace)
)
def _presence_forced_contract(ctx: _RoundLimitContext, tools_ctx: Any) -> str:
"""Arm a Presence task's ONE forced call to declare its outward delivery apart from its record.
The forced answer is the internal record (owner, review, task result); what the
conversation receives is only the nested ``presence_finish`` declaration. Until a
valid declaration arrives with a model-final answer the arm stays ``missing``, so
no untyped internal prose becomes Presence speech (owner Q4). Not a delivery receipt.
Only a context whose final becomes a Presence result is armed: the host ceiling AND
the Presence metadata the pipeline keys that result on. A delegated child inherits
the ceiling with its contract, but it answers its parent, not a conversation.
"""
contract = getattr(tools_ctx, "task_contract", None)
metadata = getattr(tools_ctx, "task_metadata", None)
presence = metadata.get("presence") if isinstance(metadata, dict) else None
if not (isinstance(contract, dict) and isinstance(contract.get("capability_ceiling"), dict)
and isinstance(presence, dict)):
return ""
tools_ctx._presence_forced_declaration = {"status": "missing", "reason": "no presence_finish declaration"}
tools_ctx._presence_forced_pending = None
try:
from ouroboros.presence_context import presence_send_facts
from ouroboros.tool_access import canonical_data_root
# Receipts live on the canonical root; a forked execution drive holds none.
sent = presence_send_facts(canonical_data_root(tools_ctx), ctx.task_id, presence)
except Exception:
sent = "unknown (receipts unreadable)"
handoff = getattr(tools_ctx, "_swarm_handoff_attempt", None)
scheduled = isinstance(handoff, dict) and str(handoff.get("status") or "") == "scheduled"
return (
"\n\n[PRESENCE_DELIVERY]\n"
"This task answers a Presence conversation. Your answer is its internal record for the "
"owner and review; the people in the conversation receive only what you declare. Return "
"exactly one JSON object and no other text: "
'{"delivery_control":"replace","full_answer":"<complete internal record>",'
'"presence_finish":{"outcome":"message","message":"<new text for the conversation>"}}'
+ (' ("delivery_control":"keep" without full_answer keeps the current answer as the record)'
if _loop()._live_delivery_candidate(ctx) is not None else "")
+ ". Outcomes: message = new useful speech on their subject, an honest partial included; "
"silent = nothing new needs saying; tool_delivered = the substantive result already reached "
"them through a transport tool; deferred = acknowledge work that was actually scheduled"
+ (" (it was)" if scheduled else " (none was)") + ". Execution, review, provider and helper "
"failures, limits and internal ids are not conversation content; keep them in full_answer. "
"An early acknowledgement is not the promised result and an uncertain send may not have "
"landed. Without a valid presence_finish nothing new is sent. Sends confirmed for this task "
f"so far: {sent}."
)
def _read_presence_declaration(tools_ctx: Any, extracted: str) -> None:
"""Record the nested ``presence_finish`` of an armed forced body without rewriting that body.
The resolver keeps reading the original bytes (``envelope_keys`` admits this one
key), so the parser's duplicate-key evidence still reaches every rail, the
acceptance subject included. A declaration is valid only in an envelope that
repeats no key anywhere. Each read records its non-speaking verdict at once; a
valid declaration speaks only once its answer becomes the model final.
"""
from ouroboros.loop_delivery import _parse_delivery_control_object
from ouroboros.observability import strip_protocol_fence
from ouroboros.tools.presence import PRESENCE_OUTCOMES
tools_ctx._presence_forced_pending = None
tools_ctx._presence_forced_declaration = {"status": "missing", "reason": "no presence_finish declaration"}
parsed, duplicate = _parse_delivery_control_object(strip_protocol_fence(extracted))
if not duplicate and (not isinstance(parsed, dict) or "presence_finish" not in parsed):
return
value = parsed.get("presence_finish") if isinstance(parsed, dict) else None
reason = ""
if duplicate or getattr(parsed, "has_duplicate_keys", False):
reason = "the forced envelope repeats a key"
elif not isinstance(value, dict) or not set(value) <= {"outcome", "message"} \
or value.get("outcome") not in PRESENCE_OUTCOMES or not isinstance(value.get("message", ""), str):
reason = "presence_finish must be one {outcome, message} object with a known outcome"
else:
outcome, message = value["outcome"], value.get("message", "").strip()
handoff = getattr(tools_ctx, "_swarm_handoff_attempt", None)
if outcome == "message" and not message:
reason = "message needs nonblank conversational text"
elif outcome == "silent" and message:
reason = "silent carries no text"
elif outcome == "deferred" and not (isinstance(handoff, dict) and handoff.get("status") == "scheduled"):
reason = "deferred needs work that was actually scheduled"
if reason:
tools_ctx._presence_forced_declaration = {"status": "invalid", "reason": reason}
else:
tools_ctx._presence_forced_pending = {
"status": "declared", "outcome": value["outcome"], "message": value.get("message", "").strip()}
def _forced_subject_prompt(ctx: _RoundLimitContext, llm_trace: Dict[str, Any]) -> str:
"""Capture the source before pricing/sending, never after a reply arrives."""
from ouroboros.loop_acceptance import capture_acceptance_observation, acceptance_observation_prompt
@ -1012,14 +1103,18 @@ def _resolve_forced_delivery_control(
"""Resolve forced control; returns text, degradation, retained, replaced."""
if tools_ctx is None or not extracted:
return extracted, "", False, False
presence_armed = isinstance(getattr(tools_ctx, "_presence_forced_declaration", None), dict)
if presence_armed:
_read_presence_declaration(tools_ctx, extracted)
candidate = getattr(tools_ctx, "_delivery_candidate", None)
armed = bool(getattr(tools_ctx, "_delivery_control_required", False)) or (
armed = presence_armed or bool(getattr(tools_ctx, "_delivery_control_required", False)) or (
isinstance(candidate, _loop().DeliveryCandidate)
and _loop()._delivery_replace_required(candidate)
)
resolved, retained, degraded, consumed, replaced = (
_loop()._resolve_forced_delivery_control_body(
extracted, candidate, armed=armed,
envelope_keys=("presence_finish",) if presence_armed else (),
)
)
if consumed:
@ -1194,6 +1289,8 @@ def _forced_final_answer(
set_terminal_host_notice(ctx.accumulated_usage, plan_suffix, _loop()._forced_orphan_note(ctx))
full_text = extracted
ctx.accumulated_usage["terminal_origin"] = TERMINAL_ORIGIN_MODEL_FINAL
if isinstance(getattr(tools_ctx, "_presence_forced_pending", None), dict):
tools_ctx._presence_forced_declaration = tools_ctx._presence_forced_pending
candidate = _publish_model_forced_candidate(
ctx, llm_trace, full_text, reason_code,
degraded_reason=control_degraded,

View file

@ -252,11 +252,14 @@ def write_task_message(
msg_id: Optional[str] = None,
review_feedback: Optional[Dict[str, Any]] = None,
relation: str = "",
sender_origin: Optional[Dict[str, Any]] = None,
) -> bool:
"""Write an addressed task-tree message without forging owner provenance.
``relation`` is the peer_task sender's typed place relative to the
recipient (``sibling`` / ``parent``); stored only when non-empty.
``sender_origin`` is the sending run's host-recorded origin (a Presence
room/event), never the author of the words it quotes.
"""
if provenance not in TASK_MESSAGE_PROVENANCES:
@ -275,6 +278,8 @@ def write_task_message(
entry["relayed_from_task_id"] = str(relayed_from_task_id)
if str(relation or ""):
entry["relation"] = str(relation)
if sender_origin:
entry["sender_origin"] = {str(key): str(value) for key, value in dict(sender_origin).items()}
if provenance == "system" and isinstance(review_feedback, dict):
entry["review_feedback"] = dict(review_feedback)
try:
@ -376,7 +381,10 @@ def deliver_task_message(
elif provenance == PROVENANCE_INDEPENDENT_TASK:
# A peer root's own words: never the ancestor fallback, which would
# place a stranger above the recipient in its tree.
prefix = f"[Message from independent task {source}]"
origin = entry.get("sender_origin") if isinstance(entry.get("sender_origin"), dict) else {}
prefix = f"[Message from independent task {source}" + (
"; that task's run started from " + json.dumps(origin, ensure_ascii=False, sort_keys=True)
+ ", which does not make it the author of any words it quotes]" if origin else "]")
elif provenance == PROVENANCE_PEER_TASK:
# A contribution from inside the tree without authority over the
# recipient: the stamped relation names the sender's place, so a
@ -691,6 +699,8 @@ def drain_owner_entries(
# left out of the projection it would never be delivered.
if str(entry.get("relation") or ""):
drained["relation"] = str(entry["relation"])
if isinstance(entry.get("sender_origin"), dict) and entry.get("sender_origin"):
drained["sender_origin"] = dict(entry["sender_origin"]) # rendered beside the words
entries.append(drained)
if _read_status is not None:
_read_status["complete"] = complete

View file

@ -8,6 +8,16 @@ from dataclasses import dataclass
from pathlib import Path, PurePosixPath
from typing import Any, Mapping
from ouroboros.dialogue_provenance import ( # the provenance predicate; authority re-exports it
PRESENCE_OWN_WORK_SCOPE,
presence_caller_binding,
presence_effective_hops,
presence_effective_related,
presence_queue_task,
presence_record_binding,
presence_related_work,
presence_target_record,
)
from ouroboros.presence_capabilities import (
PresenceArgumentBinding,
PresenceProfileResolution,
@ -19,6 +29,12 @@ from ouroboros.tool_capabilities import COGNITIVE_MEMORY_TOOL_NAMES
PRESENCE_CEILING_SCHEMA_VERSION = 1
_SHA256_LEN = 64
# The own-work baseline (owner Q1/Q2): a Presence mind may find, read and message
# the independent work it started from ANY conversation of its own binding. The two
# readers arrive bound to that scope (the host overwrites the argument); steer_task
# is narrowed by its handler for every Presence caller. A profile that selected one
# of these names keeps that grant, so a selected global reader stays global.
_OWN_WORK_BASELINE = (("get_task_result", True), ("recent_tasks", True), ("steer_task", False))
class PresenceAuthorityError(ValueError):
@ -331,6 +347,12 @@ def build_presence_capability_ceiling(
for name in sorted(COGNITIVE_MEMORY_TOOL_NAMES)
if name not in selected
)
scope = (PresenceArgumentBinding(("presence_scope",), "static", static_value=PRESENCE_OWN_WORK_SCOPE),)
tools.extend(
PresenceToolGrant(name, scope if scoped else ())
for name, scoped in _OWN_WORK_BASELINE
if name not in selected
)
provisional = PresenceCapabilityCeiling(
skill_name=_text(skill_name, "skill_name"),
skill_content_hash=_sha(skill_content_hash, "skill_content_hash"),
@ -518,8 +540,68 @@ def presence_ceiling_allows_binding(ceiling: PresenceCapabilityCeiling, binding:
return False
def presence_effective_refusal(ctx: Any, task_id: str, effective: Any, *, drive_root: Any = None,
same_tree: bool = False) -> str:
"""Refusal when an admitted ``task_id``'s effective projection reaches work the caller may not address.
A retry successor replaces the requested record's content, so each hop is judged
by the same rule before anything is projected; the foreign id is not named.
"""
for hop in presence_effective_hops(task_id, effective):
if presence_work_refusal(ctx, hop, drive_root=drive_root, same_tree=same_tree):
return (
f"⚠️ PRESENCE_CAPABILITY_BLOCKED: task {task_id}'s effective result continues in work "
"that was not started from this Presence binding; a Presence turn reads only that work."
)
return ""
def presence_work_refusal(ctx: Any, task_id: str, *, drive_root: Any = None, same_tree: bool = False) -> str:
"""Refusal text when a Presence caller may not address ``task_id``, else ``""``.
A non-Presence caller is never narrowed here. ``presence_target_record``
decides, and a record naming another binding refuses.
``same_tree`` also admits the caller's own task tree (its lineage reads).
``drive_root`` defaults to the caller's task-status root.
"""
binding = presence_caller_binding(ctx)
if binding is None:
return ""
if drive_root is None:
metadata = getattr(ctx, "task_metadata", None)
drive_root = ((metadata.get("budget_drive_root") if isinstance(metadata, Mapping) else "")
or getattr(ctx, "budget_drive_root", "") or ctx.drive_root)
target = str(task_id or "").strip()
record = presence_target_record(drive_root, target)
if isinstance(record, Mapping) and presence_related_work(binding, record):
return ""
if same_tree and isinstance(record, Mapping):
metadata = getattr(ctx, "task_metadata", None)
own_id = str(getattr(ctx, "task_id", "") or "")
own_root = str((metadata or {}).get("root_task_id") or own_id) if isinstance(metadata, Mapping) else own_id
if own_id and (target == own_id or str(record.get("parent_task_id") or "") == own_id
or str(record.get("root_task_id") or "") == own_root):
return ""
return (
f"⚠️ PRESENCE_CAPABILITY_BLOCKED: task {target or '?'} is not independent work started from "
"this Presence binding (same binding_id, any of its conversations); a Presence turn "
"addresses only that work."
)
__all__ = [
"PRESENCE_CEILING_SCHEMA_VERSION",
"PRESENCE_OWN_WORK_SCOPE",
"presence_caller_binding",
"presence_effective_refusal",
"presence_effective_related",
"presence_queue_task",
"presence_record_binding",
"presence_related_work",
"presence_target_record",
"presence_work_refusal",
"PresenceAuthorityError",
"PresenceCapabilityCeiling",
"PresenceResourceGrant",

View file

@ -83,8 +83,102 @@ def _previous_turn_line(previous: Mapping[str, Any]) -> str:
f"outcome {previous.get('outcome')}, delivery {previous.get('delivery') or 'unknown'}): {body}.{work}")
def build_presence_context_section(drive_root: Path, value: Any) -> str:
"""Render host-authored presence context, including declared full KB topics."""
_OWN_WORK_PAGE = 5
def _own_work_section(drive_root: Path, value: Mapping[str, Any], task_id: str) -> str:
"""The first page of independent work this binding started, from the scoped reader."""
from ouroboros.tools.recent_tasks import recent_tasks_page
binding = str(value.get("binding_id") or "").strip()
if not binding:
return ""
page = recent_tasks_page(Path(drive_root), limit=_OWN_WORK_PAGE, binding=binding,
exclude=str(task_id or ""), restricted=True)
here = str((value.get("event") or {}).get("conversation_key") or "")
lines = []
for row in page.get("tasks") or []:
origin = row.get("presence_origin") if isinstance(row.get("presence_origin"), Mapping) else {}
key = str(origin.get("conversation_key") or "")
where = ("this conversation" if key and key == here
else f"conversation {key}" if key else "a conversation the queue row names")
preview = " ".join(str(row.get("result_preview") or "").split())[:200]
lines.append(
f"- {row.get('task_id')} [{row.get('status') or 'unknown'}"
+ (f", cancel {row['cancel_state']}" if row.get("cancel_state") else "")
+ (", its result row is unreadable" if row.get("result_row") == "unreadable" else "") + f"] from {where}: "
+ json.dumps(" ".join(str(row.get("description") or "").split())[:200], ensure_ascii=False)
+ (f"; result preview {json.dumps(preview, ensure_ascii=False)}"
if preview and row.get("status") in {"completed", "failed", "cancelled"} else "")
+ (f"; {row['effective_result']}" if row.get("effective_result") else "")
)
gap = page.get("read_gap") if isinstance(page.get("read_gap"), Mapping) else {}
if gap:
lines.append("(some results are unreadable now: "
+ ("the result root" if gap.get("result_root") else
f"{gap.get('unattributed_unreadable_rows')} row(s) no record attributes")
+ "; this binding's work may be among them)")
if page.get("error"):
lines.append(f"(listing unavailable now: {page['error'].get('code')}; page again with recent_tasks)")
elif page.get("remaining"):
lines.append(f"({page['remaining']} more: recent_tasks(presence_scope=\"own_binding\", "
f"offset={page['offset'] + page['returned']}, snapshot=\"{page['snapshot']}\"))")
if not lines:
return ""
return (
"## Work started from this binding (host-authored facts)\n\n"
"Independent work this Presence binding started, from this or another of its conversations, "
"newest first (queued work without a result row leads). Being listed says nothing about whether "
"its result reached anyone. Where your tools include them, get_task_result(task_id, "
"presence_scope=\"own_binding\") reads one exactly, steer_task gives it new facts as this "
"task's own message, and presence_cancel_work requests its cancellation. Work of other "
"bindings and the owner's own tasks are not addressable from here.\n\n"
+ "\n".join(lines)
)
def presence_send_facts(drive_root: Path, task_id: str, value: Any) -> str:
"""Transport sends of this task confirmed so far in its own conversation, or honest unknown."""
from ouroboros.presence_runner import _live_task_rows, _turn_sends
value = value if isinstance(value, Mapping) else {}
key = str((value.get("event") or {}).get("conversation_key") or "")
if not key or value.get("delivery_reporting_version") != 1:
return "unknown (this transport reports no delivery receipts)"
sent, uncertain = _turn_sends(_live_task_rows(Path(drive_root), str(task_id or ""), key))
if sent is None:
return "unknown (no readable receipt coverage for this task)"
said = " / ".join(json.dumps(text, ensure_ascii=False) for text in sent if text) or "none"
return said + (f"; {uncertain} more part(s) have an uncertain outcome and may have landed" if uncertain else "")
def presence_finish_not_accepted_note(ctx: Any, completion: Mapping[str, Any]) -> str:
"""Tell the round that follows an invalidated finish what is void and what is already sent."""
from ouroboros.tool_access import canonical_data_root
metadata = getattr(ctx, "task_metadata", None)
try: # the host records receipts on the canonical root, never on a forked execution drive
sent = presence_send_facts(canonical_data_root(ctx), str(getattr(ctx, "task_id", "") or ""),
metadata.get("presence") if isinstance(metadata, Mapping) else None)
except Exception:
sent = "unknown (receipts unreadable)"
return (
f"[PRESENCE_FINISH_NOT_ACCEPTED]\npresence_finish({completion.get('outcome')}) was not accepted: "
"finalization asked for the work above first, so it no longer decides what the conversation "
f"receives. Sends confirmed for this task so far: {sent}. When the work is done, finish again: "
"tool_delivered or silent when the substantive result already reached the conversation and "
"nothing new needs saying, message only for new useful speech. Review, helper and limit "
"details are not conversation content."
)
def build_presence_context_section(drive_root: Path, value: Any, task_id: str = "", *,
status_root: Path | None = None) -> str:
"""Render host-authored presence context, including declared full KB topics.
``status_root`` is the canonical task root the own-work catalogue reads (a forked
execution drive holds only its own worker rows); it defaults to ``drive_root``.
"""
if not isinstance(value, Mapping):
return ""
@ -170,6 +264,12 @@ def build_presence_context_section(drive_root: Path, value: Any) -> str:
f"The host lost an earlier attempt of this event before it finished; {detail}. "
"Do not resend what was already delivered."
)
try:
own_work = _own_work_section(Path(status_root or drive_root), value, task_id)
except Exception:
own_work = "## Work started from this binding (host-authored facts)\n\nUnavailable now; page it with recent_tasks."
if own_work:
parts.append(own_work)
parts += [
"## Current presence event (host-authored facts)\n\n"
+ json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True, default=str),
@ -178,4 +278,7 @@ def build_presence_context_section(drive_root: Path, value: Any) -> str:
return "\n\n".join(parts)
__all__ = ["build_presence_context_section", "frame_presence_user_content"]
__all__ = [
"build_presence_context_section", "frame_presence_user_content",
"presence_finish_not_accepted_note", "presence_send_facts",
]

View file

@ -121,6 +121,15 @@ def build_presence_result_event(task: dict[str, Any], text: str, ctx: Any, *, te
completion = completion if (
isinstance(completion, dict) and getattr(ctx, "_presence_completion_accepted", False)
) else {}
# A forced final declares its outward delivery beside the internal record; without
# a valid declaration the record never becomes conversation speech (owner Q4).
declared = getattr(ctx, "_presence_forced_declaration", None)
if not completion and isinstance(declared, dict):
completion = declared if declared.get("status") == "declared" else {"outcome": "silent"}
# Only a message/deferred body is speech; a tool_delivered note stays context even
# when owed work below turns the outcome into deferred.
text = str(completion.get("message") or "") if completion.get("outcome") in {"message", "deferred"} else ""
note = str(completion.get("message") or "") if completion.get("outcome") == "tool_delivered" else ""
outcome = str(completion.get("outcome") or "message").strip()
handoff = getattr(ctx, "_swarm_handoff_attempt", None)
handoff = handoff if isinstance(handoff, dict) else {}
@ -141,6 +150,8 @@ def build_presence_result_event(task: dict[str, Any], text: str, ctx: Any, *, te
metadata["presence_result_text"] = result_text
if work_ref:
metadata["presence_work_ref"] = work_ref
if not getattr(ctx, "_presence_completion_accepted", False) and isinstance(declared, dict):
metadata["presence_declaration"] = {key: declared[key] for key in ("status", "reason") if declared.get(key)}
task["metadata"] = metadata
return {
"type": "presence_result",
@ -148,7 +159,7 @@ def build_presence_result_event(task: dict[str, Any], text: str, ctx: Any, *, te
"outcome": outcome,
"text": result_text,
# The accepted presence_finish message; a tool_delivered note is context, never speech.
"message": result_text or (str(completion.get("message") or "") if outcome == "tool_delivered" else ""),
"message": result_text or note,
"work_ref": work_ref,
"ts": utc_now_iso(),
}

View file

@ -415,6 +415,7 @@ def get_tools() -> List[ToolEntry]:
"focus_source_sha256": {"type": "string", "default": "", "description": "With include_focus_source: select the retained source by the sha256 the roster row quoted, so a later focus of the same author cannot substitute its evidence."},
"source_start_char": {"type": "integer", "description": "Inclusive character offset for the requested canonical source range."},
"source_end_char": {"type": "integer", "description": "Exclusive character offset for the requested canonical source range. A range outside the source returns no text: the answer names complete_chars and the range received, and is an argument error."},
"presence_scope": {"type": "string", "enum": ["own_binding"], "description": "Presence tasks only: read just independent work started from this Presence binding (any of its conversations) or this task's own tree."},
}},
}, _get_task_result),
ToolEntry("wait_task", {

View file

@ -189,12 +189,15 @@ def _record_promotion_admission_stub(ctx: ToolContext, evt: Dict[str, Any], mode
duplicate root (#1160). The stub is the negative side only: ``emitted`` is not
a scheduled status, positive scheduling authority stays with the supervisor's
own receipt, and ``create_only`` initializes ABSENCE alone, so a supervisor
that already answered keeps its row byte-for-byte.
that already answered keeps its row byte-for-byte. A Presence promote's stub
carries the event's host provenance as the would-be root, exactly as the
admission writes it, so its own binding can read the pending reconciliation.
"""
from ouroboros.routing_wait import PROMOTION_ADMISSION_EMITTED
from ouroboros.task_results import STATUS_REQUESTED, write_task_result
task_id = str(evt.get("task_id") or "")
presence = evt.get("presence") if isinstance(evt.get("presence"), dict) else None
try:
write_task_result(
_routing_status_root(ctx), task_id, STATUS_REQUESTED,
@ -207,6 +210,8 @@ def _record_promotion_admission_stub(ctx: ToolContext, evt: Dict[str, Any], mode
"emitted_at": utc_now_iso(),
"transport_mode": mode,
},
**({"metadata": {"presence": dict(presence)}, "source": "presence_promote",
"delegation_role": "root", "root_task_id": task_id} if presence else {}),
)
except Exception as exc:
# The promote itself proceeds; what is lost is the reconciliation read, so

View file

@ -972,6 +972,7 @@ def _send_task_message(
No origin-bytes substitution, attachments or owner client surface.
The result says WRITTEN: the target reads it at its next checkpoint.
"""
from ouroboros.dialogue_provenance import presence_caller_binding, presence_sender_origin
from ouroboros.project_dialogue import AGENT_RECEIPT_ID_PREFIX
routing_token = uuid.uuid4().hex
@ -987,6 +988,8 @@ def _send_task_message(
"issuer": dict(issuer),
"ts": utc_now_iso(),
}
if (binding := presence_caller_binding(ctx)) is not None: # admitted only to this binding's own work (owner Q2)
evt.update(presence_binding_id=binding, sender_origin=presence_sender_origin(ctx))
mode, receipt = _emit_and_wait_for_routing(ctx, evt)
if str(receipt.get("status") or "") == "delivered":
return (

View file

@ -188,11 +188,33 @@ def _get_task_result(
include_work_order_source: bool = False, source_start_char: Any = None,
source_end_char: Any = None, include_completion_source: bool = False,
known_result_sha256: str = "", include_focus_source: bool = False, focus_source_sha256: str = "",
presence_scope: str = "",
) -> str:
"""Read a task result, or a bounded canonical work-order/completion source range."""
metadata = getattr(ctx, "task_metadata", {}) if isinstance(getattr(ctx, "task_metadata", {}), dict) else {}
status_drive_root = Path(str(metadata.get("budget_drive_root") or getattr(ctx, "budget_drive_root", "") or ctx.drive_root))
scoped = bool(str(presence_scope or "").strip())
if scoped:
from ouroboros.presence_authority import PRESENCE_OWN_WORK_SCOPE, presence_caller_binding, presence_work_refusal
# The scoped read admits exactly the work a scoped page can list, plus this
# task's own tree; everything else refuses before any record is projected.
refusal = (
"⚠️ PRESENCE_CAPABILITY_BLOCKED: presence_scope=own_binding needs a Presence task with a binding id."
if str(presence_scope).strip() != PRESENCE_OWN_WORK_SCOPE or not presence_caller_binding(ctx)
else presence_work_refusal(ctx, str(task_id or ""), drive_root=status_drive_root, same_tree=True)
)
if refusal:
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=refusal))
data = load_effective_task_result(status_drive_root, task_id)
if scoped:
from ouroboros.presence_authority import presence_effective_refusal
# The effective read may substitute a retry successor's record: it is judged too.
refusal = presence_effective_refusal(ctx, str(task_id or ""), data, drive_root=status_drive_root,
same_tree=True)
if refusal:
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=refusal))
from ouroboros.tools.recent_tasks import _restricted_actor
restricted = _restricted_actor(ctx)

View file

@ -670,6 +670,23 @@ def _cancel_task(ctx: ToolContext, task_id: str, reason: str = "") -> str:
if not own and (_is_delegated_task(ctx) or is_observe_origin(getattr(ctx, "task_metadata", {}))):
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=(f"⚠️ cancel_task: {tid} is not a child of this task — a delegated task may only cancel its own children, and a consciousness wake at the Observe level is held to the same rule.")))
from ouroboros.presence_authority import presence_caller_binding, presence_work_refusal
if not own and presence_caller_binding(ctx) is not None:
# A Presence caller stops only its own binding's work or its own tree, and
# a redirected retry is judged at its effective target too — before any intent.
from ouroboros.cancel_intents import _validated_single_cancel_target
refusal = presence_work_refusal(ctx, tid, drive_root=status_drive_root, same_tree=True)
if not refusal:
try:
effective = _validated_single_cancel_target(status_drive_root, tid)
except Exception:
effective = tid
if effective != tid:
refusal = presence_work_refusal(ctx, effective, drive_root=status_drive_root, same_tree=True)
if refusal:
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=refusal))
# Durable cancel intent — the ONE ingress (phase A, owner batch-4 1=A). The
# canonical status never carries intent: the supervisor's cancellation
# custody claims this intent, tears the task down, and settles the terminal

View file

@ -300,9 +300,10 @@ def _initiate_presence(
def _cancel_presence_work(ctx: ToolContext, work_ref: str, reason: str = "") -> str:
"""Cancel only work correlated to this exact presence binding/conversation."""
"""Cancel work started from this presence binding (any of its conversations) or this turn's own tree."""
from ouroboros.task_results import load_task_result, validate_task_id
from ouroboros.presence_authority import presence_work_refusal
from ouroboros.task_results import validate_task_id
from ouroboros.tool_access import canonical_data_root
from ouroboros.tools.join_ledger import _cancel_task
@ -310,21 +311,10 @@ def _cancel_presence_work(ctx: ToolContext, work_ref: str, reason: str = "") ->
task_id = validate_task_id(work_ref)
except ValueError as exc:
return _publish_tool_result(ctx, ToolResult(status="error", code="TOOL_ARG_ERROR", text=(f"ERROR: PRESENCE_WORK_REF_INVALID: {exc}")))
current_meta = getattr(ctx, "task_metadata", {})
current = current_meta.get("presence") if isinstance(current_meta, dict) else None
stored = load_task_result(canonical_data_root(ctx), task_id) or {}
target_meta = stored.get("metadata") if isinstance(stored.get("metadata"), dict) else {}
target = target_meta.get("presence") if isinstance(target_meta.get("presence"), dict) else None
if not isinstance(current, dict) or not isinstance(target, dict):
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=("ERROR: PRESENCE_WORK_NOT_CORRELATED")))
current_event = current.get("event") if isinstance(current.get("event"), dict) else {}
target_event = target.get("event") if isinstance(target.get("event"), dict) else {}
if (
str(current.get("binding_id") or "") != str(target.get("binding_id") or "")
or str(current_event.get("conversation_key") or "")
!= str(target_event.get("conversation_key") or "")
):
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=("ERROR: PRESENCE_WORK_NOT_CORRELATED")))
refusal = presence_work_refusal(ctx, task_id, drive_root=canonical_data_root(ctx), same_tree=True)
if refusal or not isinstance((getattr(ctx, "task_metadata", {}) or {}).get("presence"), dict):
return _publish_tool_result(ctx, ToolResult(status="blocked", code="ACCESS_BLOCKED", text=(
"ERROR: PRESENCE_WORK_NOT_CORRELATED: " + (refusal.split(": ", 1)[-1] or "this is not a presence task."))))
return _cancel_task(ctx, task_id, reason)
@ -440,9 +430,10 @@ def get_tools() -> List[ToolEntry]:
schema={
"name": "presence_cancel_work",
"description": (
"Request cancellation of long work previously deferred from this exact "
"presence binding and conversation. The opaque work_ref is correlation, "
"not general task authority."
"Request cancellation of independent work started from this presence "
"binding, in this or another of its conversations (or of this turn's own "
"children). The result is a request receipt, not proof the work stopped; "
"work of another binding or the owner's own tasks is refused."
),
"parameters": {
"type": "object",

View file

@ -10,11 +10,19 @@ from typing import Any, Dict, List
from ouroboros.tools.registry import ToolContext, ToolEntry
from ouroboros.outcomes import normalize_outcome_axes
from ouroboros.task_status import effective_task_result
from ouroboros.dialogue_provenance import is_presence_task
from ouroboros.dialogue_provenance import (
PRESENCE_OWN_WORK_SCOPE,
is_presence_task,
presence_caller_binding,
presence_effective_related,
presence_provenance_from_task,
presence_related_work,
)
_MAX_TASKS = 20
_PREVIEW_CHARS = 800
_ORIGIN_KEYS = ("provider", "conversation_id", "thread_id", "conversation_key", "source_event_id")
def _coerce_limit(value: Any) -> int:
@ -49,11 +57,18 @@ def _task_record(
drive_root: pathlib.Path,
include_results: bool,
include_traces: bool,
binding: str | None = None,
) -> tuple[Dict[str, Any] | None, Dict[str, str] | None]:
data, error = _read_json(path)
if data is None:
raw, error = _read_json(path)
if raw is None:
return None, {"path": str(path), "error": error}
data = effective_task_result(drive_root, data)
data = effective_task_result(drive_root, raw)
withheld = False
if binding is not None:
# A scoped page lists this binding's row; a retry successor it redirects to is judged too.
withheld = not presence_effective_related(binding, str(raw.get("task_id") or path.stem), data,
drive_root=drive_root)
data = raw if withheld else data
result = str(data.get("result") or "")
from ouroboros.cost_projection import cost_projection
@ -94,10 +109,19 @@ def _task_record(
record["result"] = result
if include_traces:
record["trace_summary"] = str(data.get("trace_summary") or "")
if data.get("cancel_state"):
record["cancel_state"] = str(data["cancel_state"]) # requested is not stopped
origin = presence_provenance_from_task(data)
if origin:
# Which of the binding's conversations started it: the source room is a fact
# for the reader, never a reply address or a public disclosure.
record["presence_origin"] = {key: origin[key] for key in _ORIGIN_KEYS if origin.get(key)}
if withheld:
record["effective_result"] = "withheld: it continues in work not started from this binding"
return record, None
def _running_tasks(drive_root: pathlib.Path) -> List[Dict[str, Any]]:
def _running_tasks(drive_root: pathlib.Path, binding: str | None = None) -> List[Dict[str, Any]]:
snapshot, _error = _read_json(drive_root / "state" / "queue_snapshot.json")
snapshot = snapshot or {}
running = snapshot.get("running")
@ -107,15 +131,68 @@ def _running_tasks(drive_root: pathlib.Path) -> List[Dict[str, Any]]:
for item in running:
if not isinstance(item, dict):
continue
task = item.get("task") if isinstance(item.get("task"), dict) else {}
if binding is not None and not presence_related_work(binding, task):
continue # a scoped page lists no foreign running work
rows.append({
"task_id": str(item.get("id") or item.get("task_id") or ""),
"status": "running",
"description": str(item.get("text") or item.get("description") or ""),
"description": str(item.get("text") or item.get("description")
or task.get("description") or task.get("text") or ""),
"ts": str(item.get("ts") or snapshot.get("ts") or ""),
})
return rows
def _presence_scope_inventory(
drive_root: pathlib.Path, task_dir: pathlib.Path, binding: str, exclude: str,
) -> tuple[List[Dict[str, Any]], Dict[str, Any]]:
"""This binding's own work (queue-only active rows first, then result files newest first) and its read gap.
The memo's scalar binding fact selects files without a second read; a row
without Presence provenance may be established only by the queue's own task
metadata (a legacy pending promotion), and a row naming another binding stays out.
Only a READABLE result row replaces a queue row: an unreadable one leaves the
queued work listed, and unreadable rows nothing attributes are counted, never dropped.
"""
from ouroboros.gateway.task_list_scan import raw_result_facts
gap: Dict[str, Any] = {}
try:
facts, malformed = raw_result_facts(task_dir)
except OSError:
facts, malformed = {}, []
gap["result_root"] = "unreadable" # queued rows remain; no result row could be read
snapshot, _error = _read_json(drive_root / "state" / "queue_snapshot.json")
queued: Dict[str, tuple[str, Dict[str, Any]]] = {}
for status in ("running", "pending"):
for item in (snapshot or {}).get(status) or []:
task = item.get("task") if isinstance(item, dict) and isinstance(item.get("task"), dict) else {}
task_id = str(item.get("id") or task.get("id") or "") if isinstance(item, dict) else ""
if task_id and task_id not in queued:
queued[task_id] = (status, task)
selected: set[str] = set()
for name, row in facts.items():
task_id = row.get("task_id") or row.get("id") or name[:-5]
record = ({"metadata": {"presence": {"binding_id": row["presence_binding_id"]}},
"delegation_role": row.get("delegation_role"), "parent_task_id": row.get("parent_task_id")}
if row.get("presence_binding_id") else queued.get(task_id, ("", {}))[1])
if task_id != exclude and presence_related_work(binding, record):
selected.add(name)
unreadable = set(malformed)
queue_only = [
{"queue_task_id": task_id, "status": status,
"description": str(task.get("description") or task.get("text") or ""),
**({"result_row": "unreadable"} if f"{task_id}.json" in unreadable else {})}
for task_id, (status, task) in queued.items()
if f"{task_id}.json" not in facts and task_id != exclude and presence_related_work(binding, task)
]
unattributed = sum(1 for name in unreadable if name[:-5] not in queued) # a queue row attributed the rest
if unattributed:
gap["unattributed_unreadable_rows"] = unattributed # any of them may be this binding's work
return queue_only + [row for row in _task_file_inventory(task_dir) if row["name"] in selected], gap
def _task_file_inventory(task_dir: pathlib.Path) -> List[Dict[str, Any]]:
inventory: List[Dict[str, Any]] = []
if not task_dir.is_dir():
@ -141,13 +218,17 @@ def _recent_tasks_snapshot(
*,
include_results: bool,
include_traces: bool,
binding: str | None = None,
) -> str:
query: Dict[str, Any] = {
"include_results": bool(include_results),
"include_traces": bool(include_traces),
}
if binding is not None:
query["presence_binding"] = binding # another binding or scope never continues this cursor
payload = {
"schema_version": 1,
"query": {
"include_results": bool(include_results),
"include_traces": bool(include_traces),
},
"query": query,
"files": inventory,
}
encoded = json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
@ -161,14 +242,45 @@ def _handle_recent_tasks(
snapshot: str = "",
include_results: bool = False,
include_traces: bool = False,
presence_scope: str = "",
**_kwargs: Any,
) -> str:
"""Return recent completed task summaries from the canonical task root."""
from ouroboros.tool_access import canonical_data_root
drive_root = canonical_data_root(ctx)
binding = None
if str(presence_scope or "").strip():
binding = presence_caller_binding(ctx)
if str(presence_scope).strip() != PRESENCE_OWN_WORK_SCOPE or not binding:
return json.dumps({"ok": False, "host_code": "TOOL_ARG_ERROR", "error": {
"code": "PRESENCE_SCOPE_UNAVAILABLE",
"message": "presence_scope=own_binding needs a Presence task with a binding id.",
}}, ensure_ascii=False)
page = recent_tasks_page(
canonical_data_root(ctx), limit=limit, offset=offset, snapshot=snapshot,
include_results=bool(include_results), include_traces=bool(include_traces),
restricted=_restricted_actor(ctx), binding=binding,
exclude=str(getattr(ctx, "task_id", "") or "") if binding is not None else "",
)
if binding is not None:
page["presence_scope"] = {"scope": PRESENCE_OWN_WORK_SCOPE, "binding_id": binding}
return json.dumps(page, ensure_ascii=False, indent=2)
def recent_tasks_page(
drive_root: pathlib.Path,
*,
limit: Any = 5,
offset: Any = 0,
snapshot: str = "",
include_results: bool = False,
include_traces: bool = False,
restricted: bool = False,
binding: str | None = None,
exclude: str = "",
) -> Dict[str, Any]:
"""One stable page; ``binding`` filters to that Presence binding's own work BEFORE paging."""
task_dir = drive_root / "task_results"
restricted = _restricted_actor(ctx)
task_limit = _coerce_limit(limit)
try:
skip = max(0, int(offset or 0))
@ -180,23 +292,42 @@ def _handle_recent_tasks(
inventory: List[Dict[str, Any]] = []
current_snapshot = ""
stable = False
read_gap: Dict[str, Any] = {}
def _inventory() -> List[Dict[str, Any]]:
if binding is None:
return _task_file_inventory(task_dir)
rows, gap = _presence_scope_inventory(drive_root, task_dir, binding, exclude)
read_gap.clear()
read_gap.update(gap)
return rows
for _attempt in range(2):
tasks = []
unreadable_tasks = []
before = _task_file_inventory(task_dir)
before = _inventory()
current_snapshot = _recent_tasks_snapshot(
before,
include_results=bool(include_results),
include_traces=bool(include_traces),
binding=binding,
)
selected = before[skip:skip + task_limit]
for item in selected:
if item.get("queue_task_id"):
tasks.append({
"task_id": str(item["queue_task_id"]), "status": str(item.get("status") or ""),
"description": str(item.get("description") or ""), "source": "queue_snapshot",
**({"result_row": item["result_row"]} if item.get("result_row") else {}),
})
continue
path = task_dir / str(item["name"])
record, error = _task_record(
path,
drive_root=drive_root,
include_results=bool(include_results),
include_traces=bool(include_traces),
binding=binding,
)
if record is not None:
if restricted:
@ -206,7 +337,7 @@ def _handle_recent_tasks(
tasks.append(record)
elif error is not None:
unreadable_tasks.append(error)
inventory = _task_file_inventory(task_dir)
inventory = _inventory()
stable = before == inventory
if stable:
break
@ -214,7 +345,7 @@ def _handle_recent_tasks(
returned = min(task_limit, max(0, total - skip))
remaining = max(0, total - skip - returned)
base = {
"running": _running_tasks(drive_root),
"running": _running_tasks(drive_root, binding),
"tasks": tasks,
"unreadable_tasks": unreadable_tasks,
"source": {"reader": "recent_tasks", "root": "canonical_task_results"},
@ -229,10 +360,12 @@ def _handle_recent_tasks(
"snapshot": current_snapshot,
"include_results": bool(include_results),
"include_traces": bool(include_traces),
**({"presence_scope": PRESENCE_OWN_WORK_SCOPE} if binding is not None else {}),
} if remaining else None),
**({"read_gap": dict(read_gap)} if read_gap else {}),
}
if not stable:
return json.dumps({
return {
**base,
"tasks": [],
"unreadable_tasks": [],
@ -243,9 +376,9 @@ def _handle_recent_tasks(
"was returned; restart with offset=0 and no snapshot."
),
},
}, ensure_ascii=False, indent=2)
}
if requested_snapshot and requested_snapshot != current_snapshot:
return json.dumps({
return {
**base,
"tasks": [],
"unreadable_tasks": [],
@ -256,8 +389,8 @@ def _handle_recent_tasks(
"returned; restart with offset=0 and no snapshot."
),
},
}, ensure_ascii=False, indent=2)
return json.dumps(base, ensure_ascii=False, indent=2)
}
return base
def _restricted_actor(ctx: ToolContext) -> bool:
@ -321,6 +454,14 @@ def get_tools() -> List[ToolEntry]:
"description": "Include each task's trace_summary.",
"default": False,
},
"presence_scope": {
"type": "string",
"enum": ["own_binding"],
"description": (
"Presence tasks only: list just the independent work started from this "
"Presence binding in any of its conversations, pending and running included."
),
},
},
"required": [],
},

View file

@ -190,6 +190,13 @@ def _presence_bound_args(ctx: Any, name: str, args: Any) -> tuple[dict[str, Any]
"⚠️ PRESENCE_CAPABILITY_BLOCKED: "
f"{name!r} is outside this presence task's positive capability ceiling."
)
if ceiling is not None and name == "forward_to_worker":
# A selected forward keeps its own tree and reaches only this binding's work.
from ouroboros.presence_authority import presence_work_refusal
refusal = presence_work_refusal(ctx, str(bound.get("task_id") or ""), same_tree=True)
if refusal:
return {}, refusal
return bound, ""
except Exception as exc:
return {}, f"⚠️ PRESENCE_ARGUMENT_BINDING_BLOCKED: {exc}"

View file

@ -527,6 +527,11 @@ def _promote_chat_to_task_outcome(evt: Dict[str, Any], ctx: Any) -> Dict[str, An
else "Task is scheduled, but its owner-facing routing receipt was not confirmed."
),
attachment_manifest=list(outcome.get("attachment_manifest") or []),
# A Presence promotion's host-carried provenance is canonical from
# admission, so its binding finds, polls and controls the work while
# it is still queued (the worker's running write keeps the same value).
**({"metadata": {"presence": dict(evt["presence"])}, "source": "presence_promote"}
if isinstance(evt.get("presence"), dict) and evt.get("presence") else {}),
)
admission = stored.get("promotion_admission") if isinstance(stored, dict) else {}
if (

View file

@ -53,6 +53,13 @@ def _task_issued(evt: Dict[str, Any]) -> bool:
return str(_issuer(evt).get("kind") or "") == "task"
def _presence_target_related(evt: Dict[str, Any], task: Dict[str, Any]) -> bool:
"""A Presence sender's live target is independent work of its own binding."""
from ouroboros.dialogue_provenance import presence_related_work
return presence_related_work(str(evt.get("presence_binding_id") or ""), task)
def _refuse_steering_while_cancelling(
ctx: Any,
evt: Dict[str, Any],
@ -250,6 +257,8 @@ def _handle_steer_task(evt: Dict[str, Any], ctx: Any) -> None:
refusal = "target_unknown"
elif str(task.get("delegation_role") or "") == "subagent":
refusal = "subagent_target"
elif task_issued and "presence_binding_id" in evt and not _presence_target_related(evt, task):
refusal = "presence_work_not_related"
elif not task_issued and not _owner_lane_allows(ctx, task, target, chat_id):
refusal = "chat_mismatch"
else:
@ -373,6 +382,7 @@ def _handle_steer_task(evt: Dict[str, Any], ctx: Any) -> None:
if not write_task_message(
drive, message, target, source_task_id=issuer_task_id,
provenance=PROVENANCE_INDEPENDENT_TASK, msg_id=msg_id,
sender_origin=evt.get("sender_origin") if isinstance(evt.get("sender_origin"), dict) else None,
):
raise OSError("task mailbox append was not durable")
delivered = True

View file

@ -135,12 +135,15 @@ def test_admission_freezes_reviewed_behavior_runtime_digests_and_authority(tmp_p
assert admission.destination == binding.destination
assert admission.capability_ceiling.skill_name == "community-helper"
# The reviewed profile selected chat_history; the rest is the constant
# cognitive baseline every admitted conversation carries.
# cognitive baseline plus the own-work baseline every new ceiling carries.
assert [grant.name for grant in admission.capability_ceiling.tool_grants] == [
"chat_history",
"get_task_result",
"knowledge_list",
"knowledge_read",
"knowledge_write",
"recent_tasks",
"steer_task",
"update_identity",
"update_scratchpad",
]

View file

@ -59,16 +59,27 @@ def test_ceiling_compiles_exact_tools_scripts_resources_and_digest():
)
# The profile selected chat_history and one script; the cognitive baseline
# (own memory, no new authority) is compiled in beside them.
# (own memory, no new authority) and the own-work baseline (this binding's
# readers, host-bound to its scope, and steer_task) are compiled in beside them.
assert [grant.name for grant in ceiling.tool_grants] == [
"chat_history",
"get_task_result",
"knowledge_list",
"knowledge_read",
"knowledge_write",
"recent_tasks",
"skill_exec",
"steer_task",
"update_identity",
"update_scratchpad",
]
scoped = {grant.name: [(item.argument_path, item.static_value) for item in grant.bindings]
for grant in ceiling.tool_grants if grant.name in {"get_task_result", "recent_tasks", "steer_task"}}
assert scoped == {
"get_task_result": [(("presence_scope",), "own_binding")],
"recent_tasks": [(("presence_scope",), "own_binding")],
"steer_task": [],
}
script = next(grant for grant in ceiling.tool_grants if grant.name == "skill_exec")
assert [(item.argument_path, item.static_value) for item in script.bindings] == [
(("skill",), "calendar"),
@ -230,9 +241,12 @@ def test_registry_filters_schema_dispatch_and_resolved_targets(tmp_path):
"presence_cancel_work",
"read_file",
"chat_history",
"get_task_result",
"knowledge_list",
"knowledge_read",
"knowledge_write",
"recent_tasks",
"steer_task",
"update_identity",
"update_scratchpad",
}

View file

@ -54,11 +54,14 @@ def _ceiling(*selections):
)
_OWN_WORK = {"get_task_result", "recent_tasks", "steer_task"}
def test_profile_without_tool_selections_still_carries_its_own_memory():
ceiling = _ceiling()
assert [grant.name for grant in ceiling.tool_grants] == sorted(COGNITIVE_MEMORY_TOOL_NAMES)
assert all(grant.bindings == () for grant in ceiling.tool_grants)
assert [grant.name for grant in ceiling.tool_grants] == sorted(COGNITIVE_MEMORY_TOOL_NAMES | _OWN_WORK)
assert all(grant.bindings == () for grant in ceiling.tool_grants if grant.name in COGNITIVE_MEMORY_TOOL_NAMES)
for name in COGNITIVE_MEMORY_TOOL_NAMES:
assert presence_ceiling_allows_tool(ceiling, name)
@ -84,7 +87,7 @@ def test_a_selected_baseline_tool_keeps_the_profile_authored_bindings():
assert [(item.argument_path, item.static_value) for item in grant.bindings] == [
(("scope",), "global"),
]
assert [item.name for item in ceiling.tool_grants] == sorted(COGNITIVE_MEMORY_TOOL_NAMES)
assert [item.name for item in ceiling.tool_grants] == sorted(COGNITIVE_MEMORY_TOOL_NAMES | _OWN_WORK)
ctx = ToolContext(
repo_dir=None,
@ -137,7 +140,7 @@ def test_admitted_external_turn_writes_global_knowledge_and_nothing_else(tmp_pat
_select_history(data, skill_dir)
admission = _admit(data, _binding(data))
assert [grant.name for grant in admission.capability_ceiling.tool_grants] == sorted(
COGNITIVE_MEMORY_TOOL_NAMES
COGNITIVE_MEMORY_TOOL_NAMES | _OWN_WORK
)
seen: dict[str, object] = {}

View file

@ -210,18 +210,35 @@ def test_ordinary_empty_and_failed_silent_outcomes_remain_failed():
assert outcome["outcome_axes"]["execution"]["status"] in {"failed", "infra_failed"}
def test_pending_children_still_require_absorption(turn, tmp_path):
_DECLARED = json.dumps({"delivery_control": "replace", "full_answer": "Best available; child1 still running",
"presence_finish": {"outcome": "message", "message": "Here is what I have so far."}})
@pytest.mark.parametrize("forced,outcome,spoken", [
# Owner Q4: the forced answer is the internal record; undeclared prose never becomes speech.
("Best available; child1 still running", "silent", ""),
(_DECLARED, "message", "Here is what I have so far."),
])
def test_pending_children_still_require_absorption(turn, tmp_path, forced, outcome, spoken):
from ouroboros.task_results import write_task_result, STATUS_RUNNING
_registry, calls, run = turn
# A real turn carries its Presence metadata; the ceiling alone does not arm the protocol.
_registry._ctx.task_metadata = {"presence": {"binding_id": "a" * 32, "event": {"conversation_key": "k"}}}
write_task_result(tmp_path, "child1", STATUS_RUNNING, parent_task_id="parent1",
root_task_id="parent1", delegation_role="subagent", role="reviewer", result="Still running")
text, usage, _trace = run([
_call("silent", "Premature"),
{"content": '{"delivery_control":"keep"}'},
{"content": "Best available"}, {"content": "Best available"},
{"content": "Best available"}, {"content": forced},
])
assert len(calls) > 1
assert usage["reason_code"] == "children_unabsorbed"
assert "presence_completion_outcome" not in usage
assert build_presence_result_event({"id": "parent1"}, text, _registry._ctx, terminal_origin=usage.get("terminal_origin", ""))["outcome"] == "message"
assert text == "Best available; child1 still running" # the internal record keeps the child facts
assert "[PRESENCE_DELIVERY]" in str(calls[-1][-1]["content"])
assert "name the unabsorbed" not in str(calls[-1][-1]["content"])
task = {"id": "parent1"}
result = build_presence_result_event(task, text, _registry._ctx, terminal_origin=usage.get("terminal_origin", ""))
assert (result["outcome"], result["text"]) == (outcome, spoken)
assert task["metadata"]["presence_declaration"]["status"] == ("declared" if spoken else "missing")

View file

@ -0,0 +1,292 @@
"""A Presence forced final keeps its internal record apart from what the conversation receives.
Owner Q4: execution, review and helper diagnostics stay with the owner. The ONE
forced model call declares its outward delivery beside the record; without a valid
declaration nothing new is spoken, and a declared useful reply is spoken even when
the run itself failed. Deterministic fake-model replay; no transport sends anything.
"""
from __future__ import annotations
import json
import queue
from types import SimpleNamespace
import pytest
from starlette.testclient import TestClient
from ouroboros import agent_task_pipeline as pipeline, loop
from ouroboros.gateway.host_service import create_host_service_app
from ouroboros.presence_authority import presence_ceiling_payload
from ouroboros.presence_runner import _cached_result
from ouroboros.task_results import load_task_result
from ouroboros.tools.registry import ToolRegistry
from ouroboros.utils import append_jsonl
from tests.test_host_service_api import _seed_presence_behavior, _seed_token
from tests.test_presence_completion import _call
from tests.test_presence_runner import _admission
KEY = "telegram:bot-1:room-1:0"
RECORD = "Internal record: helper child-7 failed with provider 400; review not run; figures verified for Q1 only."
def _presence(binding="a" * 32, version=1):
return {"binding_id": binding, "delivery_reporting_version": version,
"event": {"conversation_key": KEY, "conversation_id": "room-1"}}
def _forced(outcome=None, message=None, **extra):
body = {"delivery_control": "replace", "full_answer": RECORD, **extra}
if outcome is not None:
body["presence_finish"] = {"outcome": outcome, **({"message": message} if message is not None else {})}
return json.dumps(body)
def _read():
return {"role": "assistant", "content": None, "tool_calls": [{
"id": "read", "type": "function", "function": {"name": "chat_history", "arguments": "{}"},
}]}
def _run(root, monkeypatch, forced, *, task=None, presence=True, handoff=None, first=None, ceiling=None,
metadata=None, traces=None):
"""Round limit after one tool round, then the ONE forced call; real pipeline after it."""
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.setenv("OUROBOROS_MAX_ROUNDS", "1")
registry = ToolRegistry(repo_dir=root, drive_root=root)
ctx = registry._ctx
ctx.is_direct_chat = True
ctx.task_metadata = {"inline_max_rounds": 1, **({"presence": _presence()} if presence else {}), **(metadata or {})}
if presence if ceiling is None else ceiling:
ctx.task_contract = {"capability_ceiling": presence_ceiling_payload(_admission().capability_ceiling)}
if handoff:
ctx._swarm_handoff_attempt = handoff
registry.override_handler("chat_history", lambda *_a, **_kw: "Synthetic history")
calls, replies = [], iter([first or _read(), {"content": forced}])
def respond(_llm, messages, *_a, **_k):
calls.append([dict(row) for row in messages])
return next(replies), 0.0
monkeypatch.setattr(loop, "call_llm_with_retry", respond)
task = task or {"id": "presence-loop", "type": "presence", "_presence_turn": True, "chat_id": 7,
"text": "Please help", "metadata": {"presence": _presence()}}
task["_skip_post_task_synthesis"] = True
text, usage, trace = loop.run_llm_loop(
[{"role": "user", "content": "Please help"}], registry,
SimpleNamespace(default_model=lambda: "test-model"), root / "logs",
lambda *_a, **_kw: None, queue.Queue(), task_id=task["id"], drive_root=root,
)
events = []
pipeline.emit_task_results(SimpleNamespace(drive_root=root, repo_dir=root), None, None,
events, task, text, usage, trace, 0.0, root / "logs", ctx=ctx)
result = next((row for row in events if row["type"] == "presence_result"), None)
if traces is not None:
traces.append((trace, events))
return result, load_task_result(root, task["id"]), calls, text
@pytest.mark.parametrize("forced,outcome,spoken,status", [
(_forced("message", "Q1 figures are ready; Q2 is still coming."), "message",
"Q1 figures are ready; Q2 is still coming.", "declared"),
(_forced("silent", ""), "silent", "", "declared"),
(_forced("tool_delivered", "sent the table via the transport tool"), "tool_delivered", "", "declared"),
(RECORD, "silent", "", "missing"), # untyped internal prose is never speech
(_forced("message", " "), "silent", "", "invalid"),
(_forced("maybe", "hi"), "silent", "", "invalid"),
(_forced("silent", "but also this"), "silent", "", "invalid"),
(_forced("deferred", "on it"), "silent", "", "invalid"), # nothing was scheduled
(json.dumps({"delivery_control": "replace", "full_answer": RECORD,
"presence_finish": {"outcome": "message", "message": "x", "to": "y"}}), "silent", "", "invalid"),
], ids=["message", "silent", "tool_delivered", "prose", "blank_message", "unknown_outcome", "silent_with_text",
"unscheduled_deferred", "extra_key"])
def test_forced_final_speaks_only_what_it_declares(tmp_path, monkeypatch, forced, outcome, spoken, status):
result, stored, calls, text = _run(tmp_path, monkeypatch, forced)
assert len(calls) == 2 # the one forced call; no repair or polishing round
assert "[PRESENCE_DELIVERY]" in str(calls[-1][-1]["content"])
assert (result["outcome"], result["text"]) == (outcome, spoken)
assert stored["reason_code"] == "round_limit"
assert stored["outcome_axes"]["execution"]["status"] != "ok" # speech is declared, not read off status
assert stored["metadata"]["presence_declaration"]["status"] == status
assert stored["metadata"]["presence_result_text"] == spoken
assert _cached_result(tmp_path, "presence-loop").text == spoken
assert text == RECORD and stored["result"].startswith(RECORD) # the record survives beside it
if outcome == "tool_delivered":
assert result["message"] == "sent the table via the transport tool" # a note, never speech
def test_a_malformed_control_body_speaks_nothing_and_keeps_the_host_fallback(tmp_path, monkeypatch):
duplicate = ('{"delivery_control": "replace", "full_answer": "%s", "presence_finish": {"outcome": "silent"}, '
'"presence_finish": {"outcome": "message", "message": "dup"}}' % RECORD)
result, stored, calls, _text = _run(tmp_path, monkeypatch, duplicate)
assert len(calls) == 2
assert (result["outcome"], result["text"]) == ("silent", "")
assert stored["terminal_origin"] != "model_final" # the host fallback is never spoken
assert stored["metadata"]["presence_declaration"] == {"status": "invalid",
"reason": "the forced envelope repeats a key"}
def test_early_acknowledgement_is_shown_and_does_not_suppress_the_result(tmp_path, monkeypatch):
logs = tmp_path / "logs"
append_jsonl(logs / "chat.jsonl", {"task_id": "presence-loop", "direction": "in", "text": "Please help"})
append_jsonl(logs / "chat.jsonl", {
"task_id": "presence-loop", "type": "presence_delivery", "text": "Looking into it now.",
"transport": {"conversation_key": KEY, "delivery": {"state": "delivered", "delivery_id": "d1", "part_id": "0"}},
})
append_jsonl(logs / "chat.jsonl", {
"task_id": "presence-loop", "type": "presence_delivery", "text": "Partial table",
"transport": {"conversation_key": KEY, "delivery": {"state": "uncertain", "delivery_id": "d2", "part_id": "0"}},
})
result, _stored, calls, _text = _run(tmp_path, monkeypatch, _forced("message", "Here are the Q1 figures."))
prompt = str(calls[-1][-1]["content"])
assert '"Looking into it now."' in prompt and "1 more part(s) have an uncertain outcome" in prompt
assert "An early acknowledgement is not the promised result" in prompt
assert (result["outcome"], result["text"]) == ("message", "Here are the Q1 figures.")
def test_declared_useful_partial_survives_failure_and_owed_child_stays_pollable(tmp_path, monkeypatch):
handoff = {"status": "scheduled", "task_id": "later-work"}
result, stored, calls, _text = _run(tmp_path, monkeypatch, _forced("deferred", "Started the full audit."),
handoff=handoff)
assert "deferred = acknowledge work that was actually scheduled (it was)" in str(calls[-1][-1]["content"])
assert (result["outcome"], result["text"], result["work_ref"]) == ("deferred", "Started the full audit.", "later-work")
silent, _stored, _calls, _text = _run(tmp_path / "silent", monkeypatch, RECORD, handoff=handoff)
# No declaration: nothing new is said, but the admitted child is still owed.
assert (silent["outcome"], silent["text"], silent["work_ref"]) == ("deferred", "", "later-work")
def test_ordinary_forced_final_is_unchanged(tmp_path, monkeypatch):
task = {"id": "owner-task", "type": "task", "chat_id": 1, "text": "Please help"}
_result, stored, calls, text = _run(tmp_path, monkeypatch, RECORD, task=task, presence=False)
assert "[PRESENCE_DELIVERY]" not in str(calls[-1][-1]["content"])
assert text == RECORD and stored["result"].startswith(RECORD)
assert "presence_declaration" not in (stored.get("metadata") or {})
def test_promoted_work_result_reaches_the_work_endpoint_as_declared(tmp_path, monkeypatch):
_seed_token(tmp_path, skill="telegram-bot", token="presence-token",
permissions=["presence"], manifest_permissions=["presence"])
binding = _seed_presence_behavior(tmp_path)
task = {"id": "promoted-9", "type": "task", "chat_id": 7, "text": "Compile", "delegation_role": "root",
"root_task_id": "promoted-9", "metadata": {"presence": _presence(binding)}}
_run(tmp_path, monkeypatch, _forced("message", "The figures you asked for: 41 and 43."), task=task)
with TestClient(create_host_service_app(tmp_path)) as client:
body = client.get("/presence/work/promoted-9", params={"binding_id": binding},
headers={"X-Skill-Token": "presence-token"}).json()
assert (body["outcome"], body["text"]) == ("message", "The figures you asked for: 41 and 43.")
assert RECORD not in json.dumps(body)
@pytest.mark.parametrize("feedback", [True, False])
def test_invalidated_finish_is_named_void_with_this_tasks_confirmed_sends(tmp_path, monkeypatch, feedback):
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
registry._ctx.task_contract = {"capability_ceiling": presence_ceiling_payload(_admission().capability_ceiling)}
registry._ctx.task_metadata = {"presence": _presence()}
registry._ctx.is_direct_chat = True
logs = tmp_path / "logs"
append_jsonl(logs / "chat.jsonl", {"task_id": "parent1", "direction": "in", "text": "Please help"})
append_jsonl(logs / "chat.jsonl", {
"task_id": "parent1", "type": "presence_delivery", "text": "The full answer, sent by tool.",
"transport": {"conversation_key": KEY, "delivery": {"state": "delivered", "delivery_id": "d1", "part_id": "0"}},
})
reviews, calls = [], []
def review(**kwargs):
reviews.append(kwargs["content"])
if feedback and len(reviews) == 1: # the panel's feedback is new transcript content
kwargs["messages"].append({"role": "user", "content": "[REVIEW] Also confirm the Q2 figures."})
return len(reviews) == 1 # the first finish is held for another pass
def respond(_llm, messages, *_a, **_k):
calls.append([dict(row) for row in messages])
return ([_call("tool_delivered", ""), _call("tool_delivered", "")][len(calls) - 1]), 0.0
monkeypatch.setattr(loop, "_run_task_acceptance_review_once", review)
monkeypatch.setattr(loop, "call_llm_with_retry", respond)
_text, usage, _trace = loop.run_llm_loop(
[{"role": "user", "content": "Please help"}], registry,
SimpleNamespace(default_model=lambda: "test-model"), logs,
lambda *_a, **_kw: None, queue.Queue(), task_id="parent1", drive_root=tmp_path,
)
assert len(calls) == 2
notes = [row for row in calls[1] if "[PRESENCE_FINISH_NOT_ACCEPTED]" in str(row.get("content"))]
assert "[PRESENCE_FINISH_NOT_ACCEPTED]" not in json.dumps(calls[0])
assert usage["presence_completion_outcome"] == "tool_delivered" # the fresh finish is accepted
if not feedback: # nothing new was said: the turn parks as before, with no added reminder
assert notes == []
return
assert len(notes) == 1
assert '"The full answer, sent by tool."' in str(notes[0]["content"])
# --- repair pass: arming identity, duplicate evidence, internal notes --------------
def test_a_child_inheriting_only_the_ceiling_keeps_its_ordinary_forced_final(tmp_path, monkeypatch):
task = {"id": "child-3", "type": "task", "chat_id": 7, "text": "Check Q1", "delegation_role": "subagent",
"parent_task_id": "promoted-9", "root_task_id": "promoted-9"}
traces = []
result, stored, calls, text = _run(tmp_path, monkeypatch, RECORD, task=task, presence=False, ceiling=True,
metadata={"delegation_role": "subagent", "parent_task_id": "promoted-9"},
traces=traces)
assert "[PRESENCE_DELIVERY]" not in str(calls[-1][-1]["content"]) # a child answers its parent
assert result is None and text == RECORD and stored["result"].startswith(RECORD)
assert "presence_declaration" not in (stored.get("metadata") or {})
assert [event["type"] for event in traces[0][1]].count("presence_result") == 0
def test_duplicate_subject_evidence_survives_the_presence_arm_and_declares_nothing(tmp_path, monkeypatch):
duplicated = ('{"delivery_control": "replace", "full_answer": "%s", "acceptance_subject": '
'{"owner_source_sha256": "aaa", "owner_source_sha256": "bbb"}, '
'"presence_finish": {"outcome": "message", "message": "Here are the figures."}}' % RECORD)
traces = []
result, stored, calls, text = _run(tmp_path, monkeypatch, duplicated, traces=traces)
assert len(calls) == 2
# The resolver saw the original bytes: the ambiguous subject is refused as a duplicate,
# not silently collapsed to its last value and judged as a different source.
assert traces[0][0]["forced_acceptance_subject"] == {
"applied": False, "reason": "acceptance_subject requires an exact owner source and optional criteria/tool indices"}
assert (result["outcome"], result["text"]) == ("silent", "")
assert stored["metadata"]["presence_declaration"] == {"status": "invalid",
"reason": "the forced envelope repeats a key"}
assert text == RECORD # the record keeps its ordinary degraded-subject handling
def test_a_duplicate_inside_the_declaration_voids_only_the_declaration(tmp_path, monkeypatch):
body = ('{"delivery_control": "replace", "full_answer": "%s", '
'"presence_finish": {"outcome": "message", "message": "hi", "outcome": "silent"}}' % RECORD)
result, stored, _calls, text = _run(tmp_path, monkeypatch, body)
assert (result["outcome"], result["text"]) == ("silent", "")
assert stored["metadata"]["presence_declaration"]["status"] == "invalid"
assert text == RECORD and stored["terminal_origin"] == "model_final"
def test_a_tool_delivered_note_is_never_speech_even_when_owed_work_defers_the_turn(tmp_path, monkeypatch):
note = "sent the table via the transport tool; helper child-7 failed"
handoff = {"status": "scheduled", "task_id": "later-work"}
result, stored, _calls, _text = _run(tmp_path, monkeypatch, _forced("tool_delivered", note), handoff=handoff)
assert (result["outcome"], result["text"], result["work_ref"]) == ("deferred", "", "later-work")
assert result["message"] == note # context for the next turn, never a reply body
assert stored["metadata"]["presence_result_text"] == ""
replay = _cached_result(tmp_path, "presence-loop")
assert (replay.outcome, replay.text, replay.work_ref) == ("deferred", "", "later-work")
assert note not in json.dumps(stored["metadata"])
def test_forced_prompt_and_facts_read_the_canonical_root_on_a_forked_drive(tmp_path, monkeypatch):
canonical = tmp_path / "canonical"
append_jsonl(canonical / "logs" / "chat.jsonl", {"task_id": "presence-loop", "direction": "in", "text": "x"})
append_jsonl(canonical / "logs" / "chat.jsonl", {
"task_id": "presence-loop", "type": "presence_delivery", "text": "Canonical receipt",
"transport": {"conversation_key": KEY, "delivery": {"state": "delivered", "delivery_id": "d1", "part_id": "0"}},
})
append_jsonl(tmp_path / "logs" / "chat.jsonl", {"task_id": "presence-loop", "direction": "in", "text": "x"})
_result, _stored, calls, _text = _run(tmp_path, monkeypatch, _forced("silent", ""),
metadata={"budget_drive_root": str(canonical)})
prompt = str(calls[-1][-1]["content"])
assert 'Sends confirmed for this task so far: "Canonical receipt"' in prompt

View file

@ -0,0 +1,554 @@
"""A Presence binding's own work: scoped discovery, exact reads and control (owner Q1-Q3).
Related work is independent work started from the same nonempty binding id, from
any of its conversations or threads. Other bindings, owner roots, inline turns,
delegated children and rows without Presence provenance are never attributed.
"""
from __future__ import annotations
import json
import types
import pytest
from ouroboros.presence_authority import (
PresenceAuthorityError,
build_presence_capability_ceiling,
presence_ceiling_from_payload,
presence_ceiling_payload,
)
from ouroboros.presence_capabilities import PresenceToolTarget
from ouroboros.task_results import load_task_result, write_task_result
from ouroboros.tools.registry import ToolContext, ToolRegistry
from ouroboros.utils import atomic_write_json
from tests.test_presence_authority import _resolution
BINDING = "a" * 32
OTHER = "b" * 32
HERE = "slack:T1:D1:0"
THREAD = "slack:T1:D1:1712.5"
ROOM = "slack:T1:C9:0"
def _presence(binding=BINDING, key=HERE):
provider, account, conversation, thread = key.split(":")
return {"binding_id": binding, "event": {
"conversation_key": key, "provider": provider, "account_id": account,
"conversation_id": conversation, "thread_id": "" if thread == "0" else thread,
"source_event_id": f"evt-{conversation}-{thread}",
}}
def _work(root, task_id, status, *, binding=BINDING, key=HERE, **fields):
metadata = {"presence": _presence(binding, key)} if binding else {}
fields.setdefault("delegation_role", "root")
write_task_result(root, task_id, status, metadata=metadata, description=f"goal of {task_id}", **fields)
def _ceiling(*targets):
return build_presence_capability_ceiling(
skill_name="community-helper", skill_content_hash="c" * 64,
state_fingerprint="d" * 64, resolution=_resolution(*targets),
)
def _registry(root, ceiling=None, *, binding=BINDING, key=HERE, task_id="presence-turn-1"):
ctx = ToolContext(
repo_dir=root, drive_root=root, task_id=task_id,
task_contract={"capability_ceiling": presence_ceiling_payload(ceiling or _ceiling())},
task_metadata={"presence": _presence(binding, key)},
)
registry = ToolRegistry(repo_dir=root, drive_root=root)
registry.set_context(ctx)
return registry, ctx
def _queue(root, *, pending=(), running=()):
atomic_write_json(root / "state" / "queue_snapshot.json", {
"pending": [{"id": task["id"], "task": task} for task in pending],
"running": [{"id": task["id"], "task": task} for task in running],
})
def _page_ids(registry, **args):
"""Every id reachable by following ``next`` from the first page."""
seen, page = [], json.loads(registry.execute("recent_tasks", {"limit": 2, **args}))
while True:
assert "error" not in page, page
seen += [row["task_id"] for row in page["tasks"]]
if not page["next"]:
return seen, page
page = json.loads(registry.execute("recent_tasks", page["next"]))
def test_same_binding_work_from_other_threads_is_paged_and_nothing_else(tmp_path):
_work(tmp_path, "queued-here", "scheduled")
_work(tmp_path, "running-thread", "running", key=THREAD)
_work(tmp_path, "done-room", "completed", key=ROOM, result="The report is ready.")
_work(tmp_path, "done-here", "completed", result="Earlier answer")
_work(tmp_path, "failed-thread", "failed", key=THREAD)
# Never attributed: another binding, the owner's root, an inline turn,
# a delegated child, a row with no provenance at all.
_work(tmp_path, "foreign", "running", binding=OTHER)
_work(tmp_path, "owner-root", "running", binding="")
_work(tmp_path, "presence-inline", "completed", delegation_role=None)
_work(tmp_path, "child", "running", delegation_role="subagent", parent_task_id="running-thread")
write_task_result(tmp_path, "unknown", "completed", description="no provenance")
# A legacy row whose canonical record predates provenance is established by
# the queue's own task metadata; a queue claim never overrides another binding.
write_task_result(tmp_path, "legacy-queued", "scheduled", delegation_role="root")
_work(tmp_path, "conflict", "scheduled", binding=OTHER)
queue_rows = [
{"id": "legacy-queued", "delegation_role": "root", "metadata": {"presence": _presence(key=ROOM)}},
{"id": "conflict", "delegation_role": "root", "metadata": {"presence": _presence()}},
{"id": "queue-only", "delegation_role": "root", "description": "not yet recorded",
"metadata": {"presence": _presence(key=THREAD)}},
]
running_rows = [
{"id": "running-thread", "delegation_role": "root", "description": "thread work",
"metadata": {"presence": _presence(key=THREAD)}},
{"id": "foreign", "delegation_role": "root", "metadata": {"presence": _presence(OTHER)}},
]
_queue(tmp_path, pending=queue_rows, running=running_rows)
registry, _ctx = _registry(tmp_path)
ids, last = _page_ids(registry) # the model supplies no scope: the host binds it
assert sorted(ids) == sorted([
"queued-here", "running-thread", "done-room", "done-here", "failed-thread",
"legacy-queued", "queue-only",
])
assert len(ids) == len(set(ids)) and ids[0] == "queue-only" # queued work without a row leads
assert last["presence_scope"] == {"scope": "own_binding", "binding_id": BINDING}
assert [row["task_id"] for row in last["running"]] == ["running-thread"]
first = json.loads(registry.execute("recent_tasks", {"limit": 20}))
by_id = {row["task_id"]: row for row in first["tasks"]}
assert by_id["done-room"]["presence_origin"]["conversation_key"] == ROOM
assert by_id["done-room"]["result_preview"] == "The report is ready."
assert by_id["queue-only"] == {"task_id": "queue-only", "status": "pending",
"description": "not yet recorded", "source": "queue_snapshot"}
# A changed inventory never continues an old cursor into a mixed page.
stale = json.loads(registry.execute("recent_tasks", {"limit": 2}))
_work(tmp_path, "new-work", "scheduled", key=ROOM)
moved = json.loads(registry.execute("recent_tasks", stale["next"]))
assert moved["error"]["code"] == "RECENT_TASKS_SNAPSHOT_CHANGED" and moved["tasks"] == []
def test_an_empty_or_foreign_binding_attributes_nothing_and_owner_reads_stay_whole(tmp_path):
_work(tmp_path, "mine", "completed")
_work(tmp_path, "theirs", "completed", binding=OTHER)
registry, _ctx = _registry(tmp_path, binding="")
refused = json.loads(registry.execute("recent_tasks", {}))
assert refused["error"]["code"] == "PRESENCE_SCOPE_UNAVAILABLE"
owner = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
owner.set_context(ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="owner-turn"))
everything = json.loads(owner.execute("recent_tasks", {"limit": 20}))
assert {row["task_id"] for row in everything["tasks"]} == {"mine", "theirs"}
assert "presence_scope" not in everything
def test_exact_read_admits_own_work_and_tree_and_refuses_the_rest(tmp_path):
_work(tmp_path, "done-room", "completed", key=ROOM, result="Full report text")
_work(tmp_path, "foreign", "completed", binding=OTHER, result="Not yours")
_work(tmp_path, "owner-root", "completed", binding="", result="Owner work")
write_task_result(tmp_path, "my-child", "completed", parent_task_id="presence-turn-1",
root_task_id="presence-turn-1", delegation_role="subagent", result="Child result")
registry, _ctx = _registry(tmp_path)
assert "Full report text" in registry.execute("get_task_result", {"task_id": "done-room"})
assert "Child result" in registry.execute("get_task_result", {"task_id": "my-child"})
for task_id in ("foreign", "owner-root", "never-existed"):
refused = registry.execute("get_task_result", {"task_id": task_id})
assert "is not independent work started from this Presence binding" in refused and "Not yours" not in refused
# A profile that explicitly selected the global readers keeps that grant.
selected = _ceiling(PresenceToolTarget("builtin", "get_task_result"),
PresenceToolTarget("builtin", "recent_tasks"))
global_reader, _ctx = _registry(tmp_path, selected)
assert "Not yours" in global_reader.execute("get_task_result", {"task_id": "foreign"})
listed = json.loads(global_reader.execute("recent_tasks", {"limit": 20}))
assert {"foreign", "owner-root"} <= {row["task_id"] for row in listed["tasks"]}
# ...and the model may still narrow it on purpose.
narrowed = json.loads(global_reader.execute("recent_tasks", {"limit": 20, "presence_scope": "own_binding"}))
assert [row["task_id"] for row in narrowed["tasks"]] == ["done-room"]
def test_old_frozen_ceilings_keep_their_digest_and_gain_nothing(tmp_path):
payload = presence_ceiling_payload(_ceiling())
old = json.loads(json.dumps(payload))
old["tools"] = [tool for tool in old["tools"]
if tool["name"] not in {"get_task_result", "recent_tasks", "steer_task"}]
with pytest.raises(PresenceAuthorityError):
presence_ceiling_from_payload(old) # stripping the grants is not a valid frozen ceiling
from ouroboros.presence_authority import _digest
old["digest"] = _digest(old) # a genuinely older ceiling, compiled before the baseline
ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="presence-old",
task_contract={"capability_ceiling": old},
task_metadata={"presence": _presence()})
registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
registry.set_context(ctx)
names = {schema["function"]["name"] for schema in registry.schemas()}
assert not names & {"get_task_result", "recent_tasks", "steer_task"}
assert "PRESENCE_CAPABILITY_BLOCKED" in registry.execute("steer_task", {"task_id": "x", "message": "y"})
assert "presence_cancel_work" in names # the intrinsic control is unchanged
# --- steering: task-authored, own binding only, pending included ------------------
def _supervisor(root, *, pending=(), running=()):
return types.SimpleNamespace(
DRIVE_ROOT=root, PENDING=list(pending), bridge=None,
RUNNING={task["id"]: {"task": task, "started_at": 1.0} for task in running},
send_with_budget=lambda *_a, **_k: None,
)
def _steering_turn(root, supervisor_ctx, emitted):
from supervisor.events import _handle_steer_task
def _dispatch(event):
emitted.append(event)
_handle_steer_task(event, supervisor_ctx)
return types.SimpleNamespace(
pending_events=[], event_queue=types.SimpleNamespace(put_nowait=_dispatch),
current_chat_id=4242, drive_root=root, task_id="presence-turn-1", is_direct_chat=True,
task_metadata={"presence": _presence(), "source": "presence", "client_message_id": "evt-D1-0"},
last_owner_delivery=None,
)
def test_presence_steer_reaches_pending_and_running_own_work_as_task_authored_text(tmp_path, monkeypatch):
import supervisor.queue as queue
from ouroboros.owner_mailbox import deliver_task_message, drain_owner_entries
from ouroboros.tools.control import _steer_task
monkeypatch.setattr(queue, "DRIVE_ROOT", str(tmp_path))
queued = {"id": "queued-work", "delegation_role": "root", "chat_id": 77,
"metadata": {"presence": _presence(key=THREAD)}}
running = {"id": "running-work", "delegation_role": "root", "chat_id": 78,
"metadata": {"presence": _presence(key=ROOM)}}
foreign = {"id": "foreign-work", "delegation_role": "root", "chat_id": 79,
"metadata": {"presence": _presence(OTHER)}}
owner = {"id": "owner-work", "delegation_role": "root", "chat_id": 1, "metadata": {}}
fence = {"root_task_id": "running-work", "status": "active", "owner_message_generation": 3}
monkeypatch.setitem(queue.ACCEPTANCE_FENCES, "running-work", fence)
emitted = []
turn = _steering_turn(tmp_path, _supervisor(tmp_path, pending=[queued, foreign], running=[running, owner]),
emitted)
exact = "Alex says: use the March figures, not February."
for target in ("queued-work", "running-work"):
out = _steer_task(turn, target, exact)
assert "written to its mailbox" in out and "not as owner text" in out
[entry] = drain_owner_entries(tmp_path, target)
assert entry["text"] == exact
assert (entry["provenance"], entry["source_task_id"]) == ("independent_task", "presence-turn-1")
# The run's origin rides beside the words; it is never offered as their author.
assert entry["sender_origin"] == {"provider": "slack", "account_id": "T1", "conversation_id": "D1",
"source_event_id": "evt-D1-0"}
rendered = []
deliver_task_message(entry, target, None, rendered.append)
assert rendered[0].startswith("[Message from independent task presence-turn-1; that task's run started from ")
assert "does not make it the author of any words it quotes]" in rendered[0]
assert rendered[0].endswith("\n" + exact) and "[Message from my human]" not in rendered[0]
assert fence["owner_message_generation"] == 3 # a task's words supersede no reviewed answer
assert all(evt["presence_binding_id"] == BINDING for evt in emitted)
assert all(evt["issuer"]["kind"] == "task" for evt in emitted)
for target in ("foreign-work", "owner-work"):
refused = _steer_task(turn, target, "stop")
assert "STEER_REJECTED" in refused and "presence_work_not_related" in refused
assert drain_owner_entries(tmp_path, target) == []
def test_an_ordinary_task_still_messages_any_listed_root(tmp_path, monkeypatch):
import supervisor.queue as queue
from ouroboros.owner_mailbox import drain_owner_entries
from ouroboros.tools.control import _steer_task
monkeypatch.setattr(queue, "DRIVE_ROOT", str(tmp_path))
owner = {"id": "owner-work", "delegation_role": "root", "chat_id": 1, "metadata": {}}
emitted = []
turn = _steering_turn(tmp_path, _supervisor(tmp_path, running=[owner]), emitted)
turn.task_metadata = {}
turn.task_id = "managed-root"
assert "written to its mailbox" in _steer_task(turn, "owner-work", "status please")
assert "presence_binding_id" not in emitted[0] and "sender_origin" not in emitted[0]
[entry] = drain_owner_entries(tmp_path, "owner-work")
assert entry["provenance"] == "independent_task" and "sender_origin" not in entry
# --- cancellation: request receipts, own binding only, before any intent -----------
def _cancel_ctx(root, *, binding=BINDING, key=HERE):
return types.SimpleNamespace(
pending_events=[], event_queue=None, drive_root=root, task_id="presence-turn-1",
task_metadata={"presence": _presence(binding, key)}, current_chat_id=4242,
task_contract={"capability_ceiling": presence_ceiling_payload(_ceiling())},
)
def _intents(root):
path = root / "state" / "cancel_intents.json"
return json.loads(path.read_text(encoding="utf-8")) if path.exists() else {}
def test_presence_cancel_requests_own_pending_work_from_another_thread(tmp_path):
from ouroboros.tools.presence import get_tools
_work(tmp_path, "queued-thread", "scheduled", key=THREAD, root_task_id="queued-thread")
_queue(tmp_path, pending=[{"id": "queued-thread", "delegation_role": "root",
"metadata": {"presence": _presence(key=THREAD)}}])
cancel = next(item for item in get_tools() if item.name == "presence_cancel_work").handler
out = cancel(_cancel_ctx(tmp_path), "queued-thread", "the person withdrew the request")
assert out.startswith("Cancel requested: queued-thread")
assert "cancel_state=pending" in out # a request receipt, never a claim that it stopped
assert "queued-thread" in json.dumps(_intents(tmp_path))
assert load_task_result(tmp_path, "queued-thread")["status"] == "scheduled"
def test_selected_cancel_and_forward_refuse_foreign_work_before_any_effect(tmp_path):
from ouroboros.tools.join_ledger import _cancel_task
_work(tmp_path, "foreign", "running", binding=OTHER)
_work(tmp_path, "owner-root", "running", binding="")
_work(tmp_path, "mine", "running", key=ROOM)
ctx = _cancel_ctx(tmp_path)
for target in ("foreign", "owner-root"):
refused = _cancel_task(ctx, target, "stop")
assert "is not independent work started from this Presence binding" in refused
assert _intents(tmp_path) == {}
selected = _ceiling(PresenceToolTarget("builtin", "forward_to_worker"))
registry, _ctx = _registry(tmp_path, selected)
assert "is not independent work started from this Presence binding" in registry.execute(
"forward_to_worker", {"task_id": "foreign", "message": "stop"})
assert "is not independent work started from this Presence binding" not in registry.execute(
"forward_to_worker", {"task_id": "mine", "message": "new fact"})
# --- canonical provenance from admission; the work endpoint and context read it ----
def test_scheduled_presence_promotion_is_canonical_and_pollable_while_queued(tmp_path, monkeypatch):
from starlette.testclient import TestClient
import supervisor.workers as workers
from ouroboros.gateway.host_service import create_host_service_app
from supervisor.events import _handle_promote_chat_to_task
from tests.test_host_service_api import _seed_presence_behavior, _seed_token
monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path)
_seed_token(tmp_path, skill="telegram-bot", token="presence-token",
permissions=["presence"], manifest_permissions=["presence"])
binding = _seed_presence_behavior(tmp_path)
pending = []
handler_ctx = types.SimpleNamespace(
DRIVE_ROOT=tmp_path, WORKERS={0: types.SimpleNamespace()}, PENDING=pending, bridge=None,
append_jsonl=lambda *_a, **_k: None, persist_queue_snapshot=lambda **_k: True,
enqueue_task=lambda task: pending.append(dict(task)) or pending[-1],
load_state=lambda: {"owner_chat_id": 1},
)
presence = _presence(binding, THREAD)
event = {"type": "promote_chat_to_task", "task_id": "promoted-1", "routing_token": "tok-1",
"objective": "Compile the figures", "chat_id": 4242, "client_message_id": "evt-D1-1712.5",
"project_id": "", "workspace_root": "", "source": "", "presence": presence,
"task_contract": {"capability_ceiling": presence_ceiling_payload(_ceiling())}}
outcome = _handle_promote_chat_to_task(event, handler_ctx)
assert outcome["status"] == "scheduled", outcome
stored = load_task_result(tmp_path, "promoted-1")
assert stored["status"] == "scheduled" and stored["metadata"]["presence"] == presence
assert pending[0]["metadata"]["presence"] == presence # the queue and the record agree
with TestClient(create_host_service_app(tmp_path)) as client:
polled = client.get("/presence/work/promoted-1", params={"binding_id": binding},
headers={"X-Skill-Token": "presence-token"})
foreign = client.get("/presence/work/promoted-1", params={"binding_id": OTHER},
headers={"X-Skill-Token": "presence-token"})
assert polled.status_code == 202 and polled.json()["status"] == "pending"
assert foreign.status_code in {403, 404}
def test_context_lists_own_work_from_other_conversations_after_the_pointer_moved(tmp_path):
from ouroboros.presence_context import build_presence_context_section
from ouroboros.cancel_intents import request_cancel
_work(tmp_path, "done-room", "completed", key=ROOM, result="The report is ready.")
_work(tmp_path, "queued-here", "scheduled")
_work(tmp_path, "running-thread", "running", key=THREAD)
request_cancel(tmp_path, "running-thread", reason="withdrawn", source="agent_tool")
_work(tmp_path, "foreign", "completed", binding=OTHER, result="Other binding")
_work(tmp_path, "promoted-self", "running")
value = {**_presence(), "instructions": "Be useful.", "previous_turn": {
"task_id": "presence-later", "outcome": "silent", "finished_at": "2026-09-24T12:00:00+00:00",
"work_ref": "", # a later turn replaced the pointer; the work is still found
}}
section = build_presence_context_section(tmp_path, value, "promoted-self")
own = section.split("## Work started from this binding (host-authored facts)", 1)[1]
assert "done-room [completed] from conversation slack:T1:C9:0" in own
assert "The report is ready." in own
assert "queued-here [scheduled] from this conversation" in own
assert "running-thread [running, cancel pending] from conversation slack:T1:D1:1712.5" in own
assert "foreign" not in own and "promoted-self" not in own
assert "says nothing about whether its result reached anyone" in own
# --- repair pass: read gaps, effective redirects, unconfirmed promotion, forked roots ----
def test_an_unreadable_result_row_leaves_queued_own_work_listed_and_the_gap_counted(tmp_path):
from ouroboros.presence_context import build_presence_context_section
task_dir = tmp_path / "task_results"
task_dir.mkdir(parents=True)
# A torn row of own queued work, a torn row nothing attributes, and a torn row
# the queue says is another binding's: only the first is this binding's work.
for name in ("queued-torn", "loose-torn", "foreign-torn"):
(task_dir / f"{name}.json").write_text("{not json", encoding="utf-8")
_work(tmp_path, "readable-queued", "scheduled")
_queue(tmp_path, pending=[
{"id": "queued-torn", "delegation_role": "root", "description": "compile the figures",
"metadata": {"presence": _presence(key=THREAD)}},
{"id": "foreign-torn", "delegation_role": "root", "metadata": {"presence": _presence(OTHER)}},
{"id": "readable-queued", "delegation_role": "root", "metadata": {"presence": _presence()}},
])
registry, _ctx = _registry(tmp_path)
page = json.loads(registry.execute("recent_tasks", {"limit": 20}))
rows = {row["task_id"]: row for row in page["tasks"]}
assert set(rows) == {"queued-torn", "readable-queued"} # the readable row replaced its queue row once
assert rows["queued-torn"] == {"task_id": "queued-torn", "status": "pending", "description": "compile the figures",
"source": "queue_snapshot", "result_row": "unreadable"}
assert rows["readable-queued"]["status"] == "scheduled" and "result_row" not in rows["readable-queued"]
assert page["read_gap"] == {"unattributed_unreadable_rows": 1} # foreign-torn is attributed by its queue row
section = build_presence_context_section(tmp_path, {**_presence(), "instructions": "Be useful."}, "turn-x")
assert "queued-torn [pending, its result row is unreadable]" in section
assert "1 row(s) no record attributes; this binding's work may be among them" in section
# Nothing torn: no gap is claimed.
for name in ("queued-torn", "loose-torn", "foreign-torn"):
(task_dir / f"{name}.json").unlink()
clean = json.loads(registry.execute("recent_tasks", {"limit": 20}))
assert "read_gap" not in clean and {row["task_id"] for row in clean["tasks"]} == {"queued-torn", "readable-queued"}
def test_an_effective_redirect_is_judged_before_projection_and_own_retries_still_read(tmp_path):
# Own work whose retry successor is another binding's: neither reader projects it.
_work(tmp_path, "mine-redirected", "interrupted", superseded_by="theirs-retry", result="own interrupted row")
_work(tmp_path, "theirs-retry", "completed", binding=OTHER, result="FOREIGN SUCCESSOR BODY")
# A real same-binding retry (the reaper copies metadata, role and lineage): read through.
_work(tmp_path, "mine-timed-out", "interrupted", superseded_by="mine-retry", result="timed out")
_work(tmp_path, "mine-retry", "completed", root_task_id="mine-timed-out", original_task_id="mine-timed-out",
supersedes_task_id="mine-timed-out", result="Retry finished the report.")
registry, _ctx = _registry(tmp_path)
refused = registry.execute("get_task_result", {"task_id": "mine-redirected"})
assert "effective result continues in work that was not started from this Presence binding" in refused
assert "FOREIGN SUCCESSOR BODY" not in refused and "theirs-retry" not in refused
assert "Retry finished the report." in registry.execute("get_task_result", {"task_id": "mine-timed-out"})
page = json.loads(registry.execute("recent_tasks", {"limit": 20, "include_results": True}))
rows = {row["task_id"]: row for row in page["tasks"]}
assert "FOREIGN SUCCESSOR BODY" not in json.dumps(page) and "theirs-retry" not in rows
assert rows["mine-redirected"]["status"] == "interrupted"
assert rows["mine-redirected"]["effective_result"].startswith("withheld")
# The own retry reads through (the effective row names its successor, as it always has).
retried = [row for row in page["tasks"] if row["task_id"] == "mine-retry"]
assert len(retried) == 2 and all(row["result"] == "Retry finished the report." for row in retried)
assert not any("effective_result" in row for row in retried)
# The owner's unscoped reader keeps the ordinary effective projection.
owner = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
owner.set_context(ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="owner-turn"))
assert "FOREIGN SUCCESSOR BODY" in owner.execute("get_task_result", {"task_id": "mine-redirected"})
def test_an_unconfirmed_presence_promote_is_readable_as_pending_by_its_own_binding(tmp_path, monkeypatch):
from ouroboros.tools import control_events
from ouroboros.tools.control_routing import _promote_chat_to_task
from ouroboros.tools.control_task_results import _get_task_result
monkeypatch.setattr(control_events, "_PROMOTE_CONFIRM_TIMEOUT_SEC", 0.05)
monkeypatch.setattr(control_events, "_PROMOTE_CONFIRM_POLL_SEC", 0.005)
monkeypatch.setattr("ouroboros.config.DATA_DIR", tmp_path)
emitted = []
presence = _presence(key=THREAD)
ctx = types.SimpleNamespace(
pending_events=[], event_queue=types.SimpleNamespace(put_nowait=emitted.append), current_chat_id=4242,
drive_root=tmp_path, budget_drive_root=str(tmp_path), project_id="", task_id="presence-turn-1",
task_metadata={"presence": presence}, is_direct_chat=True,
task_contract={"capability_ceiling": presence_ceiling_payload(_ceiling())},
)
out = _promote_chat_to_task(ctx, "Compile the figures", predecessor_task_id="")
assert out.startswith("⚠️ PROMOTE_UNCONFIRMED")
task_id = emitted[0]["task_id"]
stub = load_task_result(tmp_path, task_id)
# Emitted, not scheduled: the supervisor alone grants that; the provenance is the event's own.
assert (stub["status"], stub["promotion_admission"]["status"]) == ("requested", "emitted")
assert stub["metadata"]["presence"] == presence and stub["delegation_role"] == "root"
read = _get_task_result(ctx, task_id, presence_scope="own_binding")
assert "admission pending since" in read and "PRESENCE_CAPABILITY_BLOCKED" not in read
registry, _registry_ctx = _registry(tmp_path)
listed = json.loads(registry.execute("recent_tasks", {"limit": 20}))
assert [(row["task_id"], row["status"]) for row in listed["tasks"]] == [(task_id, "requested")]
stranger, _stranger_ctx = _registry(tmp_path, binding=OTHER)
assert "PRESENCE_CAPABILITY_BLOCKED" in stranger.execute("get_task_result", {"task_id": task_id})
# An ordinary promote's stub is unchanged: no Presence provenance is invented.
owner_emitted = []
owner_ctx = types.SimpleNamespace(
pending_events=[], event_queue=types.SimpleNamespace(put_nowait=owner_emitted.append), current_chat_id=1,
drive_root=tmp_path, budget_drive_root=str(tmp_path), project_id="", task_metadata={}, task_id="",
)
_promote_chat_to_task(owner_ctx, "Owner work", workspace="none", predecessor_task_id="")
owner_stub = load_task_result(tmp_path, owner_emitted[0]["task_id"])
assert "presence" not in (owner_stub.get("metadata") or {}) and "delegation_role" not in owner_stub
def test_a_forked_promoted_root_reads_its_bindings_work_and_sends_from_the_canonical_root(tmp_path):
from ouroboros.agent import Env
from ouroboros.context import build_llm_messages
from ouroboros.memory import Memory
from ouroboros.presence_context import presence_finish_not_accepted_note
from ouroboros.utils import append_jsonl
from tests.test_doc_context import _make_env_and_memory
canonical_env, _memory = _make_env_and_memory(tmp_path)
canonical, child = canonical_env.drive_root, tmp_path / "child-drive"
for sub in ("memory/knowledge", "logs", "state"):
(child / sub).mkdir(parents=True, exist_ok=True)
_work(canonical, "done-room", "completed", key=ROOM, result="The canonical report.")
_work(child, "child-drive-decoy", "completed", key=ROOM, result="decoy") # worker-local rows only
presence = {**_presence(), "instructions": "Be useful.", "delivery_reporting_version": 1}
for root, text in ((canonical, "Canonical sent reply"), (child, "Child drive decoy send")):
append_jsonl(root / "logs" / "chat.jsonl", {"task_id": "promoted-self", "direction": "in", "text": "x"})
append_jsonl(root / "logs" / "chat.jsonl", {
"task_id": "promoted-self", "type": "presence_delivery", "text": text,
"transport": {"conversation_key": HERE, "delivery": {"state": "delivered", "delivery_id": "d", "part_id": "0"}},
})
env = Env(repo_dir=canonical_env.repo_dir, drive_root=child, budget_drive_root=canonical)
task = {"id": "promoted-self", "type": "task", "text": "Compile", "delegation_role": "root",
"_presence_origin": True, "budget_drive_root": str(canonical), "metadata": {"presence": presence}}
messages, _ = build_llm_messages(env=env, memory=Memory(child, repo_dir=env.repo_dir), task=task)
rendered = json.dumps(messages, ensure_ascii=False)
assert "done-room [completed] from conversation slack:T1:C9:0" in rendered
assert "child-drive-decoy" not in rendered
ctx = types.SimpleNamespace(drive_root=child, budget_drive_root=str(canonical), task_id="promoted-self",
task_metadata={"presence": presence, "budget_drive_root": str(canonical)})
note = presence_finish_not_accepted_note(ctx, {"outcome": "tool_delivered"})
assert '"Canonical sent reply"' in note and "decoy" not in note

View file

@ -1,5 +1,6 @@
"""External speech follows terminal authorship through execution and replay."""
import json
import queue
from types import SimpleNamespace
@ -19,12 +20,12 @@ from tests.test_presence_failed_handoff import _failed_parent
from tests.test_presence_runner import _admission
def _run_loop(root, monkeypatch, responses, *, held=False):
def _run_loop(root, monkeypatch, responses, *, held=False, presence=None):
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.setenv("OUROBOROS_MAX_ROUNDS", "1")
registry = ToolRegistry(repo_dir=root, drive_root=root)
registry._ctx.is_direct_chat = True
registry._ctx.task_metadata = {"inline_max_rounds": 1}
registry._ctx.task_metadata = {"inline_max_rounds": 1, **({"presence": presence} if presence else {})}
registry._ctx.task_contract = {"capability_ceiling": presence_ceiling_payload(_admission().capability_ceiling)}
registry.override_handler("chat_history", lambda *_a, **_kw: "Synthetic history")
calls, held_origins = [], []
@ -65,7 +66,12 @@ def _read_response():
@pytest.mark.parametrize("authored", [False, True])
def test_real_round_limit_delivers_only_the_current_authored_final(tmp_path, monkeypatch, authored):
reply = "I found the record; the remaining check is incomplete." if authored else ""
result, stored, calls, _held = _run_loop(tmp_path, monkeypatch, [_read_response(), {"content": reply}])
forced = json.dumps({"delivery_control": "replace", "full_answer": reply,
"presence_finish": {"outcome": "message", "message": reply}}) if authored else ""
# A real turn's context carries its Presence metadata; that, with the ceiling, arms the forced call.
presence = {"binding_id": "1" * 32, "event": {"conversation_key": "telegram:bot-1:room-1:topic-1"}}
result, stored, calls, _held = _run_loop(tmp_path, monkeypatch, [_read_response(), {"content": forced}],
presence=presence)
assert len(calls) == 2 and stored["reason_code"] == "round_limit"
assert stored["terminal_origin"] == ("model_final" if authored else "host_notice")
assert result["outcome"] == ("message" if authored else "silent")

View file

@ -177,7 +177,7 @@ def test_initiate_presence_resolves_binding_and_reports_actual_delivery(monkeypa
assert captured["event"].actor["kind"] == "proactive_initiation"
def test_presence_cancel_work_accepts_only_same_binding_and_conversation(monkeypatch, tmp_path) -> None:
def test_presence_cancel_work_accepts_own_binding_work_from_any_of_its_conversations(monkeypatch, tmp_path) -> None:
ctx = _ctx(tmp_path)
ctx.task_metadata = {
"presence": {
@ -185,28 +185,33 @@ def test_presence_cancel_work_accepts_only_same_binding_and_conversation(monkeyp
"event": {"conversation_key": "telegram:bot-1:room-1"},
}
}
atomic_write_json(
tmp_path / "task_results" / "presence-work-1.json",
{
"_schema_version": 1,
"task_id": "presence-work-1",
"status": "running",
"metadata": {
"presence": {
"binding_id": "1" * 32,
"event": {"conversation_key": "telegram:bot-1:room-1"},
}
},
},
)
def row(task_id, *, binding="1" * 32, key="telegram:bot-1:room-1", **fields):
atomic_write_json(tmp_path / "task_results" / f"{task_id}.json", {
"_schema_version": 1, "task_id": task_id, "status": "running",
"metadata": {"presence": {"binding_id": binding, "event": {"conversation_key": key}}} if binding else {},
**fields,
})
row("promoted-work-1", delegation_role="root", root_task_id="promoted-work-1")
row("other-thread-work", key="telegram:bot-1:room-1:thread-9", delegation_role="root")
row("foreign-binding-work", binding="2" * 32, delegation_role="root")
row("presence-inline-turn") # an inline turn is not deferred work
row("owner-root", binding="", delegation_role="root")
monkeypatch.setattr(
"ouroboros.tools.join_ledger._cancel_task",
lambda _ctx, task_id, reason="": f"cancel:{task_id}:{reason}",
)
entry = next(item for item in get_tools() if item.name == "presence_cancel_work")
assert entry.handler(ctx, "presence-work-1", "no longer needed") == (
"cancel:presence-work-1:no longer needed"
assert entry.handler(ctx, "promoted-work-1", "no longer needed") == (
"cancel:promoted-work-1:no longer needed"
)
# Owner Q1: the same nonempty binding, a different thread or room of it.
assert entry.handler(ctx, "other-thread-work") == "cancel:other-thread-work:"
ctx.task_metadata["presence"]["event"]["conversation_key"] = "telegram:bot-1:other"
assert entry.handler(ctx, "presence-work-1") == "ERROR: PRESENCE_WORK_NOT_CORRELATED"
assert entry.handler(ctx, "promoted-work-1") == "cancel:promoted-work-1:"
for foreign in ("foreign-binding-work", "presence-inline-turn", "owner-root", "never-existed"):
assert entry.handler(ctx, foreign).startswith("ERROR: PRESENCE_WORK_NOT_CORRELATED")
ctx.task_metadata["presence"]["binding_id"] = "" # an empty binding never compares equal
assert entry.handler(ctx, "promoted-work-1").startswith("ERROR: PRESENCE_WORK_NOT_CORRELATED")

View file

@ -0,0 +1,163 @@
"""A forced Presence final under the Host contract the installed transport actually honours.
The fake transport below mirrors the installed Telegram adapter without importing
Hub code: ``submit`` is ``custody.record_submission`` (a message/deferred turn body
is queued, a deferred turn registers its work_ref) and ``poll`` is
``runtime.process_one_work`` over ``host.poll`` (pending keeps the work; a terminal
result must echo the polled work_ref, closes the work, and only a ``message`` body
is sent). The real /presence/turn and /presence/work endpoints, loop and pipeline
run with a scripted model; nothing leaves the process.
Boundary, not covered here: a polled late result that is itself ``deferred`` (owed
work scheduling further work) is stored truthfully (body and nested work_ref), but
that consumer sends no deferred late body and closes the work, and the poll
response must echo the polled work_ref, so the nested obligation cannot reach it
without a transport protocol change.
"""
from __future__ import annotations
import queue
from types import SimpleNamespace
import pytest
from starlette.testclient import TestClient
from ouroboros import agent_task_pipeline as pipeline, loop
from ouroboros.gateway.host_service import create_host_service_app
from ouroboros.presence_runner import PresenceTurnGate, run_presence_turn
from ouroboros.task_results import write_task_result
from ouroboros.tools.registry import ToolRegistry
from tests.test_host_service_api import _seed_presence_behavior, _seed_token
from tests.test_presence_forced_delivery import RECORD, _forced, _read
_TOKEN = "presence-token"
_OUTCOMES = {"message", "silent", "tool_delivered", "deferred"}
PARTIAL = "Q1 is ready: 41. Q2 is still being checked."
RESULT = "Q2 is ready too: 43."
NOTE = "sent the table via the transport tool; helper child-7 failed with provider 400"
class _InstalledTransportContract:
def __init__(self, client: TestClient, binding: str) -> None:
self.client, self.binding = client, binding
self.outbox: list[str] = []
self.open_work: list[str] = []
def submit(self, body: dict) -> None:
assert body["ok"] is True and body["status"] == "completed" and body["outcome"] in _OUTCOMES
if body["text"] and body["outcome"] in {"message", "deferred"}:
self.outbox.append(body["text"])
if body["outcome"] == "deferred":
assert body["work_ref"], "deferred submission requires work_ref"
self.open_work.append(body["work_ref"])
def poll(self) -> None:
for work_ref in list(self.open_work):
body = self.client.get(f"/presence/work/{work_ref}", params={"binding_id": self.binding},
headers={"X-Skill-Token": _TOKEN}).json()
if body["status"] == "pending":
continue
assert body["status"] in {"completed", "failed", "cancelled"} and body["outcome"] in _OUTCOMES
assert body.get("work_ref") in {"", work_ref}, "presence Host returned a different work_ref"
self.open_work.remove(work_ref)
if body["outcome"] == "message" and body["text"]:
self.outbox.append(body["text"])
def _scripted_agent(root, monkeypatch, forced, *, handoff=None):
"""The real loop and pipeline for the task it is handed: one tool round, then the ONE forced call."""
class Agent:
def handle_task(self, task):
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.setenv("OUROBOROS_MAX_ROUNDS", "1")
registry = ToolRegistry(repo_dir=root, drive_root=root)
ctx = registry._ctx
ctx.is_direct_chat = bool(task.get("_is_direct_chat"))
ctx.task_metadata = {**task["metadata"], "inline_max_rounds": 1}
ctx.task_contract = dict(task["task_contract"])
if handoff:
ctx._swarm_handoff_attempt = dict(handoff)
registry.override_handler("chat_history", lambda *_a, **_kw: "Synthetic history")
replies = iter([_read(), {"content": forced}])
monkeypatch.setattr(loop, "call_llm_with_retry", lambda *_a, **_k: (next(replies), 0.0))
task["_skip_post_task_synthesis"] = True
text, usage, trace = loop.run_llm_loop(
[{"role": "user", "content": task["text"]}], registry,
SimpleNamespace(default_model=lambda: "test-model"), root / "logs",
lambda *_a, **_kw: None, queue.Queue(), task_id=task["id"], drive_root=root,
)
events = []
pipeline.emit_task_results(SimpleNamespace(drive_root=root, repo_dir=root), None, None,
events, task, text, usage, trace, 0.0, root / "logs", ctx=ctx)
return events
return Agent()
def _turn(client, binding):
return client.post("/presence/turn", headers={"X-Skill-Token": _TOKEN}, json={
"binding_id": binding,
"event": {
"source_event_id": "telegram:bot-1:42", "provider": "telegram", "account_id": "bot-1",
"conversation_id": "room-1", "thread_id": "topic-1", "conversation_key": "ignored",
"actor": {"platform_actor_id": "user-7"}, "conversation": {"title": "Community"},
"message": {"message_id": "42"}, "text": "Please compile Q1 and Q2",
},
}).json()
@pytest.mark.parametrize("turn_forced,spoken_now", [
(_forced("message", PARTIAL), [PARTIAL]), # a declared partial beside owed work
(_forced("deferred", PARTIAL), [PARTIAL]),
(_forced("tool_delivered", NOTE), []), # the note is context, never a reply
(RECORD, []), # an undeclared internal record says nothing new
], ids=["message", "deferred", "tool_delivered", "undeclared"])
def test_the_installed_transport_gets_the_partial_now_and_the_owed_result_later(
tmp_path, monkeypatch, turn_forced, spoken_now):
_seed_token(tmp_path, skill="telegram-bot", token=_TOKEN, permissions=["presence"],
manifest_permissions=["presence"])
binding = _seed_presence_behavior(tmp_path)
repo = tmp_path / "repo"
repo.mkdir()
captured = {}
def runner(**kwargs):
def factory(**_kwargs):
agent = _scripted_agent(tmp_path, monkeypatch, turn_forced,
handoff={"status": "scheduled", "task_id": "work-9"})
original = agent.handle_task
def handle(task):
captured["turn"] = task
# The admitted promotion the handoff names, canonical from admission.
write_task_result(tmp_path, "work-9", "scheduled", delegation_role="root",
root_task_id="work-9", description="Compile Q2",
metadata={"presence": dict(task["metadata"]["presence"])})
return original(task)
agent.handle_task = handle
return agent
return run_presence_turn(repo_dir=repo, drive_root=tmp_path, agent_factory=factory,
gate=PresenceTurnGate(1), **kwargs)
with TestClient(create_host_service_app(tmp_path, presence_runner=runner)) as client:
transport = _InstalledTransportContract(client, binding)
transport.submit(_turn(client, binding))
assert transport.outbox == spoken_now and transport.open_work == ["work-9"] # speech AND custody
transport.poll() # the owed work has not run: nothing is sent and it stays owed
assert transport.outbox == spoken_now and transport.open_work == ["work-9"]
turn = captured["turn"]
work = {"id": "work-9", "type": "task", "chat_id": turn["chat_id"], "text": "Compile Q2",
"delegation_role": "root", "root_task_id": "work-9", "_presence_origin": True,
"metadata": {"presence": dict(turn["metadata"]["presence"])},
"task_contract": dict(turn["task_contract"])}
_scripted_agent(tmp_path, monkeypatch, _forced("message", RESULT)).handle_task(work)
transport.poll()
assert transport.outbox == spoken_now + [RESULT] and transport.open_work == []
assert not any(RECORD in text or NOTE in text for text in transport.outbox)

View file

@ -143,7 +143,10 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# after the long-work continuity merge, 308678 after the Claudexor 3.14.0 pin, #1264);
# the run-origin sentences replace the "request decides what the run was for" and the
# owner-turn descriptions (+144 on the merged base) rather than appending to them.
"docs/architecture/06-agent-core.md": 308900,
# 308900 -> 309100 (TZ2 own work): the task-message sentence names the Presence
# sender's own-binding boundary and the delivery-control paragraph the Presence
# forced declaration it resolves around; both clauses extend existing sentences.
"docs/architecture/06-agent-core.md": 309100,
"docs/architecture/07-configuration.md": 36991,
# 18947 -> 19287: CI failure collection now documents diagnostic desktop builds while release remains gated.
# 19287 -> 20560 (#1215): three contracts the chapter had no older text for — the
@ -183,7 +186,11 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# conversation key, placeholder re-run and its lost-attempt facts, presence-local liveness,
# previous-turn pointer and its replay repair, split in-flight budgets, silent orphaned work,
# presence room label); the base sat 2 bytes under.
"docs/architecture/12-host-service-companions-and-chat-ids.md": 12500,
# 12500 -> 13400 (TZ2 own work): the Presence paragraph gains two contracts it had
# no text for — what a binding's own work is and which readers/controls reach it
# (replacing the conversation-exact cancel sentence), and the forced-final split
# between the internal record and the declared reply; the base sat 10 bytes under.
"docs/architecture/12-host-service-companions-and-chat-ids.md": 13400,
# 7764 -> 8600 (#1195): the fresh selected-subject + immutable peer projection
# execution check (`skill_peer_inventory.py`, `skill_conflicts.py`) replaces
# whole-inventory hashing; the chapter had no description of that seam to swap out.
@ -209,7 +216,9 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# chapter had 5 bytes left. Sized to the text: 5 bytes of margin.
# 94520 -> 94900: the delegated-lane bullet names the worktree ops lock rule
# (issue #1241: no tree walk or per-file git process under the lock).
"docs/development/06-rules-by-change-class.md": 94900,
# 94900 -> 95000 (TZ2 own work): the Presence bullets replace the conversation-exact
# cancel clause with the own-binding rule and name the forced declaration.
"docs/development/06-rules-by-change-class.md": 95000,
"docs/development/07-managed-update-rule.md": 4166,
"docs/development/08-mutation-attribution-rule.md": 2899,
"docs/development/09-process-custody-rule.md": 10028,