项目文件夹

文件
wehub-resource-sync 7a0da7932b
OSV-Scanner (Scheduled) / scan-scheduled (push) Failing after 0s
Create Release / test-gate (push) Has been cancelled
Create Release / release-gate (push) Has been cancelled
Create Release / ci-gate (push) Has been cancelled
Create Release / version-check (push) Has been cancelled
Create Release / e2e-test-gate (push) Has been cancelled
Create Release / responsive-test-gate (push) Has been cancelled
Create Release / compat-test-gate (push) Has been cancelled
Create Release / compose-integration-gate (push) Has been cancelled
Create Release / vulture-gate (push) Has been cancelled
Create Release / build (push) Has been cancelled
Create Release / provenance (push) Has been cancelled
Create Release / prerelease-docker (push) Has been cancelled
Create Release / publish-docker (push) Has been cancelled
Create Release / create-release (push) Has been cancelled
Create Release / cleanup-changelog (push) Has been cancelled
Create Release / trigger-pypi (push) Has been cancelled
Create Release / monitor-pypi (push) Has been cancelled
Create Release / Clean up orphan prerelease tags and signatures (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [research-form] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [research-metrics] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [research-workflow] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [settings-core] (push) Has been cancelled
CodeQL Advanced / Analyze (javascript-typescript) (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [history-news] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [library] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [link-analytics] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [chat-core] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [chat-lifecycle] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [error-benchmark] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [settings-pages] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) (push) Has been cancelled
Docker Tests (Consolidated) / Accessibility Tests (push) Has been cancelled
Docker Tests (Consolidated) / LLM Unit Tests (push) Has been cancelled
Docker Tests (Consolidated) / LLM Example Tests (push) Has been cancelled
Docker Tests (Consolidated) / Production Image Smoke Test (push) Has been cancelled
Docker Tests (Consolidated) / Infrastructure Tests (push) Has been cancelled
OSSF Scorecard / OSSF Security Scorecard Analysis (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [mobile] (push) Has been cancelled
Backwards Compatibility / Verify Encryption Constants (push) Has been cancelled
Backwards Compatibility / PyPI Version Compatibility (push) Has been cancelled
Backwards Compatibility / Database Migration Tests (push) Has been cancelled
CodeQL Advanced / Analyze (python) (push) Has been cancelled
Docker Tests (Consolidated) / detect-changes (push) Has been cancelled
Docker Tests (Consolidated) / Build Test Image (push) Has been cancelled
Docker Tests (Consolidated) / All Pytest Tests + Coverage (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [accessibility] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [api-crud] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [auth-login] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [auth-pages] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [auth-register] (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:08:55 +08:00

1231 行
45 KiB
Python

"""Tests for the unified unindexed-document reconciler.
These drive ``BackgroundJobScheduler._reconcile_unindexed_documents`` (the
self-healing reconciler that indexes EVERY unindexed document — both in-
collection library uploads missed by PR #3939's drop-on-saturation auto-index
queue AND research downloads that were never ingested) and
``_schedule_reconciler`` (the job lifecycle) with focused mocks, mirroring the
mocking style of ``test_download_service_3827_fix.py``.
The reconciler handles two cases per tick under one batch cap:
(a) in-collection unindexed docs (DocumentCollection.indexed == False)
(b) research orphans (Document.research_id set, no default-library link)
and is gated by ``sweep_library_collections OR generate_rag`` (the generate_rag
arm preserves the legacy research-download indexing behaviour).
"""
from contextlib import contextmanager
from datetime import datetime, UTC
from unittest.mock import MagicMock, patch
from local_deep_research.scheduler.background import (
_LIBRARY_SWEEP_BATCH,
BackgroundJobScheduler,
DocumentSchedulerSettings,
)
MODULE = "local_deep_research.scheduler.background"
# get_default_library_id is imported at module scope in background.py, so patch
# the name bound there (not its definition module).
DEFAULT_LIB = f"{MODULE}.get_default_library_id"
# get_rag_service is imported lazily inside the reconciler, so patch it at its
# definition module.
FACTORY = (
"local_deep_research.research_library.services.rag_service_factory"
".get_rag_service"
)
def _fresh_scheduler():
"""Return a fresh scheduler singleton with a registered test user."""
BackgroundJobScheduler._instance = None
with patch(f"{MODULE}.BackgroundScheduler"):
scheduler = BackgroundJobScheduler()
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "testpass")
return scheduler
def _db_returning(case_a_rows, case_b_rows=None):
"""Build a MagicMock db whose query chain returns ``case_a_rows`` on the
first ``.all()`` (case-a in-collection query) and ``case_b_rows`` on the
second (case-b research-orphan query).
case_a_rows: list of (document_id, collection_id) tuples.
case_b_rows: list of (document_id,) tuples; defaults to empty.
"""
if case_b_rows is None:
case_b_rows = []
db = MagicMock()
query = MagicMock()
query.join.return_value = query
query.outerjoin.return_value = query
query.filter.return_value = query
query.order_by.return_value = query
query.limit.return_value = query
# First .all() => case (a) rows, second .all() => case (b) rows.
query.all.side_effect = [case_a_rows, case_b_rows]
db.query.return_value = query
return db, query
def _rag_service():
"""Return a MagicMock RAG service usable as a context manager that reports
successful indexing.
"""
rag = MagicMock()
rag.__enter__ = MagicMock(return_value=rag)
rag.__exit__ = MagicMock(return_value=False)
rag.index_document.return_value = {"status": "success", "chunk_count": 3}
return rag
@contextmanager
def _patched_session(db):
"""Patch the lazily-imported get_user_db_session to yield ``db`` and stub
SettingsManager + the egress backstop arming + the default-library lookup.
"""
@contextmanager
def fake_get_user_db_session(*a, **kw):
yield db
with (
patch(
"local_deep_research.database.session_context.get_user_db_session",
side_effect=fake_get_user_db_session,
),
patch(
"local_deep_research.settings.manager.SettingsManager",
return_value=MagicMock(),
),
patch.object(
BackgroundJobScheduler,
"_arm_egress_backstop",
lambda *a, **kw: None,
),
patch(DEFAULT_LIB, return_value="default-lib"),
):
yield
def test_reconciler_indexes_in_collection_unindexed_when_enabled():
"""Case (a): with sweep ON, every in-collection unindexed doc is indexed via
the per-collection RAG service with force_reindex=False; one service per
distinct collection.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
db, _ = _db_returning(
[("doc-1", "coll-A"), ("doc-2", "coll-A"), ("doc-3", "coll-B")],
case_b_rows=[],
)
rag_service = _rag_service()
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service) as mock_factory,
):
scheduler._reconcile_unindexed_documents("testuser")
# Each in-collection unindexed doc indexed exactly once, never reindexed.
assert rag_service.index_document.call_count == 3
for call in rag_service.index_document.call_args_list:
assert call.kwargs["force_reindex"] is False
# One RAG service built per distinct collection (A, B). No orphans here so
# the default-library service is not built.
built_collections = {
c.kwargs.get("collection_id") for c in mock_factory.call_args_list
}
assert built_collections == {"coll-A", "coll-B"}
def test_reconciler_ingests_and_indexes_research_orphan():
"""Case (b): a research orphan (research_id set, no default-library link) is
ingested into the default library AND indexed via index_document into the
default-library collection. index_document's ensure_in_collection does the
ingest, so a single call ingests + indexes.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# No in-collection work (case a empty), one research orphan (case b).
db, _ = _db_returning([], case_b_rows=[("orphan-doc",)])
rag_service = _rag_service()
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service) as mock_factory,
):
scheduler._reconcile_unindexed_documents("testuser")
# The orphan was indexed into the DEFAULT library collection.
rag_service.index_document.assert_called_once()
call = rag_service.index_document.call_args
assert call.kwargs["document_id"] == "orphan-doc"
assert call.kwargs["collection_id"] == "default-lib"
assert call.kwargs["force_reindex"] is False
# The RAG service for the orphan path was built for the default library.
assert any(
c.kwargs.get("collection_id") == "default-lib"
for c in mock_factory.call_args_list
)
def test_reconciler_runs_for_generate_rag_only_no_regression():
"""No-regression path: with sweep_library_collections OFF but the legacy
generate_rag ON, the reconciler still runs and indexes research orphans —
preserving the behaviour of the retired _process_user_documents RAG block.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(
sweep_library_collections=False, generate_rag=True
)
assert settings.sweep_library_collections is False
db, _ = _db_returning([], case_b_rows=[("orphan-doc",)])
rag_service = _rag_service()
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service),
):
scheduler._reconcile_unindexed_documents("testuser")
# generate_rag alone is enough to drive the reconciler.
rag_service.index_document.assert_called_once()
assert (
rag_service.index_document.call_args.kwargs["collection_id"]
== "default-lib"
)
def test_reconciler_is_noop_when_both_disabled():
"""With BOTH gating settings OFF (default), the reconciler returns before
opening a DB session or building any RAG service.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings() # both False
assert settings.sweep_library_collections is False
assert settings.generate_rag is False
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
patch(
"local_deep_research.database.session_context.get_user_db_session",
) as mock_session,
patch(FACTORY) as mock_factory,
):
scheduler._reconcile_unindexed_documents("testuser")
mock_session.assert_not_called()
mock_factory.assert_not_called()
def test_reconciler_is_noop_when_scheduler_disabled():
"""When the document scheduler is disabled (enabled=False) the reconciler
must no-op even if sweep_library_collections is on — the already-live job
can keep firing after a runtime disable until the next reschedule, and the
setting promises it only runs while the scheduler is enabled.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(
enabled=False, sweep_library_collections=True
)
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
patch(
"local_deep_research.database.session_context.get_user_db_session",
) as mock_session,
patch(FACTORY) as mock_factory,
):
scheduler._reconcile_unindexed_documents("testuser")
mock_session.assert_not_called()
mock_factory.assert_not_called()
def test_reconciler_filters_to_unindexed_only_no_work():
"""The case-(a) query filters on DocumentCollection.indexed == False so an
already-indexed document is never selected. With both queries empty, no RAG
service is built (idempotency: already-indexed docs are never touched).
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
db, query = _db_returning([], case_b_rows=[]) # nothing to do
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY) as mock_factory,
):
scheduler._reconcile_unindexed_documents("testuser")
# Queries were issued (filters applied) but selected nothing.
assert query.filter.called
mock_factory.assert_not_called()
def test_reconciler_caps_each_case_at_its_own_batch_budget():
"""Each case has its OWN _LIBRARY_SWEEP_BATCH budget (independent, NOT a
shared/leftover budget). Even when case (a) returns a full batch, the
research-orphan path still runs with its own LIMIT(_LIBRARY_SWEEP_BATCH), so
total work is bounded at 2 x _LIBRARY_SWEEP_BATCH per tick.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# Case (a) returns a full batch; case (b) also has rows to process.
rows = [(f"doc-{i}", "coll-A") for i in range(_LIBRARY_SWEEP_BATCH)]
db, query = _db_returning(rows, case_b_rows=[("orphan-doc",)])
rag_service = _rag_service()
rag_service.index_document.return_value = {
"status": "success",
"chunk_count": 1,
}
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service),
):
scheduler._reconcile_unindexed_documents("testuser")
# BOTH queries are capped at the module batch constant (independent budgets,
# not a single shared one). With case (a) at full batch the orphan query
# must STILL use the full _LIBRARY_SWEEP_BATCH limit (no `remaining`
# subtraction), so .limit is called with the constant for both cases.
assert query.limit.call_args_list == [
((_LIBRARY_SWEEP_BATCH,), {}),
((_LIBRARY_SWEEP_BATCH,), {}),
]
# Case (a) docs (BATCH) plus the one orphan were all indexed: the orphan was
# NOT starved by a full case (a).
assert rag_service.index_document.call_count == _LIBRARY_SWEEP_BATCH + 1
indexed_collections = {
c.kwargs["collection_id"]
for c in rag_service.index_document.call_args_list
}
assert indexed_collections == {"coll-A", "default-lib"}
def test_reconciler_per_doc_error_does_not_crash_tick_or_leak_password(
loguru_caplog_full,
):
"""A per-doc indexing error is logged and the loop continues with the next
doc; the password is never present in the captured log output.
Uses ``loguru_caplog_full`` (which re-enables the package's loguru logging,
disabled by default in tests) so the leak assertion is NOT vacuous — the
error line really is captured and checked for the secret.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
db, _ = _db_returning(
[("doc-1", "coll-A"), ("doc-2", "coll-A")], case_b_rows=[]
)
rag_service = _rag_service()
# First doc raises an error containing the password; second succeeds.
rag_service.index_document.side_effect = [
RuntimeError("boom with secret testpass inside"),
{"status": "success", "chunk_count": 2},
]
with loguru_caplog_full.at_level("DEBUG"):
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service),
):
# Must not raise even though the first doc errored.
scheduler._reconcile_unindexed_documents("testuser")
# Both docs were attempted (the error didn't abort the loop).
assert rag_service.index_document.call_count == 2
text = loguru_caplog_full.text
# The per-doc error handler ran (path exercised) ...
assert "[RECONCILER] Failed to index document" in text
# ... and the plaintext password never appears in the log output.
assert "testpass" not in text, "Password leaked into reconciler log output"
def test_schedule_reconciler_adds_job_when_sweep_enabled():
"""_schedule_reconciler registers the interval job (with max_instances=1)
and tracks it in the session's scheduled-jobs set when sweep is ON.
"""
scheduler = _fresh_scheduler()
session_info = scheduler.user_sessions["testuser"]
settings = DocumentSchedulerSettings(
sweep_library_collections=True, interval_seconds=1800
)
# remove_job has no prior job -> raise JobLookupError so the add path runs.
from apscheduler.jobstores.base import JobLookupError
scheduler.scheduler.remove_job.side_effect = JobLookupError("library_sweep")
scheduler._schedule_reconciler("testuser", settings, session_info)
assert scheduler.scheduler.add_job.called
kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert kwargs["id"] == "testuser_library_sweep"
assert kwargs["trigger"] == "interval"
assert kwargs["seconds"] == 1800
assert kwargs["max_instances"] == 1
assert "testuser_library_sweep" in session_info["scheduled_jobs"]
def test_schedule_reconciler_adds_job_when_generate_rag_enabled():
"""The legacy generate_rag setting alone also schedules the reconciler
(no-regression: research downloads keep getting indexed).
"""
scheduler = _fresh_scheduler()
session_info = scheduler.user_sessions["testuser"]
settings = DocumentSchedulerSettings(
sweep_library_collections=False,
generate_rag=True,
interval_seconds=1800,
)
from apscheduler.jobstores.base import JobLookupError
scheduler.scheduler.remove_job.side_effect = JobLookupError("library_sweep")
scheduler._schedule_reconciler("testuser", settings, session_info)
assert scheduler.scheduler.add_job.called
assert (
scheduler.scheduler.add_job.call_args.kwargs["id"]
== "testuser_library_sweep"
)
assert "testuser_library_sweep" in session_info["scheduled_jobs"]
def test_schedule_reconciler_removes_job_when_both_disabled():
"""When BOTH gating settings are OFF, no job is added and any existing one is
removed from both the scheduler and the session's tracked-jobs set.
"""
scheduler = _fresh_scheduler()
session_info = scheduler.user_sessions["testuser"]
session_info["scheduled_jobs"].add("testuser_library_sweep")
settings = DocumentSchedulerSettings(
sweep_library_collections=False, generate_rag=False
)
scheduler._schedule_reconciler("testuser", settings, session_info)
scheduler.scheduler.remove_job.assert_called_once_with(
"testuser_library_sweep"
)
scheduler.scheduler.add_job.assert_not_called()
assert "testuser_library_sweep" not in session_info["scheduled_jobs"]
# ---------------------------------------------------------------------------
# reschedule_document_jobs: the settings-save seam that makes a
# document_scheduler.* toggle take effect on the next tick instead of only
# after the user logs out and back in.
# ---------------------------------------------------------------------------
def test_reschedule_document_jobs_creates_reconciler_on_enable():
"""Toggling sweep ON re-creates the {username}_library_sweep job for an
active user (no re-login needed)."""
scheduler = _fresh_scheduler()
scheduler.is_running = True
settings = DocumentSchedulerSettings(
enabled=True, sweep_library_collections=True, interval_seconds=1800
)
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
result = scheduler.reschedule_document_jobs("testuser")
assert result is True
scheduled_ids = {
call.kwargs.get("id")
for call in scheduler.scheduler.add_job.call_args_list
}
assert "testuser_library_sweep" in scheduled_ids
def test_reschedule_document_jobs_removes_reconciler_on_disable():
"""Toggling both gating settings OFF tears the reconciler job down on the
next save (not only on the next login)."""
scheduler = _fresh_scheduler()
scheduler.is_running = True
scheduler.user_sessions["testuser"]["scheduled_jobs"].add(
"testuser_library_sweep"
)
settings = DocumentSchedulerSettings(
enabled=True,
sweep_library_collections=False,
generate_rag=False,
interval_seconds=1800,
)
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
result = scheduler.reschedule_document_jobs("testuser")
assert result is True
scheduler.scheduler.remove_job.assert_any_call("testuser_library_sweep")
add_ids = {
call.kwargs.get("id")
for call in scheduler.scheduler.add_job.call_args_list
}
assert "testuser_library_sweep" not in add_ids
assert (
"testuser_library_sweep"
not in scheduler.user_sessions["testuser"]["scheduled_jobs"]
)
def test_reschedule_document_jobs_noop_for_inactive_user():
"""A user the scheduler isn't tracking is skipped (their jobs are rebuilt
from current settings on their next login) — returns False, no scheduling.
"""
scheduler = _fresh_scheduler()
scheduler.is_running = True
result = scheduler.reschedule_document_jobs("ghost")
assert result is False
scheduler.scheduler.add_job.assert_not_called()
def test_reschedule_document_jobs_noop_when_scheduler_not_running():
"""No rescheduling is attempted when the scheduler isn't running."""
scheduler = _fresh_scheduler()
scheduler.is_running = False
result = scheduler.reschedule_document_jobs("testuser")
assert result is False
scheduler.scheduler.add_job.assert_not_called()
def test_reschedule_helper_only_fires_for_document_scheduler_keys():
"""The settings-save helper reschedules ONLY when a document_scheduler.* key
changed, so an unrelated settings save never churns the document jobs."""
from local_deep_research.web.services.settings_service import (
reschedule_document_jobs_if_needed,
)
with patch(f"{MODULE}.get_background_job_scheduler") as get_sched:
sched = get_sched.return_value
# Unrelated key -> the scheduler is never even looked up.
reschedule_document_jobs_if_needed("u", ["llm.temperature"])
get_sched.assert_not_called()
# A document_scheduler.* key -> reschedule that user's document jobs.
reschedule_document_jobs_if_needed(
"u", ["document_scheduler.sweep_library_collections"]
)
sched.reschedule_document_jobs.assert_called_once_with("u")
# Missing / empty inputs are no-ops (return before touching the scheduler).
reschedule_document_jobs_if_needed(None, ["document_scheduler.enabled"])
reschedule_document_jobs_if_needed("u", [])
# ---------------------------------------------------------------------------
# Finding 1: case (b) must not be starved when case (a) is full of failures
# ---------------------------------------------------------------------------
def test_case_b_not_starved_when_case_a_full_batch_all_fail():
"""Regression guard for the case-(b) starvation bug.
Case (a) returns a FULL batch of in-collection docs that ALL fail to index
(index_document raises every time, so none flip to indexed=True), AND there
is a research orphan pending. With the old shared/leftover budget
(``remaining = _LIBRARY_SWEEP_BATCH - len(unindexed)`` == 0) the orphan path
would never run. With independent per-case budgets it MUST still run and
index the orphan.
Mutation check: revert the fix (gate case (b) on
``remaining = _LIBRARY_SWEEP_BATCH - len(unindexed)`` and ``if remaining >
0``) and this test fails because the orphan is never indexed.
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# Case (a): a full batch, all in coll-A. Case (b): one orphan.
case_a = [(f"doc-{i}", "coll-A") for i in range(_LIBRARY_SWEEP_BATCH)]
db, _ = _db_returning(case_a, case_b_rows=[("orphan-doc",)])
# One service object, used for BOTH the per-collection path (coll-A) and the
# default-library path. Case-(a) docs always raise; the orphan succeeds.
rag_service = _rag_service()
def index_side_effect(*, document_id, collection_id, force_reindex):
if collection_id == "coll-A":
raise RuntimeError(f"index failed for {document_id}")
return {"status": "success", "chunk_count": 1}
rag_service.index_document.side_effect = index_side_effect
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service),
):
scheduler._reconcile_unindexed_documents("testuser")
# The orphan was indexed into the default library DESPITE case (a) filling a
# whole batch with failures: case (b) is not starved.
orphan_calls = [
c
for c in rag_service.index_document.call_args_list
if c.kwargs["collection_id"] == "default-lib"
]
assert len(orphan_calls) == 1
assert orphan_calls[0].kwargs["document_id"] == "orphan-doc"
# ---------------------------------------------------------------------------
# Finding 2: real in-memory DB so the actual query predicates are exercised
# ---------------------------------------------------------------------------
def _seed_document(
session,
source_type_id,
*,
doc_id,
text_content="some text",
research_id=None,
):
"""Insert one Document row with all NOT-NULL columns satisfied."""
from local_deep_research.database.models.library import (
Document,
DocumentStatus,
)
doc = Document(
id=doc_id,
source_type_id=source_type_id,
research_id=research_id,
document_hash=(doc_id + "x" * 64)[:64],
file_size=10,
file_type="txt",
text_content=text_content,
status=DocumentStatus.COMPLETED,
)
session.add(doc)
return doc
def _real_db_session():
"""Build an in-memory SQLite session with the library schema created and a
SourceType + two collections (a default library + a normal collection)
seeded. Returns (session, source_type_id, default_lib_id, coll_id).
"""
import uuid
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from local_deep_research.database.models import Base
from local_deep_research.database.models.library import (
Collection,
SourceType,
)
engine = create_engine("sqlite:///:memory:")
Base.metadata.create_all(engine)
session = sessionmaker(bind=engine)()
st = SourceType(
id=str(uuid.uuid4()),
name="research_download",
display_name="Research Download",
)
default_lib = Collection(
id="default-lib",
name="Library",
collection_type="default_library",
is_default=True,
)
other = Collection(
id="coll-A",
name="Coll A",
collection_type="user_collection",
)
session.add_all([st, default_lib, other])
session.commit()
return session, st.id, "default-lib", "coll-A"
def test_case_a_predicate_selects_only_unindexed_with_text_real_db():
"""Finding 2 (a): with a REAL DB, the case-(a) query selects an unindexed
in-collection doc with text, and excludes an already-indexed doc and a
null-text doc. A wrong predicate (e.g. indexed.is_(True)) fails here.
"""
from local_deep_research.database.models.library import DocumentCollection
session, st_id, default_lib_id, coll_id = _real_db_session()
# Doc 1: unindexed, has text, in coll-A -> SHOULD be selected.
_seed_document(session, st_id, doc_id="doc-unindexed", text_content="hi")
# Doc 2: already indexed in coll-A -> should NOT be selected.
_seed_document(session, st_id, doc_id="doc-indexed", text_content="hi")
# Doc 3: unindexed but NULL text -> should NOT be selected.
_seed_document(session, st_id, doc_id="doc-notext", text_content=None)
session.add_all(
[
DocumentCollection(
document_id="doc-unindexed",
collection_id=coll_id,
indexed=False,
),
DocumentCollection(
document_id="doc-indexed",
collection_id=coll_id,
indexed=True,
),
DocumentCollection(
document_id="doc-notext",
collection_id=coll_id,
indexed=False,
),
]
)
session.commit()
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
rag_service = _rag_service()
indexed_ids = []
rag_service.index_document.side_effect = (
lambda *, document_id, collection_id, force_reindex: (
indexed_ids.append(document_id)
or {"status": "success", "chunk_count": 1}
)
)
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(session),
patch(FACTORY, return_value=rag_service),
):
scheduler._reconcile_unindexed_documents("testuser")
# Only the unindexed, text-bearing, in-collection doc was selected for case
# (a). (doc-indexed excluded by indexed.is_(False); doc-notext excluded by
# text_content.isnot(None).)
case_a_indexed = [
i for i in indexed_ids if i in {"doc-unindexed", "doc-indexed"}
]
assert case_a_indexed == ["doc-unindexed"]
assert "doc-notext" not in indexed_ids
def test_case_b_predicate_selects_research_orphans_real_db():
"""Finding 2 (b): with a REAL DB, the case-(b) outerjoin/predicate selects a
research orphan (research_id set, text not null, NO default-library link)
and excludes a research doc already linked in the default library. Inverting
``DocumentCollection.id.is_(None)`` -> ``isnot(None)`` fails here.
"""
from local_deep_research.database.models.library import DocumentCollection
session, st_id, default_lib_id, coll_id = _real_db_session()
# Orphan: research_id set, text, NO link in default library -> selected.
_seed_document(
session, st_id, doc_id="orphan", text_content="t", research_id="r1"
)
# Already in default library -> NOT selected by case (b).
_seed_document(
session, st_id, doc_id="in-lib", text_content="t", research_id="r2"
)
session.add(
DocumentCollection(
document_id="in-lib",
collection_id=default_lib_id,
indexed=True,
)
)
session.commit()
scheduler = _fresh_scheduler()
# generate_rag-only: the no-regression path that depends on case (b).
settings = DocumentSchedulerSettings(
sweep_library_collections=False, generate_rag=True
)
rag_service = _rag_service()
indexed_into_default = []
rag_service.index_document.side_effect = (
lambda *, document_id, collection_id, force_reindex: (
indexed_into_default.append((document_id, collection_id))
or {"status": "success", "chunk_count": 1}
)
)
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(session),
patch(FACTORY, return_value=rag_service),
):
scheduler._reconcile_unindexed_documents("testuser")
default_targets = [
doc_id
for (doc_id, coll) in indexed_into_default
if coll == "default-lib"
]
# The orphan was ingested into the default library; the already-linked doc
# was NOT re-selected by the orphan query.
assert default_targets == ["orphan"]
# ---------------------------------------------------------------------------
# Finding 3: per-collection isolation + orphan/per-collection password redaction
# ---------------------------------------------------------------------------
def test_per_collection_error_isolated_other_collections_still_indexed():
"""Finding 3 (i): when coll-B's get_rag_service raises, the tick completes,
coll-A's docs are still indexed, and the shared session was rolled back via
safe_rollback (the per-collection coll_error handler).
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
db, _ = _db_returning(
[("doc-a", "coll-A"), ("doc-b", "coll-B")], case_b_rows=[]
)
good_service = _rag_service()
def factory_side_effect(username, *, collection_id, db_password):
if collection_id == "coll-B":
raise RuntimeError("coll-B factory boom")
return good_service
rollback_calls = []
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, side_effect=factory_side_effect),
patch(
f"{MODULE}.safe_rollback",
side_effect=lambda *a, **kw: rollback_calls.append(a),
),
):
# Must not raise even though coll-B blew up.
scheduler._reconcile_unindexed_documents("testuser")
# coll-A's doc was indexed despite coll-B failing (cross-collection
# isolation).
indexed_docs = {
c.kwargs["document_id"]
for c in good_service.index_document.call_args_list
}
assert indexed_docs == {"doc-a"}
# The per-collection handler rolled back the poisoned session at least once.
assert len(rollback_calls) >= 1
def test_per_collection_error_does_not_leak_password(loguru_caplog_full):
"""Finding 3 (ii): a per-collection failure whose message contains the
sentinel password must not leak it into the log output (coll_error path uses
redact_secrets).
Uses ``loguru_caplog_full`` because the package disables loguru by default
in tests; without re-enabling it the capture would be empty and the leak
assertion would pass vacuously (it also captures the exception block).
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
db, _ = _db_returning([("doc-a", "coll-A")], case_b_rows=[])
def factory_side_effect(username, *, collection_id, db_password):
raise RuntimeError("collection failure with secret testpass inside")
with loguru_caplog_full.at_level("DEBUG"):
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, side_effect=factory_side_effect),
):
scheduler._reconcile_unindexed_documents("testuser")
text = loguru_caplog_full.text
# The per-collection error handler ran (proving the path was exercised).
assert "[RECONCILER] Failed to index collection" in text
# ...and the plaintext password never reached the log output.
assert "testpass" not in text, (
"Password leaked into per-collection error log output"
)
def test_orphan_path_error_does_not_leak_password(loguru_caplog_full):
"""Finding 3 (ii): an orphan-path failure whose message contains the
sentinel password must not leak it (orphan_error path uses redact_secrets).
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# No case (a) work; one orphan whose service build raises.
db, _ = _db_returning([], case_b_rows=[("orphan-doc",)])
def factory_side_effect(username, *, collection_id, db_password):
raise RuntimeError("orphan failure with secret testpass inside")
with loguru_caplog_full.at_level("DEBUG"):
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, side_effect=factory_side_effect),
):
scheduler._reconcile_unindexed_documents("testuser")
text = loguru_caplog_full.text
assert "[RECONCILER] Failed to index research" in text
assert "testpass" not in text, (
"Password leaked into orphan error log output"
)
# ---------------------------------------------------------------------------
# Finding 4: partial case-(a) budget; case (b) still uses its own full budget
# ---------------------------------------------------------------------------
def test_case_a_selection_is_randomized_not_deterministic_real_db():
"""Finding 1: the case-(a) query uses ``order_by(func.random())`` (not a
stable id order) so a block of permanently-failing low-id rows can't pin the
LIMIT slots every tick and starve indexable higher-id rows forever.
Seed MORE than ``_LIBRARY_SWEEP_BATCH`` selectable in-collection unindexed
rows and run the case-(a) SELECTION several times WITHOUT letting any row
flip to indexed=True between draws (``index_document`` is a no-op stub that
only records which ids were drawn). With a pool of (BATCH + 25) and a LIMIT
of BATCH, the probability that two random draws return the identical
BATCH-subset is astronomically small (C(75, 50) is ~2.5e19 possible
subsets), so "not all draws identical" is a deterministic pass in practice.
Mutation check: change BOTH ``order_by(func.random())`` to a stable
``order_by(DocumentCollection.document_id)`` / ``order_by(Document.id)`` and
this test fails because every draw returns the SAME id set.
"""
from local_deep_research.database.models.library import DocumentCollection
session, st_id, default_lib_id, coll_id = _real_db_session()
# Seed a pool LARGER than the batch so the LIMIT(BATCH) must subsample it.
pool_size = _LIBRARY_SWEEP_BATCH + 25
links = []
for i in range(pool_size):
# Zero-padded ids so a stable id-order would be a well-defined, fixed
# prefix (making the mutation's determinism unambiguous).
doc_id = f"doc-{i:04d}"
_seed_document(session, st_id, doc_id=doc_id, text_content="hi")
links.append(
DocumentCollection(
document_id=doc_id,
collection_id=coll_id,
indexed=False,
)
)
session.add_all(links)
session.commit()
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# Capture the set of document ids drawn on the most recent reconcile run.
# CRITICAL: the stub does NOT flip indexed=True, so every draw sees the same
# candidate pool — any change in the selected set is due to random ordering,
# not to rows leaving the pool.
last_draw: set = set()
rag_service = _rag_service()
def index_side_effect(*, document_id, collection_id, force_reindex):
last_draw.add(document_id)
return {"status": "success", "chunk_count": 1}
rag_service.index_document.side_effect = index_side_effect
draws = []
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(session),
patch(FACTORY, return_value=rag_service),
):
# Up to 4 draws for extra robustness: assert they are NOT all identical.
for _ in range(4):
last_draw = set()
scheduler._reconcile_unindexed_documents("testuser")
draws.append(frozenset(last_draw))
# Sanity: each draw selected exactly the LIMIT (the pool is larger).
for draw in draws:
assert len(draw) == _LIBRARY_SWEEP_BATCH
# Randomized order_by => the draws are NOT all the same BATCH-subset. A
# stable order_by(id) would make every draw identical (mutation fails here).
assert len(set(draws)) > 1, (
"case-(a) selection appears deterministic across draws; "
"order_by(func.random()) regressed to a stable order"
)
def test_case_b_orphan_selection_is_randomized_not_deterministic_real_db():
"""Finding 4: case-(b) (research orphans) ALSO uses order_by(func.random())
so a block of permanently-failing low-id orphans can't pin the LIMIT slots
every tick and starve the rest. The case-(a) randomization test above seeds
only case-(a) rows (research_id=None) so it does NOT guard case (b); this
seeds > _LIBRARY_SWEEP_BATCH research orphans.
Mutation check: change case (b)'s order_by(func.random()) to a stable
order_by(Document.id) and this test fails (every draw identical) while the
case-(a) test still passes.
"""
session, st_id, default_lib_id, coll_id = _real_db_session()
# Seed a pool of research ORPHANS (research_id set, text, and NO
# default-library DocumentCollection link) larger than the batch. With no
# link they are excluded from case (a) (which inner-joins DocumentCollection)
# and selected only by case (b).
pool_size = _LIBRARY_SWEEP_BATCH + 25
for i in range(pool_size):
_seed_document(
session,
st_id,
doc_id=f"orphan-{i:04d}",
text_content="hi",
research_id=f"r{i:04d}",
)
session.commit()
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# No-op stub: does NOT create the default-library link, so every orphan
# stays an orphan and the candidate pool is stable across draws — any
# variation in the selected set is due to random ordering alone.
last_draw: set = set()
rag_service = _rag_service()
def index_side_effect(*, document_id, collection_id, force_reindex):
last_draw.add(document_id)
return {"status": "success", "chunk_count": 1}
rag_service.index_document.side_effect = index_side_effect
draws = []
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(session),
patch(FACTORY, return_value=rag_service),
):
for _ in range(4):
last_draw = set()
scheduler._reconcile_unindexed_documents("testuser")
draws.append(frozenset(last_draw))
for draw in draws:
assert len(draw) == _LIBRARY_SWEEP_BATCH
assert len(set(draws)) > 1, (
"case-(b) orphan selection appears deterministic across draws; "
"order_by(func.random()) regressed to a stable order"
)
def test_partial_case_a_budget_orphan_uses_own_full_budget():
"""Finding 4: when case (a) returns a PARTIAL batch (0 < N < BATCH), the
orphan query still uses its OWN full _LIBRARY_SWEEP_BATCH limit (independent
budgets), and total index_document work is bounded by N + (orphans found).
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# Partial case-(a) batch: 0 < n < _LIBRARY_SWEEP_BATCH so the OLD shared
# budget would have reduced the orphan limit to (BATCH - n).
n = 3
case_a = [(f"doc-{i}", "coll-A") for i in range(n)]
db, query = _db_returning(case_a, case_b_rows=[("orphan-doc",)])
rag_service = _rag_service()
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
_patched_session(db),
patch(FACTORY, return_value=rag_service),
):
scheduler._reconcile_unindexed_documents("testuser")
# Both queries used the full per-case budget — the orphan query's limit is
# NOT reduced by the partial case (a) count (no `remaining` subtraction).
assert query.limit.call_args_list == [
((_LIBRARY_SWEEP_BATCH,), {}),
((_LIBRARY_SWEEP_BATCH,), {}),
]
# Total work = N case-(a) docs + 1 orphan, all bounded.
assert rag_service.index_document.call_count == n + 1
# ---------------------------------------------------------------------------
# Finding 2: the egress backstop is armed for the scheduled reconcile run
# ---------------------------------------------------------------------------
def test_reconciler_arms_egress_backstop_in_body():
"""Finding 2: once past the gate, the reconciler arms the PEP-578 audit-hook
egress backstop exactly once via ``_arm_egress_backstop(settings_manager,
username)`` — defense-in-depth parity with an interactive run, since embedding
providers may make network calls during indexing.
Drives the reconciler fully THROUGH the body (enabled + sweep ON) with both
the case-(a) and case-(b) queries returning EMPTY so it does minimal work,
and asserts the backstop was armed once with (settings_manager, username).
Uses a spy (MagicMock) for ``_arm_egress_backstop`` rather than the no-op
lambda the shared helper installs, so the call is actually observed.
Mutation check: delete the ``self._arm_egress_backstop(...)`` line and this
test fails (the spy is never called).
"""
scheduler = _fresh_scheduler()
settings = DocumentSchedulerSettings(sweep_library_collections=True)
# Both queries empty => the reconciler still passes the gate, opens the
# session, arms the backstop, then finds no work.
db, _ = _db_returning([], case_b_rows=[])
@contextmanager
def fake_get_user_db_session(*a, **kw):
yield db
# Capture the SettingsManager instance constructed inside the body so we can
# assert it is exactly what was passed to the backstop.
settings_manager_instance = MagicMock()
with (
patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=settings,
),
patch(
"local_deep_research.database.session_context.get_user_db_session",
side_effect=fake_get_user_db_session,
),
patch(
"local_deep_research.settings.manager.SettingsManager",
return_value=settings_manager_instance,
),
# autospec=True preserves the descriptor so the spy is invoked as a
# bound method (``self`` is the implicit first arg).
patch.object(
BackgroundJobScheduler,
"_arm_egress_backstop",
autospec=True,
) as arm_spy,
patch(DEFAULT_LIB, return_value="default-lib"),
patch(FACTORY) as mock_factory,
):
scheduler._reconcile_unindexed_documents("testuser")
# The backstop was armed exactly once with the in-body SettingsManager and
# the username (``self`` is the scheduler since autospec keeps it bound).
arm_spy.assert_called_once_with(
scheduler, settings_manager_instance, "testuser"
)
# No work to do, so no RAG service was built.
mock_factory.assert_not_called()