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