diff --git a/asset/arena/quantus-arena-exporter.py b/asset/arena/quantus-arena-exporter.py index 31a492a..4db2c10 100755 --- a/asset/arena/quantus-arena-exporter.py +++ b/asset/arena/quantus-arena-exporter.py @@ -152,7 +152,8 @@ class Telemetry: self.errors = 0 # node_id -> (name, peer id). Kept across reconnects: attributions made # while connected must still render a name after a feed hiccup. - self.names = {} + self.names = {} # node_id -> (name, peer_id); node ids are per connection + self.peer_names = {} # peer_id -> latest name; peer ids survive restarts and renames # Block imports near the tip, by hash: who reported it first and how far # ahead of the next reporter (ms). The feed stamps the first reporter of # a hash with propagation_time 0 and later reporters with their delay. @@ -178,9 +179,12 @@ class Telemetry: "cores": hw.get("core_count") or 0, "vm": bool(hw.get("is_virtual_machine")), } + name = str(d[0] or "") if len(d) > 0 else "" + peer = str(d[4] or "") if len(d) > 4 else "" with self.lock: - self.names[nid] = (str(d[0] or "") if len(d) > 0 else "", - str(d[4] or "") if len(d) > 4 else "") + self.names[nid] = (name, peer) + if peer: + self.peer_names[peer] = name elif code == 4: # RemovedNode with self.lock: self.nodes.pop(payload, None) @@ -253,6 +257,22 @@ class Telemetry: with self.lock: return self.names.get(nid, ("", "")) + def node_key(self, nid): + """Stable identity for votes: the peer id when telemetry gave one + (survives restarts and renames, which change the node id), else the + node id itself.""" + with self.lock: + peer = self.names.get(nid, ("", ""))[1] + return ("peer", peer) if peer else ("nid", nid) + + def key_name(self, key): + """Current (name, peer_id) for a vote key.""" + kind, v = key + with self.lock: + if kind == "peer": + return self.peer_names.get(v, ""), v + return self.names.get(v, ("", "")) + def render(self, g, out): with self.lock: nodes = dict(self.nodes) @@ -345,6 +365,10 @@ ATTRIBUTION_LEAD_MS = 20 ATTRIBUTION_SETTLE_S = 8.0 ATTRIBUTION_MIN_ATTEMPTS = 3 ATTRIBUTION_MIN_CONFIDENCE = 0.6 +# Votes are kept for the last this many blocks per author, so a node that +# restarts or renames (a new telemetry node id) takes over the attribution +# after a few blocks instead of having to outvote its own history. +ATTRIBUTION_WINDOW = 20 GAP_BUCKETS = [0.5, 1, 2, 3, 5, 8, 12, 16, 24, 32, 48, 64, 128] # Block timestamp minus the moment this exporter first saw the block, seconds. @@ -485,7 +509,7 @@ class Arena: # first-seen replaces it (lair/quantus#11). self.arrival_hist = {} # preimage -> Hist(GAP_BUCKETS) # Author attribution via telemetry first import (lair/quantus#11): - # preimage -> {"attempts": n, "nodes": {node_id: n}}; blocks whose first + # preimage -> {"attempts": n, "attributed": n, "votes": deque of vote keys}; blocks whose first # reporter did not lead by ATTRIBUTION_LEAD_MS count as attempts only. self.attr = {} self.pending_attr = deque() # (block hash, preimage, seen at) @@ -571,12 +595,16 @@ class Arena: name, _ = self.telemetry.node_name(fi[0]) LOG.debug("attribution: %s by %s… first import %r lead %s ms", block_hash[:10], pre[:10], name, fi[1]) with self.lock: - a = self.attr.setdefault(pre, {"attempts": 0, "nodes": {}}) + a = self.attr.setdefault(pre, {"attempts": 0, "attributed": 0, + "votes": deque(maxlen=ATTRIBUTION_WINDOW)}) a["attempts"] += 1 + key = None if fi is not None: nid, lead = fi if lead is not None and lead >= ATTRIBUTION_LEAD_MS: - a["nodes"][nid] = a["nodes"].get(nid, 0) + 1 + key = self.telemetry.node_key(nid) + a["attributed"] += 1 + a["votes"].append(key) def author_display(self, pre): """(display, node_name, peer_id, confidence). Node name when the @@ -584,11 +612,12 @@ class Arena: short = f"{pre[:8]}…{pre[-4:]}" if len(pre) > 14 else pre with self.lock: a = self.attr.get(pre) - if not a or a["attempts"] < ATTRIBUTION_MIN_ATTEMPTS or not a["nodes"]: + votes = [k for k in a["votes"] if k is not None] if a else [] + if not a or len(a["votes"]) < ATTRIBUTION_MIN_ATTEMPTS or not votes: return short, "", "", 0.0 - nid, n = max(a["nodes"].items(), key=lambda kv: kv[1]) - conf = n / a["attempts"] - name, peer = self.telemetry.node_name(nid) if self.telemetry else ("", "") + key, n = max(((k, votes.count(k)) for k in set(votes)), key=lambda kv: kv[1]) + conf = n / len(a["votes"]) + name, peer = self.telemetry.key_name(key) if self.telemetry else ("", "") if conf < ATTRIBUTION_MIN_CONFIDENCE or not name: return short, "", "", conf return name, name, peer, conf @@ -761,7 +790,7 @@ class Arena: out.append(f'quantus_author_node_info{{preimage="{pre}",display="{esc(display)}",node_name="{esc(name)}",peer_id="{peer}",self="{is_self}"}} {conf:.3f}') with self.lock: attempted = {pre: a["attempts"] for pre, a in self.attr.items()} - attributed = {pre: sum(a["nodes"].values()) for pre, a in self.attr.items()} + attributed = {pre: a["attributed"] for pre, a in self.attr.items()} out.append("# HELP quantus_author_attribution_attempts_total Tip blocks checked against telemetry first imports, by author.") out.append("# TYPE quantus_author_attribution_attempts_total counter") out.append("# HELP quantus_author_attribution_attributed_total Tip blocks whose first reporter led by enough to be called the author, by author.")