File size: 4,941 Bytes
7124acf
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
"""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