from dataclasses import replace import json,sqlite3,time,threading import pytest from fastapi.testclient import TestClient from aim_webgui.config import Settings from aim_webgui.auth.service import Auth from aim_webgui.journal import Journal,Capture,project from aim_webgui.app import create_app from aim_webgui.workflows import Workflows from aim_webgui.errors import WebError from evidence_fixtures import job,ev @pytest.fixture def local(tmp_path): s=Settings(state_dir=tmp_path/'state',public_url='https://aim.example.test');a=Auth(s);a.bootstrap() with a.store.transaction()as db:db.execute('UPDATE users SET must_change_password=0');uid=db.execute('SELECT id FROM users').fetchone()[0] return s,a,uid def batch(*items):return [(*project(e,{'test01.example'}),time.time())for e in items] def test_committed_snapshot_duplicate_sequences_and_metadata_only(local): s,a,u=local;ident=job(a,u);j=Journal(s);j.begin(ident) e=ev(1,raw='SYNTHETIC-SECRET-CANARY',msg={'password':'HIDDEN'}) j.append(ident,batch(e,e,ev(3,'host_result'))) snap=j.snapshot(u,ident) assert snap['cursor']==2 and len(snap['events'])==2 and snap['dropped_events']==0 assert snap['checkpoint']['observed_task_starts']==1 assert snap['checkpoint']['hosts'][0]['task_id']=='t1' with a.store.read()as db:assert 'SYNTHETIC-SECRET-CANARY'not in ''.join(r[0] for r in db.execute('SELECT payload FROM job_progress_events')) j.append(ident,batch(e));assert j.snapshot(u,ident)['cursor']==2 def test_retention_cursor_gap_and_checkpoint(local): s,a,u=local;s=replace(s,journal_max_events=100);ident=job(a,u);j=Journal(s);j.begin(ident) j.append(ident,batch(*(ev(i)for i in range(1,251)))) snap=j.snapshot(u,ident,after=0) assert snap['first_cursor']==151 and snap['cursor']==250 and snap['omitted_events']==150 and snap['gap_before'] assert snap['checkpoint']['observed_task_starts']==250 and len(snap['checkpoint']['tasks'])==64 assert len(j.snapshot(u,ident,before=201,limit=20)['events'])==20 with pytest.raises(WebError):j.snapshot(u,ident,after=9999) def test_byte_limit_and_deletion(local): s,a,u=local;s=replace(s,journal_max_bytes=65536);ident=job(a,u,status='failed');j=Journal(s);j.begin(ident) j.append(ident,batch(*(ev(i,label='X'*200)for i in range(1,400))),closed=True) snap=j.snapshot(u,ident);assert snap['retained_bytes']<=65536 and snap['omitted_events']>0 Workflows(s).delete_jobs(u,[ident]) with a.store.read()as db: assert db.execute('SELECT COUNT(*) FROM job_progress_events').fetchone()[0]==0 assert db.execute('SELECT COUNT(*) FROM job_progress_state').fetchone()[0]==0 def test_capture_independent_of_viewers_and_end(local): s,a,u=local;ident=job(a,u);c=Capture(s,ident,['test01.example']) for i in range(1,31):c.submit(ev(i)) assert c.close() snap=Journal(s).snapshot(u,ident);assert len(snap['events'])==30 and snap['capture_state']=='closed' with a.store.transaction()as db:db.execute("UPDATE jobs SET status='successful' WHERE id=?",(ident,)) token,_=a.new_session(u) with TestClient(create_app(s),base_url=s.public_url)as client: client.cookies.set(s.cookie_name,token) first=client.get('/api/v2/runs/'+ident+'/progress').json() replay=client.get('/api/v2/runs/'+ident+'/progress/stream?after=0',headers={'Last-Event-ID':'20'}) assert replay.status_code==200 and replay.text.count('event: line\n')==10 assert 'id: 21\n'in replay.text and 'id: 20\n'not in replay.text assert 'event: end'in replay.text second=client.get('/api/v2/runs/'+ident+'/progress').json();assert first['cursor']==second['cursor']==30 def test_owner_scope_active_delete_and_crash_marker(local): s,a,u=local;a.create_user('viewer','Synthetic-password-value-2026','viewer') with a.store.transaction()as db:db.execute('UPDATE users SET must_change_password=0');v=db.execute("SELECT id FROM users WHERE username='viewer'").fetchone()[0] ident=job(a,u);j=Journal(s);j.begin(ident);j.append(ident,batch(ev(1))) with pytest.raises(WebError):j.snapshot(v,ident) with pytest.raises(WebError):Workflows(s).delete_jobs(u,[ident]) with a.store.transaction()as db:db.execute("UPDATE jobs SET status='interrupted' WHERE id=?",(ident,)) assert j.snapshot(u,ident)['capture_interrupted'] def test_queue_pressure_is_bounded_and_disclosed(local,monkeypatch): s,a,u=local;ident=job(a,u) monkeypatch.setattr('aim_webgui.journal.QUEUE_EVENTS',2) original=Journal.append def slow(self,*args,**kwargs):time.sleep(.05);return original(self,*args,**kwargs) monkeypatch.setattr(Journal,'append',slow) c=Capture(s,ident,['test01.example']);start=time.monotonic() for i in range(1,301):c.submit(ev(i)) assert time.monotonic()-start<1 c.close();snap=Journal(s).snapshot(u,ident) assert snap['dropped_events']>0 and snap['capture_state']=='closed_with_gaps' def test_storage_failure_does_not_throw_from_sink(local,monkeypatch): s,a,u=local;ident=job(a,u) def fail(*args,**kwargs):raise sqlite3.OperationalError('SYNTHETIC disk full; never logged') monkeypatch.setattr(Journal,'append',fail) c=Capture(s,ident,['test01.example']);c.submit(ev(1));c.close() with a.store.transaction()as db:db.execute("UPDATE jobs SET status='failed' WHERE id=?",(ident,)) assert Journal(s).snapshot(u,ident)['capture_interrupted'] def test_core_result_event_is_not_final_success(local): s,a,u=local;ident=job(a,u);j=Journal(s);j.begin(ident) j.append(ident,batch(ev(1,'result',result={'status':'succeeded','operation_result':{'raw':'NOT-STORED'}}))) snap=j.snapshot(u,ident);assert snap['job_status']=='running' and not snap['terminal'] assert snap['checkpoint']['result_event_observed'] and 'NOT-STORED'not in json.dumps(snap) def test_legacy_and_new_no_events_are_explicit(local): s,a,u=local;ident=job(a,u,status='successful') snap=Journal(s).snapshot(u,ident);assert not snap['available'] and 'not captured'in snap['message'] def test_wrong_identity_and_unknown_additive_payloads(local): s,a,u=local;ident=job(a,u);c=Capture(s,ident,['test01.example']) c.submit(ev(1,'host_result',host='someone-else.example')) c.submit(ev(2,'progress',status='ok',msg='NO-SAVE')) c.submit(ev(3,'task_started',environment={'secret':'NO-SAVE'})) c.close();snap=Journal(s).snapshot(u,ident) assert snap['dropped_events']==1 and len(snap['events'])==1 and 'NO-SAVE'not in json.dumps(snap) def test_two_live_authorized_streams_share_committed_cursor_and_revocation(local): s,a,u=local;ident=job(a,u);j=Journal(s);j.begin(ident);j.append(ident,batch(ev(1))) tokens=[a.new_session(u)[0] for _ in range(2)];responses=[];errors=[] def view(token): try: with TestClient(create_app(s),base_url=s.public_url)as client: client.cookies.set(s.cookie_name,token) responses.append(client.get('/api/v2/runs/'+ident+'/progress/stream').text) except Exception as e:errors.append(type(e).__name__) threads=[threading.Thread(target=view,args=(t,),daemon=True)for t in tokens] for t in threads:t.start() time.sleep(.3);j.append(ident,batch(ev(2),ev(3))) time.sleep(.3) with a.store.transaction()as db:db.execute("UPDATE jobs SET status='successful' WHERE id=?",(ident,)) for t in threads:t.join(timeout=5) assert not errors and len(responses)==2 for text in responses: assert text.count('event: line\n')==3 and text.count('id: 3\n')==1 and 'event: end'in text # A separate active stream must stop when its server-side session is revoked. ident=job(a,u);j.begin(ident);responses.clear() token=a.new_session(u)[0];thread=threading.Thread(target=view,args=(token,),daemon=True);thread.start() time.sleep(.3) with a.store.transaction()as db:db.execute('DELETE FROM sessions') thread.join(timeout=5) assert responses and '"clear":true'in responses[0] assert not thread.is_alive()