"""Fencing fixture for the cancellation question (tantive thread 1589). A sink that (a) refuses any effect carrying a stale generation, (b) deduplicates on a stable effect_id, and (c) keeps an append-only hash chain of what it did, so a stranger can read the order rather than trust either party's clock. """ import hashlib, json def h(*parts): return hashlib.sha256('\n'.join(str(p) for p in parts).encode()).hexdigest() class Sink: def __init__(self): self.gen = 1 # current fence self.effects = {} # effect_id -> record (one commit per intended effect) self.log = [] # append-only chain self.head = '0'*64 def advance(self, why): self.gen += 1 self._append('fence', f'gen={self.gen} why={why}') return self.gen def _append(self, kind, detail, effect_id=None): e = {'i': display_seq(), 'kind': kind, 'effect_id': effect_id, 'gen': self.gen, 'detail': detail, 'prev': self.head} e['hash'] = h(e['prev'], e['i'], kind, effect_id, e['gen'], detail) self.head = e['hash']; self.log.append(e); return e def guarded_write(self, effect_id, token_gen, payload): """The only door to a side effect. One atomic step, as a conditional UPDATE would be.""" if token_gen != self.gen: self._append('refused_fenced', f'carried gen={token_gen}, sink at gen={self.gen}', effect_id) return ('FENCED_OUT', None) if effect_id in self.effects: # stable id, not the generation self._append('deduplicated', f'same effect_id already committed', effect_id) return ('DEDUPLICATED', self.effects[effect_id]) rec = {'effect_id': effect_id, 'gen': self.gen, 'payload_hash': h(payload), 'receipt': None} e = self._append('committed', f'payload {rec["payload_hash"][:16]}', effect_id) rec['receipt'] = e['hash']; self.effects[effect_id] = rec self.gen += 1 # a commit advances the fence too self._append('fence', f'gen={self.gen} why=effect committed') return ('COMMITTED_WITH_EFFECT', rec) def classify(self, cancel_gen, effect_id=None): """cancel_gen is the generation the operator's cancellation advanced the fence to. TOO_LATE the effect for that task committed at an older fence, so it stands. CANCELLED the sink attests its fence reached the boundary and holds no effect for it. UNKNOWN the sink cannot show it crossed the boundary, or an effect escaped it.""" if effect_id is not None and effect_id in self.effects: return ('TOO_LATE', self.effects[effect_id]) if self.gen < cancel_gen: return ('UNKNOWN', None) return ('CANCELLED', None) seq = {'n': 0} def display_seq(): seq['n'] += 1; return seq['n'] print('A. worker reads at gen 7, cancellation advances the fence, stale retry:') s = Sink(); s.gen = 7; s._append('fence', 'gen=7 start') print(' read:', 'gen', 7) s.advance('operator cancelled') print(' write(effect=E1, gen=7) ->', s.guarded_write('E1', 7, b'payload')[0]) print(' cancel at boundary=8 ->', s.classify(8, 'E1')[0]) print('B. same intended effect retried at the current fence, then retried again:') print(' write(effect=E1, gen=8) ->', s.guarded_write('E1', 8, b'payload')[0]) print(' write(effect=E1, gen=9) ->', s.guarded_write('E1', 9, b'payload')[0], '(one commit, second returns the record)') print(' commits on E1:', sum(1 for e in s.log if e['kind']=='committed')) print('C. opposite ordering, commit wins the race:') t = Sink(); t.gen = 7; t._append('fence','gen=7 start') print(' write(effect=E2, gen=7) ->', t.guarded_write('E2', 7, b'other')[0]) print(' cancel at boundary=8 ->', t.classify(8, 'E2')[0]) print(' log kinds:', [e['kind'] for e in t.log]) print('D. a sink that never crossed the boundary cannot say CANCELLED:') u = Sink(); u.gen = 5; u._append('fence','gen=5') print(' cancel at boundary=8 ->', u.classify(8, 'E3')[0]) print('E. chain check, scenario C:') prev = '0'*64; ok = True for e in t.log: if h(e['prev'], e['i'], e['kind'], e['effect_id'], e['gen'], e['detail']) != e['hash']: ok = False if e['prev'] != prev: ok = False prev = e['hash'] print(' entries', len(t.log), 'chain verifies', ok, 'head', prev[:16])