Stamp lineage on every tool row and task_id on the two direct-chat rows

`tools.jsonl` rows copied `root_task_id` / `delegation_role` from task
metadata, so a direct root — whose metadata carries neither — logged rows
with no lineage at all while pooled roots and children did (h7 M2). The
row writer now resolves lineage through the ONE resolver already used by
usage, cost and acceptance (`task_results.resolve_task_lineage`: a direct
root is its own root, `delegation_role=root`), and copies only the
parent id and depth from metadata. Rows without a task stay as they were.

The `proactive_message` and `schedule_task_from_direct_chat` event rows
carried no `task_id`. They now do, which makes them visible in the task's
own log stream: `gateway/logs.py` filters `/api/logs/{name}?task_id=` by
task/parent/root id, so a proactive owner message or a subagent scheduled
from a direct turn appears in that task's Logs tail instead of only in
the global events log. `_schedule_task` stays at its 300-line function
cap by dropping a local import the module already had.
This commit is contained in:
Ouroboros 2026-09-21 12:58:42 +03:00
parent b04395e0a5
commit 288ae40f4b
5 changed files with 123 additions and 4 deletions

View file

@ -125,7 +125,7 @@ scanned data-relative path to be covered by a row here (count-anchored both ways
| `logs/chat.jsonl` | `supervisor/message_bus.py` (+presence, project summaries) | `direction` + optional `type`; no version — accepted (projection replayed by chain-aware readers) | rotated 800 KB → `archive/chat_*.jsonl`; archive chain WARN at 100 MB | newest generation lost; consolidation cursor reports gap (recoverable) |
| `logs/progress.jsonl` | `supervisor/message_bus.py` (+plan review) | `type: send_message`, `is_progress` | rotated 800 KB → `archive/progress_*.jsonl`; 8 MB WARN = rotation broken | current segment lost; readers archive-chain-aware |
| `logs/events.jsonl` | ~60 modules via `append_jsonl` (+`delegate_custody.emit`) | universal `type` discriminator — accepted (per-type payloads owned by emitters) | rotated 800 KB → `archive/events_*.jsonl`; custody readers (replay, fault tail-scan, `complete_custody_rows`, settled-terminal chain cursor, legacy-usage import, swarm rollup, worker-boot verify) are chain-aware; delegated start rows reference their complete raw replay envelope in `observability/blobs/**` (blob before event), with inline legacy rows still readable; `task_received` omits only an equal metadata mirror of its top-level task contract; 8 MB live WARN = rotation broken; 100 MB chain WARN = replay degradation | delegated-run custody destroyed (chain incl. archive segments): open runs invisible/unreapable; lineage, citations, legacy-usage source lost |
| `logs/tools.jsonl` | `ouroboros/loop_tool_execution.py` (+budget-drive mirror) | `type: tool_call`, untruncated args | rotated 800 KB → `archive/tools_*.jsonl`; tail readers (api_logs_tail, task_events) archive-backfill; 8 MB WARN = rotation broken | untruncated tool record + `result_ref` pointers lost |
| `logs/tools.jsonl` | `ouroboros/loop_tool_execution.py` (+budget-drive mirror) | `type: tool_call` + task lineage (`root_task_id`/`delegation_role` from `resolve_task_lineage`, a direct root is its own root), untruncated args | rotated 800 KB → `archive/tools_*.jsonl`; tail readers (api_logs_tail, task_events) archive-backfill; 8 MB WARN = rotation broken | untruncated tool record + `result_ref` pointers lost |
| `logs/supervisor.jsonl` | supervisor family, `process_custody`, gateway control, server shutdown | `type` (+secondary `event_type`) | rotated 800 KB → `archive/supervisor_*.jsonl` + 8 MB tripwire; tail readers (`memory.read_jsonl_tail`, api_logs_tail) archive-backfill | reap receipts, rescue disclosures, shutdown causes lost |
| `logs/task_reflections.jsonl` | `ouroboros/reflection.py` (+ project-scoped copy under `projects/<id>/logs/`) | full rows unversioned; pointer rows `type: project_reflection_pointer` | rotated 800 KB → `archive/task_reflections_*.jsonl` + 8 MB tripwire; tail-20 read archive-backfills; project-scoped copies follow project retention (never age-pruned) | inter-task memory-carry signal lost |
| `logs/containment_faults.jsonl` | `ouroboros/delegate_custody.py` (mirrored to events.jsonl) | `type` ∈ CONTAINMENT_FAULT/RESOLVED joined on run_id | unbounded BY DESIGN — read whole so an open fault never ages out — accepted | health invariant degrades to the 4 MB events tail scan (the regression this file fixed) |

View file

@ -22,6 +22,7 @@ from ouroboros.config import (
)
from ouroboros.deadline_utils import deadline_remaining_sec
from ouroboros.observability import new_call_id, persist_call
from ouroboros.task_results import resolve_task_lineage
from ouroboros.tool_capabilities import (
FOREGROUND_MUTATIVE_TOOLS,
PARALLEL_SAFE_ENQUEUE_TOOLS,
@ -162,7 +163,14 @@ def _tool_task_metadata(tools: ToolRegistry) -> Dict[str, Any]:
def _append_tool_log(tools: ToolRegistry, drive_logs: pathlib.Path, payload: Dict[str, Any]) -> None:
meta = _tool_task_metadata(tools)
for key in ("parent_task_id", "root_task_id", "delegation_role", "task_depth"):
if task_id := str(payload.get("task_id") or "").strip():
# ONE lineage resolver (a direct root is its own root), so every task row
# carries root_task_id/delegation_role for the task log stream readers.
lineage = resolve_task_lineage(task_id, metadata=meta)
payload["root_task_id"] = lineage["root_task_id"]
if role := lineage["delegation_role"] or ("root" if lineage["is_root_task"] else ""):
payload["delegation_role"] = role
for key in ("parent_task_id", "task_depth"):
if meta.get(key) not in (None, ""):
payload[key] = meta.get(key)
append_jsonl(drive_logs / "tools.jsonl", payload)

View file

@ -250,6 +250,7 @@ def _send_user_message(ctx: ToolContext, text: str, reason: str = "") -> str:
append_jsonl(ctx.drive_logs() / "events.jsonl", {
"ts": utc_now_iso(),
"type": "proactive_message",
"task_id": str(getattr(ctx, "task_id", "") or ""),
"reason": reason,
"transport_mode": mode,
"text_preview": text[:200],

View file

@ -672,18 +672,18 @@ def _schedule_task(ctx: ToolContext, internal: Dict[str, Any] | None = None, /,
current_depth=current_depth, new_depth=new_depth, max_depth=max_depth,
)
current_task_id = str(getattr(ctx, "task_id", "") or "")
if getattr(ctx, 'is_direct_chat', False):
from ouroboros.utils import append_jsonl
try:
append_jsonl(ctx.drive_logs() / "events.jsonl", {
"ts": utc_now_iso(),
"type": "schedule_task_from_direct_chat",
"task_id": current_task_id,
"description": objective[:200],
"warning": "schedule_subagent called from direct chat context — potential duplicate work",
})
except Exception:
pass
current_task_id = str(getattr(ctx, "task_id", "") or "")
parent_task_id = str(current_task_id or metadata.get("parent_task_id") or "").strip()
root_task_id_seed = str(metadata.get("root_task_id") or current_task_id or "").strip()
session_id = str(metadata.get("session_id") or "")

View file

@ -0,0 +1,110 @@
"""M2: every ``tools.jsonl`` row of a task carries ``root_task_id`` and
``delegation_role`` through the ONE lineage resolver
(``task_results.resolve_task_lineage``; a direct root is its own root), and the
``proactive_message`` / ``schedule_task_from_direct_chat`` event rows carry the
``task_id`` that makes them visible in the task's own log stream
(``gateway/logs.py`` filters rows by task/parent/root id). Rows that belong to
no task stay exactly as they were.
"""
from __future__ import annotations
import json
import pathlib
import queue
import types
import pytest
from ouroboros.loop_tool_execution import _append_tool_log
def _rows(path: pathlib.Path) -> list[dict]:
return [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines() if line.strip()]
def _tools(meta: dict | None, **attrs) -> types.SimpleNamespace:
ctx = types.SimpleNamespace(task_metadata=meta if meta is not None else {}, **attrs)
return types.SimpleNamespace(_ctx=ctx)
@pytest.mark.parametrize("shape,meta,expected", [
("direct_root", {}, {"root_task_id": "t-self", "delegation_role": "root"}),
("pooled_root", {"delegation_role": "root", "root_task_id": "t-self"},
{"root_task_id": "t-self", "delegation_role": "root"}),
("child", {"parent_task_id": "t-parent", "root_task_id": "t-root", "delegation_role": "subagent"},
{"root_task_id": "t-root", "delegation_role": "subagent", "parent_task_id": "t-parent"}),
("headless_child", {"parent_task_id": "t-parent", "root_task_id": "t-root", "delegation_role": "subagent",
"headless_child_drive_root": "/elsewhere", "task_depth": 2},
{"root_task_id": "t-root", "delegation_role": "subagent", "parent_task_id": "t-parent", "task_depth": 2}),
])
def test_every_task_row_carries_lineage_from_the_one_resolver(tmp_path, shape, meta, expected):
logs = tmp_path / "logs"
logs.mkdir()
_append_tool_log(_tools(meta), logs, {"type": "tool_call", "tool": "read_file", "task_id": "t-self"})
[row] = _rows(logs / "tools.jsonl")
for key, value in expected.items():
assert row[key] == value, (shape, key, row)
if "parent_task_id" not in expected:
assert "parent_task_id" not in row, (shape, row) # a root invents no parent
def test_budget_root_mirror_row_carries_the_same_lineage(tmp_path):
logs = tmp_path / "child" / "logs"
logs.mkdir(parents=True)
budget = tmp_path / "budget"
meta = {"parent_task_id": "t-parent", "root_task_id": "t-root", "delegation_role": "subagent",
"budget_drive_root": str(budget)}
_append_tool_log(_tools(meta), logs, {"type": "tool_call", "tool": "run_command", "task_id": "t-child"})
for path in (logs / "tools.jsonl", budget / "logs" / "tools.jsonl"):
[row] = _rows(path)
assert row["root_task_id"] == "t-root" and row["delegation_role"] == "subagent", path
def test_rows_without_a_task_stay_as_today(tmp_path):
logs = tmp_path / "logs"
logs.mkdir()
_append_tool_log(_tools({}), logs, {"type": "tool_call", "tool": "read_file", "task_id": ""})
[row] = _rows(logs / "tools.jsonl")
assert "root_task_id" not in row and "delegation_role" not in row, row
class _Queue:
def __init__(self):
self.items: list = []
def put_nowait(self, item):
self.items.append(item)
put = put_nowait
def test_proactive_message_row_carries_the_task_id(tmp_path):
from ouroboros.tools.control import _send_user_message
ctx = types.SimpleNamespace(
current_chat_id=7, pending_events=[], drive_root=None, task_id="t-live",
task_metadata={"root_task_id": "t-live"}, event_queue=_Queue(),
drive_logs=lambda: tmp_path,
)
assert "OK" in _send_user_message(ctx, "a word while I work", reason="progress")
[row] = [r for r in _rows(tmp_path / "events.jsonl") if r["type"] == "proactive_message"]
assert row["task_id"] == "t-live", row
assert row["reason"] == "progress" and row["text_preview"] == "a word while I work"
def test_schedule_from_direct_chat_row_carries_the_task_id(tmp_path, monkeypatch):
from ouroboros.tools.control import _schedule_task
from ouroboros.tools.registry import ToolContext
from tests._shared import configure_test_subagent
subagent_id = configure_test_subagent(monkeypatch)
repo, drive = tmp_path / "repo", tmp_path / "data"
repo.mkdir()
(drive / "logs").mkdir(parents=True)
ctx = ToolContext(repo_dir=repo, drive_root=drive, task_id="direct-1", is_direct_chat=True,
current_chat_id=5, event_queue=queue.Queue())
out = _schedule_task(ctx, subagent_id=subagent_id, objective="scout X", expected_output="Y")
assert "queued" in out, out
[row] = [r for r in _rows(drive / "logs" / "events.jsonl") if r["type"] == "schedule_task_from_direct_chat"]
assert row["task_id"] == "direct-1", row
assert row["description"] == "scout X" and "duplicate" in row["warning"]