File size: 6,608 Bytes
30ea0e9
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
"""Shared list-query grammar (§16.2) over read-model records.

One grammar across ``GET /v1/messages``, ``/v1/results``, ``/v1/agents`` and
``/v1/inbox/{handle}``: filename-tier filters (``agent``, ``since``/``until``)
prune before any content is touched; frontmatter/content filters run over the
cached records; then order → cursor/limit → expand.
"""
from __future__ import annotations

import re
from datetime import datetime, timezone

from app.errors import InvalidQuery
from app.models import MessageListing, MessageRecord
from app.naming import agent_from_filename
from app.read_model import Record
from app.validation import validate_agent_id


STAMP_LEN = len("YYYYMMDD-HHmmss-mmm")

_COMPACT_RE = re.compile(r"^\d{8}(?:-\d{6}(?:-\d{3})?)?$")

VERIFICATION_STATES = ("pending", "valid", "invalid")


def filename_stamp(filename: str) -> str:
    """The server-stamped chronological prefix of a message/result filename."""
    return filename[:STAMP_LEN]


def normalize_stamp(value: str, *, param: str) -> str:
    """Accept ISO 8601 or compact ``YYYYMMDD[-HHmmss[-mmm]]``; return a compact
    stamp comparable against the server-stamped filename prefix (UTC)."""
    v = value.strip()
    if _COMPACT_RE.match(v):
        if len(v) == 8:
            return v + "-000000-000"
        if len(v) == 15:
            return v + "-000"
        return v
    try:
        dt = datetime.fromisoformat(v.replace("Z", "+00:00"))
    except ValueError:
        raise InvalidQuery(
            f"`{param}` must be ISO 8601 or YYYYMMDD-HHmmss[-mmm], got {value!r}"
        )
    if dt.tzinfo is not None:
        dt = dt.astimezone(timezone.utc)
    return dt.strftime("%Y%m%d-%H%M%S-") + f"{dt.microsecond // 1000:03d}"


def parse_verification_param(value: str | None) -> set[str] | None:
    """CSV of verification states, e.g. ``valid,pending``. None → no filter."""
    if value is None:
        return None
    states = {s.strip() for s in value.split(",") if s.strip()}
    bad = states - set(VERIFICATION_STATES)
    if bad or not states:
        raise InvalidQuery(
            f"`verification` must be a CSV of {VERIFICATION_STATES}, got {value!r}"
        )
    return states


def apply_filters(
    records: list[Record],
    *,
    agent: str | None = None,
    since: str | None = None,
    until: str | None = None,
    fm_eq: dict[str, str] | None = None,
    q: str | None = None,
) -> list[Record]:
    """``agent``/``since``/``until`` are answerable from filenames alone;
    ``fm_eq`` matches frontmatter values by string equality; ``q`` is a
    case-insensitive substring over frontmatter+body."""
    out: list[Record] = []
    for r in records:
        if agent is not None and agent_from_filename(r.filename) != agent:
            continue
        if since is not None and filename_stamp(r.filename) < since:
            continue
        if until is not None and filename_stamp(r.filename) > until:
            continue
        if fm_eq is not None:
            if any(str(r.frontmatter.get(k, "")) != want for k, want in fm_eq.items()):
                continue
        if q is not None and not _q_match(r, q):
            continue
        out.append(r)
    return out


def _q_match(r: Record, q: str) -> bool:
    ql = q.lower()
    if ql in r.body.lower():
        return True
    return any(ql in f"{k}: {v}".lower() for k, v in r.frontmatter.items())


def effective_limit(limit: int | None, expand: bool, cap: int) -> int | None:
    """Expanded pages are capped so one call can't serialize the whole corpus."""
    if not expand:
        return limit
    if limit is None or limit <= 0 or limit > cap:
        return cap
    return limit


def paginate(
    records: list[Record],
    *,
    order: str,
    limit: int | None,
    after: str | None,
    before: str | None,
) -> tuple[list[Record], str | None]:
    """Cursor + slice over filtered records (ascending filename order in).

    ``after``/``before`` are exclusive filename bounds. Returns the page and a
    ``next`` cursor (the page's last filename) when more matches remain in the
    traversal direction — pass it back as ``after`` for asc, ``before`` for
    desc.
    """
    if after is not None:
        records = [r for r in records if r.filename > after]
    if before is not None:
        records = [r for r in records if r.filename < before]
    ordered = list(reversed(records)) if order == "desc" else records
    if limit is not None and 0 < limit < len(ordered):
        page = ordered[:limit]
        return page, page[-1].filename
    return ordered, None


def list_message_like(
    records: list[Record],
    *,
    agent: str | None,
    since: str | None,
    until: str | None,
    type_: str | None,
    via: str | None,
    q: str | None,
    expand: bool,
    limit: int | None,
    order: str,
    after: str | None,
    before: str | None,
    expand_cap: int,
) -> MessageListing:
    """The full §16.2 pipeline for message-shaped folders (board and inboxes)."""
    if agent is not None:
        validate_agent_id(agent)
    fm_eq: dict[str, str] = {}
    if type_ is not None:
        fm_eq["type"] = type_
    if via is not None:
        fm_eq["via"] = via
    filtered = apply_filters(
        records,
        agent=agent,
        since=normalize_stamp(since, param="since") if since is not None else None,
        until=normalize_stamp(until, param="until") if until is not None else None,
        fm_eq=fm_eq or None,
        q=q,
    )
    page, next_cursor = paginate(
        filtered,
        order="desc" if order == "desc" else "asc",
        limit=effective_limit(limit, expand, expand_cap),
        after=after,
        before=before,
    )
    items: list[str] | list[MessageRecord]
    if expand:
        items = [
            MessageRecord(
                filename=r.filename,
                frontmatter=r.frontmatter,
                body=r.body,
                reasons=r.reasons,
            )
            for r in page
        ]
    else:
        items = [r.filename for r in page]
    return MessageListing(
        count=len(records),
        matched=len(filtered),
        items=items,
        next=next_cursor,
        # The cursor to persist verbatim (WATCH_DESIGN.md §4.4): the newest
        # filename ON THIS PAGE, computed here so no client ever has to scan
        # author-controlled record content for a maximum. Independent of
        # `order` (asc pages end on it, desc pages start on it) and of `next`
        # (which is a pagination handle, not a read position).
        cursor=max((r.filename for r in page), default=None),
    )