Drain verified prepared PCM tails after terminal replay model

This commit is contained in:
firestar5683
2026-09-20 12:07:43 -07:00
parent f2dbf78a75
commit 75c2a8b8a1
3 changed files with 145 additions and 5 deletions
+8 -3
View File
@@ -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')
+52 -1
View File
@@ -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):
+85 -1
View File
@@ -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()