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`)