diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py index a801e34766..4f9a4b7b58 100644 --- a/loopx/capabilities/manager_context/roundtrip.py +++ b/loopx/capabilities/manager_context/roundtrip.py @@ -63,6 +63,8 @@ "initial_delivery_receipt_unavailable", "original_route_or_return_delivery_unavailable", "delivery_state_unreadable", + "return_transport_unavailable", + "manager_return_payload_conflict", } class ReturnResolutionBlocked(ValueError, RuntimeError): @@ -597,6 +599,92 @@ def _retry_state(state, now, *, error): return result +def _resolve_delivery_sender(external_sender): + """Resolve a manager return transport to ``(sender, attempt_aware)``. + + Two transport shapes are legitimate: an object exposing a callable + ``send_with_attempt`` (the attempt-aware protocol) and a bare callable + transport. The predicate is shape-based on purpose — the transports + share no base class, so an ``isinstance`` check would reject valid + transports. ``(None, False)`` means the transport cannot deliver: + callers must record ``return_transport_unavailable`` before invoking + instead of letting the call raise ``TypeError`` mid-loop. + """ + sender = getattr(external_sender, "send_with_attempt", None) + if callable(sender): + return sender, True + if callable(external_sender): + return external_sender, False + return None, False + + +_MANAGER_RETURN_PAYLOAD_FIELDS = ( + "role", + "text", + "turn_id", + "origin", + "attachments", + "goal_draft", +) + +# Canonical per-field comparison defaults for the delivery-side replay +# check. The append path drops falsy optional fields together with their +# column, so rows written through the store never carry literal empty +# values; externally written or migrated rows can. Folding a missing key +# and an explicit ``None`` onto each field's canonical empty value keeps +# the comparison closed over those stored shapes: ``attachments`` is +# list-valued and normalizes to ``[]``, the scalar fields to ``None``. +_MANAGER_RETURN_PAYLOAD_DEFAULTS = { + "role": None, + "text": None, + "turn_id": None, + "origin": None, + "attachments": [], + "goal_draft": None, +} + + +def _normalized_return_field(field, value): + """Fold a missing column and an explicit ``None`` into one value.""" + if value is not None: + return value + return _MANAGER_RETURN_PAYLOAD_DEFAULTS[field] + + +def _manager_return_payload(text, turn): + """The transcript payload this pump is about to append for ``turn``.""" + return { + "role": "agent", + "text": text, + "turn_id": turn["turn_id"], + "origin": "manager_followup", + } + + +def _replayed_return_payload_conflict(store, session_id, message_id, incoming): + """First divergent field if ``message_id`` is already recorded with + different payload, else ``None``. + + Same-id replay at the append layer is intentional replay-ignore + semantics: the store returns the existing row without comparing. A + silent merge is safe for generic appends but not for manager return + delivery, where the same handoff id arriving with a different + conclusion must surface as a delivery fault instead of leaving the + old text standing behind a delivered receipt. The comparison therefore + lives here, on the delivery path, before anything is sent or appended. + """ + for row in store.messages(session_id): + if row.get("message_id") != message_id: + continue + for field in _MANAGER_RETURN_PAYLOAD_FIELDS: + existing = _normalized_return_field(field, row.get(field)) + candidate = _normalized_return_field(field, incoming.get(field)) + if existing != candidate: + return field + return None + return None + + def _drain_exact(root, registry, store, external_sender, *, now, cancelled): processed = 0 for path in iter_result_paths(_root(root) / "replies"): @@ -650,6 +738,36 @@ def _drain_exact(root, registry, store, external_sender, *, now, cancelled): try: if cancelled(): return processed + conflict_field = _replayed_return_payload_conflict( + store, + route["session_id"], + mid, + _manager_return_payload(text, turn), + ) + if conflict_field is not None: + # A payload conflict is a data fault, not a transport fault: + # the transcript already holds this handoff message_id with + # different content. Record the existing explicit_unverified + # status with the typed error code instead of delivering or + # retrying; a retry can never fix a payload mismatch and + # would only multiply the divergence, while a silent merge + # would leave the old text behind a delivered receipt. + logging.getLogger(__name__).warning( + "Manager return payload conflict on %s: %s", + mid, + conflict_field, + ) + _write_exact_return_state( + root, + registry, + context, + { + "status": "explicit_unverified", + "error": "manager_return_payload_conflict", + }, + ) + processed += 1 + continue store.append_message( route["session_id"], role="agent", @@ -786,11 +904,24 @@ def record_attempt(value): preserve_admission=True, ) - sender = getattr(external_sender, "send_with_attempt", None) + sender, attempt_aware = _resolve_delivery_sender(external_sender) + if sender is None: + _write_exact_return_state( + root, + registry, + context, + _retry_state( + state, + now, + error="return_transport_unavailable", + ), + ) + processed += 1 + continue sent = ( sender(route, session, turn, text, record_attempt) - if callable(sender) - else external_sender(route, session, turn, text) + if attempt_aware + else sender(route, session, turn, text) ) if sent.get("reply_verified") is not True: if sent.get("external_write_performed") is True: @@ -962,6 +1093,34 @@ def drain(root, registry, store, external_sender, *, now=None, cancelled=lambda: mid = "handoff." + _hash([row["request_id"], path.stem]) if cancelled(): return processed + conflict_field = _replayed_return_payload_conflict( + store, + route["session_id"], + mid, + _manager_return_payload(text, turn), + ) + if conflict_field is not None: + # Same data-fault rule as the exact loop: the transcript + # already holds this handoff message_id with different + # content. Record the existing explicit_unverified status + # with the typed error code instead of delivering or + # folding the conflict into the transport retry path + # below; a retry can never fix a payload mismatch and + # would only multiply the divergence. + logging.getLogger(__name__).warning( + "Manager return payload conflict on %s: %s", + mid, + conflict_field, + ) + _write( + state_path, + { + "status": "explicit_unverified", + "error": "manager_return_payload_conflict", + }, + ) + processed += 1 + continue store.append_message( route["session_id"], role="agent", @@ -1116,11 +1275,29 @@ def record_attempt(value): }, ) - sender = getattr(external_sender, "send_with_attempt", None) + sender, attempt_aware = _resolve_delivery_sender(external_sender) + if sender is None: + attempts = int(state.get("attempts", 0)) + 1 + _write( + state_path, + { + "status": "retry_pending", + "attempts": attempts, + "error": "return_transport_unavailable", + "retry_at": ( + now + + timedelta( + seconds=min(300, 5 * 2 ** min(attempts, 6)) + ) + ).isoformat(), + }, + ) + processed += 1 + continue sent = ( sender(route, session, turn, text, record_attempt) - if callable(sender) - else external_sender(route, session, turn, text) + if attempt_aware + else sender(route, session, turn, text) ) if sent.get("reply_verified") is not True: if sent.get("external_write_performed") is True: diff --git a/loopx/chat_store.py b/loopx/chat_store.py index 94f0dc9901..c92ba7fa37 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -561,6 +561,12 @@ def latest_session( agent_id: str, channel_id: str | None = None, ) -> dict[str, Any] | None: + """Return the newest resumable session row on one explicit channel. + + Omitting ``channel_id`` pins the exact ``goal.`` channel: + a single namespace key, not a search over that goal's recorded + conversations. Use ``list_sessions`` for intentional discovery. + """ candidates = self.resumable_session_candidates( goal_id=goal_id, agent_id=agent_id, @@ -575,7 +581,12 @@ def resumable_session_candidates( agent_id: str, channel_id: str | None = None, ) -> list[dict[str, Any]]: - """Return matching storage facts without deciding Goal identity.""" + """Return matching storage facts without deciding Goal identity. + + Omitting ``channel_id`` pins the exact ``goal.`` channel: + a single namespace key, not a search over that goal's recorded + conversations. Use ``list_sessions`` for intentional discovery. + """ return [ candidate @@ -594,7 +605,14 @@ def session_candidates( agent_id: str, channel_id: str | None = None, ) -> list[dict[str, Any]]: - """Return matching Session records without lifecycle filtering.""" + """Return matching Session records without lifecycle filtering. + + Omitting ``channel_id`` pins the exact ``goal.`` channel: + a single namespace key, not a search over that goal's recorded + conversations. Use ``list_sessions`` for intentional discovery. + Raises ``ValueError`` when ``channel_id`` and ``goal_id`` are both + omitted. + """ if goal_id is None and channel_id is None: raise ValueError("channel_id is required when goal_id is omitted") diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 27d95ee802..f48d7aebb9 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -247,7 +247,7 @@ }, { "site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1", - "line": 879, + "line": 1010, "column": 22, "kind": "codec_read", "api": "load_project_registry", diff --git a/tests/test_chat_store_input_validation.py b/tests/test_chat_store_input_validation.py index 7a8c2fd9a2..c537ec0196 100644 --- a/tests/test_chat_store_input_validation.py +++ b/tests/test_chat_store_input_validation.py @@ -3,6 +3,7 @@ import pytest +import loopx.chat_store as chat_store from loopx.chat_store import CHAT_SESSION_SCHEMA_VERSION, ChatSessionStore @@ -20,3 +21,113 @@ def test_session_id_cannot_escape_sessions_directory(tmp_path: Path) -> None: with pytest.raises(ValueError, match="session_id"): store.load_session("..") + + +def _dedup_store(tmp_path: Path) -> tuple[ChatSessionStore, str]: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + return store, session_id + + +def test_same_message_id_same_payload_returns_existing_row(tmp_path: Path) -> None: + store, session_id = _dedup_store(tmp_path) + first = store.append_message( + session_id, + role="agent", + text="durable completion", + turn_id="turn-one", + message_id="replay-same", + ) + # Append-side replay-ignore is intentional base semantics: a same-id + # retry stays idempotent at the transcript layer, whatever the payload. + replayed = store.append_message( + session_id, + role="agent", + text="durable completion", + turn_id="turn-one", + message_id="replay-same", + ) + assert replayed == first + assert store.messages(session_id) == [first] + + +def test_same_message_id_different_text_replay_returns_existing_row( + tmp_path: Path, +) -> None: + store, session_id = _dedup_store(tmp_path) + first = store.append_message( + session_id, + role="agent", + text="original text", + message_id="replay-text", + ) + # The store deliberately ignores same-id replays without comparing + # payloads: conflict detection for manager return delivery lives on the + # drain side, so the append path must stay replay-ignore (this pins the + # base semantics the generic transcript tests rely on). + replayed = store.append_message( + session_id, + role="agent", + text="rewritten text", + message_id="replay-text", + ) + assert replayed == first + assert store.messages(session_id) == [first] + + +def test_same_message_id_different_role_replay_returns_existing_row( + tmp_path: Path, +) -> None: + store, session_id = _dedup_store(tmp_path) + first = store.append_message( + session_id, + role="user", + text="hello there", + message_id="replay-role", + ) + replayed = store.append_message( + session_id, + role="agent", + text="hello there", + message_id="replay-role", + ) + assert replayed == first + assert store.messages(session_id) == [first] + + +def test_only_created_at_difference_is_not_a_conflict( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store, session_id = _dedup_store(tmp_path) + stamps = iter( + [ + "2026-01-01T00:00:00+00:00", + "2026-01-02T00:00:00+00:00", + ] + ) + monkeypatch.setattr(chat_store, "utc_now", lambda: next(stamps)) + first = store.append_message( + session_id, + role="user", + text="identical body", + message_id="replay-time", + ) + replayed = store.append_message( + session_id, + role="user", + text="identical body", + message_id="replay-time", + ) + assert replayed == first + assert replayed["created_at"] == "2026-01-01T00:00:00+00:00" + assert store.messages(session_id) == [first] diff --git a/tests/test_manager_context_roundtrip.py b/tests/test_manager_context_roundtrip.py index 64939a1cfd..9803a08cf6 100644 --- a/tests/test_manager_context_roundtrip.py +++ b/tests/test_manager_context_roundtrip.py @@ -992,3 +992,316 @@ def sender(_route, _session, _turn, payload): assert sent[0].startswith("协作回复 · worker") assert "Private receiver rationale" not in json.dumps(sent) assert all(not row.get("collaboration") for row in snapshot["messages"]) + + +def test_unresolvable_return_transport_records_retry_without_fake_receipt(flow): + root, registry, store, create = flow + session, _, receipt = create(True) + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Checked") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Checked with the current evidence and found no gap.", + ) + + class InertTransport: + """Neither attempt-aware nor callable: cannot deliver anything.""" + + drain(root, registry, store, InertTransport()) + state = json.loads( + (_root(root) / "replies" / rid / "conclusion.delivery.json").read_text() + ) + assert state["status"] == "retry_pending" + assert state["error"] == "return_transport_unavailable" + assert "delivered_at" not in state and "message_id" not in state + projected = reply_status(root, receipt)[0] + assert projected["status"] == "retry_pending" + assert projected["error"] == "return_transport_unavailable" + + +def test_attempt_only_transport_still_delivers_without_call_protocol(flow): + root, registry, store, create = flow + session, _, receipt = create(True) + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Checked") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Recorded the validation result for the existing plan.", + ) + + class Transport: + """Attempt-aware only: the instance itself is not callable.""" + + def __init__(self): + self.calls = 0 + + def send_with_attempt(self, route, session, turn, text, record_attempt): + self.calls += 1 + return {"reply_verified": True, "idempotency_key": "sha256:provider-proof"} + + transport = Transport() + drain(root, registry, store, transport) + assert transport.calls == 1 + state = json.loads( + (_root(root) / "replies" / rid / "conclusion.delivery.json").read_text() + ) + assert state["status"] == "delivered" + assert state["provider_receipt"] == "sha256:provider-proof" + assert reply_status(root, receipt)[0]["status"] == "delivered" + + +def test_private_channel_completes_transcript_only_without_sender(flow): + root, registry, store, create = flow + session, turn, receipt = create() + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Private deliberation") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Completed the private check against the existing plan.", + ) + + def unexpected(*_): + raise AssertionError("private channel must deliver without any transport") + + drain(root, registry, store, unexpected) + state = json.loads( + (_root(root) / "replies" / rid / "conclusion.delivery.json").read_text() + ) + assert state["status"] == "delivered" + returned = [ + row + for row in store.messages(session["session_id"]) + if row.get("origin") == "manager_followup" + ] + assert len(returned) == 1 and returned[0]["turn_id"] == turn["turn_id"] + assert reply_status(root, receipt)[0]["status"] == "delivered" + + +def test_payload_conflict_records_terminal_state_without_retry(flow): + root, registry, store, create = flow + session, turn, receipt = create(True) + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Checked") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Checked with the current evidence and found no gap.", + ) + # Pre-seed the transcript with the same handoff message_id carrying a + # different payload. A payload conflict is a data fault, not a transport + # fault: drain must stop in the existing explicit_unverified status with + # the typed error code — no new status word, so every existing reader + # keeps its semantics — instead of folding the conflict into the + # transport retry path. + conflicting_mid = "handoff." + _hash([rid, "conclusion"]) + store.append_message( + session["session_id"], + role="agent", + text="An earlier, different version of this conclusion.", + turn_id=turn["turn_id"], + origin="manager_followup", + message_id=conflicting_mid, + ) + deliveries = [] + + def recording_transport(route, session, turn, text): + deliveries.append(text) + return {"reply_verified": True, "idempotency_key": "sha256:provider-proof"} + + drain(root, registry, store, recording_transport) + state_path = _root(root) / "replies" / rid / "conclusion.delivery.json" + state = json.loads(state_path.read_text()) + assert state["status"] == "explicit_unverified" + assert state["error"] == "manager_return_payload_conflict" + assert "delivered_at" not in state and "message_id" not in state + assert "attempts" not in state and "retry_at" not in state + assert deliveries == [] + # The recorded state is terminal: a second pass neither retries the + # transport nor rewrites the record. + drain(root, registry, store, recording_transport) + assert json.loads(state_path.read_text()) == state + assert deliveries == [] + projected = reply_status(root, receipt)[0] + assert projected["status"] == "explicit_unverified" + assert projected["error"] == "manager_return_payload_conflict" + + +def test_stored_empty_attachments_list_replays_without_delivery_conflict(flow): + root, registry, store, create = flow + session, turn, receipt = create() + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Private deliberation") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Completed the private check against the existing plan.", + ) + + def unexpected(*_): + raise AssertionError("private channel must deliver without any transport") + + # First pass writes the transcript row (no ``attachments`` key) and the + # delivered receipt. + drain(root, registry, store, unexpected) + messages_path = ( + store.root / "sessions" / session["session_id"] / "messages.jsonl" + ) + rows = [ + json.loads(line) + for line in messages_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + handoff_rows = [row for row in rows if row.get("message_id", "").startswith("handoff.")] + assert len(handoff_rows) == 1 and "attachments" not in handoff_rows[0] + # An externally written or migrated row can carry a literal empty + # ``attachments`` list where the delivery payload has no key at all: + # field normalization folds both onto the canonical empty list, so the + # same-id replay still delivers instead of reporting a phantom payload + # conflict (a bare ``.get()`` comparison would raise one here). + for row in rows: + if row.get("message_id", "").startswith("handoff."): + row["attachments"] = [] + messages_path.write_text( + "".join(json.dumps(row) + "\n" for row in rows), + encoding="utf-8", + ) + state_path = _root(root) / "replies" / rid / "conclusion.delivery.json" + state_path.unlink() + drain(root, registry, store, unexpected) + state = json.loads(state_path.read_text()) + assert state["status"] == "delivered" + stored = [ + row + for row in store.messages(session["session_id"]) + if row.get("message_id", "").startswith("handoff.") + ] + assert len(stored) == 1 and stored[0]["attachments"] == [] + + +def test_exact_loop_payload_conflict_sets_terminal_state_without_retry( + flow, + monkeypatch, +): + root, registry, store, create = flow + session, turn, receipt = create(True) + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Checked") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Checked with the current evidence and found no gap.", + ) + conflicting_mid = "handoff." + _hash([rid, "conclusion"]) + store.append_message( + session["session_id"], + role="agent", + text="An earlier, different version of this conclusion.", + turn_id=turn["turn_id"], + origin="manager_followup", + message_id=conflicting_mid, + ) + # The exact admission machinery needs a source-session registry fixture + # this module does not build; stub the context construction and the + # settlement writer so the exact loop's own delivery-side conflict + # check is what runs against the real divergent row above. + import loopx.capabilities.manager_context.roundtrip as return_roundtrip + from loopx.control_plane.projects.registry_codec import ( + SOURCE_SESSION_PROFILE_ID, + ) + + settlements = [] + + def fake_settlement(root, registry, context, value, **kwargs): + settlements.append(dict(value)) + + monkeypatch.setattr( + return_roundtrip, "_write_exact_return_state", fake_settlement + ) + + def fake_context(root, registry, store, path, state_path, now): + reply = return_roundtrip._read(path) + return { + "reply": {**reply, "phase": path.stem}, + "row": {"request_id": rid, "agent_id": "worker"}, + "route": {"session_id": session["session_id"]}, + "session": session, + "turn": turn, + "store": store, + "state_path": state_path, + "state": {}, + "token": "fixture", + "source_id": "lark:om_fixture_source", + } + + monkeypatch.setattr(return_roundtrip, "_exact_return_context", fake_context) + profile = json.loads(registry.read_text()) + profile["profile_id"] = SOURCE_SESSION_PROFILE_ID + registry.write_text(json.dumps(profile)) + deliveries = [] + + def recording_transport(route, session, turn, text): + deliveries.append(text) + return {"reply_verified": True, "idempotency_key": "sha256:provider-proof"} + + drain(root, registry, store, recording_transport) + assert deliveries == [] + assert len(settlements) == 1 + assert settlements[0] == { + "status": "explicit_unverified", + "error": "manager_return_payload_conflict", + } + + +def test_explicit_channel_return_completes_transcript_only_without_sender(flow): + root, registry, store, create = flow + session, turn, receipt = create() + assert session["channel_id"] == "manager" + rid = receipt["request_id"] + acknowledge(root, "research", "worker", rid, "adopt", "Checked") + report( + root, + "research", + "worker", + rid, + "conclusion", + "Completed the check against the existing plan.", + ) + + def unexpected(*_): + raise AssertionError( + "explicit private channel must complete without any transport" + ) + + drain(root, registry, store, unexpected) + state = json.loads( + (_root(root) / "replies" / rid / "conclusion.delivery.json").read_text() + ) + assert state["status"] == "delivered" + returned = [ + row + for row in store.messages(session["session_id"]) + if row.get("origin") == "manager_followup" + ] + assert len(returned) == 1 and returned[0]["turn_id"] == turn["turn_id"] + assert reply_status(root, receipt)[0]["status"] == "delivered"