fix(compression): rotation path clones the concurrent tail into the child
CI caught the sibling site the in-place fix missed: legacy (non-in-place) compression rotates via publish_compression_child, where a mid-summary append previously stranded in the closed parent. Same watermark + pure-SQL column clone as archive_and_compact, with session_id rewritten to the child. Lineage-guard test flipped to pin the appends-flow-freely contract; rotation watermark tests added (tail follows the child; None = historical behavior).
This commit is contained in:
@@ -3524,6 +3524,7 @@ def compress_context(
|
||||
profile_name=_profile_for_child,
|
||||
compression_lock_holder=_lock_holder,
|
||||
require_compression_lease=_lock_holder is not None,
|
||||
watermark=_commit_watermark,
|
||||
)
|
||||
agent.session_id = new_session_id
|
||||
try:
|
||||
|
||||
@@ -5590,12 +5590,21 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
|
||||
profile_name: str = None,
|
||||
compression_lock_holder: str = None,
|
||||
require_compression_lease: bool = True,
|
||||
watermark: Optional[int] = None,
|
||||
) -> None:
|
||||
"""Atomically close a parent and publish its durable compression child.
|
||||
|
||||
The parent closure, child row, and compacted handoff become visible in
|
||||
one transaction. Readers can therefore observe either the live parent or
|
||||
a complete child, never an ended parent with a missing/empty child.
|
||||
|
||||
Concurrent-append safety (#75316): when *watermark* is provided (the
|
||||
parent's :meth:`get_active_message_watermark` captured at compression
|
||||
start), parent rows that arrived during the slow summary call
|
||||
(``id > watermark``) are cloned into the child AFTER the handoff —
|
||||
same pure-SQL column clone as :meth:`archive_and_compact`, with the
|
||||
session id rewritten — so a mid-compression append survives rotation
|
||||
instead of stranding in the closed parent.
|
||||
"""
|
||||
def _do(conn):
|
||||
lock_row = conn.execute(
|
||||
@@ -5663,6 +5672,39 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
|
||||
total_messages, total_tool_calls = self._insert_message_rows(
|
||||
conn, child_session_id, messages
|
||||
)
|
||||
if watermark is not None:
|
||||
# Clone the parent's concurrent tail (rows landed after the
|
||||
# watermark) into the child, after the handoff. Column-exact
|
||||
# except id/session_id; originals stay in the (closed) parent
|
||||
# for lineage recovery.
|
||||
tail_rows = conn.execute(
|
||||
"SELECT id, tool_calls FROM messages "
|
||||
"WHERE session_id = ? AND active = 1 AND id > ? ORDER BY id",
|
||||
(parent_session_id, int(watermark)),
|
||||
).fetchall()
|
||||
if tail_rows:
|
||||
tail_ids = [int(r["id"]) for r in tail_rows]
|
||||
placeholders = ",".join("?" for _ in tail_ids)
|
||||
clone_cols = [
|
||||
c for c in self._message_column_names(conn)
|
||||
if c not in ("id", "session_id", "active", "compacted")
|
||||
]
|
||||
col_list = ", ".join(clone_cols)
|
||||
conn.execute(
|
||||
f"INSERT INTO messages ({col_list}, session_id, active, compacted) "
|
||||
f"SELECT {col_list}, ?, 1, 0 FROM messages "
|
||||
f"WHERE id IN ({placeholders}) ORDER BY id",
|
||||
[child_session_id, *tail_ids],
|
||||
)
|
||||
total_messages += len(tail_ids)
|
||||
for r in tail_rows:
|
||||
raw = r["tool_calls"]
|
||||
if raw:
|
||||
try:
|
||||
parsed = json.loads(raw) if isinstance(raw, str) else raw
|
||||
total_tool_calls += len(parsed) if isinstance(parsed, list) else 0
|
||||
except (TypeError, ValueError):
|
||||
pass
|
||||
conn.execute(
|
||||
"UPDATE sessions SET message_count = ?, tool_call_count = ? WHERE id = ?",
|
||||
(total_messages, total_tool_calls, child_session_id),
|
||||
|
||||
@@ -305,16 +305,21 @@ def test_publish_compression_child_rejects_lost_or_expired_lease(db: SessionDB)
|
||||
def test_compression_lease_blocks_non_owner_but_allows_owner_flush(
|
||||
db: SessionDB,
|
||||
) -> None:
|
||||
"""Contract flipped by the watermark commit (#75316): a live lease no
|
||||
longer fences ordinary appends — both the owner's flush and a concurrent
|
||||
turn land immediately, and the commit-side watermark decides what
|
||||
survives compaction (see test_compression_watermark_commit.py)."""
|
||||
db.create_session("leased", source="webui")
|
||||
assert db.try_acquire_compression_lock("leased", "winner", ttl_seconds=60)
|
||||
|
||||
with pytest.raises(RuntimeError, match="being compressed"):
|
||||
db.append_message("leased", "user", "late stale turn")
|
||||
|
||||
db.append_message("leased", "user", "late concurrent turn")
|
||||
db.append_message(
|
||||
"leased",
|
||||
"assistant",
|
||||
"winner flush",
|
||||
compression_lock_holder="winner",
|
||||
)
|
||||
assert [m["content"] for m in db.get_messages("leased")] == ["winner flush"]
|
||||
assert [m["content"] for m in db.get_messages("leased")] == [
|
||||
"late concurrent turn",
|
||||
"winner flush",
|
||||
]
|
||||
|
||||
@@ -226,3 +226,54 @@ class TestConcurrentAppendDuringCompaction:
|
||||
assert not append_err, f"append died during commit race: {append_err}"
|
||||
contents = [r["content"] for r in db.get_messages("sess1")]
|
||||
assert "racer" in contents, "racing append was lost"
|
||||
|
||||
|
||||
class TestRotationPathWatermark:
|
||||
"""Legacy (non-in-place) compression rotates to a child session —
|
||||
the concurrent tail must follow the rotation instead of stranding in
|
||||
the closed parent."""
|
||||
|
||||
def test_tail_clones_into_the_child(self, db: SessionDB) -> None:
|
||||
_seed(db)
|
||||
watermark = db.get_active_message_watermark("sess1")
|
||||
assert db.try_acquire_compression_lock("sess1", "rotator") is True
|
||||
db.append_message("sess1", role="user", content="mid-rotation steer")
|
||||
|
||||
db.publish_compression_child(
|
||||
parent_session_id="sess1",
|
||||
child_session_id="child1",
|
||||
source="test",
|
||||
messages=SUMMARY,
|
||||
compression_lock_holder="rotator",
|
||||
require_compression_lease=True,
|
||||
watermark=watermark,
|
||||
)
|
||||
|
||||
child = db.get_messages_as_conversation("child1")
|
||||
assert [m["content"] for m in child] == [
|
||||
SUMMARY[0]["content"],
|
||||
SUMMARY[1]["content"],
|
||||
"mid-rotation steer",
|
||||
]
|
||||
info = db.get_session("child1")
|
||||
assert info["message_count"] == 3
|
||||
# Parent keeps its copy for lineage recovery; parent is closed.
|
||||
parent_info = db.get_session("sess1")
|
||||
assert parent_info["end_reason"] == "compression"
|
||||
|
||||
def test_no_watermark_keeps_historical_rotation(self, db: SessionDB) -> None:
|
||||
_seed(db)
|
||||
assert db.try_acquire_compression_lock("sess1", "rotator") is True
|
||||
db.append_message("sess1", role="user", content="stranded either way")
|
||||
db.publish_compression_child(
|
||||
parent_session_id="sess1",
|
||||
child_session_id="child1",
|
||||
source="test",
|
||||
messages=SUMMARY,
|
||||
compression_lock_holder="rotator",
|
||||
require_compression_lease=True,
|
||||
)
|
||||
child = db.get_messages_as_conversation("child1")
|
||||
assert [m["content"] for m in child] == [
|
||||
SUMMARY[0]["content"], SUMMARY[1]["content"],
|
||||
]
|
||||
|
||||
Reference in New Issue
Block a user