From a0a247d18a1f93b8ada1ba66c02e65646173f647 Mon Sep 17 00:00:00 2001 From: Chris McDonough Date: Mon, 1 Jun 2026 06:30:42 -0400 Subject: [PATCH] Batch sync_state writes during poller sweeps MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Each discovered file previously triggered a separate sync.upsert() call with its own lock acquire + SQLite commit (fsync). On a sweep finding 1,000 files this meant 1,000 individual commits. Collect sync_state rows into a list during the sweep and flush them in a single SyncStateRepo.batch_upsert() call at the end — one lock acquisition, one commit, one fsync. --- .../haiku/rag/ingester/pollers/base.py | 24 ++++----- .../haiku/rag/ingester/queue/repository.py | 31 +++++++++++ tests/ingester/test_queue.py | 51 +++++++++++++++++++ 3 files changed, 93 insertions(+), 13 deletions(-) diff --git a/haiku_rag_slim/haiku/rag/ingester/pollers/base.py b/haiku_rag_slim/haiku/rag/ingester/pollers/base.py index 43fbcdb6..b4e6b8b0 100644 --- a/haiku_rag_slim/haiku/rag/ingester/pollers/base.py +++ b/haiku_rag_slim/haiku/rag/ingester/pollers/base.py @@ -129,11 +129,13 @@ class BasePoller: SourceEventKind.DELETE: 0, SourceEventKind.UNCHANGED: 0, } + sync_batch: list[tuple[str, str, str | None, str | None, bool]] = [] async for event in self.source.discover( since=revisions, known_uris=known ): counts[event.kind] += 1 - await self._handle_event(event) + await self._handle_event(event, sync_batch) + await self._sync.batch_upsert(sync_batch) self._breaker.record_success() self._last_polled_at = datetime.now(UTC) self._last_skip_reason = None @@ -163,7 +165,11 @@ class BasePoller: ) return False - async def _handle_event(self, event: SourceEvent) -> None: + async def _handle_event( + self, + event: SourceEvent, + sync_batch: list[tuple[str, str, str | None, str | None, bool]], + ) -> None: if event.kind is SourceEventKind.UPSERT: await self._jobs.enqueue( event.source_id, @@ -176,19 +182,11 @@ class BasePoller: # Don't write revision to sync_state here — the worker writes it # after a successful ingestion. last_seen_at gets bumped to keep # orphan detection accurate. - await self._sync.upsert( - event.source_id, - event.uri, - revision=None, - content_hash=None, - ) + sync_batch.append((event.source_id, event.uri, None, None, False)) elif event.kind is SourceEventKind.UNCHANGED: # Touch last_seen_at without changing the stored revision. - await self._sync.upsert( - event.source_id, - event.uri, - revision=event.revision, - content_hash=None, + sync_batch.append( + (event.source_id, event.uri, event.revision, None, False) ) elif event.kind is SourceEventKind.DELETE: if not self.config.delete_orphans: diff --git a/haiku_rag_slim/haiku/rag/ingester/queue/repository.py b/haiku_rag_slim/haiku/rag/ingester/queue/repository.py index 799c0529..60070f46 100644 --- a/haiku_rag_slim/haiku/rag/ingester/queue/repository.py +++ b/haiku_rag_slim/haiku/rag/ingester/queue/repository.py @@ -504,6 +504,37 @@ class SyncStateRepo: pass await self._conn.commit() + async def batch_upsert( + self, + rows: list[tuple[str, str, str | None, str | None, bool]], + ) -> None: + """Batch insert-or-update sync_state rows in a single transaction. + Each tuple is (source_id, uri, revision, content_hash, ingested).""" + if not rows: + return + now = _utcnow_iso() + async with self._lock: + for source_id, uri, revision, content_hash, ingested in rows: + ingested_at = now if ingested else None + await self._conn.execute( + """ + INSERT INTO sync_state ( + source_id, uri, revision, content_hash, + last_seen_at, last_ingested_at + ) + VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT(source_id, uri) DO UPDATE SET + revision = COALESCE(excluded.revision, revision), + content_hash = COALESCE(excluded.content_hash, content_hash), + last_seen_at = excluded.last_seen_at, + last_ingested_at = COALESCE( + excluded.last_ingested_at, last_ingested_at + ) + """, + (source_id, uri, revision, content_hash, now, ingested_at), + ) + await self._conn.commit() + async def delete(self, source_id: str, uri: str) -> None: async with self._lock: async with self._conn.execute( diff --git a/tests/ingester/test_queue.py b/tests/ingester/test_queue.py index 4ca3ae6c..1577d805 100644 --- a/tests/ingester/test_queue.py +++ b/tests/ingester/test_queue.py @@ -834,3 +834,54 @@ async def test_sync_state_upsert_replaces_revision_when_provided(sync): assert row.revision == "v2" assert row.content_hash == "hash-v2" assert row.last_ingested_at is not None + + +@pytest.mark.asyncio +async def test_sync_state_batch_upsert_inserts_multiple_rows(sync): + """batch_upsert writes many rows in a single transaction.""" + await sync.batch_upsert([ + ("s", "u1", "rev1", "hash1", False), + ("s", "u2", "rev2", "hash2", False), + ("s", "u3", "rev3", None, True), + ]) + assert await sync.get_revision_snapshot("s") == { + "u1": "rev1", + "u2": "rev2", + "u3": "rev3", + } + row = await sync.get_row("s", "u3") + assert row is not None + assert row.last_ingested_at is not None + + +@pytest.mark.asyncio +async def test_sync_state_batch_upsert_updates_existing(sync): + """batch_upsert applies ON CONFLICT update semantics like upsert().""" + await sync.upsert("s", "u1", revision="old", content_hash="old-hash") + await sync.batch_upsert([ + ("s", "u1", "new", "new-hash", False), + ]) + row = await sync.get_row("s", "u1") + assert row is not None + assert row.revision == "new" + assert row.content_hash == "new-hash" + + +@pytest.mark.asyncio +async def test_sync_state_batch_upsert_preserves_revision_when_none(sync): + """batch_upsert with revision=None leaves existing revision in place.""" + await sync.upsert("s", "u1", revision="keep", content_hash="keep-hash") + await sync.batch_upsert([ + ("s", "u1", None, None, False), + ]) + row = await sync.get_row("s", "u1") + assert row is not None + assert row.revision == "keep" + assert row.content_hash == "keep-hash" + + +@pytest.mark.asyncio +async def test_sync_state_batch_upsert_empty_is_noop(sync): + """batch_upsert with an empty list does nothing.""" + await sync.batch_upsert([]) + assert await sync.get_revision_snapshot("s") == {}