File size: 9,013 Bytes
a1dd5ba
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
"""The motion vector channel: motion = temporal DIFFERENCE of appearance.

The Marengo lesson is that motion must be its own vector, never averaged
into appearance. The obvious implementations were measured and lost:
  X-CLIP (video-native, cross-frame attention)    close/open AUC 0.67
  VLM swap-contrast (query-time, 2 model calls)   0.86 (2B) / 0.91 (7B)
  ctx adapter (caption-trained)                   ~chance
  THIS: delta-appearance + query-swap difference  0.98 / 0.98

An event's motion vector is normalize(app(end) - app(start)) in SigLIP
space, from the per-frame vectors ALREADY on disk — ingest cost is a
numpy mean. A directional query's motion vector is
normalize(embed(q) - embed(swap(q))) — the same swap-contrast idea, moved
from VLM logits into embedding space, so it runs index-only in
microseconds. Both sides subtract their shared appearance component, so
what remains on each side is the DIRECTION of state change, and they
were trained (by SigLIP, incidentally) to point the same way.

Limits, measured: works when the state change is visually large
(drawer open<->close: 0.98); chance when a small object moves
(put-in/take-out: 0.46-0.50) — those queries are carried by the
appearance and caption-lexical channels instead. Non-directional queries
have no swap and the channel ABSTAINS (RRF median rank, no veto).
"""
from __future__ import annotations

import numpy as np
import pyarrow as pa


def index_motion(store, events_table="context_events", edge_s=1.5,
                 frames_table="frame_vectors"):
    """One motion vector per event span -> `motion_vectors` table.

    Event spans come from the gate-segmented event index so every retrieval
    unit carries both an appearance-family vector and a motion vector —
    the multi-vector layout, aligned. Rebuild = atomic replace commit.
    """
    import time
    import uuid

    from .embeddings import _vec_table
    from .log import FileEntry
    from .store import write_parquet
    t_start = time.time()

    ev = store.table(events_table).scan()
    fv, fvecs = _vec_table(store, frames_table)
    fs = np.asarray(fv.column("stream").to_pylist())
    ft = np.asarray([int(v) for v in fv.column("ts").to_pylist()])
    edge = int(edge_s * 1e9)

    # per-stream sorted timelines + searchsorted per event. The full-array
    # boolean-mask version was O(events x frames): at 100 h that is 50k
    # events over 1.8M frames = 16.5 MINUTES for what is 25 s of means.
    by_stream = {}
    for s in np.unique(fs):
        m = np.where(fs == s)[0]
        o = np.argsort(ft[m], kind="stable")
        by_stream[s] = (ft[m][o], fvecs[m[o]])

    rows_s, rows_a, rows_b, rows_v = [], [], [], []
    for s, a, b in zip(ev.column("stream").to_pylist(),
                       (int(v) for v in ev.column("ts").to_pylist()),
                       (int(v) for v in ev.column("t1").to_pylist())):
        if s not in by_stream:
            continue
        tt, vv = by_stream[s]
        lo, hi = np.searchsorted(tt, [a, b + 1])
        if hi - lo < 6:
            continue
        t0, t1 = int(tt[lo]), int(tt[hi - 1])
        e0 = np.searchsorted(tt, t0 + edge, side="right")
        e1 = np.searchsorted(tt, t1 - edge, side="left")
        v0 = vv[lo:max(e0, lo + 1)].mean(0)
        v1 = vv[min(e1, hi - 1):hi].mean(0)
        d = v1 - v0
        n = float(np.linalg.norm(d))
        if n < 1e-6:
            continue
        rows_s.append(s)
        rows_a.append(a)
        rows_b.append(b)
        rows_v.append((d / n).astype(np.float32))
    if not rows_v:
        return {"events": 0}

    dim = len(rows_v[0])
    tbl = pa.table({
        "ts": pa.array(rows_a, pa.int64()),
        "t1": pa.array(rows_b, pa.int64()),
        "stream": pa.array(rows_s),
        "vector": pa.array([v.tolist() for v in rows_v],
                           pa.list_(pa.float32(), dim)),
    })
    tab = store.table("motion_vectors")
    meta = {"model": "delta-appearance", "dim": dim, "edge_s": edge_s,
            "events_table": events_table}
    try:
        existing = tab.state().files
    except Exception:
        existing = []
    if existing:
        fname = f"part-{uuid.uuid4().hex[:12]}.parquet"
        path = tab.dir / fname
        write_parquet(tbl, path)
        tsv = tbl.column("ts")
        version = tab.log.commit(
            op="replace", kind="embeddings", schema=str(tbl.schema),
            add=[FileEntry(fname, len(tbl), path.stat().st_size,
                           tsv[0].as_py(), tsv[-1].as_py())],
            remove=[f.path for f in existing], meta=meta)
    else:
        version = tab.append(tbl, kind="embeddings", meta=meta)
    return {"events": len(tbl), "dim": dim, "version": version,
            "seconds": round(time.time() - t_start, 2)}


