PeerTube 8.3.1 runs on the shared Postgres and Redis, transcoding off. Its scenario passes 24 checks: a persona follows a channel; a new video comes as the channel's Announce and reaches the persona's home as a boost, playable through the media proxy with byte ranges; comments both ways thread; a like and its undo count; renaming and deletion arrive; the unfollow; statistics. The town's seeder also takes reruns on the same accounts in its stride: a Lemmy community already made is found, a circle member already approved asks nothing again. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
424 lines
17 KiB
Python
424 lines
17 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))
|
|
# a post quoting or answering one of its own round goes once that one is made
|
|
these = [s for s in steps if s["round"] == r]
|
|
made_here = {s["ref"] for s in these}
|
|
later = [s for s in these if {s["args"].get("quote"), s["args"].get("reply_to")} & made_here]
|
|
self.parallel([s for s in these if s not in later])
|
|
if later:
|
|
self.parallel(later)
|
|
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):
|
|
# a follow that came back accepted needs no answer: the target was unlocked after all, or a run before this one
|
|
# on the same accounts already accepted it
|
|
if any(f.get("type") == "relation" and f["verb"] == "follow" and f["actor"] == s["args"]["target"]
|
|
and f["target"] == s["actor"] and f["state"] == "accepted" for f in self.ledger.facts):
|
|
self.ledger.add(type="relation", n=s["n"], verb="accept", actor=s["actor"], target=s["args"]["target"],
|
|
state="accepted", already=True)
|
|
return
|
|
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
|
|
# a member an earlier run on these accounts approved asks nothing again
|
|
if podman.mongo(f"db.Follower.findOne({{LocalActorId: '{group['id']}', ActorURI: '{member_uri}', IsAccepted: true}}, {{_id: 1}})"):
|
|
self.ledger.add(type="relation", n=s["n"], verb="approve", actor=s["actor"], group=s["args"]["group"],
|
|
member=member, state="accepted", already=True)
|
|
return
|
|
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:]))
|