"""Regression: lazy-retry run-journal recovery across multiple session reads.

The scenario this test pins down:

1. A WebUI live response stream stops mid-turn. On the first sidecar repair
   attempt the run-journal for the dead stream is NOT visible yet (page-cache loss,
   un-fsynced writes, slow network FS, etc.) so
   `_append_journaled_partial_output` returns False.
2. Pre-fix the repair path baked a permanent "no agent output was recovered"
   marker into the session and never looked at the journal again — even
   after the journaled tokens appeared on disk on a later read.
3. With the fix, the repair instead leaves a `_pending_journal_recovery`
   flag on the marker; the next `get_session()` call lazily re-runs the
   recovery, promotes the marker wording, and threads the journaled
   assistant text/tools into the transcript in the correct chronological
   position.
"""
from concurrent.futures import ThreadPoolExecutor
import threading
import time

import pytest

import api.models as models
import api.config as config
import api.profiles as profiles
import api.streaming as streaming  # noqa: F401  imported for fixture parity
from api.models import (
    Session,
    _apply_core_sync_or_error_marker,
    merge_session_messages_append_only,
)
from api.run_journal import append_run_event


# ── Fixtures (shape mirrors test_session_sidecar_repair.py) ────────────────


@pytest.fixture(autouse=True)
def _isolate_session_dir(tmp_path, monkeypatch):
    session_dir = tmp_path / "sessions"
    session_dir.mkdir()
    index_file = session_dir / "_index.json"
    monkeypatch.setattr(models, "SESSION_DIR", session_dir)
    monkeypatch.setattr(models, "SESSION_INDEX_FILE", index_file)
    models.SESSIONS.clear()
    yield session_dir, index_file
    models.SESSIONS.clear()


@pytest.fixture(autouse=True)
def _isolate_stream_state():
    config.STREAMS.clear()
    config.CANCEL_FLAGS.clear()
    config.AGENT_INSTANCES.clear()
    config.STREAM_PARTIAL_TEXT.clear()
    config.ACTIVE_RUNS.clear()
    yield
    config.STREAMS.clear()
    config.CANCEL_FLAGS.clear()
    config.AGENT_INSTANCES.clear()
    config.STREAM_PARTIAL_TEXT.clear()
    config.ACTIVE_RUNS.clear()


@pytest.fixture(autouse=True)
def _isolate_agent_locks():
    config.SESSION_AGENT_LOCKS.clear()
    models._JOURNAL_RETRY_LOCKS.clear()
    yield
    config.SESSION_AGENT_LOCKS.clear()
    models._JOURNAL_RETRY_LOCKS.clear()


@pytest.fixture()
def hermes_home(tmp_path, monkeypatch):
    home = tmp_path / "hermes_home"
    home.mkdir()
    (home / "sessions").mkdir()
    monkeypatch.setenv("HERMES_HOME", str(home))
    monkeypatch.setattr(profiles, "_DEFAULT_HERMES_HOME", home)
    return home


