"""Six controlled stop/restart schedules; SQLite state plus an in-memory broker model.""" import csv,json,sqlite3 from pathlib import Path class Stop(Exception): pass def case(name,write_mode='atomic',write_stop=False,relay_stop=None,deduplicate=True): db=sqlite3.connect(':memory:') db.executescript('CREATE TABLE orders(id TEXT PRIMARY KEY); CREATE TABLE outbox(id TEXT PRIMARY KEY, sent INTEGER NOT NULL DEFAULT 0);') trace=[];broker=[] try: db.execute('BEGIN');db.execute("INSERT INTO orders VALUES('order-1')") if write_mode=='atomic':db.execute("INSERT INTO outbox(id) VALUES('event-1')") if write_stop:raise Stop() db.commit();trace.append('business transaction committed') except Stop: db.rollback();trace.append('stop before commit: rolled back') def relay(stop=None): for (event,) in db.execute('SELECT id FROM outbox WHERE sent=0 ORDER BY id').fetchall(): if stop=='mark-first': db.execute('UPDATE outbox SET sent=1 WHERE id=?',(event,));db.commit();trace.append('marked sent before publishing');raise Stop() broker.append(event);trace.append('broker accepted '+event) if stop=='after-publish':raise Stop() db.execute('UPDATE outbox SET sent=1 WHERE id=?',(event,));db.commit();trace.append('marked sent') if write_mode=='direct':trace.append('stop between database commit and direct publication') elif relay_stop=='before-relay':trace.append('stop before relay starts') else: try:relay(relay_stop) except Stop:trace.append('relay stopped') pending_before_restart=db.execute('SELECT count(*) FROM outbox WHERE sent=0').fetchone()[0] relay();trace.append('restart scan complete') consumer=sqlite3.connect(':memory:');consumer.executescript('CREATE TABLE receipts(id TEXT PRIMARY KEY); CREATE TABLE effects(event TEXT);') for event in broker: with consumer: if deduplicate: inserted=consumer.execute('INSERT OR IGNORE INTO receipts VALUES(?)',(event,)).rowcount if not inserted:trace.append('consumer skipped duplicate');continue consumer.execute('INSERT INTO effects VALUES(?)',(event,)) result=dict(case=name,orders=db.execute('SELECT count(*) FROM orders').fetchone()[0],pending_before_restart=pending_before_restart,deliveries=len(broker),effects=consumer.execute('SELECT count(*) FROM effects').fetchone()[0],pending_final=db.execute('SELECT count(*) FROM outbox WHERE sent=0').fetchone()[0],trace=trace) db.close();consumer.close();return result def run(): rows=[case('Direct-write gap',write_mode='direct'),case('Stop before commit',write_stop=True),case('Stop before relay',relay_stop='before-relay'),case('Publish then stop',relay_stop='after-publish'),case('Mark then stop',relay_stop='mark-first'),case('No consumer dedup',relay_stop='after-publish',deduplicate=False)] expected=[(1,0,0,0),(0,0,0,0),(1,1,1,1),(1,1,2,1),(1,0,0,0),(1,1,2,2)] assert [(r['orders'],r['pending_before_restart'],r['deliveries'],r['effects']) for r in rows]==expected assert all(r['pending_final']==0 for r in rows) return {'scope':'Controlled sequential stops; no process kill, real broker, network, concurrency or external side effect','rows':rows} if __name__=='__main__': p=Path(__file__).parent;r=run();assert r==run();(p/'results.json').write_text(json.dumps(r,indent=2)+'\n') with (p/'results.csv').open('w',newline='') as f: fields=['case','orders','pending_before_restart','deliveries','effects','pending_final'];w=csv.DictWriter(f,fields,extrasaction='ignore');w.writeheader();w.writerows(r['rows']) print(json.dumps(r))