codelion commited on
Commit
a866fc9
·
verified ·
1 Parent(s): 9327c55

Record one summary per visit instead of every click; hourly writes to data/incoming; monthly compaction

Browse files
DEPLOY.md CHANGED
@@ -9,7 +9,7 @@ Prerequisites: `hf auth login` as a member of `mlx-community` with write access,
9
  ```bash
10
  .venv/bin/python -m pytest
11
  EXPLORER_SINK=local EXPLORER_FLUSH_SECONDS=20 .venv/bin/uvicorn app.main:app --port 7860
12
- # walk through the UI; events land in .runtime/local_dataset/data/events/
13
  docker build --platform linux/amd64 -t mlx-model-explorer . \
14
  && docker run --rm -p 7860:7860 -e EXPLORER_SINK=local mlx-model-explorer
15
  ```
@@ -21,7 +21,7 @@ docker build --platform linux/amd64 -t mlx-model-explorer . \
21
  .venv/bin/python scripts/deploy.py create # refuses if the Space exists; add --allow-existing-dataset if the dataset was created earlier
22
  ```
23
 
24
- This creates `mlx-community/mlx-model-explorer-data` (private dataset, card uploaded) and `codelion/mlx-model-explorer` (private Docker Space with `EXPLORER_SINK=hub`, `DATASET_REPO`, `EXPLORER_FLUSH_SECONDS=600`).
25
 
26
  > **Why the Space isn't in `mlx-community`:** Hugging Face now requires a Team or Enterprise plan for an organization to run Docker or Gradio Spaces on free CPU (`402 Payment Required`). The dataset can live in the org. Once an org admin has a plan or a hardware grant, deploy there with `SPACE_REPO=mlx-community/mlx-model-explorer scripts/deploy.py create --allow-existing-dataset` followed by `upload`.
27
 
@@ -82,6 +82,17 @@ hf upload mlx-community/mlx-model-explorer static_embed . --repo-type space
82
 
83
  Deployed as private `mlx-community/mlx-model-explorer` on 2026-09-14. The app renders inside the Hub page. Make the Static Space public at launch with `hf repos settings mlx-community/mlx-model-explorer --type space --no-private` (or from its Settings page).
84
 
 
 
 
 
 
 
 
 
 
 
 
85
  ## Updating
86
 
87
  Change code, run the tests, `docker build` locally, then `scripts/deploy.py upload -m "what changed"`. Events are append-only, so redeploys never touch existing data. On shutdown the app flushes buffered events, and anything that fails to upload is retried from the on-disk spool.
 
9
  ```bash
10
  .venv/bin/python -m pytest
11
  EXPLORER_SINK=local EXPLORER_FLUSH_SECONDS=20 .venv/bin/uvicorn app.main:app --port 7860
12
+ # walk through the UI; events land in .runtime/local_dataset/data/incoming/
13
  docker build --platform linux/amd64 -t mlx-model-explorer . \
14
  && docker run --rm -p 7860:7860 -e EXPLORER_SINK=local mlx-model-explorer
15
  ```
 
21
  .venv/bin/python scripts/deploy.py create # refuses if the Space exists; add --allow-existing-dataset if the dataset was created earlier
22
  ```
23
 
24
+ This creates `mlx-community/mlx-model-explorer-data` (private dataset, card uploaded) and `codelion/mlx-model-explorer` (private Docker Space with `EXPLORER_SINK=hub`, `DATASET_REPO`, `EXPLORER_FLUSH_SECONDS=3600`).
25
 
26
  > **Why the Space isn't in `mlx-community`:** Hugging Face now requires a Team or Enterprise plan for an organization to run Docker or Gradio Spaces on free CPU (`402 Payment Required`). The dataset can live in the org. Once an org admin has a plan or a hardware grant, deploy there with `SPACE_REPO=mlx-community/mlx-model-explorer scripts/deploy.py create --allow-existing-dataset` followed by `upload`.
27
 
 
82
 
83
  Deployed as private `mlx-community/mlx-model-explorer` on 2026-09-14. The app renders inside the Hub page. Make the Static Space public at launch with `hf repos settings mlx-community/mlx-model-explorer --type space --no-private` (or from its Settings page).
84
 
85
+ ## Monthly compaction (run from any machine with write access)
86
+
87
+ Early each month, merge the previous month's hourly shards into one file:
88
+
89
+ ```bash
90
+ .venv/bin/python scripts/compact.py # dry run: shows shards, rows before/after
91
+ .venv/bin/python scripts/compact.py --apply # one commit: add data/events/YYYY-MM.parquet, delete the month's data/incoming shards
92
+ ```
93
+
94
+ It keeps the newest `session` row per visit, never edits row values, refuses the current month, and is safe to re-run (late shards are merged into the existing monthly file).
95
+
96
  ## Updating
97
 
98
  Change code, run the tests, `docker build` locally, then `scripts/deploy.py upload -m "what changed"`. Events are append-only, so redeploys never touch existing data. On shutdown the app flushes buffered events, and anything that fails to upload is retried from the on-disk spool.
README.md CHANGED
@@ -46,14 +46,15 @@ app/ FastAPI
46
  memory.py weights + KV cache + overhead, GPU-usable memory, fit classes
47
  recommend.py RecommendationEngine interface + HeuristicEngine v1
48
  events.py strict event schema, bounds, note scrubbing, plausibility flags
49
- sink.py append-only Parquet shards: local directory (dev) or Hub dataset (prod)
50
- stats.py k-anonymous aggregates + per-model community signals
51
  bench/ mlx_explorer_bench.py (runs locally with mlx-lm; submits only with --submit)
52
  data/ catalogue_snapshot.json.gz (fallback when the Hub API is unreachable)
 
53
  tests/ pytest suite
54
  ```
55
 
56
- Events are buffered and written as **one Parquet shard per flush** (default every 10 minutes, only when there is data) under `data/events/YYYY/MM/DD/`. That caps Hub commits at about 144 a day. Shards are never rewritten, and a failed upload stays spooled on disk and retries with backoff.
57
 
58
  ## Configuration
59
 
@@ -63,8 +64,9 @@ Events are buffered and written as **one Parquet shard per flush** (default ever
63
  | `DATASET_REPO` | `mlx-community/mlx-model-explorer-data` | dataset for `hub` mode |
64
  | `HF_TOKEN` | (none) | **Space secret**, a fine-grained token with write access to the dataset only. It is never used for catalogue requests |
65
  | `HF_READ_TOKEN` | (none) | optional token for Hub reads (rate limits only) |
66
- | `EXPLORER_FLUSH_SECONDS` | `600` | flush interval |
67
- | `EXPLORER_DATA_PREFIX` | `data/events` | path prefix inside the dataset (tests use `data/_test`) |
 
68
  | `EXPLORER_ORGS` | `mlx-community` | organizations to list |
69
  | `EXPLORER_RATE_PER_MIN` | `120` | event requests per client per minute |
70
 
 
46
  memory.py weights + KV cache + overhead, GPU-usable memory, fit classes
47
  recommend.py RecommendationEngine interface + HeuristicEngine v1
48
  events.py strict event schema, bounds, note scrubbing, plausibility flags
49
+ sink.py buffered Parquet shards into data/incoming: local directory (dev) or Hub dataset (prod)
50
+ stats.py k-anonymous aggregates (one count per visit) + per-model community signals
51
  bench/ mlx_explorer_bench.py (runs locally with mlx-lm; submits only with --submit)
52
  data/ catalogue_snapshot.json.gz (fallback when the Hub API is unreachable)
53
+ scripts/ deploy.py (create/upload/publish), compact.py (monthly compaction)
54
  tests/ pytest suite
55
  ```
56
 
57
+ The page records **one summary row per visit** (final filters, what was searched, hardware class, models opened/compared/visited) plus explicit contributions (feedback, MLX and browser benchmarks). Individual clicks are never logged. The Space buffers rows and writes one Parquet shard to `data/incoming/YYYY/MM/DD/` at most hourly, or sooner after 2,000 rows. A failed upload stays spooled on disk and retries with backoff. Once a month, `scripts/compact.py` merges the previous month into `data/events/YYYY-MM.parquet` (see DEPLOY.md).
58
 
59
  ## Configuration
60
 
 
64
  | `DATASET_REPO` | `mlx-community/mlx-model-explorer-data` | dataset for `hub` mode |
65
  | `HF_TOKEN` | (none) | **Space secret**, a fine-grained token with write access to the dataset only. It is never used for catalogue requests |
66
  | `HF_READ_TOKEN` | (none) | optional token for Hub reads (rate limits only) |
67
+ | `EXPLORER_FLUSH_SECONDS` | `3600` | maximum time between writes |
68
+ | `EXPLORER_FLUSH_ROWS` | `2000` | write early once this many rows are buffered |
69
+ | `EXPLORER_DATA_PREFIX` | `data/incoming` | where the Space writes shards |
70
  | `EXPLORER_ORGS` | `mlx-community` | organizations to list |
71
  | `EXPLORER_RATE_PER_MIN` | `120` | event requests per client per minute |
72
 
app/events.py CHANGED
@@ -1,9 +1,16 @@
1
  """Event schema, validation and server-side enrichment.
2
 
 
 
 
 
 
 
 
3
  Clients send a small, typed event. Anything outside the schema is rejected
4
- (`extra="forbid"`), numbers are bounded, free text is capped and scrubbed, and the
5
- server fills in catalogue facts itself instead of trusting the client. Values that
6
- are allowed but implausible are kept and flagged, so dataset users can filter them.
7
  """
8
 
9
  from __future__ import annotations
@@ -15,15 +22,16 @@ from typing import Literal
15
 
16
  from pydantic import BaseModel, ConfigDict, Field, field_validator
17
 
18
- from .memory import RAM_CLASSES, CONTEXTS, bits_per_weight
19
  from .parsing import PARAM_BUCKETS, QUANT_BUCKETS
20
 
21
- SCHEMA_VERSION = 1
22
 
23
- EventType = Literal[
24
- "search", "filter", "model_view", "model_select", "model_click", "compare",
25
- "hardware_test", "browser_benchmark", "feedback", "mlx_benchmark_submission",
26
- ]
 
27
 
28
  CHIP_RE = re.compile(r"^Apple M[1-9]( (Pro|Max|Ultra))?$")
29
  MODEL_ID_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,95}/[A-Za-z0-9][A-Za-z0-9._-]{0,95}$")
@@ -36,7 +44,8 @@ _URL = re.compile(r"(https?://|www\.)\S+", re.I)
36
  _LONG_DIGITS = re.compile(r"\+?\d[\d ()-]{7,}\d")
37
 
38
  MAX_NOTES = 280
39
- MAX_BATCH = 50
 
40
 
41
 
42
  def clean_notes(text: str | None) -> str | None:
@@ -51,21 +60,37 @@ def clean_notes(text: str | None) -> str | None:
51
  return text[:MAX_NOTES] or None
52
 
53
 
 
 
 
 
 
 
 
 
 
 
 
54
  class ClientEvent(BaseModel):
55
  """What the browser (or the benchmark script) is allowed to send."""
56
 
57
  model_config = ConfigDict(extra="forbid", str_max_length=200)
58
 
59
- event_type: EventType
60
  session_id: str | None = None
61
 
 
62
  model_family: str | None = Field(None, max_length=40)
63
  parameter_bucket: Literal[tuple(PARAM_BUCKETS + ["MoE", "unknown"])] | None = None # type: ignore[valid-type]
64
- quantization: Literal[tuple(QUANT_BUCKETS)] | None = None # type: ignore[valid-type]
65
  target_context: int | None = None
66
  priority: Literal["balanced", "quality", "speed", "memory", "long_context"] | None = None
67
  sort: Literal["recommended", "popular", "recent", "community", "all"] | None = None
68
  result_count: int | None = Field(None, ge=0, le=100_000)
 
 
 
 
69
 
70
  hardware_source: Literal["detected", "confirmed", "none"] | None = None
71
  hardware_memory_class: int | None = None
@@ -78,13 +103,16 @@ class ClientEvent(BaseModel):
78
  os_family: Literal["macos", "ios", "windows", "linux", "android", "other"] | None = None
79
  cpu_cores: int | None = Field(None, ge=1, le=256)
80
 
81
- selected_model: str | None = None
82
- selected_model_rank: int | None = Field(None, ge=0, le=10_000)
83
- # what the user was shown for that model at the time (for studying recommendations)
84
- recommendation_score: float | None = Field(None, ge=0, le=100)
85
- fit_class: Literal["Comfortable", "Likely", "Borderline", "Unlikely"] | None = None
86
- compare_models: list[str] | None = Field(None, max_length=5)
 
87
 
 
 
88
  tried: Literal["yes", "no", "planning"] | None = None
89
  quality_rating: Literal["excellent", "good", "acceptable"] | None = None
