From e2b6b5c88d6336107cef83f856c96f6c7b959303 Mon Sep 17 00:00:00 2001 From: Prajakta Date: Fri, 21 Aug 2026 16:01:09 +0530 Subject: [PATCH 1/4] fix: lock knowledge_queue rows with SKIP LOCKED to prevent duplicate consumption DbKnowledgeSource read unconsumed knowledge_queue rows with a plain SELECT and no row lock. queue_runner.run_librarian_queue holds one open transaction across the read, the full retrieval/rerank pipeline, and the write-back, committing only at the end. Two concurrent runs could therefore both read the same unconsumed rows, both pay for the expensive pipeline work, and both persist a decision envelope for the same chunk before either reached mark_consumed. Add SELECT ... FOR UPDATE SKIP LOCKED to DbKnowledgeSource._query(), matching the with_for_update() pattern already used in db.py:set_user_resource_selection. A second concurrent reader now excludes rows the first is holding instead of blocking on them or re-reading them. Postgres-only in effect (a no-op on SQLite), so existing SQLite-backed tests are unaffected. Fixes #1025 --- .../tests/librarian/knowledge_source_test.py | 50 +++++++++++++++++++ .../utils/librarian/knowledge_source.py | 14 ++++++ 2 files changed, 64 insertions(+) diff --git a/application/tests/librarian/knowledge_source_test.py b/application/tests/librarian/knowledge_source_test.py index 27ce9c901..506b409db 100644 --- a/application/tests/librarian/knowledge_source_test.py +++ b/application/tests/librarian/knowledge_source_test.py @@ -12,8 +12,10 @@ import json import os import tempfile +import threading import unittest from datetime import datetime, timezone +from typing import List from application import create_app, sqla from application.database.db import KnowledgeQueueItem as KnowledgeQueueRow @@ -167,6 +169,54 @@ def test_unmodellable_row_is_skipped_not_fatal(self) -> None: self.assertEqual([i.id for i in items], ["a"]) + def test_concurrent_readers_skip_locked_rows(self) -> None: + """Two Module C workers must never both claim the same row. + + ``queue_runner.run_librarian_queue`` reads a batch, then runs the full + retrieval/rerank pipeline on it, then finally commits. If a second + worker's read is not fenced off from the first worker's still-open + transaction, both would process (and both would persist a decision + for) the same chunk. FOR UPDATE SKIP LOCKED must make the second + worker's read exclude rows the first is holding, rather than block on + them (which would just delay the double-processing) or read them + again. Postgres-only: SKIP LOCKED is a no-op on SQLite (there is + nothing to skip — FOR UPDATE itself is ignored), so this reproduces + the real bug only against Postgres, same as the existing + ``user_model_test`` row-lock test. + """ + if "postgresql" not in str(sqla.engine.url): + self.skipTest("row-lock serialization requires Postgres (SKIP LOCKED)") + + sqla.session.add_all([_row("a"), _row("b")]) + sqla.session.commit() + + # Worker 1: read (and thereby lock) both rows, then hold the + # transaction open -- exactly queue_runner.py's shape, which does not + # commit until the whole batch, LLM calls included, has finished. + worker1_ids = [i.id for i in DbKnowledgeSource(sqla.session).items()] + self.assertEqual(sorted(worker1_ids), ["a", "b"]) + + worker2_ids: List[str] = [] + + def worker2() -> None: + with self.app.app_context(): + try: + items = list(DbKnowledgeSource(sqla.session).items()) + worker2_ids.extend(i.id for i in items) + finally: + sqla.session.remove() + + t = threading.Thread(target=worker2) + t.start() + t.join(timeout=5) + + # Worker 2 must see neither row: both are still locked by worker 1's + # open transaction, so SKIP LOCKED excludes them instead of blocking + # or (worse) reading and reprocessing them a second time. + self.assertEqual(worker2_ids, []) + + sqla.session.rollback() + if __name__ == "__main__": unittest.main() diff --git a/application/utils/librarian/knowledge_source.py b/application/utils/librarian/knowledge_source.py index 87939dc9d..2b06324fc 100644 --- a/application/utils/librarian/knowledge_source.py +++ b/application/utils/librarian/knowledge_source.py @@ -84,6 +84,18 @@ class DbKnowledgeSource(KnowledgeSource): Rows are ordered by ``created_at`` then ``id``: the timestamp alone is not unique (B inserts a batch inside one transaction), and an unstable order would make a ``limit``ed run non-reproducible. + + **Concurrency.** ``items()`` claims every row it yields with + ``SELECT ... FOR UPDATE SKIP LOCKED`` (Postgres only — a no-op on SQLite, + matching ``db.py``'s existing ``with_for_update()`` use). The lock is held + for the life of the caller's transaction, i.e. until + ``queue_runner.run_librarian_queue`` commits after ``mark_consumed`` — which + is the whole point: without it, two concurrent runs (a retry overlapping a + scheduled pass, two orchestrator workers) would both read the same + unconsumed rows, both pay for the retrieval/rerank work, and both persist a + decision envelope for the same chunk before either reaches + ``mark_consumed``. SKIP LOCKED means a second reader simply excludes rows + the first is holding instead of blocking on them or re-reading them. """ def __init__( @@ -119,6 +131,8 @@ def _query(self) -> object: query = query.order_by(KnowledgeQueueRow.created_at, KnowledgeQueueRow.id) if self._limit is not None: query = query.limit(self._limit) + # Claim the batch: see the concurrency note on the class docstring. + query = query.with_for_update(skip_locked=True) return query def items(self) -> Iterator[KnowledgeQueueItem]: From f54e7b0ceb20a1b57354f1b043bc187856ded7b6 Mon Sep 17 00:00:00 2001 From: Prajakta Date: Fri, 21 Aug 2026 16:17:25 +0530 Subject: [PATCH 2/4] test: assert worker2 actually completes and surface its exceptions Addresses CodeRabbit review comment: t.join(timeout=5) alone doesn't confirm the thread finished -- if SKIP LOCKED failed and worker2 blocked instead, the test would still pass silently. Now asserts the thread is not alive, captures/asserts no exceptions from the worker, and moves the rollback into a finally block so a failed assertion can't leak a held lock into the next test. --- .../tests/librarian/knowledge_source_test.py | 65 +++++++++++-------- 1 file changed, 39 insertions(+), 26 deletions(-) diff --git a/application/tests/librarian/knowledge_source_test.py b/application/tests/librarian/knowledge_source_test.py index 506b409db..ad291d3c7 100644 --- a/application/tests/librarian/knowledge_source_test.py +++ b/application/tests/librarian/knowledge_source_test.py @@ -190,32 +190,45 @@ def test_concurrent_readers_skip_locked_rows(self) -> None: sqla.session.add_all([_row("a"), _row("b")]) sqla.session.commit() - # Worker 1: read (and thereby lock) both rows, then hold the - # transaction open -- exactly queue_runner.py's shape, which does not - # commit until the whole batch, LLM calls included, has finished. - worker1_ids = [i.id for i in DbKnowledgeSource(sqla.session).items()] - self.assertEqual(sorted(worker1_ids), ["a", "b"]) - - worker2_ids: List[str] = [] - - def worker2() -> None: - with self.app.app_context(): - try: - items = list(DbKnowledgeSource(sqla.session).items()) - worker2_ids.extend(i.id for i in items) - finally: - sqla.session.remove() - - t = threading.Thread(target=worker2) - t.start() - t.join(timeout=5) - - # Worker 2 must see neither row: both are still locked by worker 1's - # open transaction, so SKIP LOCKED excludes them instead of blocking - # or (worse) reading and reprocessing them a second time. - self.assertEqual(worker2_ids, []) - - sqla.session.rollback() + try: + # Worker 1: read (and thereby lock) both rows, then hold the + # transaction open -- exactly queue_runner.py's shape, which does + # not commit until the whole batch, LLM calls included, has + # finished. + worker1_ids = [i.id for i in DbKnowledgeSource(sqla.session).items()] + self.assertEqual(sorted(worker1_ids), ["a", "b"]) + + worker2_ids: List[str] = [] + worker2_errors: List[BaseException] = [] + + def worker2() -> None: + with self.app.app_context(): + try: + items = list(DbKnowledgeSource(sqla.session).items()) + worker2_ids.extend(i.id for i in items) + except BaseException as exc: # noqa: BLE001 - surface below + worker2_errors.append(exc) + finally: + sqla.session.remove() + + t = threading.Thread(target=worker2) + t.start() + t.join(timeout=5) + + # A still-running thread means SKIP LOCKED failed to exclude the + # locked rows and worker2 is blocked waiting on them instead -- + # that is a failure, not a pass, so confirm it actually finished. + self.assertFalse(t.is_alive(), "worker2 did not finish -- it is blocked") + self.assertEqual(worker2_errors, []) + + # Worker 2 must see neither row: both are still locked by worker + # 1's open transaction, so SKIP LOCKED excludes them instead of + # blocking or (worse) reading and reprocessing them a second time. + self.assertEqual(worker2_ids, []) + finally: + # Release worker 1's row locks regardless of outcome, so a failed + # assertion above cannot leak a held lock into the next test. + sqla.session.rollback() if __name__ == "__main__": From bd2e1527ac3dee902b39afccb6a0a57000a1839c Mon Sep 17 00:00:00 2001 From: prajakta G kamble Date: Tue, 1 Sep 2026 20:29:07 +0530 Subject: [PATCH 3/4] Enhance type hinting and comments in tests git add application/tests/librarian/knowledge_source_test.py git commit -m "test: rejoin worker2 after rollback on the assertion-failure path Addresses CodeRabbit review comment: if the t.is_alive() assertion fails (worker2 genuinely blocked), the finally block's rollback releases worker 1's locks and lets worker2 resume -- but tearDown() runs right after and drops the tables, racing worker2's still-running query. Join the thread again after rollback so it is guaranteed to finish before teardown." git push origin fix/1025-knowledge-queue-row-locking Signed-off-by: prajakta G kamble --- application/tests/librarian/knowledge_source_test.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/application/tests/librarian/knowledge_source_test.py b/application/tests/librarian/knowledge_source_test.py index ad291d3c7..ae87ae8ff 100644 --- a/application/tests/librarian/knowledge_source_test.py +++ b/application/tests/librarian/knowledge_source_test.py @@ -15,7 +15,7 @@ import threading import unittest from datetime import datetime, timezone -from typing import List +from typing import List, Optional from application import create_app, sqla from application.database.db import KnowledgeQueueItem as KnowledgeQueueRow @@ -187,9 +187,10 @@ def test_concurrent_readers_skip_locked_rows(self) -> None: if "postgresql" not in str(sqla.engine.url): self.skipTest("row-lock serialization requires Postgres (SKIP LOCKED)") - sqla.session.add_all([_row("a"), _row("b")]) + sqla.session.add_all([_row("a"), _row("b")]) sqla.session.commit() + t: Optional[threading.Thread] = None try: # Worker 1: read (and thereby lock) both rows, then hold the # transaction open -- exactly queue_runner.py's shape, which does @@ -229,7 +230,11 @@ def worker2() -> None: # Release worker 1's row locks regardless of outcome, so a failed # assertion above cannot leak a held lock into the next test. sqla.session.rollback() - + # If worker2 was still blocked when an assertion above failed, + # the rollback just now unblocks it -- join again so it is fully + # finished before tearDown() drops the tables out from under it. + if t is not None: + t.join(timeout=5) if __name__ == "__main__": unittest.main() From 55839a720f3e6454042e3c1869b21cb887b18b7a Mon Sep 17 00:00:00 2001 From: prajakta G kamble Date: Wed, 2 Sep 2026 14:29:12 +0530 Subject: [PATCH 4/4] Update knowledge_source_test.py test: rejoin worker2 after rollback on the assertion-failure path Addresses CodeRabbit review comment: if the t.is_alive() assertion fails (worker2 genuinely blocked), the finally block's rollback releases worker 1's locks and lets worker2 resume -- but tearDown() runs right after and drops the tables, racing worker2's still-running query. Join the thread again after rollback so it is guaranteed to finish before teardown. Signed-off-by: prajakta G kamble --- application/tests/librarian/knowledge_source_test.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/application/tests/librarian/knowledge_source_test.py b/application/tests/librarian/knowledge_source_test.py index ae87ae8ff..17dc8f9b6 100644 --- a/application/tests/librarian/knowledge_source_test.py +++ b/application/tests/librarian/knowledge_source_test.py @@ -187,7 +187,7 @@ def test_concurrent_readers_skip_locked_rows(self) -> None: if "postgresql" not in str(sqla.engine.url): self.skipTest("row-lock serialization requires Postgres (SKIP LOCKED)") - sqla.session.add_all([_row("a"), _row("b")]) + sqla.session.add_all([_row("a"), _row("b")]) sqla.session.commit() t: Optional[threading.Thread] = None @@ -236,5 +236,7 @@ def worker2() -> None: if t is not None: t.join(timeout=5) + if __name__ == "__main__": unittest.main() +