From e27bbe17fd5d8f06fc3ccb4a8de97a0775a1439c Mon Sep 17 00:00:00 2001 From: Christian Cleberg Date: Sat, 11 Apr 2026 22:48:38 -0500 Subject: add durable actor queueing and staggered discovery imports --- .gitignore | 2 + API.md | 1 + README.md | 18 +++++ .../20260411_0006_actor_queue_scheduling.py | 80 +++++++++++++++++++ pyproject.toml | 3 + src/srht_contrib/config.py | 3 + src/srht_contrib/jobs/poller.py | 43 +++++++++- src/srht_contrib/models.py | 5 ++ src/srht_contrib/scripts/__init__.py | 1 + src/srht_contrib/scripts/enqueue_actors.py | 86 ++++++++++++++++++++ tests/test_ingestion.py | 93 ++++++++++++++++++++-- tests/test_migrations.py | 2 + 12 files changed, 328 insertions(+), 9 deletions(-) create mode 100644 alembic/versions/20260411_0006_actor_queue_scheduling.py create mode 100644 src/srht_contrib/scripts/__init__.py create mode 100644 src/srht_contrib/scripts/enqueue_actors.py diff --git a/.gitignore b/.gitignore index 3ee03bd..2006f36 100644 --- a/.gitignore +++ b/.gitignore @@ -29,3 +29,5 @@ dist/ .DS_Store .idea/ .vscode/ + +srht_usernames.txt diff --git a/API.md b/API.md index 8d05885..0117df2 100644 --- a/API.md +++ b/API.md @@ -52,6 +52,7 @@ Background polling: - When `ENABLE_SCHEDULER=true`, the service runs one poll immediately at startup and then continues polling on `POLL_INTERVAL_SECONDS`. - The scheduler always seeds `DEFAULT_ACTOR` as a known actor. - Public contribution reads register additional actors for later background polling. +- The scheduler only processes due actors, up to `DISCOVERY_BATCH_SIZE` per pass. - Manual polling remains available through `POST /api/contributions/poll`. Repository names: diff --git a/README.md b/README.md index 656f1c7..695d222 100644 --- a/README.md +++ b/README.md @@ -83,6 +83,9 @@ Environment variables: - `DATABASE_URL`: defaults to `sqlite:///./srht_contrib.db` - `DEFAULT_ACTOR`: actor used by the scheduled poll job - `POLL_INTERVAL_SECONDS`: scheduler interval in seconds +- `DISCOVERY_BATCH_SIZE`: max number of due actors to process per scheduler pass +- `INDEXED_ACTOR_REPOLL_SECONDS`: how long to wait before re-polling an already indexed actor +- `DISCOVERY_ERROR_BACKOFF_SECONDS`: base retry delay after a failed scheduled poll - `ACTOR_ALIASES_JSON`: optional JSON object for actor/email/display-name alias mapping - `GIT_TRACKED_REPOSITORIES`: optional JSON array of repository names or `owner/repo` strings to union into git polling @@ -97,6 +100,9 @@ GIT_SRHT_ENDPOINT=https://git.sr.ht/query DATABASE_URL=sqlite:///./srht_contrib.db DEFAULT_ACTOR=~your-user POLL_INTERVAL_SECONDS=900 +DISCOVERY_BATCH_SIZE=5 +INDEXED_ACTOR_REPOLL_SECONDS=21600 +DISCOVERY_ERROR_BACKOFF_SECONDS=3600 ACTOR_ALIASES_JSON={"~your-user":["you@example.com","Your Name"]} GIT_TRACKED_REPOSITORIES=["your-repo","~your-user/your-site"] ``` @@ -167,6 +173,18 @@ Example response: 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. +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`. + +## Bulk Enqueue Without Immediate Indexing + +To durably queue a large username list without polling it immediately: + +```bash +srht-enqueue-actors srht_usernames.txt --stagger-seconds 300 +``` + +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. + 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: - `"Hutch"` for a repository owned by `DEFAULT_ACTOR` diff --git a/alembic/versions/20260411_0006_actor_queue_scheduling.py b/alembic/versions/20260411_0006_actor_queue_scheduling.py new file mode 100644 index 0000000..a4fb9c1 --- /dev/null +++ b/alembic/versions/20260411_0006_actor_queue_scheduling.py @@ -0,0 +1,80 @@ +"""tracked actor queue scheduling fields""" + +from __future__ import annotations + +from alembic import op +import sqlalchemy as sa +from sqlalchemy import inspect + + +revision = "20260411_0006" +down_revision = "20260411_0005" +branch_labels = None +depends_on = None + + +def _column_names(table_name: str) -> set[str]: + return {column["name"] for column in inspect(op.get_bind()).get_columns(table_name)} + + +def upgrade() -> None: + columns = _column_names("tracked_actors") + + if "discovery_state" not in columns: + op.add_column( + "tracked_actors", + sa.Column("discovery_state", sa.String(length=32), nullable=False, server_default="queued"), + ) + if "queued_for_discovery_at" not in columns: + op.add_column("tracked_actors", sa.Column("queued_for_discovery_at", sa.DateTime(timezone=True), nullable=True)) + if "next_poll_after" not in columns: + op.add_column("tracked_actors", sa.Column("next_poll_after", sa.DateTime(timezone=True), nullable=True)) + if "last_claimed_at" not in columns: + op.add_column("tracked_actors", sa.Column("last_claimed_at", sa.DateTime(timezone=True), nullable=True)) + if "poll_attempts" not in columns: + op.add_column( + "tracked_actors", + sa.Column("poll_attempts", sa.Integer(), nullable=False, server_default="0"), + ) + + op.execute( + sa.text( + """ + UPDATE tracked_actors + SET discovery_state = CASE + WHEN last_poll_status IS NOT NULL THEN last_poll_status + ELSE 'queued' + END + """ + ) + ) + op.execute( + sa.text( + """ + UPDATE tracked_actors + SET queued_for_discovery_at = COALESCE(queued_for_discovery_at, last_requested_at, last_polled_at, CURRENT_TIMESTAMP) + """ + ) + ) + op.execute( + sa.text( + """ + UPDATE tracked_actors + SET next_poll_after = COALESCE(next_poll_after, last_polled_at, CURRENT_TIMESTAMP) + """ + ) + ) + + +def downgrade() -> None: + columns = _column_names("tracked_actors") + if "poll_attempts" in columns: + op.drop_column("tracked_actors", "poll_attempts") + if "last_claimed_at" in columns: + op.drop_column("tracked_actors", "last_claimed_at") + if "next_poll_after" in columns: + op.drop_column("tracked_actors", "next_poll_after") + if "queued_for_discovery_at" in columns: + op.drop_column("tracked_actors", "queued_for_discovery_at") + if "discovery_state" in columns: + op.drop_column("tracked_actors", "discovery_state") diff --git a/pyproject.toml b/pyproject.toml index 8eb5497..4030ab0 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -24,6 +24,9 @@ dev = [ "pytest>=8.2,<9.0", ] +[project.scripts] +srht-enqueue-actors = "srht_contrib.scripts.enqueue_actors:main" + [tool.setuptools] package-dir = {"" = "src"} 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() diff --git a/tests/test_ingestion.py b/tests/test_ingestion.py index 53848e7..974dac6 100644 --- a/tests/test_ingestion.py +++ b/tests/test_ingestion.py @@ -1,10 +1,12 @@ -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta +from pathlib import Path from sqlalchemy import select from srht_contrib.config import Settings from srht_contrib.jobs.poller import PollerService from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository +from srht_contrib.scripts.enqueue_actors import enqueue_actors from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.git import GitIngestionService, GitPollResult from srht_contrib.services.todo import TodoIngestionService, TodoPollResult @@ -64,7 +66,7 @@ class EmptyGitService: POLL_INTERVAL_SECONDS=60, ) - def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult: + def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None, db=None) -> GitPollResult: return GitPollResult(events=[], cursor="2026-03-31T00:00:00+00:00") def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: @@ -156,7 +158,7 @@ class QueueShrinkingGitService: POLL_INTERVAL_SECONDS=60, ) - def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult: + def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None, db=None) -> GitPollResult: return GitPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat()) 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) assert second_inserted == 0 assert state is not None assert len(todo_service.calls) == 2 - assert todo_service.calls[1].isoformat() == "2026-03-30T00:00:00+00:00" + assert todo_service.calls[1].isoformat() == "2026-03-30T23:00:00+00:00" 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) - tracked_actors = db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() - assert results == {"~default": 0, "~known": 1} + assert results == {"~default": 1, "~known": 0} assert [actor.actor for actor in tracked_actors] == ["~default", "~known"] assert all(actor.last_poll_status == "indexed" for actor in tracked_actors) 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: assert deleted == 1 assert remaining == ["todo:recent"] + + +def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None: + settings = make_settings(DISCOVERY_BATCH_SIZE=2, INDEXED_ACTOR_REPOLL_SECONDS=3600) + todo_service = RecordingTodoService(events_by_call=[[], []]) + git_service = EmptyGitService() + poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings) + + now = datetime.now(tz=UTC) + db_session.add_all( + [ + TrackedActor( + actor="~a", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=3), + next_poll_after=now - timedelta(minutes=3), + recent_backfill_status="completed", + ), + TrackedActor( + actor="~b", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=2), + next_poll_after=now - timedelta(minutes=2), + recent_backfill_status="completed", + ), + TrackedActor( + actor="~c", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=1), + next_poll_after=now - timedelta(minutes=1), + recent_backfill_status="completed", + ), + ] + ) + db_session.commit() + + results = poller.poll_tracked_actors(db_session) + + assert set(results) == {"~a", "~b"} + actors = { + actor.actor: actor + for actor in db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() + } + assert actors["~a"].discovery_state == "indexed" + assert actors["~b"].discovery_state == "indexed" + assert actors["~c"].discovery_state == "queued" + assert actors["~a"].poll_attempts == 1 + assert actors["~b"].poll_attempts == 1 + assert actors["~c"].poll_attempts == 0 + + +def test_enqueue_actors_staggers_without_polling(tmp_path, monkeypatch) -> None: + database_path = tmp_path / "enqueue.db" + username_path = tmp_path / "srht_usernames.txt" + username_path.write_text("alice\nbob\nalice\n~carol\n", encoding="utf-8") + monkeypatch.setenv("DATABASE_URL", f"sqlite:///{database_path}") + monkeypatch.setenv("SRHT_TOKEN", "test-token") + monkeypatch.setenv("DEFAULT_ACTOR", "~ccleberg") + + from srht_contrib.db import Base, make_engine, make_session_factory + + settings = Settings() + engine = make_engine(settings) + Base.metadata.create_all(bind=engine) + session_factory = make_session_factory(settings) + queued_at = datetime(2026, 4, 11, 12, 0, tzinfo=UTC) + + inserted = enqueue_actors(Path(username_path), stagger_seconds=60, start_at=queued_at) + + with session_factory() as db: + actors = db.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() + + assert inserted == 3 + assert [actor.actor for actor in actors] == ["~alice", "~bob", "~carol"] + assert all(actor.discovery_state == "queued" for actor in actors) + assert actors[0].next_poll_after == queued_at.replace(tzinfo=None) + assert actors[1].next_poll_after == (queued_at + timedelta(seconds=60)).replace(tzinfo=None) + assert actors[2].next_poll_after == (queued_at + timedelta(seconds=120)).replace(tzinfo=None) diff --git a/tests/test_migrations.py b/tests/test_migrations.py index a3091fc..3ce842b 100644 --- a/tests/test_migrations.py +++ b/tests/test_migrations.py @@ -111,6 +111,7 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None: inspector = inspect(create_engine(database_url)) columns = {column["name"]: column for column in inspector.get_columns("tracked_repositories")} + tracked_actor_columns = {column["name"] for column in inspector.get_columns("tracked_actors")} unique_constraints = {constraint["name"] for constraint in inspector.get_unique_constraints("tracked_repositories")} with create_engine(database_url).connect() as connection: 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: assert "discovered_repositories" in inspector.get_table_names() assert "tracked_actors" in inspector.get_table_names() assert "service_backfill_states" in inspector.get_table_names() + assert {"discovery_state", "queued_for_discovery_at", "next_poll_after", "last_claimed_at", "poll_attempts"} <= tracked_actor_columns def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: -- cgit v1.2.3