90
  failure_reason: Literal["too_slow", "too_much_memory", "low_quality", "didnt_run"] | None = None
@@ -112,6 +140,14 @@ class ClientEvent(BaseModel):
112
  mlx_lm_version: str | None = None
113
  macos_major: int | None = Field(None, ge=11, le=40)
114
 
 
 
 
 
 
 
 
 
115
  @field_validator("session_id")
116
  @classmethod
117
  def _session(cls, v):
@@ -119,22 +155,38 @@ class ClientEvent(BaseModel):
119
  raise ValueError("session_id must be 16-32 lowercase hex chars")
120
  return v
121
 
122
- @field_validator("selected_model")
123
  @classmethod
124
  def _model(cls, v):
125
  if v is not None and not MODEL_ID_RE.match(v):
126
  raise ValueError("invalid model id")
127
  return v
128
 
129
- @field_validator("compare_models")
130
  @classmethod
131
- def _models(cls, v):
 
 
 
 
 
 
 
 
 
 
132
  if v is not None:
133
- for m in v:
134
- if not MODEL_ID_RE.match(m):
135
- raise ValueError("invalid model id in compare_models")
 
136
  return v
137
 
 
 
 
 
 
138
  @field_validator("target_context")
139
  @classmethod
140
  def _ctx(cls, v):
@@ -183,28 +235,34 @@ class EventBatch(BaseModel):
183
 
184
  # Column order for the dataset. Every row has every column (null when unused).
185
  COLUMNS: dict[str, str] = {
 
186
  "timestamp": "string", "schema_version": "int", "app_version": "string", "event_type": "string",
187
  "session_id": "string",
188
- "model_family": "string", "model_name": "string", "parameter_bucket": "string", "quantization": "string",
189
- "quant_bits": "float", "quant_mixed": "bool",
190
  "target_context": "int", "priority": "string", "sort": "string", "result_count": "int",
 
 
191
  "hardware_source": "string", "hardware_memory_class": "int", "hardware_confirmed": "bool",
192
  "webgpu_available": "bool", "webgpu_score": "float", "gpu_capability_class": "string",
193
  "gpu_vendor": "string", "gpu_arch": "string", "browser_family": "string", "os_family": "string",
194
  "cpu_cores": "int",
195
- "selected_model": "string", "selected_model_rank": "int", "compare_models": "list",
196
- "hf_downloads_at_selection": "int", "hf_likes_at_selection": "int",
197
- "recommendation_score": "float", "fit_class": "string", "engine_version": "string",
198
- "pipeline": "string",
 
 
 
199
  "tried": "string", "outcome": "string", "quality_rating": "string", "failure_reason": "string",
200
  "reported_ram_gb": "int", "reported_mac_model": "string", "reported_tokens_per_second": "float",
201
  "reported_context": "int", "notes": "string",
 
202
  "benchmark_type": "string", "benchmark_version": "string", "benchmark_duration_ms": "int",
203
  "prompt_tokens": "int", "generation_tokens": "int", "prompt_tps": "float", "generation_tps": "float",
204
  "ttft_ms": "float", "peak_memory_gb": "float", "chip": "string",
205
  "perplexity": "float", "perplexity_stderr": "float", "eval_dataset": "string", "eval_tokens": "int",
206
- "mlx_version": "string",
207
- "mlx_lm_version": "string", "macos_major": "int",
208
  "suspicious_flags": "list",
209
  }
210
 
@@ -213,8 +271,10 @@ def utc_now_seconds() -> str:
213
  return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
214
 
215
 
216
- def to_row(ev: ClientEvent, catalogue=None, app_version: str = "dev", engine_version: str | None = None) -> dict:
217
- """Validated client event -> dataset row, with server-side facts and suspicion flags."""
 
 
218
  d = ev.model_dump()
219
  row = {k: None for k in COLUMNS}
220
  for k, v in d.items():