def motion_scores(store, qv, qv_swap):
    """(stream, t0, t1, score) per event: motion vectors against the
    query-direction vector. Pure matmul over a version-cached matrix."""
    from .embeddings import _vec_table
    tbl, vecs = _vec_table(store, "motion_vectors")
    qd = np.asarray(qv, np.float32) - np.asarray(qv_swap, np.float32)
    n = float(np.linalg.norm(qd))
    if n < 1e-6:
        return []
    sc = vecs @ (qd / n)
    return list(zip(tbl.column("stream").to_pylist(),
                    (int(v) for v in tbl.column("ts").to_pylist()),
                    (int(v) for v in tbl.column("t1").to_pylist()),
                    sc.tolist()))


_MOT_IDX = {}


def motion_lookup(store, qv, qv_swap):
    """lookup(stream, t0, t1) -> score of the recording CONTAINING the span,
    or nan. The per-stream sorted span index is cached per table version;
    a query pays one matmul + searchsorted per segment. The Python-dict
    scan this replaces cost ~70 ms per directional query at 50k recordings
    (measured at the 100 h scale test)."""
    from .embeddings import _vec_table
    ver = store.table("motion_vectors").state().version
    key = (str(store.dir), ver)
    if key not in _MOT_IDX:
        tbl, _ = _vec_table(store, "motion_vectors")
        ss = np.asarray(tbl.column("stream").to_pylist())
        sa = np.asarray([int(v) for v in tbl.column("ts").to_pylist()])
        sb = np.asarray([int(v) for v in tbl.column("t1").to_pylist()])
        idx = {}
        for s in np.unique(ss):
            m = np.where(ss == s)[0]
            o = np.argsort(sa[m], kind="stable")
            idx[s] = (sa[m][o], sb[m][o], m[o])
        if len(_MOT_IDX) > 8:
            _MOT_IDX.clear()
        _MOT_IDX[key] = idx
    idx = _MOT_IDX[key]
    _, vecs = _vec_table(store, "motion_vectors")
    qd = np.asarray(qv, np.float32) - np.asarray(qv_swap, np.float32)
    n = float(np.linalg.norm(qd))
    if n < 1e-6:
        return lambda s, a, b: float("nan")
    sc = vecs @ (qd / n)

    def lookup(s, a, b):
        if s not in idx:
            return float("nan")
        t0s, t1s, rows = idx[s]
        j = int(np.searchsorted(t0s, a, side="right")) - 1
        if j >= 0 and b <= int(t1s[j]) + 1:
            return float(sc[rows[j]])
        return float("nan")
    return lookup


def motion_candidates(store, qv, qv_swap, streams, t0s, top=24):
    """DIRECTION-FIRST RECALL, content-gated. Given a broad appearance-
    plausible window set (their streams + start times), return the top
    recordings ranked by motion score as (stream, rec_t0, rec_t1, score).

    Why the gate: unrestricted motion recall was measured useless — the
    tiny direction cosines (~0.07) drown corpus-wide in random tabletop
    deltas. Why recall at all: at the 100 h scale, appearance ordering
    alone never surfaces true close/open recordings into the candidate
    pool (put-in clips carry a stronger 'drawer' signal), so ranking-only
    motion had nothing correct to rank. Appearance answers WHAT is
    plausible; motion picks WHICH of those changed the right way."""
    from .embeddings import _vec_table
    ver = store.table("motion_vectors").state().version
    key = (str(store.dir), ver)
    if key not in _MOT_IDX:
        motion_lookup(store, qv, qv_swap)      # builds the cache
    idx = _MOT_IDX[key]
    _, vecs = _vec_table(store, "motion_vectors")
    qd = np.asarray(qv, np.float32) - np.asarray(qv_swap, np.float32)
    n = float(np.linalg.norm(qd))
    if n < 1e-6:
        return []
    sc = vecs @ (qd / n)

    streams = np.asarray(streams)
    t0s = np.asarray(t0s)
    hit_rows = set()
    for s in np.unique(streams):
        if s not in idx:
            continue
        r0, r1, rows = idx[s]
        w = t0s[streams == s]
        j = np.searchsorted(r0, w, side="right") - 1
        ok = j >= 0
        hit_rows.update(int(r) for r in np.unique(rows[j[ok]]))
    if not hit_rows:
        return []
    hits = sorted(hit_rows, key=lambda r: -sc[r])[:top]
    tbl, _ = _vec_table(store, "motion_vectors")
    ss = tbl.column("stream")
    sa = tbl.column("ts")
    sb = tbl.column("t1")
    return [(ss[r].as_py(), int(sa[r].as_py()), int(sb[r].as_py()),
             float(sc[r])) for r in hits]