Reproduced inclusive polling consumption defect
Pinned sourceb6960706a6ea1797a63e901c34f06cabd0f7e7a5: dashboard.py#L424-L520 and inbox.py#L90-L171. Dashboard uses latest received_at as next since boundary; reader includes entries equal to since. Consumer unconditionally appends/counts returned boundary entries again.
One unchanged raw stored fixture followed by two actual poll-method calls produces two displayed rows, two history/export-source rows and counter2 instead of one consumption. Transport/agent counters double too. Appending a distinct entry at the same timestamp produces first/first/second. Thus a simplistic strict timestamp cutoff would lose legitimate later boundary arrivals.
Expected: previously consumed entry is not new activity, while distinct same-timestamp arrivals remain visible. Actual: idle polls duplicate rows/activity. Suggested fix: deduplicate consumed boundary identities at the dashboard consumer, preserving reader semantics and same-timestamp distinct arrivals. No signature/security or repeated notification claim.
Source links:
|
specific = self._route_table_id(str(row.get("transport", "")).lower(), str(row.get("kind", ""))) |
|
if specific != "tbl-all": |
|
self._add_row(specific, row_t) |
|
self._visible_rows.append(row) |
|
if len(self._visible_rows) > 2000: |
|
self._visible_rows = self._visible_rows[-2000:] |
|
|
|
def _rebuild_filtered_view(self) -> None: |
|
self._clear_rows() |
|
self._visible_rows = [] |
|
for row in self._history_rows: |
|
if _row_matches_query(row, self._filter_query): |
|
self._display_row(row) |
|
|
|
def _poll_inbox(self) -> None: |
|
entries = read_inbox(since=self._last_ts, limit=500) |
|
if not entries: |
|
return |
|
for entry in entries: |
|
rts = float(entry.get("received_at") or 0.0) |
|
if rts > self._last_ts: |
|
self._last_ts = rts |
|
|
|
row = _entry_to_row(entry) |
|
self._history_rows.append(row) |
|
if len(self._history_rows) > 5000: |
|
self._history_rows = self._history_rows[-5000:] |
|
|
|
transport = str(row.get("transport", "")).lower() |
|
self._count_today += 1 |
|
self._transport_counter[transport] += 1 |
|
self._agent_counter[str(row.get("agent", "unknown"))] += 1 |
|
|
|
if _row_matches_query(row, self._filter_query): |
|
self._display_row(row) |
|
|
|
rtc = row.get("rtc_value") |
|
kind = str(row.get("kind") or "").lower() |
|
high_value = isinstance(rtc, float) and rtc >= 5 |
|
mayday = kind == "mayday" |
|
if high_value or mayday: |
|
label = str(row.get("agent", "unknown")) |
|
if rtc is not None: |
|
self.notify(f"{kind.upper()} from {label} ({rtc:g} RTC)", severity="warning", timeout=4) |
|
else: |
|
self.notify(f"{kind.upper()} from {label}", severity="warning", timeout=4) |
|
if sound: |
|
print("\a", end="", flush=True) |
|
|
|
def _poll_api(self) -> None: |
|
self._api_state = fetch_beacon_snapshot( |
|
api_base_url=api_base_url, |
|
timeout_s=8.0, |
|
session=self._http, |
|
) |
|
|
|
def _refresh_sidebar(self) -> None: |
|
top_agents = self._agent_counter.most_common(5) |
|
lines = [ |
|
"[b]Beacon Network[/b]", |
|
"", |
|
f"Pings today: {self._count_today}", |
|
f"Visible rows: {len(self._visible_rows)}", |
|
f"Filter: {self._filter_query or '(none)'}", |
|
"", |
|
"Beacon API:", |
|
f"- base: {api_base_url}", |
|
f"- status: {'OK' if self._api_state.get('ok') else 'degraded'}", |
|
f"- agents: {self._api_state.get('agents_count', 0)}", |
|
f"- contracts: {self._api_state.get('contracts_count', 0)}", |
|
f"- reputation: {self._api_state.get('reputation_count', 0)}", |
|
] |
|
if self._api_state.get("errors"): |
|
lines.append(f"- last_error: {self._api_state['errors'][0]}") |
|
|
|
lines.append("") |
|
lines.append("Transports:") |
|
for t in [ |
|
"udp", |
|
"webhook", |
|
"discord", |
|
"bottube", |
|
"rustchain", |
|
"moltbook", |
|
"clawcities", |
|
"clawsta", |
|
"fourclaw", |
|
"pinchedin", |
|
"clawtasks", |
|
"clawnews", |
|
]: |
|
n = self._transport_counter.get(t, 0) |
|
marker = "[green]*[/green]" if n > 0 else "[red]*[/red]" |
|
lines.append(f"{marker} {t}: {n}") |
|
|
|
lines.append("") |
|
lines.append("Top agents:") |
and
|
|
|
return keys |
|
|
|
|
|
def read_inbox( |
|
*, |
|
kind: Optional[str] = None, |
|
agent_id: Optional[str] = None, |
|
since: Optional[float] = None, |
|
unread_only: bool = False, |
|
limit: Optional[int] = None, |
|
) -> List[Dict[str, Any]]: |
|
"""Read and filter inbox entries from inbox.jsonl. |
|
|
|
Each entry is enriched with: |
|
- verified: True/False/None (signature verification result) |
|
- is_read: bool (whether this nonce was marked read) |
|
""" |
|
path = _inbox_path() |
|
if not path.exists(): |
|
return [] |
|
|
|
known_keys = load_known_keys() |
|
read_nonces = _read_nonces() |
|
results: List[Dict[str, Any]] = [] |
|
|
|
for line in path.read_text(encoding="utf-8").splitlines(): |
|
line = line.strip() |
|
if not line: |
|
continue |
|
try: |
|
entry = json.loads(line) |
|
except Exception: |
|
continue |
|
|
|
# Extract envelopes from the entry. |
|
envelopes = entry.get("envelopes", []) |
|
if not envelopes and entry.get("text"): |
|
envelopes = decode_envelopes(entry["text"]) |
|
|
|
# Process each envelope in the entry. |
|
for env in envelopes: |
|
# Auto-learn keys (with TTL tracking). |
|
known_keys = _learn_key_from_envelope(env, known_keys) |
|
|
|
# Verify signature. |
|
verified = verify_envelope(env, known_keys=known_keys) |
|
nonce = env.get("nonce", "") |
|
is_read = nonce in read_nonces if nonce else False |
|
|
|
enriched = dict(entry) |
|
enriched["envelope"] = env |
|
enriched["verified"] = verified |
|
enriched["is_read"] = is_read |
|
|
|
# Apply filters. |
|
if kind and env.get("kind") != kind: |
|
continue |
|
if agent_id and env.get("agent_id") != agent_id: |
|
continue |
|
if since and entry.get("received_at", 0) < since: |
|
continue |
|
if unread_only and is_read: |
|
continue |
|
|
|
results.append(enriched) |
|
|
|
# If no envelopes, include the raw entry (e.g., plain text UDP). |
|
if not envelopes: |
|
enriched = dict(entry) |
|
enriched["envelope"] = None |
|
enriched["verified"] = None |
|
enriched["is_read"] = False |
|
|
|
if kind or agent_id: |
|
continue # Can't filter raw entries by kind/agent_id. |
|
if since and entry.get("received_at", 0) < since: |
|
continue |
|
|
|
results.append(enriched) |
|
|
|
# Save any updated keys (last_seen timestamps, etc.) |
.
Duplicate exact polling searches0; broader dashboard records were feature/docs/packaging/base-URL work, no boundary-repeat report identified.
Reproduction and controls
Python3 stdlib on macOS arm64. Save pinned dashboard.py, inbox.py and this script as repro.py together in an isolated directory, run python3 repro.py. Original reader/poll/display/API/sidebar function ASTs execute unchanged; widgets and HTTP are local fakes. Inbox is a temporary JSONL fixture; known-key/state seams return empty values or no-op. No Textual app, UDP, quick-send, identities/keys, ~/.beacon, live API, account or production writes.
import ast
import json
import time
from pathlib import Path
from collections import Counter
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
from types import SimpleNamespace
ROOT = Path(__file__).parent
ns = dict(globals(), requests=SimpleNamespace(Session=lambda: None), sound=False, api_base_url='https://fixture.invalid/beacon', DEFAULT_API_BASE_URL='https://fixture.invalid/beacon')
def load_functions(filename, names):
tree = ast.parse((ROOT / filename).read_text())
funcs = [node for node in ast.walk(tree) if isinstance(node, ast.FunctionDef) and node.name in names]
assert len(funcs) == len(names), (filename, names)
exec(compile(ast.fix_missing_locations(ast.Module(body=funcs, type_ignores=[])), filename, 'exec'), ns)
# Only original function ASTs execute: no module imports, app construction, identity, network or send path.
load_functions('inbox.py', {'read_inbox'})
fixture = ROOT / 'fixture.jsonl'
ns.update(_inbox_path=lambda: fixture, load_known_keys=lambda: {}, _read_nonces=lambda: set(), save_known_keys=lambda keys: None, decode_envelopes=lambda text: [])
load_functions('dashboard.py', {'_format_ts', '_short_agent', '_as_text', '_rtc_tip', '_transport_tag', '_entry_to_row', '_row_matches_query', '_normalize_api_rows', 'fetch_beacon_snapshot', '_route_table_id', '_add_row', '_display_row', '_poll_inbox', '_poll_api', '_refresh_sidebar'})
class Widget:
def __init__(self): self.rows = {}; self.text = ''
@property
def row_count(self): return len(self.rows)
def add_row(self, *row): self.rows[len(self.rows)] = row
def update(self, text): self.text = text
class Dashboard:
def __init__(self):
self._last_ts = 0; self._count_today = 0
self._transport_counter = Counter(); self._agent_counter = Counter()
self._history_rows = []; self._visible_rows = []; self._filter_query = ''
self._api_state = {}; self.widgets = {}; self.notifications = []
def query_one(self, name, kind): return self.widgets.setdefault(name, Widget())
def notify(self, *args, **kwargs): self.notifications.append((args, kwargs))
for name in ('_route_table_id', '_add_row', '_display_row', '_poll_inbox', '_poll_api', '_refresh_sidebar'):
setattr(Dashboard, name, ns[name])
ns.update(DataTable=Widget, Static=Widget)
def write_entries(entries): fixture.write_text(''.join(json.dumps(e) + '\n' for e in entries))
a = {'received_at': 1000, 'platform': 'udp', 'from': 'fixture-a', 'text': 'first'}
b = {'received_at': 1000, 'platform': 'udp', 'from': 'fixture-b', 'text': 'second'}
write_entries([a]); d = Dashboard(); d._poll_inbox(); d._poll_inbox()
print('UNCHANGED: expected pings/history/visible/table=1/1/1/1; actual=' + '/'.join(map(str, [d._count_today, len(d._history_rows), len(d._visible_rows), d.widgets['#tbl-all'].row_count])))
print('UNCHANGED counters:', dict(d._transport_counter), dict(d._agent_counter))
assert d._count_today == 2 and len(d._history_rows) == 2
write_entries([a]); d = Dashboard(); d._poll_inbox(); write_entries([a, b]); d._poll_inbox()
print('EQUAL_TIMESTAMP_APPEND: expected distinct pings=2; actual pings=%s messages=%s' % (d._count_today, [r['message'] for r in d._history_rows]))
print('INCLUSIVE_READER: since=1000 returns', [e['text'] for e in ns['read_inbox'](since=1000)])
assert [r['message'] for r in d._history_rows] == ['first', 'first', 'second']
class Response:
status_code = 200
def __init__(self, payload): self.payload = payload
def json(self):
if self.payload == 'BAD': raise ValueError('invalid JSON fixture')
return self.payload
class Session:
def __init__(self, payloads): self.payloads = iter(payloads); self.calls = []
def request(self, method, url, timeout):
self.calls.append((method, url, timeout)); return Response(next(self.payloads))
for label, payloads in [('ALL_INVALID', ['BAD'] * 3), ('VALID_EMPTY', [[], [], []]), ('ONE_INVALID', [[{'id': 'fixture-agent'}], 'BAD', [{'id': 'fixture-reputation'}]])]:
d = Dashboard(); d._http = Session(payloads); d._poll_api(); d._refresh_sidebar()
state = d._api_state
print(label + ': ok=%s errors=%s counts=%s' % (state['ok'], state['errors'], [state[k + '_count'] for k in ('agents', 'contracts', 'reputation')]))
print(label + ' sidebar:', ' | '.join(line for line in d.widgets['#sidebar'].text.splitlines() if line.startswith(('- status:', '- agents:', '- contracts:', '- reputation:', '- last_error:'))))
assert len(d._http.calls) == 3
assert state['ok'] is True and state['errors'] == []
if label == 'ONE_INVALID': assert state['agents_count'] == state['reputation_count'] == 1
print('Reproduced both defects through original consumer methods; no external runtime I/O.')
Executed output (exit0 checks observed failure signatures, not corrected behavior):
UNCHANGED: expected pings/history/visible/table=1/1/1/1; actual=2/2/2/2
UNCHANGED counters: {'udp': 2} {'fixture-a': 2}
EQUAL_TIMESTAMP_APPEND: expected distinct pings=2; actual pings=3 messages=['first', 'first', 'second']
INCLUSIVE_READER: since=1000 returns ['first', 'second']
ALL_INVALID: ok=True errors=[] counts=[0, 0, 0]
ALL_INVALID sidebar: - status: OK | - agents: 0 | - contracts: 0 | - reputation: 0
VALID_EMPTY: ok=True errors=[] counts=[0, 0, 0]
VALID_EMPTY sidebar: - status: OK | - agents: 0 | - contracts: 0 | - reputation: 0
ONE_INVALID: ok=True errors=[] counts=[1, 0, 1]
ONE_INVALID sidebar: - status: OK | - agents: 1 | - contracts: 0 | - reputation: 1
Reproduced both defects through original consumer methods; no external runtime I/O.
AI assistance disclosed. Bounded novelty searches and local repros do not assert deployed incidents or accepted awards.
Reproduced inclusive polling consumption defect
Pinned sourceb6960706a6ea1797a63e901c34f06cabd0f7e7a5: dashboard.py#L424-L520 and inbox.py#L90-L171. Dashboard uses latest received_at as next since boundary; reader includes entries equal to since. Consumer unconditionally appends/counts returned boundary entries again.
One unchanged raw stored fixture followed by two actual poll-method calls produces two displayed rows, two history/export-source rows and counter2 instead of one consumption. Transport/agent counters double too. Appending a distinct entry at the same timestamp produces first/first/second. Thus a simplistic strict timestamp cutoff would lose legitimate later boundary arrivals.
Expected: previously consumed entry is not new activity, while distinct same-timestamp arrivals remain visible. Actual: idle polls duplicate rows/activity. Suggested fix: deduplicate consumed boundary identities at the dashboard consumer, preserving reader semantics and same-timestamp distinct arrivals. No signature/security or repeated notification claim.
Source links:
beacon-skill/beacon_skill/dashboard.py
Lines 424 to 520 in b696070
beacon-skill/beacon_skill/inbox.py
Lines 90 to 171 in b696070
Duplicate exact polling searches0; broader dashboard records were feature/docs/packaging/base-URL work, no boundary-repeat report identified.
Reproduction and controls
Python3 stdlib on macOS arm64. Save pinned
dashboard.py,inbox.pyand this script asrepro.pytogether in an isolated directory, runpython3 repro.py. Original reader/poll/display/API/sidebar function ASTs execute unchanged; widgets and HTTP are local fakes. Inbox is a temporary JSONL fixture; known-key/state seams return empty values or no-op. No Textual app, UDP, quick-send, identities/keys, ~/.beacon, live API, account or production writes.Executed output (exit0 checks observed failure signatures, not corrected behavior):
AI assistance disclosed. Bounded novelty searches and local repros do not assert deployed incidents or accepted awards.