@@ -223,8 +283,9 @@ def to_row(ev: ClientEvent, catalogue=None, app_version: str = "dev", engine_ver
223
  row["timestamp"] = utc_now_seconds()
224
  row["schema_version"] = SCHEMA_VERSION
225
  row["app_version"] = app_version
226
- row["engine_version"] = engine_version
227
  row["hardware_confirmed"] = ev.hardware_source == "confirmed" if ev.hardware_source else None
 
 
228
  if ev.tried == "yes":
229
  if ev.failure_reason:
230
  row["outcome"] = "problem"
@@ -247,9 +308,11 @@ def to_row(ev: ClientEvent, catalogue=None, app_version: str = "dev", engine_ver
247
  row["hf_downloads_at_selection"] = rec.downloads
248
  row["hf_likes_at_selection"] = rec.likes
249
  row["pipeline"] = rec.pipeline
250
- if ev.compare_models and catalogue is not None:
251
- if any(m not in catalogue.by_id for m in ev.compare_models):
252
- flags.append("unknown_model_in_compare")
 
 
253
 
254
  flags += plausibility_flags(ev, rec)
255
  row["suspicious_flags"] = sorted(set(flags))
@@ -276,13 +339,15 @@ def plausibility_flags(ev: ClientEvent, rec=None) -> list[str]:
276
  ram = ev.reported_ram_gb or ev.hardware_memory_class
277
  if ev.peak_memory_gb and ram and ev.peak_memory_gb > ram * 1.05:
278
  flags.append("peak_memory_exceeds_ram")
279
- if ev.event_type == "mlx_benchmark_submission":
280
  if not (ev.selected_model and (ev.generation_tps or ev.perplexity) and ev.benchmark_version):
281
  flags.append("incomplete_benchmark")
282
  if ev.benchmark_type not in (None, "mlx_lm"):
283
  flags.append("benchmark_type_mismatch")
284
  if ev.event_type == "browser_benchmark" and ev.benchmark_type == "mlx_lm":
285
  flags.append("benchmark_type_mismatch")
 
 
286
  if ev.perplexity is not None and not (ev.eval_dataset and ev.eval_tokens and ev.selected_model):
287
  flags.append("incomplete_quality")
288
  if ev.perplexity is not None and ev.perplexity > 1000:
@@ -292,3 +357,22 @@ def plausibility_flags(ev: ClientEvent, rec=None) -> list[str]:
292
  if ev.failure_reason and ev.quality_rating:
293
  flags.append("conflicting_feedback")
294
  return flags
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
  """Event schema, validation and server-side enrichment.
2
 
3
+ The dataset records outcomes, not clicks:
4
+
5
+ * ``session``: one row per visit (re-sent with the latest state; the newest row per
6
+ session_id wins), holding the final filters, what was searched, hardware class and
7
+ which models were opened, compared or visited.
8
+ * ``feedback``, ``mlx_benchmark``, ``browser_benchmark``: explicit contributions.
9
+
10
  Clients send a small, typed event. Anything outside the schema is rejected
11
+ (``extra="forbid"``), numbers are bounded, free text is capped and scrubbed, and the
12
+ server fills in catalogue facts itself instead of trusting the client. Values that are
13
+ allowed but implausible are kept and flagged, so dataset users can filter them.
14
  """
15
 
16
  from __future__ import annotations
 
22
 
23
  from pydantic import BaseModel, ConfigDict, Field, field_validator
24
 
25
+ from .memory import CONTEXTS, RAM_CLASSES, bits_per_weight
26
  from .parsing import PARAM_BUCKETS, QUANT_BUCKETS
27
 
28
+ SCHEMA_VERSION = 2
29
 
30
+ EventType = Literal["session", "feedback", "mlx_benchmark", "browser_benchmark"]
31
+ # Accepted from older clients and the published benchmark script, stored under the new name.
32
+ ALIASES = {"mlx_benchmark_submission": "mlx_benchmark"}
33
+ # Clickstream events from pages cached before the redesign: accepted and dropped.
34
+ LEGACY_DROPPED = {"search", "filter", "model_view", "model_select", "model_click", "compare", "hardware_test"}
35
 
36
  CHIP_RE = re.compile(r"^Apple M[1-9]( (Pro|Max|Ultra))?$")
37
  MODEL_ID_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,95}/[A-Za-z0-9][A-Za-z0-9._-]{0,95}$")
 
44
  _LONG_DIGITS = re.compile(r"\+?\d[\d ()-]{7,}\d")
45
 
46
  MAX_NOTES = 280
47
+ MAX_BATCH = 20
48
+ QuantBucket = Literal[tuple(QUANT_BUCKETS)] # type: ignore[valid-type]
49
 
50
 
51
  def clean_notes(text: str | None) -> str | None:
 
60
  return text[:MAX_NOTES] or None
61
 
62
 
63
+ def _model_ids(v, cap):
64
+ if v is None:
65
+ return v
66
+ if len(v) > cap:
67
+ raise ValueError(f"at most {cap} model ids")
68
+ for m in v:
69
+ if not MODEL_ID_RE.match(m):
70
+ raise ValueError("invalid model id")
71
+ return list(dict.fromkeys(v)) # de-duplicate, keep order
72
+
73
+
74
  class ClientEvent(BaseModel):
75
  """What the browser (or the benchmark script) is allowed to send."""
76
 
77
  model_config = ConfigDict(extra="forbid", str_max_length=200)
78
 
79
+ event_type: str
80
  session_id: str | None = None
81
 
82
+ # the visit's final query
83
  model_family: str | None = Field(None, max_length=40)
84
  parameter_bucket: Literal[tuple(PARAM_BUCKETS + ["MoE", "unknown"])] | None = None # type: ignore[valid-type]
85
+ quantization: QuantBucket | None = None
86
  target_context: int | None = None
87
  priority: Literal["balanced", "quality", "speed", "memory", "long_context"] | None = None
88
  sort: Literal["recommended", "popular", "recent", "community", "all"] | None = None
89
  result_count: int | None = Field(None, ge=0, le=100_000)
90
+ # everything the visit looked at along the way
91
+ families_searched: list[str] | None = Field(None, max_length=5)
92
+ quants_searched: list[QuantBucket] | None = Field(None, max_length=5)
93
+ distinct_queries: int | None = Field(None, ge=0, le=10_000)
94
 
95
  hardware_source: Literal["detected", "confirmed", "none"] | None = None
96
  hardware_memory_class: int | None = None
 
103
  os_family: Literal["macos", "ios", "windows", "linux", "android", "other"] | None = None
104
  cpu_cores: int | None = Field(None, ge=1, le=256)
105
 
106
+ # what the visit engaged with
107
+ top_model: str | None = None
108
+ top_model_score: float | None = Field(None, ge=0, le=100)
109
+ top_model_fit: Literal["Comfortable", "Likely", "Borderline", "Unlikely"] | None = None
110
+ models_viewed: list[str] | None = None
111
+ models_compared: list[str] | None = None
112
+ models_clicked: list[str] | None = None
113
 
114
+ # feedback and benchmarks are about one model
115
+ selected_model: str | None = None
116
  tried: Literal["yes", "no", "planning"] | None = None
117
  quality_rating: Literal["excellent", "good", "acceptable"] | None = None
118
  failure_reason: Literal["too_slow", "too_much_memory", "low_quality", "didnt_run"] | None = None
 
140
  mlx_lm_version: str | None = None
141
  macos_major: int | None = Field(None, ge=11, le=40)
142
 
143
+ @field_validator("event_type")
144
+ @classmethod
145
+ def _type(cls, v):
146
+ v = ALIASES.get(v, v)
147
+ if v not in ("session", "feedback", "mlx_benchmark", "browser_benchmark") and v not in LEGACY_DROPPED:
148
+ raise ValueError("unknown event_type")
149
+ return v
150
+
151
  @field_validator("session_id")
152
  @classmethod
153
  def _session(cls, v):
 
155
  raise ValueError("session_id must be 16-32 lowercase hex chars")
156
  return v
157
 
158
+ @field_validator("selected_model", "top_model")
159
  @classmethod
160
  def _model(cls, v):
161
  if v is not None and not MODEL_ID_RE.match(v):
162
  raise ValueError("invalid model id")
163
  return v
164
 
165
+ @field_validator("models_viewed")
166
  @classmethod
167
+ def _viewed(cls, v):
168
+ return _model_ids(v, 10)
169
+
170
+ @field_validator("models_compared", "models_clicked")
171
+ @classmethod
172
+ def _compared(cls, v):
173
+ return _model_ids(v, 5)
174
+
175
+ @field_validator("families_searched")
176
+ @classmethod
177
+ def _families(cls, v):
178
  if v is not None:
179
+ for f in v:
180
+ if not SHORT_TOKEN_RE.match(f):
181
+ raise ValueError("invalid family")
182
+ v = list(dict.fromkeys(v))
183
  return v
184
 
185
+ @field_validator("quants_searched")
186
+ @classmethod
187
+ def _quants(cls, v):
188
+ return list(dict.fromkeys(v)) if v is not None else v
189
+
190
  @field_validator("target_context")
191
  @classmethod
192
  def _ctx(cls, v):
 
235
 
236
  # Column order for the dataset. Every row has every column (null when unused).
237
  COLUMNS: dict[str, str] = {
238
+ # every row
239
  "timestamp": "string", "schema_version": "int", "app_version": "string", "event_type": "string",
240
  "session_id": "string",
241
+ # session: final query and what was searched
242
+ "model_family": "string", "parameter_bucket": "string", "quantization": "string",
243
  "target_context": "int", "priority": "string", "sort": "string", "result_count": "int",
244
+ "families_searched": "list", "quants_searched": "list", "distinct_queries": "int",
245
+ # session + benchmarks: hardware class
246
  "hardware_source": "string", "hardware_memory_class": "int", "hardware_confirmed": "bool",
247
  "webgpu_available": "bool", "webgpu_score": "float", "gpu_capability_class": "string",
248
  "gpu_vendor": "string", "gpu_arch": "string", "browser_family": "string", "os_family": "string",
249
  "cpu_cores": "int",
250
+ # session: engagement
251
+ "top_model": "string", "top_model_score": "float", "top_model_fit": "string", "engine_version": "string",
252
+ "models_viewed": "list", "models_compared": "list", "models_clicked": "list",
253
+ # feedback / mlx_benchmark: the model, with catalogue facts filled in by the server
254
+ "selected_model": "string", "model_name": "string", "quant_bits": "float", "quant_mixed": "bool",
255
+ "pipeline": "string", "hf_downloads_at_selection": "int", "hf_likes_at_selection": "int",
256
+ # feedback
257
  "tried": "string", "outcome": "string", "quality_rating": "string", "failure_reason": "string",
258
  "reported_ram_gb": "int", "reported_mac_model": "string", "reported_tokens_per_second": "float",
259
  "reported_context": "int", "notes": "string",
260
+ # benchmarks
261
  "benchmark_type": "string", "benchmark_version": "string", "benchmark_duration_ms": "int",
262
  "prompt_tokens": "int", "generation_tokens": "int", "prompt_tps": "float", "generation_tps": "float",
263
  "ttft_ms": "float", "peak_memory_gb": "float", "chip": "string",
264
  "perplexity": "float", "perplexity_stderr": "float", "eval_dataset": "string", "eval_tokens": "int",
265
+ "mlx_version": "string", "mlx_lm_version": "string", "macos_major": "int",
 
266
  "suspicious_flags": "list",
267
  }
268
 
 
271
  return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
272
 
273
 
274
+ def to_row(ev: ClientEvent, catalogue=None, app_version: str = "dev", engine_version: str | None = None) -> dict | None:
275
+ """Validated client event -> dataset row (None for legacy clickstream events, which are dropped)."""
276
+ if ev.event_type in LEGACY_DROPPED:
277
+ return None
278
  d = ev.model_dump()
279
  row = {k: None for k in COLUMNS}
280
  for k, v in d.items():
 
283
  row["timestamp"] = utc_now_seconds()
284
  row["schema_version"] = SCHEMA_VERSION
285
  row["app_version"] = app_version
 
286
  row["hardware_confirmed"] = ev.hardware_source == "confirmed" if ev.hardware_source else None
287
+ if ev.event_type == "session":
288
+ row["engine_version"] = engine_version
289
  if ev.tried == "yes":
290
  if ev.failure_reason:
291
  row["outcome"] = "problem"
 
308
  row["hf_downloads_at_selection"] = rec.downloads
309
  row["hf_likes_at_selection"] = rec.likes
310
  row["pipeline"] = rec.pipeline
311
+ if catalogue is not None and ev.event_type == "session":
312
+ listed = [m for m in [ev.top_model, *(ev.models_viewed or []), *(ev.models_compared or []),
313
+ *(ev.models_clicked or [])] if m]
314
+ if any(m not in catalogue.by_id for m in listed):
315
+ flags.append("unknown_model_in_session")
316
 
317
  flags += plausibility_flags(ev, rec)
318
  row["suspicious_flags"] = sorted(set(flags))
 
339
  ram = ev.reported_ram_gb or ev.hardware_memory_class
340
  if ev.peak_memory_gb and ram and ev.peak_memory_gb > ram * 1.05:
341
  flags.append("peak_memory_exceeds_ram")
342
+ if ev.event_type == "mlx_benchmark":
343
  if not (ev.selected_model and (ev.generation_tps or ev.perplexity) and ev.benchmark_version):
344
  flags.append("incomplete_benchmark")
345
  if ev.benchmark_type not in (None, "mlx_lm"):
346
  flags.append("benchmark_type_mismatch")
347
  if ev.event_type == "browser_benchmark" and ev.benchmark_type == "mlx_lm":
348
  flags.append("benchmark_type_mismatch")
349
+ if ev.event_type == "feedback" and not ev.selected_model:
350
+ flags.append("incomplete_feedback")
351
  if ev.perplexity is not None and not (ev.eval_dataset and ev.eval_tokens and ev.selected_model):
352
  flags.append("incomplete_quality")
353
  if ev.perplexity is not None and ev.perplexity > 1000:
 
357
  if ev.failure_reason and ev.quality_rating:
358
  flags.append("conflicting_feedback")
359
  return flags
360
+
361
+
362
+ def latest_sessions(rows: list[dict]) -> list[dict]:
363
+ """Collapse re-sent session rows to the newest per session_id; other rows pass through.
364
+
365
+ Timestamps have one-second precision, so rows later in the input win ties (the
366
+ sink and the compactor both preserve arrival order).
367
+ """
368
+ out, latest = [], {}
369
+ for i, r in enumerate(rows):
370
+ if r.get("event_type") == "session" and r.get("session_id"):
371
+ key = r["session_id"]
372
+ prev = latest.get(key)
373
+ if prev is None or (r.get("timestamp") or "") >= (rows[prev].get("timestamp") or ""):
374
+ latest[key] = i
375
+ else:
376
+ out.append(r)
377
+ out.extend(rows[i] for i in sorted(latest.values()))
378
+ return sorted(out, key=lambda r: r.get("timestamp") or "")
app/main.py CHANGED
@@ -270,9 +270,11 @@ def create_app(state: State | None = None) -> FastAPI:
270
  if not st.collection_enabled:
271
  return {"accepted": 0}
272
  rows = [to_row(ev, st.catalogue, APP_VERSION, st.engine.version) for ev in batch.events]
273
- st.sink.add(rows)
274
- st.stats.invalidate()
275
- return {"accepted": len(rows), "flags": [r["suspicious_flags"] for r in rows]}
 
 
276
 
277
  @app.get("/api/stats")
278
  def stats():
 
270
  if not st.collection_enabled:
271
  return {"accepted": 0}
272
  rows = [to_row(ev, st.catalogue, APP_VERSION, st.engine.version) for ev in batch.events]
273
+ kept = [r for r in rows if r is not None] # legacy clickstream events from cached pages are dropped
274
+ if kept:
275
+ st.sink.add(kept)
276
+ st.stats.invalidate()
277
+ return {"accepted": len(kept), "flags": [r["suspicious_flags"] if r else ["dropped_legacy_event"] for r in rows]}
278
 
279
  @app.get("/api/stats")
280
  def stats():
app/sink.py CHANGED
@@ -1,9 +1,11 @@
1
- """Append-only event storage.
2
-
3
- Rows are buffered in memory and mirrored to a local spool file. On each flush the
4
- buffer becomes ONE new Parquet shard (one Hub commit), so a busy day costs at most
5
- ~144 commits at the default 10-minute interval. Existing shards are never
6
- rewritten. If an upload fails, the rows stay spooled and go out with the next flush.
 
 
7
  """
8
 
9
  from __future__ import annotations
@@ -29,6 +31,9 @@ _PA = {"string": pa.string(), "int": pa.int64(), "float": pa.float64(), "bool":
29
  "list": pa.list_(pa.string())}
30
  SCHEMA = pa.schema([(k, _PA[t]) for k, t in COLUMNS.items()])
31
  MAX_BUFFER = 100_000
 
 
 
32
 
33
 
34
  def rows_to_parquet_bytes(rows: list[dict]) -> bytes:
@@ -47,7 +52,7 @@ def shard_path(prefix: str, instance: str, when: datetime | None = None) -> str:
47
  class EventSink:
48
  kind = "base"
49
 
50
- def __init__(self, spool_dir: Path, flush_seconds: float = 600, prefix: str = "data/events"):
51
  self.prefix = prefix
52
  self.flush_seconds = flush_seconds
53
  self.instance = uuid.uuid4().hex[:8]
@@ -94,6 +99,9 @@ class EventSink:
94
  f.write(json.dumps(r) + "\n")
95
  except OSError:
96
  log.warning("spool write failed", exc_info=True)
 
 
 
97
 
98
  @property
99
  def pending(self) -> int:
@@ -183,9 +191,10 @@ class LocalParquetSink(EventSink):
183
 
184
  def read_existing(self) -> list[dict]:
185
  rows = []
186
- for p in sorted((self.root / self.prefix).rglob("*.parquet")):
 
187
  try:
188
- rows.extend(pq.read_table(p).to_pylist())
189
  except Exception:
190
  log.warning("unreadable shard %s", p)
191
  return rows
@@ -216,15 +225,18 @@ class HubParquetSink(EventSink):
216
 
217
  rows = []
218
  try:
219
- files = [f for f in self.api.list_repo_files(self.repo_id, repo_type="dataset")
220
- if f.startswith(self.prefix.strip("/") + "/") and f.endswith(".parquet")]
 
 
 
221
  except Exception as e:
222
  self.last_error = f"read: {type(e).__name__}"
223
  return rows
224
  for f in files:
225
  try:
226
  local = hf_hub_download(self.repo_id, f, repo_type="dataset", token=self.api.token)
227
- rows.extend(pq.read_table(local).to_pylist())
228
  except Exception:
229
  log.warning("unreadable shard %s", f)
230
  return rows
@@ -242,10 +254,15 @@ class NullSink(EventSink):
242
  pass
243
 
244
 
 
 
 
 
 
245
  def make_sink(base_dir: Path) -> EventSink:
246
  mode = os.environ.get("EXPLORER_SINK", "local")
247
- flush = float(os.environ.get("EXPLORER_FLUSH_SECONDS", "600"))
248
- prefix = os.environ.get("EXPLORER_DATA_PREFIX", "data/events")
249
  if mode == "hub":
250
  repo = os.environ.get("DATASET_REPO", "mlx-community/mlx-model-explorer-data")
251
  return HubParquetSink(repo, os.environ.get("HF_TOKEN"), spool_dir=base_dir / "_spool",
 
1
+ """Event storage.
2
+
3
+ The Space buffers rows in memory, mirrored to a local spool file, and writes them as ONE
4
+ new Parquet shard under ``data/incoming/`` per flush: hourly, or sooner once the buffer
5
+ reaches ``FLUSH_ROWS``. That is at most ~24 small files a day. Once a month
6
+ ``scripts/compact.py`` merges a finished month into ``data/events/YYYY-MM.parquet`` and
7
+ removes its incoming shards in the same commit. If an upload fails, the rows stay spooled
8
+ and go out with the next flush.
9
  """
10
 
11
  from __future__ import annotations
 
31
  "list": pa.list_(pa.string())}
32
  SCHEMA = pa.schema([(k, _PA[t]) for k, t in COLUMNS.items()])
33
  MAX_BUFFER = 100_000
34
+ FLUSH_ROWS = int(os.environ.get("EXPLORER_FLUSH_ROWS", "2000"))
35
+ INCOMING = "data/incoming"
36
+ MONTHLY = "data/events"
37
 
38
 
39
  def rows_to_parquet_bytes(rows: list[dict]) -> bytes:
 
52
  class EventSink:
53
  kind = "base"
54
 
55
+ def __init__(self, spool_dir: Path, flush_seconds: float = 3600, prefix: str = INCOMING):
56
  self.prefix = prefix
57
  self.flush_seconds = flush_seconds
58
  self.instance = uuid.uuid4().hex[:8]
 
99
  f.write(json.dumps(r) + "\n")
100
  except OSError:
101
  log.warning("spool write failed", exc_info=True)
102
+ big = len(self._buf) >= FLUSH_ROWS
103
+ if big and not self._flush_lock.locked():
104
+ threading.Thread(target=self.flush, daemon=True, name="sink-size-flush").start()
105
 
106
  @property
107
  def pending(self) -> int:
 
191
 
192
  def read_existing(self) -> list[dict]:
193
  rows = []
194
+ paths = sorted((self.root / MONTHLY).glob("*.parquet")) + sorted((self.root / self.prefix).rglob("*.parquet"))
195
+ for p in dict.fromkeys(paths):
196
  try:
197
+ rows.extend(_aligned(pq.read_table(p).to_pylist()))
198
  except Exception:
199
  log.warning("unreadable shard %s", p)
200
  return rows
 
225
 
226
  rows = []
227
  try:
228
+ all_files = self.api.list_repo_files(self.repo_id, repo_type="dataset")
229
+ monthly = sorted(f for f in all_files if f.startswith(MONTHLY + "/") and f.count("/") == 2
230
+ and f.endswith(".parquet"))
231
+ incoming = sorted(f for f in all_files if f.startswith(self.prefix.strip("/") + "/") and f.endswith(".parquet"))
232
+ files = list(dict.fromkeys(monthly + incoming))
233
  except Exception as e:
234
  self.last_error = f"read: {type(e).__name__}"
235
  return rows
236
  for f in files:
237
  try:
238
  local = hf_hub_download(self.repo_id, f, repo_type="dataset", token=self.api.token)
239
+ rows.extend(_aligned(pq.read_table(local).to_pylist()))
240
  except Exception:
241
  log.warning("unreadable shard %s", f)
242
  return rows
 
254
  pass
255
 
256
 
257
+ def _aligned(rows: list[dict]) -> list[dict]:
258
+ """Rows in the current column set (older shards may lack newer columns)."""
259
+ return [{k: r.get(k) for k in COLUMNS} for r in rows]
260
+
261
+
262
  def make_sink(base_dir: Path) -> EventSink:
263
  mode = os.environ.get("EXPLORER_SINK", "local")
264
+ flush = float(os.environ.get("EXPLORER_FLUSH_SECONDS", "3600"))
265
+ prefix = os.environ.get("EXPLORER_DATA_PREFIX", INCOMING)
266
  if mode == "hub":
267
  repo = os.environ.get("DATASET_REPO", "mlx-community/mlx-model-explorer-data")
268
  return HubParquetSink(repo, os.environ.get("HF_TOKEN"), spool_dir=base_dir / "_spool",
app/stats.py CHANGED
@@ -1,6 +1,7 @@
1
  """Aggregate statistics. Only counts and shares leave this module, never rows.
2
 
3
- Any bucket seen fewer than K_MIN times is folded into "other", so a rare
 
4
  configuration can't be traced back to one visitor.
5
  """
6
 
@@ -12,9 +13,11 @@ import time
12
  from collections import Counter, defaultdict
13
  from typing import Callable, Iterable
14
 
 
15
  from .recommend import EVAL_DATASET, CommunitySignal
16
 
17
  K_MIN = 5
 
18
 
19
 
20
  def _suppress(counter: Counter, k: int = K_MIN, top: int = 15) -> list[dict]:
@@ -30,62 +33,46 @@ def _suppress(counter: Counter, k: int = K_MIN, top: int = 15) -> list[dict]:
30
  return items
31
 
32
 
 
 
 
 
 
33
  def compute(rows: Iterable[dict], k: int = K_MIN) -> dict:
34
- rows = [r for r in rows if not r.get("suspicious_flags")]
 
35
  by_type = Counter(r.get("event_type") for r in rows)
36
- sessions = len({r.get("session_id") for r in rows if r.get("session_id")})
37
- # A "query" is a distinct configuration a session looked at (explicit search, or re-running
38
- # it after changing sort/hardware). Counting each once keeps re-renders from inflating shares.
39
- seen, searches = set(), []
40
- for r in rows:
41
- if r.get("event_type") in ("search", "filter"):
42
- key = (r.get("session_id") or id(r), r.get("model_family"), r.get("parameter_bucket"),
43
- r.get("quantization"), r.get("target_context"), r.get("priority"), r.get("hardware_memory_class"))
44
- if key not in seen:
45
- seen.add(key)
46
- searches.append(r)
47
-
48
- quant = Counter(r.get("quantization") for r in searches if r.get("quantization"))
49
- qtotal = sum(quant.values())
50
- share = lambda b: round(quant.get(b, 0) / qtotal, 4) if qtotal else None
51
- contexts = [r["target_context"] for r in searches if r.get("target_context")]
52
-
53
- hw_seen, hw = set(), []
54
- for r in rows: # one count per session and memory class
55
- if r.get("hardware_memory_class"):
56
- key = (r.get("session_id") or id(r), r["hardware_memory_class"])
57
- if key not in hw_seen:
58
- hw_seen.add(key)
59
- hw.append(r)
60
  return {
61
  "k_min": k,
62
  "totals": {
63
- "events": len(rows),
64
- "sessions": sessions,
65
- "searches": by_type.get("search", 0),
66
- "distinct_queries": len(searches),
67
- "model_views": by_type.get("model_view", 0),
68
- "model_clicks": by_type.get("model_click", 0),
69
- "hardware_tests": by_type.get("hardware_test", 0),
70
  "browser_benchmarks": by_type.get("browser_benchmark", 0),
71
  "feedback": by_type.get("feedback", 0),
72
- "mlx_benchmarks": by_type.get("mlx_benchmark_submission", 0),
73
  },
74
- "quant_share": {"2-bit": share("2-bit"), "3-bit": share("3-bit"), "4-bit": share("4-bit"),
75
- "5-bit": share("5-bit"), "6-bit": share("6-bit"), "8-bit": share("8-bit")},
76
  "average_target_context": round(statistics.mean(contexts)) if len(contexts) >= k else None,
77
- "families": _suppress(Counter(r.get("model_family") for r in searches), k),
78
- "quantizations": _suppress(quant, k),
79
- "sizes": _suppress(Counter(r.get("parameter_bucket") for r in searches), k),
80
- "ram_classes": _suppress(Counter(r.get("hardware_memory_class") for r in hw), k),
81
- "contexts": _suppress(Counter(r.get("target_context") for r in searches), k),
82
- "priorities": _suppress(Counter(r.get("priority") for r in searches), k),
83
- "selected_models": _suppress(Counter(r.get("selected_model") for r in rows
84
- if r.get("event_type") in ("model_select", "model_click")), k),
85
- "gpu_capability": _suppress(Counter(r.get("gpu_capability_class") for r in rows
86
- if r.get("event_type") == "hardware_test"), k),
87
- "webgpu_available": _suppress(Counter(r.get("webgpu_available") for r in rows
88
- if r.get("event_type") == "hardware_test"), k),
89
  }
90
 
91
 
@@ -101,7 +88,7 @@ def community_signals(rows: Iterable[dict]) -> dict[str, CommunitySignal]:
101
  continue
102
  if r.get("event_type") == "feedback" and r.get("tried") == "yes":
103
  fb[m].append(r)
104
- elif r.get("event_type") == "mlx_benchmark_submission":
105
  if r.get("generation_tps"):
106
  bench[m].append(float(r["generation_tps"]))
107
  if r.get("perplexity") and r.get("eval_dataset") == EVAL_DATASET:
 
1
  """Aggregate statistics. Only counts and shares leave this module, never rows.
2
 
3
+ Sessions are counted once (the newest row per session_id), flagged rows are left out,
4
+ and any bucket seen fewer than K_MIN times is folded into "other", so a rare
5
  configuration can't be traced back to one visitor.
6
  """
7
 
 
13
  from collections import Counter, defaultdict
14
  from typing import Callable, Iterable
15
 
16
+ from .events import latest_sessions
17
  from .recommend import EVAL_DATASET, CommunitySignal
18
 
19
  K_MIN = 5
20
+ BENCH_TYPES = ("mlx_benchmark", "mlx_benchmark_submission")
21
 
22
 
23
  def _suppress(counter: Counter, k: int = K_MIN, top: int = 15) -> list[dict]:
 
33
  return items
34
 
35
 
36
+ def _present(sessions: list[dict], key: str) -> Counter:
37
+ """Shares among visits that set the field (unset isn't a category)."""
38
+ return Counter(s[key] for s in sessions if s.get(key) is not None)
39
+
40
+
41
  def compute(rows: Iterable[dict], k: int = K_MIN) -> dict:
42
+ rows = latest_sessions([r for r in rows if not r.get("suspicious_flags")])
43
+ sessions = [r for r in rows if r.get("event_type") == "session"]
44
  by_type = Counter(r.get("event_type") for r in rows)
45
+
46
+ quants = Counter(q for s in sessions for q in (s.get("quants_searched") or []))
47
+ qtotal = sum(v for key, v in quants.items() if key != "unknown")
48
+ share = lambda b: round(quants.get(b, 0) / qtotal, 4) if qtotal else None
49
+ contexts = [s["target_context"] for s in sessions if s.get("target_context")]
50
+ engaged = [s for s in sessions if s.get("models_viewed") or s.get("models_compared") or s.get("models_clicked")]
51
+ picked = Counter(m for s in sessions for m in dict.fromkeys(
52
+ [*(s.get("models_compared") or []), *(s.get("models_clicked") or [])]))
53
+ hw = [s for s in sessions if s.get("gpu_capability_class")]
54
+
 
 
 
 
 
 
 
 
 
 
 
 
 
 
55
  return {
56
  "k_min": k,
57
  "totals": {
58
+ "sessions": len(sessions),
59
+ "engaged_sessions": len(engaged),
60
+ "hardware_checks": len(hw),
 
 
 
 
61
  "browser_benchmarks": by_type.get("browser_benchmark", 0),
62
  "feedback": by_type.get("feedback", 0),
63
+ "mlx_benchmarks": sum(by_type.get(t, 0) for t in BENCH_TYPES),
64
  },
65
+ "quant_share": {b: share(b) for b in ("2-bit", "3-bit", "4-bit", "5-bit", "6-bit", "8-bit")},
 
66
  "average_target_context": round(statistics.mean(contexts)) if len(contexts) >= k else None,
67
+ "families": _suppress(Counter(f for s in sessions for f in (s.get("families_searched") or [])), k),
68
+ "quantizations": _suppress(quants, k),
69
+ "sizes": _suppress(_present(sessions, "parameter_bucket"), k),
70
+ "ram_classes": _suppress(_present(sessions, "hardware_memory_class"), k),
71
+ "contexts": _suppress(_present(sessions, "target_context"), k),
72
+ "priorities": _suppress(_present(sessions, "priority"), k),
73
+ "selected_models": _suppress(picked, k),
74
+ "gpu_capability": _suppress(Counter(s.get("gpu_capability_class") for s in hw), k),
75
+ "webgpu_available": _suppress(Counter(s.get("webgpu_available") for s in hw), k),
 
 
 
76
  }
77
 
78
 
 
88
  continue
89
  if r.get("event_type") == "feedback" and r.get("tried") == "yes":
90
  fb[m].append(r)
91
+ elif r.get("event_type") in BENCH_TYPES:
92
  if r.get("generation_tps"):
93
  bench[m].append(float(r["generation_tps"]))
94
  if r.get("perplexity") and r.get("eval_dataset") == EVAL_DATASET:
bench/mlx_explorer_bench.py CHANGED
@@ -155,7 +155,7 @@ def run(args) -> dict:
155
  med = lambda xs: statistics.median(xs)
156
  mac = platform.mac_ver()[0]
157
  return result | {
158
- "event_type": "mlx_benchmark_submission",
159
  "benchmark_type": "mlx_lm",
160
  "benchmark_version": BENCHMARK_VERSION,
161
  "selected_model": args.model,
 
155
  med = lambda xs: statistics.median(xs)
156
  mac = platform.mac_ver()[0]
157
  return result | {
158
+ "event_type": "mlx_benchmark",
159
  "benchmark_type": "mlx_lm",
160
  "benchmark_version": BENCHMARK_VERSION,
161
  "selected_model": args.model,
static/app.js CHANGED
@@ -76,23 +76,35 @@
76
  sid = [...b].map((x) => x.toString(16).padStart(2, "0")).join("");
77
  try { sessionStorage.setItem("mme-sid", sid); } catch (e) {}
78
  }
79
- let queue = [];
80
  let enabled = true;
 
 
81
 
82
- function send(beacon) {
83
- if (!queue.length || optedOut || !enabled) { queue = []; return; }
84
- const batch = queue.splice(0, 50);
85
- const body = JSON.stringify({ events: batch });
86
  if (beacon && navigator.sendBeacon) {
87
  navigator.sendBeacon("/api/events", new Blob([body], { type: "application/json" }));
88
  } else {
89
  fetch("/api/events", { method: "POST", headers: { "content-type": "application/json" }, body, keepalive: true }).catch(() => {});
90
  }
91
- if (queue.length) send(beacon);
92
  }
93
- setInterval(() => send(false), 10000);
94
- document.addEventListener("visibilitychange", () => { if (document.visibilityState === "hidden") send(true); });
95
- window.addEventListener("pagehide", () => send(true));
 
 
 
 
 
 
 
 
 
 
 
 
 
 
96
 
97
  return {
98
  gpc,
@@ -103,15 +115,16 @@
103
  if (v) { localStorage.setItem("mme-optout", "1"); localStorage.removeItem("mme-optin"); }
104
  else { localStorage.removeItem("mme-optout"); if (gpc) localStorage.setItem("mme-optin", "1"); }
105
  } catch (e) {}
106
- if (optedOut) queue = [];
107
  },
108
  setEnabled(v) { enabled = v; },
109
- track(type, fields) {
 
 
 
110
  if (optedOut || !enabled) return;
111
  const ev = { event_type: type, session_id: sid };
112
  for (const [k, v] of Object.entries(fields || {})) if (v !== undefined && v !== null && v !== "") ev[k] = v;
113
- queue.push(ev);
114
- if (queue.length >= 20) send(false);
115
  },
116
  // Submissions the user explicitly asked to send go out immediately, with a result.
117
  async submit(type, fields) {
@@ -154,6 +167,43 @@
154
  $$(`#query input[name="${name}"]`).forEach((x) => (x.checked = false));
155
  }
156
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
157
  function currentRam() {
158
  const v = parseInt(radioValue("ram"), 10);
159
  return Number.isFinite(v) ? v : null;
@@ -350,12 +400,7 @@
350
  state.results = append ? state.results.concat(data.results) : data.results;
351
  renderResults(data, append);
352
  syncUrl(q);
353
- if (!append) {
354
- Telemetry.track(reason, {
355
- model_family: q.family, parameter_bucket: q.size, quantization: q.quant, target_context: q.context,
356
- priority: q.priority, sort: q.sort, result_count: data.total, ...hwFields(),
357
- });
358
- }
359
  } catch (e) {
360
  if (id !== state.requestId) return;
361
  $("#summary").textContent = "";
@@ -439,13 +484,8 @@
439
  return r.reasons.find((t) => !skip.test(t)) || r.reasons[0] || "";
440
  }
441
 
442
- function shown(r) {
443
- return r ? { recommendation_score: Math.round(r.score * 10) / 10, fit_class: r.memory && r.memory.fit } : {};
444
- }
445
-
446
- function trackClick(id, rank, r) {
447
- const q = state.lastQuery || readQuery();
448
- Telemetry.track("model_click", { selected_model: id, selected_model_rank: rank, target_context: q.context, priority: q.priority, ...shown(r), ...hwFields() });
449
  }
450
 
451
  // ------------------------------------------------------------------ compare
@@ -457,8 +497,7 @@
457
  return;
458
  }
459
  state.compare.set(r.model.id, r);
460
- const q = state.lastQuery || readQuery();
461
- Telemetry.track("model_select", { selected_model: r.model.id, selected_model_rank: rank, target_context: q.context, priority: q.priority, ...shown(r), ...hwFields() });
462
  } else {
463
  state.compare.delete(r.model.id);
464
  }
