wip: preserve Presence TZ2 corrections pending Q4 and full review

This commit is contained in:
Ouroboros 2026-09-25 09:37:03 +03:00
parent 7ed80c8a79
commit 7d8e742d12
25 changed files with 508 additions and 76 deletions

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, 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. A delegated descendant inherits the ceiling plus only `metadata.presence_binding_authority`, never the speaker's `metadata.presence` (no forced reply): its readers, steer and forward reach that work and its own tree, root included, and the supervisor also fences its steer by its live row. 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. A delegated descendant inherits the ceiling plus only `metadata.presence_binding_authority`, never the speaker's `metadata.presence` (no forced reply): its readers, steer and forward reach that work and its own tree, root included; its promote and that root's follow-ups carry the same carrier (`presence_root_carrier`): related work that never speaks. Steer, like read/cancel, trusts a canonical row over a live one; an unstamped steer is fenced by the sender's live row. 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, 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.
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-authored text stays owner-side; an ordinary final's own wording is not screened. A tool-send finish note stays a separate `finish_note` in the previous-turn pointer even when owed work changes its outcome to deferred; it is never prior speech. 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

@ -54,7 +54,7 @@ Enforcement: `tests/test_protected_artifacts_policy.py` and `tests/test_acceptan
- 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). A Presence caller reads, steers and cancels only independent roots of its own nonempty binding (`presence_authority.presence_work_refusal`); a delegated descendant is one through the inherited `metadata.presence_binding_authority` alone, never the speaker's `metadata.presence`; 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.
- 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`); a delegated descendant is one through the inherited `metadata.presence_binding_authority` alone, never the speaker's `metadata.presence`; and a forced final speaks only its nested `presence_finish` declaration. Promotion and `schedule_followup` copy one Presence carrier (`presence_root_carrier`: speaker metadata or a descendant's binding), 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

