From db96c8ced7da9a74f4e0258adb51a4c1361958d2 Mon Sep 17 00:00:00 2001 From: FalconOrtiz Date: Mon, 7 Sep 2026 16:41:12 -0700 Subject: [PATCH] fix: admit Bot Chat deliveries through the live session owner Adapt FalconOrtiz's owner-mailbox proposal from #101564 onto the current notification poller and topical modules. The durable mailbox is cross-process ingress only: the existing owner admits its normal prompt turn after the current turn and human FIFO clear. Retain immutable receipts, stable admission identities, capability and lease fencing, compression lineage, and disable blind recovery replay of imported turns. The original stale server hunks and expiring receipt protocol were rebuilt rather than cherry-picked: the current facade decomposition and durable busy admission contract differ. Credit the earlier owner-mailbox work in #100544 and durable producer work in #100319. Co-authored-by: fangliquanflq Co-authored-by: 686f6c61 --- contributors/emails/falcon.ortiz11@gmail.com | 2 + tests/tools/test_bot_live_owner_delivery.py | 95 ++++++++ tests/tui_gateway/test_auto_continue.py | 4 +- .../test_bot_live_owner_delivery.py | 89 +++++++ tools/bot_live_delivery.py | 223 ++++++++++++++++++ tui_gateway/prompt_turn.py | 12 +- tui_gateway/session_auto_continue.py | 2 + tui_gateway/session_lifecycle.py | 5 +- tui_gateway/session_notifications.py | 55 +++++ tui_gateway/turn_marker.py | 9 +- 10 files changed, 486 insertions(+), 10 deletions(-) create mode 100644 contributors/emails/falcon.ortiz11@gmail.com create mode 100644 tests/tools/test_bot_live_owner_delivery.py create mode 100644 tests/tui_gateway/test_bot_live_owner_delivery.py create mode 100644 tools/bot_live_delivery.py diff --git a/contributors/emails/falcon.ortiz11@gmail.com b/contributors/emails/falcon.ortiz11@gmail.com new file mode 100644 index 0000000000..caffee26b4 --- /dev/null +++ b/contributors/emails/falcon.ortiz11@gmail.com @@ -0,0 +1,2 @@ +FalconOrtiz +# PR #101564 live-owner ingress salvage diff --git a/tests/tools/test_bot_live_owner_delivery.py b/tests/tools/test_bot_live_owner_delivery.py new file mode 100644 index 0000000000..eacecbc72f --- /dev/null +++ b/tests/tools/test_bot_live_owner_delivery.py @@ -0,0 +1,95 @@ +"""Durable mailbox invariants, using real disk and exec boundaries.""" +import json +import os +import subprocess +import sys + +import pytest + + +@pytest.mark.parametrize("terminal_status", ["settled", "failed", "cancelled"]) +def test_delivery_is_idempotent_fenced_and_permanent(tmp_path, terminal_status): + from tools import bot_live_delivery as mailbox + + owner = dict(profile_home=str(tmp_path.resolve()), session_id="chat", + lease_id="lease", live_session_id="live") + delivery_id = "a" * 32 + queued = mailbox.deliver_to_live_owner(tmp_path, owner, "hello", delivery_id=delivery_id) + assert queued["status"] == "queued" + assert mailbox.deliver_to_live_owner(tmp_path, owner, "hello", delivery_id=delivery_id) == queued + with pytest.raises(ValueError): + mailbox.deliver_to_live_owner(tmp_path, owner, "different", delivery_id=delivery_id) + assert mailbox.claim_pending_delivery(tmp_path, dict(owner, lease_id="other")) is None + assert mailbox.claim_pending_delivery(tmp_path, dict(owner, live_session_id="other")) is None + script = ( + "import json,sys; from tools.bot_live_delivery import claim_pending_delivery; " + "print(json.dumps(claim_pending_delivery(sys.argv[1],json.loads(sys.argv[2]))))" + ) + children = [subprocess.Popen([sys.executable, "-c", script, str(tmp_path), json.dumps(owner)], + stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, + stderr=subprocess.PIPE, text=True) for _ in range(2)] + results = [] + for child in children: + out, err = child.communicate(timeout=30) + assert child.returncode == 0, err + results.append(json.loads(out)) + claims = [r for r in results if r is not None] + assert len(claims) == 1 and claims[0]["message"] == "hello" + assert mailbox.read_delivery_result(tmp_path, delivery_id)["status"] == "claimed" + assert mailbox.claim_pending_delivery(tmp_path, owner) is None + receipt = mailbox.complete_delivery(tmp_path, delivery_id, status=terminal_status, reply="answer") + assert mailbox.read_delivery_result(tmp_path, delivery_id) == receipt + assert mailbox.complete_delivery(tmp_path, delivery_id, status=terminal_status, reply="answer") == receipt + with pytest.raises(ValueError): + mailbox.complete_delivery(tmp_path, delivery_id, status=terminal_status, reply="rewrite") + assert mailbox.deliver_to_live_owner(tmp_path, owner, "hello", delivery_id=delivery_id) == receipt + assert mailbox.claim_pending_delivery(tmp_path, owner) is None + if os.name != "nt": + for path in (tmp_path / "runtime" / mailbox.DELIVERY_DIR_NAME).iterdir(): + assert path.stat().st_mode & 0o077 == 0 + + +def test_fifo_survives_clock_rollback(tmp_path, monkeypatch): + from tools import bot_live_delivery as mailbox + + owner = dict(profile_home=str(tmp_path.resolve()), session_id="chat", + lease_id="lease", live_session_id="live") + for timestamp, message in ((100, "first"), (90, "second")): + monkeypatch.setattr(mailbox.time, "time_ns", lambda: timestamp) + mailbox.deliver_to_live_owner(tmp_path, owner, message) + assert mailbox.claim_pending_delivery(tmp_path, owner)["message"] == "first" + assert mailbox.claim_pending_delivery(tmp_path, owner)["message"] == "second" + + +@pytest.mark.parametrize("capable", [True, False]) +def test_only_canonical_capable_owner_receives_across_compression(tmp_path, capable): + from hermes_state import SessionDB + from hermes_cli.active_sessions import try_acquire_active_session, transfer_active_session + from tools import bot_live_delivery as mailbox + + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="chat", source="cli") + db.set_session_title("chat", "Bot Chat") + meta = dict(live_session_id="live", bot_live_delivery_consumer=capable) + lease, refusal = try_acquire_active_session(session_id="chat", surface="desktop", config={}, + registry_home=tmp_path, metadata=meta) + assert refusal is None + try: + owner = mailbox.find_canonical_live_owner(tmp_path) + if not capable: + assert owner is None + return + assert owner["lease_id"] == lease.lease_id + queued = mailbox.deliver_to_live_owner(tmp_path, owner, "before compression") + db.end_session("chat", "compression") + db.create_session(session_id="tip", source="cli", parent_session_id="chat") + assert transfer_active_session(lease, session_id="tip", metadata=meta) + current = mailbox.find_canonical_live_owner(tmp_path) + assert current["session_id"] == "tip" + claim = mailbox.claim_pending_delivery(tmp_path, current) + assert claim["delivery_id"] == queued["delivery_id"] + assert claim["session_id"] == "chat" + assert mailbox.claim_pending_delivery(tmp_path, current) is None + finally: + lease.release() + db.close() diff --git a/tests/tui_gateway/test_auto_continue.py b/tests/tui_gateway/test_auto_continue.py index 922cd62961..213a925db4 100644 --- a/tests/tui_gateway/test_auto_continue.py +++ b/tests/tui_gateway/test_auto_continue.py @@ -181,12 +181,12 @@ def test_interrupt_racing_marker_write_cannot_leave_recovery_state( session = _session(agent=agent, running=True) _patch_local_interrupt(monkeypatch, session) - def write_after_stop(home, key, prompt, *, attempts=0): + def write_after_stop(home, key, prompt, *, attempts=0, auto_continue=True): response = server._methods["session.interrupt"]( "stop-during-write", {"session_id": "runtime-race"} ) assert response["result"]["status"] == "interrupted" - record_turn_start(home, key, prompt, attempts=attempts) + record_turn_start(home, key, prompt, attempts=attempts, auto_continue=auto_continue) monkeypatch.setattr(server, "record_turn_start", write_after_stop) diff --git a/tests/tui_gateway/test_bot_live_owner_delivery.py b/tests/tui_gateway/test_bot_live_owner_delivery.py new file mode 100644 index 0000000000..9484d49dc5 --- /dev/null +++ b/tests/tui_gateway/test_bot_live_owner_delivery.py @@ -0,0 +1,89 @@ +"""Imported turns retain their receipt and cannot bypass the local FIFO.""" +import threading +from types import SimpleNamespace + +from tui_gateway.method_ctx import rebind +from tui_gateway import session_notifications, session_auto_continue +from tui_gateway.turn_marker import record_turn_start, read_turn_marker + + +def test_refused_input_commits_failed_mailbox_receipt(tmp_path): + import contextvars + import logging + import time + from tui_gateway import prompt_turn + from tools import bot_live_delivery as mailbox + + owner = dict(profile_home=str(tmp_path.resolve()), session_id="chat", + lease_id="lease", live_session_id="live") + queued = mailbox.deliver_to_live_owner(tmp_path, owner, "refused input") + mailbox.claim_pending_delivery(tmp_path, owner) + agent = SimpleNamespace(session_id="chat") + session = dict(agent=agent, session_key="chat", history_lock=threading.RLock(), running=True) + retired = [] + noop = lambda *args, **kwargs: None + submit = rebind(prompt_turn._run_prompt_submit, { + "threading": threading, "time": time, "logger": logging.getLogger(__name__), + "_sessions_lock": threading.RLock(), "_sessions": {}, + "_admit_prompt_turn": lambda *args: ([], agent), + "_emit": noop, "bind_transport": noop, "reset_transport": noop, + "_current_runtime_session_record": contextvars.ContextVar("refused_turn"), + "_TurnRun": prompt_turn._TurnRun, + "_record_turn_marker": lambda *args, **kwargs: "marker", + "_prepare_turn_input": lambda *args: None, + "_finish_turn": noop, "_clear_inflight_turn": noop, + "_retire_turn_marker": lambda *args: retired.append(args), + "_emit_settled_session_info": noop, + }) + def terminal(outcome): + mailbox.complete_delivery(tmp_path, queued["id"], status=outcome["status"], + error=outcome.get("error", "")) + assert submit(None, "live", session, "refused input", terminal_callback=terminal) + session["_run_thread"].join(timeout=5) + assert not session["_run_thread"].is_alive() + assert mailbox.read_delivery_result(tmp_path, queued["id"])["status"] == "failed" + assert retired and session["running"] is False + + +def test_imported_crash_marker_never_autocontinues(tmp_path): + record_turn_start(tmp_path, "chat", "imported", auto_continue=False) + marker = read_turn_marker(tmp_path, "chat") + assert marker["auto_continue"] is False + schedule = rebind(session_auto_continue._maybe_schedule_auto_continue, { + "_session_home": lambda session: tmp_path, + "read_turn_marker": read_turn_marker, + }) + assert schedule("live", {}, "chat") is None + + +def test_local_work_blocks_mailbox_claim_without_consuming_envelope(monkeypatch, tmp_path): + import tools.bot_live_delivery as mailbox + owner = {"lease_id": "lease", "live_session_id": "live", "session_id": "chat"} + pending = [{"id": "receipt", "message": "imported"}] + monkeypatch.setattr(mailbox, "find_canonical_live_owner", lambda home: owner) + monkeypatch.setattr(mailbox, "claim_pending_delivery", lambda home, pinned: pending.pop(0)) + receipts = [] + monkeypatch.setattr(mailbox, "complete_delivery", lambda *args, **kwargs: receipts.append((args, kwargs))) + submitted = [] + def submit(rid, sid, session, text, **kwargs): + submitted.append(text) + kwargs["terminal_callback"]({"status": "settled", "text": "reply"}) + return True + poll = rebind(session_notifications._poll_bot_live_delivery_once, { + "_session_home": lambda session: tmp_path, + "_run_prompt_submit": submit, + "_notif_release_turn": lambda session: session.update(running=False), + }) + session = {"history_lock": threading.RLock(), "agent": object(), "session_key": "chat", + "active_session_lease": SimpleNamespace(lease_id="lease", released=False)} + for blocker in ("running", "queued_prompt", "queued_prompts", "_auto_continue_scheduled"): + session[blocker] = True + assert poll("live", session) is False + assert pending and not submitted + session.pop(blocker) + assert poll("other-live", session) is False + assert pending + assert poll("live", session) is True + assert submitted == ["imported"] and not pending + assert receipts[0][0][1] == "receipt" + assert receipts[0][1]["reply"] == "reply" diff --git a/tools/bot_live_delivery.py b/tools/bot_live_delivery.py new file mode 100644 index 0000000000..7a19537790 --- /dev/null +++ b/tools/bot_live_delivery.py @@ -0,0 +1,223 @@ +"""Durable, at-most-once handoff to an existing Bot Chat owner. + +Adapted from FalconOrtiz's live-owner mailbox (#101564). A single private +record advances queued -> claimed -> terminal under a process-shared lock. +Claims never expire: a crashed consumer leaves an inspectable unknown outcome, +not permission to execute the same input again. Receipts are permanent. +""" +from __future__ import annotations + +import json +import os +import re +import tempfile +import time +import uuid +from contextlib import contextmanager +from pathlib import Path +from typing import Any + +from hermes_cli.active_sessions import _FileLock + +DELIVERY_DIR_NAME = "bot_live_delivery" +_OWNER_KEYS = ("profile_home", "session_id", "lease_id", "live_session_id") +_TERMINAL = frozenset({"settled", "failed", "cancelled", "ambiguous"}) + + +def find_canonical_live_owner(profile_home: Path | str) -> dict[str, Any] | None: + """Resolve exact Bot Chat's compression tip without creating/migrating its DB. + + Capability advertisement is mandatory; old Desktop/TUI processes must not + receive work they cannot consume. Registry errors propagate, failing closed. + """ + from hermes_cli.active_sessions import active_session_registry_snapshot + from hermes_state import SessionDB + + home = Path(profile_home).resolve() + if not (home / "state.db").is_file(): + return None + db = SessionDB(db_path=home / "state.db", read_only=True) + try: + row = db.get_session_by_title("Bot Chat") + session_id = db.get_compression_tip(row["id"]) if row else None + finally: + db.close() + if not session_id: + return None + for entry in active_session_registry_snapshot(registry_home=home): + meta = entry.get("metadata") or {} + if (entry["session_id"] == session_id + and meta.get("bot_live_delivery_consumer") is True + and meta.get("live_session_id")): + return dict(profile_home=str(home), session_id=session_id, + lease_id=entry["lease_id"], live_session_id=meta["live_session_id"]) + return None + + +def _owner(home: Path | str, owner: dict[str, Any]) -> dict[str, str]: + pinned = {key: owner.get(key) for key in _OWNER_KEYS} + if not all(isinstance(value, str) and value for value in pinned.values()): + raise ValueError("owner requires profile_home, session_id, lease_id and live_session_id") + if pinned["profile_home"] != str(Path(home).resolve()): + raise ValueError("owner belongs to a different profile home") + return pinned + + +def _delivery_id(value: str) -> str: + if not isinstance(value, str) or re.fullmatch(r"[0-9a-f]{32,64}", value) is None: + raise ValueError("delivery id must be 32 to 64 lowercase hex characters") + return value + + +def _root(home: Path | str) -> Path: + return Path(home).resolve() / "runtime" / DELIVERY_DIR_NAME + + +def _fsync_dir(path: Path) -> None: + # Windows cannot open directories with os.open; file fsync still applies. + if os.name == "nt": + return + fd = os.open(path, os.O_RDONLY) + try: + os.fsync(fd) + finally: + os.close(fd) + + +@contextmanager +def _locked(home: Path | str): + root = _root(home) + root.parent.mkdir(parents=True, exist_ok=True) + root.mkdir(mode=0o700, exist_ok=True) + root.chmod(0o700) + _fsync_dir(root.parent) + _fsync_dir(root.parent.parent) + lock = root / ".lock" + fd = os.open(lock, os.O_CREAT | os.O_WRONLY, 0o600) + os.close(fd) + with _FileLock(lock): + yield root + + +def _read(path: Path) -> dict[str, Any] | None: + try: + return json.loads(path.read_text(encoding="utf-8")) + except FileNotFoundError: + return None + + +def _write(path: Path, record: dict[str, Any]) -> None: + fd, temporary = tempfile.mkstemp(dir=path.parent, prefix=".delivery-") + try: + with os.fdopen(fd, "w", encoding="utf-8") as stream: + json.dump(record, stream, ensure_ascii=False, sort_keys=True) + stream.flush() + os.fsync(stream.fileno()) + os.replace(temporary, path) + _fsync_dir(path.parent) + finally: + Path(temporary).unlink(missing_ok=True) + + +def deliver_to_live_owner( + profile_home: Path | str, owner: dict[str, Any], message: str, + *, delivery_id: str | None = None, +) -> dict[str, Any]: + """Return durable admission immediately, without waiting for the owner. + + Retry with the same id AND pinned owner/message to inspect the existing + state. Reusing an id with a different payload is an error, never an overwrite. + """ + pinned = _owner(profile_home, owner) + if not isinstance(message, str): + raise ValueError("message must be a string") + key = _delivery_id(delivery_id if delivery_id is not None else uuid.uuid4().hex) + with _locked(profile_home) as root: + path = root / f"{key}.json" + existing = _read(path) + if existing is not None: + if existing["owner"] != pinned or existing["message"] != message: + raise ValueError("delivery id already belongs to a different payload") + return existing + # Wall time can roll back. Permanent receipts retain the admission + # high-water mark, allocated while holding the cross-process lock. + sequence = max((record.get("sequence", record["created_at"]) + for candidate in root.glob("*.json") + if (record := _read(candidate)) is not None), default=0) + 1 + record = dict(delivery_id=key, id=key, owner=pinned, **pinned, + message=message, status="queued", created_at=time.time_ns(), + sequence=sequence) + _write(path, record) + return record + + +def _matches(home: Path | str, record: dict, owner: dict) -> bool: + pinned = record["owner"] + if any(pinned[key] != owner[key] for key in ("profile_home", "lease_id", "live_session_id")): + return False + if pinned["session_id"] == owner["session_id"]: + return True + from hermes_state import SessionDB + + db = SessionDB(db_path=Path(home) / "state.db", read_only=True) + try: + return db.get_compression_tip(pinned["session_id"]) == owner["session_id"] + finally: + db.close() + + +def claim_pending_delivery( + profile_home: Path | str, owner: dict[str, Any], +) -> dict[str, Any] | None: + """Claim oldest matching input exactly once; caller supplies its current lease. + + A lease transfer across compression is accepted only along the original + stored session's compression chain. A new lease/live session cannot steal it. + Caller must hold its normal turn-admission guard before invoking this. + """ + current = _owner(profile_home, owner) + if not _root(profile_home).is_dir(): + return None + with _locked(profile_home) as root: + pending = [] + for path in root.glob("*.json"): + record = _read(path) + if record is not None and record["status"] == "queued" and _matches(profile_home, record, current): + pending.append(record) + if not pending: + return None + record = min(pending, key=lambda item: ( + item.get("sequence", item["created_at"]), item["delivery_id"])) + record.update(status="claimed", claimed_at=time.time_ns()) + _write(root / f"{record['delivery_id']}.json", record) + return record + + +def complete_delivery( + profile_home: Path | str, delivery_id: str, *, status: str, + reply: str = "", error: str = "", reason: str = "", +) -> dict[str, Any]: + """Persist an immutable terminal receipt; duplicate identical completion is safe.""" + key = _delivery_id(delivery_id) + if status not in _TERMINAL: + raise ValueError("invalid terminal delivery status") + outcome = dict(status=status, reply=reply, error=error, reason=reason) + with _locked(profile_home) as root: + path = root / f"{key}.json" + record = _read(path) + if record is None: + raise FileNotFoundError(f"delivery not found: {key}") + if record["status"] in _TERMINAL: + if any(record.get(k) != v for k, v in outcome.items()): + raise ValueError("delivery already has a different terminal receipt") + return record + if record["status"] != "claimed": + raise ValueError("delivery must be claimed before completion") + record.update(outcome, completed_at=time.time_ns()) + _write(path, record) + return record + + +def read_delivery_result(profile_home: Path | str, delivery_id: str) -> dict[str, Any] | None: + """Read admission/claim/terminal state without waiting or deleting its receipt.""" + return _read(_root(profile_home) / f"{_delivery_id(delivery_id)}.json") diff --git a/tui_gateway/prompt_turn.py b/tui_gateway/prompt_turn.py index 0b82e37ef6..987aa66ebd 100644 --- a/tui_gateway/prompt_turn.py +++ b/tui_gateway/prompt_turn.py @@ -115,7 +115,7 @@ def _admit_prompt_turn( return images, agent -def _record_turn_marker(session: dict, text: Any) -> str: +def _record_turn_marker(session: dict, text: Any, *, auto_continue: bool = True) -> str: """Write the durable crash marker; returns the session key it was written under (compression can rotate session_key mid-turn). A surviving marker means the process died mid-turn. The key is published before the disk write so an interrupt racing startup can retire @@ -127,7 +127,8 @@ def _record_turn_marker(session: dict, text: Any) -> str: if isinstance(marker_text, str) and marker_text.strip(): with session["history_lock"]: session["_active_turn_marker_key"] = marker_key - record_turn_start(marker_home, marker_key, marker_text, attempts=marker_attempt) + record_turn_start(marker_home, marker_key, marker_text, attempts=marker_attempt, + auto_continue=auto_continue) with session["history_lock"]: marker_cancelled = bool(session.get("_turn_cancel_requested")) if marker_cancelled: @@ -779,11 +780,16 @@ def _run_prompt_submit( st = _TurnRun( session["agent"], session.pop("one_turn_model_restore", None), terminal_callback, receipt_committed=terminal_callback is None) - st.marker_key = _record_turn_marker(session, text) + st.marker_key = _record_turn_marker(session, text, auto_continue=terminal_callback is None) goal_followup = None try: prepared = _prepare_turn_input(sid, session, st, text, images) if prepared is None: + if st.terminal_callback is not None and not st.receipt_attempted: + st.receipt_attempted = True + st.terminal_callback({ + "status": "failed", "text": "", "error": "Context injection refused."}) + st.receipt_committed = True return prompt, run_message, cols, streamer = prepared _invoke_agent( diff --git a/tui_gateway/session_auto_continue.py b/tui_gateway/session_auto_continue.py index 291d494cd0..f31bdf0feb 100644 --- a/tui_gateway/session_auto_continue.py +++ b/tui_gateway/session_auto_continue.py @@ -65,6 +65,8 @@ def _maybe_schedule_auto_continue(sid: str, session: dict, session_key: str) -> home = _session_home(session) if (marker := read_turn_marker(home, session_key)) is None: return None + if not marker.get("auto_continue", True): + return None # The mailbox owns recovery and receipt identity for imported turns. enabled, freshness_secs, max_attempts = _auto_continue_config() age = time.time() - marker["started_at"] if not enabled or age > freshness_secs or marker["attempts"] >= max_attempts: diff --git a/tui_gateway/session_lifecycle.py b/tui_gateway/session_lifecycle.py index 6800e1cd1a..90bbf31783 100644 --- a/tui_gateway/session_lifecycle.py +++ b/tui_gateway/session_lifecycle.py @@ -31,7 +31,7 @@ def _claim_active_session_slot( from hermes_cli.active_sessions import try_acquire_active_session return try_acquire_active_session( session_id=session_key, surface=surface, config=_load_cfg(), registry_home=profile_home, - metadata={"live_session_id": live_session_id}, + metadata={"live_session_id": live_session_id, "bot_live_delivery_consumer": True}, track_liveness=str(surface or "").strip().lower() == "desktop") except Exception as exc: logger.warning("Failed to claim active session slot: %s", exc) @@ -138,7 +138,8 @@ def _transfer_active_session_slot(sid: str, session: dict, *, new_session_id: st return True try: from hermes_cli.active_sessions import transfer_active_session - if transfer_active_session(lease, session_id=new_session_id, metadata={"live_session_id": sid}): + if transfer_active_session(lease, session_id=new_session_id, metadata={ + "live_session_id": sid, "bot_live_delivery_consumer": True}): return True except Exception: logger.debug("Failed to transfer active session slot", exc_info=True) diff --git a/tui_gateway/session_notifications.py b/tui_gateway/session_notifications.py index 246880691c..8f3bc45e52 100644 --- a/tui_gateway/session_notifications.py +++ b/tui_gateway/session_notifications.py @@ -489,6 +489,57 @@ def _notif_handle_ready(sid, session, events, emitted, registry, fmt, deferred, _notif_dispatch_completions(sid, session, completions, registry, deferred) +def _poll_bot_live_delivery_once(sid: str, session: dict) -> bool: + """Run one durable envelope only after local FIFO/continuations yield the idle boundary.""" + from tools.bot_live_delivery import claim_pending_delivery, complete_delivery, find_canonical_live_owner + + home = _session_home(session) + with session["history_lock"]: + if any(session.get(key) for key in ( + "running", "_closing", "_finalized", "queued_prompt", "queued_prompts", + "_auto_continue_scheduled")) or session.get("agent") is None: + return False + lease = session.get("active_session_lease") + if lease is None or getattr(lease, "released", False): + return False + owner = find_canonical_live_owner(home) + if (not owner or owner.get("lease_id") != lease.lease_id + or owner.get("live_session_id") != sid + or owner.get("session_id") != session.get("session_key")): + return False + # The mailbox matches each envelope to this pinned lease/live id and compression lineage. + claimed = claim_pending_delivery(home, owner) + if claimed is None: + return False + session["running"] = True + + delivery_id = str(claimed["id"]) + + def terminal_receipt(terminal: dict) -> None: + status = str(terminal.get("status") or "failed") + error = str(terminal.get("error") or "") + reason = "cancelled" if status == "cancelled" else "" + if status not in {"settled", "cancelled"}: + from tools.bot_failure_reasons import classify_agent_error + reason = classify_agent_error(error) + # Let a failed write propagate: the turn must not retire its crash marker without its receipt. + complete_delivery(home, delivery_id, status=status, + reply=str(terminal.get("text") or "") if status == "settled" else "", + error=error, reason=reason) + + try: + started = _run_prompt_submit(f"__bot_dm__{delivery_id}", sid, session, claimed["message"], + image_paths=[], terminal_callback=terminal_receipt) + except Exception as exc: + _notif_release_turn(session) + terminal_receipt({"status": "failed", "error": str(exc)}) + raise + if not started: + _notif_release_turn(session) + terminal_receipt({"status": "failed", "error": "live session owner could not start the delivery turn"}) + return started + + def _notification_poller_loop(stop_event: threading.Event, sid: str, session: dict) -> None: """Daemon thread (started by _init_session()) that drains the process-global completion_queue for this session (ownership routing: _notif_handle_event) and polls ``kanban_notify_subs`` every ``_KANBAN_POLL_SECONDS`` — the @@ -507,6 +558,10 @@ def _notification_poller_loop(stop_event: threading.Event, sid: str, session: di last_kanban_poll = last_loop_poll = 0.0 while not stop_event.is_set() and not session.get("_finalized"): now = time.monotonic() + try: + _poll_bot_live_delivery_once(sid, session) + except Exception: + logger.warning("Bot live-owner delivery poll failed", exc_info=True) # /loop and /heartbeat wakeup drivers: fire a due tick for THIS session while idle (same claim-under-lock # as kanban dispatch). An active non-parked /goal owns the idle boundary and defers the loop tick. if now - last_loop_poll >= _LOOP_POLL_SECONDS: diff --git a/tui_gateway/turn_marker.py b/tui_gateway/turn_marker.py index a5cc237412..e17891feed 100644 --- a/tui_gateway/turn_marker.py +++ b/tui_gateway/turn_marker.py @@ -83,13 +83,15 @@ def _update(home: Path | str, session_key: str, mutate, what: str) -> None: logger.debug("failed to %s turn marker for %s", what, session_key, exc_info=True) -def record_turn_start(home: Path | str, session_key: str, prompt: str, *, attempts: int = 0) -> None: +def record_turn_start(home: Path | str, session_key: str, prompt: str, *, attempts: int = 0, + auto_continue: bool = True) -> None: """Persist the marker for a turn that is about to run. ``attempts`` = how many auto-continues led to this run (0 for a user-initiated turn); the crash-loop breaker reads it back on the next resume.""" if not session_key or not prompt: return now = time.time() - entry = {"attempts": max(0, int(attempts)), "prompt": prompt[:_MAX_PROMPT_CHARS], "started_at": now} + entry = {"attempts": max(0, int(attempts)), "prompt": prompt[:_MAX_PROMPT_CHARS], "started_at": now, + "auto_continue": bool(auto_continue)} _update(home, session_key, lambda entries: {**_prune(entries, now), session_key: entry}, "record") @@ -109,6 +111,7 @@ def read_turn_marker(home: Path | str, session_key: str) -> dict[str, Any] | Non prompt = str(entry.get("prompt") or "") if isinstance(entry, dict) else "" if not prompt.strip(): return None - return {"attempts": max(0, int(entry.get("attempts") or 0)), "prompt": prompt, "started_at": _started_at(entry)} + return {"attempts": max(0, int(entry.get("attempts") or 0)), "prompt": prompt, "started_at": _started_at(entry), + "auto_continue": bool(entry.get("auto_continue", True))} except Exception: return None