@@ -492,7 +531,6 @@
492
  function openCompare() {
493
  const items = [...state.compare.values()];
494
  const q = state.lastQuery || readQuery();
495
- Telemetry.track("compare", { compare_models: items.map((x) => x.model.id), target_context: q.context, priority: q.priority, ...hwFields() });
496
  const scale = scaleFor(q.ram_gb, items.map((r) => r.memory.total_gb));
497
  const qt = (r) => r.quant || r.model.quant;
498
  const rows = [
@@ -527,7 +565,7 @@
527
  $("#d-title").textContent = id.split("/")[1];
528
  $("#d-body").replaceChildren(el("p", { class: "fine" }, "Loading model details"));
529
  if (!dlg.open) dlg.showModal();
530
- Telemetry.track("model_view", { selected_model: id, selected_model_rank: rank ?? null, target_context: q.context, priority: q.priority, ...shown(r), ...hwFields() });
531
  const params = new URLSearchParams({ context: q.context, priority: q.priority });
532
  if (q.ram_gb) { params.set("ram_gb", q.ram_gb); if (q.ram_source) params.set("ram_source", q.ram_source); }
533
  let d;
@@ -620,6 +658,7 @@
620
  const tried = ($("input[name=tried]:checked", f) || {}).value;
621
  if (!tried) { out.textContent = "Choose Yes, No or Planning to."; return; }
622
  const fields = { selected_model: m.id, tried, target_context: q.context, priority: q.priority, ...hwFields() };
 
623
  if (tried === "yes") {
624
  const oc = ($("input[name=outcome]:checked", f) || {}).value;
625
  if (oc && oc.startsWith("q:")) fields.quality_rating = oc.slice(2);
@@ -678,12 +717,6 @@
678
  }
679
  if (msgs.length) { $("#hw-error").textContent = msgs.join(" "); $("#hw-error").hidden = false; }
680
  $("#full-bench").disabled = !hw.webgpu_available;
681
- Telemetry.track("hardware_test", {
682
- ...hwFields(), hardware_source: "detected", webgpu_score: hw.quick_score,
683
- benchmark_type: hw.quick_score != null ? "webgpu_quick" : null,
684
- benchmark_version: hw.quick_score != null ? window.HW.BENCH_VERSION : null,
685
- benchmark_duration_ms: hw.duration_ms,
686
- });
687
  explore({ reason: "filter" });
688
  }
689
 
@@ -697,7 +730,7 @@
697
  try {
698
  const res = await window.HW.fullBenchmark((p) => { prog.value = p; });
699
  $("#bench-out").textContent = `WebGPU compute score: ${res.score.toLocaleString()}. A relative browser GPU score for grouping hardware, not MLX speed or tokens per second.`;
700
- Telemetry.track("browser_benchmark", {
701
  ...hwFields(), webgpu_score: res.score, benchmark_type: "webgpu_full", benchmark_version: res.version,
702
  benchmark_duration_ms: res.duration_ms, gpu_capability_class: window.HW.capabilityClass(state.hw && state.hw.quick_score),
703
  });
@@ -742,7 +775,7 @@
742
  if (!obj || typeof obj !== "object" || Array.isArray(obj)) { out.textContent = "Paste a single JSON object, the block the script printed."; return; }
743
  const fields = {};
744
  for (const k of BENCH_KEYS) if (k in obj) fields[k] = obj[k];
745
- const res = await Telemetry.submit("mlx_benchmark_submission", fields);
746
  out.textContent = res.ok
747
  ? res.flags.length ? `Result recorded and flagged for review: ${res.flags.join(", ")}.` : "Result recorded anonymously. Thank you."
748
  : res.message;
 
76
  sid = [...b].map((x) => x.toString(16).padStart(2, "0")).join("");
77
  try { sessionStorage.setItem("mme-sid", sid); } catch (e) {}
78
  }
 
79
  let enabled = true;
80
+ let buildSession = () => null; // set by the app: returns the visit's current summary or null
81
+ let lastSent = "";
82
 
83
+ function post(events, beacon) {
84
+ const body = JSON.stringify({ events });
 
 
85
  if (beacon && navigator.sendBeacon) {
86
  navigator.sendBeacon("/api/events", new Blob([body], { type: "application/json" }));
87
  } else {
88
  fetch("/api/events", { method: "POST", headers: { "content-type": "application/json" }, body, keepalive: true }).catch(() => {});
89
  }
 
90
  }
91
+
92
+ // One row per visit: re-sent only when the summary changed since the last send.
93
+ function sendSession(beacon) {
94
+ if (optedOut || !enabled) return;
95
+ const summary = buildSession();
96
+ if (!summary) return;
97
+ const ev = { event_type: "session", session_id: sid };
98
+ for (const [k, v] of Object.entries(summary)) {
99
+ if (v !== undefined && v !== null && v !== "" && !(Array.isArray(v) && !v.length)) ev[k] = v;
100
+ }
101
+ const key = JSON.stringify(ev);
102
+ if (key === lastSent) return;
103
+ lastSent = key;
104
+ post([ev], beacon);
105
+ }
106
+ document.addEventListener("visibilitychange", () => { if (document.visibilityState === "hidden") sendSession(true); });
107
+ window.addEventListener("pagehide", () => sendSession(true));
108
 
109
  return {
110
  gpc,
 
115
  if (v) { localStorage.setItem("mme-optout", "1"); localStorage.removeItem("mme-optin"); }
116
  else { localStorage.removeItem("mme-optout"); if (gpc) localStorage.setItem("mme-optin", "1"); }
117
  } catch (e) {}
 
118
  },
