From 75c2a8b8a15b57fc52e5e7edd9d07f5655ee1c5e Mon Sep 17 00:00:00 2001 From: firestar5683 <168790843+firestar5683@users.noreply.github.com> Date: Sun, 20 Sep 2026 12:07:43 -0700 Subject: [PATCH] Drain verified prepared PCM tails after terminal replay model --- roadscore/prototype/mac_showcase.py | 11 ++- roadscore/prototype/prepared_core.py | 53 +++++++++++++- roadscore/prototype/test_prepared_core.py | 86 ++++++++++++++++++++++- 3 files changed, 145 insertions(+), 5 deletions(-) diff --git a/roadscore/prototype/mac_showcase.py b/roadscore/prototype/mac_showcase.py index e08973c191..668e88f021 100644 --- a/roadscore/prototype/mac_showcase.py +++ b/roadscore/prototype/mac_showcase.py @@ -188,10 +188,11 @@ def audio_worker(a): from core import Conductor from demo_engagement import DemoEngagement from operator_output import PresentationDelay - from prepared_core import load_archive, initial_frame, PreparedPresentation + from prepared_core import load_archive, initial_frame, PreparedPresentation, RecordedTail from prepared_clock import PreparedClock from stream_clock_bridge import StreamClockBridge audio, rate, meta = load_archive(a.score_archive, a.route) + tail = RecordedTail(meta['archived_tail']) processor = PreparedPresentation(a.score_archive, rate) playback_clock = PreparedClock(rate) sync = ReplayFollower(a.route,meta['first_model_ns'],duration=len(audio)/rate) if a.follow_playhead else None @@ -282,6 +283,7 @@ def audio_worker(a): (a.out/'prepared_ready').write_text('ready') while not done: sm.update(50);now=time.monotonic();controls.poll() + tail.observe(int(sm.logMonoTime['modelV2']),now,position) for name in received: if sm.updated[name]:received[name]=now latest=max(sm.logMonoTime.values()) @@ -339,7 +341,10 @@ def audio_worker(a): write_json(a.out/'status.json',snapshot) trace.write(json.dumps({**snapshot,'audio_s':shared['audio_s']})+'\n');last_status=now if started and now-started>=a.duration:break - if started and now-received['modelV2']>2 and (position or 0)/rate < len(audio)/rate-2:raise RuntimeError('Replay model stream stopped before prepared audio ended') + if (started and now-received['modelV2']>2 + and ((position or 0)/rate < len(audio)/rate-2 or tail.terminal_wall is not None) + and not tail.allows_stale(int(sm.logMonoTime['modelV2']),now)): + raise RuntimeError('Replay model stream stopped before prepared audio ended') finally: try:write_json(a.out/'audio_drained.json',dict(wall=time.monotonic(),drained=not errors)) finally: @@ -351,7 +356,7 @@ def audio_worker(a): finally: trace.close() if sync_trace is not None:sync_trace.close() - write_json(a.out/'prepared_summary.json',dict(generation_invoked=False,source=str(a.score_archive),first_source_frame=first_frame,last_source_frame=position,sample_rate=rate,portaudio_flags=flags,max_clock_error_seconds=max_drift,prepared_clock=playback_clock.snapshot(),stream_clock_bridge=bridge.snapshot() if bridge else None,callback_errors=errors,clock_errors=clock_errors,muted=a.muted,session_id=session,manual_scope='isolated replay display and presentation only')) + write_json(a.out/'prepared_summary.json',dict(generation_invoked=False,source=str(a.score_archive),first_source_frame=first_frame,last_source_frame=position,sample_rate=rate,portaudio_flags=flags,max_clock_error_seconds=max_drift,prepared_clock=playback_clock.snapshot(),stream_clock_bridge=bridge.snapshot() if bridge else None,callback_errors=errors,clock_errors=clock_errors,archived_tail=meta['archived_tail'],muted=a.muted,session_id=session,manual_scope='isolated replay display and presentation only')) if errors:raise RuntimeError(errors[0]) if abs(playback_clock.snapshot()['current_post_error_seconds'])>.05:raise RuntimeError('Prepared audio clock drift remained above 50 ms after recovery') diff --git a/roadscore/prototype/prepared_core.py b/roadscore/prototype/prepared_core.py index 741b0f4f49..2548cae381 100644 --- a/roadscore/prototype/prepared_core.py +++ b/roadscore/prototype/prepared_core.py @@ -16,6 +16,53 @@ from rhythm_timeline import RhythmTimeline from signal_shaker import SignalShaker, ShakerGrid +def archive_tail_proof(bridge, launch, origin, duration, offset): + """Recognize a bounded recorded tail after the bridge's final modelV2. + + replay_bridge records last_t only from modelV2, relative to origin_ns. The + original replay supervisor must also have confirmed the actual route EOF. + Missing or invalid optional evidence grants no exemption from stream loss. + """ + if (not isinstance(bridge, dict) or launch.get('end_reason') != 'native final segment exhausted' + or bridge.get('route') != launch.get('route') or 'failure' not in bridge or bridge['failure'] is not None + or type(bridge.get('origin_ns')) is not int or bridge['origin_ns'] != origin.get('first_model_ns') + or type(bridge.get('messages')) is not int or bridge['messages'] <= 0 + or bridge.get('clock') != 'zero at first model received; source logMonoTime remains unchanged'): + return None + last = bridge.get('last_t') + if any(type(value) not in (int,float) or not math.isfinite(value) for value in (last,duration,offset)) or last < 0: + return None + tail = duration + offset - last + if not 0 < tail <= 10: + return None + return dict(terminal_model_ns=bridge['origin_ns']+round(last*1e9),terminal_route_t=last, + tail_seconds=tail,source='bridge.json modelV2 endpoint at confirmed replay EOF') + + +class RecordedTail: + def __init__(self, proof): + self.proof = proof + self.terminal_wall = self.progress_wall = self.position = None + + def matches(self, model_ns): + # At most 1 ns accounts for floating-point reconstruction of last_t. + return bool(self.proof and type(model_ns) is int and abs(model_ns-self.proof['terminal_model_ns']) <= 1) + + def observe(self, model_ns, wall, position): + if not self.proof:return + if self.matches(model_ns) and self.terminal_wall is None:self.terminal_wall = wall + if type(position) is int and (self.position is None or position > self.position): + self.position, self.progress_wall = position, wall + + def allows_stale(self, model_ns, wall): + # Duplicate terminal messages never refresh this deadline. The extra second + # covers the existing 750 ms source-clock guard and callback quantization. + # PCM must continue advancing, and its original file length is unchanged. + return bool(self.matches(model_ns) and self.terminal_wall is not None and self.progress_wall is not None + and 0 <= wall-self.terminal_wall <= self.proof['tail_seconds']+1. + and 0 <= wall-self.progress_wall <= 1.) + + def load_archive(path, route): path = Path(path) launch = json.loads((path/'launch.json').read_text()) @@ -37,7 +84,11 @@ def load_archive(path, route): offset = first['callback_wall'] + first['dac_delay'] - origin['host_received_wall'] if abs(offset) > 2: raise ValueError('Unexpected original audio clock offset') - return audio, rate, {**launch, **origin, 'audio_offset': offset, 'duration': len(audio)/rate} + try:bridge = json.loads((path/'bridge.json').read_text()) + except (OSError,ValueError):bridge = None + duration = len(audio)/rate + tail = archive_tail_proof(bridge,launch,origin,duration,offset) + return audio, rate, {**launch, **origin, 'audio_offset': offset, 'duration': duration, 'archived_tail': tail} def initial_frame(meta, mono_ns, dac_since_receipt, rate): diff --git a/roadscore/prototype/test_prepared_core.py b/roadscore/prototype/test_prepared_core.py index 8f6edf05eb..e5d50a9faa 100644 --- a/roadscore/prototype/test_prepared_core.py +++ b/roadscore/prototype/test_prepared_core.py @@ -1,10 +1,12 @@ +import ast import json from pathlib import Path import tempfile +from types import SimpleNamespace import unittest import numpy as np import soundfile as sf -from prepared_core import load_archive, initial_frame, PreparedPresentation +from prepared_core import load_archive, initial_frame, PreparedPresentation, archive_tail_proof, RecordedTail from replay_ui_controls import isolated_replay @@ -23,6 +25,17 @@ class PreparedTests(unittest.TestCase): np.testing.assert_array_equal(pcm,self.pcm) self.assertEqual(initial_frame(meta,1000000000,.15,rate),0) self.assertEqual(initial_frame(meta,1500000000,.15,rate),24000) + self.assertIsNone(meta['archived_tail']) + def test_optional_terminal_proof_keeps_original_pcm(self): + bridge=dict(route='fixture',origin_ns=1000000000,last_t=.7,failure=None,messages=100, + clock='zero at first model received; source logMonoTime remains unchanged') + (self.path/'bridge.json').write_text(json.dumps(bridge)) + pcm,rate,meta=load_archive(self.path,'fixture') + np.testing.assert_array_equal(pcm,self.pcm) + self.assertEqual(meta['archived_tail']['terminal_model_ns'],1700000000) + self.assertAlmostEqual(meta['archived_tail']['tail_seconds'],.45) + (self.path/'bridge.json').write_text('incomplete JSON') + self.assertIsNone(load_archive(self.path,'fixture')[2]['archived_tail']) def test_wrong_route_and_final_pcm_rejected(self): with self.assertRaises(ValueError):load_archive(self.path,'other') (self.path/'launch.json').write_text(json.dumps(dict(route='fixture',render_mode='current',end_reason='native final segment exhausted'))) @@ -89,4 +102,75 @@ class PreparedTests(unittest.TestCase): self.assertEqual(result.shape,core.shape) self.assertTrue(np.isfinite(result).all()) + +class ArchivedTailTests(unittest.TestCase): + def setUp(self): + self.origin={'first_model_ns':36517892685238} + self.launch={'route':'fixture','end_reason':'native final segment exhausted'} + self.bridge=dict(route='fixture',origin_ns=self.origin['first_model_ns'],last_t=300.000663318, + failure=None,messages=84032,clock='zero at first model received; source logMonoTime remains unchanged') + def proof(self,bridge=None,launch=None,duration=306.2): + return archive_tail_proof(self.bridge if bridge is None else bridge,self.launch if launch is None else launch, + self.origin,duration,.173852003) + def test_measured_terminal_model_and_recorded_tail(self): + proof=self.proof() + self.assertEqual(proof['terminal_model_ns'],36817893348556) + self.assertAlmostEqual(proof['tail_seconds'],6.373188685) + for delta in (-1,0,1): + tail=RecordedTail(proof);tail.observe(proof['terminal_model_ns']+delta,300.,100) + tail.observe(proof['terminal_model_ns']+delta,303.,200) + self.assertTrue(tail.allows_stale(proof['terminal_model_ns']+delta,303.)) + for delta in (-2,2,-50_000_000): + self.assertFalse(tail.allows_stale(proof['terminal_model_ns']+delta,303.)) + for age in (-1,float('nan'),proof['tail_seconds']+1.001): + self.assertFalse(tail.allows_stale(proof['terminal_model_ns'],300.+age)) + self.assertFalse(RecordedTail(None).allows_stale(proof['terminal_model_ns'],303.)) + def test_duplicate_endpoint_cannot_extend_deadline_and_pcm_must_progress(self): + proof=self.proof();mono=proof['terminal_model_ns'];tail=RecordedTail(proof) + tail.observe(mono,300.,100) + tail.observe(mono,301.,200) + self.assertEqual(tail.terminal_wall,300.) + self.assertTrue(tail.allows_stale(mono,301.)) + tail.observe(mono,302.1,200) + self.assertFalse(tail.allows_stale(mono,302.1)) + # A bounded backwards DAC correction cannot masquerade as forward progress. + tail.observe(mono,302.2,150) + self.assertFalse(tail.allows_stale(mono,302.2)) + tail.observe(mono,302.3,201) + self.assertTrue(tail.allows_stale(mono,302.3)) + tail.observe(mono,308.,300) + self.assertFalse(tail.allows_stale(mono,308.)) + def test_incomplete_foreign_failed_or_excessive_tail_proof_is_rejected(self): + for patch in ({'route':'other'},{'origin_ns':100},{'failure':'lost pipe'}, + {'last_t':float('nan')},{'last_t':True},{'last_t':-1},{'messages':0},{'clock':'unknown'}): + with self.subTest(patch=patch):self.assertIsNone(self.proof({**self.bridge,**patch})) + self.assertIsNone(self.proof({key:value for key,value in self.bridge.items() if key!='failure'})) + self.assertIsNone(self.proof(launch={**self.launch,'end_reason':'requested duration'})) + self.assertIsNone(self.proof(duration=311.)) + self.assertIsNone(self.proof(duration=299.)) + def test_real_worker_stream_loss_gate_accepts_only_bounded_exact_endpoint(self): + path=Path(__file__).with_name('mac_showcase.py') + worker=next(node for node in ast.parse(path.read_text()).body if isinstance(node,ast.FunctionDef) and node.name=='audio_worker') + guard=next(node for node in ast.walk(worker) if isinstance(node,ast.If) and len(node.body)==1 + and isinstance(node.body[0],ast.Raise) + and 'Replay model stream stopped before prepared audio ended' in ast.unparse(node)) + code=compile(ast.Module(body=[guard],type_ignores=[]),str(path),'exec') + proof=self.proof() + tail=RecordedTail(proof);tail.observe(proof['terminal_model_ns'],300.,300*48000) + tail.observe(proof['terminal_model_ns'],303.,303*48000) + context=dict(started=1.,now=303.,received={'modelV2':300.},position=303*48000,rate=48000, + audio=range(round(306.2*48000)),tail=tail, + sm=SimpleNamespace(logMonoTime={'modelV2':proof['terminal_model_ns']})) + exec(code,context) + context['sm'].logMonoTime['modelV2']-=50_000_000 + with self.assertRaisesRegex(RuntimeError,'stream stopped'):exec(code,context) + context['sm'].logMonoTime['modelV2']=proof['terminal_model_ns'] + context['now']=308. + with self.assertRaisesRegex(RuntimeError,'stream stopped'):exec(code,context) + context['now']=303.;context['tail']=RecordedTail(None) + with self.assertRaisesRegex(RuntimeError,'stream stopped'):exec(code,context) + # A stuck callback cannot hide forever inside the old final-two-second grace. + context.update(now=305.,position=305*48000,tail=tail) + with self.assertRaisesRegex(RuntimeError,'stream stopped'):exec(code,context) + if __name__=='__main__':unittest.main()