Something went wrong. Try again.
Central platform for European atproto.<cc> country community websites atproto.eu
community atproto
Something went wrong. Try again.
Python
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725# /// script# dependencies = []# ///"""Union the NL-stats feeders into the public `nl-stats.json` the /stats page reads.
Each feeder writes a `<source>.dids` file (one DID per line) into a shared directoryon infra we control: - dutchPoster : nl-live.json.dids (the firehose collector, estimate-nl-users.py) Windowed, so it is accumulated into a persistent dutch-posters.dids master (union across runs) that only grows. Pass --no-master to skip. - nlHandle : nl-handles.dids (the PLC .nl-handle scan, nl-handles-scan.py) - starterpack : starterpack.dids (Dutch starterpack members -- feeder TBD) - sifaNl : sifa-nl.dids (Sifa ID NL-location accounts -- feeder TBD)
This unions those DIDs so overlap is deduplicated, not double-counted: `identified.total`is the size of the union (the headline floor), and `bySource` holds the per-feeder counts(non-exclusive -- one DID can be in several). It also folds in the activity estimate fromthe collector's full result. Output is COUNTS ONLY: the DID lists stay private on the hostthat generates them, and only `nl-stats.json` (which carries no DIDs) is committed to the public repo.
The output matches src/data/nl-stats.ts: total is always >= the largest single source and<= the summed sources (both hold for any real union), so the site's build-time guard passes.
Usage: uv run scripts/nl-stats-aggregate.py [--dir DIR] [--out nl-stats.json] [--estimate nl-live.json] [--dutch-poster P] [--nl-handle P] [--starterpack P] [--sifa-nl P] uv run scripts/nl-stats-aggregate.py --self-test
Run it where the .dids live, then commit the resulting nl-stats.json tosrc/data/nl-stats.json. Idempotent; safe to re-run. Missing feeder files count as 0."""import argparseimport datetimeimport jsonimport osimport statisticsimport sysimport urllib.request
# CONFIRMING sources answer "is this account in the Netherlands?" with evidence a person# chose: a .nl handle, a curated starterpack, a declared Sifa location, a company domain.# Their union is the identified floor.CONFIRMING_SOURCES = ["nlHandle", "starterpack", "sifaNl", "sifaCompany"]
# SIGNAL sources do not confirm a country and are NOT unioned into the floor. dutchPoster is# "posted something in Dutch", and Dutch is spoken in Belgium too -- roughly a third of Dutch# speakers live there. Counting it as confirmation contradicted the project's own rule# ("Language never confirms a country") and silently absorbed Flemish accounts into a Dutch# total, while the estimate layer below multiplies by NL_SHARE_OF_DUTCH precisely because it# knows they are in there. It is reported separately, with its overlap with the floor.SIGNAL_SOURCES = ["dutchPoster"]
# Every feeder we read a file for, in the order the schema lists them.SOURCES = ["nlHandle", "dutchPoster", "starterpack", "sifaNl", "sifaCompany"]# When the shared attribution pass is in play the per-source breakdown MUST be computed from the# same files it read, or the floor can contain accounts no source claims and the union invariant# fails. handles.nl.dids (the 45-country scan) supersedes nl-handles.dids (the old NL-only one).DEFAULT_FILE = { "nlHandle": "handles.nl.dids", "dutchPoster": "nl-live.json.dids", "starterpack": "starterpack.dids", "sifaNl": "sifa-nl.dids", "sifaCompany": "sifa-company.dids",}# Live Bluesky account total (the extrapolation anchor). It moves, so the aggregator fetches it# fresh each run from jazco. Falls back to the collector's value, then to this constant if both# are unreachable. Keep the constant roughly in step with src/data/nl-estimate.ts.# Tie-break order for crediting which detector first identified an account. See# record_first_seen: arrival signal first, catch-up scans after.# Confirming sources only: an account we cannot place in the Netherlands should not get a# "first seen in the NL community" date from having posted in Dutch.FIRST_SEEN_PRECEDENCE = ["nlHandle", "starterpack", "sifaNl", "sifaCompany"]SCRIPT_VERSION = 2 # bump when a change alters what the numbers mean, not on every editJAZCO_STATS = "https://bsky-search.jazco.io/stats"DEFAULT_BSKY_TOTAL = 46_184_219
def fetch_bluesky_total(): """Current Bluesky account total from jazco, fetched fresh so blueskyActive tracks the moving target instead of freezing on the fallback. Returns None on any failure.""" try: req = urllib.request.Request( JAZCO_STATS, headers={"accept": "application/json", "user-agent": "atproto.nl-nl-stats/1.0"} ) with urllib.request.urlopen(req, timeout=15) as r: n = json.load(r).get("total_users") return int(n) if isinstance(n, (int, float)) and n > 0 else None except Exception: return None
def read_dids(path): """Read a .dids file into a set of DIDs. Missing file -> empty set.""" if not path or not os.path.exists(path): return set() out = set() with open(path) as f: for line in f: did = line.strip() if did and not did.startswith("#"): out.add(did) return out
def accumulate_master(window_path, master_path): """Merge a windowed feeder set into a persistent master and write it back.
The Dutch-poster feeder is windowed: each collector run only sees who posted Dutch during that run, and the count resets when the collector restarts. Unioning every run's DIDs into a master `dutch-posters.dids` gives the "ever seen posting Dutch, never forget" floor the /stats page wants: it only grows. The other feeders are current-state snapshots (a `.nl` handle that lapses SHOULD drop out), so they are used as-is and never accumulated. """ merged = read_dids(master_path) | read_dids(window_path) with open(master_path, "w") as f: for did in sorted(merged): f.write(did + "\n") return merged
def load_json(path): if not path or not os.path.exists(path): return None try: with open(path) as f: return json.load(f) except (OSError, ValueError): return None
def iso_mtime(path): if not path or not os.path.exists(path): return None return ( datetime.datetime.fromtimestamp(os.path.getmtime(path), datetime.timezone.utc) .strftime("%Y-%m-%d") )
# The anchored activity model, kept in sync with src/data/nl-estimate.ts and# /community-size. The Dutch post share swings a lot over a day (it peaks in the# NL evening), so a short collector window gives a noisy, misleading estimate.# Until the collector has a representative window, we publish this stable anchor so# /stats and /community-size agree; once the window is long enough we switch to live.ANCHOR_ESTIMATE = { "dutchPostSharePct": 0.761, "blueskyActive": 46_184_219, "windowMinutes": 1100, "modelledNl": 230_000,}MIN_LIVE_WINDOW_MIN = 1440 # 24h: only trust the live share past a full day
# Error bounds are derived from how much the daily Dutch-post share actually moves,# not from a guess. Each completed collector day contributes one sample; below this# many samples we publish no range at all rather than a fake one.MIN_SHARE_SAMPLES = 5NL_SHARE_OF_DUTCH = 0.65 # NL's share of all Dutch speakers, the model's last multiplier
def build_estimate(collector, bsky_override=None): """The /stats estimate tier: the anchored model, or the live one once the collector window is representative (>= 24h and a plausible share). `bsky_override` is the freshly-fetched Bluesky total (see aggregate()); kept as a param so this stays pure/testable.""" c = collector or {} share = c.get("nl_share_pct") window_min = round((c.get("window_seconds") or 0) / 60) if not share or window_min < MIN_LIVE_WINDOW_MIN or not (0.3 <= share <= 2.0): return dict(ANCHOR_ESTIMATE) bsky = bsky_override or c.get("bluesky_total_users") or DEFAULT_BSKY_TOTAL dutch_accounts = bsky * share / 100 # NL is ~0.65 of all Dutch speakers; round to the nearest 10k like /community-size. modelled = round(dutch_accounts * 0.65 / 10_000) * 10_000 return { "dutchPostSharePct": share, "blueskyActive": bsky, "windowMinutes": window_min, "modelledNl": modelled, }
def exclusive_counts(dids_by_source): """How many DIDs each source found that NO other source found.
This is what makes overlap showable without set algebra: a bar per source split into "only here" and "also seen elsewhere". The sets are already in memory, so it costs nothing. """ out = {} for src, own in dids_by_source.items(): others = set() for other_src, other in dids_by_source.items(): if other_src != src: others |= other out[src] = len(own - others) return out
# --- daily share samples, error bounds, history -------------------------------# The three pieces the /community-size redesign needs, all counts-only so they can# live in the public repo alongside nl-stats.json.
def read_share_samples(path): """Read the daily Dutch-post-share samples (one JSON object per line).
One sample per UTC date, each from a completed collector day. Malformed lines are skipped rather than fatal: this file is appended to by a cron job on another host, and a half-written line must never break the build. """ out = [] if not path or not os.path.exists(path): return out with open(path) as f: for line in f: line = line.strip() if not line or line.startswith("#"): continue try: row = json.loads(line) except ValueError: continue if isinstance(row.get("date"), str) and isinstance(row.get("sharePct"), (int, float)): out.append(row) # One sample per date; a re-run of the same day replaces the earlier line. by_date = {} for row in out: by_date[row["date"]] = row return [by_date[d] for d in sorted(by_date)]
def record_share_sample(path, collector, date=None): """Append today's completed collector window as one share sample. Returns the row.
Only a window of at least MIN_LIVE_WINDOW_MIN counts: a partial day is not a sample of a day. Returns None when there is nothing worth recording. """ c = collector or {} share = c.get("nl_share_pct") window_min = round((c.get("window_seconds") or 0) / 60) if not share or window_min < MIN_LIVE_WINDOW_MIN or not (0.3 <= share <= 2.0): return None row = { "date": date or datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%d"), "sharePct": share, "windowMinutes": window_min, "dutchPosts": c.get("nl_posts"), "totalPosts": c.get("total_posts"), } existing = {r["date"]: r for r in read_share_samples(path)} existing[row["date"]] = row with open(path, "w") as f: for d in sorted(existing): f.write(json.dumps(existing[d], sort_keys=True) + "\n") return row
def derive_bounds(estimate, samples, identified_total): """Error bounds from the observed day-to-day spread of the Dutch post share.
The share is the only input that actually moves; the account total and the NL share of Dutch speakers are near-constant on this timescale. So the honest range is what the daily samples say the share could be, pushed through the same model.
Hard floor: the low end can never sit below the accounts we have actually identified. We can point at those; no model may claim fewer.
Returns {} until there are enough samples -- no range beats a made-up range. """ shares = [r["sharePct"] for r in samples] if len(shares) < MIN_SHARE_SAMPLES: return {} mean = statistics.fmean(shares) sd = statistics.stdev(shares) bsky = estimate.get("blueskyActive") or DEFAULT_BSKY_TOTAL
def to_accounts(share_pct): return round(bsky * share_pct / 100 * NL_SHARE_OF_DUTCH / 10_000) * 10_000
lo = to_accounts(max(mean - 1.96 * sd, 0.01)) hi = to_accounts(mean + 1.96 * sd) lo = max(lo, identified_total) # never below what we can point at if hi <= lo: return {} return { "modelledLo": lo, "modelledHi": hi, "shareSamples": len(shares), "shareMeanPct": round(mean, 3), "shareStdevPct": round(sd, 3), }
def read_first_seen(path): """did -> {firstSeen, source}. The date we FIRST identified an account as Dutch.""" out = {} if not path or not os.path.exists(path): return out with open(path) as f: for line in f: line = line.strip() if not line or line.startswith("#"): continue try: row = json.loads(line) except ValueError: continue if row.get("did") and row.get("firstSeen"): out.setdefault(row["did"], row) return out
def record_first_seen(path, dids_by_source, date): """Append a line for every DID we had not identified before. Returns per-source counts.
This is what makes the history interpretable later. A rising floor is two different claims tangled together -- more Dutch accounts exist, or we got better at spotting them -- and a cumulative count cannot separate them. Knowing WHEN each account first appeared, and WHICH detector found it first, lets a chart show arrivals per day rather than a total that only ever climbs.
The distinction it buys, concretely: a DID first seen by the PLC handle scan is almost always catch-up (the account existed for years; we only just looked), while a DID first seen by the firehose is closer to a real arrival (it posted Dutch today). Same +1 on the floor, opposite meanings.
Append-only, one line per DID ever. A DID is never re-dated. """ known = read_first_seen(path) new_by_source = {src: 0 for src in SOURCES} fresh = [] # When several detectors see a DID in the same run the credit is a tie, so the # order is fixed deliberately rather than left to whatever SOURCES happens to be. # The firehose comes first because it is the only real-time signal: it means the # account posted Dutch during the window. The rest are standing-state scans that # would have found the account whenever we happened to look. for src in FIRST_SEEN_PRECEDENCE: for did in sorted(dids_by_source.get(src) or ()): if did in known: continue known[did] = {"did": did, "firstSeen": date, "source": src} fresh.append(known[did]) new_by_source[src] += 1 if fresh: with open(path, "a") as f: for row in fresh: f.write(json.dumps(row, sort_keys=True) + "\n") return new_by_source, len(known)
def append_history(path, stats, coverage): """Append one counts-only line per aggregator run to a JSONL history.
Append-only and idempotent per run: re-running on the same day adds another line rather than rewriting, because a same-day re-run is a real second observation.
Deliberately carries `coverage`: without it a jump in the floor is uninterpretable -- did the community grow, did a scan finish, or did we add starter packs? Consumers must treat missing dates as gaps and never interpolate across them; the collector host can be down for a day and that is not a dip to zero. """ ident = stats["identified"] row = { "generatedAt": stats["generatedAt"], "date": stats["generatedAt"][:10], "identifiedTotal": ident["total"], "bySource": ident["bySource"], "exclusiveBySource": ident.get("exclusiveBySource", {}), "summedSources": ident.get("summedSources"), "estimate": stats["estimate"], "coverage": coverage, # Arrivals, not the running total: how many DIDs were identified for the first # time in this run, and by which detector. A total that only grows says nothing # about whether the community grew. "newBySource": (stats.get("_newBySource") or {}), "newTotal": sum((stats.get("_newBySource") or {}).values()), } with open(path, "a") as f: f.write(json.dumps(row, sort_keys=True) + "\n") return row
def build_coverage(paths, collector, sources_meta): """What was and wasn't fully measured this run, so a later jump stays readable.""" secs = (collector or {}).get("window_seconds") or 0 return { "collectorWindowMinutes": round(secs / 60), "collectorComplete": round(secs / 60) >= MIN_LIVE_WINDOW_MIN, "starterpackPacks": (sources_meta.get("starterpack") or {}).get("packs", 0), "feedersPresent": sorted(src for src in SOURCES if paths.get(src) and os.path.exists(paths[src])), "scriptVersion": SCRIPT_VERSION, }
def aggregate(paths, collector_path, now, share_samples=None, abandoned=None, attributed_path=None): """paths: {source -> dids file path}. Returns the nl-stats.json dict.
`abandoned` is the set of DIDs whose repo we have given up on: either the PDS says there is no repo, or several separate runs got no answer at all. They are removed from every source and from the union, on the same principle as the lapsed-handle prune: the floor is what we can point at TODAY, and we cannot point at an account that is not there. The DIDs stay on record in nl-lexicons.jsonl so the collectors know to stop asking. """ abandoned = abandoned or set() dids_by_source = {src: read_dids(paths.get(src)) - abandoned for src in SOURCES} # Fall back to the retired NL-only scan if the multi-country one has not run yet. if not dids_by_source.get("nlHandle") and paths.get("nlHandleLegacy"): dids_by_source["nlHandle"] = read_dids(paths["nlHandleLegacy"]) - abandoned
# The floor comes from the shared European attribution pass when it is available # (attributed.nl.dids, written by eu-stats-aggregate.py), so this page and the NL row on # atproto.eu are the same number by construction. Two aggregators computing "accounts in # the Netherlands" independently is how you end up publishing two different figures for # one country, which is exactly what happened. # # The difference is real, not cosmetic: attribution gives each account to ONE country, so # someone with a .nl handle and a declared location in Germany counts as German. A local # union would claim them here too. attributed = read_dids(attributed_path) - abandoned if attributed_path else set() if attributed: union = attributed else: union = set() for src in CONFIRMING_SOURCES: union |= dids_by_source[src]
# Per-source counts are of accounts INSIDE the floor, so a source cannot claim someone the # attribution gave to another country. Still non-exclusive: one account can have several. by_source = {src: len(dids_by_source[src] & union) for src in CONFIRMING_SOURCES} # Language signals: reported, never unioned. `alsoConfirmed` is the honest interesting # number -- how many Dutch posters some other source independently places in the NL. signals = { src: { "count": len(dids_by_source[src]), "alsoConfirmed": len(dids_by_source[src] & union), "signalOnly": len(dids_by_source[src] - union), } for src in SIGNAL_SOURCES } collector = load_json(collector_path) estimate = build_estimate(collector, fetch_bluesky_total()) bounds = derive_bounds(estimate, share_samples or [], len(union)) if bounds: estimate = dict(estimate, **bounds)
sources_meta = {} for src in SOURCES: meta = {"count": len(dids_by_source[src]), "asOf": iso_mtime(paths.get(src))} if src == "dutchPoster": secs = (collector or {}).get("window_seconds") if secs: meta["window"] = f"~{round(secs / 3600)}h firehose sample" if src == "starterpack": sp_path = paths.get(src) or "" companion = os.path.splitext(sp_path)[0] + ".json" if sp_path else None packs = (load_json(companion) or {}).get("packs") meta["packs"] = packs or 0 sources_meta[src] = meta
exclusive = exclusive_counts({k: dids_by_source[k] & union for k in CONFIRMING_SOURCES}) return { "generatedAt": now, "excludedAbandoned": len(abandoned), "identified": { "total": len(union), "bySource": by_source, # Per source: found ONLY by that source. summed - total is the overlap. "exclusiveBySource": exclusive, "summedSources": sum(by_source.values()), }, "attributionSource": "shared" if attributed else "local-union", "signals": signals, "estimate": estimate, "sources": sources_meta, }
def assert_union_invariants(stats): """The same guard src/data/nl-stats.ts runs at build time.""" counts = list(stats["identified"]["bySource"].values()) # confirming sources only total = stats["identified"]["total"] assert total >= max([0, *counts]), f"total {total} < largest source {max([0, *counts])}" assert total <= sum(counts), f"total {total} > summed sources {sum(counts)}"
def now_iso(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def self_test(): import tempfile
d = tempfile.mkdtemp()
def w(name, dids): p = os.path.join(d, name) with open(p, "w") as f: f.write("\n".join(dids) + "\n") return p
# a,b,c posted Dutch; b,c,e have .nl handles. The floor is the CONFIRMING union only, # so it is {b,c,e} = 3 -- 'a' posted Dutch and nothing else, which does not place it in # the Netherlands. Of the 3 Dutch posters, 2 are independently confirmed. paths = { "dutchPoster": w("nl-live.json.dids", ["did:plc:a", "did:plc:b", "did:plc:c"]), "nlHandle": w("nl-handles.dids", ["did:plc:b", "did:plc:c", "did:plc:e"]), "starterpack": None, "sifaNl": None, } est = os.path.join(d, "nl-live.json") with open(est, "w") as f: json.dump({"nl_share_pct": 0.941, "bluesky_total_users": None, "window_seconds": 21600}, f)
stats = aggregate(paths, est, "2026-08-27T00:00:00Z") assert stats["identified"]["total"] == 3, stats["identified"] assert stats["identified"]["bySource"] == { "nlHandle": 3, "starterpack": 0, "sifaNl": 0, "sifaCompany": 0, }, stats["identified"]["bySource"] # Language is reported, never unioned into the floor. assert "dutchPoster" not in stats["identified"]["bySource"], stats["identified"]["bySource"] assert stats["signals"]["dutchPoster"] == { "count": 3, "alsoConfirmed": 2, # b, c "signalOnly": 1, # a -- posted Dutch, nothing places it in the NL }, stats["signals"] assert_union_invariants(stats) # 3 >= 3 and 3 <= 3 assert stats["attributionSource"] == "local-union", stats["attributionSource"]
# With a shared attribution set, the floor IS that set -- not a locally recomputed union. # 'e' has an .nl handle but the European pass gave it to another country, so it drops out # here too, and nlHandle reports 2 rather than 3. attr = w("attributed.nl.dids", ["did:plc:b", "did:plc:c"]) shared = aggregate(paths, est, "2026-08-27T00:00:00Z", attributed_path=attr) assert shared["identified"]["total"] == 2, shared["identified"] assert shared["identified"]["bySource"]["nlHandle"] == 2, shared["identified"]["bySource"] assert shared["attributionSource"] == "shared", shared["attributionSource"] assert_union_invariants(shared) # estimate: a 6h window is too short -> anchored model (coherent with community-size) assert stats["estimate"] == ANCHOR_ESTIMATE, stats["estimate"] assert stats["sources"]["dutchPoster"]["window"] == "~6h firehose sample" assert stats["sources"]["dutchPoster"]["count"] == 3
# A representative window (>= 24h) with a plausible share switches to the live estimate. live = os.path.join(d, "live.json") with open(live, "w") as f: json.dump({"nl_share_pct": 0.8, "bluesky_total_users": 46_000_000, "window_seconds": 172800}, f) est_live = build_estimate(load_json(live)) assert est_live["windowMinutes"] == 2880, est_live assert est_live["modelledNl"] == round(46_000_000 * 0.8 / 100 * 0.65 / 10_000) * 10_000, est_live
# Master accumulation: two windows with partial overlap grow the master, never shrink. master = os.path.join(d, "dutch-posters.dids") m1 = accumulate_master(w("win1.dids", ["did:plc:a", "did:plc:b"]), master) assert m1 == {"did:plc:a", "did:plc:b"}, m1 m2 = accumulate_master(w("win2.dids", ["did:plc:b", "did:plc:c"]), master) # overlap on b assert m2 == {"did:plc:a", "did:plc:b", "did:plc:c"}, m2 # a fresh window that lost 'a' must NOT shrink the master (never forget) m3 = accumulate_master(w("win3.dids", ["did:plc:c"]), master) assert m3 == {"did:plc:a", "did:plc:b", "did:plc:c"}, m3 assert read_dids(master) == m3
# Abandoned DIDs drop out of every source and out of the union. ab = aggregate(paths, est, "2026-08-27T00:00:00Z", abandoned={"did:plc:a", "did:plc:e"}) assert ab["identified"]["total"] == 2, ab["identified"] # floor was 3: e removed assert ab["identified"]["bySource"]["nlHandle"] == 2, ab["identified"]["bySource"] # Abandoned DIDs drop out of the signal counts too. assert ab["signals"]["dutchPoster"]["count"] == 2, ab["signals"] assert ab["excludedAbandoned"] == 2, ab assert_union_invariants(ab)
# Exclusive counts: b,c are seen by both feeders, so each source's "only here" count # excludes them. a is poster-only, e is handle-only. ex = exclusive_counts({ "dutchPoster": {"did:plc:a", "did:plc:b", "did:plc:c"}, "nlHandle": {"did:plc:b", "did:plc:c", "did:plc:e"}, }) assert ex == {"dutchPoster": 1, "nlHandle": 1}, ex
# Share samples: one row per date, a re-run of the same date replaces it, and a # window shorter than a day is not a sample at all. sp = os.path.join(d, "shares.jsonl") assert record_share_sample(sp, {"nl_share_pct": 0.8, "window_seconds": 3600}) is None record_share_sample(sp, {"nl_share_pct": 0.8, "window_seconds": 90000}, date="2026-08-01") record_share_sample(sp, {"nl_share_pct": 0.9, "window_seconds": 90000}, date="2026-08-01") rows = read_share_samples(sp) assert len(rows) == 1 and rows[0]["sharePct"] == 0.9, rows with open(sp, "a") as f: f.write("{not json\n") # a half-written cron line must not break the build assert len(read_share_samples(sp)) == 1
# Bounds: too few samples -> no range published at all. est = {"blueskyActive": 46_000_000, "modelledNl": 230_000} assert derive_bounds(est, rows, 9_000) == {} many = [{"date": f"2026-08-{i:02d}", "sharePct": v} for i, v in enumerate([0.70, 0.75, 0.80, 0.85, 0.90, 0.78], start=2)] b = derive_bounds(est, many, 9_000) assert b["modelledLo"] < est["modelledNl"] < b["modelledHi"], b assert b["shareSamples"] == 6, b # Hard floor: the low end can never fall below the accounts we identified. b2 = derive_bounds(est, many, 300_000) assert b2 == {} or b2["modelledLo"] >= 300_000, b2
# First-seen: a DID is dated once and never re-dated, and the count reports arrivals. # Language is NOT in FIRST_SEEN_PRECEDENCE, so a Dutch-posting account gets no arrival # date from that alone -- 'a' is invisible here even though it is in the signal set. fs = os.path.join(d, "first-seen.jsonl") sets = {"dutchPoster": {"did:plc:a", "did:plc:b"}, "nlHandle": {"did:plc:b", "did:plc:c"}} new1, total1 = record_first_seen(fs, sets, "2026-08-01") assert total1 == 2 and new1["nlHandle"] == 2, (new1, total1) assert new1.get("dutchPoster", 0) == 0, new1 # b was first seen via nlHandle, so nothing must re-date it later. new2, total2 = record_first_seen(fs, sets, "2026-08-02") assert total2 == 2 and sum(new2.values()) == 0, (new2, total2) sets["nlHandle"].add("did:plc:d") new3, total3 = record_first_seen(fs, sets, "2026-08-03") assert total3 == 3 and new3["nlHandle"] == 1, (new3, total3) assert read_first_seen(fs)["did:plc:b"]["firstSeen"] == "2026-08-01", "never re-date a DID" assert "did:plc:a" not in read_first_seen(fs), "language alone must not date an account"
# History: append-only, one line per run, gaps stay gaps (no interpolation here). hp = os.path.join(d, "history.jsonl") append_history(hp, stats, {"scriptVersion": SCRIPT_VERSION}) append_history(hp, stats, {"scriptVersion": SCRIPT_VERSION}) with open(hp) as f: lines = [json.loads(x) for x in f if x.strip()] assert len(lines) == 2 and lines[0]["identifiedTotal"] == 3, lines assert lines[0]["date"] == "2026-08-27", lines[0] print("self-test OK")
def main(): ap = argparse.ArgumentParser(description="Union NL-stats feeders into nl-stats.json") ap.add_argument("--dir", default=".", help="directory holding the .dids + collector json") ap.add_argument("--out", default="nl-stats.json") ap.add_argument("--estimate", default=None, help="collector full-result json (default nl-live.json)") for src in SOURCES: ap.add_argument(f"--{src}", dest=src, default=None, help=f"{src} .dids path") ap.add_argument("--master", default=None, help="persistent Dutch-poster master (default dutch-posters.dids)") ap.add_argument("--abandoned", default=None, help="DIDs to exclude (default nl-abandoned.dids)") ap.add_argument("--attributed", default=None, help="attributed.nl.dids from eu-stats-aggregate.py; makes this page and " "the NL row on atproto.eu the same number (default: auto-detect in --dir)") ap.add_argument("--first-seen", default=None, help="append-only per-DID first-seen log") ap.add_argument("--history", default=None, help="append-only run log (default nl-stats-history.jsonl)") ap.add_argument("--share-samples", default=None, help="daily share samples (default nl-share-samples.jsonl)") ap.add_argument("--no-history", action="store_true", help="don't append to the history log") ap.add_argument("--no-master", action="store_true", help="count only the current window, don't accumulate") ap.add_argument("--self-test", action="store_true") args = ap.parse_args()
if args.self_test: self_test() return
paths = {src: getattr(args, src) or os.path.join(args.dir, DEFAULT_FILE[src]) for src in SOURCES} estimate_path = args.estimate or os.path.join(args.dir, "nl-live.json")
# Accumulate the windowed Dutch-poster feed into a persistent master so the floor # only grows across collector restarts, then count the master (not just the window). if not args.no_master: master_path = args.master or os.path.join(args.dir, "dutch-posters.dids") accumulate_master(paths["dutchPoster"], master_path) paths["dutchPoster"] = master_path
abandoned_path = args.abandoned or os.path.join(args.dir, "nl-abandoned.dids") abandoned = read_dids(abandoned_path) samples_path = args.share_samples or os.path.join(args.dir, "nl-share-samples.jsonl") collector = load_json(estimate_path) # A completed collector day is one observation of how much Dutch gets posted. The # spread across days is where the estimate's error bounds come from, so record it # before aggregating -- today's day counts toward today's bounds. sample = record_share_sample(samples_path, collector) samples = read_share_samples(samples_path)
attributed_path = args.attributed if attributed_path is None: auto = os.path.join(args.dir, "attributed.nl.dids") attributed_path = auto if os.path.exists(auto) else None stats = aggregate(paths, estimate_path, now_iso(), share_samples=samples, abandoned=abandoned, attributed_path=attributed_path) assert_union_invariants(stats)
first_seen_path = args.first_seen or os.path.join(args.dir, "nl-first-seen.jsonl") dids_by_source = {src: read_dids(paths.get(src)) - abandoned for src in SOURCES} # Fall back to the retired NL-only scan if the multi-country one has not run yet. if not dids_by_source.get("nlHandle") and paths.get("nlHandleLegacy"): dids_by_source["nlHandle"] = read_dids(paths["nlHandleLegacy"]) - abandoned new_by_source, known_total = record_first_seen( first_seen_path, dids_by_source, stats["generatedAt"][:10] ) stats["_newBySource"] = new_by_source
if not args.no_history: history_path = args.history or os.path.join(args.dir, "nl-stats-history.jsonl") append_history(history_path, stats, build_coverage(paths, collector, stats["sources"]))
published = {k: v for k, v in stats.items() if not k.startswith("_")} with open(args.out, "w") as f: f.write(json.dumps(published, indent=2) + "\n") ident = stats["identified"] print( f"Wrote {args.out}: {ident['total']} identified " + "(" + ", ".join(f"{k} {v}" for k, v in ident["bySource"].items()) + "); " + "signals: " + ", ".join( f"{k} {v['count']} ({v['alsoConfirmed']} also confirmed)" for k, v in stats.get("signals", {}).items()) + "; " f"estimate ~{stats['estimate']['modelledNl']}" + (f"; {len(abandoned)} abandoned excluded" if abandoned else "") + (f" [{stats['estimate']['modelledLo']}-{stats['estimate']['modelledHi']}]" if "modelledLo" in stats["estimate"] else f" (no bounds: {len(samples)}/{MIN_SHARE_SAMPLES} share samples)") + f"; {sum(new_by_source.values())} first-seen (of {known_total} ever)" + (f"; recorded share sample for {sample['date']}" if sample else ""), file=sys.stderr, )
if __name__ == "__main__": main()