119
  setEnabled(v) { enabled = v; },
120
+ setSessionBuilder(fn) { buildSession = fn; },
121
+ sendSession,
122
+ // A contribution the page makes on the user's behalf (the GPU test): one row, sent now.
123
+ contribute(type, fields) {
124
  if (optedOut || !enabled) return;
125
  const ev = { event_type: type, session_id: sid };
126
  for (const [k, v] of Object.entries(fields || {})) if (v !== undefined && v !== null && v !== "") ev[k] = v;
127
+ post([ev], false);
 
128
  },
129
  // Submissions the user explicitly asked to send go out immediately, with a result.
130
  async submit(type, fields) {
 
167
  $$(`#query input[name="${name}"]`).forEach((x) => (x.checked = false));
168
  }
169
 
170
+ // ------------------------------------------------------------------ visit summary (the only analytics row)
171
+ const visit = { families: [], quants: [], queries: new Set(), query: null, top: null,
172
+ viewed: [], compared: [], clicked: [], engagedSent: false };
173
+ const pushCapped = (arr, v, cap) => { if (v && !arr.includes(v)) { arr.push(v); if (arr.length > cap) arr.shift(); } };
174
+
175
+ function noteQuery(q, data) {
176
+ visit.query = { ...q, result_count: data.total };
177
+ if (q.family) pushCapped(visit.families, q.family, 5);
178
+ if (q.quant) pushCapped(visit.quants, q.quant, 5);
179
+ visit.queries.add([q.family, q.size, q.quant, q.context, q.priority, q.ram_gb].join("|"));
180
+ const top = data.results && data.results[0];
181
+ if (top && q.sort === "recommended") {
182
+ visit.top = { id: top.model.id, score: Math.round(top.score * 10) / 10, fit: top.memory && top.memory.fit };
183
+ }
184
+ }
185
+
186
+ function noteModel(kind, id) {
187
+ pushCapped(visit[kind], id, kind === "viewed" ? 10 : 5);
188
+ if (!visit.engagedSent) { // the first real interaction sends the summary right away
189
+ visit.engagedSent = true;
190
+ Telemetry.sendSession(false);
191
+ }
192
+ }
193
+
194
+ Telemetry.setSessionBuilder(() => {
195
+ const q = visit.query;
196
+ if (!q) return null;
197
+ return {
198
+ model_family: q.family, parameter_bucket: q.size, quantization: q.quant, target_context: q.context,
199
+ priority: q.priority, sort: q.sort, result_count: q.result_count,
200
+ families_searched: visit.families, quants_searched: visit.quants, distinct_queries: visit.queries.size,
201
+ ...hwFields(), webgpu_score: state.hw ? state.hw.quick_score : null,
202
+ top_model: visit.top && visit.top.id, top_model_score: visit.top && visit.top.score, top_model_fit: visit.top && visit.top.fit,
203
+ models_viewed: visit.viewed, models_compared: visit.compared, models_clicked: visit.clicked,
204
+ };
205
+ });
206
+
207
  function currentRam() {
208
  const v = parseInt(radioValue("ram"), 10);
209
  return Number.isFinite(v) ? v : null;
 
400
  state.results = append ? state.results.concat(data.results) : data.results;
401
  renderResults(data, append);
402
  syncUrl(q);
403
+ if (!append) noteQuery(q, data);
 
 
 
 
 
404
  } catch (e) {
405
  if (id !== state.requestId) return;
406
  $("#summary").textContent = "";
 
484
  return r.reasons.find((t) => !skip.test(t)) || r.reasons[0] || "";
485
  }
486
 
487
+ function trackClick(id) {
488
+ noteModel("clicked", id);
 
 
 
 
 
489
  }
490
 
491
  // ------------------------------------------------------------------ compare
 
497
  return;
498
  }
499
  state.compare.set(r.model.id, r);
500
+ noteModel("compared", r.model.id);
 
501
  } else {
502
  state.compare.delete(r.model.id);
503
  }
 
