Xiangyi Li commited on
Commit
bab7845
·
1 Parent(s): 9bcfb03

Run records: newline-only JSONL, catalogs, linked records.

Browse files

JSONL splits on newlines only (a JSON string may hold U+2028, U+2029 or U+0085) and a line that does not parse is
counted, logged and reported as a finding instead of dropped silently; the search index is read the same way. The
importers' catalogs at the public root add each run's 'sampled' (which attempts were kept). Public records of the
same source run (same url) link to each other. A sampling-only run counts the tasks of its eval step. A dataset
value such as 'none (training rollouts only)' names no dataset, and an eval set a run names without logging a score
still has its dataset page.

Files changed (4) hide show
  1. app_api.py +1 -1
  2. fwruns.py +35 -12
  3. results_api.py +20 -4
  4. test_public_runs.py +26 -0
app_api.py CHANGED
@@ -275,7 +275,7 @@ def fw_run(name: str):
275
  pool, held = results_api.datasets_of(r)
276
  steps = len({row.get('step') for row in m if row.get('step') is not None}) or len(m) # distinct steps in the whole log
277
  return {**r, 'metrics': fwruns.thin(m, METRIC_ROWS, lambda row: any(k in row for k in evals)), 'metrics_total': len(m), 'metrics_steps': steps, 'datasets': {'training': pool, 'heldout': held},
278
- 'token_ids': fwruns.token_ids(name) if r.get('steps') else []}
279
 
280
 
281
  @router.get('/fwruns/{name}/metrics', response_class=PlainTextResponse)
 
275
  pool, held = results_api.datasets_of(r)
276
  steps = len({row.get('step') for row in m if row.get('step') is not None}) or len(m) # distinct steps in the whole log
277
  return {**r, 'metrics': fwruns.thin(m, METRIC_ROWS, lambda row: any(k in row for k in evals)), 'metrics_total': len(m), 'metrics_steps': steps, 'datasets': {'training': pool, 'heldout': held},
278
+ 'token_ids': fwruns.token_ids(name) if r.get('steps') else [], 'linked': results_api.linked_of(r)}
279
 
280
 
281
  @router.get('/fwruns/{name}/metrics', response_class=PlainTextResponse)
fwruns.py CHANGED
@@ -22,7 +22,7 @@ A run's steps follow AC2's vocabulary: eval@k is the held-out suite sampled afte
22
  groups optimizer step k+1 trained on (fw_sync rebuilds the batches from completion order; a last train@k may be a batch
23
  the run sampled but never trained on). Each step holds its tasks and each task its attempts (traces).