def _make_dead_stream_session(
    session_id: str,
    *,
    stream_id: str,
    existing_msgs_count: int = 96,
    pending_text: str = (
        "[IMPORTANT: Background process polling. "
        "Continue the user's prior request.]"
    ),
):
    """Build a session that mirrors the production bug: lots of prior history,
    pending_user_message set, an active_stream_id pointing at a dead stream,
    and pending_started_at populated."""
    messages = []
    for i in range(existing_msgs_count // 2):
        messages.append({"role": "user", "content": f"q{i}"})
        messages.append({"role": "assistant", "content": f"a{i}"})
    s = Session(session_id=session_id, title="Lost-response repro", messages=messages)
    s.pending_user_message = pending_text
    s.pending_started_at = 1779237637  # production-shaped value
    s.active_stream_id = stream_id
    return s


def _make_pending_retry_session(session_id: str, *, stream_id: str):
    s = Session(session_id=session_id, title="Pending retry", messages=[
        {"role": "user", "content": "q", "timestamp": 1},
        {
            "role": "assistant",
            "content": models._INTERRUPTED_PENDING_RETRY_WORDING,
            "timestamp": 2,
            "_error": True,
            "type": "interrupted",
            "_pending_journal_recovery": True,
            "_journal_retry_stream_id": stream_id,
            "_journal_retry_attempts": 0,
            "_journal_retry_first_seen_ts": int(models.time.time()),
        },
    ])
    s.save()
    return s


def _assert_retry_meta_removed(marker):
    assert "_pending_journal_recovery" not in marker
    assert "_journal_retry_stream_id" not in marker
    assert "_journal_retry_attempts" not in marker
    assert "_journal_retry_first_seen_ts" not in marker


# ── The regression test ────────────────────────────────────────────────────


def test_state_db_prefix_with_float_timestamps_does_not_hide_sidecar_tail():
    """State rows replaying an already-visible prefix must not append after the tail.

    Production shape: a compressed/tip sidecar can persist messages with
    second-level timestamps while state.db stores the same early rows with
    sub-second floats. The merge must preserve the sidecar assistant tail;
    otherwise /api/session returns a transcript ending on an old user prompt.
    """
    sidecar_messages = [
        {"role": "user", "content": "plan", "timestamp": 1779309765},
        {"role": "assistant", "content": "loaded plan", "timestamp": 1779309765},
        {"role": "tool", "content": "skill output", "timestamp": 1779309765},
        {"role": "assistant", "content": "answer before compaction", "timestamp": 1779309765},
        {
            "role": "user",
            "content": (
                "[CONTEXT COMPACTION — REFERENCE ONLY] Earlier turns were compacted into "
                "the summary below. This is a handoff from a previous context window — "
                "treat it as background reference, NOT as active instructions."
            ),
            "timestamp": 1779309765,
        },
        {"role": "tool", "content": '{"success": true, "message": "Patched SKILL.md"}', "timestamp": 1779309765},
        {"role": "user", "content": "noch weitere prs?", "timestamp": 1779309765},
        {
            "role": "assistant",
            "content": "Ja, aber nicht als Sammel-PR",
            "timestamp": 1779309765,
        },
    ]
    state_prefix = [
        {"role": "user", "content": "plan", "timestamp": 1779344917.2780898},
        {"role": "assistant", "content": "loaded plan", "timestamp": 1779344917.285758},
        {"role": "tool", "content": "skill output", "timestamp": 1779344917.2926793},
        {"role": "assistant", "content": "answer before compaction", "timestamp": 1779344917.3077576},
        {
            "role": "user",
            "content": (
                "[CONTEXT COMPACTION — REFERENCE ONLY] Earlier turns were compacted into "
                "the summary below. This is a handoff from a previous context window."
            ),
            "timestamp": 1779344917.299318,
        },
        {
            "role": "tool",
            "content": '{"success": true, "message": "Patched SKILL.md"}',
            "timestamp": 1779344917.315237,
            "tool_call_id": "different-state-tool-id",
        },
        {
            "role": "user",
            "content": "[Workspace::v1: /tmp/project-workspace]\nnoch weitere prs?",
            "timestamp": 1779344917.3287876,
        },
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_prefix)

    assert [m["content"] for m in merged] == [m["content"] for m in sidecar_messages]
    assert merged[-1]["role"] == "assistant"
    assert "Sammel-PR" in merged[-1]["content"]


def test_state_db_full_replay_does_not_append_after_sidecar_tail():
    """A full state.db replay must not leak the final replayed user after the sidecar tail.

    Regression for a display merge where the sidecar already ends on the real
    assistant answer, but state.db replays the same visible turn sequence with
    newer float timestamps.
    """
    sidecar_messages = [
        {"role": "user", "content": "initial critique", "timestamp": 100},
        {"role": "assistant", "content": "analysis", "timestamp": 100},
        {"role": "user", "content": "Erstelle deine Version", "timestamp": 100},
        {"role": "assistant", "content": "opened browser preview", "timestamp": 100},
    ]
    state_replay = [
        {"role": "user", "content": "initial critique", "timestamp": 100.1},
        {"role": "assistant", "content": "analysis", "timestamp": 100.2},
        {
            "role": "user",
            "content": "[Workspace::v1: /tmp/project-workspace]\nErstelle deine Version",
            "timestamp": 100.3,
        },
        {"role": "assistant", "content": "opened browser preview", "timestamp": 100.4},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_replay)

    assert [m["content"] for m in merged] == [m["content"] for m in sidecar_messages]
    assert merged[-1]["role"] == "assistant"
    assert merged[-1]["content"] == "opened browser preview"


def test_state_db_only_user_prompt_before_sidecar_tail_is_reinserted_chronologically():
    """A missing state.db-only user prompt must not be hidden by a newer sidecar tail.

    Interruption/restart recovery can persist a newer assistant/error tail in
    the WebUI sidecar while the user prompt that triggered it only exists in
    state.db. The append-only reconciliation should preserve that state-backed
    user turn and place it before the newer sidecar tail instead of treating
    every state row older than the sidecar tail as replay.
    """
    sidecar_messages = [
        {"role": "user", "content": "old user", "timestamp": 100.0},
        {"role": "assistant", "content": "old assistant", "timestamp": 101.0},
        {"role": "assistant", "content": "newer recovery tail", "timestamp": 103.0},
    ]
    state_messages = [
        {"role": "user", "content": "old user", "timestamp": 100.0},
        {"role": "assistant", "content": "old assistant", "timestamp": 101.0},
        {"role": "user", "content": "missing prompt", "timestamp": 102.0},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)

    assert [m["content"] for m in merged] == [
        "old user",
        "old assistant",
        "missing prompt",
        "newer recovery tail",
    ]


def test_state_db_only_user_prompt_equal_timestamp_stays_before_assistant_tail():
    """Same-timestamp rescued user prompts must not render after the answer."""
    sidecar_messages = [
        {"role": "assistant", "content": "answer", "timestamp": 100.0},
    ]
    state_messages = [
        {"role": "user", "content": "the question", "timestamp": 100.0},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)

    assert [m["content"] for m in merged] == ["the question", "answer"]


def test_state_db_rescue_same_second_multiturn_keeps_order_no_user_assistant_inversion():
    """Regression (#4216): when an EARLIER user/assistant turn, a recovery tail,
    and the rescued state-only prompt all share one timestamp, the rescue must
    NOT insert the prompt before the earlier assistant (which would produce
    user, user, assistant — re-ordering the already-answered turn). The rescued
    prompt belongs after the earlier turn, before the recovery tail."""
    sidecar_messages = [
        {"role": "user", "content": "old user", "timestamp": 100.0},
        {"role": "assistant", "content": "old assistant", "timestamp": 100.0},
        {"role": "user", "content": "recovery tail", "timestamp": 100.0},
    ]
    state_messages = [
        {"role": "user", "content": "missing prompt", "timestamp": 100.0},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)
    contents = [m["content"] for m in merged]

    # The rescued prompt lands after the answered turn, not before its assistant.
    assert contents == ["old user", "old assistant", "missing prompt", "recovery tail"], contents
    # No user message is inverted ahead of an assistant that already answered it:
    # 'old assistant' must come before BOTH later user turns.
    assert contents.index("old assistant") < contents.index("missing prompt")
    assert contents.index("old assistant") < contents.index("recovery tail")


def test_state_db_only_user_prompt_does_not_split_tool_call_and_result_pair():
    """Chronological rescue must preserve assistant tool_call/tool result adjacency."""
    sidecar_messages = [
        {
            "role": "assistant",
            "content": "I will inspect it",
            "timestamp": 100.0,
            "tool_calls": [{"id": "call_1", "type": "function", "function": {"name": "read_file"}}],
        },
        {"role": "tool", "content": "tool result", "timestamp": 101.0, "tool_call_id": "call_1"},
        {"role": "assistant", "content": "final answer", "timestamp": 103.0},
    ]
    state_messages = [
        {"role": "user", "content": "missing prompt", "timestamp": 100.5},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)

    assert [m["content"] for m in merged] == [
        "I will inspect it",
        "tool result",
        "missing prompt",
        "final answer",
    ]


def test_state_db_rescue_equal_timestamp_does_not_split_tool_call_result_pair():
    """Regression (#4216): when the rescued prompt and an assistant(tool_calls)
    -> tool result block share one timestamp, the rescue must NOT land between
    the tool call and its result (which would corrupt tool context / 400). It
    belongs after the complete tool pair."""
    sidecar_messages = [
        {"role": "user", "content": "u", "timestamp": 100.0},
        {
            "role": "assistant",
            "content": "call",
            "timestamp": 100.0,
            "tool_calls": [{"id": "c1", "type": "function", "function": {"name": "f"}}],
        },
        {"role": "tool", "content": "res", "timestamp": 100.0, "tool_call_id": "c1"},
        {"role": "user", "content": "recovery tail", "timestamp": 100.0},
    ]
    state_messages = [
        {"role": "user", "content": "missing", "timestamp": 100.0},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)
    contents = [m["content"] for m in merged]

    assert contents == ["u", "call", "res", "missing", "recovery tail"], contents
    # The assistant tool_call is immediately followed by its tool result.
    roles = [m["role"] for m in merged]
    ai = roles.index("assistant")
    assert roles[ai + 1] == "tool", f"tool pair split: {roles}"


def test_state_db_rescue_equal_timestamp_does_not_split_multi_tool_result_block():
    """Regression (#4216): a multi-tool assistant turn (2+ tool_calls -> 2+
    adjacent tool results) sharing the rescue timestamp must not have the
    rescued prompt land between the tool results — the whole contiguous tool
    block stays intact and the prompt lands after it."""
    sidecar_messages = [
        {"role": "user", "content": "u", "timestamp": 100.0},
        {
            "role": "assistant",
            "content": "call",
            "timestamp": 100.0,
            "tool_calls": [
                {"id": "c1", "type": "function", "function": {"name": "f"}},
                {"id": "c2", "type": "function", "function": {"name": "g"}},
            ],
        },
        {"role": "tool", "content": "r1", "timestamp": 100.0, "tool_call_id": "c1"},
        {"role": "tool", "content": "r2", "timestamp": 100.0, "tool_call_id": "c2"},
        {"role": "user", "content": "recovery", "timestamp": 100.0},
    ]
    state_messages = [
        {"role": "user", "content": "missing", "timestamp": 100.0},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)
    roles = [m["role"] for m in merged]
    contents = [m["content"] for m in merged]

    assert contents == ["u", "call", "r1", "r2", "missing", "recovery"], contents
    # Both tool results stay contiguous immediately after the assistant.
    ai = roles.index("assistant")
    assert roles[ai + 1] == "tool" and roles[ai + 2] == "tool", f"multi-tool block split: {roles}"


def test_state_db_only_user_prompt_before_sidecar_start_is_not_resurrected_without_watermark():
    """No-watermark context/compression reads must not revive old pre-sidecar users."""
    sidecar_messages = [
        {"role": "assistant", "content": "compacted summary", "timestamp": 200.0},
        {"role": "assistant", "content": "current answer", "timestamp": 201.0},
    ]
    state_messages = [
        {"role": "user", "content": "compressed-out old prompt", "timestamp": 100.0},
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_messages)

    assert [m["content"] for m in merged] == [
        "compacted summary",
        "current answer",
    ]


def test_state_db_middle_segment_replay_does_not_append_after_sidecar_tail():
    """A replayed state.db segment from the middle must not be appended after the tail."""
    sidecar_messages = [
        {"role": "user", "content": "older setup", "timestamp": 100},
        {"role": "assistant", "content": "older answer", "timestamp": 100},
        {"role": "assistant", "content": "analysis before request", "timestamp": 100},
        {"role": "user", "content": "Erstelle deine Version", "timestamp": 100},
        {"role": "assistant", "content": "opened browser preview", "timestamp": 100},
    ]
    state_middle_replay = [
        {"role": "assistant", "content": "analysis before request", "timestamp": 100.1},
        {
            "role": "user",
            "content": "[Workspace::v1: /tmp/project-workspace]\nErstelle deine Version",
            "timestamp": 100.2,
        },
    ]

    merged = merge_session_messages_append_only(sidecar_messages, state_middle_replay)

    assert [m["content"] for m in merged] == [m["content"] for m in sidecar_messages]
    assert merged[-1]["role"] == "assistant"
    assert merged[-1]["content"] == "opened browser preview"


def test_interrupted_recovery_markers_do_not_claim_restart_as_fact():
    """A stale live worker is not always a WebUI process restart.

    Broken SSE connections, browser disconnects, lost worker bookkeeping, and
    real restarts all enter the same recovery marker path. User-visible wording
    must describe the generic interruption instead of asserting a process
    restart that systemd evidence may later disprove.
    """
    marker_texts = [
        models._INTERRUPTED_RECOVERED_WORDING,
        models._INTERRUPTED_NO_OUTPUT_WORDING,
        models._INTERRUPTED_PENDING_RETRY_WORDING,
        models._INTERRUPTED_NEUTRAL_WORDING,
    ]

    for text in marker_texts:
        assert "Response interrupted" in text
        assert "process restarted" not in text
        assert "before this turn finished" in text


def test_interrupted_marker_distinguishes_real_process_restart(monkeypatch):
    monkeypatch.setattr(config, "SERVER_START_TIME", 2000.0)
    marker = models._interrupted_recovery_marker(
        recovered_output=False,
        stream_id="stream_crash",
        pending_started_at=1000.0,
    )

    assert marker["interruption_cause"] == "process_restart"
    assert "WebUI process started after this turn began" in marker["content"]
    assert "process restarted" not in marker["content"]


def test_interrupted_marker_distinguishes_stream_run_split_brain(monkeypatch):
    monkeypatch.setattr(config, "SERVER_START_TIME", 1000.0)
    config.ACTIVE_RUNS["stream_split"] = {"session_id": "sid", "phase": "running"}

    marker = models._interrupted_recovery_marker(
        recovered_output=False,
        stream_id="stream_split",
        pending_started_at=2000.0,
    )

    assert marker["interruption_cause"] == "stream_run_split_brain"
    assert "stream was gone but the worker registry still listed the run" in marker["content"]


def test_interrupted_marker_distinguishes_lost_worker_bookkeeping(monkeypatch):
    monkeypatch.setattr(config, "SERVER_START_TIME", 1000.0)

    marker = models._interrupted_recovery_marker(
        recovered_output=False,
        stream_id="stream_lost",
        pending_started_at=2000.0,
    )

    assert marker["interruption_cause"] == "lost_worker_bookkeeping"
    assert "worker bookkeeping no longer had an active run" in marker["content"]


def test_messages_js_names_browser_sse_disconnect_separately():
    repo = models.Path(__file__).parent.parent
    js = (repo / "static" / "messages.js").read_text(encoding="utf-8")

    assert "Connection interrupted" in js
    assert "browser lost the live SSE connection" in js
    assert "Connection lost" not in js


def test_server_treats_broken_pipe_as_client_disconnect_not_500():
    server_py = (models.Path(__file__).parent.parent / "server.py").read_text(encoding="utf-8")

    # server.py now uses the centralized _CLIENT_DISCONNECT_ERRORS tuple from api.helpers
    assert "_CLIENT_DISCONNECT_ERRORS" in server_py
    assert "do not convert it into a misleading server 500" in server_py


def test_lost_response_recovered_on_second_read(hermes_home):
    sid = "9f14583f0e4e4444aaaa111122223333"
    stream_id = "7c8b4108d52b4aba9af362d3a54f47ac"

    # ── Stage 1: simulate page-cache loss — sidecar repair runs while the
    # run-journal for this stream is empty/absent on disk.
    s = _make_dead_stream_session(sid, stream_id=stream_id)
    s.save()
    core_path = hermes_home / "sessions" / f"session_{sid}.json"

    result = _apply_core_sync_or_error_marker(
        s, core_path, stream_id_for_recheck=stream_id,
    )
    assert result is True

    # Marker should carry the lazy-retry flag and *not* the permanent
    # "no agent output was recovered" wording yet.
    last = s.messages[-1]
    assert last.get("_error") is True
    assert last.get("type") == "interrupted"
    assert last.get("_pending_journal_recovery") is True, (
        "First repair pass should defer the recovery decision via a "
        "_pending_journal_recovery flag so a later read can self-heal."
    )
    assert last.get("_journal_retry_stream_id") == stream_id
    assert last.get("_journal_retry_attempts") == 0
    assert isinstance(last.get("_journal_retry_first_seen_ts"), int)
    assert "no agent output was recovered" not in last["content"]
    # pending fields cleared regardless of journal visibility
    assert s.pending_user_message is None
    assert s.active_stream_id is None
    assert s.pending_started_at is None

    # ── Stage 2: the journaled events become visible on disk.
    append_run_event(sid, stream_id, "token", {"text": "Checking GitHub first."})
    append_run_event(
        sid,
        stream_id,
        "tool",
        {
            "name": "terminal",
            "preview": "gh pr list --repo nesquena/hermes-webui",
            "args": {"command": "gh pr list --repo nesquena/hermes-webui"},
        },
    )
    append_run_event(
        sid,
        stream_id,
        "tool_complete",
        {"name": "terminal", "duration": 1.2, "is_error": False},
    )
    append_run_event(
        sid, stream_id, "token", {"text": " The first PR scan completed."},
    )

    # Pin session into the LRU cache and call get_session — this is the
    # production path that triggers lazy retry.
    models.SESSIONS[sid] = s
    reloaded = models.get_session(sid)
    assert reloaded is s

    contents = [m.get("content", "") for m in s.messages]
    # The marker self-healed:
    assert any("recovered from the run journal" in c for c in contents), (
        "After journaled tokens become readable, the marker must promote to "
        "the recovered-output wording."
    )
    assert not any("no agent output was recovered" in c for c in contents)

    # The journaled assistant text and tool card landed BEFORE the marker
    # so chronological order in the transcript is preserved.
    marker_idx = next(
        i for i, m in enumerate(s.messages)
        if m.get("type") == "interrupted" and m.get("_error")
    )
    recovered_msgs = [
        m for m in s.messages[:marker_idx]
        if m.get("_recovered_from_run_journal") is True
    ]
    assert recovered_msgs, "recovered assistant content must sit above the marker"
    recovered_text = " ".join(m.get("content", "") for m in recovered_msgs)
    assert "Checking GitHub first." in recovered_text
    assert "first PR scan completed" in recovered_text

    # Tool card lives in session.tool_calls and points at one of the
    # recovered assistant indices.
    assert s.tool_calls, "journaled tool should be materialized"
    assert s.tool_calls[-1]["name"] == "terminal"
    assert s.tool_calls[-1]["done"] is True

    # Flag and meta cleaned up after promotion.
    promoted = s.messages[marker_idx]
    assert "_pending_journal_recovery" not in promoted
    assert "_journal_retry_stream_id" not in promoted
    assert "_journal_retry_attempts" not in promoted
    assert "_journal_retry_first_seen_ts" not in promoted


def test_concurrent_get_session_serializes_lazy_journal_retry(hermes_home, monkeypatch):
    sid = "retry_lock_sid"
    stream_id = "retry_lock_stream"
    s = _make_pending_retry_session(sid, stream_id=stream_id)
    models.SESSIONS[sid] = s

    entered = threading.Event()
    release = threading.Event()
    counter_lock = threading.Lock()
    calls = 0

    def slow_retry(session, *, preserve_arriving_budget=False):
        nonlocal calls
        with counter_lock:
            calls += 1
        entered.set()
        assert release.wait(timeout=2), "test timed out waiting to release retry body"
        return False

    monkeypatch.setattr(models, "_retry_journal_recovery_in_place", slow_retry)

    with ThreadPoolExecutor(max_workers=2) as executor:
        first = executor.submit(models.get_session, sid)
        assert entered.wait(timeout=2), "first caller did not enter retry helper"
        second = executor.submit(models.get_session, sid)
        assert second.result(timeout=2) is s
        release.set()
        assert first.result(timeout=2) is s

    assert calls == 1


def test_still_arriving_journal_does_not_consume_retry_budget(hermes_home, monkeypatch):
    sid = "retry_arriving_sid"
    stream_id = "retry_arriving_stream"
    s = _make_pending_retry_session(sid, stream_id=stream_id)
    models.SESSIONS[sid] = s

    monkeypatch.setattr(models, "_append_journaled_partial_output", lambda *a, **kw: False)
    monkeypatch.setattr(models, "_journal_is_still_arriving", lambda *a, **kw: True)

    for _ in range(20):
        assert models.get_session(sid) is s

    marker = s.messages[-1]
    assert marker["_journal_retry_attempts"] == 0
    assert marker["_pending_journal_recovery"] is True
    assert marker["content"] == models._INTERRUPTED_PENDING_RETRY_WORDING


def test_sealed_empty_journal_consumes_retry_budget_and_demotes_at_max(hermes_home, monkeypatch):
    sid = "retry_sealed_sid"
    stream_id = "retry_sealed_stream"
    s = _make_pending_retry_session(sid, stream_id=stream_id)
    s.messages[-1]["_journal_retry_attempts"] = models._JOURNAL_RETRY_MAX_ATTEMPTS - 1
    append_run_event(sid, stream_id, "stream_end", {})
    models.SESSIONS[sid] = s

    assert models.get_session(sid) is s

    marker = s.messages[-1]
    assert marker["content"] == models._INTERRUPTED_NEUTRAL_WORDING
    _assert_retry_meta_removed(marker)
    assert not any(m.get("_recovered_from_run_journal") for m in s.messages)


def test_marker_demotes_after_max_attempts_with_sealed_empty_journal(hermes_home, monkeypatch):
    sid = "retry_max_sid"
    stream_id = "retry_max_stream"
    s = _make_pending_retry_session(sid, stream_id=stream_id)
    s.messages[-1]["_journal_retry_attempts"] = models._JOURNAL_RETRY_MAX_ATTEMPTS - 1
    append_run_event(sid, stream_id, "stream_end", {})
    models.SESSIONS[sid] = s

    assert models.get_session(sid) is s

    marker = s.messages[-1]
    assert marker["content"] == models._INTERRUPTED_NEUTRAL_WORDING
    _assert_retry_meta_removed(marker)
    assert not any(m.get("_recovered_from_run_journal") for m in s.messages)


def test_marker_demotes_after_giveup_seconds(hermes_home, monkeypatch):
    base = 1_779_000_000
    monkeypatch.setattr(models.time, "time", lambda: base)
    sid = "retry_age_sid"
    stream_id = "retry_age_stream"
    s = _make_pending_retry_session(sid, stream_id=stream_id)
    s.messages[-1]["_journal_retry_first_seen_ts"] = (
        base - models._JOURNAL_RETRY_GIVEUP_SECONDS - 1
    )
    models.SESSIONS[sid] = s

    append_calls = 0

    def append_should_not_run(*args, **kwargs):
        nonlocal append_calls
        append_calls += 1
        return False

    monkeypatch.setattr(models, "_append_journaled_partial_output", append_should_not_run)

    assert models.get_session(sid) is s

    marker = s.messages[-1]
    assert marker["content"] == models._INTERRUPTED_NEUTRAL_WORDING
    _assert_retry_meta_removed(marker)
    assert append_calls == 0


def test_repair_stale_pending_skips_pre_compression_snapshot_parent(hermes_home):
    """Archived compression parents must not get synthetic interrupt markers."""
    s = _make_dead_stream_session("compressed_parent", stream_id="dead-stream")
    s.pre_compression_snapshot = True
    original_messages = list(s.messages)

    assert models._repair_stale_pending(s) is False

    assert s.messages == original_messages
    assert s.active_stream_id == "dead-stream"
    assert s.pending_user_message


def test_repair_stale_pending_skips_parent_when_continuation_exists(hermes_home):
    """Compression old→new rotation owns the turn in the child, not the old parent."""
    parent = _make_dead_stream_session("compression_parent", stream_id="rotated-stream")
    child = Session(
        session_id="compression_child",
        title="Continuation",
        parent_session_id="compression_parent",
        # Pin the production regression: older code could accidentally save the
        # child with pre_compression_snapshot=True, but its parent link still
        # proves the parent must not be repaired as a lost standalone turn.
        pre_compression_snapshot=True,
        messages=[
            {"role": "user", "content": "ok, push beide", "timestamp": 10},
            {"role": "assistant", "content": "done", "timestamp": 11},
        ],
    )
    child.save()
    original_messages = list(parent.messages)

    assert models._repair_stale_pending(parent) is False

    assert parent.messages == original_messages
    assert parent.active_stream_id == "rotated-stream"
    assert parent.pending_user_message

def test_get_session_syncs_sidecar_from_newer_state_db_even_when_stream_not_terminal(monkeypatch):
    """Refresh should not lose text that already reached state.db.

    Repro shape from a source-WebUI run: sidecar JSON still has
    active_stream_id/pending_user_message and ends at the previous completed
    turn, while the underlying Hermes agent has already written the current
    user/assistant/tool rows to state.db. The run journal has no terminal done,
    so stale-pending repair deliberately does not fire; read-side reconciliation
    must still sync the sidecar from state.db.
    """
    sid = "state_newer_sid"
    stream_id = "stream_live_no_done"
    s = Session(
        session_id=sid,
        title="State newer",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        context_messages=[
            {"role": "user", "content": "old question"},
            {"role": "assistant", "content": "old answer"},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        pending_started_at=102.0,
    )
    s.save()
    models.SESSIONS.pop(sid, None)

    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": 102.0},
        {"role": "assistant", "content": "visible text that disappeared", "timestamp": 103.0},
        {"role": "tool", "content": "tool result", "timestamp": 104.0},
        {"role": "assistant", "content": "latest live progress", "timestamp": 105.0},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": 105.0},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    loaded = models.get_session(sid)

    assert loaded.active_stream_id is None
    assert loaded.pending_user_message is None
    assert [m["content"] for m in loaded.messages] == [m["content"] for m in state_messages]
    assert loaded.context_messages[-1]["content"] == "latest live progress"

    reloaded = Session.load(sid)
    assert reloaded.active_stream_id is None
    assert reloaded.pending_user_message is None
    assert reloaded.messages[-1]["content"] == "latest live progress"


def test_get_session_sync_keeps_pending_image_prompt_when_agent_row_is_proven(monkeypatch, hermes_home):
    sid = "state_newer_pending_image_sid"
    stream_id = "dead_pending_image_stream"
    timestamp = time.time() - models._REPAIR_STALE_PENDING_GRACE_SECONDS - 5
    prompt = "Describe the attached image"
    attachment = {"name": "sample.png", "mime": "image/png", "is_image": True}
    agent_only_note = "Agent-only recall note"
    agent_user_content = [
        {"type": "text", "text": prompt},
        {"type": "image_url", "image_url": {"url": "data:image/png;base64,YWJj"}},
        {"type": "text", "text": agent_only_note},
    ]
    previous_messages = [
        {"role": "user", "content": "old question", "timestamp": timestamp - 2},
        {"role": "assistant", "content": "old answer", "timestamp": timestamp - 1},
    ]
    state_messages = [
        *previous_messages,
        {"role": "user", "content": agent_user_content, "timestamp": timestamp},
        {"role": "assistant", "content": "recovered assistant tail", "timestamp": timestamp + 1},
    ]
    s = Session(
        session_id=sid,
        title="Pending image recovery",
        messages=previous_messages,
        context_messages=previous_messages,
        active_stream_id=stream_id,
        pending_user_message=prompt,
        pending_attachments=[attachment],
        pending_started_at=timestamp,
        pending_user_source="webui",
        _webui_pending_user_timestamp_identity=(stream_id, timestamp),
    )
    s.save()
    models.SESSIONS.pop(sid, None)

    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {
            "message_count": len(state_messages),
            "last_message_at": timestamp + 1,
        },
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    models.get_session(sid)
    reloaded = Session.load(sid)

    submitted_rows = [
        message for message in reloaded.messages
        if message.get("role") == "user" and message.get("content") == prompt
    ]
    assert len(submitted_rows) == 1
    assert submitted_rows[0]["timestamp"] == timestamp
    assert submitted_rows[0]["attachments"] == [attachment]
    assert reloaded.messages[-1] == state_messages[-1]
    assert agent_only_note not in str(reloaded.messages)

    context_rows = models.reconciled_state_db_messages_for_session(
        reloaded,
        prefer_context=True,
        state_messages=state_messages,
    )
    enriched_rows = [
        message for message in context_rows
        if message.get("role") == "user" and message.get("timestamp") == timestamp
    ]
    assert len(enriched_rows) == 1
    assert enriched_rows[0]["content"] == agent_user_content
    assert agent_only_note in str(enriched_rows[0]["content"])


def test_state_db_pending_image_row_alone_does_not_skip_stale_journal_recovery(
    monkeypatch, hermes_home,
):
    sid = "state_pending_image_without_output_sid"
    stream_id = "dead_pending_image_without_output_stream"
    timestamp = time.time() - models._REPAIR_STALE_PENDING_GRACE_SECONDS - 5
    prompt = "Describe the attached image"
    attachment = {"name": "sample.png", "mime": "image/png", "is_image": True}
    old_messages = [
        {"role": "user", "content": "old question", "timestamp": timestamp - 2},
        {"role": "assistant", "content": "old answer", "timestamp": timestamp - 1},
    ]
    state_messages = [
        *old_messages,
        {
            "role": "user",
            "content": [
                {"type": "text", "text": prompt},
                {"type": "image_url", "image_url": {"url": "data:image/png;base64,YWJj"}},
            ],
            "timestamp": timestamp,
        },
    ]
    session = Session(
        session_id=sid,
        title="Pending image without output",
        messages=old_messages,
        context_messages=old_messages,
        active_stream_id=stream_id,
        pending_user_message=prompt,
        pending_attachments=[attachment],
        pending_started_at=timestamp,
        pending_user_source="webui",
        _webui_pending_user_timestamp_identity=(stream_id, timestamp),
    )
    session.save()
    models.SESSIONS[sid] = session

    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {
            "message_count": len(state_messages),
            "last_message_at": timestamp,
        },
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )
    append_run_event(sid, stream_id, "token", {"text": "Recovered partial text."})

    # Cached get_session runs state.db self-heal, then the real stale-stream
    # cleanup must still get the pending turn so it can replay the run journal.
    loaded = models.get_session(sid)
    assert loaded is session
    assert loaded.active_stream_id == stream_id
    assert loaded.pending_user_message == prompt
    assert loaded.pending_attachments == [attachment]

    import api.routes as routes

    assert routes._clear_stale_stream_state(loaded) is True

    reloaded = Session.load(sid)
    submitted_rows = [
        message for message in reloaded.messages
        if message.get("role") == "user" and message.get("content") == prompt
    ]
    assert len(submitted_rows) == 1
    assert submitted_rows[0]["attachments"] == [attachment]
    assert any(
        message.get("_recovered_from_run_journal")
        and "Recovered partial text." in message.get("content", "")
        for message in reloaded.messages
    )
    markers = [
        message for message in reloaded.messages
        if message.get("_error") and message.get("type") == "interrupted"
    ]
    assert len(markers) == 1
    assert "**Response interrupted.**" in markers[0]["content"]
    assert "partial output above was recovered" in markers[0]["content"]


def test_get_session_does_not_sync_while_stream_is_still_live(monkeypatch):
    """A still-running worker owns its own writeback; do not race it.

    If the sidecar's active_stream_id is still present in STREAMS/ACTIVE_RUNS,
    the turn is live and will merge the agent result + clear pending state when
    it ends. Read-side reconciliation must skip so it cannot drop the active
    stream mid-run and make the terminal writeback look stale.
    """
    sid = "state_newer_live_stream_sid"
    stream_id = "stream_still_live"
    s = Session(
        session_id=sid,
        title="Live stream",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        pending_started_at=102.0,
    )
    s.save()
    models.SESSIONS.pop(sid, None)

    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": 102.0},
        {"role": "assistant", "content": "partial live text", "timestamp": 103.0},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": 103.0},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    # Mark the stream as a live in-process worker.
    config.ACTIVE_RUNS[stream_id] = {"session_id": sid, "stream_id": stream_id}
    try:
        loaded = models.get_session(sid)
    finally:
        config.ACTIVE_RUNS.pop(stream_id, None)

    # The live turn keeps owning its writeback: pending state is untouched and
    # the sidecar transcript is not force-synced from the mid-run state.db rows.
    assert loaded.active_stream_id == stream_id
    assert loaded.pending_user_message == "new request"
    assert [m["content"] for m in loaded.messages] == ["old question", "old answer"]


def test_sync_save_failure_does_not_mutate_session(monkeypatch):
    """If the durable snapshot save raises, the live object must stay untouched.

    The sync mutates the shared/cached Session only AFTER persisting a private
    snapshot. A failing save must leave active_stream_id/pending fields intact
    so later reads still see the unfinished stream that is still on disk.
    Tests the helper directly to isolate it from _repair_stale_pending, which
    has its own (separate) dead-stream pending cleanup. (greptile P1)
    """
    sid = "state_sync_save_fail_sid"
    stream_id = "stream_save_fail_no_done"
    attachment = {"name": "sample.png", "mime": "image/png", "is_image": True}
    s = Session(
        session_id=sid,
        title="Save failure",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        pending_attachments=[attachment],
        pending_started_at=102.0,
        pending_user_source="webui",
        _webui_pending_user_timestamp_identity=(stream_id, 102.0),
    )
    s.save()

    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": 102.0},
        {"role": "assistant", "content": "recovered text", "timestamp": 103.0},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": 103.0},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    # Force the snapshot save to fail.
    def boom(self, *args, **kwargs):
        raise OSError("disk full")

    monkeypatch.setattr(Session, "save", boom)

    result = models._sync_sidecar_from_state_db_if_newer(s)

    # Sync reports failure and the live object is NOT mutated.
    assert result is False
    assert s.active_stream_id == stream_id
    assert s.pending_user_message == "new request"
    assert s.pending_attachments == [attachment]
    assert [m["content"] for m in s.messages] == ["old question", "old answer"]
    reloaded = Session.load(sid)
    assert reloaded.active_stream_id == stream_id
    assert reloaded.pending_user_message == "new request"
    assert reloaded.pending_attachments == s.pending_attachments
    assert [m["content"] for m in reloaded.messages] == ["old question", "old answer"]


def test_sync_persists_recovered_state_db_tail_when_stream_dead(monkeypatch):
    """Happy path: dead stream + newer state.db tail is written back durably.

    Mutates the live object only after the snapshot save succeeds, and the
    reconciled transcript reaches both memory and disk.
    """
    sid = "state_sync_ok_sid"
    stream_id = "stream_dead_no_done"
    s = Session(
        session_id=sid,
        title="Recover tail",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        pending_started_at=102.0,
    )
    s.save()

    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": 102.0},
        {"role": "assistant", "content": "recovered tail", "timestamp": 103.0},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": 103.0},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    result = models._sync_sidecar_from_state_db_if_newer(s)

    assert result is True
    assert s.active_stream_id is None
    assert s.pending_user_message is None
    assert [m["content"] for m in s.messages] == [m["content"] for m in state_messages]

    reloaded = Session.load(sid)
    assert reloaded.active_stream_id is None
    assert reloaded.messages[-1]["content"] == "recovered tail"


def test_sync_skips_during_registration_window_recent_pending(monkeypatch):
    """Registration-window race: a just-submitted, not-yet-registered stream
    must NOT be cleared by the self-heal.

    The request handler persists active_stream_id + pending_started_at to the
    sidecar a moment BEFORE the worker thread registers the stream in
    STREAMS/ACTIVE_RUNS. In that window the stream is absent from the registries
    but the turn is legitimately starting. A *recent* pending_started_at must
    keep the self-heal off so it can't drop a turn that is about to run.
    (maintainer-requested regression)
    """
    sid = "state_registration_window_sid"
    stream_id = "stream_just_submitted_not_registered"
    s = Session(
        session_id=sid,
        title="Registration window",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        # Submitted just now: inside the registration grace window.
        pending_started_at=time.time(),
    )
    s.save()

    # state.db is ahead (the new turn's worker has started writing rows), which
    # would otherwise satisfy the "newer transcript" gate.
    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": time.time()},
        {"role": "assistant", "content": "first streamed token", "timestamp": time.time()},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": time.time()},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )
    # The stream is NOT in STREAMS/ACTIVE_RUNS yet (registration not done).
    config.STREAMS.pop(stream_id, None)
    config.ACTIVE_RUNS.pop(stream_id, None)

    result = models._sync_sidecar_from_state_db_if_newer(s)

    # Self-heal must back off: the turn is still starting up.
    assert result is False
    assert s.active_stream_id == stream_id
    assert s.pending_user_message == "new request"
    assert [m["content"] for m in s.messages] == ["old question", "old answer"]

    reloaded = Session.load(sid)
    assert reloaded.active_stream_id == stream_id
    assert reloaded.pending_user_message == "new request"


def test_sync_backs_off_when_session_lock_is_held(monkeypatch):
    """Lock contention: self-heal must not write while another writer holds the
    per-session lock.

    The reconcile + sidecar write happen under the per-session agent lock so a
    concurrent worker/checkpoint save can't be clobbered. If the lock is already
    held (the worker is mid-writeback), the self-heal must back off entirely and
    leave the sidecar untouched rather than racing a full-record overwrite.
    """
    sid = "state_lock_contended_sid"
    stream_id = "stream_dead_but_locked"
    s = Session(
        session_id=sid,
        title="Lock contended",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        pending_started_at=102.0,  # old → past grace window
    )
    s.save()

    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": 102.0},
        {"role": "assistant", "content": "recovered tail", "timestamp": 103.0},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": 103.0},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    # Simulate a concurrent holder of the per-session lock (e.g. the worker's own
    # finalize path) so the non-blocking acquire fails.
    lock = models._get_session_agent_lock(sid)
    assert lock.acquire(blocking=False)
    try:
        result = models._sync_sidecar_from_state_db_if_newer(s)
    finally:
        lock.release()

    # Self-heal backed off; nothing was written and the live object is untouched.
    assert result is False
    assert s.active_stream_id == stream_id
    assert s.pending_user_message == "new request"
    assert [m["content"] for m in s.messages] == ["old question", "old answer"]

    reloaded = Session.load(sid)
    assert reloaded.active_stream_id == stream_id
    assert [m["content"] for m in reloaded.messages] == ["old question", "old answer"]


def test_sync_revalidates_against_concurrent_disk_write_under_lock(monkeypatch):
    """Under-lock reload: a newer concurrent sidecar write is not clobbered.

    The self-heal reloads the session under the lock and writes back the
    reconciled result based on that fresh load, not a snapshot captured before
    the lock. If a concurrent writer rotated the stream / cleared pending on disk
    after our pre-lock read, the under-lock recheck must observe it and back off
    rather than overwriting the whole record with our stale view.
    """
    sid = "state_concurrent_disk_write_sid"
    stream_id = "stream_pre_lock"
    s = Session(
        session_id=sid,
        title="Concurrent disk write",
        messages=[
            {"role": "user", "content": "old question", "timestamp": 100.0},
            {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        ],
        active_stream_id=stream_id,
        pending_user_message="new request",
        pending_started_at=102.0,
    )
    s.save()

    state_messages = [
        {"role": "user", "content": "old question", "timestamp": 100.0},
        {"role": "assistant", "content": "old answer", "timestamp": 101.0},
        {"role": "user", "content": "new request", "timestamp": 102.0},
        {"role": "assistant", "content": "recovered tail", "timestamp": 103.0},
    ]
    monkeypatch.setattr(
        models,
        "get_state_db_session_summary",
        lambda sid_arg, profile=None: {"message_count": len(state_messages), "last_message_at": 103.0},
    )
    monkeypatch.setattr(
        models,
        "get_state_db_session_messages",
        lambda sid_arg, **kwargs: list(state_messages),
    )

    # Between the caller's pre-lock view (s, still pointing at stream_id) and the
    # under-lock reload, a concurrent worker finalized the turn on disk: the
    # sidecar now has a DIFFERENT active_stream_id. The under-lock recheck must
    # detect the mismatch and refuse to clobber.
    disk = Session.load(sid)
    disk.active_stream_id = "rotated_stream_after_compression"
    disk.save(touch_updated_at=False)

    result = models._sync_sidecar_from_state_db_if_newer(s)

    assert result is False
    # On-disk record keeps the concurrent writer's value — not overwritten.
    reloaded = Session.load(sid)
    assert reloaded.active_stream_id == "rotated_stream_after_compression"
