Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
189 changes: 183 additions & 6 deletions loopx/capabilities/manager_context/roundtrip.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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"):
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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:
Expand Down
22 changes: 20 additions & 2 deletions loopx/chat_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.<goal_id>`` 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,
Expand All @@ -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.<goal_id>`` channel:
a single namespace key, not a search over that goal's recorded
conversations. Use ``list_sessions`` for intentional discovery.
"""

return [
candidate
Expand All @@ -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.<goal_id>`` 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")
Expand Down
2 changes: 1 addition & 1 deletion loopx/semantics/project_registry_io_manifest_v1.json
Original file line number Diff line number Diff line change
Expand Up @@ -247,7 +247,7 @@
},
{
"site": "loopx/capabilities/manager_context/roundtrip.py::<module>.drain::codec_read:load_project_registry#1",
"line": 879,
"line": 1010,
"column": 22,
"kind": "codec_read",
"api": "load_project_registry",
Expand Down
Loading
Loading