# /// script # dependencies = ["websockets"] # /// """Sample the firehose for European-language posting activity. Not a census — the protocol has no country field. We sample the public Jetstream firehose for a window and measure, per language: - posts tagged with that language (record.langs) - distinct accounts posting it (a hard floor of *active* speakers) Plus total create-post events, a handle-suffix tally, and Bluesky's current account total (bsky.jazco.dev/stats) as the extrapolation anchor. LANGUAGE IS NOT COUNTRY. These buckets feed the *estimate* layer only, weighted by speaker distribution — nl spans NL/BE, de spans DE/AT/CH, fr spans FR/BE/LU/CH, and so on. Nothing here confirms that an account belongs to a country; that comes from ccTLD handles, self-declared location, starterpacks, Sifa, or a member record. Do not promote these DIDs to "confirmed" anywhere downstream. One firehose, many buckets: adding a language costs a set, not another collector, so widening from NL to all of Europe does not multiply the running cost. English is deliberately excluded from author tracking (`--langs`) — it is spoken everywhere, so its author set would be both enormous and meaningless; its post count is still recorded. Activity peaks in local evening hours, so a SHORT sample over- or under-counts depending on when it runs. For representative figures run it for a full day or more; this script auto-reconnects across the whole window and checkpoints its progress. Post counts are kept for every language in --langs; DID sets only for --author-langs, which defaults to --langs minus a few non-EU-dominated languages (see NO_AUTHOR_TRACKING). In byLang, `distinct_authors: null` means "not tracked", never "zero". Usage: uv run scripts/estimate-nl-users.py [seconds] [out.json] [--langs nl,de,fr] [--author-langs nl,de] uv run scripts/estimate-nl-users.py 86400 nl-24h.json # a full day uv run scripts/estimate-nl-users.py 2592000 nl-live.json # 30 days (a long, capped run) uv run scripts/estimate-nl-users.py 0 nl-live.json # run forever (uncapped) On any long run it keeps updating .partial every CHECKPOINT_EVERY seconds and rewrites the full (with handle resolution) every RESOLVE_EVERY seconds, so a live tracker can read it. Every checkpoint it also writes: .dids distinct Dutch-poster DIDs (unchanged; the NL aggregator's feed) ..dids distinct DIDs per tracked language Keep the .dids files private (host only, never committed). Prefer a capped duration (e.g. 30 days) over 0/forever so it can't run unbounded. Run it under cron/systemd or in a container; see scripts/README.md. """ import asyncio, json, sys, time, urllib.request, urllib.parse JETSTREAM = ("wss://jetstream2.us-east.bsky.network/subscribe" "?wantedCollections=app.bsky.feed.post") # European languages worth bucketing. Excludes 'en' on purpose (see the module docstring). # Cheap to extend: each entry costs one set of DIDs, not another firehose connection. DEFAULT_LANGS = ( "nl,fr,de,es,it,pt,pl,sv,da,fi,no,nb,nn,is,ga,cy,mt,el,cs,sk,sl,hr,sr,bs,mk," "sq,ro,bg,hu,et,lv,lt,uk,ru,tr,ca,eu,gl,lb,fo" ) # Languages whose posts we count but whose AUTHOR SETS we deliberately do not keep. # Each of these is dominated by a non-European population — pt by Brazil, es by Latin # America, ru/tr/uk largely outside the EU — so their DID sets run to millions while the # European slice the estimate would extract is a couple of percent. Keeping them would # cost hundreds of MB resident and a multi-second .dids rewrite every checkpoint, to # sharpen a number the Eurobarometer weighting discounts to near zero anyway. # Spain and Portugal are better served by ca/eu/gl, which are near-unique to Iberia and # stay tracked. Override with --author-langs when a specific question needs one of these. NO_AUTHOR_TRACKING = {"pt", "es", "ru", "tr"} # The language the legacy .dids feed and the nl_* JSON keys refer to. The NL # aggregator reads that exact filename, so it stays put while the rest generalises. LEGACY_LANG = "nl" CHECKPOINT_EVERY = 300 # seconds between partial writes RESOLVE_EVERY = 3600 # in forever mode, resolve handles + write the full result this often def parse_args(argv): """Positional [seconds] [out.json], plus optional --langs / --author-langs. Positional order is kept for compatibility with the deployed scheduled job and the docs.""" langs = DEFAULT_LANGS author_langs = None rest = [] i = 0 while i < len(argv): a = argv[i] for flag, setter in (("--langs", "langs"), ("--author-langs", "author_langs")): if a == flag and i + 1 < len(argv): val, i = argv[i + 1], i + 2 break if a.startswith(flag + "="): val, i = a.split("=", 1)[1], i + 1 break else: rest.append(a) i += 1 continue if setter == "langs": langs = val else: author_langs = val duration = int(rest[0]) if rest else 720 out = rest[1] if len(rest) > 1 else None tracked = [x.strip().lower() for x in langs.split(",") if x.strip()] if LEGACY_LANG not in tracked: tracked.insert(0, LEGACY_LANG) # the legacy feed must never silently vanish if author_langs is None: authored = [l for l in tracked if l not in NO_AUTHOR_TRACKING] else: authored = [x.strip().lower() for x in author_langs.split(",") if x.strip()] if LEGACY_LANG not in authored: authored.insert(0, LEGACY_LANG) return duration, out, tracked, authored DURATION, OUT, LANGS, AUTHOR_LANGS = parse_args(sys.argv[1:]) FOREVER = DURATION <= 0 # 0 (or negative) = run indefinitely # urllib announces itself as "Python-urllib/3.x", which bsky-search.jazco.io rejects with # a 403. That silently nulled bluesky_total_users -- the extrapolation anchor -- from at # least 2026-08-30 until it was caught on 2026-09-02. Always send a real User-Agent. USER_AGENT = "atproto.nl-stats/1.0 (+https://atproto.nl)" def fetch_json(url, timeout=20): req = urllib.request.Request(url, headers={"User-Agent": USER_AGENT}) with urllib.request.urlopen(req, timeout=timeout) as r: return json.load(r) def bluesky_total(): """Current Bluesky account total (extrapolation anchor).""" try: return fetch_json("https://bsky-search.jazco.io/stats").get("total_users") except Exception as e: print(f" bluesky total fetch failed: {e}", file=sys.stderr) return None def resolve_handles(dids): """Resolve DIDs to handles in batches via the public AppView.""" handles = {} base = "https://public.api.bsky.app/xrpc/app.bsky.actor.getProfiles" dids = list(dids) for i in range(0, len(dids), 25): batch = dids[i:i + 25] qs = "&".join("actors=" + urllib.parse.quote(d) for d in batch) try: data = fetch_json(f"{base}?{qs}") for p in data.get("profiles", []): handles[p["did"]] = p.get("handle", "") except Exception as e: print(f" resolve batch {i} failed: {e}", file=sys.stderr) return handles async def sample(on_checkpoint=None): import websockets authors = {lang: set() for lang in AUTHOR_LANGS} # lang -> distinct DIDs (bounded set) posts = {lang: 0 for lang in LANGS} # lang -> post count (all tracked) total_posts = 0 cursor = None start = time.monotonic() last_ckpt = start print(f"Sampling Jetstream for {'forever' if FOREVER else str(DURATION) + 's'}, " f"{len(LANGS)} languages...", file=sys.stderr) # Reconnect across the whole window; resume with a cursor so gaps stay small. while FOREVER or time.monotonic() - start < DURATION: url = JETSTREAM + (f"&cursor={cursor}" if cursor else "") try: async with websockets.connect(url, max_size=None) as ws: while FOREVER or time.monotonic() - start < DURATION: try: msg = await asyncio.wait_for(ws.recv(), timeout=10) except asyncio.TimeoutError: continue ev = json.loads(msg) if ev.get("time_us"): cursor = ev["time_us"] if ev.get("kind") != "commit": continue c = ev.get("commit") or {} if c.get("operation") != "create": continue rec = c.get("record") or {} if rec.get("$type") != "app.bsky.feed.post": continue total_posts += 1 # A post may declare several languages; it counts for each tracked one. for lang in (rec.get("langs") or []): # Normalise regional tags: "nl-BE" and "pt-BR" bucket as nl / pt. base_lang = str(lang).split("-")[0].lower() if base_lang in posts: posts[base_lang] += 1 if base_lang in authors: authors[base_lang].add(ev.get("did")) now = time.monotonic() if on_checkpoint and now - last_ckpt >= CHECKPOINT_EVERY: last_ckpt = now on_checkpoint(total_posts, posts, authors, now - start) except Exception as ex: print(f" reconnecting after: {ex}", file=sys.stderr) await asyncio.sleep(2) return total_posts, posts, authors, time.monotonic() - start def write_json(obj, path): out = json.dumps(obj, indent=2) if path: with open(path, "w") as f: f.write(out + "\n") return out def write_dids(authors, path): """Dump distinct DIDs (one per line) for the stats aggregator. These are WINDOWED sets: DIDs seen posting a language during this run. The aggregator unions them across runs into a master file, so the "ever seen" floor grows and never forgets. The windowed count stays separate in the JSON, so "distinct this window" is never conflated with "ever seen". Keep the .dids files private (host only, never committed); only the aggregated counts go public. """ if not path: return with open(path, "w") as f: for did in sorted(authors): if did: f.write(did + "\n") def write_all_dids(authors, out): """.dids stays the Dutch feed the NL aggregator already reads; every tracked language also gets its own ..dids for the per-country pipelines.""" if not out: return write_dids(authors.get(LEGACY_LANG, set()), out + ".dids") for lang, dids in authors.items(): write_dids(dids, f"{out}.{lang}.dids") def suffix_tally(handles): """Tally handle suffixes: every ccTLD/gTLD seen, plus bsky.social called out. Country-code suffixes are a genuine confirmation signal (unlike language), so this is worth keeping wide rather than collapsing everything into "other". """ tally = {"nl": 0, "be": 0, "bsky.social": 0} # keys kept for output compatibility for h in handles.values(): if h.endswith(".bsky.social"): tally["bsky.social"] += 1 continue tld = h.rsplit(".", 1)[-1].lower() if "." in h else "" tally[tld] = tally.get(tld, 0) + 1 return dict(sorted(tally.items(), key=lambda kv: -kv[1])) def build_result(elapsed, total_posts, posts, authors): """Full result: resolve handles for the legacy language + per-language counts. Handle resolution stays limited to the legacy set: resolving every European-language author would mean tens of thousands of AppView calls per checkpoint. The per-language DIDs are written to disk regardless, so any downstream job can resolve at its own pace. """ legacy_authors = authors.get(LEGACY_LANG, set()) handles = resolve_handles(legacy_authors) suffix = suffix_tally(handles) nl_posts = posts.get(LEGACY_LANG, 0) return { "window_seconds": round(elapsed), "bluesky_total_users": bluesky_total(), "total_posts": total_posts, "nl_posts": nl_posts, "nl_share_pct": round(nl_posts / max(total_posts, 1) * 100, 3), "distinct_nl_authors": len(legacy_authors), "resolved": len(handles), "handle_suffix": suffix, "nl_handle_share_pct": round(suffix.get("nl", 0) / max(len(handles), 1) * 100, 1), # distinct_authors is null where author tracking is off (see NO_AUTHOR_TRACKING): # null means "not counted", never "zero speakers". "byLang": { lang: { "posts": posts[lang], "distinct_authors": len(authors[lang]) if lang in authors else None, "share_pct": round(posts[lang] / max(total_posts, 1) * 100, 3), } for lang in sorted(posts, key=lambda l: -posts[l]) if posts[lang] }, } def main(): last_full = [0.0] def checkpoint(total, posts, authors, elapsed): if not OUT: return nl_posts = posts.get(LEGACY_LANG, 0) write_json({ "partial": True, "window_seconds": round(elapsed), "total_posts": total, "nl_posts": nl_posts, "nl_share_pct": round(nl_posts / max(total, 1) * 100, 3), "distinct_nl_authors": len(authors.get(LEGACY_LANG, set())), "byLang": { lang: {"posts": posts[lang], "distinct_authors": len(authors[lang]) if lang in authors else None} for lang in sorted(posts, key=lambda l: -posts[l]) if posts[lang] }, }, OUT + ".partial") # Dump the DID sets every checkpoint so the aggregators always have a fresh feed. write_all_dids(authors, OUT) live = sum(1 for lang in posts if posts[lang]) print(f" checkpoint: {round(elapsed)}s, {nl_posts}/{total} nl posts, " f"{len(authors.get(LEGACY_LANG, set()))} distinct nl, {live} langs seen", file=sys.stderr) # On long runs the end-of-run write is far off (or never, if forever), so periodically # resolve handles and rewrite the full result (this briefly blocks the firehose read). if elapsed - last_full[0] >= RESOLVE_EVERY: last_full[0] = elapsed write_json(build_result(elapsed, total, dict(posts), {k: set(v) for k, v in authors.items()}), OUT) print(f" full result written ({len(authors.get(LEGACY_LANG, set()))} nl authors)", file=sys.stderr) total_posts, posts, authors, elapsed = asyncio.run(sample(checkpoint)) # Reached only for a finite run (forever mode loops until killed). print(f"Resolving {len(authors.get(LEGACY_LANG, set()))} distinct Dutch-posting accounts...", file=sys.stderr) write_all_dids(authors, OUT) print(write_json(build_result(elapsed, total_posts, posts, authors), OUT)) if __name__ == "__main__": main()