diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/srht_contrib/config.py | 3 | ||||
| -rw-r--r-- | src/srht_contrib/jobs/poller.py | 43 | ||||
| -rw-r--r-- | src/srht_contrib/models.py | 5 | ||||
| -rw-r--r-- | src/srht_contrib/scripts/__init__.py | 1 | ||||
| -rw-r--r-- | src/srht_contrib/scripts/enqueue_actors.py | 86 |
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() |
