Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 8 additions & 5 deletions plugins-dist/everos-memory/raven_everos/backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -146,9 +146,9 @@ def _jsonify(obj: Any) -> Any:
_STORE_TIMEOUT_S: float = 10.0
# ...and an append is not flat work: EverOS may carve a boundary out of any
# add, which runs a model, so the cost follows how much is handed over. A turn
# passes a handful of messages and lands well inside the floor; a bulk import
# passes up to a hundred at once and did not, which read as a dead service and
# failed every source behind it.
# passes a handful of messages and lands well inside the floor; a writer that
# hands over more and says so (``metadata["bulk"]``) skips this estimate for
# the extraction budget instead.
_STORE_TIMEOUT_PER_MESSAGE_S: float = 0.5

# Shutdown's total budget for flushing every session left with buffered-but-
Expand Down Expand Up @@ -1242,8 +1242,11 @@ async def store(

# A per-turn append must not hold a turn open; a final flush is the call
# that makes EverOS extract, which is what the six-minute budget was
# sized for. One number for both silently overrode the other.
budget = _MEMORIZE_TIMEOUT_S if is_final else _store_budget(len(payload))
# sized for. One number for both silently overrode the other. A bulk
# write is neither a turn nor a flush: nothing waits on it, and EverOS
# extracts on the add itself, so it takes the extraction budget outright.
bulk = bool(metadata and metadata.get("bulk"))
budget = _MEMORIZE_TIMEOUT_S if is_final or bulk else _store_budget(len(payload))
# Marked before the call, not after: if this is cancelled mid-flight
# the add may already have landed, and the safe direction is one
# redundant flush rather than content that is never extracted.
Expand Down
9 changes: 7 additions & 2 deletions raven/importer/orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,11 @@
# message_id from (session_id, timestamp_ms, index-within-batch), so those
# boundaries are part of the id: two messages sharing a millisecond collide,
# and one is dropped, if they land at the same index in different batches.
_BATCH_MSG_LIMIT = 100
# Fifty, not a hundred: EverOS extracts on every add, and that cost is
# superlinear in the message count -- against a real service a 15-message
# batch took 12s and a 52-message batch 24s, while a batch of 100 ran past the
# six-minute extraction budget and failed every memory-file source.
_BATCH_MSG_LIMIT = 50
Comment thread
gloryfromca marked this conversation as resolved.
_BATCH_CHAR_LIMIT = 30_000


Expand Down Expand Up @@ -200,7 +204,8 @@ async def _feed_session(backend: MemoryBackend, session: ImportSession) -> None:

async def _flush(*, is_final: bool) -> None:
nonlocal batch, batch_chars
metadata: dict[str, Any] = {"is_final": is_final}
# bulk: nothing waits on an import write; the backend budgets it as extraction.
metadata: dict[str, Any] = {"is_final": is_final, "bulk": True}
_log_store_request(session.session_id, batch, metadata, batch_chars)
landed = await backend.store(session.session_id, batch, metadata=metadata)
if landed is False:
Expand Down
17 changes: 6 additions & 11 deletions tests/integration/test_import_e2e.py
Original file line number Diff line number Diff line change
Expand Up @@ -237,7 +237,7 @@ async def test_full_pipeline_memory_files(scanner: ClaudeCodeScanner, tmp_path:

@pytest.mark.asyncio
async def test_batching_large_conversation(scanner: ClaudeCodeScanner, tmp_path: Path) -> None:
"""160 messages -> multiple store calls with is_final only on the last."""
"""160 messages -> four store calls of 50, 50, 50 and 10, is_final only on the last, every one bulk."""
results = await scanner.scan()
items = _items_of_kind(scanner, results, SourceKind.CONVERSATION, source_key="sess-large")
assert len(items) == 1
Expand All @@ -248,16 +248,11 @@ async def test_batching_large_conversation(scanner: ClaudeCodeScanner, tmp_path:
summary = await run_import(items, backend, state)

assert summary.submitted == 1
assert len(backend.store_calls) == 2

first, second = backend.store_calls
assert len(first["messages"]) == 100
assert first["metadata"]["is_final"] is False
assert len(second["messages"]) == 60
assert second["metadata"]["is_final"] is True

total_messages = len(first["messages"]) + len(second["messages"])
assert total_messages == 160
sizes = [len(call["messages"]) for call in backend.store_calls]
assert sizes == [50, 50, 50, 10]
assert [call["metadata"]["is_final"] for call in backend.store_calls] == [False, False, False, True]
assert all(call["metadata"]["bulk"] is True for call in backend.store_calls)
assert sum(sizes) == 160


@pytest.mark.asyncio
Expand Down
26 changes: 26 additions & 0 deletions tests/test_everos_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -1646,6 +1646,32 @@ async def _spy(coro, timeout=None):
assert seen == [mod._store_budget(100)]
assert seen[0] > mod._STORE_TIMEOUT_S * 2

async def test_a_bulk_write_gets_the_extraction_budget_however_small(self, monkeypatch) -> None:
"""The importer marks its appends ``bulk``: nothing waits on them, and
EverOS extracts on the add itself, so a per-message estimate is the
wrong shape -- a fifty-message batch measured 24s against a real
service and a hundred ran past six minutes. Only the extraction budget
holds that, and it must not depend on the slice being large.
"""
from raven_everos import backend as mod

seen: list[float] = []

async def _spy(coro, timeout=None):
seen.append(timeout)
return await coro

monkeypatch.setattr(mod.asyncio, "wait_for", _spy)
adapter = MagicMock()
adapter.memorize = AsyncMock(return_value=None)
b = self._backend(adapter)

await b.store("s", [{"role": "user", "content": "x"}], metadata={"is_final": False, "bulk": True})

assert seen == [mod._MEMORIZE_TIMEOUT_S]
adapter.memorize.assert_awaited_once()
assert adapter.memorize.await_args.kwargs["is_final"] is False

async def test_a_final_flush_gets_the_extraction_budget(self, monkeypatch) -> None:
from raven_everos import backend as mod

Expand Down
39 changes: 23 additions & 16 deletions tests/test_importer_orchestrator.py
Original file line number Diff line number Diff line change
Expand Up @@ -250,18 +250,26 @@ async def test_a_dropped_write_is_treated_exactly_like_a_raised_one(self, tmp_pa
class TestBatching:
@pytest.mark.asyncio
async def test_msg_count_limit(self, tmp_path: Path) -> None:
"""150 messages -> 2 batches (100 + 50)."""
"""120 messages -> 3 batches (50 + 50 + 20), only the last one final, every one bulk."""
state = ImportState(path=tmp_path / "state.json")
backend = FakeBackend()
scanner = FakeScanner({"k1": _session(n_msgs=150, session_id="s1", content="x")})
scanner = FakeScanner({"k1": _session(n_msgs=120, session_id="s1", content="x")})

await run_import([(scanner, _scan_result("k1"))], backend, state)

assert len(backend.calls) == 2
assert len(backend.calls[0]["messages"]) == 100
assert backend.calls[0]["metadata"]["is_final"] is False
assert len(backend.calls[1]["messages"]) == 50
assert backend.calls[1]["metadata"]["is_final"] is True
assert [len(c["messages"]) for c in backend.calls] == [50, 50, 20]
assert [c["metadata"]["is_final"] for c in backend.calls] == [False, False, True]
assert all(c["metadata"]["bulk"] is True for c in backend.calls)

def test_a_batch_stays_inside_the_zone_everos_extracts_linearly(self) -> None:
"""Fifty is a ceiling, not a tuning knob. EverOS extracts on every add
and the cost is superlinear in the count: against a real service a
15-message batch took 12s and a 52-message batch 24s, while a batch of
100 ran past the six-minute extraction budget and failed every
memory-file source it belonged to."""
from raven.importer.orchestrator import _BATCH_MSG_LIMIT

assert _BATCH_MSG_LIMIT <= 50

@pytest.mark.asyncio
async def test_char_limit_fallback(self, tmp_path: Path) -> None:
Expand All @@ -280,10 +288,10 @@ async def test_char_limit_fallback(self, tmp_path: Path) -> None:

@pytest.mark.asyncio
async def test_is_final_only_on_last_batch(self, tmp_path: Path) -> None:
"""Exactly 100 messages -> 1 batch with is_final=True."""
"""Exactly 50 messages -> 1 batch with is_final=True."""
state = ImportState(path=tmp_path / "state.json")
backend = FakeBackend()
scanner = FakeScanner({"k1": _session(n_msgs=100, session_id="s1", content="x")})
scanner = FakeScanner({"k1": _session(n_msgs=50, session_id="s1", content="x")})

await run_import([(scanner, _scan_result("k1"))], backend, state)

Expand Down Expand Up @@ -349,19 +357,18 @@ async def test_no_tool_fields_when_absent(self, tmp_path: Path) -> None:

class TestMetadata:
@pytest.mark.asyncio
async def test_metadata_contains_is_final_only(self, tmp_path: Path) -> None:
"""app_id/project_id are deliberately omitted so EverOS defaults
to 'default'/'default', matching the daily recall partition."""
async def test_metadata_marks_the_write_bulk_and_names_no_owner(self, tmp_path: Path) -> None:
"""``bulk`` tells the backend nothing waits on this append, so it may
take its extraction budget; app_id/project_id are deliberately omitted
so EverOS defaults to 'default'/'default', matching the daily recall
partition."""
state = ImportState(path=tmp_path / "state.json")
backend = FakeBackend()
scanner = FakeScanner({"k1": _session(n_msgs=1, session_id="s1")})

await run_import([(scanner, _scan_result("k1"))], backend, state)

meta = backend.calls[0]["metadata"]
assert "app_id" not in meta
assert "project_id" not in meta
assert meta["is_final"] is True
assert backend.calls[0]["metadata"] == {"is_final": True, "bulk": True}


class TestOnProgress:
Expand Down
Loading