"""Sampled single-node first-observation study. No RPC or wallet functionality here.""" import sqlite3,json,math WINDOW=7*86400 class Observer: def __init__(self,path,fraction=16): self.db=sqlite3.connect(path);self.db.row_factory=sqlite3.Row;self.fraction=fraction;self.previous=None self.db.executescript("""PRAGMA journal_mode=WAL; CREATE TABLE IF NOT EXISTS cohort(txid TEXT PRIMARY KEY,first_seen INTEGER NOT NULL,state TEXT NOT NULL,closed_at INTEGER,confirmed_at INTEGER,block_height INTEGER,block_hash TEXT); CREATE INDEX IF NOT EXISTS cohort_time ON cohort(first_seen); CREATE TABLE IF NOT EXISTS polls(at INTEGER PRIMARY KEY,state TEXT NOT NULL,covered_seconds INTEGER NOT NULL,baseline_excluded INTEGER NOT NULL); """) # A process restart loses observation continuity. Recheck prior confirmed blocks. self.db.execute("UPDATE cohort SET state='censored_gap',closed_at=COALESCE((SELECT MAX(at) FROM polls),first_seen) WHERE state='pending'") self.db.execute("UPDATE cohort SET state='recheck_required' WHERE state='confirmed'");self.db.commit() def selected(self,txid):return int(txid[:8],16)%self.fraction==0 def last_tip(self):return self.previous['tip'] if self.previous else None def censor(self,now):self.db.execute("UPDATE cohort SET state='censored_gap',closed_at=? WHERE state='pending'",(now,)) def unavailable(self,now,state): self.censor(now);self.previous=None;self.db.execute("UPDATE cohort SET state='recheck_required' WHERE state='confirmed'");self.db.execute('INSERT OR REPLACE INTO polls VALUES(?,?,0,0)',(now,state));self.db.commit() def observe(self,now,tip,mempool,block=None,reorg=False): """tip=(height,hash), block={'hash','previousblockhash','tx'}, IDs fully validated upstream.""" old=self.previous;gap=old is None or not 0=? ORDER BY first_seen,txid',(since,))) def summary(self,now,node): counts={r[0]:r[1] for r in self.db.execute('SELECT state,COUNT(*) FROM cohort WHERE first_seen>=? GROUP BY state',(now-WINDOW,))} durations=[r[0] for r in self.db.execute("SELECT confirmed_at-first_seen AS duration FROM cohort WHERE first_seen>=? AND state='confirmed' ORDER BY duration",(now-WINDOW,))] bins=[0]*5 for seconds in durations:bins[0 if seconds<=60 else 1 if seconds<=150 else 2 if seconds<=300 else 3 if seconds<=600 else 4]+=1 def q(f): if len(durations)<20:return None idx=(len(durations)-1)*f;lo=math.floor(idx);hi=math.ceil(idx);return round(durations[lo]+(durations[hi]-durations[lo])*(idx-lo),1) coverage=self.db.execute('SELECT MIN(at),MAX(at),COALESCE(SUM(covered_seconds),0),COALESCE(SUM(baseline_excluded),0) FROM polls WHERE at>=?',(now-WINDOW,)).fetchone() days=[dict(r) for r in self.db.execute("SELECT strftime('%Y-%m-%d',at,'unixepoch') AS day,COUNT(*) AS polls,SUM(covered_seconds) AS covered_seconds,SUM(baseline_excluded) AS baseline_exclusions FROM polls WHERE at>=? GROUP BY day ORDER BY day",(now-WINDOW,))] return {'schema':1,'generated_at':now,'window_start':now-WINDOW,'sample_denominator':self.fraction,'counts':counts,'histogram':bins,'confirmed_median_seconds':q(.5),'confirmed_p90_seconds':q(.9),'coverage':{'first_poll':coverage[0],'last_poll':coverage[1],'covered_seconds':coverage[2],'baseline_exclusions':coverage[3]},'days':days,'node':node}