Files
SocialPub/tools/pasture/town/seed.py
T
thepraandClaude Opus 5.5 4b6d261739 Town: decePub's cells in the report, a Hollo pair, kinder waits
The report merges decePub's end-to-end cells (out/client.jsonl, the
latest result of each test) into the matrix under the observer
"decepub", and names its run without the workstation's paths. A second
spec, hollo-pair, seeds PrivaPub and Hollo alone (258 of 261 checks pass;
the three left found the emoji variation selector bug fixed in
PrivaPub), and Hollo accounts may join circles. The seeder waits up to
90 s for a follow request, two of a sender's retries after a refused
first delivery (Hollo answers its first concurrent deliveries 500 while
it creates its tables). town.sh drive prints its JSON unescaped.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
2026-10-05 01:24:17 +02:00

406 lines
16 KiB
Python

"""Replays a plan against the pasture and writes what happened to out/<run>/ledger.jsonl, one fact per line:
accounts made, groups made, follows and their state, objects with their URIs, interactions, mutations, errors.
A reply waits for its parent to arrive on the replier's server (the round before had time to settle); if it never
arrives the parent is fetched and the ledger says `arrived_by: fetch` for that server, so a delivery that failed shows
up as a cell in the report rather than as a broken thread. An error never stops the seeding: it is a ledger line."""
import json
import os
import sys
import threading
import time
import traceback
from concurrent.futures import ThreadPoolExecutor
import gen
import state
from core import podman
from dialects import driver
from dialects.base import Account, PostSpec, Unsupported
OUT = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "out")
class Ledger:
def __init__(self, run_dir):
self.path = os.path.join(run_dir, "ledger.jsonl")
self.lock = threading.Lock()
self.facts = []
if os.path.exists(self.path):
with open(self.path) as f:
self.facts = [json.loads(line) for line in f if line.strip()]
def add(self, **fact):
fact["t"] = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
with self.lock:
self.facts.append(fact)
with open(self.path, "a") as f:
f.write(json.dumps(fact, ensure_ascii=False, default=str) + "\n")
def done(self):
return {f["n"] for f in self.facts if "n" in f and f.get("type") != "error"}
class Seeder:
def __init__(self, spec, run_dir, only=None):
self.spec = spec
snap_path = os.path.join(run_dir, "plan.json")
if os.path.exists(snap_path):
with open(snap_path) as f:
self.steps = json.load(f)["steps"]
else:
planner = gen.Planner(spec)
planner.plan()
self.steps = planner.steps
with open(snap_path, "w") as f:
json.dump(gen.snapshot(planner), f, ensure_ascii=False)
self.ledger = Ledger(run_dir)
self.only = only
self.uris = {} # ref -> uri
self.origin = {} # ref -> platform
self.groups = {} # ref -> {kind, host, id, code, handle, ...}
self.timing = spec.get("time", {})
for f in self.ledger.facts:
if f.get("type") == "object":
self.uris[f["ref"]] = f["uri"]
self.origin[f["ref"]] = f["origin"]
elif f.get("type") == "group":
self.groups[f["ref"]] = f["group"]
def acct(self, key):
return gen.Planner.acct(key)
def session(self, key):
platform, user = key.split("/", 1)
return state.session(platform, user)
def run(self):
done = self.ledger.done()
phases = ["accounts", "groups", "graph", "content", "interact", "mutate"]
for phase in phases:
steps = [s for s in self.steps if s["phase"] == phase and s["n"] not in done]
if self.only and phase not in self.only:
continue
if not steps:
continue
started = time.time()
print(f"{phase}: {len(steps)} steps", flush=True)
if phase == "accounts":
self.accounts(steps)
elif phase == "graph" and any(k.startswith("lemmy/") for k in self.accounts_planned()):
self.warm_lemmy()
self.parallel(steps, ordered=True)
elif phase == "content":
# follows, joins and their Accepts travel through queues: posts made before they land reach nobody new
time.sleep(self.timing.get("graphSettleSeconds", 30))
rounds = sorted({s["round"] for s in steps})
for r in rounds:
if r > 0:
time.sleep(self.timing.get("roundSettleSeconds", 20))
self.parallel([s for s in steps if s["round"] == r])
else:
if phase in ("interact", "mutate"):
time.sleep(self.timing.get("roundSettleSeconds", 20))
self.parallel(steps, ordered=phase in ("graph", "groups"))
print(f"{phase}: done in {time.time() - started:.0f}s", flush=True)
return 0
def accounts_planned(self):
return [s["actor"] for s in self.steps if s["verb"] == "account"]
def warm_lemmy(self):
"""One account on every other server, looked up from Lemmy, so its send workers exist before the first Follow."""
keys = self.accounts_planned()
first = {}
for k in keys:
first.setdefault(k.split("/")[0], k)
lemmy = next(k for k in keys if k.startswith("lemmy/"))
try:
ready = driver("lemmy").warm(self.session(lemmy), [self.acct(k) for p, k in first.items() if p != "lemmy"])
self.ledger.add(type="warmup", platform="lemmy", ready=ready)
except Exception as e:
self.ledger.add(type="error", verb="warmup", actor=lemmy, error=repr(e)[:400])
def parallel(self, steps, ordered=False):
"""Steps of one actor run in plan order; different actors run side by side (four per platform at most)."""
if ordered:
# the graph's accept and approve steps need the follow before them, whoever made it
for s in steps:
self.do(s)
return
by_actor = {}
for s in steps:
by_actor.setdefault(s["actor"], []).append(s)
with ThreadPoolExecutor(max_workers=12) as pool:
for actor_steps in by_actor.values():
pool.submit(lambda ss: [self.do(s) for s in ss], actor_steps)
def do(self, s):
try:
getattr(self, f"v_{s['verb']}")(s)
except Exception as e:
kind = "unsupported" if isinstance(e, Unsupported) else "error"
self.ledger.add(type=kind, n=s["n"], verb=s["verb"], actor=s["actor"], ref=s.get("ref"), error=repr(e)[:600])
if kind == "error" and os.environ.get("TOWN_DEBUG"):
traceback.print_exc()
# -- accounts
def accounts(self, steps):
by_platform = {}
for s in steps:
by_platform.setdefault(s["args"]["platform"], []).append(s)
for platform, ss in by_platform.items():
accounts = [Account(**{k: v for k, v in s["args"].items()}) for s in ss]
for a in accounts:
a.fields = [tuple(f) for f in a.fields]
try:
sessions = state.provision(platform, accounts)
except Exception as e:
for s in ss:
self.ledger.add(type="error", n=s["n"], verb="account", actor=s["actor"], error=repr(e)[:600])
traceback.print_exc()
continue
d = driver(platform)
def profile(pair):
s, sess = pair
try:
d.update_profile(sess, sess.account)
self.ledger.add(type="account", n=s["n"], key=s["actor"], actor_uri=sess.actor_uri, local_id=sess.local_id,
profile=True)
except Unsupported:
self.ledger.add(type="account", n=s["n"], key=s["actor"], actor_uri=sess.actor_uri, local_id=sess.local_id,
profile=False)
except Exception as e:
# the account exists; only its profile is missing, which the checker can tell from profile=False
self.ledger.add(type="account", n=s["n"], key=s["actor"], actor_uri=sess.actor_uri, local_id=sess.local_id,
profile=False, profile_error=repr(e)[:300])
with ThreadPoolExecutor(max_workers=6) as pool:
list(pool.map(profile, zip(ss, sessions)))
# -- groups
def v_group(self, s):
platform = s["actor"].split("/")[0]
sess = self.session(s["actor"])
kind, name = s["args"]["kind"], s["args"]["name"]
if platform == "privapub":
g = driver("privapub").group(sess, name, name.replace("_", " ").title(), community=kind == "community")
group = {"kind": kind, "host": "privapub", "id": g["id"], "code": g.get("invitationCode"),
"handle": f"{name}@privapub.test", "owner": s["actor"]}
else:
c = driver("lemmy").community(sess, name, name.replace("_", " ").title())
group = {"kind": kind, "host": "lemmy", "id": c["id"], "handle": f"{name}@lemmy.test", "owner": s["actor"],
"ap_id": c["ap_id"]}
self.groups[s["ref"]] = group
self.ledger.add(type="group", n=s["n"], ref=s["ref"], group=group)
# -- the graph
def v_follow(self, s):
platform = s["actor"].split("/")[0]
outcome = driver(platform).follow(self.session(s["actor"]), self.acct(s["args"]["target"]))
self.ledger.add(type="relation", n=s["n"], verb="follow", actor=s["actor"], target=s["args"]["target"], state=outcome)
def _wait_pending(self, s):
platform = s["actor"].split("/")[0]
d, sess, follower = driver(platform), self.session(s["actor"]), self.acct(s["args"]["target"])
# a refused first delivery is retried on the sender's backoff (16 s, then 31 s): wait out two of them
for _ in range(45):
if follower.lower() in (p.lower() for p in d.pending(sess)):
return d, sess, follower
time.sleep(2)
raise TimeoutError(f"{follower}'s follow request never reached {s['actor']}")
def v_accept(self, s):
d, sess, follower = self._wait_pending(s)
d.accept(sess, follower)
self.ledger.add(type="relation", n=s["n"], verb="accept", actor=s["actor"], target=s["args"]["target"], state="accepted")
def v_reject(self, s):
d, sess, follower = self._wait_pending(s)
d.reject(sess, follower)
self.ledger.add(type="relation", n=s["n"], verb="reject", actor=s["actor"], target=s["args"]["target"], state="rejected")
def v_join(self, s):
group = self.groups[s["args"]["group"]]
platform = s["actor"].split("/")[0]
sess = self.session(s["actor"])
how = None
if platform == "privapub" and group["host"] == "privapub":
pp = driver("privapub")
pp.http.post(pp.base + "/clientapi/group/join", ok={200, 201, 204},
headers={"Authorization": f"Bearer {pp.jwt(sess.account.root, sess.account.password)}"},
json={"avatarId": sess.local_id, "invitationCode": group["code"]})
how = "invitation"
elif platform == "lemmy":
lm = driver("lemmy")
cid = group["id"] if group["host"] == "lemmy" else lm.resolve_community(sess, "!" + group["handle"])
lm.follow_community(sess, cid)
how = "community-follow"
else:
how = driver(platform).follow(sess, group["handle"])
self.ledger.add(type="relation", n=s["n"], verb="join", actor=s["actor"], group=s["args"]["group"], state=how)
def v_approve(self, s):
group = self.groups[s["args"]["group"]]
member = s["args"]["member"]
if member.startswith("privapub/"):
self.ledger.add(type="relation", n=s["n"], verb="approve", actor=s["actor"], group=s["args"]["group"],
member=member, state="local")
return
member_uri = self.session(member).actor_uri
for _ in range(30):
found = podman.mongo(f"db.Follower.findOne({{LocalActorId: '{group['id']}', IsAccepted: false}}, {{ActorURI: 1}})")
pending = podman.mongo(f"db.Follower.find({{LocalActorId: '{group['id']}', IsAccepted: false}}, {{ActorURI: 1}}).toArray()") or []
if any(p.get("ActorURI") == member_uri for p in pending):
break
time.sleep(1)
else:
raise TimeoutError(f"{member}'s request to join {s['args']['group']} never arrived")
sess = self.session(s["actor"])
pp = driver("privapub")
pp.http.post(pp.base + "/clientapi/group/approve", ok={200, 201, 204},
headers={"Authorization": f"Bearer {pp.jwt(sess.account.root, sess.account.password)}"},
json={"avatarId": sess.local_id, "groupId": group["id"], "memberActorURI": member_uri})
self.ledger.add(type="relation", n=s["n"], verb="approve", actor=s["actor"], group=s["args"]["group"],
member=member, state="accepted")
# -- content
def arrived(self, platform, sess, ref, wait_seconds):
"""The local id of `ref` on the actor's server once it is there; None (and a fetch to come) if it never comes."""
uri = self.uris[ref]
if self.origin.get(ref) == platform:
return True
d = driver(platform)
deadline = time.time() + wait_seconds
while True:
try:
if d.local_status_id(sess, uri):
return True
except Exception:
pass
if time.time() >= deadline:
return False
time.sleep(2)
def v_post(self, s):
platform = s["actor"].split("/")[0]
sess = self.session(s["actor"])
a = dict(s["args"])
fetched = []
spec = PostSpec(text=a["text"], visibility=a["visibility"], cw=a.get("cw"), tags=a.get("tags", []),
media=a.get("media", []), poll=a.get("poll"), language=a.get("language"), title=a.get("title"),
location=a.get("location"), mentions=[self.acct(k) for k in a.get("mentions", [])])
for field, key in (("reply_to_uri", "reply_to"), ("quote_uri", "quote")):
if a.get(key):
if a[key] not in self.uris:
raise LookupError(f"{a[key]} was never made")
if not self.arrived(platform, sess, a[key], self.timing.get("arrivalSeconds", 20)):
fetched.append(a[key])
setattr(spec, field, self.uris[a[key]])
if a.get("group"):
group = self.groups[a["group"]]
if platform == "lemmy":
lm = driver("lemmy")
spec.group = group["id"] if group["host"] == "lemmy" else lm.resolve_community(sess, "!" + group["handle"])
if spec.visibility == "community":
spec.visibility = "public"
else:
spec.group = group["id"]
made = driver(platform).post(sess, spec)
self.uris[s["ref"]] = made.uri
self.origin[s["ref"]] = platform
self.ledger.add(type="object", n=s["n"], ref=s["ref"], origin=platform, author=s["actor"], uri=made.uri,
local_id=made.local_id, fetched_on={platform: fetched} if fetched else {})
def _target(self, s):
platform = s["actor"].split("/")[0]
ref = s["args"]["target"]
if ref not in self.uris:
raise LookupError(f"{ref} was never made")
sess = self.session(s["actor"])
arrived = self.arrived(platform, sess, ref, self.timing.get("arrivalSeconds", 20))
return platform, sess, self.uris[ref], arrived
def _interaction(self, s, call):
platform, sess, uri, arrived = self._target(s)
call(driver(platform), sess, uri)
self.ledger.add(type="interaction", n=s["n"], verb=s["verb"], actor=s["actor"], ref=s["args"]["target"],
arrived_by="delivery" if arrived else "fetch", **{k: v for k, v in s["args"].items() if k != "target"})
def v_like(self, s):
self._interaction(s, lambda d, sess, uri: d.like(sess, uri))
def v_downvote(self, s):
self._interaction(s, lambda d, sess, uri: d.downvote(sess, uri))
def v_boost(self, s):
self._interaction(s, lambda d, sess, uri: d.boost(sess, uri))
def v_bookmark(self, s):
self._interaction(s, lambda d, sess, uri: d.bookmark(sess, uri))
def v_react(self, s):
self._interaction(s, lambda d, sess, uri: d.react(sess, uri, s["args"]["emoji"]))
def v_vote(self, s):
self._interaction(s, lambda d, sess, uri: d.vote(sess, uri, s["args"]["choices"]))
# -- mutations
def v_edit(self, s):
platform = s["actor"].split("/")[0]
spec = PostSpec(text=s["args"]["text"], cw=s["args"].get("cw"), mentions=[self.acct(k) for k in s["args"].get("mentions", [])])
driver(platform).edit(self.session(s["actor"]), self.uris[s["args"]["target"]], spec)
self.ledger.add(type="mutation", n=s["n"], verb="edit", actor=s["actor"], ref=s["args"]["target"], text=s["args"]["text"])
def v_delete(self, s):
platform = s["actor"].split("/")[0]
driver(platform).delete(self.session(s["actor"]), self.uris[s["args"]["target"]])
self.ledger.add(type="mutation", n=s["n"], verb="delete", actor=s["actor"], ref=s["args"]["target"])
def v_block(self, s):
platform = s["actor"].split("/")[0]
driver(platform).block(self.session(s["actor"]), self.acct(s["args"]["target"]))
self.ledger.add(type="relation", n=s["n"], verb="block", actor=s["actor"], target=s["args"]["target"], state="blocking")
def v_mute(self, s):
platform = s["actor"].split("/")[0]
driver(platform).mute(self.session(s["actor"]), self.acct(s["args"]["target"]))
self.ledger.add(type="relation", n=s["n"], verb="mute", actor=s["actor"], target=s["args"]["target"], state="muting")
def v_report(self, s):
platform = s["actor"].split("/")[0]
uris = [self.uris[r] for r in s["args"]["refs"] if r in self.uris]
driver(platform).report(self.session(s["actor"]), self.acct(s["args"]["target"]), uris, s["args"]["comment"])
self.ledger.add(type="report", n=s["n"], actor=s["actor"], target=s["args"]["target"], refs=s["args"]["refs"])
def run_dir(name, new=False):
os.makedirs(OUT, exist_ok=True)
latest = os.path.join(OUT, f"{name}.latest")
if not new and os.path.exists(latest):
with open(latest) as f:
return f.read().strip()
d = os.path.join(OUT, f"{name}-{time.strftime('%Y%m%d-%H%M%S')}")
os.makedirs(d, exist_ok=True)
with open(latest, "w") as f:
f.write(d)
return d
def main(args):
spec = gen.load(args[0] if args else "village")
resume = "--resume" in args
only = None
for a in args:
if a.startswith("--only="):
only = a.split("=", 1)[1].split(",")
d = run_dir(spec["name"], new=not resume and not only)
with open(os.path.join(d, "spec.json"), "w") as f:
json.dump(spec, f, indent=1)
print(f"run: {d}")
return Seeder(spec, d, only).run()
if __name__ == "__main__":
sys.exit(main(sys.argv[1:]))