diff options
| author | Christian Cleberg <[email protected]> | 2026-04-11 22:48:38 -0500 |
|---|---|---|
| committer | Christian Cleberg <[email protected]> | 2026-04-11 22:48:38 -0500 |
| commit | e27bbe17fd5d8f06fc3ccb4a8de97a0775a1439c (patch) | |
| tree | 7ce854d65009decbe8631175f233f3b2c9df6373 /src/srht_contrib/scripts | |
| parent | 533866679755bd6e7a97cfa0f050eaa832b0b373 (diff) | |
| download | hutch-stats-e27bbe17fd5d8f06fc3ccb4a8de97a0775a1439c.tar.gz hutch-stats-e27bbe17fd5d8f06fc3ccb4a8de97a0775a1439c.tar.bz2 hutch-stats-e27bbe17fd5d8f06fc3ccb4a8de97a0775a1439c.zip | |
add durable actor queueing and staggered discovery imports
Diffstat (limited to 'src/srht_contrib/scripts')
| -rw-r--r-- | src/srht_contrib/scripts/__init__.py | 1 | ||||
| -rw-r--r-- | src/srht_contrib/scripts/enqueue_actors.py | 86 |
2 files changed, 87 insertions, 0 deletions
diff --git a/src/srht_contrib/scripts/__init__.py b/src/srht_contrib/scripts/__init__.py new file mode 100644 index 0000000..8b13789 --- /dev/null +++ b/src/srht_contrib/scripts/__init__.py @@ -0,0 +1 @@ + diff --git a/src/srht_contrib/scripts/enqueue_actors.py b/src/srht_contrib/scripts/enqueue_actors.py new file mode 100644 index 0000000..6735881 --- /dev/null +++ b/src/srht_contrib/scripts/enqueue_actors.py @@ -0,0 +1,86 @@ +from __future__ import annotations + +import argparse +from datetime import UTC, datetime, timedelta +from pathlib import Path + +from sqlalchemy import select + +from srht_contrib.config import Settings +from srht_contrib.db import make_session_factory +from srht_contrib.models import TrackedActor + + +def _iter_usernames(path: Path) -> list[str]: + usernames: list[str] = [] + seen: set[str] = set() + for raw_line in path.read_text(encoding="utf-8").splitlines(): + username = raw_line.strip() + if not username or username.startswith("#"): + continue + if not username.startswith("~"): + username = f"~{username}" + if username in seen: + continue + seen.add(username) + usernames.append(username) + return usernames + + +def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: datetime | None = None) -> int: + settings = Settings() + session_factory = make_session_factory(settings) + usernames = _iter_usernames(username_file) + queued_at = start_at or datetime.now(tz=UTC) + inserted = 0 + + with session_factory() as db: + for index, actor in enumerate(usernames): + next_poll_after = queued_at + timedelta(seconds=index * stagger_seconds) + tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) + if tracked_actor is None: + tracked_actor = TrackedActor( + actor=actor, + is_active=True, + discovery_state="queued", + queued_for_discovery_at=queued_at, + next_poll_after=next_poll_after, + recent_backfill_status="pending", + ) + db.add(tracked_actor) + inserted += 1 + continue + + tracked_actor.is_active = True + if tracked_actor.queued_for_discovery_at is None: + tracked_actor.queued_for_discovery_at = queued_at + if tracked_actor.last_polled_at is None and tracked_actor.discovery_state != "indexed": + tracked_actor.discovery_state = "queued" + tracked_actor.next_poll_after = next_poll_after + db.commit() + + return inserted + + +def main() -> None: + parser = argparse.ArgumentParser(description="Durably enqueue SourceHut actors without polling them immediately.") + parser.add_argument( + "username_file", + nargs="?", + default="srht_usernames.txt", + help="Path to a newline-delimited SourceHut username file.", + ) + parser.add_argument( + "--stagger-seconds", + type=int, + default=300, + help="Seconds to space out each actor's first eligible poll time.", + ) + args = parser.parse_args() + + inserted = enqueue_actors(Path(args.username_file), stagger_seconds=args.stagger_seconds) + print(f"Enqueued {inserted} new actors from {args.username_file}.") + + +if __name__ == "__main__": + main() |
