diff --git a/CHANGELOG.md b/CHANGELOG.md index 534142f5a..7c7be9b02 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,8 @@ metadata and the backend fallback mirror it. ### Fixed +- Completed call history remains terminal when an earlier callback saves late (#2640) — thanks @rudycelekli! + - MCP speech tools wait through model loading and progress-extended CPU renders instead of timing out before the backend (#2609) ## [0.5.7] — 2026-10-05 diff --git a/backend/services/telephony/calls.py b/backend/services/telephony/calls.py index 2b94a3dd6..21e591e9d 100644 --- a/backend/services/telephony/calls.py +++ b/backend/services/telephony/calls.py @@ -259,11 +259,13 @@ def store_save(session: CallSession) -> None: session.recording_path, session.error, session.created_at, session.started_at, session.ended_at, session.duration_s if session.finalized else None, ) - try: - with db_conn() as conn: - conn.execute(_UPSERT_SQL, values) - except Exception: # noqa: BLE001 — a DB hiccup must not drop a live call - logger.warning("Could not save call %s", session.id, exc_info=True) + # Keep the snapshot and its commit in one session operation. A delayed + # older save must not overwrite the terminal record from a later one. + try: + with db_conn() as conn: + conn.execute(_UPSERT_SQL, values) + except Exception: # noqa: BLE001 — a DB hiccup must not drop a live call + logger.warning("Could not save call %s", session.id, exc_info=True) def _row_record(row, *, transcript: bool) -> dict: diff --git a/docs/integrations/calls.md b/docs/integrations/calls.md index e593ae868..9439a0717 100644 --- a/docs/integrations/calls.md +++ b/docs/integrations/calls.md @@ -196,3 +196,5 @@ turn the other person talked over. resampled to 8 kHz μ-law and sent as 20 ms frames. On barge-in (about 200 ms of the other person's speech while the agent talks), VoiceStudio sends Twilio a `clear` to drop buffered audio. + +Call-history snapshots commit in session order. A delayed provider callback cannot overwrite the completed record after the call ends; other calls retain their own independent records. diff --git a/tests/test_call_agent.py b/tests/test_call_agent.py index 028666aa6..bc68ddfbd 100644 --- a/tests/test_call_agent.py +++ b/tests/test_call_agent.py @@ -951,3 +951,120 @@ def test_migration_adds_call_sessions_and_matches_the_base_schema(tmp_path, monk fresh = [(r[1], r[2].upper(), r[3], r[5]) for r in canon.execute("PRAGMA table_info(call_sessions)")] norm = lambda cols: [(n, {"FLOAT": "REAL"}.get(t, t), nn, pk) for n, t, nn, pk in cols] # noqa: E731 assert norm(migrated) == norm(fresh) + + +def test_call_store_save_cannot_publish_an_older_snapshot_after_completion(tmp_path, monkeypatch): + """Provider callbacks and finalization must commit in session order.""" + from contextlib import contextmanager + from core import db + from services.telephony import calls + + path = tmp_path / "call-order.db" + with sqlite3.connect(path) as conn: + conn.executescript(db._BASE_SCHEMA) + monkeypatch.setattr(db, "DB_PATH", path) + real_db_conn = db.db_conn + first_entered = threading.Event() + release_first = threading.Event() + finalized = threading.Event() + latest_attempted_lock = threading.Event() + latest_blocked_on_lock = threading.Event() + errors = [] + + @contextmanager + def delayed_first_save(): + # Delay the first writer before opening its real SQLite connection, + # exactly where a blocked thread-pool callback can lag a later save. + if threading.current_thread().name == "old-call-save": + first_entered.set() + if not release_first.wait(10): + raise TimeoutError("first save was not released") + with real_db_conn() as conn: + yield conn + + monkeypatch.setattr(db, "db_conn", delayed_first_save) + session = calls.CallSession(id="call-order", direction="outbound", + remote_number="+15550001111", from_number=FROM, brief="test", profile_id="voice") + + real_session_lock = session._lock + + class ObservedLock: + def __enter__(self): + if threading.current_thread().name == "latest-call-save": + acquired = real_session_lock.acquire(blocking=False) + if not acquired: + latest_blocked_on_lock.set() + latest_attempted_lock.set() + if acquired: + return self + real_session_lock.acquire() + return self + + def __exit__(self, *exc): + real_session_lock.release() + + session._lock = ObservedLock() + + def save_old(): + try: + calls.store_save(session) + except BaseException as exc: + errors.append(exc) + + def complete(): + try: + calls.finish_unconnected(session, "completed") + except BaseException as exc: + errors.append(exc) + finally: + finalized.set() + + old = threading.Thread(target=save_old, name="old-call-save") + latest = threading.Thread(target=complete, name="latest-call-save") + old.start() + try: + assert first_entered.wait(5) + latest.start() + # Confirm the latest writer actually contests the lock while the old + # SQLite save is held. A delayed thread must not miss the interleaving. + assert latest_attempted_lock.wait(5) + blocked_before_release = latest_blocked_on_lock.is_set() + if blocked_before_release: + assert not finalized.is_set() + else: + # On the unfixed source, prove the later snapshot commits before + # releasing the old writer so the lost terminal row is reproducible. + assert finalized.wait(5) + finally: + release_first.set() + old.join(5) + if latest.ident is not None: + latest.join(5) + assert not old.is_alive() and not latest.is_alive() + assert not errors + assert session.finalized is True + with real_db_conn() as conn: + row = conn.execute("SELECT status, outcome, ended_at FROM call_sessions WHERE id=?", + (session.id,)).fetchone() + assert row["status"] == "completed" + assert row["outcome"] == "not_done" + assert row["ended_at"] is not None + assert blocked_before_release + + +def test_call_store_save_keeps_independent_session_records(tmp_path, monkeypatch): + from core import db + from services.telephony import calls + + path = tmp_path / "call-independent.db" + with sqlite3.connect(path) as conn: + conn.executescript(db._BASE_SCHEMA) + monkeypatch.setattr(db, "DB_PATH", path) + for call_id in ("one", "two"): + session = calls.CallSession(id=call_id, direction="outbound", + remote_number="+15550001111", from_number=FROM, brief="test", profile_id="voice") + calls.store_save(session) + calls.finish_unconnected(session, "completed") + with db.db_conn() as conn: + rows = conn.execute("SELECT id, status FROM call_sessions ORDER BY id").fetchall() + assert [tuple(row) for row in rows] == [("one", "completed"), ("two", "completed")]