531
  function openCompare() {
532
  const items = [...state.compare.values()];
533
  const q = state.lastQuery || readQuery();
 
534
  const scale = scaleFor(q.ram_gb, items.map((r) => r.memory.total_gb));
535
  const qt = (r) => r.quant || r.model.quant;
536
  const rows = [
 
565
  $("#d-title").textContent = id.split("/")[1];
566
  $("#d-body").replaceChildren(el("p", { class: "fine" }, "Loading model details"));
567
  if (!dlg.open) dlg.showModal();
568
+ noteModel("viewed", id);
569
  const params = new URLSearchParams({ context: q.context, priority: q.priority });
570
  if (q.ram_gb) { params.set("ram_gb", q.ram_gb); if (q.ram_source) params.set("ram_source", q.ram_source); }
571
  let d;
 
658
  const tried = ($("input[name=tried]:checked", f) || {}).value;
659
  if (!tried) { out.textContent = "Choose Yes, No or Planning to."; return; }
660
  const fields = { selected_model: m.id, tried, target_context: q.context, priority: q.priority, ...hwFields() };
661
+ delete fields.cpu_cores;
662
  if (tried === "yes") {
663
  const oc = ($("input[name=outcome]:checked", f) || {}).value;
664
  if (oc && oc.startsWith("q:")) fields.quality_rating = oc.slice(2);
 
717
  }
718
  if (msgs.length) { $("#hw-error").textContent = msgs.join(" "); $("#hw-error").hidden = false; }
719
  $("#full-bench").disabled = !hw.webgpu_available;
 
 
 
 
 
 
720
  explore({ reason: "filter" });
721
  }
722
 
 
730
  try {
731
  const res = await window.HW.fullBenchmark((p) => { prog.value = p; });
732
  $("#bench-out").textContent = `WebGPU compute score: ${res.score.toLocaleString()}. A relative browser GPU score for grouping hardware, not MLX speed or tokens per second.`;
733
+ Telemetry.contribute("browser_benchmark", {
734
  ...hwFields(), webgpu_score: res.score, benchmark_type: "webgpu_full", benchmark_version: res.version,
735
  benchmark_duration_ms: res.duration_ms, gpu_capability_class: window.HW.capabilityClass(state.hw && state.hw.quick_score),
736
  });
 
775
  if (!obj || typeof obj !== "object" || Array.isArray(obj)) { out.textContent = "Paste a single JSON object, the block the script printed."; return; }
776
  const fields = {};
777
  for (const k of BENCH_KEYS) if (k in obj) fields[k] = obj[k];
778
+ const res = await Telemetry.submit("mlx_benchmark", fields);
779
  out.textContent = res.ok
780
  ? res.flags.length ? `Result recorded and flagged for review: ${res.flags.join(", ")}.` : "Result recorded anonymously. Thank you."
781
  : res.message;
static/index.html CHANGED
@@ -151,7 +151,7 @@
151
  </div>
152
  <details>
153
  <summary>Paste a result you already ran</summary>
154
- <textarea id="bench-json" rows="6" maxlength="4000" spellcheck="false" placeholder='{"event_type": "mlx_benchmark_submission", ...}'></textarea>
155
  <button type="button" id="bench-submit" class="solid">Submit result</button>
156
  <p id="bench-submit-out" class="fine" aria-live="polite"></p>
157
  </details>
@@ -194,7 +194,7 @@
194
  <h2 id="priv-h">Privacy and data</h2>
195
  <p>We collect anonymous model-selection and optional benchmark data to improve MLX Model Explorer and community recommendations. It's published as the <a id="dataset-link" href="https://huggingface.co/datasets/mlx-community/mlx-model-explorer-data" target="_blank" rel="noopener">mlx-model-explorer-data</a> dataset.</p>
196
  <ul>
197
- <li><strong>Collected:</strong> the filters you choose, models you open or compare, coarse hardware class (GPU vendor and architecture, core count, browser and OS family), test scores, and feedback you send.</li>
198
  <li><strong>Never collected:</strong> name, email, IP address, location, cookies, user-agent strings or device fingerprints. Links, emails and phone numbers are stripped from notes.</li>
199
  <li><strong>Session ID:</strong> random, lives only in this tab, and is gone when you close it.</li>
200
  <li>The <a href="/stats">stats page</a> only shows groups of five or more.</li>
 
151
  </div>
152
  <details>
153
  <summary>Paste a result you already ran</summary>
154
+ <textarea id="bench-json" rows="6" maxlength="4000" spellcheck="false" placeholder='{"event_type": "mlx_benchmark", ...}'></textarea>
155
  <button type="button" id="bench-submit" class="solid">Submit result</button>
156
  <p id="bench-submit-out" class="fine" aria-live="polite"></p>
157
  </details>
 
194
  <h2 id="priv-h">Privacy and data</h2>
195
  <p>We collect anonymous model-selection and optional benchmark data to improve MLX Model Explorer and community recommendations. It's published as the <a id="dataset-link" href="https://huggingface.co/datasets/mlx-community/mlx-model-explorer-data" target="_blank" rel="noopener">mlx-model-explorer-data</a> dataset.</p>
196
  <ul>
197
+ <li><strong>Collected:</strong> one summary per visit (the filters you ended on, the families and quantizations you looked at, which models you opened, compared or visited, and a coarse hardware class: memory, GPU vendor and architecture, core count, browser and OS family), plus test scores and feedback you choose to send. Individual clicks aren't recorded.</li>
198
  <li><strong>Never collected:</strong> name, email, IP address, location, cookies, user-agent strings or device fingerprints. Links, emails and phone numbers are stripped from notes.</li>
199
  <li><strong>Session ID:</strong> random, lives only in this tab, and is gone when you close it.</li>
200
  <li>The <a href="/stats">stats page</a> only shows groups of five or more.</li>
static/stats.html CHANGED
@@ -24,11 +24,11 @@
24
  <main class="shell">
25
  <section class="stats-hero">
26
  <h1>What the MLX community is trying to run</h1>
27
- <p class="lede">Anonymous totals from MLX Model Explorer. Groups smaller than <span id="kmin">5</span> are combined into "other", flagged submissions are left out, and no individual record is shown.</p>
28
  </section>
29
  <div class="totals" id="totals"></div>
30
  <section class="chart" aria-labelledby="q-h">
31
- <h2 id="q-h">Quantization people search for</h2>
32
  <div class="qshare" id="qshare" role="img" aria-label="Quantization share"></div>
33
  <p class="scale-note" id="qshare-note"></p>
34
  <p class="scale-note" id="avgctx"></p>
 
24
  <main class="shell">
25
  <section class="stats-hero">
26
  <h1>What the MLX community is trying to run</h1>
27
+ <p class="lede">Anonymous totals from MLX Model Explorer. Each visit counts once. Groups smaller than <span id="kmin">5</span> are combined into "other", flagged submissions are left out, and no individual record is shown.</p>
28
  </section>
29
  <div class="totals" id="totals"></div>
30
  <section class="chart" aria-labelledby="q-h">
31
+ <h2 id="q-h">Quantizations people look at</h2>
32
  <div class="qshare" id="qshare" role="img" aria-label="Quantization share"></div>
33
  <p class="scale-note" id="qshare-note"></p>
34
  <p class="scale-note" id="avgctx"></p>
static/stats.js CHANGED
@@ -9,13 +9,12 @@
9
  return n;
10
  }
11
  const LABELS = {
12
- events: "Events", sessions: "Sessions", searches: "Searches", distinct_queries: "Distinct queries", model_views: "Models opened",
13
- model_clicks: "Hugging Face visits", hardware_tests: "Hardware checks", browser_benchmarks: "GPU tests",
14
- feedback: "Reports", mlx_benchmarks: "MLX benchmarks",
15
  };
16
  const CHARTS = [
17
  ["families", "Model families"], ["sizes", "Model sizes"], ["ram_classes", "Mac memory (GB)"],
18
- ["contexts", "Target context"], ["priorities", "What matters most"], ["selected_models", "Most compared and visited models"],
19
  ["gpu_capability", "Browser GPU class"], ["webgpu_available", "WebGPU available"],
20
  ];
21
  const QCOLORS = { "2-bit": "#8a4f9e", "3-bit": "#7a5ab8", "4-bit": "#2f5d8c", "5-bit": "#3f7f86", "6-bit": "#2e7d4f", "8-bit": "#62802a" };
@@ -68,7 +67,7 @@
68
  seg.title = `${k}: ${Math.round(v * 100)}%`;
69
  return seg;
70
  }));
71
- $("#qshare-note").textContent = q.length ? q.map(([k, v]) => `${k} ${Math.round(v * 100)}%`).join(", ") : "Not enough searches with a quantization filter yet.";
72
  $("#avgctx").textContent = s.average_target_context ? `Average target context: ${Math.round(s.average_target_context / 1024)}k tokens.` : "";
73
  $("#charts").replaceChildren(...CHARTS.map(([key, title]) => {
74
  const c = el("section", "chart");
 
9
  return n;
10
  }
11
  const LABELS = {
12
+ sessions: "Visits", engaged_sessions: "Visits that opened a model", hardware_checks: "Hardware checks",
13
+ browser_benchmarks: "GPU tests", feedback: "Reports", mlx_benchmarks: "MLX benchmarks",
 
14
  };
15
  const CHARTS = [
16
  ["families", "Model families"], ["sizes", "Model sizes"], ["ram_classes", "Mac memory (GB)"],
17
+ ["contexts", "Target context"], ["priorities", "What matters most"], ["selected_models", "Most compared or visited models"],
18
  ["gpu_capability", "Browser GPU class"], ["webgpu_available", "WebGPU available"],
19
  ];
20
  const QCOLORS = { "2-bit": "#8a4f9e", "3-bit": "#7a5ab8", "4-bit": "#2f5d8c", "5-bit": "#3f7f86", "6-bit": "#2e7d4f", "8-bit": "#62802a" };
 
67
  seg.title = `${k}: ${Math.round(v * 100)}%`;
68
  return seg;
69
  }));
70
+ $("#qshare-note").textContent = q.length ? q.map(([k, v]) => `${k} ${Math.round(v * 100)}%`).join(", ") : "Not enough visits with a quantization filter yet.";
71
  $("#avgctx").textContent = s.average_target_context ? `Average target context: ${Math.round(s.average_target_context / 1024)}k tokens.` : "";