24
  """
25
- import gzip, json, math, os, re, statistics, threading, time
26
  from concurrent.futures import ThreadPoolExecutor
27
  from pathlib import Path
28
 
@@ -32,6 +32,7 @@ DEFAULTS = {'fireworks': {'source': 'fireworks', 'org': 'BenchFlow', 'project':
32
  PREFIX = ROOTS['fireworks']['prefix'] # the account snapshot and project reports live under the Fireworks root
33
  TTL = 60
34
  _cache, _lock = {}, threading.Lock()
 
35
 
36
 
37
  def _repo():
@@ -137,13 +138,32 @@ def finite(v):
137
  return v
138
 
139
 
140
- def _jsonl(raw):
141
- out = []
142
- for line in (raw or '').splitlines():
 
 
 
143
  if not line.strip(): continue
144
  try: out.append(finite(json.loads(line)))
145
- except ValueError: continue # a torn last line while a sync writes the file
146
- return out
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
147
 
148
 
149
  def _attempt(a, seen):
@@ -171,7 +191,7 @@ def runs():
171
  summary = {k: v for k, v in r.items() if k not in ('steps', 'metrics', 'totals')}
172
  training = r.get('kind', 'training') == 'training'
173
  # A task gives GRPO a signal only if its scored attempts differ; a screen exists to count those tasks.
174
- groups = [g for s in trains for g in s['groups']]
175
  scored = [{a['reward'] for a in g['attempts'] if a.get('reward') is not None} for g in groups]
176
  summary.update(steps=len(trains) if training else None, train_reward=trains[-1]['mean'] if trains and training else None, train_key=trains[-1]['key'] if trains and training else None,
177
  heldout=(f"{evals[0]['mean']:.3f}" + (f" → {evals[-1]['mean']:.3f}" if len(evals) > 1 else '')) if evals else None,
@@ -208,8 +228,10 @@ def run(name):
208
  except ValueError: return None
209
  if not isinstance(meta, dict): return None
210
  seen = {}
211
- attempts = [_attempt(a, seen) for a in _jsonl(_read(f'{name}/attempts.jsonl', root)) if isinstance(a, dict)]
212
- metrics = [m for m in _jsonl(_read(f'{name}/metrics.jsonl', root)) if isinstance(m, dict)]
 
 
213
  steps = {}
214
  for a in attempts:
215
  steps.setdefault((a['kind'], a['version']), []).append(a)
@@ -225,9 +247,10 @@ def run(name):
225
  ref = run(meta['baseline_run']) or {}
226
  ev = [st for st in ref.get('steps') or [] if st['kind'] == 'eval']
227
  if ev: baseline = {'run': meta['baseline_run'], **{k: v for k, v in ev[0].items() if k != 'groups'}}
228
- base = {**DEFAULTS[root], **{k: v for k, v in meta.items() if v not in (None, '')}}
229
  if not base.get('project'): base['project'] = slug(base.get('org')) if base.get('org') else 'public-runs'
230
- return {**meta, **base, 'id': name, 'root': root, 'project': slug(base['project']), 'steps': out, 'metrics': metrics, 'totals': _stats(attempts), 'baseline': baseline}
 
231
  return _cached(('run', name), load, ROOTS[root]['ttl'])
232
 
233
 
@@ -326,7 +349,7 @@ def _index(name):
326
  raw = _read_bytes(f'{name}/search.jsonl.gz', root)
327
  if raw:
328
  out = {}
329
- for line in gzip.decompress(raw).decode('utf-8', 'replace').splitlines():
330
  if line.strip(): x = json.loads(line); out[x['id']] = x['parts']
331
  return out
332
  local = _local(root)
 
22
  groups optimizer step k+1 trained on (fw_sync rebuilds the batches from completion order; a last train@k may be a batch
23
  the run sampled but never trained on). Each step holds its tasks and each task its attempts (traces).
24
  """
25
+ import gzip, json, logging, math, os, re, statistics, threading, time
26
  from concurrent.futures import ThreadPoolExecutor
27
  from pathlib import Path
28
 
 
32
  PREFIX = ROOTS['fireworks']['prefix'] # the account snapshot and project reports live under the Fireworks root
33
  TTL = 60
34
  _cache, _lock = {}, threading.Lock()
35
+ log = logging.getLogger('fwruns')
36
 
37
 
38
  def _repo():
 
138
  return v
139
 
140
 
141
+ def _jsonl(raw, what=''):
142
+ """(rows, unreadable): one JSON object per line. Lines split on newlines only, since a JSON string may hold U+2028,
143
+ U+2029 or U+0085, which str.splitlines also splits on. A line that does not parse (a torn last line while a sync
144
+ writes the file, say) is counted and logged, not dropped silently."""
145
+ out, bad = [], 0
146
+ for line in (raw or '').split('\n'):
147
  if not line.strip(): continue
148
  try: out.append(finite(json.loads(line)))
