Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions backend/src/apis/app_api/documents/services/document_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
360 changes: 359 additions & 1 deletion backend/src/apis/app_api/kb_migration/ingestion_consumer.py

Large diffs are not rendered by default.

11 changes: 11 additions & 0 deletions backend/src/apis/app_api/kb_sync/records.py
Original file line number Diff line number Diff line change
Expand Up @@ -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] = {}
Expand All @@ -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

Expand Down
9 changes: 7 additions & 2 deletions backend/src/apis/app_api/kb_sync/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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(
Expand Down Expand Up @@ -299,14 +302,16 @@ 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,
source_etag=etag,
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(
Expand Down
160 changes: 159 additions & 1 deletion backend/src/apis/shared/kb_backend/byte_cap.py
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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.

Expand Down
Loading