72
  $("#charts").replaceChildren(...CHARTS.map(([key, title]) => {
73
  const c = el("section", "chart");
tests/test_api.py CHANGED
@@ -110,22 +110,28 @@ def test_model_detail(client):
110
  def test_events_ingest_to_parquet(client, tmp_path):
111
  sid = "abcdef0123456789"
112
  batch = {"events": [
113
- {"event_type": "search", "session_id": sid, "model_family": "Qwen", "quantization": "4-bit",
114
- "target_context": 32768, "priority": "balanced", "hardware_source": "confirmed", "hardware_memory_class": 36},
115
- {"event_type": "model_select", "session_id": sid, "selected_model": "mlx-community/Qwen3-8B-4bit",
116
- "selected_model_rank": 0, "target_context": 32768, "hardware_memory_class": 36},
117
  {"event_type": "feedback", "session_id": sid, "selected_model": "mlx-community/Qwen3-8B-4bit",
118
  "tried": "yes", "quality_rating": "good", "notes": "worked fine, email me x@y.com"},
 
119
  ]}
120
  r = client.post("/api/events", json=batch)
121
- assert r.status_code == 202 and r.json()["accepted"] == 3
 
 
 
 
122
  assert client.st.sink.flush()
123
  rows = client.st.sink.read_existing()
124
- assert [x["event_type"] for x in rows] == ["search", "model_select", "feedback"]
125
- assert rows[1]["hf_downloads_at_selection"] == 90000
126
- assert "x@y.com" not in rows[2]["notes"]
127
  stats = client.get("/api/stats").json()
128
- assert stats["totals"]["events"] == 3 and stats["totals"]["sessions"] == 1
 
129
  assert stats["families"] == [{"key": "other", "count": 1, "share": 1.0}] # k-anonymity folds tiny buckets
130
 
131
 
@@ -166,19 +172,21 @@ def test_collection_disabled(client):
166
  assert r.json() == {"accepted": 0} and client.st.sink.pending == 0
167
 
168
 
169
- def test_stats_dedupe_and_suppression():
170
  from app.stats import compute
171
  rows = []
172
  for i in range(6):
173
  sid = f"{i:016x}"
174
- for _ in range(3): # the same query re-rendered three times counts once
175
- rows.append({"event_type": "filter", "session_id": sid, "model_family": "Qwen", "quantization": "4-bit",
176
- "target_context": 8192, "hardware_memory_class": 36, "suspicious_flags": []})
177
- rows.append({"event_type": "search", "session_id": "ffffffffffffffff", "model_family": "Rare",
178
- "quantization": "3-bit", "suspicious_flags": []})
179
- rows.append({"event_type": "search", "model_family": "Flagged", "suspicious_flags": ["unknown_model"]})
 
 
180
  s = compute(rows)
181
- assert s["totals"]["distinct_queries"] == 7
182
  assert s["families"] == [{"key": "Qwen", "count": 6, "share": 0.8571}, {"key": "other", "count": 1, "share": 0.1429}]
183
  assert s["ram_classes"][0] == {"key": 36, "count": 6, "share": 1.0}
184
- assert all(f["key"] != "Flagged" for f in s["families"])
 
110
  def test_events_ingest_to_parquet(client, tmp_path):
111
  sid = "abcdef0123456789"
112
  batch = {"events": [
113
+ {"event_type": "session", "session_id": sid, "model_family": "Qwen", "quantization": "4-bit",
114
+ "target_context": 32768, "priority": "balanced", "hardware_source": "confirmed", "hardware_memory_class": 36,
115
+ "families_searched": ["Qwen"], "quants_searched": ["4-bit"], "distinct_queries": 2,
116
+ "models_viewed": ["mlx-community/Qwen3-8B-4bit"]},
117
  {"event_type": "feedback", "session_id": sid, "selected_model": "mlx-community/Qwen3-8B-4bit",
118
  "tried": "yes", "quality_rating": "good", "notes": "worked fine, email me x@y.com"},
119
+ {"event_type": "search", "session_id": sid, "model_family": "Qwen"},
120
  ]}
121
  r = client.post("/api/events", json=batch)
122
+ assert r.status_code == 202 and r.json()["accepted"] == 2
123
+ assert r.json()["flags"][2] == ["dropped_legacy_event"]
124
+ # the visit re-sends its summary later: the newest row wins in stats
125
+ client.post("/api/events", json={"events": [dict(batch["events"][0], distinct_queries=5,
126
+ models_compared=["mlx-community/Qwen3-14B-4bit"])]})
127
  assert client.st.sink.flush()
128
  rows = client.st.sink.read_existing()
129
+ assert [x["event_type"] for x in rows] == ["session", "feedback", "session"]
130
+ assert rows[1]["hf_downloads_at_selection"] == 90000 and "x@y.com" not in rows[1]["notes"]
131
+ client.st.stats.invalidate()
132
  stats = client.get("/api/stats").json()
133
+ assert stats["totals"]["sessions"] == 1 and stats["totals"]["engaged_sessions"] == 1
134
+ assert stats["totals"]["feedback"] == 1
135
  assert stats["families"] == [{"key": "other", "count": 1, "share": 1.0}] # k-anonymity folds tiny buckets
136
 
137
 
 
172
  assert r.json() == {"accepted": 0} and client.st.sink.pending == 0
173
 
174
 
175
+ def test_stats_one_count_per_visit_and_suppression():
176
  from app.stats import compute
177
  rows = []
178
  for i in range(6):
179
  sid = f"{i:016x}"
180
+ for n in range(3): # each visit re-sent its summary three times
181
+ rows.append({"event_type": "session", "session_id": sid, "timestamp": f"2026-09-14T01:0{n}:00Z",
182
+ "families_searched": ["Qwen"], "quants_searched": ["4-bit"], "target_context": 8192,
183
+ "hardware_memory_class": 36, "suspicious_flags": []})
184
+ rows.append({"event_type": "session", "session_id": "f" * 16, "timestamp": "2026-09-14T02:00:00Z",
185
+ "families_searched": ["Rare"], "quants_searched": ["3-bit"], "suspicious_flags": []})
186
+ rows.append({"event_type": "session", "session_id": "e" * 16, "families_searched": ["Flagged"],
187
+ "suspicious_flags": ["unknown_model_in_session"]})
188
  s = compute(rows)
189
+ assert s["totals"]["sessions"] == 7
190
  assert s["families"] == [{"key": "Qwen", "count": 6, "share": 0.8571}, {"key": "other", "count": 1, "share": 0.1429}]
191
  assert s["ram_classes"][0] == {"key": 36, "count": 6, "share": 1.0}
192
+ assert s["quant_share"]["4-bit"] == 0.8571
tests/test_compact.py ADDED
@@ -0,0 +1,43 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import importlib.util
2
+ from pathlib import Path
3
+
4
+ from app.events import COLUMNS
5
+
6
+ spec = importlib.util.spec_from_file_location("compact", Path(__file__).resolve().parent.parent / "scripts" / "compact.py")
7
+ compact = importlib.util.module_from_spec(spec)
8
+ spec.loader.exec_module(compact)
9
+
10
+
11
+ def test_plan_selects_only_that_month():
12
+ files = [
13
+ "README.md",
14
+ "data/incoming/2026/09/01/a-1.parquet", "data/incoming/2026/09/30/b-2.parquet",
15
+ "data/incoming/2026/10/01/c-3.parquet", "data/events/2026-08.parquet", "data/events/2026-09.parquet",
16
+ ]
17
+ incoming, existing, target = compact.plan(files, "2026-09")
18
+ assert incoming == ["data/incoming/2026/09/01/a-1.parquet", "data/incoming/2026/09/30/b-2.parquet"]
19
+ assert existing == "data/events/2026-09.parquet" and target == "data/events/2026-09.parquet"
20
+ assert compact.plan(files, "2026-11") == ([], None, "data/events/2026-11.parquet")
21
+
22
+
23
+ def test_compact_rows_merges_aligns_and_dedupes():
24
+ old_monthly = [{"event_type": "session", "session_id": "a" * 16, "timestamp": "2026-09-02T00:00:00Z",
25
+ "distinct_queries": 1}]
26
+ shard1 = [{"event_type": "session", "session_id": "a" * 16, "timestamp": "2026-09-02T00:10:00Z",
27
+ "distinct_queries": 3},
28
+ {"event_type": "mlx_benchmark", "selected_model": "mlx-community/x-4bit", "timestamp": "2026-09-02T00:11:00Z"}]
29
+ shard2 = [{"event_type": "session", "session_id": "b" * 16, "timestamp": "2026-09-03T00:00:00Z"}]
30
+ rows = compact.compact_rows([old_monthly, shard1, shard2])
31
+ assert len(rows) == 3 and all(set(r) == set(COLUMNS) for r in rows)
32
+ assert [r["distinct_queries"] for r in rows if r["session_id"] == "a" * 16] == [3]
33
+ assert [r["timestamp"] for r in rows] == sorted(r["timestamp"] for r in rows)
34
+
35
+
36
+ def test_previous_month_and_refuses_current():
37
+ from datetime import datetime, timezone
38
+ assert compact.previous_month(datetime(2026, 1, 15, tzinfo=timezone.utc)) == "2025-12"
39
+ assert compact.previous_month(datetime(2026, 9, 14, tzinfo=timezone.utc)) == "2026-08"
40
+ import pytest
41
+ now = datetime.now(timezone.utc).strftime("%Y-%m")
42
+ with pytest.raises(SystemExit):
43
+ compact.main(["--month", now])
tests/test_events.py CHANGED
@@ -24,26 +24,65 @@ def ev(**kw):
24
  return ClientEvent(**{"event_type": "search", "session_id": SID, **kw})
25
 
26
 
27
- def test_valid_search_row_has_all_columns_and_no_pii():
28
- row = to_row(ev(model_family="Qwen", parameter_bucket="8-15B", quantization="4-bit",
29
  target_context=32768, priority="balanced", hardware_source="confirmed",
30
- hardware_memory_class=36), CAT)
 
 
31
  assert set(row) == set(COLUMNS)
 
 
32
  assert row["hardware_confirmed"] is True and row["suspicious_flags"] == []
33
  assert row["timestamp"].endswith("Z") and "." not in row["timestamp"]
34
- for forbidden in ("ip", "user_agent", "email", "name", "cookie"):
35
  assert forbidden not in row
36
 
37
 
38
- def test_server_fills_catalogue_facts():
39
- row = to_row(ev(event_type="model_select", selected_model="mlx-community/Qwen3-8B-4bit",
40
- selected_model_rank=0), CAT)
41
- assert row["hf_downloads_at_selection"] == 42 and row["quant_bits"] == 4
42
  assert row["model_family"] == "Qwen" and row["parameter_bucket"] == "8-15B"
43
 
44
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
45
  @pytest.mark.parametrize("bad", [
46
  {"event_type": "delete_everything"},
 
47
  {"event_type": "search", "session_id": "NOT-HEX"},
48
  {"event_type": "search", "target_context": 12345},
49
  {"event_type": "search", "hardware_memory_class": 17},
@@ -56,7 +95,6 @@ def test_server_fills_catalogue_facts():
56
  {"event_type": "search", "ip": "1.2.3.4"},
57
  {"event_type": "search", "model_family": "<img src=x>"},
58
  {"event_type": "search", "benchmark_version": "1.0; rm -rf"},
59
- {"event_type": "search", "compare_models": ["a/b"] * 6},
60
  ])
61
  def test_rejects_malformed(bad):
62
  with pytest.raises(ValidationError):
@@ -67,7 +105,7 @@ def test_batch_limits():
67
  with pytest.raises(ValidationError):
68
  EventBatch(events=[])
69
  with pytest.raises(ValidationError):
70
- EventBatch(events=[{"event_type": "search"}] * 51)
71
 
72
 
73
  def test_notes_are_scrubbed_and_capped():
@@ -80,20 +118,20 @@ def test_notes_are_scrubbed_and_capped():
80
 
81
 
82
  def test_flags_unknown_model_and_implausible_tps():
83
- row = to_row(ev(event_type="mlx_benchmark_submission", selected_model="mlx-community/Nope-1B",
84
  generation_tps=10, benchmark_version="1", benchmark_type="mlx_lm"), CAT)
85
  assert "unknown_model" in row["suspicious_flags"]
86
- row = to_row(ev(event_type="mlx_benchmark_submission", selected_model="mlx-community/Llama-3.3-70B-Instruct-4bit",
87
  generation_tps=1500, benchmark_version="1", benchmark_type="mlx_lm"), CAT)
88
  assert "implausible_tps" in row["suspicious_flags"]
89
- row = to_row(ev(event_type="mlx_benchmark_submission", selected_model="mlx-community/Qwen3-8B-4bit",
90
  generation_tps=45, benchmark_version="1", benchmark_type="mlx_lm", peak_memory_gb=5.1,
91
  reported_ram_gb=36), CAT)
92
  assert row["suspicious_flags"] == []
93
 
94
 
95
  def test_flags_other_inconsistencies():
96
- row = to_row(ev(event_type="mlx_benchmark_submission", peak_memory_gb=40, reported_ram_gb=16), CAT)
97
  assert {"peak_memory_exceeds_ram", "incomplete_benchmark"} <= set(row["suspicious_flags"])
98
  row = to_row(ev(event_type="feedback", tried="yes", quality_rating="good", failure_reason="too_slow"), CAT)
99
  assert "conflicting_feedback" in row["suspicious_flags"]
@@ -103,17 +141,17 @@ def test_flags_other_inconsistencies():
103
 
104
 
105
  def test_quality_submission_validation():
106
- ok = to_row(ev(event_type="mlx_benchmark_submission", selected_model="mlx-community/Qwen3-8B-4bit",
107
  benchmark_type="mlx_lm", benchmark_version="mlxbench-1", perplexity=9.8, perplexity_stderr=0.2,
108
  eval_dataset="wikitext2-test-v1", eval_tokens=16368), CAT)
109
  assert ok["suspicious_flags"] == [] and ok["perplexity"] == 9.8
110
- bad = to_row(ev(event_type="mlx_benchmark_submission", selected_model="mlx-community/Qwen3-8B-4bit",
111
  benchmark_version="mlxbench-1", perplexity=5000, eval_dataset="wikitext2-test-v1", eval_tokens=16368), CAT)
112
  assert "implausible_perplexity" in bad["suspicious_flags"]
113
- inc = to_row(ev(event_type="mlx_benchmark_submission", selected_model="mlx-community/Qwen3-8B-4bit",
114
  benchmark_version="mlxbench-1", perplexity=9.8), CAT)
115
  assert "incomplete_quality" in inc["suspicious_flags"]
116
  with pytest.raises(ValidationError):
117
- ev(event_type="mlx_benchmark_submission", perplexity=0.5)
118
  with pytest.raises(ValidationError):
119
- ev(event_type="mlx_benchmark_submission", perplexity=9.8, eval_dataset="my-own-text")
 
24
  return ClientEvent(**{"event_type": "search", "session_id": SID, **kw})
25
 
26
 
27
+ def test_session_row_has_all_columns_and_no_pii():
28
+ row = to_row(ev(event_type="session", model_family="Qwen", parameter_bucket="8-15B", quantization="4-bit",
29
  target_context=32768, priority="balanced", hardware_source="confirmed",
30
+ hardware_memory_class=36, families_searched=["Qwen", "Gemma", "Qwen"], quants_searched=["4-bit"],
31
+ distinct_queries=3, top_model="mlx-community/Qwen3-8B-4bit", top_model_score=88.2,
32
+ top_model_fit="Comfortable", models_viewed=["mlx-community/Qwen3-8B-4bit"]), CAT)
33
  assert set(row) == set(COLUMNS)
34
+ assert row["schema_version"] == 2 and row["event_type"] == "session"
35
+ assert row["families_searched"] == ["Qwen", "Gemma"] # de-duplicated
36
  assert row["hardware_confirmed"] is True and row["suspicious_flags"] == []
37
  assert row["timestamp"].endswith("Z") and "." not in row["timestamp"]
38
+ for forbidden in ("ip", "user_agent", "email", "cookie"):
39
  assert forbidden not in row
40
 
41
 
42
+ def test_feedback_gets_catalogue_facts():
43
+ row = to_row(ev(event_type="feedback", selected_model="mlx-community/Qwen3-8B-4bit", tried="yes",
44
+ quality_rating="good"), CAT)
45
+ assert row["hf_downloads_at_selection"] == 42 and row["quant_bits"] == 4 and row["outcome"] == "worked"
46
  assert row["model_family"] == "Qwen" and row["parameter_bucket"] == "8-15B"
47
 
48
 
49
+ def test_legacy_clickstream_events_are_dropped_and_old_bench_name_is_accepted():
50
+ assert to_row(ev(event_type="search", model_family="Qwen"), CAT) is None
51
+ assert to_row(ev(event_type="model_view", selected_model="mlx-community/Qwen3-8B-4bit"), CAT) is None
52
+ row = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Qwen3-8B-4bit",
53
+ generation_tps=40, benchmark_version="mlxbench-1"), CAT)
54
+ assert row["event_type"] == "mlx_benchmark" and row["suspicious_flags"] == []
55
+
56
+
57
+ def test_session_lists_are_capped_and_checked():
58
+ with pytest.raises(ValidationError):
59
+ ev(event_type="session", models_viewed=[f"mlx-community/m{i}" for i in range(11)])
60
+ with pytest.raises(ValidationError):
61
+ ev(event_type="session", models_compared=["../etc/passwd"])
62
+ with pytest.raises(ValidationError):
63
+ ev(event_type="session", families_searched=["a", "b", "c", "d", "e", "f"])
64
+ row = to_row(ev(event_type="session", models_clicked=["mlx-community/Nope-1B"]), CAT)
65
+ assert "unknown_model_in_session" in row["suspicious_flags"]
66
+
67
+
68
+ def test_latest_sessions_keeps_newest_row_per_visit():
69
+ from app.events import latest_sessions
70
+ rows = [
71
+ {"event_type": "session", "session_id": "a" * 16, "timestamp": "2026-09-14T01:00:00Z", "distinct_queries": 1},
72
+ {"event_type": "feedback", "session_id": "a" * 16, "timestamp": "2026-09-14T01:00:30Z"},
73
+ {"event_type": "session", "session_id": "a" * 16, "timestamp": "2026-09-14T01:05:00Z", "distinct_queries": 4},
74
+ {"event_type": "session", "session_id": "b" * 16, "timestamp": "2026-09-14T01:05:00Z", "distinct_queries": 1},
75
+ {"event_type": "session", "session_id": "b" * 16, "timestamp": "2026-09-14T01:05:00Z", "distinct_queries": 2},
76
+ ]
77
+ out = latest_sessions(rows)
78
+ sessions = {r["session_id"]: r["distinct_queries"] for r in out if r["event_type"] == "session"}
79
+ assert sessions == {"a" * 16: 4, "b" * 16: 2} # later arrival wins a same-second tie
80
+ assert sum(1 for r in out if r["event_type"] == "feedback") == 1 and len(out) == 3
81
+
82
+
83
  @pytest.mark.parametrize("bad", [
84
  {"event_type": "delete_everything"},
85
+ {"event_type": "session", "compare_models": ["a/b"]},
86
  {"event_type": "search", "session_id": "NOT-HEX"},
87
  {"event_type": "search", "target_context": 12345},
88
  {"event_type": "search", "hardware_memory_class": 17},
 
95
  {"event_type": "search", "ip": "1.2.3.4"},
96
  {"event_type": "search", "model_family": "<img src=x>"},
97
  {"event_type": "search", "benchmark_version": "1.0; rm -rf"},
 
98
  ])
99
  def test_rejects_malformed(bad):
100
  with pytest.raises(ValidationError):
 
105
  with pytest.raises(ValidationError):
106
  EventBatch(events=[])
107
  with pytest.raises(ValidationError):
108
+ EventBatch(events=[{"event_type": "session"}] * 21)
109
 
110
 
111
  def test_notes_are_scrubbed_and_capped():
 
118
 
119
 
120
  def test_flags_unknown_model_and_implausible_tps():
121
+ row = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Nope-1B",
122
  generation_tps=10, benchmark_version="1", benchmark_type="mlx_lm"), CAT)