149
+ except ValueError: bad += 1
150
+ if bad: log.warning('%s: %d line(s) could not be read', what, bad)
151
+ return out, bad
152
+
153
+
154
+ def catalog():
155
+ """What the importers' catalogs (_catalog*.json at the public root) say about each run, by run id: which attempts
156
+ were kept ('sampled') and which texts were cut. {} when there is no catalog."""
157
+ def load():
158
+ local, out = _local('public'), {}
159
+ names = sorted(p.name for p in local.glob('_catalog*.json')) if local and local.is_dir() else [] if local else sorted(x.path.split('/')[-1] for x in _tree('', 'public') if hasattr(x, 'size') and x.path.split('/')[-1].startswith('_catalog'))
160
+ for f in names:
161
+ try: runs_ = json.loads(_read(f, 'public') or '{}').get('runs') or []
162
+ except ValueError: continue
163
+ for r in runs_:
164
+ if isinstance(r, dict) and r.get('id'): out.setdefault(r['id'], {}).update({k: r[k] for k in ('sampled', 'texts_cut') if r.get(k)})
165
+ return out
166
+ return _cached('catalog', load, ROOTS['public']['ttl'])
167
 
168
 
169
  def _attempt(a, seen):
 
191
  summary = {k: v for k, v in r.items() if k not in ('steps', 'metrics', 'totals')}
192
  training = r.get('kind', 'training') == 'training'
193
  # A task gives GRPO a signal only if its scored attempts differ; a screen exists to count those tasks.
194
+ groups = [g for s in (trains if training else [s for s in steps if s['kind'] == 'eval'] or trains) for g in s['groups']] # a sampling-only run's tasks are those of its eval step
195
  scored = [{a['reward'] for a in g['attempts'] if a.get('reward') is not None} for g in groups]
196
  summary.update(steps=len(trains) if training else None, train_reward=trains[-1]['mean'] if trains and training else None, train_key=trains[-1]['key'] if trains and training else None,
197
  heldout=(f"{evals[0]['mean']:.3f}" + (f" → {evals[-1]['mean']:.3f}" if len(evals) > 1 else '')) if evals else None,
 
228
  except ValueError: return None
229
  if not isinstance(meta, dict): return None
230
  seen = {}
231
+ rows, bad_a = _jsonl(_read(f'{name}/attempts.jsonl', root), f'{name}/attempts.jsonl')
232
+ mrows, bad_m = _jsonl(_read(f'{name}/metrics.jsonl', root), f'{name}/metrics.jsonl')
233
+ attempts = [_attempt(a, seen) for a in rows if isinstance(a, dict)]
234
+ metrics = [m for m in mrows if isinstance(m, dict)]
235
  steps = {}
236
  for a in attempts:
237
  steps.setdefault((a['kind'], a['version']), []).append(a)
 
247
  ref = run(meta['baseline_run']) or {}
248
  ev = [st for st in ref.get('steps') or [] if st['kind'] == 'eval']
249
  if ev: baseline = {'run': meta['baseline_run'], **{k: v for k, v in ev[0].items() if k != 'groups'}}
250
+ base = {**DEFAULTS[root], **({k: v for k, v in catalog().get(name, {}).items()} if root == 'public' else {}), **{k: v for k, v in meta.items() if v not in (None, '')}}
251
  if not base.get('project'): base['project'] = slug(base.get('org')) if base.get('org') else 'public-runs'
252
+ unreadable = {f: n for f, n in (('attempts.jsonl', bad_a), ('metrics.jsonl', bad_m)) if n}
253
+ return {**meta, **base, 'id': name, 'root': root, 'project': slug(base['project']), 'steps': out, 'metrics': metrics, 'totals': _stats(attempts), 'baseline': baseline, 'unreadable': unreadable}
254
  return _cached(('run', name), load, ROOTS[root]['ttl'])
255
 
256
 
 
349
  raw = _read_bytes(f'{name}/search.jsonl.gz', root)
350
  if raw:
351
  out = {}
352
+ for line in gzip.decompress(raw).decode('utf-8', 'replace').split('\n'): # newlines only: a text may hold U+2028
353
  if line.strip(): x = json.loads(line); out[x['id']] = x['parts']
354
  return out
355
  local = _local(root)
results_api.py CHANGED
@@ -284,11 +284,14 @@ def as_list(v):
284
  return [str(x) for x in v if x] if isinstance(v, list) else []
285
 
286
 
 
 
 
287
  def datasets_of(r):
288
- """(training data, held-out data) of a run as dataset refs {id, title}. fw_sync calls the training pool the run's
289
- collection; a public record calls it dataset."""
290
- pool = (as_list(r.get('collection')) or as_list(r.get('dataset')) or [None])[0]
291
- held = suite_for(r.get('eval_tasks'), (as_list(r.get('eval_suite')) or [None])[0])
292
  return ({'id': slug(pool), 'title': pool} if pool else None), held
293
 
294
 
@@ -352,6 +355,15 @@ def run_item(r, snap):
352
  'href': run_href(r)}
