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.
This commit is contained in:
faligam
2026-09-07 14:04:48 +02:00
parent 610c716a97
commit 56d6415e68
5 changed files with 109 additions and 20 deletions

View File

@@ -7,6 +7,32 @@
## Aktueller Stand — 2026-09-07 ## 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 ### Evidenzbewertung und sichtbarer Fallback-Bericht in die Pipeline integriert — 2026-09-07
Ein Browser-Run hatte bereits abgerufene Quellen, aber keine Evidenz, keinen Ein Browser-Run hatte bereits abgerufene Quellen, aber keine Evidenz, keinen

View File

@@ -193,6 +193,7 @@ class SourceListItem(BaseModel):
domain: str = "" domain: str = ""
source_type: str | None = None source_type: str | None = None
retrieved_at: str = "" retrieved_at: str = ""
error: str = ""
class SourcesResponse(BaseModel): class SourcesResponse(BaseModel):
@@ -585,6 +586,10 @@ async def _run_pipeline(run: _ResearchRunState, request: ResearchRequest, budget
run.state = ResearchRunState.FAILED.value run.state = ResearchRunState.FAILED.value
run.metadata["error"] = result.get("error", "Pipeline failed") run.metadata["error"] = result.get("error", "Pipeline failed")
run.metadata["failed_step"] = result.get("failed_step", "unknown") 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) metrics.increment(C_RESEARCH_FAILED_TOTAL)
logger.warning( logger.warning(
"[pipeline] run_id=%s failed: %s", "[pipeline] run_id=%s failed: %s",
@@ -702,6 +707,7 @@ async def get_sources(research_id: str) -> SourcesResponse:
domain=s.get("domain", ""), domain=s.get("domain", ""),
source_type=s.get("source_type"), source_type=s.get("source_type"),
retrieved_at=s.get("retrieved_at", _now()), retrieved_at=s.get("retrieved_at", _now()),
error=s.get("error", ""),
) )
for s in run.sources for s in run.sources
] ]

View File

@@ -77,7 +77,14 @@ class CrawlerManager:
if result.status in (FetchStatus.BLOCKED, FetchStatus.ERROR, FetchStatus.TIMEOUT, if result.status in (FetchStatus.BLOCKED, FetchStatus.ERROR, FetchStatus.TIMEOUT,
FetchStatus.SIZE_LIMIT_EXCEEDED, FetchStatus.CONTENT_TYPE_BLOCKED): 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 = result.content
content_type = result.content_type content_type = result.content_type
@@ -165,7 +172,11 @@ class CrawlerManager:
# ------------------------------------------------------------------ # ------------------------------------------------------------------
@staticmethod @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.""" """Create a NormalizedDocument representing an error."""
# Keep the error result structurally identical to a successfully # Keep the error result structurally identical to a successfully
# normalized document. ``from_text`` hashes the (empty) content, so # normalized document. ``from_text`` hashes the (empty) content, so
@@ -178,6 +189,7 @@ class CrawlerManager:
metadata={ metadata={
"error": error, "error": error,
"content_type": "error", "content_type": "error",
"download_bytes": download_bytes,
}, },
extraction_tool="", extraction_tool="",
) )

View File

@@ -468,6 +468,8 @@ class ResearchOrchestrator:
"state": self._state_machine.current_state.value, "state": self._state_machine.current_state.value,
"report": None, "report": None,
"plan": self._plan, "plan": self._plan,
"sources": self._sources,
"claims": [claim.model_dump(mode="json") for claim in self._claims],
} }
# Zeit-Tracking # Zeit-Tracking
@@ -488,6 +490,8 @@ class ResearchOrchestrator:
"state": self._state_machine.current_state.value, "state": self._state_machine.current_state.value,
"report": None, "report": None,
"plan": self._plan, "plan": self._plan,
"sources": self._sources,
"claims": [claim.model_dump(mode="json") for claim in self._claims],
} }
# Alle Schritte erfolgreich → COMPLETED # Alle Schritte erfolgreich → COMPLETED
@@ -721,21 +725,24 @@ class ResearchOrchestrator:
"data": {}, "data": {},
} }
remaining_download_bytes = self._budget_config.max_total_download_bytes - int( # Fetch sources one by one with the *remaining* byte budget. A
usage["max_total_download_bytes"]["usage"] # former equal split assigned a quick run's 500 KB across 15 URLs
) # (about 33 KB each), turning otherwise readable web pages into
max_size_per_url = 0 # size-limit errors before extraction could start.
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
self._sources = [] 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 = { src = {
"id": str(uuid4()), "id": str(uuid4()),
"url": doc.url or "", "url": doc.url or "",
@@ -752,9 +759,22 @@ class ResearchOrchestrator:
downloaded_bytes = int(doc.metadata.get("download_bytes", len(doc.text.encode("utf-8")))) 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.increment_download_bytes(downloaded_bytes)
self._budget_tracker.record_time_elapsed()
successful = [s for s in self._sources if not s.get("error")] successful = [s for s in self._sources if not s.get("error")]
metrics.increment(C_SOURCES_FETCHED_TOTAL, len(self._sources)) metrics.increment(C_SOURCES_FETCHED_TOTAL, len(self._sources))
logger.info("Fetching complete: %d/%d sources successfully fetched", len(successful), 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 { return {
"success": True, "success": True,
"data": {"sources": self._sources}, "data": {"sources": self._sources},

View File

@@ -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} "https://example.org/two", "two", metadata={"download_bytes": 180}
) )
crawler = Mock() 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 orchestrator._crawler = crawler
response = await orchestrator._step_fetching() response = await orchestrator._step_fetching()
assert response["success"] is True assert response["success"] is True
crawler.fetch_and_extract_many.assert_awaited_once_with( assert crawler.fetch_and_extract_many.await_args_list[0].args == (["https://example.org/one"],)
["https://example.org/one", "https://example.org/two"], max_size=500 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() usage = orchestrator.budget_tracker.get_usage()
assert usage["max_sources"]["usage"] == 2 assert usage["max_sources"]["usage"] == 2
assert usage["max_total_download_bytes"]["usage"] == 300 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()) 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: def test_hard_budget_allows_exactly_the_configured_limit() -> None:
tracker = BudgetTracker(HardBudgetConfig(max_sources=2)) tracker = BudgetTracker(HardBudgetConfig(max_sources=2))
tracker.increment_sources(2) tracker.increment_sources(2)