123
  assert "unknown_model" in row["suspicious_flags"]
124
+ row = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Llama-3.3-70B-Instruct-4bit",
125
  generation_tps=1500, benchmark_version="1", benchmark_type="mlx_lm"), CAT)
126
  assert "implausible_tps" in row["suspicious_flags"]
127
+ row = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Qwen3-8B-4bit",
128
  generation_tps=45, benchmark_version="1", benchmark_type="mlx_lm", peak_memory_gb=5.1,
129
  reported_ram_gb=36), CAT)
130
  assert row["suspicious_flags"] == []
131
 
132
 
133
  def test_flags_other_inconsistencies():
134
+ row = to_row(ev(event_type="mlx_benchmark", peak_memory_gb=40, reported_ram_gb=16), CAT)
135
  assert {"peak_memory_exceeds_ram", "incomplete_benchmark"} <= set(row["suspicious_flags"])
136
  row = to_row(ev(event_type="feedback", tried="yes", quality_rating="good", failure_reason="too_slow"), CAT)
137
  assert "conflicting_feedback" in row["suspicious_flags"]
 
141
 
142
 
143
  def test_quality_submission_validation():
144
+ ok = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Qwen3-8B-4bit",
145
  benchmark_type="mlx_lm", benchmark_version="mlxbench-1", perplexity=9.8, perplexity_stderr=0.2,
146
  eval_dataset="wikitext2-test-v1", eval_tokens=16368), CAT)
147
  assert ok["suspicious_flags"] == [] and ok["perplexity"] == 9.8
148
+ bad = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Qwen3-8B-4bit",
149
  benchmark_version="mlxbench-1", perplexity=5000, eval_dataset="wikitext2-test-v1", eval_tokens=16368), CAT)
150
  assert "implausible_perplexity" in bad["suspicious_flags"]
151
+ inc = to_row(ev(event_type="mlx_benchmark", selected_model="mlx-community/Qwen3-8B-4bit",
152
  benchmark_version="mlxbench-1", perplexity=9.8), CAT)
153
  assert "incomplete_quality" in inc["suspicious_flags"]
154
  with pytest.raises(ValidationError):
155
+ ev(event_type="mlx_benchmark", perplexity=0.5)
156
  with pytest.raises(ValidationError):
157
+ ev(event_type="mlx_benchmark", perplexity=9.8, eval_dataset="my-own-text")
tests/test_sink.py CHANGED
@@ -6,7 +6,7 @@ from app.sink import HubParquetSink, LocalParquetSink, rows_to_parquet_bytes
6
 
7
 
8
  def rows(n):
9
- return [to_row(ClientEvent(event_type="search", model_family="Qwen", quantization="4-bit")) for _ in range(n)]
10
 
11
 
12
  def test_batch_becomes_one_shard(tmp_path):
@@ -14,14 +14,14 @@ def test_batch_becomes_one_shard(tmp_path):
14
  s.add(rows(7))
15
  s.add(rows(5))
16
  assert s.flush()
17
- shards = list((tmp_path / "data/events").rglob("*.parquet"))
18
  assert len(shards) == 1
19
  t = pq.read_table(shards[0])
20
  assert t.num_rows == 12 and t.column_names == list(COLUMNS)
21
- assert s.flush() and len(list((tmp_path / "data/events").rglob("*.parquet"))) == 1 # empty flush = no shard
22
  s.add(rows(1))
23
  s.flush()
24
- assert len(list((tmp_path / "data/events").rglob("*.parquet"))) == 2 # appends, never rewrites
25
  assert len(s.read_existing()) == 13
26
 
27
 
@@ -50,7 +50,7 @@ def test_hub_failure_keeps_spool_and_retries(tmp_path):
50
  assert not s2.flush()
51
  assert s2.flush() and s2.pending == 0 and s2.status()["healthy"]
52
  assert len(api.commits) == 1 and len(api.commits[0]) == 1
53
- assert api.commits[0][0].startswith("data/events/") and api.commits[0][0].endswith(".parquet")
54
  s3 = HubParquetSink("org/data", None, spool_dir=tmp_path, api=api, flush_seconds=999)
55
  assert s3.pending == 0
56
 
@@ -58,13 +58,38 @@ def test_hub_failure_keeps_spool_and_retries(tmp_path):
58
  def test_local_refuses_overwrite(tmp_path):
59
  s = LocalParquetSink(tmp_path, flush_seconds=999)
60
  data = rows_to_parquet_bytes(rows(1))
61
- s._write_shard("data/events/x.parquet", data)
62
  with pytest.raises(FileExistsError):
63
- s._write_shard("data/events/x.parquet", data)
64
 
65
 
66
  def test_dataset_schema_is_deliberate():
67
  """Changing columns changes the public dataset schema: bump this list on purpose, and
68
  migrate or document existing shards (see the dataset card's versioning section)."""
69
- assert len(COLUMNS) == 62
70
  assert {"perplexity", "perplexity_stderr", "eval_dataset", "eval_tokens"} <= set(COLUMNS)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
6
 
7
 
8
  def rows(n):
9
+ return [to_row(ClientEvent(event_type="session", model_family="Qwen", quantization="4-bit")) for _ in range(n)]
10
 
11
 
12
  def test_batch_becomes_one_shard(tmp_path):
 
14
  s.add(rows(7))
15
  s.add(rows(5))
16
  assert s.flush()
17
+ shards = list((tmp_path / "data/incoming").rglob("*.parquet"))
18
  assert len(shards) == 1
19
  t = pq.read_table(shards[0])
20
  assert t.num_rows == 12 and t.column_names == list(COLUMNS)
21
+ assert s.flush() and len(list((tmp_path / "data/incoming").rglob("*.parquet"))) == 1 # empty flush = no shard
22
  s.add(rows(1))
23
  s.flush()
24
+ assert len(list((tmp_path / "data/incoming").rglob("*.parquet"))) == 2 # appends, never rewrites
25
  assert len(s.read_existing()) == 13
26
 
27
 
 
50
  assert not s2.flush()
51
  assert s2.flush() and s2.pending == 0 and s2.status()["healthy"]
52
  assert len(api.commits) == 1 and len(api.commits[0]) == 1
53
+ assert api.commits[0][0].startswith("data/incoming/") and api.commits[0][0].endswith(".parquet")
54
  s3 = HubParquetSink("org/data", None, spool_dir=tmp_path, api=api, flush_seconds=999)
55
  assert s3.pending == 0
56
 
 
58
  def test_local_refuses_overwrite(tmp_path):
59
  s = LocalParquetSink(tmp_path, flush_seconds=999)
60
  data = rows_to_parquet_bytes(rows(1))
61
+ s._write_shard("data/incoming/x.parquet", data)
62
  with pytest.raises(FileExistsError):
63
+ s._write_shard("data/incoming/x.parquet", data)
64
 
65
 
66
  def test_dataset_schema_is_deliberate():
67
  """Changing columns changes the public dataset schema: bump this list on purpose, and
68
  migrate or document existing shards (see the dataset card's versioning section)."""
69
+ assert len(COLUMNS) == 67
70
  assert {"perplexity", "perplexity_stderr", "eval_dataset", "eval_tokens"} <= set(COLUMNS)
71
+
72
+
73
+ def test_size_triggered_flush(tmp_path, monkeypatch):
74
+ import time
75
+ import app.sink as sink_mod
76
+ monkeypatch.setattr(sink_mod, "FLUSH_ROWS", 5)
77
+ s = LocalParquetSink(tmp_path, flush_seconds=999)
78
+ s.add(rows(3))
79
+ assert s.pending == 3
80
+ s.add(rows(3)) # crosses FLUSH_ROWS -> flushes in the background
81
+ for _ in range(50):
82
+ if s.pending == 0:
83
+ break
84
+ time.sleep(0.05)
85
+ assert s.pending == 0 and len(list((tmp_path / "data/incoming").rglob("*.parquet"))) == 1
86
+
87
+
88
+ def test_reads_monthly_and_incoming(tmp_path):
89
+ from app.sink import rows_to_parquet_bytes as rb
90
+ (tmp_path / "data/events").mkdir(parents=True)
91
+ (tmp_path / "data/events/2026-08.parquet").write_bytes(rb(rows(4)))
92
+ s = LocalParquetSink(tmp_path, flush_seconds=999)
93
+ s.add(rows(2))
94
+ s.flush()
95
+ assert len(s.read_existing()) == 6