summaryrefslogtreecommitdiff
path: root/src/srht_contrib
diff options
context:
space:
mode:
authorChristian Cleberg <[email protected]>2026-04-11 22:48:38 -0500
committerChristian Cleberg <[email protected]>2026-04-11 22:48:38 -0500
commite27bbe17fd5d8f06fc3ccb4a8de97a0775a1439c (patch)
tree7ce854d65009decbe8631175f233f3b2c9df6373 /src/srht_contrib
parent533866679755bd6e7a97cfa0f050eaa832b0b373 (diff)
downloadhutch-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')
-rw-r--r--src/srht_contrib/config.py3
-rw-r--r--src/srht_contrib/jobs/poller.py43
-rw-r--r--src/srht_contrib/models.py5
-rw-r--r--src/srht_contrib/scripts/__init__.py1
-rw-r--r--src/srht_contrib/scripts/enqueue_actors.py86
5 files changed, 134 insertions, 4 deletions
diff --git a/src/srht_contrib/config.py b/src/srht_contrib/config.py
index 80884c1..2cab56e 100644
--- a/src/srht_contrib/config.py
+++ b/src/srht_contrib/config.py
@@ -37,6 +37,9 @@ class Settings(BaseSettings):
poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS")
sync_overlap_hours: int = Field(default=1, alias="SYNC_OVERLAP_HOURS")
srht_request_delay_seconds: float = Field(default=0.5, alias="SRHT_REQUEST_DELAY_SECONDS")
+ discovery_batch_size: int = Field(default=5, alias="DISCOVERY_BATCH_SIZE")
+ indexed_actor_repoll_seconds: int = Field(default=21600, alias="INDEXED_ACTOR_REPOLL_SECONDS")
+ discovery_error_backoff_seconds: int = Field(default=3600, alias="DISCOVERY_ERROR_BACKOFF_SECONDS")
git_repo_discovery_ttl_seconds: int = Field(default=3600, alias="GIT_REPO_DISCOVERY_TTL_SECONDS")
actor_aliases_json: dict[str, list[str]] = Field(
default_factory=dict,
diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py
index b2fe70f..5e62b55 100644
--- a/src/srht_contrib/jobs/poller.py
+++ b/src/srht_contrib/jobs/poller.py
@@ -4,7 +4,7 @@ import copy
import logging
from datetime import UTC, datetime, timedelta
-from sqlalchemy import select
+from sqlalchemy import or_, select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
@@ -31,6 +31,7 @@ class PollerService:
) -> None:
self.todo_service = todo_service
self.git_service = git_service
+ self.settings = settings
self._sync_overlap = timedelta(hours=settings.sync_overlap_hours)
def poll_all(self, db: Session, actor: str) -> int:
@@ -57,9 +58,27 @@ class PollerService:
actors = db.scalars(
select(TrackedActor.actor)
.where(TrackedActor.is_active.is_(True))
- .order_by(TrackedActor.last_requested_at.is_(None), TrackedActor.last_requested_at.desc(), TrackedActor.actor)
+ .where(
+ or_(
+ TrackedActor.next_poll_after.is_(None),
+ TrackedActor.next_poll_after <= datetime.now(tz=UTC),
+ )
+ )
+ .order_by(
+ TrackedActor.next_poll_after.is_(None).desc(),
+ TrackedActor.next_poll_after,
+ TrackedActor.queued_for_discovery_at,
+ TrackedActor.actor,
+ )
+ .limit(self.settings.discovery_batch_size)
).all()
for actor in actors:
+ claimed_actor = self.track_actor_request(db, actor, update_last_requested=False)
+ claimed_actor.discovery_state = "in_progress"
+ claimed_actor.last_claimed_at = datetime.now(tz=UTC)
+ claimed_actor.poll_attempts += 1
+ db.add(claimed_actor)
+ db.commit()
try:
results[actor] = self.poll_all(db, actor)
except SourceHutClientError:
@@ -73,17 +92,25 @@ class PollerService:
def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor:
tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
+ now = datetime.now(tz=UTC)
if tracked_actor is None:
tracked_actor = TrackedActor(
actor=actor,
is_active=True,
+ discovery_state="queued",
+ queued_for_discovery_at=now,
+ next_poll_after=now,
recent_backfill_status="pending",
)
db.add(tracked_actor)
tracked_actor.is_active = True
+ if tracked_actor.queued_for_discovery_at is None:
+ tracked_actor.queued_for_discovery_at = now
+ if tracked_actor.next_poll_after is None:
+ tracked_actor.next_poll_after = now
if update_last_requested:
- tracked_actor.last_requested_at = datetime.now(tz=UTC)
+ tracked_actor.last_requested_at = now
db.flush()
return tracked_actor
@@ -175,8 +202,16 @@ class PollerService:
tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
tracked_actor.last_poll_status = status
tracked_actor.last_poll_error = error
+ now = datetime.now(tz=UTC)
if status == "indexed":
- tracked_actor.last_polled_at = datetime.now(tz=UTC)
+ tracked_actor.discovery_state = "indexed"
+ tracked_actor.last_polled_at = now
+ tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds)
+ elif status == "error":
+ tracked_actor.discovery_state = "error"
+ tracked_actor.next_poll_after = now + timedelta(
+ seconds=self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1)
+ )
db.add(tracked_actor)
db.flush()
diff --git a/src/srht_contrib/models.py b/src/srht_contrib/models.py
index e0725ff..441b6ec 100644
--- a/src/srht_contrib/models.py
+++ b/src/srht_contrib/models.py
@@ -67,6 +67,11 @@ class TrackedActor(Base):
id: Mapped[int] = mapped_column(Integer, primary_key=True)
actor: Mapped[str] = mapped_column(String(255), nullable=False)
is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
+ discovery_state: Mapped[str] = mapped_column(String(32), nullable=False, default="queued")
+ queued_for_discovery_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
+ next_poll_after: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
+ last_claimed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
+ poll_attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
last_requested_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
last_poll_status: Mapped[str | None] = mapped_column(String(32), nullable=True)
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()