ishaq101's picture
/fix planner and unstructured pipeline (#23)
7124acf
Raw History Blame Contribute Delete
4.94 kB
"""Figure markers → short-lived signed URLs (UNSTRUCTURED_V2_PLAN U4).
Answers carry `![caption](asset://<asset_id>)`, never a URL. The FE sends the ids
it is about to render; this returns a 1-hour presigned GET per id it can prove
the caller may see, and `None` for every other id. Nothing here is persisted.
Resolution, adapted from the reference design Rifqi supplied (2026-09-22):
1. Accept only the canonical id shape — 16 lowercase hex (`assets.asset_id`).
Never a storage key, a URL, or free text.
2. The document must be bound to the caller's analysis (`scope.bound_documents`,
read with the caller's `user_id`, so another tenant's analysis yields nothing).
3. The storage key comes ONLY from a stored v2 chunk row of that document whose
`cmetadata.assets` lists the id and whose `user_id` is the caller — and it must
sit under that document's own prefix (`{user_id}/{document_id}/assets/`), a
second check that a row can only ever point at its own document's objects.
4. Sign for 3600 s. Log ids, never the URL or the key.
⚠️ `user_id` is caller-supplied — the surface is unauthenticated until DEV_PLAN
#43 (decision U-D2). The checks above narrow access to "whoever holds a valid
user_id + analysis_id pair".
Never-throw: a failed read or a failed signature degrades to `None` for the
affected ids, logged with `degraded_seam`, never a 500.
"""
from __future__ import annotations
import re
from datetime import UTC, datetime, timedelta
from typing import Any
from sqlalchemy import text
from src.document.scope import bound_document_ids
from src.middlewares.logging import get_logger
logger = get_logger("image_token")
URL_TTL_SECONDS = 3600
_ASSET_ID = re.compile(r"[0-9a-f]{16}")
_LOOKUP_SQL = text("""
SELECT e.cmetadata->'data'->>'document_id' AS document_id,
a->>'asset_id' AS asset_id,
a->>'storage_key' AS storage_key
FROM langchain_pg_embedding e
JOIN langchain_pg_collection c ON e.collection_id = c.uuid
CROSS JOIN LATERAL jsonb_array_elements(
COALESCE(e.cmetadata->'assets', '[]'::jsonb)
) AS a
WHERE c.name = 'documents'
AND e.cmetadata->>'parser_version' = 'v2'
AND e.cmetadata->>'user_id' = :user_id
AND e.cmetadata->'data'->>'document_id' = ANY(:document_ids)
AND a->>'asset_id' = ANY(:asset_ids)
""")
async def _lookup_keys(
user_id: str, document_ids: list[str], asset_ids: list[str]
) -> list[tuple[str, str, str]]:
from src.db.postgres.connection import engine
async with engine.connect() as conn:
rows = await conn.execute(
_LOOKUP_SQL,
{"user_id": user_id, "document_ids": document_ids, "asset_ids": asset_ids},
)
return [(r.document_id, r.asset_id, r.storage_key) for r in rows]
async def resolve_image_urls(
user_id: str,
analysis_id: str,
asset_ids: list[Any],
*,
storage: Any = None,
scope_fn=bound_document_ids,
lookup_fn=_lookup_keys,
) -> tuple[dict[str, str | None], str]:
"""Return ({asset_id: url | None}, expires_at ISO-8601 UTC). Never raises."""
expires_at = (datetime.now(UTC) + timedelta(seconds=URL_TTL_SECONDS)).isoformat()
urls: dict[str, str | None] = {str(a): None for a in asset_ids}
valid = sorted({a for a in asset_ids if isinstance(a, str) and _ASSET_ID.fullmatch(a)})
if not valid:
return urls, expires_at
doc_ids = await scope_fn(user_id, analysis_id)
if not doc_ids:
return urls, expires_at
try:
rows = await lookup_fn(user_id, doc_ids, valid)
except Exception as e:
logger.warning(
"image key lookup failed",
analysis_id=analysis_id,
n_ids=len(valid),
error=repr(e),
degraded_seam="image_token_lookup",
)
return urls, expires_at
if storage is None:
from src.storage.object_storage.supabase_s3 import object_storage
storage = object_storage
for document_id, asset_id, key in rows:
if urls.get(asset_id) is not None or not key:
continue
if not key.startswith(f"{user_id}/{document_id}/assets/"):
logger.warning(
"storage key outside its document prefix — refused",
document_id=document_id,
asset_id=asset_id,
)
continue
try:
urls[asset_id] = storage.presigned_get_url(key, expires_in=URL_TTL_SECONDS)
except Exception as e:
logger.warning(
"presign failed",
asset_id=asset_id,
error=repr(e),
degraded_seam="image_token_presign",
)
logger.info(
"image urls resolved",
analysis_id=analysis_id,
requested=len(asset_ids),
signed=sum(1 for v in urls.values() if v),
)
return urls, expires_at