"""Sampled single-node first-observation study. No RPC or wallet functionality here."""
import sqlite3,json,math
WINDOW=7*86400
RETENTION=90*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<now-old['at']<=30
  if old and tip!=old['tip']:
   if not block or tip[0]!=old['tip'][0]+1 or block['previousblockhash']!=old['tip'][1]:gap=True
  if reorg:
   gap=True;self.db.execute("UPDATE cohort SET state='recheck_required' WHERE state='confirmed'")
  excluded=len(mempool) if gap else 0
  if gap:self.censor(now)
  else:
   pending={r['txid'] for r in self.db.execute("SELECT txid FROM cohort WHERE state='pending'")}
   included=set(block['tx']) if block else set()
   for txid in pending & included:
    self.db.execute("UPDATE cohort SET state='confirmed',closed_at=?,confirmed_at=?,block_height=?,block_hash=? WHERE txid=?",(now,now,tip[0],tip[1],txid))
   for txid in pending-included-mempool:
    self.db.execute("UPDATE cohort SET state='left_mempool',closed_at=? WHERE txid=?",(now,txid))
   for txid in mempool-old['mempool']:
    if self.selected(txid):self.db.execute("INSERT OR IGNORE INTO cohort(txid,first_seen,state) VALUES(?,?,'pending')",(txid,now))
  self.db.execute('INSERT OR REPLACE INTO polls VALUES(?,?,?,?)',(now,'baseline' if gap else 'collecting',0 if gap else now-old['at'],excluded))
  self.previous={'at':now,'tip':tip,'mempool':mempool};self.db.commit()
 def recheck_heights(self):return [r[0] for r in self.db.execute("SELECT DISTINCT block_height FROM cohort WHERE state='recheck_required' LIMIT 50")]
 def recheck(self,canonical):
  for height,blockhash in canonical.items():
   self.db.execute("UPDATE cohort SET state=CASE WHEN block_hash=? THEN 'confirmed' ELSE 'censored_reorg' END WHERE state='recheck_required' AND block_height=?",(blockhash,height))
  self.db.commit()
 def prune(self,now):
  self.db.execute('DELETE FROM cohort WHERE first_seen<?',(now-RETENTION,));self.db.execute('DELETE FROM polls WHERE at<?',(now-RETENTION,));self.db.commit()
 def rows(self,since=0):return (dict(r) for r in self.db.execute('SELECT * FROM cohort WHERE first_seen>=? 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}
