From 56d6415e68b67ec2543f08c12d759a5d886e1d04 Mon Sep 17 00:00:00 2001 From: faligam Date: Mon, 7 Sep 2026 14:04:48 +0200 Subject: [PATCH] fix(crawler): prioritize usable source retrieval over equal byte splits Fetch candidate pages sequentially with the full remaining download budget instead of dividing a quick run's 500 KB equally across every URL. This avoids aborting otherwise readable pages at roughly 33 KB before claim extraction can inspect them. Account for bytes consumed by failed downloads, preserve fetch errors in source responses, and fail runs with no successfully extracted document using an actionable error. Keep attempted sources on failed runs for diagnostics and cover the allocation and no-usable-source paths with regression tests. --- HANDOFF.md | 26 ++++++++++++++ src/nsct/api/rest_research.py | 6 ++++ src/nsct/crawler/manager.py | 16 +++++++-- src/nsct/orchestration/orchestrator.py | 48 ++++++++++++++++++-------- tests/test_backend_pipeline_repairs.py | 33 +++++++++++++++--- 5 files changed, 109 insertions(+), 20 deletions(-) diff --git a/HANDOFF.md b/HANDOFF.md index d16795e..b659311 100644 --- a/HANDOFF.md +++ b/HANDOFF.md @@ -7,6 +7,32 @@ ## Aktueller Stand — 2026-09-07 +### Claim-Ausfall durch zu kleine Seitenlimits repariert — 2026-09-07 + +Der Browser-Run `ab7d8263-cb6a-4db0-97da-4bee88408711` zeigte Quellen, aber +keine Claims. Das API-Log bestätigte: `No successful sources ... for claim +extraction`. Die angezeigten Quellen waren Fehlerdokumente ohne extrahierten +Inhalt, nicht verwendbare Dokumente. Das war kein CAPTCHA- oder +Erreichbarkeitsproblem: Browser konnten die URLs laden. + +Ursache: Ein Quick-Run teilt sein 500-KB-Downloadbudget durch bis zu 15 URLs +und setzte dadurch ein hartes Limit von nur ca. 33 KB pro Seite. Der Crawler +markiert größere, sonst lesbare HTML-Seiten als `SIZE_LIMIT_EXCEEDED`; die +frühere Quellen-API verbarg diesen Fehler zudem. + +- Quellen werden nun nacheinander mit dem jeweils gesamten verbleibenden + Bytebudget abgerufen. Damit erhält die erste brauchbare Seite genügend + Platz für die Extraktion; tatsächlich abgerufene Bytes werden auch bei + Fehlern budgetiert. +- Wenn keine Quelle erfolgreich extrahiert wird, endet der Run eindeutig als + fehlgeschlagen und nennt die konkreten Abruffehler. +- Der Quellen-Endpunkt liefert das Fehlerfeld mit, damit die Oberfläche keine + Fehlerdokumente mehr als uneingeschränkt nutzbare Quellen ausgibt. + +Validierung: 50 gezielte Crawler- und Pipeline-Tests, Syntaxprüfung und +Diff-Check erfolgreich. Nach Deployment eine neue Recherche starten; der +vorherige In-Memory-Run kann nicht nachträglich mit Inhalten ergänzt werden. + ### Evidenzbewertung und sichtbarer Fallback-Bericht in die Pipeline integriert — 2026-09-07 Ein Browser-Run hatte bereits abgerufene Quellen, aber keine Evidenz, keinen diff --git a/src/nsct/api/rest_research.py b/src/nsct/api/rest_research.py index 4a7fd6f..5e14bed 100644 --- a/src/nsct/api/rest_research.py +++ b/src/nsct/api/rest_research.py @@ -193,6 +193,7 @@ class SourceListItem(BaseModel): domain: str = "" source_type: str | None = None retrieved_at: str = "" + error: str = "" class SourcesResponse(BaseModel): @@ -585,6 +586,10 @@ async def _run_pipeline(run: _ResearchRunState, request: ResearchRequest, budget run.state = ResearchRunState.FAILED.value run.metadata["error"] = result.get("error", "Pipeline failed") run.metadata["failed_step"] = result.get("failed_step", "unknown") + # Preserve attempted sources on a fetch failure. Their public + # error fields explain why no claims could be produced. + run.sources = result.get("sources", []) + run.claims = result.get("claims", []) metrics.increment(C_RESEARCH_FAILED_TOTAL) logger.warning( "[pipeline] run_id=%s failed: %s", @@ -702,6 +707,7 @@ async def get_sources(research_id: str) -> SourcesResponse: domain=s.get("domain", ""), source_type=s.get("source_type"), retrieved_at=s.get("retrieved_at", _now()), + error=s.get("error", ""), ) for s in run.sources ] diff --git a/src/nsct/crawler/manager.py b/src/nsct/crawler/manager.py index 9ba10ea..72a1e94 100644 --- a/src/nsct/crawler/manager.py +++ b/src/nsct/crawler/manager.py @@ -77,7 +77,14 @@ class CrawlerManager: if result.status in (FetchStatus.BLOCKED, FetchStatus.ERROR, FetchStatus.TIMEOUT, FetchStatus.SIZE_LIMIT_EXCEEDED, FetchStatus.CONTENT_TYPE_BLOCKED): - return self._error_doc(url, error=result.error or "Fetch failed") + # A failed response can still have consumed bytes (notably a page + # exceeding its cap). Retain that fact so the orchestrator cannot + # silently spend the download budget again on later URLs. + return self._error_doc( + url, + error=result.error or "Fetch failed", + download_bytes=len(result.content), + ) content = result.content content_type = result.content_type @@ -165,7 +172,11 @@ class CrawlerManager: # ------------------------------------------------------------------ @staticmethod - def _error_doc(url: str, error: str) -> NormalizedDocument: + def _error_doc( + url: str, + error: str, + download_bytes: int = 0, + ) -> NormalizedDocument: """Create a NormalizedDocument representing an error.""" # Keep the error result structurally identical to a successfully # normalized document. ``from_text`` hashes the (empty) content, so @@ -178,6 +189,7 @@ class CrawlerManager: metadata={ "error": error, "content_type": "error", + "download_bytes": download_bytes, }, extraction_tool="", ) diff --git a/src/nsct/orchestration/orchestrator.py b/src/nsct/orchestration/orchestrator.py index 275fa46..127a988 100644 --- a/src/nsct/orchestration/orchestrator.py +++ b/src/nsct/orchestration/orchestrator.py @@ -468,6 +468,8 @@ class ResearchOrchestrator: "state": self._state_machine.current_state.value, "report": None, "plan": self._plan, + "sources": self._sources, + "claims": [claim.model_dump(mode="json") for claim in self._claims], } # Zeit-Tracking @@ -488,6 +490,8 @@ class ResearchOrchestrator: "state": self._state_machine.current_state.value, "report": None, "plan": self._plan, + "sources": self._sources, + "claims": [claim.model_dump(mode="json") for claim in self._claims], } # Alle Schritte erfolgreich → COMPLETED @@ -721,21 +725,24 @@ class ResearchOrchestrator: "data": {}, } - remaining_download_bytes = self._budget_config.max_total_download_bytes - int( - usage["max_total_download_bytes"]["usage"] - ) - max_size_per_url = 0 - if self._budget_config.max_total_download_bytes > 0: - max_size_per_url = max(1, remaining_download_bytes // len(urls)) - - self._budget_tracker.increment_sources(len(urls)) - - docs = await crawler.fetch_and_extract_many(urls, max_size=max_size_per_url) - self._budget_tracker.record_time_elapsed() - - # In sources-Dicts umwandeln + # Fetch sources one by one with the *remaining* byte budget. A + # former equal split assigned a quick run's 500 KB across 15 URLs + # (about 33 KB each), turning otherwise readable web pages into + # size-limit errors before extraction could start. self._sources = [] - for i, doc in enumerate(docs): + for url in urls: + usage = self._budget_tracker.get_usage() + remaining_download_bytes = self._budget_config.max_total_download_bytes - int( + usage["max_total_download_bytes"]["usage"] + ) + if self._budget_config.max_total_download_bytes > 0 and remaining_download_bytes <= 0: + logger.info("Download budget exhausted after %d source attempts", len(self._sources)) + break + + self._budget_tracker.increment_sources() + max_size = remaining_download_bytes if self._budget_config.max_total_download_bytes > 0 else 0 + docs = await crawler.fetch_and_extract_many([url], max_size=max_size) + doc = docs[0] src = { "id": str(uuid4()), "url": doc.url or "", @@ -752,9 +759,22 @@ class ResearchOrchestrator: downloaded_bytes = int(doc.metadata.get("download_bytes", len(doc.text.encode("utf-8")))) self._budget_tracker.increment_download_bytes(downloaded_bytes) + self._budget_tracker.record_time_elapsed() + successful = [s for s in self._sources if not s.get("error")] metrics.increment(C_SOURCES_FETCHED_TOTAL, len(self._sources)) logger.info("Fetching complete: %d/%d sources successfully fetched", len(successful), len(self._sources)) + if not successful: + errors = sorted({str(source["error"]) for source in self._sources if source.get("error")}) + detail = "; ".join(errors[:3]) or "no usable document content" + return { + "success": False, + "error": f"No source could be fetched and extracted: {detail}", + "data": {"sources": self._sources}, + "sources": self._sources, + "source_count": len(self._sources), + "successful_count": 0, + } return { "success": True, "data": {"sources": self._sources}, diff --git a/tests/test_backend_pipeline_repairs.py b/tests/test_backend_pipeline_repairs.py index d867e2a..33514be 100644 --- a/tests/test_backend_pipeline_repairs.py +++ b/tests/test_backend_pipeline_repairs.py @@ -163,15 +163,16 @@ def test_fetching_uses_actual_bytes_and_respects_remaining_source_budget() -> No "https://example.org/two", "two", metadata={"download_bytes": 180} ) crawler = Mock() - crawler.fetch_and_extract_many = AsyncMock(return_value=[doc_one, doc_two]) + crawler.fetch_and_extract_many = AsyncMock(side_effect=[[doc_one], [doc_two]]) orchestrator._crawler = crawler response = await orchestrator._step_fetching() assert response["success"] is True - crawler.fetch_and_extract_many.assert_awaited_once_with( - ["https://example.org/one", "https://example.org/two"], max_size=500 - ) + assert crawler.fetch_and_extract_many.await_args_list[0].args == (["https://example.org/one"],) + assert crawler.fetch_and_extract_many.await_args_list[0].kwargs == {"max_size": 1_000} + assert crawler.fetch_and_extract_many.await_args_list[1].args == (["https://example.org/two"],) + assert crawler.fetch_and_extract_many.await_args_list[1].kwargs == {"max_size": 880} usage = orchestrator.budget_tracker.get_usage() assert usage["max_sources"]["usage"] == 2 assert usage["max_total_download_bytes"]["usage"] == 300 @@ -179,6 +180,30 @@ def test_fetching_uses_actual_bytes_and_respects_remaining_source_budget() -> No asyncio.run(run()) +def test_fetching_fails_with_source_errors_when_no_document_is_usable() -> None: + async def run() -> None: + orchestrator = ResearchOrchestrator( + _config(), uuid4(), "test", budget_config=HardBudgetConfig(max_sources=2, max_total_download_bytes=1_000) + ) + orchestrator._search_results = [{"url": "https://example.org/blocked"}] + failed_doc = NormalizedDocument.from_text( + "https://example.org/blocked", + "", + metadata={"error": "Download exceeded 1000 bytes", "download_bytes": 1_000}, + ) + crawler = Mock() + crawler.fetch_and_extract_many = AsyncMock(return_value=[failed_doc]) + orchestrator._crawler = crawler + + response = await orchestrator._step_fetching() + + assert response["success"] is False + assert "Download exceeded 1000 bytes" in response["error"] + assert response["sources"][0]["error"] == "Download exceeded 1000 bytes" + + asyncio.run(run()) + + def test_hard_budget_allows_exactly_the_configured_limit() -> None: tracker = BudgetTracker(HardBudgetConfig(max_sources=2)) tracker.increment_sources(2)