From def4441de1e1849a9c04286f297d6e954c4bc823 Mon Sep 17 00:00:00 2001 From: Phil Merrell Date: Sun, 27 Sep 2026 10:22:54 -0600 Subject: [PATCH] fix(kb): re-ingest a synced source that changed on a managed knowledge base MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A KB sync that finds its Drive file (or crawled page) changed overwrites the document's S3 object in place and relies on the ObjectCreated event to re-ingest. On a managed knowledge base that event reaches the ingestion consumer, which found a complete row with byteCapSettled and returned "already-settled" — the early exit that makes redelivery idempotent. The new bytes were never ingested, and the size change never reached the byte ledger. The sync now writes stagedContentHash before staging. The consumer re-ingests while it differs from the ingestedContentHash the last re-ingest recorded: it reserves the size growth against the cap before submitting (a file that grew past the cap is refused and its previous version keeps serving; the next sync retries), submits once per version, waits out Bedrock statuses that predate the submission, and on completion re-stamps committedBytes and commits or refunds the difference. A per-version claim on the DOC# row keeps all of it idempotent under EventBridge redelivery; a delete mid-flight returns the growth reservation. Legacy knowledge bases are routed away before any of this, as before. Co-Authored-By: Claude Opus 5.5 --- .../documents/services/document_service.py | 15 + .../kb_migration/ingestion_consumer.py | 360 +++++++++++- backend/src/apis/app_api/kb_sync/records.py | 11 + backend/src/apis/app_api/kb_sync/worker.py | 9 +- .../src/apis/shared/kb_backend/byte_cap.py | 160 +++++- .../lambdas/test_kb_sync_managed_reingest.py | 536 ++++++++++++++++++ backend/tests/lambdas/test_kb_sync_worker.py | 5 + docs/specs/assistant-kb-sync.md | 12 + 8 files changed, 1104 insertions(+), 4 deletions(-) create mode 100644 backend/tests/lambdas/test_kb_sync_managed_reingest.py diff --git a/backend/src/apis/app_api/documents/services/document_service.py b/backend/src/apis/app_api/documents/services/document_service.py index 7f6113f09..094e0e0d7 100644 --- a/backend/src/apis/app_api/documents/services/document_service.py +++ b/backend/src/apis/app_api/documents/services/document_service.py @@ -151,10 +151,25 @@ async def settle_bytes_on_delete(document: Document, previous_status: Optional[s Refund runs first. A row can hold both markers only because a consumer settled it, and then it has nothing left to release. + + Independently, a document deleted while a KB sync's changed source was being + re-ingested holds the growth reserved for that version + (``byte_cap.claim_reingest``). The consumer's completion is refused on a + ``deleting`` row, so this is the only path left to return it. """ from apis.shared.kb_backend import byte_cap assistant_id = document.assistant_id + try: + reingest_reserved = byte_cap.release_reingest_once(assistant_id, document.document_id) + if reingest_reserved: + byte_cap.release(assistant_id, assistant_id, reingest_reserved) + except Exception as e: # noqa: BLE001 - a bookkeeping failure must not break the delete + logger.error( + f"Failed to release re-ingest reservation for document {document.document_id}: {e}", + exc_info=True, + ) + try: refunded = byte_cap.refund_once(assistant_id, document.document_id) if refunded: diff --git a/backend/src/apis/app_api/kb_migration/ingestion_consumer.py b/backend/src/apis/app_api/kb_migration/ingestion_consumer.py index 384a21268..29656ed19 100644 --- a/backend/src/apis/app_api/kb_migration/ingestion_consumer.py +++ b/backend/src/apis/app_api/kb_migration/ingestion_consumer.py @@ -35,6 +35,17 @@ and ``retrievableAt`` separately so the gap stays measurable instead of becoming folklore. +A changed source is not a redelivery +------------------------------------ +A KB sync (Drive file or web re-crawl) that finds its source changed overwrites the +document's S3 object in place, so its ``ObjectCreated`` event arrives for a document +that is already ``complete`` with its bytes settled — exactly what a redelivery looks +like. The sync marks the difference: it writes ``stagedContentHash`` before the +overwrite, and :func:`staged_version_to_reingest` compares that with the +``ingestedContentHash`` the last re-ingest recorded. A version still owed is +re-ingested and its size difference settled against the byte cap +(:func:`_reingest_changed_document`); anything else keeps the settled early exit. + Import boundary --------------- Raw DynamoDB table access rather than importing ``apis.shared.assistants``, whose @@ -413,6 +424,7 @@ def wait_until_indexed( timeout_seconds: Optional[float] = None, interval_seconds: Optional[float] = None, sleep: Any = time.sleep, + not_before: Optional[str] = None, ) -> Tuple[str, Optional[str]]: """Poll Bedrock's document status until it settles, or give up. @@ -420,6 +432,10 @@ def wait_until_indexed( outcome, not an error: the caller leaves the document non-terminal and lets redelivery come back to it, by which time indexing has usually finished. + ``not_before`` is for a document Bedrock already holds a version of: a status + last updated at or before it describes that previous version, not the one just + submitted, and is waited out like an in-flight one (:func:`_is_stale`). + Why a *bounded* in-invocation wait rather than pure redelivery: most documents index in a few seconds, and making every one of them wait for an EventBridge retry would add a minute of latency to the common case for the sake of the rare @@ -438,12 +454,30 @@ def wait_until_indexed( deadline = time.monotonic() + timeout_seconds status, updated_at = document_status(backend, kb_ref, document_id) - while _still_working(status, document_id) and time.monotonic() < deadline: + while ( + _still_working(status, document_id) or _is_stale(updated_at, not_before) + ) and time.monotonic() < deadline: sleep(interval_seconds) status, updated_at = document_status(backend, kb_ref, document_id) return status, updated_at +def _is_stale(updated_at: Optional[str], not_before: Optional[str]) -> bool: + """Whether a Bedrock status predates ``not_before`` and so is not about it. + + A status with no timestamp cannot be placed and is taken at face value — the + same answer the consumer gave before this check existed. + """ + if not updated_at or not not_before: + return False + from apis.shared.timestamps import from_iso + + try: + return from_iso(updated_at) <= from_iso(not_before) + except (TypeError, ValueError): + return False + + def _still_working(status: str, document_id: str) -> bool: """Whether to keep waiting on ``status``. @@ -706,6 +740,321 @@ def _commit_settled(assistant_id: str, document_id: str, real: int, reserved: in byte_cap.release(assistant_id, assistant_id, reserved - real) +def staged_version_to_reingest(doc_row: Optional[Dict[str, Any]]) -> Optional[str]: + """The source version a KB sync staged over this document and it has yet to + ingest, or ``None``. + + A KB sync that finds its source changed writes ``stagedContentHash`` BEFORE it + overwrites the S3 object (``kb_sync.records.update_document_sync_fields``), + and a completed re-ingest stamps the same hash as ``ingestedContentHash``. A + difference between the two is therefore a version still owed to the + knowledge base, and an event that finds them equal is a redelivery. Nothing + else writes the staged hash, so a document never synced has none and keeps the + settled early exit exactly as before. + + Only a terminal document qualifies. One still being ingested for the first + time is handled by the ordinary path, which reads the object as it is now. + ``failed`` qualifies too: new bytes deserve a fresh attempt, and a document + that failed on the byte cap or in Bedrock is otherwise stuck for good. + """ + if not doc_row or doc_row.get("status") not in (STATUS_COMPLETE, STATUS_FAILED): + return None + staged = doc_row.get("stagedContentHash") + if not staged or staged == doc_row.get("ingestedContentHash"): + return None + return str(staged) + + +def _reingest_changed_document( + bucket: str, + key: str, + assistant_id: str, + document_id: str, + filename: str, + doc_row: Dict[str, Any], + record: Optional[Dict[str, Any]], + staged: str, +) -> Dict[str, Any]: + """Re-ingest a settled document whose source a KB sync replaced in place. + + The first ingestion's path cannot do this. Its byte accounting is claimed once + per document (``settle_once``), which is what makes a redelivery harmless, and + Bedrock already reports the document ``INDEXED`` — so that path would neither + submit the new bytes nor count them. Here each staged version is handled once + instead, keyed on its hash (``byte_cap.claim_reingest``): + + 1. **Reserve the growth, then claim.** The S3 size is measured now and only + the difference from ``committedBytes`` is reserved, against the cap, before + anything is submitted. A file that grew past the cap is refused while the + previous version is still indexed and still served + (:func:`_record_reingest_over_cap`), rather than replacing it and then + failing the document. A redelivery finds the claim and reserves nothing. + 2. **Submit once.** ``reingestSubmittedAt`` records the submission, so a + redelivery waits on the ingestion already running instead of restarting it. + 3. **Wait for the NEW version.** Bedrock already holds a document under this + id, so a status it last updated before the submission is the old version's + and is waited out (``not_before``). + 4. **Complete.** One conditional write re-stamps ``committedBytes``, records + the version as ingested and drops the claim; the ledger then moves by the + difference (``byte_cap.settle_reingest``). + + The status stays ``complete`` throughout, so the document keeps answering from + its previous version until the new one has replaced it. + """ + import asyncio + + from apis.shared.kb_backend import byte_cap + from apis.shared.kb_backend.managed_backend import ManagedKbBackend + from apis.shared.kb_backend.protocol import DocumentSource + + summary: Dict[str, Any] = {"routed": "managed", "document_id": document_id} + + if doc_row.get("reingestHash") != staged: + try: + _claim_reingest(assistant_id, document_id, bucket, key, doc_row, record, staged) + except byte_cap.ByteCapExceeded: + _record_reingest_over_cap(assistant_id, document_id, staged) + return {**summary, "ingested": False, "note": "byte-cap-exceeded"} + # Re-read: the claim may have been won by a concurrent delivery of this + # version, or lost to a delete or a newer version. + doc_row = _get_doc_row(assistant_id, document_id) or {} + if doc_row.get("reingestHash") != staged or doc_row.get("status") == STATUS_DELETING: + logger.info( + f"document {document_id} was deleted or re-staged before its changed " + f"source could be re-ingested; leaving it to that change" + ) + return {**summary, "ingested": False, "note": "reingest-superseded"} + + backend = ManagedKbBackend(bucket=bucket) + submitted_at = doc_row.get("reingestSubmittedAt") + if submitted_at: + logger.info( + f"document {document_id}'s changed source was already submitted at " + f"{submitted_at}; waiting on that ingestion rather than restarting it" + ) + else: + submitted_at = _now_iso() + source = DocumentSource(document_id=document_id, filename=filename, s3_key=key) + # A submit that raises leaves the claim and its reservation in place for + # redelivery. The document is not failed: its previous version is intact. + asyncio.run(backend.ingest(assistant_id, source)) + _mark_reingest_submitted(assistant_id, document_id, staged, submitted_at) + + status, bedrock_updated_at = wait_until_indexed( + backend, assistant_id, document_id, not_before=submitted_at + ) + stale = _is_stale(bedrock_updated_at, submitted_at) + + if status in DOC_STATUSES_FAILED and not stale: + logger.error(f"re-ingestion of document {document_id} became {status}") + released = byte_cap.release_reingest_once(assistant_id, document_id, staged) + if released: + byte_cap.release(assistant_id, assistant_id, released) + set_document_terminal( + assistant_id, document_id, STATUS_FAILED, + error=f"the knowledge base reports this document as {status}", + ) + return {**summary, "ingested": True, "status": status} + + if stale or status not in (DOC_STATUS_INDEXED, *DOC_STATUSES_PARTIAL): + raise IngestionRoutingError( + f"document {document_id}'s changed source is {status} after waiting; " + f"leaving it for redelivery to confirm indexing" + ) + + indexed_at = bedrock_updated_at or _now_iso() + retrievable_at = wait_until_retrievable(backend, assistant_id, document_id) + if retrievable_at is None: + raise IngestionRoutingError( + f"document {document_id}'s changed source is INDEXED but was not " + f"retrievable within the poll window; leaving it for redelivery" + ) + + if not _complete_reingest( + assistant_id, document_id, staged, doc_row, indexed_at, retrievable_at + ): + return {**summary, "ingested": True, "note": "reingest-superseded"} + return { + **summary, + "ingested": True, + "note": "reingested", + "indexedAt": indexed_at, + "retrievableAt": retrievable_at, + } + + +def _claim_reingest( + assistant_id: str, + document_id: str, + bucket: str, + key: str, + doc_row: Dict[str, Any], + record: Optional[Dict[str, Any]], + staged: str, +) -> None: + """Reserve a staged version's growth and claim its re-ingest (step 1). + + Raises ``ByteCapExceeded`` with nothing reserved or claimed. Reserve comes + before the claim so a crash between the two leaks a reservation — the safe + direction — rather than leaving a claim that commits bytes nobody reserved. + """ + from apis.shared.kb_backend import byte_cap + + # A claim for an older version never completed; this change supersedes it. + # Return its reservation first, so it cannot count against this one's cap. + stale_claim = doc_row.get("reingestHash") + if stale_claim: + stale = byte_cap.release_reingest_once(assistant_id, document_id, str(stale_claim)) + if stale: + byte_cap.release(assistant_id, assistant_id, stale) + + real = byte_cap.object_size_bytes(bucket, key) + growth = max(real - int(doc_row.get("committedBytes") or 0), 0) + if growth: + elevated = bool((record or {}).get("elevatedByteCap")) + byte_cap.reserve(assistant_id, assistant_id, growth, byte_cap.effective_cap(elevated)) + + claimed, superseded = byte_cap.claim_reingest(assistant_id, document_id, staged, real, growth) + if not claimed: + # A concurrent delivery of this version claimed it first, or the row left + # a re-ingestable state. Either way this reservation is not needed. + if growth: + byte_cap.release(assistant_id, assistant_id, growth) + return + if superseded: + byte_cap.release(assistant_id, assistant_id, superseded) + + +def _mark_reingest_submitted( + assistant_id: str, document_id: str, staged: str, submitted_at: str +) -> None: + """Record that this version was submitted (step 2). Best-effort: without it a + redelivery submits again, which restarts indexing but loses nothing.""" + from botocore.exceptions import ClientError + + try: + _table().update_item( + Key={"PK": f"AST#{assistant_id}", "SK": f"DOC#{document_id}"}, + UpdateExpression="SET reingestSubmittedAt = :t", + ConditionExpression="reingestHash = :h", + ExpressionAttributeValues={":t": submitted_at, ":h": staged}, + ) + except ClientError as exc: + if exc.response.get("Error", {}).get("Code") != "ConditionalCheckFailedException": + raise + + +def _complete_reingest( + assistant_id: str, + document_id: str, + staged: str, + doc_row: Dict[str, Any], + indexed_at: str, + retrievable_at: str, +) -> bool: + """Record the re-ingested version and move the ledger by its size (step 4). + + The ``DOC#`` write comes first and carries everything — the new + ``committedBytes``, the ingested hash, the status, and the claim's removal — + conditioned on the claim still being this version's and the row not being + deleted. That is the re-ingest's :func:`byte_cap.record_commit`: refused means + a delete or a newer version owns the claim's reservation now, and this path + must not touch the ledger. ``byteCapSettled`` is set too, for a ``failed`` row + that never settled, so a later delete refunds rather than releases. + """ + from botocore.exceptions import ClientError + + from apis.shared.kb_backend import byte_cap + from apis.shared.kb_backend.records import RECORD_EXISTS + + new_bytes = int(doc_row.get("reingestBytes") or 0) + reserved = int(doc_row.get("reingestReservedBytes") or 0) + try: + response = _table().update_item( + Key={"PK": f"AST#{assistant_id}", "SK": f"DOC#{document_id}"}, + UpdateExpression=( + "SET #status = :complete, updatedAt = :now, indexedAt = :indexed, " + "retrievableAt = :retrievable, committedBytes = :bytes, " + "byteCapSettled = :true, ingestedContentHash = :h " + f"REMOVE {', '.join(byte_cap.REINGEST_CLAIM_ATTRIBUTES)}, ingestionError" + ), + ConditionExpression=( + f"{RECORD_EXISTS} AND #status <> :deleting " + f"AND reingestHash = :h AND attribute_not_exists(byteCapRefunded)" + ), + ExpressionAttributeNames={"#status": "status"}, + ExpressionAttributeValues={ + ":complete": STATUS_COMPLETE, + ":deleting": STATUS_DELETING, + ":now": _now_iso(), + ":indexed": indexed_at, + ":retrievable": retrievable_at, + ":bytes": new_bytes, + ":true": True, + ":h": staged, + }, + ReturnValues="ALL_OLD", + ) + except ClientError as exc: + if exc.response.get("Error", {}).get("Code") == "ConditionalCheckFailedException": + logger.info( + f"document {document_id} was deleted or re-staged while its changed " + f"source was indexing; not recording the re-ingest" + ) + return False + raise + + previous = int((response.get("Attributes") or {}).get("committedBytes") or 0) + byte_cap.settle_reingest(assistant_id, assistant_id, previous, new_bytes, reserved) + logger.info( + f"document {document_id} re-ingested from its changed source: " + f"{previous} -> {new_bytes} bytes" + ) + return True + + +def _record_reingest_over_cap(assistant_id: str, document_id: str, staged: str) -> None: + """Refuse a changed source that would take the owner over the byte cap. + + Nothing was submitted, so the knowledge base keeps serving the previous + version and the document stays as it was. The sync's change-detection gates + (``sourceEtag``, ``contentHash``) are cleared so the next sync run stages the + source again — the only thing that would retry it once the owner has freed + space. Conditioned on the version, so a newer change is left alone. + """ + from botocore.exceptions import ClientError + + from apis.shared.kb_backend.records import RECORD_EXISTS + + logger.warning( + f"document {document_id}'s changed source would exceed the byte cap; keeping " + f"its previous version and retrying on the next sync" + ) + try: + _table().update_item( + Key={"PK": f"AST#{assistant_id}", "SK": f"DOC#{document_id}"}, + UpdateExpression=( + "SET ingestionError = :err, updatedAt = :now REMOVE sourceEtag, contentHash" + ), + ConditionExpression=( + f"{RECORD_EXISTS} AND #status <> :deleting AND stagedContentHash = :h" + ), + ExpressionAttributeNames={"#status": "status"}, + ExpressionAttributeValues={ + ":err": ( + "the updated source exceeds the knowledge base's storage limit; " + "delete unused documents or request an elevated storage tier" + ), + ":now": _now_iso(), + ":deleting": STATUS_DELETING, + ":h": staged, + }, + ) + except ClientError as exc: + if exc.response.get("Error", {}).get("Code") != "ConditionalCheckFailedException": + raise + + def handle_object(bucket: str, key: str) -> Dict[str, Any]: """Route one uploaded object. Returns a summary for logging and tests.""" from apis.shared.kb_backend.records import BORN_MANAGED, ENGINE_MANAGED @@ -779,6 +1128,15 @@ def handle_object(bucket: str, key: str) -> Dict[str, Any]: # commit drives reservedBytes negative, a second release over-credits the cap. # Return without touching anything (Requirement 12.4/12.5). doc_row = _get_doc_row(assistant_id, document_id) + staged = staged_version_to_reingest(doc_row) + if staged: + # ...unless a KB sync has since overwritten the object with a changed + # source. Then this event is the new version's, not a redelivery, and the + # settled early exit below would leave the knowledge base serving the old + # content for good. + return _reingest_changed_document( + bucket, key, assistant_id, document_id, filename, doc_row, record, staged + ) if ( doc_row and doc_row.get("byteCapSettled") diff --git a/backend/src/apis/app_api/kb_sync/records.py b/backend/src/apis/app_api/kb_sync/records.py index 4faff18ae..64a865dfc 100644 --- a/backend/src/apis/app_api/kb_sync/records.py +++ b/backend/src/apis/app_api/kb_sync/records.py @@ -81,11 +81,19 @@ def update_document_sync_fields( previous_chunk_count: Optional[int] = None, last_synced_at: Optional[str] = None, sync_policy_id: Optional[str] = None, + staged_content_hash: Optional[str] = None, ) -> None: """Targeted update of the sync-bookkeeping fields on a document record. Only sets the fields passed — safe alongside the ingestion pipeline's own targeted UpdateExpressions (which never touch these attributes). + + ``staged_content_hash`` is for a sync that is about to overwrite the + document's S3 object with changed bytes, and must be written BEFORE the + overwrite. A managed knowledge base's ingestion consumer reads it to tell + the overwrite's event from a redelivery of the original upload's + (``kb_migration.ingestion_consumer.staged_version_to_reingest``); without + it the changed bytes are never re-ingested. The legacy pipeline ignores it. """ set_parts = [] values: Dict[str, Any] = {} @@ -104,6 +112,9 @@ def update_document_sync_fields( if sync_policy_id is not None: set_parts.append("syncPolicyId = :spid") values[":spid"] = sync_policy_id + if staged_content_hash is not None: + set_parts.append("stagedContentHash = :staged") + values[":staged"] = staged_content_hash if not set_parts: return diff --git a/backend/src/apis/app_api/kb_sync/worker.py b/backend/src/apis/app_api/kb_sync/worker.py index 81428e852..a615e3650 100644 --- a/backend/src/apis/app_api/kb_sync/worker.py +++ b/backend/src/apis/app_api/kb_sync/worker.py @@ -208,7 +208,9 @@ async def _sync_drive_file(policy: SyncPolicy) -> Dict[str, Any]: return await _finish(policy, "unchanged") # Changed: stash the old chunk count for the ingestion tail-delete - # (shrinkage cleanup) BEFORE staging, then overwrite the S3 object. + # (shrinkage cleanup) and mark the version being staged — what a managed + # KB's consumer uses to tell this overwrite from a redelivery — BEFORE + # staging, then overwrite the S3 object. previous_chunk_count = int(document.get("chunkCount") or 0) records.update_document_sync_fields( assistant_id, @@ -217,6 +219,7 @@ async def _sync_drive_file(policy: SyncPolicy) -> Dict[str, Any]: content_hash=content_hash, previous_chunk_count=previous_chunk_count, last_synced_at=_now_timestamp(), + staged_content_hash=content_hash, ) _stage_to_s3(document["s3Key"], downloaded.content, downloaded.content_type) logger.info( @@ -299,7 +302,8 @@ async def _sync_web_crawl(policy: SyncPolicy) -> Dict[str, Any]: async def on_result(url: str, document_id: str, outcome: str, etag, content_hash) -> None: if outcome == "changed": # BEFORE the S3 overwrite: stash the previous chunk count for - # the ingestion tail-delete, alongside the new gate values. + # the ingestion tail-delete and the staged version for a managed + # KB's re-ingest, alongside the new gate values. records.update_document_sync_fields( assistant_id, document_id, @@ -307,6 +311,7 @@ async def on_result(url: str, document_id: str, outcome: str, etag, content_hash content_hash=content_hash, previous_chunk_count=int(web_docs[url].get("chunkCount") or 0), last_synced_at=now, + staged_content_hash=content_hash, ) elif outcome == "unchanged": records.update_document_sync_fields( diff --git a/backend/src/apis/shared/kb_backend/byte_cap.py b/backend/src/apis/shared/kb_backend/byte_cap.py index 92df7c895..adf23644f 100644 --- a/backend/src/apis/shared/kb_backend/byte_cap.py +++ b/backend/src/apis/shared/kb_backend/byte_cap.py @@ -60,6 +60,17 @@ reservation instead of committing it. Both ``ADD`` writes commute, so a refund landing before its commit still nets to zero. +Re-ingesting a changed source +----------------------------- +A KB sync that finds its source changed overwrites the settled document's S3 +object in place, so the new version's size has to replace the old one in the +ledger. Only the difference moves: growth is reserved against the cap *before* +the new version is submitted (:func:`claim_reingest`), so a file that grew past +the cap is refused while its previous version is still being served, and on +completion ``committedBytes`` is re-stamped and the difference committed or +refunded (:func:`settle_reingest`). A delete mid-flight returns the growth +reservation through :func:`release_reingest_once`. + Guards against driving a counter negative ----------------------------------------- :func:`release` refuses to take ``reservedBytes`` below zero, and :func:`refund` @@ -88,7 +99,7 @@ import logging import os from decimal import Decimal -from typing import Optional +from typing import Any, Dict, Optional, Tuple from apis.shared.kb_backend.metrics import emit_count @@ -603,6 +614,153 @@ def refund_once(assistant_id: str, document_id: str) -> int: return int(response.get("Attributes", {}).get("committedBytes") or 0) +#: The ``DOC#`` attributes one in-flight re-ingest holds (:func:`claim_reingest`). +REINGEST_CLAIM_ATTRIBUTES = ( + "reingestHash", + "reingestBytes", + "reingestReservedBytes", + "reingestSubmittedAt", +) + + +def claim_reingest( + assistant_id: str, + document_id: str, + content_hash: str, + new_bytes: int, + reserved_bytes: int, +) -> Tuple[bool, int]: + """Claim the re-ingest of one staged version of an already-settled document. + + A KB sync that finds its source changed overwrites the document's S3 object in + place. The document already settled its bytes, so :func:`settle_once` cannot + account for the new version; this claim does. It stamps the version + (``reingestHash``, the ``stagedContentHash`` the sync wrote), its S3 size and + the growth the caller has just reserved for it, conditioned on no claim for + this version existing. That makes the reservation exactly-once per version + under redelivery: the caller reserves *before* claiming, and releases its own + reservation again when the claim is refused. + + Returns ``(claimed, superseded)``. ``superseded`` is the reservation of an + older version's claim this one replaced — a re-ingest that never completed — + and the caller must release it: nothing else ever will. Only one claim can + replace it, so it is released once. + + Refused for a row that is gone, ``deleting``, or not terminal: an upload still + in flight is its first ingestion's to settle. + """ + from botocore.exceptions import ClientError + + from apis.shared.kb_backend.records import RECORD_EXISTS + + try: + response = _table().update_item( + Key={"PK": f"AST#{assistant_id}", "SK": f"DOC#{document_id}"}, + UpdateExpression=( + "SET reingestHash = :h, reingestBytes = :b, reingestReservedBytes = :r " + "REMOVE reingestSubmittedAt" + ), + ConditionExpression=( + f"{RECORD_EXISTS} AND #status IN (:complete, :failed) " + f"AND (attribute_not_exists(reingestHash) OR reingestHash <> :h)" + ), + ExpressionAttributeNames={"#status": "status"}, + ExpressionAttributeValues={ + ":h": content_hash, + ":b": Decimal(new_bytes), + ":r": Decimal(reserved_bytes), + ":complete": "complete", + ":failed": "failed", + }, + ReturnValues="UPDATED_OLD", + ) + except ClientError as exc: + if exc.response.get("Error", {}).get("Code") == "ConditionalCheckFailedException": + return False, 0 + raise + old = response.get("Attributes") or {} + if not old.get("reingestHash"): + return True, 0 + return True, int(old.get("reingestReservedBytes") or 0) + + +def release_reingest_once( + assistant_id: str, + document_id: str, + content_hash: Optional[str] = None, +) -> int: + """Drop a document's re-ingest claim and return the bytes it had reserved. + + Returns the reservation for exactly one caller, 0 for everyone else. For a + re-ingest that ends without committing — Bedrock failed it, or the owner + deleted the document mid-flight — so its growth reservation does not leak. + ``content_hash`` restricts the drop to that version's claim. + + The completing path removes the claim in the same conditional write that + stamps the new ``committedBytes``, so a delete racing a completion either + finds the claim (and releases it) or finds the new stamp (and refunds it). + """ + from botocore.exceptions import ClientError + + from apis.shared.kb_backend.records import RECORD_EXISTS + + condition = f"{RECORD_EXISTS} AND attribute_exists(reingestHash)" + values: Dict[str, Any] = {} + if content_hash is not None: + condition += " AND reingestHash = :h" + values[":h"] = content_hash + kwargs: Dict[str, Any] = { + "Key": {"PK": f"AST#{assistant_id}", "SK": f"DOC#{document_id}"}, + "UpdateExpression": "REMOVE " + ", ".join(REINGEST_CLAIM_ATTRIBUTES), + "ConditionExpression": condition, + "ReturnValues": "ALL_OLD", + } + if values: + kwargs["ExpressionAttributeValues"] = values + try: + response = _table().update_item(**kwargs) + except ClientError as exc: + if exc.response.get("Error", {}).get("Code") == "ConditionalCheckFailedException": + return 0 + raise + return int((response.get("Attributes") or {}).get("reingestReservedBytes") or 0) + + +def settle_reingest( + assistant_id: str, + app_kb_id: str, + previous_bytes: int, + new_bytes: int, + reserved_bytes: int, +) -> None: + """Move the ledger from a document's old committed size to its re-ingested one. + + Called after the ``DOC#`` row's ``committedBytes`` has been re-stamped to + ``new_bytes`` — the stamp first, for the same reason :func:`record_commit` + comes before :func:`commit`. Growth was reserved when the re-ingest was + claimed, so it is committed out of that reservation; shrinkage is refunded. + + The growth can only differ from the reservation if ``committedBytes`` moved + between the claim's read and its write — an older version's re-ingest + completing in that window. Any reservation left over is released, and growth + beyond it is counted as stored without a cap check (:func:`add_stored`): the + bytes are already indexed, and refusing to count them would only hide them. + """ + growth = new_bytes - previous_bytes + committed = min(reserved_bytes, max(growth, 0)) + commit(assistant_id, app_kb_id, committed) + if reserved_bytes > committed: + release(assistant_id, app_kb_id, reserved_bytes - committed) + if growth > committed: + logger.warning( + f"re-ingest of a document in {assistant_id} grew {growth} bytes but reserved " + f"only {reserved_bytes}; counting the difference as stored" + ) + add_stored(assistant_id, app_kb_id, growth - committed) + elif growth < 0: + refund(assistant_id, app_kb_id, -growth) + + def release_snapshot(assistant_id: str, app_kb_id: str, n_bytes: int) -> bool: """Return a migration's whole-corpus reservation, once. diff --git a/backend/tests/lambdas/test_kb_sync_managed_reingest.py b/backend/tests/lambdas/test_kb_sync_managed_reingest.py new file mode 100644 index 000000000..76731c4a0 --- /dev/null +++ b/backend/tests/lambdas/test_kb_sync_managed_reingest.py @@ -0,0 +1,536 @@ +"""A synced source that changed must be re-ingested into its MANAGED knowledge base. + +The KB sync worker (``kb_sync/worker.py``) detects a changed Drive file and +overwrites the document's existing S3 object, relying on the bucket's +``ObjectCreated`` event to re-run ingestion. That works for the legacy pipeline. +For a managed knowledge base the event reaches the ingestion consumer, which found +a ``complete`` row with ``byteCapSettled`` and returned "already-settled" — the +same early exit that makes EventBridge redelivery idempotent. The new bytes were +never ingested, and the size change never reached the byte ledger. + +These tests drive the real path end to end against moto: the real sync worker +(Drive and the token vault stubbed at its seams) stages to a real S3 object, and +the real consumer handles the resulting event, against a fake Bedrock that models +document status the way the service reports it. +""" + +from __future__ import annotations + +import asyncio +from datetime import datetime, timezone +from decimal import Decimal +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock, patch + +import boto3 +import pytest + +from apis.app_api.documents.models import DocumentProvenance +from apis.app_api.documents.services.document_service import create_document, soft_delete_document +from apis.app_api.file_sources.models import DownloadedFile +from apis.app_api.kb_migration import ingestion_consumer as ic +from apis.app_api.kb_sync import worker +from apis.shared.assistants.service import create_assistant +from apis.shared.sync_policies.service import create_sync_policy + +REGION = "us-east-1" +BUCKET = "test-reingest-docs" +USER_ID = "user-reingest" +FILE_ID = "drive-file-reingest" +CONNECTOR_ID = "google-workspace" +DOCUMENT_ID = "doc-reingest" +FILENAME = "report.pdf" + + +# ── fakes ──────────────────────────────────────────────────────────────────── +class FakeBedrock: + """ManagedKbBackend's surface, modelling per-document status like the service. + + ``ingest`` reads the object as it is in S3 at submission — Bedrock ingests from + the S3 location — and re-ingesting an id replaces the document. A submitted + document reports ``IN_PROGRESS`` on the next probe, and ``INDEXED`` (or + ``FAILED``) on the one after, stamped with the time it flipped. While ``hold`` + is set it stays ``IN_PROGRESS``. ``stale_probes`` makes the probes straight + after a submit still report the PREVIOUS version's status and timestamp, which + is what a document Bedrock already holds can look like. + """ + + def __init__(self, stale_probes=0, fail=False): + self.ingests: list = [] + self.indexed_content = None + self.hold = False + self.fail = fail + self.stale_probes = stale_probes + self._status = "NOT_FOUND" + self._updated_at = None + self._pending = None + self._probes_until_flip = 0 + self._stale_left = 0 + self._agent_client = SimpleNamespace(get_knowledge_base_documents=self._get_documents) + + # -- the private surface `document_status` reuses -------------------------- + def _agent(self): + return self._agent_client + + def _locate(self, kb_ref): + return ("KB123", "DS456") + + def _get_documents(self, **kwargs): + if self._stale_left: + self._stale_left -= 1 + elif self._pending is not None: + if self.hold or self._probes_until_flip: + self._probes_until_flip = max(self._probes_until_flip - 1, 0) + self._status = "IN_PROGRESS" + self._updated_at = datetime.now(timezone.utc) + else: + self._status = "FAILED" if self.fail else "INDEXED" + self._updated_at = datetime.now(timezone.utc) + if not self.fail: + self.indexed_content = self._pending + self._pending = None + if self._status == "NOT_FOUND": + return {"documentDetails": []} + return {"documentDetails": [{"status": self._status, "updatedAt": self._updated_at}]} + + # -- the protocol surface -------------------------------------------------- + async def ingest(self, kb_ref, source): + body = boto3.client("s3", region_name=REGION).get_object(Bucket=BUCKET, Key=source.s3_key) + self._pending = body["Body"].read() + self.ingests.append(self._pending) + self._probes_until_flip = 1 + self._stale_left = self.stale_probes + + async def search(self, kb_ref, query, top_k=5, retrieval_filter=None): + if self.indexed_content is None: + return [] + chunk = MagicMock() + chunk.metadata = {"document_id": DOCUMENT_ID} + return [chunk] + + +class FakeDriveAdapter: + def __init__(self, version, content): + self.version = version + self.content = content + + async def get_file_metadata(self, access_token, file_id): + return {"version": self.version, "trashed": False} + + async def download(self, access_token, file_id): + return DownloadedFile(content=self.content, filename=FILENAME, content_type="application/pdf") + + +# ── fixtures ───────────────────────────────────────────────────────────────── +@pytest.fixture(autouse=True) +def _fast_polls(monkeypatch): + monkeypatch.setattr(ic, "INDEXED_POLL_TIMEOUT_SECONDS", 0.2) + monkeypatch.setattr(ic, "INDEXED_POLL_INTERVAL_SECONDS", 0.001) + monkeypatch.setattr(ic, "RETRIEVABLE_POLL_TIMEOUT_SECONDS", 0.05) + monkeypatch.setattr(ic, "RETRIEVABLE_POLL_INTERVAL_SECONDS", 0.001) + + +@pytest.fixture() +def drive(monkeypatch): + """The worker's seams to Google and the vault. Returns a setter for the file.""" + provider = SimpleNamespace(provider_id=CONNECTOR_ID, scopes=[], custom_parameters=None) + + async def fake_get_provider(self, provider_id): + return provider if provider_id == CONNECTOR_ID else None + + async def fake_resolve(provider, user_id): + return "test-access-token" + + monkeypatch.setattr( + "apis.shared.oauth.provider_repository.OAuthProviderRepository.get_provider", fake_get_provider + ) + monkeypatch.setattr(worker, "_resolve_access_token", fake_resolve) + + def set_file(version, content): + adapter = FakeDriveAdapter(version, content) + monkeypatch.setattr(worker.registry, "get", lambda key: adapter if key == "google-drive" else None) + + return set_file + + +@pytest.fixture() +def kb(assistants_table, monkeypatch): + """A managed agent holding one Drive-imported document with a daily sync policy.""" + monkeypatch.setenv("S3_ASSISTANTS_DOCUMENTS_BUCKET_NAME", BUCKET) + boto3.client("s3", region_name=REGION).create_bucket(Bucket=BUCKET) + + async def _build(): + assistant = await create_assistant( + owner_id=USER_ID, owner_name="U", name="A", description="d", + instructions="i", vector_index_id="assistants-index", + ) + assistant_id = assistant.assistant_id + # Imports reserve nothing at request time, so the row declares no size. + document = await create_document( + assistant_id=assistant_id, + filename=FILENAME, + content_type="application/pdf", + size_bytes=0, + s3_key=f"assistants/{assistant_id}/documents/{DOCUMENT_ID}/{FILENAME}", + document_id=DOCUMENT_ID, + provenance=DocumentProvenance( + source_connector_id=CONNECTOR_ID, + source_adapter_key="google-drive", + source_file_id=FILE_ID, + imported_by_user_id=USER_ID, + source_etag="1", + ), + ) + policy = await create_sync_policy( + assistant_id=assistant_id, source_type="drive_file", source_ref=DOCUMENT_ID, + interval="daily", created_by_user_id=USER_ID, + ) + return assistant_id, document, policy + + assistant_id, document, policy = asyncio.run(_build()) + return SimpleNamespace( + table=assistants_table, assistant_id=assistant_id, key=document.s3_key, policy=policy + ) + + +def _seed_kb_record(kb, engine="managed", **counters): + item = { + "PK": f"AST#{kb.assistant_id}", + "SK": f"KB#{kb.assistant_id}", + "appKbId": kb.assistant_id, + "ownerUserId": USER_ID, + "awsKbId": "KB123", + "awsDataSourceId": "DS456", + } + if engine: + item["retrievalEngine"] = engine + item.update({k: Decimal(v) for k, v in counters.items()}) + kb.table.put_item(Item=item) + + +def _put(kb, content): + boto3.client("s3", region_name=REGION).put_object(Bucket=BUCKET, Key=kb.key, Body=content) + + +def _event(kb, bedrock): + """The ObjectCreated event the S3 overwrite produces, handled by the consumer.""" + with patch("apis.shared.kb_backend.managed_backend.ManagedKbBackend", return_value=bedrock): + return ic.handle_object(BUCKET, kb.key) + + +def _sync(kb): + return asyncio.run( + worker.run_sync({ + "policyId": kb.policy.policy_id, + "assistantId": kb.assistant_id, + "sourceType": "drive_file", + "sourceRef": DOCUMENT_ID, + }) + ) + + +def _doc(kb): + return kb.table.get_item(Key={"PK": f"AST#{kb.assistant_id}", "SK": f"DOC#{DOCUMENT_ID}"}).get("Item") + + +def _counters(kb): + item = kb.table.get_item(Key={"PK": f"AST#{kb.assistant_id}", "SK": f"KB#{kb.assistant_id}"}).get("Item") + return tuple(int((item or {}).get(k) or 0) for k in ("storedBytes", "reservedBytes", "totalBytes")) + + +def _delete(kb): + with patch( + "apis.shared.assistants.service.get_assistant", + new_callable=AsyncMock, + return_value=SimpleNamespace(assistant_id=kb.assistant_id, owner_id=USER_ID), + ): + asyncio.run(soft_delete_document(kb.assistant_id, DOCUMENT_ID, USER_ID)) + + +def _imported_and_complete(kb, bedrock, content=b"a" * 1000): + """The import's own first ingestion, to ``complete`` with its bytes committed.""" + _seed_kb_record(kb) + _put(kb, content) + result = _event(kb, bedrock) + assert result["ingested"] is True + assert _doc(kb)["status"] == "complete" + assert _counters(kb) == (len(content), 0, len(content)) + + +# ── the defect ─────────────────────────────────────────────────────────────── +class TestAChangedFileIsReingested: + def test_the_new_bytes_reach_the_knowledge_base(self, kb, drive): + """The repro. MUTATION GUARD: without the re-ingest branch the consumer + returns "already-settled", Bedrock is never asked, and it keeps serving the + version the file had at import.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + + drive("2", b"b" * 1500) + assert _sync(kb)["result"] == "changed" + result = _event(kb, bedrock) + + assert result.get("note") != "already-settled" + assert bedrock.ingests == [b"a" * 1000, b"b" * 1500] + assert bedrock.indexed_content == b"b" * 1500 + assert _doc(kb)["status"] == "complete" + + def test_the_ledger_grows_by_the_size_delta(self, kb, drive): + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + + drive("2", b"b" * 1500) + _sync(kb) + _event(kb, bedrock) + + assert _counters(kb) == (1500, 0, 1500) + assert _doc(kb)["committedBytes"] == 1500 + + def test_the_ledger_shrinks_by_the_size_delta(self, kb, drive): + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + + drive("2", b"c" * 400) + _sync(kb) + _event(kb, bedrock) + + assert bedrock.indexed_content == b"c" * 400 + assert _counters(kb) == (400, 0, 400) + assert _doc(kb)["committedBytes"] == 400 + + def test_a_delete_after_the_reingest_refunds_the_new_size(self, kb, drive): + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + _event(kb, bedrock) + + _delete(kb) + + assert _counters(kb) == (0, 0, 0) + + +# ── idempotency under EventBridge redelivery ───────────────────────────────── +class TestRedeliveryIsIdempotent: + def test_a_redelivery_after_the_reingest_does_nothing(self, kb, drive): + """The early exit the defect came from still does its job: once the new + version is recorded as ingested, its event is a plain redelivery.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + _event(kb, bedrock) + + result = _event(kb, bedrock) + + assert result["note"] == "already-settled" + assert len(bedrock.ingests) == 2 + assert _counters(kb) == (1500, 0, 1500) + + def test_a_redelivery_while_indexing_neither_resubmits_nor_rereserves(self, kb, drive): + """Lambda's async retry is capped at 2, so a slow re-ingest spans + deliveries. Re-submitting would restart Bedrock's work (the dev failure + the first-ingest path documents); re-reserving would leak the growth.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + + bedrock.hold = True + for _ in range(2): + with pytest.raises(ic.IngestionRoutingError): + _event(kb, bedrock) + assert _counters(kb) == (1000, 500, 1500) + assert _doc(kb)["status"] == "complete", "the previous version stops being served" + + bedrock.hold = False + result = _event(kb, bedrock) + + assert result["note"] == "reingested" + assert len(bedrock.ingests) == 2 + assert _counters(kb) == (1500, 0, 1500) + assert not any(k.startswith("reingest") for k in _doc(kb)) + + def test_a_status_from_before_the_submit_is_not_the_new_version(self, kb, drive): + """Bedrock already holds this id, so the probes straight after a + re-submit can still describe the previous version. MUTATION GUARD: drop + ``not_before`` and the consumer records the new version as ingested while + Bedrock is still serving the old one.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + + bedrock.stale_probes = 3 + _event(kb, bedrock) + + assert bedrock.indexed_content == b"b" * 1500 + assert _doc(kb)["ingestedContentHash"] == worker._sha256(b"b" * 1500) + + def test_a_newer_change_supersedes_an_unfinished_reingest(self, kb, drive): + """Two changes, the first never finishing: its growth reservation is + returned, the newer version is submitted, and only its size is counted.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + bedrock.hold = True + with pytest.raises(ic.IngestionRoutingError): + _event(kb, bedrock) + assert _counters(kb) == (1000, 500, 1500) + + drive("3", b"c" * 1200) + _sync(kb) + bedrock.hold = False + result = _event(kb, bedrock) + + assert result["note"] == "reingested" + assert bedrock.ingests[-1] == b"c" * 1200 + assert bedrock.indexed_content == b"c" * 1200 + assert _counters(kb) == (1200, 0, 1200) + assert _event(kb, bedrock)["note"] == "already-settled" + + +# ── the byte cap ───────────────────────────────────────────────────────────── +class TestTheByteCapGatesTheNewVersion: + @pytest.fixture() + def cap_1200(self, monkeypatch): + monkeypatch.setenv("MANAGED_KB_PER_OWNER_DEFAULT_BYTES", "1200") + monkeypatch.setenv("MANAGED_KB_PER_KB_CEILING_BYTES", "1200") + + def test_growth_past_the_cap_keeps_the_previous_version(self, kb, drive, cap_1200): + """Checked BEFORE submitting: once Bedrock has the new bytes the old + version is gone, and failing then would leave the owner with nothing.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + + result = _event(kb, bedrock) + + assert result["note"] == "byte-cap-exceeded" + assert bedrock.ingests == [b"a" * 1000] + assert bedrock.indexed_content == b"a" * 1000 + doc = _doc(kb) + assert doc["status"] == "complete" + assert doc["committedBytes"] == 1000 + assert "storage limit" in doc["ingestionError"] + assert _counters(kb) == (1000, 0, 1000) + + def test_the_next_sync_retries_once_there_is_room(self, kb, drive, cap_1200, monkeypatch): + """The refused version is re-staged by the next sync run (its gates were + cleared) and goes through once the owner has room.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + _event(kb, bedrock) + + assert _sync(kb)["result"] == "changed", "the refused version is never retried" + + monkeypatch.setenv("MANAGED_KB_PER_OWNER_DEFAULT_BYTES", "2000") + monkeypatch.setenv("MANAGED_KB_PER_KB_CEILING_BYTES", "2000") + result = _event(kb, bedrock) + + assert result["note"] == "reingested" + assert _counters(kb) == (1500, 0, 1500) + assert "ingestionError" not in _doc(kb) + + def test_shrinking_is_never_refused(self, kb, drive, cap_1200): + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock, content=b"a" * 1200) + + drive("2", b"b" * 900) + _sync(kb) + + assert _event(kb, bedrock)["note"] == "reingested" + assert _counters(kb) == (900, 0, 900) + + +# ── terminal outcomes ──────────────────────────────────────────────────────── +class TestTerminalOutcomes: + def test_a_delete_mid_reingest_returns_every_byte(self, kb, drive): + """The delete refunds the committed old size and releases the growth + reserved for the new one; the consumer, finding the row deleting, touches + neither.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + bedrock.hold = True + with pytest.raises(ic.IngestionRoutingError): + _event(kb, bedrock) + + _delete(kb) + assert _counters(kb) == (0, 0, 0) + + bedrock.hold = False + assert _event(kb, bedrock)["note"] == "document-deleted" + assert _counters(kb) == (0, 0, 0) + + def test_bedrock_failing_the_new_version_releases_its_growth(self, kb, drive): + """The document is failed like any other Bedrock failure. Its old size + stays committed — the object is still stored — so a delete refunds it.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + drive("2", b"b" * 1500) + _sync(kb) + + bedrock.fail = True + result = _event(kb, bedrock) + + assert result["status"] == "FAILED" + assert _doc(kb)["status"] == "failed" + assert _counters(kb) == (1000, 0, 1000) + + _delete(kb) + assert _counters(kb) == (0, 0, 0) + + def test_a_failed_document_recovers_on_a_changed_source(self, kb, drive): + """New bytes deserve a fresh attempt. A failed import reserved and + committed nothing, so its whole new size is reserved and committed.""" + bedrock = FakeBedrock(fail=True) + _seed_kb_record(kb) + _put(kb, b"a" * 1000) + _event(kb, bedrock) + assert _doc(kb)["status"] == "failed" + assert _counters(kb) == (0, 0, 0) + + bedrock.fail = False + drive("2", b"b" * 1500) + _sync(kb) + result = _event(kb, bedrock) + + assert result["note"] == "reingested" + assert _doc(kb)["status"] == "complete" + assert _doc(kb)["committedBytes"] == 1500 + assert _counters(kb) == (1500, 0, 1500) + + def test_an_unsynced_redelivery_keeps_the_settled_early_exit(self, kb): + """No staged version, no change in behaviour.""" + bedrock = FakeBedrock() + _imported_and_complete(kb, bedrock) + + assert _event(kb, bedrock)["note"] == "already-settled" + assert len(bedrock.ingests) == 1 + + +# ── legacy ─────────────────────────────────────────────────────────────────── +class TestTheLegacyPathIsUnaffected: + def test_a_legacy_kb_is_left_to_its_own_pipeline(self, kb, drive): + """The marker is written regardless of engine; the consumer still routes + a legacy document away before reading it, and counts nothing.""" + _seed_kb_record(kb, engine=None) + _put(kb, b"a" * 1000) + bedrock = FakeBedrock() + + drive("2", b"b" * 1500) + assert _sync(kb)["result"] == "changed" + result = _event(kb, bedrock) + + assert result["routed"] == "legacy" + assert bedrock.ingests == [] + doc = _doc(kb) + assert doc["stagedContentHash"] == worker._sha256(b"b" * 1500) + assert doc["previousChunkCount"] == 0, "the legacy shrinkage stash is unchanged" + assert _counters(kb) == (0, 0, 0) diff --git a/backend/tests/lambdas/test_kb_sync_worker.py b/backend/tests/lambdas/test_kb_sync_worker.py index 86197b801..bdd6f0453 100644 --- a/backend/tests/lambdas/test_kb_sync_worker.py +++ b/backend/tests/lambdas/test_kb_sync_worker.py @@ -166,6 +166,7 @@ async def test_same_bytes_advances_etag_without_staging( assert staged == [] doc = assistants_table.get_item(Key={"PK": f"AST#{assistant_id}", "SK": "DOC#doc-1"})["Item"] assert doc["sourceEtag"] == "42" # gate 1 passes next run + assert "stagedContentHash" not in doc # nothing staged, nothing to re-ingest async def test_changed_bytes_staged_with_stash(self, assistants_table, staged, token_ok, provider_ok, monkeypatch): adapter = FakeDriveAdapter(metadata={"version": "42", "trashed": False}, content=b"new bytes") @@ -183,6 +184,8 @@ async def test_changed_bytes_staged_with_stash(self, assistants_table, staged, t assert item["sourceEtag"] == "42" assert item["contentHash"] == worker._sha256(b"new bytes") assert item["previousChunkCount"] == 7 + # Marks the version a managed KB's consumer still has to re-ingest. + assert item["stagedContentHash"] == worker._sha256(b"new bytes") assert "lastSyncedAt" in item updated_policy = await get_sync_policy(assistant_id, policy.policy_id) assert updated_policy.last_result == "changed" @@ -465,6 +468,8 @@ async def test_changed_page_stashes_and_records_changed(self, assistants_table, assert item["previousChunkCount"] == 4 assert item["sourceEtag"] == '"e2"' assert item["contentHash"] == "hash2" + # Marks the version a managed KB's consumer still has to re-ingest. + assert item["stagedContentHash"] == "hash2" # crawler invoked in refresh mode without TTL finalization assert fake_crawl["captured"]["finalize_with_ttl"] is False assert fake_crawl["captured"]["settings"].max_pages == 10 diff --git a/docs/specs/assistant-kb-sync.md b/docs/specs/assistant-kb-sync.md index 80c9c0574..d84034393 100644 --- a/docs/specs/assistant-kb-sync.md +++ b/docs/specs/assistant-kb-sync.md @@ -192,6 +192,18 @@ boundary posture as today's ingestion Lambda). tail chunks — the one real gap in "reuse the pipeline as-is." Implemented as: worker stashes `previous_chunk_count` on the document record; the ingestion handler's completion step deletes the tail if present.) + **Managed knowledge bases** take the same S3 event to the managed + ingestion consumer instead, where the overwrite's event is + indistinguishable from a redelivery of the original upload's: the + document is already `complete` with its bytes settled. So the worker + also writes `stagedContentHash` before staging, and the consumer + re-ingests while it differs from the `ingestedContentHash` the last + re-ingest recorded — reserving the size growth against the byte cap + *before* submitting (a file that grew past the cap is refused and its + previous version keeps serving; the next sync retries), then settling + the difference. The same marker covers a changed page in §6.2. See + `kb_migration/ingestion_consumer.py` ("A changed source is not a + redelivery"). 6. Update `source_etag`, `content_hash`, `last_synced_at`, `last_result`. ### 6.2 Web crawl (`web_crawl`)