353
 
354
 
 
 
 
 
 
 
 
 
 
355
  def fw_project_of_model(base_model):
356
  """The project of the Fireworks runs that post-train this Fireworks model resource, if any."""
357
  return next((r['project'] for r in all_runs() if not public(r) and r.get('base_model') == base_model), None)
@@ -535,6 +547,8 @@ def dataset_index(project='all', org='all'):
535
  elif held and names: # a run that logged eval scores on its eval_suite
536
  s = series(r, names[0][1]); label = metric_label(names[0][0])
537
  add_run(ds(held, 'evaluation'), x, f"{plural(len(names), 'eval score')} logged", {'text': span(s), 'label': f'{label}, first and last logged value'})
 
 
538
  return D
539
 
540
 
@@ -1023,6 +1037,8 @@ def metric_findings(r):
1023
 
1024
  def findings_of(r, snap):
1025
  F = [finding('note', 'info', 'Notes', r['note'])] if r.get('note') else []
 
 
1026
  if traced(r): F += attempt_findings(r, snap)
1027
  if public(r) or r.get('metrics_map'): F += metric_findings(r)
1028
  return sorted(F, key=lambda f: ORDER[f['severity']])
 
284
  return [str(x) for x in v if x] if isinstance(v, list) else []
285
 
286
 
287
+ NONE = re.compile(r'^\s*(none|n/?a|null|no eval)\b', re.I) # a record's way of saying it has no such dataset
288
+
289
+
290
  def datasets_of(r):
291
+ """(training data, eval set) of a run as dataset refs {id, title}. fw_sync calls the training pool the run's
292
+ collection; a public record calls it dataset. A value such as 'none (training rollouts only)' means there is none."""
293
+ pool = next((x for x in as_list(r.get('collection')) or as_list(r.get('dataset')) if not NONE.match(x)), None)
294
+ held = suite_for(r.get('eval_tasks'), next((x for x in as_list(r.get('eval_suite')) if not NONE.match(x)), None))
295
  return ({'id': slug(pool), 'title': pool} if pool else None), held
296
 
297
 
 
355
  'href': run_href(r)}
356
 
357
 
358
+ def linked_of(r):
359
+ """Other records of the same source run: public records with the same source URL (a run's metrics and a sample of
360
+ its rollouts, imported apart)."""
361
+ url = (r.get('url') or '').strip()
362
+ if not url: return []
363
+ return [{'id': o['id'], 'title': o.get('title') or o['id'], 'href': run_href(o), 'traced': traced(o), 'metrics_only': not traced(o) and bool(o.get('metrics'))}
364
+ for o in all_runs() if o['id'] != r['id'] and (o.get('url') or '').strip() == url]
365
+
366
+
367
  def fw_project_of_model(base_model):
368
  """The project of the Fireworks runs that post-train this Fireworks model resource, if any."""
369
  return next((r['project'] for r in all_runs() if not public(r) and r.get('base_model') == base_model), None)
 
547
  elif held and names: # a run that logged eval scores on its eval_suite
548
  s = series(r, names[0][1]); label = metric_label(names[0][0])
549
  add_run(ds(held, 'evaluation'), x, f"{plural(len(names), 'eval score')} logged", {'text': span(s), 'label': f'{label}, first and last logged value'})
550
+ elif held: # named as the run's eval set, with no score in the records
551
+ add_run(ds(held, 'evaluation'), x, 'names it as its eval set', None)
552
  return D
