Commit e27bbe17fd
Verified · cmc
Layout: unified · split
.gitignore +2
| @@ -29,3 +29,5 @@ dist/ | |||
| 29 | .DS_Store | 29 | .DS_Store |
| 30 | .idea/ | 30 | .idea/ |
| 31 | .vscode/ | 31 | .vscode/ |
| 32 | |||
| 33 | srht_usernames.txt | ||
API.md +1
| @@ -52,6 +52,7 @@ Background polling: | |||
| 52 | - When `ENABLE_SCHEDULER=true`, the service runs one poll immediately at startup and then continues polling on `POLL_INTERVAL_SECONDS`. | 52 | - When `ENABLE_SCHEDULER=true`, the service runs one poll immediately at startup and then continues polling on `POLL_INTERVAL_SECONDS`. |
| 53 | - The scheduler always seeds `DEFAULT_ACTOR` as a known actor. | 53 | - The scheduler always seeds `DEFAULT_ACTOR` as a known actor. |
| 54 | - Public contribution reads register additional actors for later background polling. | 54 | - Public contribution reads register additional actors for later background polling. |
| 55 | - The scheduler only processes due actors, up to `DISCOVERY_BATCH_SIZE` per pass. | ||
| 55 | - Manual polling remains available through `POST /api/contributions/poll`. | 56 | - Manual polling remains available through `POST /api/contributions/poll`. |
| 56 | 57 | ||
| 57 | Repository names: | 58 | Repository names: |
README.md +18
| @@ -83,6 +83,9 @@ Environment variables: | |||
| 83 | - `DATABASE_URL`: defaults to `sqlite:///./srht_contrib.db` | 83 | - `DATABASE_URL`: defaults to `sqlite:///./srht_contrib.db` |
| 84 | - `DEFAULT_ACTOR`: actor used by the scheduled poll job | 84 | - `DEFAULT_ACTOR`: actor used by the scheduled poll job |
| 85 | - `POLL_INTERVAL_SECONDS`: scheduler interval in seconds | 85 | - `POLL_INTERVAL_SECONDS`: scheduler interval in seconds |
| 86 | - `DISCOVERY_BATCH_SIZE`: max number of due actors to process per scheduler pass | ||
| 87 | - `INDEXED_ACTOR_REPOLL_SECONDS`: how long to wait before re-polling an already indexed actor | ||
| 88 | - `DISCOVERY_ERROR_BACKOFF_SECONDS`: base retry delay after a failed scheduled poll | ||
| 86 | - `ACTOR_ALIASES_JSON`: optional JSON object for actor/email/display-name alias mapping | 89 | - `ACTOR_ALIASES_JSON`: optional JSON object for actor/email/display-name alias mapping |
| 87 | - `GIT_TRACKED_REPOSITORIES`: optional JSON array of repository names or `owner/repo` strings to union into git polling | 90 | - `GIT_TRACKED_REPOSITORIES`: optional JSON array of repository names or `owner/repo` strings to union into git polling |
| 88 | 91 | ||
| @@ -97,6 +100,9 @@ GIT_SRHT_ENDPOINT=https://git.sr.ht/query | |||
| 97 | DATABASE_URL=sqlite:///./srht_contrib.db | 100 | DATABASE_URL=sqlite:///./srht_contrib.db |
| 98 | DEFAULT_ACTOR=~your-user | 101 | DEFAULT_ACTOR=~your-user |
| 99 | POLL_INTERVAL_SECONDS=900 | 102 | POLL_INTERVAL_SECONDS=900 |
| 103 | DISCOVERY_BATCH_SIZE=5 | ||
| 104 | INDEXED_ACTOR_REPOLL_SECONDS=21600 | ||
| 105 | DISCOVERY_ERROR_BACKOFF_SECONDS=3600 | ||
| 100 | ACTOR_ALIASES_JSON={"~your-user":["you@example.com","Your Name"]} | 106 | ACTOR_ALIASES_JSON={"~your-user":["you@example.com","Your Name"]} |
| 101 | GIT_TRACKED_REPOSITORIES=["your-repo","~your-user/your-site"] | 107 | GIT_TRACKED_REPOSITORIES=["your-repo","~your-user/your-site"] |
| 102 | ``` | 108 | ``` |
| @@ -167,6 +173,18 @@ Example response: | |||
| 167 | 173 | ||
| 168 | Scheduled polling only runs when `ENABLE_SCHEDULER=true`. The scheduler seeds `DEFAULT_ACTOR` as an initial known actor, runs one poll immediately at startup, and public contribution reads register additional actors for later background polling and one-year backfill. | 174 | Scheduled polling only runs when `ENABLE_SCHEDULER=true`. The scheduler seeds `DEFAULT_ACTOR` as an initial known actor, runs one poll immediately at startup, and public contribution reads register additional actors for later background polling and one-year backfill. |
| 169 | 175 | ||
| 176 | The scheduler now drains actors gradually instead of polling every tracked actor on every pass. It only claims due actors, up to `DISCOVERY_BATCH_SIZE` per run, then reschedules indexed actors with `INDEXED_ACTOR_REPOLL_SECONDS` and failed actors with backoff based on `DISCOVERY_ERROR_BACKOFF_SECONDS`. | ||
| 177 | |||
| 178 | ## Bulk Enqueue Without Immediate Indexing | ||
| 179 | |||
| 180 | To durably queue a large username list without polling it immediately: | ||
| 181 | |||
| 182 | ```bash | ||
| 183 | srht-enqueue-actors srht_usernames.txt --stagger-seconds 300 | ||
| 184 | ``` | ||
| 185 | |||
| 186 | This command stores usernames in `tracked_actors`, marks them queued, and spaces out their first eligible poll time. With `--stagger-seconds 300`, a file of 15,771 users will be spread across roughly 54.8 days before becoming due for first poll. | ||
| 187 | |||
| 170 | For `git.sr.ht`, owned repositories are auto-discovered for the actor. `GIT_TRACKED_REPOSITORIES` can still be used to union in extra repositories. Entries may be either: | 188 | For `git.sr.ht`, owned repositories are auto-discovered for the actor. `GIT_TRACKED_REPOSITORIES` can still be used to union in extra repositories. Entries may be either: |
| 171 | 189 | ||
| 172 | - `"Hutch"` for a repository owned by `DEFAULT_ACTOR` | 190 | - `"Hutch"` for a repository owned by `DEFAULT_ACTOR` |
alembic/versions/20260411_0006_actor_queue_scheduling.py added +80
| @@ -0,0 +1,80 @@ | |||
| 1 | """tracked actor queue scheduling fields""" | ||
| 2 | |||
| 3 | from __future__ import annotations | ||
| 4 | |||
| 5 | from alembic import op | ||
| 6 | import sqlalchemy as sa | ||
| 7 | from sqlalchemy import inspect | ||
| 8 | |||
| 9 | |||
| 10 | revision = "20260411_0006" | ||
| 11 | down_revision = "20260411_0005" | ||
| 12 | branch_labels = None | ||
| 13 | depends_on = None | ||
| 14 | |||
| 15 | |||
| 16 | def _column_names(table_name: str) -> set[str]: | ||
| 17 | return {column["name"] for column in inspect(op.get_bind()).get_columns(table_name)} | ||
| 18 | |||
| 19 | |||
| 20 | def upgrade() -> None: | ||
| 21 | columns = _column_names("tracked_actors") | ||
| 22 | |||
| 23 | if "discovery_state" not in columns: | ||
| 24 | op.add_column( | ||
| 25 | "tracked_actors", | ||
| 26 | sa.Column("discovery_state", sa.String(length=32), nullable=False, server_default="queued"), | ||
| 27 | ) | ||
| 28 | if "queued_for_discovery_at" not in columns: | ||
| 29 | op.add_column("tracked_actors", sa.Column("queued_for_discovery_at", sa.DateTime(timezone=True), nullable=True)) | ||
| 30 | if "next_poll_after" not in columns: | ||
| 31 | op.add_column("tracked_actors", sa.Column("next_poll_after", sa.DateTime(timezone=True), nullable=True)) | ||
| 32 | if "last_claimed_at" not in columns: | ||
| 33 | op.add_column("tracked_actors", sa.Column("last_claimed_at", sa.DateTime(timezone=True), nullable=True)) | ||
| 34 | if "poll_attempts" not in columns: | ||
| 35 | op.add_column( | ||
| 36 | "tracked_actors", | ||
| 37 | sa.Column("poll_attempts", sa.Integer(), nullable=False, server_default="0"), | ||
| 38 | ) | ||
| 39 | |||
| 40 | op.execute( | ||
| 41 | sa.text( | ||
| 42 | """ | ||
| 43 | UPDATE tracked_actors | ||
| 44 | SET discovery_state = CASE | ||
| 45 | WHEN last_poll_status IS NOT NULL THEN last_poll_status | ||
| 46 | ELSE 'queued' | ||
| 47 | END | ||
| 48 | """ | ||
| 49 | ) | ||
| 50 | ) | ||
| 51 | op.execute( | ||
| 52 | sa.text( | ||
| 53 | """ | ||
| 54 | UPDATE tracked_actors | ||
| 55 | SET queued_for_discovery_at = COALESCE(queued_for_discovery_at, last_requested_at, last_polled_at, CURRENT_TIMESTAMP) | ||
| 56 | """ | ||
| 57 | ) | ||
| 58 | ) | ||
| 59 | op.execute( | ||
| 60 | sa.text( | ||
| 61 | """ | ||
| 62 | UPDATE tracked_actors | ||
| 63 | SET next_poll_after = COALESCE(next_poll_after, last_polled_at, CURRENT_TIMESTAMP) | ||
| 64 | """ | ||
| 65 | ) | ||
| 66 | ) | ||
| 67 | |||
| 68 | |||
| 69 | def downgrade() -> None: | ||
| 70 | columns = _column_names("tracked_actors") | ||
| 71 | if "poll_attempts" in columns: | ||
| 72 | op.drop_column("tracked_actors", "poll_attempts") | ||
| 73 | if "last_claimed_at" in columns: | ||
| 74 | op.drop_column("tracked_actors", "last_claimed_at") | ||
| 75 | if "next_poll_after" in columns: | ||
| 76 | op.drop_column("tracked_actors", "next_poll_after") | ||
| 77 | if "queued_for_discovery_at" in columns: | ||
| 78 | op.drop_column("tracked_actors", "queued_for_discovery_at") | ||
| 79 | if "discovery_state" in columns: | ||
| 80 | op.drop_column("tracked_actors", "discovery_state") | ||
pyproject.toml +3
| @@ -24,6 +24,9 @@ dev = [ | |||
| 24 | "pytest>=8.2,<9.0", | 24 | "pytest>=8.2,<9.0", |
| 25 | ] | 25 | ] |
| 26 | 26 | ||
| 27 | [project.scripts] | ||
| 28 | srht-enqueue-actors = "srht_contrib.scripts.enqueue_actors:main" | ||
| 29 | |||
| 27 | [tool.setuptools] | 30 | [tool.setuptools] |
| 28 | package-dir = {"" = "src"} | 31 | package-dir = {"" = "src"} |
| 29 | 32 | ||
src/srht_contrib/config.py +3
| @@ -37,6 +37,9 @@ class Settings(BaseSettings): | |||
| 37 | poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS") | 37 | poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS") |
| 38 | sync_overlap_hours: int = Field(default=1, alias="SYNC_OVERLAP_HOURS") | 38 | sync_overlap_hours: int = Field(default=1, alias="SYNC_OVERLAP_HOURS") |
| 39 | srht_request_delay_seconds: float = Field(default=0.5, alias="SRHT_REQUEST_DELAY_SECONDS") | 39 | srht_request_delay_seconds: float = Field(default=0.5, alias="SRHT_REQUEST_DELAY_SECONDS") |
| 40 | discovery_batch_size: int = Field(default=5, alias="DISCOVERY_BATCH_SIZE") | ||
| 41 | indexed_actor_repoll_seconds: int = Field(default=21600, alias="INDEXED_ACTOR_REPOLL_SECONDS") | ||
| 42 | discovery_error_backoff_seconds: int = Field(default=3600, alias="DISCOVERY_ERROR_BACKOFF_SECONDS") | ||
| 40 | git_repo_discovery_ttl_seconds: int = Field(default=3600, alias="GIT_REPO_DISCOVERY_TTL_SECONDS") | 43 | git_repo_discovery_ttl_seconds: int = Field(default=3600, alias="GIT_REPO_DISCOVERY_TTL_SECONDS") |
| 41 | actor_aliases_json: dict[str, list[str]] = Field( | 44 | actor_aliases_json: dict[str, list[str]] = Field( |
| 42 | default_factory=dict, | 45 | default_factory=dict, |
src/srht_contrib/jobs/poller.py +39 −4
| @@ -4,7 +4,7 @@ import copy | |||
| 4 | import logging | 4 | import logging |
| 5 | from datetime import UTC, datetime, timedelta | 5 | from datetime import UTC, datetime, timedelta |
| 6 | 6 | ||
| 7 | from sqlalchemy import select | 7 | from sqlalchemy import or_, select |
| 8 | from sqlalchemy.exc import IntegrityError | 8 | from sqlalchemy.exc import IntegrityError |
| 9 | from sqlalchemy.orm import Session | 9 | from sqlalchemy.orm import Session |
| 10 | 10 | ||
| @@ -31,6 +31,7 @@ class PollerService: | |||
| 31 | ) -> None: | 31 | ) -> None: |
| 32 | self.todo_service = todo_service | 32 | self.todo_service = todo_service |
| 33 | self.git_service = git_service | 33 | self.git_service = git_service |
| 34 | self.settings = settings | ||
| 34 | self._sync_overlap = timedelta(hours=settings.sync_overlap_hours) | 35 | self._sync_overlap = timedelta(hours=settings.sync_overlap_hours) |
| 35 | 36 | ||
| 36 | def poll_all(self, db: Session, actor: str) -> int: | 37 | def poll_all(self, db: Session, actor: str) -> int: |
| @@ -57,9 +58,27 @@ class PollerService: | |||
| 57 | actors = db.scalars( | 58 | actors = db.scalars( |
| 58 | select(TrackedActor.actor) | 59 | select(TrackedActor.actor) |
| 59 | .where(TrackedActor.is_active.is_(True)) | 60 | .where(TrackedActor.is_active.is_(True)) |
| 60 | .order_by(TrackedActor.last_requested_at.is_(None), TrackedActor.last_requested_at.desc(), TrackedActor.actor) | 61 | .where( |
| 62 | or_( | ||
| 63 | TrackedActor.next_poll_after.is_(None), | ||
| 64 | TrackedActor.next_poll_after <= datetime.now(tz=UTC), | ||
| 65 | ) | ||
| 66 | ) | ||
| 67 | .order_by( | ||
| 68 | TrackedActor.next_poll_after.is_(None).desc(), | ||
| 69 | TrackedActor.next_poll_after, | ||
| 70 | TrackedActor.queued_for_discovery_at, | ||
| 71 | TrackedActor.actor, | ||
| 72 | ) | ||
| 73 | .limit(self.settings.discovery_batch_size) | ||
| 61 | ).all() | 74 | ).all() |
| 62 | for actor in actors: | 75 | for actor in actors: |
| 76 | claimed_actor = self.track_actor_request(db, actor, update_last_requested=False) | ||
| 77 | claimed_actor.discovery_state = "in_progress" | ||
| 78 | claimed_actor.last_claimed_at = datetime.now(tz=UTC) | ||
| 79 | claimed_actor.poll_attempts += 1 | ||
| 80 | db.add(claimed_actor) | ||
| 81 | db.commit() | ||
| 63 | try: | 82 | try: |
| 64 | results[actor] = self.poll_all(db, actor) | 83 | results[actor] = self.poll_all(db, actor) |
| 65 | except SourceHutClientError: | 84 | except SourceHutClientError: |
| @@ -73,17 +92,25 @@ class PollerService: | |||
| 73 | 92 | ||
| 74 | def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: | 93 | def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: |
| 75 | tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) | 94 | tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) |
| 95 | now = datetime.now(tz=UTC) | ||
| 76 | if tracked_actor is None: | 96 | if tracked_actor is None: |
| 77 | tracked_actor = TrackedActor( | 97 | tracked_actor = TrackedActor( |
| 78 | actor=actor, | 98 | actor=actor, |
| 79 | is_active=True, | 99 | is_active=True, |
| 100 | discovery_state="queued", | ||
| 101 | queued_for_discovery_at=now, | ||
| 102 | next_poll_after=now, | ||
| 80 | recent_backfill_status="pending", | 103 | recent_backfill_status="pending", |
| 81 | ) | 104 | ) |
| 82 | db.add(tracked_actor) | 105 | db.add(tracked_actor) |
| 83 | 106 | ||
| 84 | tracked_actor.is_active = True | 107 | tracked_actor.is_active = True |
| 108 | if tracked_actor.queued_for_discovery_at is None: | ||
| 109 | tracked_actor.queued_for_discovery_at = now | ||
| 110 | if tracked_actor.next_poll_after is None: | ||
| 111 | tracked_actor.next_poll_after = now | ||
| 85 | if update_last_requested: | 112 | if update_last_requested: |
| 86 | tracked_actor.last_requested_at = datetime.now(tz=UTC) | 113 | tracked_actor.last_requested_at = now |
| 87 | db.flush() | 114 | db.flush() |
| 88 | return tracked_actor | 115 | return tracked_actor |
| 89 | 116 | ||
| @@ -175,8 +202,16 @@ class PollerService: | |||
| 175 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) | 202 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) |
| 176 | tracked_actor.last_poll_status = status | 203 | tracked_actor.last_poll_status = status |
| 177 | tracked_actor.last_poll_error = error | 204 | tracked_actor.last_poll_error = error |
| 205 | now = datetime.now(tz=UTC) | ||
| 178 | if status == "indexed": | 206 | if status == "indexed": |
| 179 | tracked_actor.last_polled_at = datetime.now(tz=UTC) | 207 | tracked_actor.discovery_state = "indexed" |
| 208 | tracked_actor.last_polled_at = now | ||
| 209 | tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) | ||
| 210 | elif status == "error": | ||
| 211 | tracked_actor.discovery_state = "error" | ||
| 212 | tracked_actor.next_poll_after = now + timedelta( | ||
| 213 | seconds=self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1) | ||
| 214 | ) | ||
| 180 | db.add(tracked_actor) | 215 | db.add(tracked_actor) |
| 181 | db.flush() | 216 | db.flush() |
| 182 | 217 | ||
src/srht_contrib/models.py +5
| @@ -67,6 +67,11 @@ class TrackedActor(Base): | |||
| 67 | id: Mapped[int] = mapped_column(Integer, primary_key=True) | 67 | id: Mapped[int] = mapped_column(Integer, primary_key=True) |
| 68 | actor: Mapped[str] = mapped_column(String(255), nullable=False) | 68 | actor: Mapped[str] = mapped_column(String(255), nullable=False) |
| 69 | is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) | 69 | is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) |
| 70 | discovery_state: Mapped[str] = mapped_column(String(32), nullable=False, default="queued") | ||
| 71 | queued_for_discovery_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) | ||
| 72 | next_poll_after: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) | ||
| 73 | last_claimed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) | ||
| 74 | poll_attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0) | ||
| 70 | last_requested_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) | 75 | last_requested_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) |
| 71 | last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) | 76 | last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) |
| 72 | last_poll_status: Mapped[str | None] = mapped_column(String(32), nullable=True) | 77 | last_poll_status: Mapped[str | None] = mapped_column(String(32), nullable=True) |
src/srht_contrib/scripts/__init__.py added +1
| @@ -0,0 +1 @@ | |||
| 1 | |||
src/srht_contrib/scripts/enqueue_actors.py added +86
| @@ -0,0 +1,86 @@ | |||
| 1 | from __future__ import annotations | ||
| 2 | |||
| 3 | import argparse | ||
| 4 | from datetime import UTC, datetime, timedelta | ||
| 5 | from pathlib import Path | ||
| 6 | |||
| 7 | from sqlalchemy import select | ||
| 8 | |||
| 9 | from srht_contrib.config import Settings | ||
| 10 | from srht_contrib.db import make_session_factory | ||
| 11 | from srht_contrib.models import TrackedActor | ||
| 12 | |||
| 13 | |||
| 14 | def _iter_usernames(path: Path) -> list[str]: | ||
| 15 | usernames: list[str] = [] | ||
| 16 | seen: set[str] = set() | ||
| 17 | for raw_line in path.read_text(encoding="utf-8").splitlines(): | ||
| 18 | username = raw_line.strip() | ||
| 19 | if not username or username.startswith("#"): | ||
| 20 | continue | ||
| 21 | if not username.startswith("~"): | ||
| 22 | username = f"~{username}" | ||
| 23 | if username in seen: | ||
| 24 | continue | ||
| 25 | seen.add(username) | ||
| 26 | usernames.append(username) | ||
| 27 | return usernames | ||
| 28 | |||
| 29 | |||
| 30 | def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: datetime | None = None) -> int: | ||
| 31 | settings = Settings() | ||
| 32 | session_factory = make_session_factory(settings) | ||
| 33 | usernames = _iter_usernames(username_file) | ||
| 34 | queued_at = start_at or datetime.now(tz=UTC) | ||
| 35 | inserted = 0 | ||
| 36 | |||
| 37 | with session_factory() as db: | ||
| 38 | for index, actor in enumerate(usernames): | ||
| 39 | next_poll_after = queued_at + timedelta(seconds=index * stagger_seconds) | ||
| 40 | tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) | ||
| 41 | if tracked_actor is None: | ||
| 42 | tracked_actor = TrackedActor( | ||
| 43 | actor=actor, | ||
| 44 | is_active=True, | ||
| 45 | discovery_state="queued", | ||
| 46 | queued_for_discovery_at=queued_at, | ||
| 47 | next_poll_after=next_poll_after, | ||
| 48 | recent_backfill_status="pending", | ||
| 49 | ) | ||
| 50 | db.add(tracked_actor) | ||
| 51 | inserted += 1 | ||
| 52 | continue | ||
| 53 | |||
| 54 | tracked_actor.is_active = True | ||
| 55 | if tracked_actor.queued_for_discovery_at is None: | ||
| 56 | tracked_actor.queued_for_discovery_at = queued_at | ||
| 57 | if tracked_actor.last_polled_at is None and tracked_actor.discovery_state != "indexed": | ||
| 58 | tracked_actor.discovery_state = "queued" | ||
| 59 | tracked_actor.next_poll_after = next_poll_after | ||
| 60 | db.commit() | ||
| 61 | |||
| 62 | return inserted | ||
| 63 | |||
| 64 | |||
| 65 | def main() -> None: | ||
| 66 | parser = argparse.ArgumentParser(description="Durably enqueue SourceHut actors without polling them immediately.") | ||
| 67 | parser.add_argument( | ||
| 68 | "username_file", | ||
| 69 | nargs="?", | ||
| 70 | default="srht_usernames.txt", | ||
| 71 | help="Path to a newline-delimited SourceHut username file.", | ||
| 72 | ) | ||
| 73 | parser.add_argument( | ||
| 74 | "--stagger-seconds", | ||
| 75 | type=int, | ||
| 76 | default=300, | ||
| 77 | help="Seconds to space out each actor's first eligible poll time.", | ||
| 78 | ) | ||
| 79 | args = parser.parse_args() | ||
| 80 | |||
| 81 | inserted = enqueue_actors(Path(args.username_file), stagger_seconds=args.stagger_seconds) | ||
| 82 | print(f"Enqueued {inserted} new actors from {args.username_file}.") | ||
| 83 | |||
| 84 | |||
| 85 | if __name__ == "__main__": | ||
| 86 | main() | ||
tests/test_ingestion.py +88 −5
| @@ -1,10 +1,12 @@ | |||
| 1 | from datetime import UTC, datetime | 1 | from datetime import UTC, datetime, timedelta |
| 2 | from pathlib import Path | ||
| 2 | 3 | ||
| 3 | from sqlalchemy import select | 4 | from sqlalchemy import select |
| 4 | 5 | ||
| 5 | from srht_contrib.config import Settings | 6 | from srht_contrib.config import Settings |
| 6 | from srht_contrib.jobs.poller import PollerService | 7 | from srht_contrib.jobs.poller import PollerService |
| 7 | from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository | 8 | from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository |
| 9 | from srht_contrib.scripts.enqueue_actors import enqueue_actors | ||
| 8 | from srht_contrib.schemas import NormalizedEvent | 10 | from srht_contrib.schemas import NormalizedEvent |
| 9 | from srht_contrib.services.git import GitIngestionService, GitPollResult | 11 | from srht_contrib.services.git import GitIngestionService, GitPollResult |
| 10 | from srht_contrib.services.todo import TodoIngestionService, TodoPollResult | 12 | from srht_contrib.services.todo import TodoIngestionService, TodoPollResult |
| @@ -64,7 +66,7 @@ class EmptyGitService: | |||
| 64 | POLL_INTERVAL_SECONDS=60, | 66 | POLL_INTERVAL_SECONDS=60, |
| 65 | ) | 67 | ) |
| 66 | 68 | ||
| 67 | def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult: | 69 | def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None, db=None) -> GitPollResult: |
| 68 | return GitPollResult(events=[], cursor="2026-03-31T00:00:00+00:00") | 70 | return GitPollResult(events=[], cursor="2026-03-31T00:00:00+00:00") |
| 69 | 71 | ||
| 70 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: | 72 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: |
| @@ -156,7 +158,7 @@ class QueueShrinkingGitService: | |||
| 156 | POLL_INTERVAL_SECONDS=60, | 158 | POLL_INTERVAL_SECONDS=60, |
| 157 | ) | 159 | ) |
| 158 | 160 | ||
| 159 | def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult: | 161 | def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None, db=None) -> GitPollResult: |
| 160 | return GitPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat()) | 162 | return GitPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat()) |
| 161 | 163 | ||
| 162 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: | 164 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: |
| @@ -530,7 +532,7 @@ def test_sync_overlap_reuses_cursor_window_and_suppresses_duplicates(db_session) | |||
| 530 | assert second_inserted == 0 | 532 | assert second_inserted == 0 |
| 531 | assert state is not None | 533 | assert state is not None |
| 532 | assert len(todo_service.calls) == 2 | 534 | assert len(todo_service.calls) == 2 |
| 533 | assert todo_service.calls[1].isoformat() == "2026-03-30T00:00:00+00:00" | 535 | assert todo_service.calls[1].isoformat() == "2026-03-30T23:00:00+00:00" |
| 534 | 536 | ||
| 535 | 537 | ||
| 536 | def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) -> None: | 538 | def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) -> None: |
| @@ -556,7 +558,7 @@ def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) - | |||
| 556 | 558 | ||
| 557 | tracked_actors = db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() | 559 | tracked_actors = db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() |
| 558 | 560 | ||
| 559 | assert results == {"~default": 0, "~known": 1} | 561 | assert results == {"~default": 1, "~known": 0} |
| 560 | assert [actor.actor for actor in tracked_actors] == ["~default", "~known"] | 562 | assert [actor.actor for actor in tracked_actors] == ["~default", "~known"] |
| 561 | assert all(actor.last_poll_status == "indexed" for actor in tracked_actors) | 563 | assert all(actor.last_poll_status == "indexed" for actor in tracked_actors) |
| 562 | assert all(actor.last_polled_at is not None for actor in tracked_actors) | 564 | assert all(actor.last_polled_at is not None for actor in tracked_actors) |
| @@ -655,3 +657,84 @@ def test_prune_old_events_removes_data_older_than_one_year(db_session) -> None: | |||
| 655 | 657 | ||
| 656 | assert deleted == 1 | 658 | assert deleted == 1 |
| 657 | assert remaining == ["todo:recent"] | 659 | assert remaining == ["todo:recent"] |
| 660 | |||
| 661 | |||
| 662 | def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None: | ||
| 663 | settings = make_settings(DISCOVERY_BATCH_SIZE=2, INDEXED_ACTOR_REPOLL_SECONDS=3600) | ||
| 664 | todo_service = RecordingTodoService(events_by_call=[[], []]) | ||
| 665 | git_service = EmptyGitService() | ||
| 666 | poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings) | ||
| 667 | |||
| 668 | now = datetime.now(tz=UTC) | ||
| 669 | db_session.add_all( | ||
| 670 | [ | ||
| 671 | TrackedActor( | ||
| 672 | actor="~a", | ||
| 673 | is_active=True, | ||
| 674 | discovery_state="queued", | ||
| 675 | queued_for_discovery_at=now - timedelta(minutes=3), | ||
| 676 | next_poll_after=now - timedelta(minutes=3), | ||
| 677 | recent_backfill_status="completed", | ||
| 678 | ), | ||
| 679 | TrackedActor( | ||
| 680 | actor="~b", | ||
| 681 | is_active=True, | ||
| 682 | discovery_state="queued", | ||
| 683 | queued_for_discovery_at=now - timedelta(minutes=2), | ||
| 684 | next_poll_after=now - timedelta(minutes=2), | ||
| 685 | recent_backfill_status="completed", | ||
| 686 | ), | ||
| 687 | TrackedActor( | ||
| 688 | actor="~c", | ||
| 689 | is_active=True, | ||
| 690 | discovery_state="queued", | ||
| 691 | queued_for_discovery_at=now - timedelta(minutes=1), | ||
| 692 | next_poll_after=now - timedelta(minutes=1), | ||
| 693 | recent_backfill_status="completed", | ||
| 694 | ), | ||
| 695 | ] | ||
| 696 | ) | ||
| 697 | db_session.commit() | ||
| 698 | |||
| 699 | results = poller.poll_tracked_actors(db_session) | ||
| 700 | |||
| 701 | assert set(results) == {"~a", "~b"} | ||
| 702 | actors = { | ||
| 703 | actor.actor: actor | ||
| 704 | for actor in db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() | ||
| 705 | } | ||
| 706 | assert actors["~a"].discovery_state == "indexed" | ||
| 707 | assert actors["~b"].discovery_state == "indexed" | ||
| 708 | assert actors["~c"].discovery_state == "queued" | ||
| 709 | assert actors["~a"].poll_attempts == 1 | ||
| 710 | assert actors["~b"].poll_attempts == 1 | ||
| 711 | assert actors["~c"].poll_attempts == 0 | ||
| 712 | |||
| 713 | |||
| 714 | def test_enqueue_actors_staggers_without_polling(tmp_path, monkeypatch) -> None: | ||
| 715 | database_path = tmp_path / "enqueue.db" | ||
| 716 | username_path = tmp_path / "srht_usernames.txt" | ||
| 717 | username_path.write_text("alice\nbob\nalice\n~carol\n", encoding="utf-8") | ||
| 718 | monkeypatch.setenv("DATABASE_URL", f"sqlite:///{database_path}") | ||
| 719 | monkeypatch.setenv("SRHT_TOKEN", "test-token") | ||
| 720 | monkeypatch.setenv("DEFAULT_ACTOR", "~ccleberg") | ||
| 721 | |||
| 722 | from srht_contrib.db import Base, make_engine, make_session_factory | ||
| 723 | |||
| 724 | settings = Settings() | ||
| 725 | engine = make_engine(settings) | ||
| 726 | Base.metadata.create_all(bind=engine) | ||
| 727 | session_factory = make_session_factory(settings) | ||
| 728 | queued_at = datetime(2026, 4, 11, 12, 0, tzinfo=UTC) | ||
| 729 | |||
| 730 | inserted = enqueue_actors(Path(username_path), stagger_seconds=60, start_at=queued_at) | ||
| 731 | |||
| 732 | with session_factory() as db: | ||
| 733 | actors = db.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() | ||
| 734 | |||
| 735 | assert inserted == 3 | ||
| 736 | assert [actor.actor for actor in actors] == ["~alice", "~bob", "~carol"] | ||
| 737 | assert all(actor.discovery_state == "queued" for actor in actors) | ||
| 738 | assert actors[0].next_poll_after == queued_at.replace(tzinfo=None) | ||
| 739 | assert actors[1].next_poll_after == (queued_at + timedelta(seconds=60)).replace(tzinfo=None) | ||
| 740 | assert actors[2].next_poll_after == (queued_at + timedelta(seconds=120)).replace(tzinfo=None) | ||
tests/test_migrations.py +2
| @@ -111,6 +111,7 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None: | |||
| 111 | 111 | ||
| 112 | inspector = inspect(create_engine(database_url)) | 112 | inspector = inspect(create_engine(database_url)) |
| 113 | columns = {column["name"]: column for column in inspector.get_columns("tracked_repositories")} | 113 | columns = {column["name"]: column for column in inspector.get_columns("tracked_repositories")} |
| 114 | tracked_actor_columns = {column["name"] for column in inspector.get_columns("tracked_actors")} | ||
| 114 | unique_constraints = {constraint["name"] for constraint in inspector.get_unique_constraints("tracked_repositories")} | 115 | unique_constraints = {constraint["name"] for constraint in inspector.get_unique_constraints("tracked_repositories")} |
| 115 | with create_engine(database_url).connect() as connection: | 116 | with create_engine(database_url).connect() as connection: |
| 116 | actor = connection.execute(text("SELECT actor FROM tracked_repositories WHERE id = 1")).scalar_one() | 117 | actor = connection.execute(text("SELECT actor FROM tracked_repositories WHERE id = 1")).scalar_one() |
| @@ -121,6 +122,7 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None: | |||
| 121 | assert "discovered_repositories" in inspector.get_table_names() | 122 | assert "discovered_repositories" in inspector.get_table_names() |
| 122 | assert "tracked_actors" in inspector.get_table_names() | 123 | assert "tracked_actors" in inspector.get_table_names() |
| 123 | assert "service_backfill_states" in inspector.get_table_names() | 124 | assert "service_backfill_states" in inspector.get_table_names() |
| 125 | assert {"discovery_state", "queued_for_discovery_at", "next_poll_after", "last_claimed_at", "poll_attempts"} <= tracked_actor_columns | ||
| 124 | 126 | ||
| 125 | 127 | ||
| 126 | def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: | 128 | def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: |