@ -35,12 +35,14 @@ PRESENCE_BINDING_AUTHORITY_KEY = "presence_binding_authority"
def presence_record_binding(record: Any) -> str:
"""The nonempty host binding id one task/queue record carries, else ``""``."""
"""The nonempty host binding id one task/queue record carries, else ``""``.
The one reader of both carriers: a speaker's ``metadata.presence`` and the
``metadata.presence_binding_authority`` of work a delegated descendant started.
"""
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 ""
return presence_metadata_binding(metadata) or ""
def presence_related_work(binding_id: str, record: Any) -> bool:
@ -99,6 +101,23 @@ def presence_binding_authority_metadata(parent_metadata: Any, *, task_contract:
return {} if binding is None else {PRESENCE_BINDING_AUTHORITY_KEY: {"binding_id": binding}}
def presence_root_carrier(source: Any, *, task_contract: Any = None) -> dict[str, Any]:
"""The Presence carrier an independent root started from ``source`` keeps, or ``{}``.
``source`` is the starting task's metadata or its promote event. A Presence
turn or root hands on its speaker metadata: the new root answers the same
conversation. A delegated descendant hands on only the binding it acts for:
its root is that binding's related work, never a speaker. A malformed or lost
carrier under a ceiling narrows to an empty binding. Producer and admission
both read this; ``presence_record_binding`` reads what it writes.
"""
presence = source.get("presence") if isinstance(source, Mapping) else None
if isinstance(presence, Mapping) and presence:
return {"presence": dict(presence)}
return presence_binding_authority_metadata(source, task_contract=task_contract)
def presence_sender_origin(ctx: Any) -> dict[str, str]:
"""Where a Presence caller's run started (its ``run_origin`` room/event facts).
@ -123,11 +142,14 @@ def presence_queue_task(drive_root: Any, task_id: str) -> dict[str, Any] | None:
return None
def presence_target_record(drive_root: Any, task_id: str) -> Mapping[str, Any] | None:
def presence_target_record(drive_root: Any, task_id: str, *,
queue_row: Mapping[str, Any] | None = None) -> 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.
The canonical task record decides, a malformed Presence carrier included (it
narrows to nothing); a legacy row without Presence provenance may be established
only by the queue's own task metadata — ``queue_row`` when the caller holds the
live row (the supervisor), else the persisted snapshot.
"""
from ouroboros.task_results import load_task_result
@ -138,8 +160,12 @@ def presence_target_record(drive_root: Any, task_id: str) -> Mapping[str, Any] |
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
contract = record.get("task_contract") if isinstance(record, Mapping) else None
has_ceiling = isinstance(contract, Mapping) and "capability_ceiling" in contract
if record is None or (presence_metadata_binding(record.get("metadata")) is None and not has_ceiling):
queued = {**dict(queue_row), "id": target} if isinstance(queue_row, Mapping) else (
presence_queue_task(drive_root, target))
record = queued or record
return record

View file

@ -12,7 +12,7 @@ import logging
import pathlib
from typing import Dict, List
from ouroboros.presence_authority import presence_record_binding
from ouroboros.presence_authority import presence_metadata_binding, presence_record_binding
from ouroboros.task_result_schema import (
quarantine_task_result,
task_result_schema_refusal,
@ -82,6 +82,14 @@ def raw_result_facts(results_dir: pathlib.Path, *, reader=None) -> tuple[Dict[st
# 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)
# An empty scalar is not proof of absent provenance: a malformed carrier or
# lost metadata under an inherited ceiling must never inherit a queue claim.
metadata = data.get("metadata")
contract = data.get("task_contract")
facts["presence_authority_recorded"] = (
presence_metadata_binding(metadata) is not None
or isinstance(contract, dict) and "capability_ceiling" in contract
)
facts["schema_refusal"] = task_result_schema_refusal(data)
rows[name] = facts
_RAW_TS_MEMO[key] = (signature, tuple(facts.items()))

View file

@ -24,7 +24,7 @@ import pathlib
import time
from typing import Any, Dict, List, Optional
from ouroboros.dialogue_provenance import is_presence_task
from ouroboros.dialogue_provenance import is_presence_task, presence_caller_binding
from ouroboros.focus import compact_focus, focus_fingerprint
from ouroboros.task_status import _load_queue_snapshot, queue_snapshot_observation
from ouroboros.utils import read_json_dict
@ -269,7 +269,7 @@ def maybe_append_roster_note(ctx: Any, messages: List[Dict[str, Any]], drive_roo
"_presence_origin": getattr(ctx, "_presence_origin", None)}
if (str(metadata.get("parent_task_id") or "").strip()
or str(metadata.get("delegation_role") or "") == "subagent"
or is_presence_task(actor_task)):
or is_presence_task(actor_task) or presence_caller_binding(ctx) is not None):
return False
task_id = str(getattr(ctx, "task_id", "") or "")
canonical = pathlib.Path(str(

View file

@ -11,6 +11,7 @@ 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_metadata_binding as presence_metadata_binding,
presence_effective_hops,
presence_effective_related,
presence_queue_task,

View file

@ -58,11 +58,16 @@ def _previous_turn_line(previous: Mapping[str, Any]) -> str:
sends = [str(text) for text in sends if str(text or "").strip()]
message = str(previous.get("message") or "").strip()
said = [json.dumps(text, ensure_ascii=False) for text in sends]
if previous.get("outcome") == "tool_delivered": # its message is the model's note, never speech
note = str(previous.get("finish_note") or "").strip()
if previous.get("outcome") == "tool_delivered": # legacy pointers kept the note in message
said = said or ["delivered via transport tool (content unrecorded)"]
said += [f"finish note {json.dumps(message, ensure_ascii=False)}"] if message else []
note = note or message
elif message and message not in sends:
said.append(json.dumps(message, ensure_ascii=False))
if note:
said.append(f"finish note {json.dumps(note, ensure_ascii=False)}")
if previous.get("previous_text_unverified"):
said.append("legacy previous text unverified as speech (source unavailable)")
body = " / ".join(said) or "nothing sent"
work = ""
if previous.get("work_ref"):
@ -101,7 +106,7 @@ def _own_work_section(drive_root: Path, value: Mapping[str, Any], task_id: str)
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")
else f"conversation {key}" if key else "a conversation this row does not name")
preview = " ".join(str(row.get("result_preview") or "").split())[:200]
lines.append(
f"- {row.get('task_id')} [{row.get('status') or 'unknown'}"
@ -113,10 +118,12 @@ def _own_work_section(drive_root: Path, value: Mapping[str, Any], task_id: str)
+ (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")
unread = [text for key, text in (
("result_root", "the result root"), ("queue_snapshot", "the queue snapshot (queued work)"),
("unattributed_unreadable_rows", f"{gap.get('unattributed_unreadable_rows')} row(s) no record attributes"),
) if gap.get(key)]
if unread:
lines.append("(some results are unreadable now: " + "; ".join(unread)
+ "; 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)")

View file

@ -158,8 +158,10 @@ def build_presence_result_event(task: dict[str, Any], text: str, ctx: Any, *, te
"task_id": str(task.get("id") or ""),
"outcome": outcome,
"text": result_text,
# The accepted presence_finish message; a tool_delivered note is context, never speech.
"message": result_text or note,
# Speech and the internal finish note stay separate even when owed work
# changes tool_delivered into deferred for the polling contract.
"message": result_text,
**({"finish_note": note} if note else {}),
"work_ref": work_ref,
"ts": utc_now_iso(),
}
@ -344,12 +346,35 @@ def _read_previous_turn(drive_root: Path, conversation_key: str) -> dict[str, An
return row if row.get("conversation_key") == conversation_key else None
def _previous_turn_source_view(drive_root: Path, pointer: dict[str, Any]) -> dict[str, Any]:
"""Read a legacy deferred pointer's speech from its canonical result, without rewriting history.
Old tool-delivered notes used the same `message` slot as replies. Owed work
changed their outcome to deferred, hiding that provenance. A missing source
proves neither speech nor a note; the context must say it is unverified.
"""
if (pointer.get("outcome") != "deferred" or not pointer.get("message")
or "finish_note" in pointer):
return pointer
row = load_task_result(drive_root, str(pointer.get("task_id") or "")) or {}
metadata = row.get("metadata") if isinstance(row.get("metadata"), dict) else {}
observed = metadata.get("presence_result_text")
if metadata.get("presence_outcome") == "deferred" and isinstance(observed, str):
if observed == pointer["message"]:
return pointer # source proves an authored partial reply, not a note
if not observed:
return {**pointer, "message": "", "finish_note": pointer["message"]}
return {**pointer, "message": "", "previous_text_unverified": True}
def _write_previous_turn(drive_root: Path, conversation_key: str, task_id: str, *, outcome: str, message: str,
sends: Sequence[str], work_ref: str, finished_at: str, delivery: str) -> None:
sends: Sequence[str], work_ref: str, finished_at: str, delivery: str,
finish_note: str = "") -> None:
"""Best effort: the pointer is a projection, and a turn that already answered is not failed over it."""
try:
atomic_write_json(_previous_turn_path(drive_root, conversation_key), {
"conversation_key": conversation_key, "task_id": task_id, "outcome": outcome, "message": message,
**({"finish_note": finish_note} if finish_note else {}),
"transport_sends": [text for text in sends if text], "work_ref": work_ref,
"finished_at": finished_at, "delivery": delivery,
})
@ -481,7 +506,8 @@ def _build_task(
}
previous_turn = _read_previous_turn(drive_root, event.conversation_key)
if previous_turn:
if previous_turn.get("work_ref"): # the deferred child's fate is read from its canonical row, never stored
previous_turn = _previous_turn_source_view(drive_root, previous_turn)
if previous_turn.get("work_ref"): # the deferred child's fate is read from its canonical row, never stored
work_ref = str(previous_turn["work_ref"])
child = load_task_result(drive_root, work_ref) or {}
status = str(child.get("status") or "")
@ -699,7 +725,8 @@ def run_presence_turn(
_write_previous_turn(Path(drive_root), event.conversation_key, task_id, outcome=result.outcome,
message=str(row.get("message") or result.text), sends=sends or [],
work_ref=result.work_ref, finished_at=utc_now_iso(),
delivery=_delivery_state(event.delivery_reporting_version, sends, result.text))
delivery=_delivery_state(event.delivery_reporting_version, sends, result.text),
finish_note=str(row.get("finish_note") or ""))
return result
return (gate or _configured_gate(Path(drive_root))).run(event.conversation_key, execute)

View file

@ -155,9 +155,11 @@ def resolve_project_id(task: Dict[str, Any]) -> str:
# mismatch the forked seed prepared at schedule time for an unscoped parent.
if str(task.get("delegation_role") or "") == "subagent":
return ""
metadata = task.get("metadata")
if isinstance(metadata, dict) and isinstance(metadata.get("presence"), dict) and metadata["presence"]:
# Presence changes file/process cwd, not its canonical memory scope.
from ouroboros.dialogue_provenance import presence_metadata_binding
if presence_metadata_binding(task.get("metadata")) is not None:
# Presence (a speaker, or work acting for its binding) changes file/process cwd,
# not its canonical memory scope.
return ""
workspace = str(task.get("workspace_root") or "").strip()
if workspace:

View file

@ -193,11 +193,12 @@ def _record_promotion_admission_stub(ctx: ToolContext, evt: Dict[str, Any], mode
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.dialogue_provenance import presence_root_carrier
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
carrier = presence_root_carrier(evt, task_contract=evt.get("task_contract"))
try:
write_task_result(
_routing_status_root(ctx), task_id, STATUS_REQUESTED,
@ -210,8 +211,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 {}),
**({"metadata": carrier, "source": "presence_promote",
"delegation_role": "root", "root_task_id": task_id} if carrier else {}),
)
except Exception as exc:
# The promote itself proceeds; what is lost is the reconciliation read, so

View file

@ -15,6 +15,7 @@ import uuid
from pathlib import Path
from typing import Any, Dict
from ouroboros.dialogue_provenance import presence_root_carrier
from ouroboros.tools.control_events import (
_PROMOTE_CONFIRM_TIMEOUT_SEC,
_emit_and_wait_for_routing,
@ -438,18 +439,17 @@ def _promote_chat_to_task(
"ts": utc_now_iso(),
}
metadata = getattr(ctx, "task_metadata", {})
presence = metadata.get("presence") if isinstance(metadata, dict) else None
if isinstance(presence, dict) and presence:
# A public conversation may promote long work, but it cannot choose a
# new Project/workspace/source authority. The immutable positive ceiling
# and exact return destination follow the promoted root by value.
presence_carrier = presence_root_carrier(metadata, task_contract=getattr(ctx, "task_contract", None))
if presence_carrier:
# A public conversation cannot choose a new Project/workspace/source authority; the immutable
# ceiling and return destination (a descendant's root: its binding only) follow it by value.
evt.update({
"project_id": "",
"project_name": "",
"workspace_root": "",
"workspace": "",
"source": "",
"presence": dict(presence),
**presence_carrier,
"task_contract": dict(getattr(ctx, "task_contract", {}) or {}),
})
repo_root_note = "" # Presence runs in its admitted folder, never over the repo

View file

@ -847,7 +847,7 @@ def _schedule_task(ctx: ToolContext, internal: Dict[str, Any] | None = None, /,
**intent_fields,
"subagent_envelope": envelope,
"origin_metadata": consciousness_origin_metadata(metadata), # a consciousness child: label, category, level
**presence_binding_authority_metadata(metadata, task_contract=ctx.task_contract), # binding, never speaker
**presence_binding_authority_metadata(metadata, task_contract=getattr(ctx, "task_contract", None)), # never speaker
}
_populate_subagent_event_extras(
evt, current_chat_id=current_chat_id, child_drive=child_drive,

View file

@ -34,6 +34,7 @@ from typing import Any, Dict, List
from ouroboros.consciousness_authority import consciousness_origin_metadata
from ouroboros.deadline_utils import parse_deadline_ts
from ouroboros.dialogue_provenance import presence_caller_binding, presence_root_carrier
from ouroboros.tools.arg_feedback import ignored_argument_note
from ouroboros.tools.registry import ToolContext, ToolEntry
@ -59,10 +60,8 @@ def _manage_schedules(
mutate_scheduled_task, schedule_tool_projection,
)
metadata = getattr(ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
operation = str(action or "list").strip().lower()
if metadata.get("presence"):
if presence_caller_binding(ctx) is not None: # a speaker, or work acting for its binding
return _publish_tool_result(ctx, ToolResult(
status="blocked", code="RESOURCE_CONSTRAINT_BLOCKED",
text="⚠️ RESOURCE_CONSTRAINT_BLOCKED: a Presence conversation cannot read or change owner schedules.",
@ -439,11 +438,14 @@ def _register_followup(ctx: ToolContext, task_id: str, drive_root: Any,
# A follow-up from a consciousness turn/tree starts a consciousness root: the
# origin, category and level ride the template; admission derives the rest.
record["task"]["metadata"].update(consciousness_origin_metadata(metadata_src))
presence = metadata_src.get("presence") if isinstance(metadata_src, dict) else None
# The same Presence carrier a promote keeps: a speaker's metadata, or the binding a
# promoted descendant root acts for; the ceiling rides by value and no Project is chosen.
contract = getattr(ctx, "task_contract", None)
if isinstance(presence, dict) and presence and isinstance(contract, dict):
record["task"]["metadata"]["presence"] = dict(presence)
carrier = presence_root_carrier(metadata_src, task_contract=contract)
if carrier and isinstance(contract, dict):
record["task"]["metadata"].update(carrier)
record["task"]["task_contract"] = dict(contract)
record["task"].pop("project_id", None)
from supervisor.queue import ScheduleRefused, ScheduleStoreUnreadable
try:

View file

@ -26,7 +26,7 @@ from ouroboros.project_facts import (
sanitize_project_id,
explicit_project_id_ok,
)
from ouroboros.dialogue_provenance import is_presence_task
from ouroboros.dialogue_provenance import is_presence_task, presence_caller_binding
from ouroboros.focus import normalize_focus
from ouroboros.tools.registry import ToolContext, ToolEntry
from ouroboros.utils import (
@ -55,7 +55,7 @@ def _scope_authority(ctx: ToolContext) -> tuple[str, Dict[str, Any]]:
delegation_role = str(lineage.get("delegation_role") or metadata.get("delegation_role") or "").strip()
root_task_id = str(lineage.get("root_task_id") or metadata.get("root_task_id") or "").strip()
child = bool(parent_task_id) or delegation_role == "subagent"
if is_presence_task(task):
if is_presence_task(task) or presence_caller_binding(ctx) is not None: # a speaker, or acting for its binding
return "presence", metadata
if child:
return "child", metadata

View file

@ -121,10 +121,34 @@ def _task_record(
return record, None
def _queue_snapshot(drive_root: pathlib.Path) -> tuple[Dict[str, Any], bool]:
"""The persisted queue snapshot, and whether one exists that could not be read.
Never written means nothing was ever queued; a written snapshot that cannot be
read, or lists its rows in a shape this reader cannot walk, proves no absence.
"""
path = drive_root / "state" / "queue_snapshot.json"
try:
raw = path.read_text(encoding="utf-8")
except FileNotFoundError:
return {}, False
except (OSError, UnicodeDecodeError):
return {}, True
try:
data = json.loads(raw)
except ValueError:
return {}, True
if not isinstance(data, dict) or any(not isinstance(data.get(key, []), list) for key in ("running", "pending")):
return {}, True
return data, False
def _owner_record(row: Dict[str, Any] | None, queued: Dict[str, Any]) -> Dict[str, Any]:
"""Whose work one task is: a readable result row's binding fact decides, else the queue's own task."""
if row and row.get("presence_binding_id"):
return {"metadata": {"presence": {"binding_id": row["presence_binding_id"]}},
if row and (row.get("presence_binding_id") or row.get("presence_authority_recorded")):
# A readable malformed/empty carrier (or a ceiling whose carrier was lost)
# outranks a stale queue claim: the empty binding grants no scoped read.
return {"metadata": {"presence_binding_authority": {"binding_id": row.get("presence_binding_id") or ""}},
"delegation_role": row.get("delegation_role"), "parent_task_id": row.get("parent_task_id")}
return queued
@ -172,6 +196,7 @@ def _presence_scope_inventory(
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.
An unreadable queue snapshot is a gap too: its queued work cannot be listed, not absent.
"""
from ouroboros.gateway.task_list_scan import raw_result_facts
@ -181,10 +206,12 @@ def _presence_scope_inventory(
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")
snapshot, snapshot_unreadable = _queue_snapshot(drive_root)
if snapshot_unreadable:
gap["queue_snapshot"] = "unreadable" # its queued work cannot be listed; that is not absence
queued: Dict[str, tuple[str, Dict[str, Any]]] = {}
for status in ("running", "pending"):
for item in (snapshot or {}).get(status) or []:
for item in snapshot.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:
@ -410,12 +437,12 @@ def recent_tasks_page(
def _restricted_actor(ctx: ToolContext) -> bool:
"""Children and Presence turns hold no live cross-focus catalogue."""
"""Children and Presence turns, or work acting for a binding, hold no live cross-focus catalogue."""
metadata = getattr(ctx, "task_metadata", {})
metadata = metadata if isinstance(metadata, dict) else {}
return bool(str(metadata.get("parent_task_id") or "").strip()
or str(metadata.get("delegation_role") or "") == "subagent"
or is_presence_task({"metadata": metadata}))
or is_presence_task({"metadata": metadata}) or presence_caller_binding(ctx) is not None)
def _handle_live_roots(ctx: ToolContext, limit: int = 20, offset: int = 0, snapshot: str = "", **_kwargs: Any) -> str:

View file

@ -12,6 +12,7 @@ import logging
import threading
from typing import Any, Dict, Optional
from ouroboros.dialogue_provenance import presence_root_carrier
from ouroboros.task_results import STATUS_FAILED, STATUS_SCHEDULED, write_task_result
from ouroboros.utils import utc_now_iso
@ -497,6 +498,7 @@ def _promote_chat_to_task_outcome(evt: Dict[str, Any], ctx: Any) -> Dict[str, An
else "unconfirmed"
)
transfer = outcome.pop("force_plan_transfer", None)
carrier = presence_root_carrier(evt, task_contract=evt.get("task_contract"))
stored = write_task_result(
ctx.DRIVE_ROOT,
str(outcome.get("task_id") or task_id),
@ -530,8 +532,7 @@ def _promote_chat_to_task_outcome(evt: Dict[str, Any], ctx: Any) -> Dict[str, An
# 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 {}),
**({"metadata": carrier, "source": "presence_promote"} if carrier else {}),
)
admission = stored.get("promotion_admission") if isinstance(stored, dict) else {}
if (

View file

@ -22,6 +22,7 @@ import uuid
from typing import Any, Dict, List
from ouroboros.consciousness_authority import apply_consciousness_authority
from ouroboros.contracts.task_contract import build_task_contract, normalize_allowed_resources
from ouroboros.dialogue_provenance import presence_metadata_binding
from ouroboros.schedule_contract import RESERVED_TEMPLATE_FIELDS, schedule_slug
from ouroboros.skill_loader import skill_identity_collision_names
from ouroboros.utils import atomic_write_json, in_worker_process, read_json_dict, utc_now_iso
@ -684,9 +685,10 @@ def _task_from_schedule(record: Dict[str, Any]) -> Dict[str, Any]:
if existing_contract:
task["task_contract"] = existing_contract
task["task_contract"] = build_task_contract(apply_consciousness_authority(task))
presence = metadata.get("presence")
workspace = task["task_contract"]["workspace"]
if isinstance(presence, dict) and presence and workspace["root"]:
# Presence-bound work (a speaker's follow-up, or one acting for a descendant's binding)
# runs in the folder its inherited contract admitted, with canonical shared memory.
if presence_metadata_binding(metadata) is not None and workspace["root"]:
task.update(
workspace_root=workspace["root"], workspace_mode=workspace["mode"],
memory_mode="shared",

View file

@ -58,12 +58,18 @@ def _presence_target_refused(ctx: Any, evt: Dict[str, Any], task: Dict[str, Any]
The sender is Presence by the host's stamp on the event or, failing that, by
its own live queue row (a delegated descendant's inherited binding authority),
so the fence never rests on one producer remembering to stamp the event.
so the fence never rests on one producer remembering to stamp the event; a
malformed stamp or carrier narrows to nothing. Whose work the target is follows
the read/cancel precedence: its canonical record decides, and the live row
stands in only for a record without Presence provenance.
"""
from ouroboros.dialogue_provenance import presence_metadata_binding, presence_related_work
from ouroboros.dialogue_provenance import (
presence_metadata_binding, presence_related_work, presence_target_record,
)
if "presence_binding_id" in evt:
binding = str(evt.get("presence_binding_id") or "")
stamp = evt.get("presence_binding_id")
binding = stamp.strip() if isinstance(stamp, str) else ""
else:
running = getattr(ctx, "RUNNING", None)
meta = running.get(str(_issuer(evt).get("task_id") or "")) if isinstance(running, dict) else None
@ -71,7 +77,10 @@ def _presence_target_refused(ctx: Any, evt: Dict[str, Any], task: Dict[str, Any]
binding = presence_metadata_binding(row.get("metadata")) if isinstance(row, dict) else None
if binding is None and isinstance(row, dict) and isinstance(row.get("task_contract"), dict):
binding = "" if "capability_ceiling" in row["task_contract"] else None
return binding is not None and not presence_related_work(binding, task)
if binding is None:
return False
target = str(evt.get("target_task_id") or task.get("id") or "").strip()
return not presence_related_work(binding, presence_target_record(ctx.DRIVE_ROOT, target, queue_row=task))
def _refuse_steering_while_cancelling(

View file

@ -183,6 +183,13 @@ def _promoted_force_plan_metadata(evt: dict) -> dict:
return {"metadata": {"force_plan": True, "force_plan_source": source}}
def _presence_promotion(evt: dict) -> bool:
"""A promote carrying Presence authority: a speaker's, or a delegated descendant's binding."""
from ouroboros.dialogue_provenance import presence_root_carrier
return bool(presence_root_carrier(evt, task_contract=evt.get("task_contract")))
def _promote_project_scope(evt: dict) -> str:
"""The project an admitted promote lands in: the explicit one the event carries,
else the project the OWNER MESSAGE it came from already has.
@ -199,7 +206,7 @@ def _promote_project_scope(evt: dict) -> str:
cannot choose a Project. Fail-open: an unreadable store leaves the scope exactly
as the event stated it."""
explicit = str(evt.get("project_id") or "")
if explicit or evt.get("presence") or not isinstance(evt.get("source_ref"), dict):
if explicit or _presence_promotion(evt) or not isinstance(evt.get("source_ref"), dict):
return explicit
try:
from ouroboros.projects_registry import origin_claim_lock, project_id_for_origin
@ -427,7 +434,7 @@ def promote_chat_to_task(evt: dict, ctx: Any) -> dict:
# assignment below turns that answer into the event's stated scope.
implicit_scope = (
not str(evt.get("project_id") or "")
and not evt.get("presence")
and not _presence_promotion(evt)
and isinstance(evt.get("source_ref"), dict)
)
effective_pid = evt["project_id"] = _promote_project_scope(evt)
@ -666,7 +673,7 @@ def _admit_promoted_workspace(evt: dict, ctx: Any, task: dict, *, pid: str, tid:
workspace_repair_hint,
)
if task.get("_presence_origin"):
if _presence_promotion(evt):
# Keep the admitted folder from the inherited contract, never a public
# event's replacement. Presence retains its canonical shared memory.
workspace = task["task_contract"].get("workspace") or {}

View file

@ -554,15 +554,22 @@ def _reject_promoted_after_attachment_stage(
def _apply_presence_promotion_authority(
evt: dict, task: dict, *, objective: str, expected_output: str,
) -> list[dict] | dict:
"""Preserve inherited Presence authority while rebinding the new root."""
"""Preserve inherited Presence authority while rebinding the new root.
A speaker's promote makes the root answer its conversation; a delegated
descendant's carries only the binding it acts for, so its root is that
binding's related work under the same ceiling and never a speaker.
"""
from ouroboros.dialogue_provenance import presence_root_carrier
presence = evt.get("presence") if isinstance(evt.get("presence"), dict) else None
if not presence:
return []
task["_presence_origin"] = True
task["source"] = "presence_promote"
task.setdefault("metadata", {})["presence"] = dict(presence)
contract = evt.get("task_contract") if isinstance(evt.get("task_contract"), dict) else {}
carrier = presence_root_carrier(evt, task_contract=contract)
if not carrier:
return []
if "presence" in carrier:
task["_presence_origin"] = True
task["source"] = "presence_promote"
task.setdefault("metadata", {}).update(carrier)
inherited_manifest = [
dict(row) for row in (contract.get("attachment_manifest") or [])
if isinstance(row, dict)
@ -573,6 +580,9 @@ def _apply_presence_promotion_authority(
"objective": objective,
"expected_output": expected_output,
"attachment_manifest": [],
# The new root owns its objective: a delegated promoter's claims are not its premise.
"acceptance_claims": [],
"success_criteria": [],
})
promoted_contract.pop("lineage", None)
promoted_contract.pop("attachment_manifest_ref", None)

View file

@ -114,7 +114,7 @@ def test_forced_final_speaks_only_what_it_declares(tmp_path, monkeypatch, forced
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
assert result["message"] == "" and result["finish_note"] == "sent the table via the transport tool"
def test_a_malformed_control_body_speaks_nothing_and_keeps_the_host_fallback(tmp_path, monkeypatch):
@ -285,7 +285,7 @@ def test_a_tool_delivered_note_is_never_speech_even_when_owed_work_defers_the_tu
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 result["message"] == "" and result["finish_note"] == note # context, never speech
assert stored["metadata"]["presence_result_text"] == ""
replay = _cached_result(tmp_path, "presence-loop")
assert (replay.outcome, replay.text, replay.work_ref) == ("deferred", "", "later-work")

View file

@ -0,0 +1,248 @@
"""One Presence carrier from producer to admission to reader, and the precedence steer shares.
A delegated descendant's promote and that root's follow-ups keep the binding the
descendant acts for (never the speaker's metadata) and the inherited ceiling. Steer
judges its target by the canonical record first, exactly as read and cancel do, and
the scoped reader states an unreadable queue snapshot as a gap. Deterministic; no
transport sends anything.
"""
from __future__ import annotations
import json
import types
from ouroboros.presence_authority import presence_ceiling_payload
from ouroboros.task_results import load_task_result, write_task_result
from ouroboros.tools.registry import ToolContext
from tests.test_presence_own_work import (
BINDING,
OTHER,
THREAD,
_admitted_child,
_ceiling,
_parent,
_presence,
_queue,
_registry,
_steering_turn,
_supervisor,
_work,
_worker_metadata,
)
def test_one_carrier_is_produced_for_new_roots_and_read_back_everywhere(tmp_path):
from ouroboros.dialogue_provenance import presence_record_binding, presence_root_carrier
from ouroboros.project_facts import resolve_project_id
speaker, authority = {"presence": _presence()}, {"presence_binding_authority": {"binding_id": BINDING}}
ceiling = {"capability_ceiling": presence_ceiling_payload(_ceiling())}
assert presence_root_carrier(speaker) == speaker # a speaker's root answers its conversation
assert presence_root_carrier({**authority, "source": "x"}) == authority # a descendant's: the binding only
for lost in ({"presence": {}}, {"presence_binding_authority": "bad"}):
assert presence_root_carrier(lost) == {"presence_binding_authority": {"binding_id": ""}}
assert presence_root_carrier({}, task_contract=ceiling) == {"presence_binding_authority": {"binding_id": ""}}
assert presence_root_carrier({"source": "owner"}) == {} and presence_root_carrier(None) == {}
for carrier in (speaker, authority):
assert presence_record_binding({"metadata": carrier}) == BINDING
# Presence moves cwd, never the canonical memory scope into a workspace-derived Project.
assert resolve_project_id({"workspace_root": str(tmp_path), "metadata": carrier}) == ""
assert resolve_project_id({"workspace_root": str(tmp_path), "metadata": {}}).startswith("proj_")
def test_a_presence_childs_promote_and_follow_up_stay_its_bindings_work_under_the_ceiling(tmp_path, monkeypatch):
"""A delegated descendant carries the binding, not the speaker. The real promote tool,
the real supervisor admission, the readers, steer and a follow-up of that root keep
that one carrier: the root is this binding's own work under the inherited ceiling,
speaks to no conversation and chooses no Project, workspace or source."""
import supervisor.queue as queue
import supervisor.workers as workers
from ouroboros.dialogue_provenance import is_presence_task, presence_related_work
from ouroboros.owner_mailbox import drain_owner_entries
from ouroboros.peer_roster import maybe_append_roster_note
from ouroboros.project_facts import resolve_project_id
from ouroboros.tools.control import _steer_task
from ouroboros.tools.control_routing import _promote_chat_to_task
from ouroboros.tools.followup import _handle_schedule_followup, _manage_schedules
from ouroboros.tools.project_journal import _scope_authority
from ouroboros.tools.recent_tasks import _restricted_actor
from supervisor.events import _handle_promote_chat_to_task
from tests.test_promote_chat_flow import _confirm_promote
monkeypatch.setattr(workers, "DRIVE_ROOT", tmp_path)
monkeypatch.setattr(queue, "DRIVE_ROOT", str(tmp_path))
_confirm_promote(monkeypatch)
_evt, child = _admitted_child(tmp_path, monkeypatch, _parent(tmp_path, {"presence": _presence()}))
ceiling = child["task_contract"]["capability_ceiling"]
pending, emitted = [], []
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},
)
child_ctx = types.SimpleNamespace(
pending_events=[], current_chat_id=4242, drive_root=tmp_path, budget_drive_root=str(tmp_path),
project_id="", task_id=child["id"], task_metadata=_worker_metadata(child), is_direct_chat=False,
# The promoter's own claim must not become the new root's acceptance premise.
task_contract={**child["task_contract"], "acceptance_claims": [{"claim": "The child's figures are checked."}]},
event_queue=types.SimpleNamespace(
put_nowait=lambda event: emitted.append(event) or _handle_promote_chat_to_task(event, handler_ctx)),
)
out = _promote_chat_to_task(child_ctx, "Compile the full audit", project_name="Widened",
workspace_root=str(tmp_path / "elsewhere"), source="api", predecessor_task_id="")
authority = {"binding_id": BINDING}
[evt], [root] = emitted, pending
assert out.startswith("OK: task") and "Widened" not in out, out
assert evt["presence_binding_authority"] == authority and "presence" not in evt # binding, never speaker
assert [evt[key] for key in ("project_id", "project_name", "workspace_root", "source")] == ["", "", "", ""]
assert evt["task_contract"]["capability_ceiling"] == ceiling
assert (root["id"], root["delegation_role"], root.get("parent_task_id")) == (evt["task_id"], "root", None)
assert root["metadata"]["presence_binding_authority"] == authority
assert "presence" not in root["metadata"] and not is_presence_task(root) # no forced reply or room context
assert root["task_contract"]["capability_ceiling"] == ceiling and root["task_contract"]["acceptance_claims"] == []
assert not root.get("project_id") and resolve_project_id(root) == ""
stored = load_task_result(tmp_path, root["id"])
assert stored["status"] == "scheduled" and stored["metadata"]["presence_binding_authority"] == authority
# Readers and steer of the binding reach it from another conversation; another binding does not.
registry, _turn = _registry(tmp_path, key=THREAD)
assert root["id"] in [row["task_id"] for row in json.loads(registry.execute("recent_tasks", {"limit": 20}))["tasks"]]
assert "PRESENCE_CAPABILITY_BLOCKED" not in registry.execute("get_task_result", {"task_id": root["id"]})
stranger, _stranger = _registry(tmp_path, binding=OTHER)
assert "PRESENCE_CAPABILITY_BLOCKED" in stranger.execute("get_task_result", {"task_id": root["id"]})
steered = _steer_task(_steering_turn(tmp_path, _supervisor(tmp_path, pending=[root]), []), root["id"], "Q2 is in.")
assert "written to its mailbox" in steered and len(drain_owner_entries(tmp_path, root["id"])) == 1
# A follow-up of that root keeps the same carrier and ceiling, never a Project.
root_ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id=root["id"], current_chat_id=4242,
task_contract=root["task_contract"], project_id="widened",
task_metadata={**root["metadata"], "root_task_id": root["id"], "delegation_role": "root"})
scheduled_text = _handle_schedule_followup(root_ctx, objective="Revisit the audit", run_at="2030-01-01T00:00:00Z")
assert scheduled_text.startswith("FOLLOWUP_SCHEDULED"), scheduled_text
[record] = queue.list_scheduled_tasks(tmp_path)["tasks"]
followup = queue._task_from_schedule(record)
assert followup["metadata"]["presence_binding_authority"] == authority and "presence" not in followup["metadata"]
assert followup["task_contract"]["capability_ceiling"] == ceiling
assert not followup.get("project_id") and presence_related_work(BINDING, followup)
# Work acting for a binding reads no owner schedule table, exactly as its speaker cannot.
assert "RESOURCE_CONSTRAINT_BLOCKED" in _manage_schedules(root_ctx, "list")
child_tools = types.SimpleNamespace(task_metadata=_worker_metadata(child), task_contract=child["task_contract"])
assert "RESOURCE_CONSTRAINT_BLOCKED" in _manage_schedules(child_tools, "list")
# An ordinary child's promote is unchanged: no carrier, and its project request stands.
_plain_evt, plain = _admitted_child(tmp_path, monkeypatch, _parent(tmp_path, {}, task_id="owner-root", ceiling=False))
plain_events = []
plain_ctx = types.SimpleNamespace(**{**vars(child_ctx), "task_id": plain["id"], "task_metadata": _worker_metadata(plain),
"task_contract": plain["task_contract"],
"event_queue": types.SimpleNamespace(put_nowait=plain_events.append)})
_promote_chat_to_task(plain_ctx, "Owner-side work", project_name="Chosen", predecessor_task_id="")
assert not {"presence", "presence_binding_authority"} & set(plain_events[0])
assert plain_events[0]["project_name"] == "Chosen"
# Work acting for a binding holds no other owner cross-focus view a speaker's promoted root is denied.
_queue(tmp_path, running=[{"id": "owner-work", "delegation_role": "root", "description": "Owner audit"}])
owner_root = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="owner-root-2",
task_metadata={"delegation_role": "root", "root_task_id": "owner-root-2"})
assert maybe_append_roster_note(owner_root, [], tmp_path) is True # an owner root sees its peers
assert maybe_append_roster_note(root_ctx, [], tmp_path) is False
assert _restricted_actor(root_ctx) and not _restricted_actor(owner_root)
assert (_scope_authority(root_ctx)[0], _scope_authority(owner_root)[0]) == ("presence", "root")
def test_supervisor_steering_follows_the_canonical_binding_over_a_stale_queue_row(tmp_path, monkeypatch):
"""Steer judges whose work the target is like read and cancel: the canonical record
decides (a malformed carrier narrows to nothing); the live row speaks only for a
record without Presence provenance. A malformed sender stamp narrows to nothing."""
import supervisor.queue as queue
from ouroboros.owner_mailbox import drain_owner_entries
from ouroboros.project_dialogue import AGENT_RECEIPT_ID_PREFIX
from ouroboros.tools.control import _steer_task
from supervisor.events import _handle_steer_task
monkeypatch.setattr(queue, "DRIVE_ROOT", str(tmp_path))
names = ("theirs-pending", "theirs-running", "mine-pending", "legacy-pending", "torn-carrier")
live = {name: {"id": name, "delegation_role": "root", "chat_id": 70 + index,
"metadata": {"presence": _presence()}} # every live row claims this binding
for index, name in enumerate(names)}
_work(tmp_path, "theirs-pending", "scheduled", binding=OTHER)
_work(tmp_path, "theirs-running", "running", binding=OTHER)
_work(tmp_path, "mine-pending", "scheduled", key=THREAD) # the canonical record agrees
write_task_result(tmp_path, "legacy-pending", "scheduled", delegation_role="root") # no provenance
write_task_result(tmp_path, "torn-carrier", "scheduled", delegation_role="root",
metadata={"presence": {"binding_id": 7}})
supervisor_ctx = _supervisor(tmp_path, running=[live["theirs-running"]],
pending=[live[name] for name in names if name != "theirs-running"])
turn = _steering_turn(tmp_path, supervisor_ctx, [])
for target in ("mine-pending", "legacy-pending"): # own pending work stays steerable (owner Q2)
assert "written to its mailbox" in _steer_task(turn, target, "New figures."), target
assert len(drain_owner_entries(tmp_path, target)) == 1
for target in ("theirs-pending", "theirs-running", "torn-carrier"):
refused = _steer_task(turn, target, "stop")
assert "STEER_REJECTED" in refused and "presence_work_not_related" in refused, target
assert drain_owner_entries(tmp_path, target) == []
def stamped(stamp, token): # the entries this one event adds
before = len(drain_owner_entries(tmp_path, "mine-pending"))
_handle_steer_task({
"type": "steer_task", "routing_token": token, "target_task_id": "mine-pending",
"message": "stamped", "chat_id": 4242, "client_message_id": f"{AGENT_RECEIPT_ID_PREFIX}{token}",
"issuer": {"kind": "task", "task_id": "presence-turn-1", "root_task_id": "presence-turn-1"},
"presence_binding_id": stamp,
}, supervisor_ctx)
return drain_owner_entries(tmp_path, "mine-pending")[before:]
assert [entry["text"] for entry in stamped(BINDING, "tok-own")] == ["stamped"]
for index, malformed in enumerate((None, 7, {"binding_id": BINDING}, [BINDING], "")):
assert stamped(malformed, f"tok-bad-{index}") == [], malformed
def test_an_unreadable_queue_snapshot_is_a_stated_gap_not_absent_queued_work(tmp_path):
from ouroboros.presence_context import build_presence_context_section
_work(tmp_path, "done-here", "completed", result="Earlier answer")
registry, _ctx = _registry(tmp_path)
page = json.loads(registry.execute("recent_tasks", {"limit": 20}))
assert "read_gap" not in page # never written: nothing was ever queued
snapshot = tmp_path / "state" / "queue_snapshot.json"
snapshot.parent.mkdir(parents=True, exist_ok=True)
for torn in ("{not json", json.dumps(["not", "an", "object"]),
json.dumps({"pending": {"queued-only": {}}, "running": []})):
snapshot.write_text(torn, encoding="utf-8")
page = json.loads(registry.execute("recent_tasks", {"limit": 20}))
assert page["read_gap"] == {"queue_snapshot": "unreadable"}, torn
assert [row["task_id"] for row in page["tasks"]] == ["done-here"]
section = build_presence_context_section(tmp_path, {**_presence(), "instructions": "Be useful."}, "turn-x")
assert "unreadable now: the queue snapshot (queued work); this binding's work may be among them" in section
_queue(tmp_path, pending=[{"id": "queued-only", "delegation_role": "root",
"metadata": {"presence": _presence(key=THREAD)}}])
page = json.loads(registry.execute("recent_tasks", {"limit": 20}))
assert "read_gap" not in page and {row["task_id"] for row in page["tasks"]} == {"queued-only", "done-here"}
def test_malformed_or_lost_canonical_binding_never_borrows_a_queued_claim(tmp_path):
"""A parseable result with Presence authority but no valid binding outranks stale queue A."""
from ouroboros.dialogue_provenance import presence_target_record
tasks = [
{"id": name, "delegation_role": "root", "metadata": {"presence": _presence()}}
for name in ("malformed-result", "lost-carrier")
]
_queue(tmp_path, pending=tasks[:1], running=tasks[1:])
write_task_result(tmp_path, "malformed-result", "scheduled", delegation_role="root",
metadata={"presence": {"binding_id": 7}}, description="secret other task")
write_task_result(tmp_path, "lost-carrier", "running", delegation_role="root", metadata={},
task_contract={"capability_ceiling": presence_ceiling_payload(_ceiling())},
description="lost binding task")
registry, _ctx = _registry(tmp_path)
page = json.loads(registry.execute("recent_tasks", {"limit": 20}))
assert not {"malformed-result", "lost-carrier"} & {row["task_id"] for row in page["tasks"]}
assert "lost-carrier" not in {row["task_id"] for row in page["running"]}
for row in tasks:
record = presence_target_record(tmp_path, row["id"], queue_row=row)
assert record["metadata"] != row["metadata"]
assert "PRESENCE_CAPABILITY_BLOCKED" in registry.execute("get_task_result", {"task_id": row["id"]})

View file

@ -431,6 +431,46 @@ def test_previous_turn_shows_what_a_transport_tool_delivered(tmp_path):
tmp_path, captured[-1]["metadata"]["presence"])
def test_deferred_handoff_keeps_internal_finish_note_out_of_prior_speech(tmp_path):
"""A tool-send note stays context even when owed work makes the turn deferred."""
note = "helper failed; the table was already sent"
captured = []
first = _pointer_turn(tmp_path, "note-e1", {"outcome": "deferred", "text": "", "message": "",
"finish_note": note, "work_ref": "owed-task"})
assert (first.outcome, first.text, first.work_ref) == ("deferred", "", "owed-task")
_pointer_turn(tmp_path, "note-e2", {"outcome": "silent", "text": ""}, captured=captured)
pointer = captured[-1]["metadata"]["presence"]["previous_turn"]
assert pointer["message"] == "" and pointer["finish_note"] == note
section = build_presence_context_section(tmp_path, captured[-1]["metadata"]["presence"])
assert 'finish note "helper failed; the table was already sent"' in section
assert '): "helper failed; the table was already sent"' not in section
def test_legacy_deferred_pointer_is_checked_against_canonical_reply_before_quoting(tmp_path):
"""Old deferred tool-send notes shared `message` with speech; source separates them."""
from ouroboros.presence_bindings import conversation_key
from ouroboros.presence_runner import _previous_turn_path
captured = []
note = "helper failed, result already sent"
_pointer_turn(tmp_path, "old-note", {"outcome": "deferred", "text": "", "message": note,
"work_ref": "owed"})
path = _previous_turn_path(tmp_path, conversation_key("telegram", "bot-1", "room-1", "topic-1"))
historical = path.read_bytes()
_pointer_turn(tmp_path, "after-note", {"outcome": "silent", "text": ""}, captured=captured)
previous = captured[-1]["metadata"]["presence"]["previous_turn"]
assert (previous["message"], previous["finish_note"]) == ("", note)
assert b'"finish_note"' not in historical # source remains unchanged; this is a read projection
assert 'finish note "helper failed, result already sent"' in build_presence_context_section(
tmp_path, captured[-1]["metadata"]["presence"])
_pointer_turn(tmp_path, "old-speech", {"outcome": "deferred", "text": "Still working", "message": "Still working",
"work_ref": "owed"})
_pointer_turn(tmp_path, "after-speech", {"outcome": "silent", "text": ""}, captured=captured)
spoken = captured[-1]["metadata"]["presence"]["previous_turn"]
assert spoken["message"] == "Still working" and "finish_note" not in spoken
def test_previous_turn_reports_the_fate_of_its_deferred_work(tmp_path):
""""Work continues" only while the child runs; a finished child's answer or failure is stated instead."""
from ouroboros.task_results import write_task_result

View file

@ -84,6 +84,13 @@ def test_real_round_limit_delivers_only_the_current_authored_final(tmp_path, mon
def test_exact_host_diagnostic_is_deliverable_when_the_model_authors_it(tmp_path, monkeypatch):
"""Pins a disclosed owner-Q4 limit, not a goal: speech follows terminal authorship.
Host-authored terminal bytes never speak, and a forced final speaks only its typed
declaration; an ordinary final is the model's own text in one channel, so a model
that restates a diagnostic there is heard. Only its wording could tell the two
apart, and this host adds no text filter; the prompt is the boundary there.
"""
_result, host, _calls, _held = _run_loop(tmp_path / "host", monkeypatch, [_read_response(), {"content": ""}])
result, authored, calls, _held = _run_loop(tmp_path / "author", monkeypatch, [{"content": host["result"]}])
assert len(calls) == 1 # ordinary implicit final, no presence_finish required

View file

@ -207,7 +207,12 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# between the internal record and the declared reply; the base sat 10 bytes under.
# 13400 -> 13700 (TZ2 descendant authority): one sentence the chapter lacked — a
# delegated descendant's inherited binding authority, apart from the speaker metadata.
"docs/architecture/12-host-service-companions-and-chat-ids.md": 13700,
# 13700 -> 13950 (TZ2 repair, measured 13918): that sentence now names what the
# descendant's promote/follow-up roots carry and the canonical-first steer precedence
# (replacing the live-row clause), and "host diagnostics" states its ordinary-final limit.
# 13950 -> 14100 (TZ2 review): a deferred tool-delivery finish note is
# carried separately from prior speech in the same previous-turn pointer.
"docs/architecture/12-host-service-companions-and-chat-ids.md": 14100,
# 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.
@ -237,7 +242,9 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = {
# cancel clause with the own-binding rule and name the forced declaration.
# 95000 -> 95150 (TZ2 descendant authority): the own-binding bullet names how a delegated
# descendant is a Presence caller (inherited binding authority, never speaker metadata).
"docs/development/06-rules-by-change-class.md": 95150,
# 95150 -> 95200 (TZ2 repair, measured 95191): the promotion/follow-up clause names the
# one carrier it copies instead of "the Presence metadata".
"docs/development/06-rules-by-change-class.md": 95200,
"docs/development/07-managed-update-rule.md": 4166,
"docs/development/08-mutation-attribution-rule.md": 2899,
"docs/development/09-process-custody-rule.md": 10028,