553
 
554
 
 
1037
 
1038
  def findings_of(r, snap):
1039
  F = [finding('note', 'info', 'Notes', r['note'])] if r.get('note') else []
1040
+ for f, n in (r.get('unreadable') or {}).items():
1041
+ F.append(finding(f'unreadable_{f}', 'medium', f"{plural(n, 'line')} of {f} could not be read", f"They are not valid JSON (a sync may have been writing the file), so they are left out of every count and chart."))
1042
  if traced(r): F += attempt_findings(r, snap)
1043
  if public(r) or r.get('metrics_map'): F += metric_findings(r)
1044
  return sorted(F, key=lambda f: ORDER[f['severity']])
test_public_runs.py CHANGED
@@ -279,3 +279,29 @@ def test_an_evaluations_endings_agree_with_its_pass_rate(client, roots):
279
  e = get(client, 'evals/fixture-agent-rollouts:eval@2').json()
280
  assert e['endings'] == {'passed': 2, 'failed': 2} and e['timed_out'] == {'count': 2, 'scored': 2, 'passed': 2}
281
  assert round(e['pass_rate'] * 4) == e['endings']['passed']
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
279
  e = get(client, 'evals/fixture-agent-rollouts:eval@2').json()
280
  assert e['endings'] == {'passed': 2, 'failed': 2} and e['timed_out'] == {'count': 2, 'scored': 2, 'passed': 2}
281
  assert round(e['pass_rate'] * 4) == e['endings']['passed']
282
+
283
+
284
+ def test_jsonl_splits_on_newlines_only_and_counts_what_it_cannot_read(roots):
285
+ """A JSON string may hold U+2028 or U+0085 (str.splitlines splits on them); a line that does not parse is counted."""
286
+ fw, pub = roots
287
+ d = pub / 'fixture-agent-rollouts'
288
+ rows = [json.loads(l) for l in (d / 'attempts.jsonl').read_text().splitlines()]
289
+ rows[0]['exception_message'] = 'line one
line two\u0085three'
290
+ (d / 'attempts.jsonl').write_text(''.join(json.dumps(a, ensure_ascii=False) + '\n' for a in rows) + '{"task": "torn')
291
+ fwruns._cache.clear(); R._cache.clear()
292
+ r = fwruns.run('fixture-agent-rollouts')
293
+ assert r['totals']['traces'] == len(rows) and r['unreadable'] == {'attempts.jsonl': 1}
294
+ assert any(a.get('exception_message') == 'line one
line two\u0085three' for s in r['steps'] for g in s['groups'] for a in g['attempts'])
295
+ assert {f['key'] for f in R.findings_of(r, None)} >= {'unreadable_attempts.jsonl'}
296
+
297
+
298
+ def test_a_catalog_says_which_attempts_were_kept_and_same_source_records_link(roots, client):
299
+ fw, pub = roots
300
+ (pub / '_catalog_rollouts.json').write_text(json.dumps({'runs': [{'id': 'fixture-agent-rollouts', 'sampled': 'every rollout of steps 0 and 1'}]}))
301
+ write(pub, 'fixture-agent-metrics', run_json('fixture-agent-metrics', url='https://example.org/runs/fixture-agent-rollouts', eval_suite='none (training rollouts only)'), metrics_rows(4))
302
+ fwruns._cache.clear(); R._cache.clear()
303
+ assert fwruns.run('fixture-agent-rollouts')['sampled'] == 'every rollout of steps 0 and 1'
304
+ r = get(client, 'fwruns/fixture-agent-rollouts').json()
305
+ assert [x['id'] for x in r['linked']] == ['fixture-agent-metrics'] and r['linked'][0]['metrics_only']
306
+ m = get(client, 'fwruns/fixture-agent-metrics').json()
307
+ assert [x['id'] for x in m['linked']] == ['fixture-agent-rollouts'] and m['datasets']['heldout'] is None # 'none (...)' names no eval set