From 6289a28296bd0863d373bb47c1c85038f4bc6a2b Mon Sep 17 00:00:00 2001 From: Christian Cleberg Date: Wed, 6 May 2026 23:19:54 -0500 Subject: feat: improve scheduled indexing throughput and capped backoff --- .env.example | 7 +- API.md | 1 + README.md | 12 +-- compose.yml | 5 +- src/srht_contrib/config.py | 5 +- src/srht_contrib/jobs/poller.py | 56 ++++++++++++-- src/srht_contrib/scripts/enqueue_actors.py | 4 +- tests/test_ingestion.py | 118 ++++++++++++++++++++++++++++- 8 files changed, 187 insertions(+), 21 deletions(-) diff --git a/.env.example b/.env.example index 95d9645..44a5d5e 100644 --- a/.env.example +++ b/.env.example @@ -5,7 +5,12 @@ TODO_SRHT_ENDPOINT=https://todo.sr.ht/query GIT_SRHT_ENDPOINT=https://git.sr.ht/query DATABASE_URL=sqlite:///./srht_contrib.db DEFAULT_ACTOR=~your-user -POLL_INTERVAL_SECONDS=900 +POLL_INTERVAL_SECONDS=300 +DISCOVERY_BATCH_SIZE=20 +INDEXED_ACTOR_REPOLL_SECONDS=21600 +DISCOVERY_ERROR_BACKOFF_SECONDS=3600 +DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=21600 +SRHT_REQUEST_DELAY_SECONDS=0.5 # Optional JSON object. Example: # {"~your-user":["you@example.com","Your Name"]} ACTOR_ALIASES_JSON={} diff --git a/API.md b/API.md index 43fb75f..8c810ab 100644 --- a/API.md +++ b/API.md @@ -54,6 +54,7 @@ Background polling: - 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. +- Scheduled first indexing skips bounded one-year backfill work, then drains backfill in later scheduled passes so newly requested actors become indexed sooner. - Manual polling remains available through `POST /api/contributions/poll`. Repository names: diff --git a/README.md b/README.md index 03b2cec..6f862e6 100644 --- a/README.md +++ b/README.md @@ -86,6 +86,7 @@ Environment variables: - `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 +- `DISCOVERY_ERROR_BACKOFF_MAX_SECONDS`: maximum retry delay after repeated scheduled poll failures - `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 @@ -99,10 +100,11 @@ TODO_SRHT_ENDPOINT=https://todo.sr.ht/query 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 +POLL_INTERVAL_SECONDS=300 +DISCOVERY_BATCH_SIZE=20 INDEXED_ACTOR_REPOLL_SECONDS=21600 DISCOVERY_ERROR_BACKOFF_SECONDS=3600 +DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=21600 ACTOR_ALIASES_JSON={"~your-user":["you@example.com","Your Name"]} GIT_TRACKED_REPOSITORIES=["your-repo","~your-user/your-site"] ``` @@ -173,7 +175,7 @@ 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`. +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, performs a fast first indexing pass before bounded one-year backfill work, then reschedules indexed actors with `INDEXED_ACTOR_REPOLL_SECONDS` and failed actors with capped backoff based on `DISCOVERY_ERROR_BACKOFF_SECONDS`. Clients can explicitly signal that a public contribution read is for the signed-in user's own graph by sending `prioritize_self=true` on the read request. That temporarily boosts the actor to the front of the due queue for the next indexing pass, then clears the boost after the poll completes. @@ -182,10 +184,10 @@ Clients can explicitly signal that a public contribution read is for the signed- To durably queue a large username list without polling it immediately: ```bash -srht-enqueue-actors srht_usernames.txt --stagger-seconds 300 +srht-enqueue-actors srht_usernames.txt --stagger-seconds 60 ``` -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. +This command stores usernames in `tracked_actors`, marks them queued, and spaces out their first eligible poll time. With `--stagger-seconds 60`, a file of 15,771 users will be spread across roughly 11 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: diff --git a/compose.yml b/compose.yml index 74fd9cf..feb2926 100644 --- a/compose.yml +++ b/compose.yml @@ -13,10 +13,11 @@ services: GIT_SRHT_ENDPOINT: ${GIT_SRHT_ENDPOINT:-https://git.sr.ht/query} DATABASE_URL: sqlite:////data/srht_contrib.db DEFAULT_ACTOR: ${DEFAULT_ACTOR} - POLL_INTERVAL_SECONDS: ${POLL_INTERVAL_SECONDS:-900} - DISCOVERY_BATCH_SIZE: ${DISCOVERY_BATCH_SIZE:-5} + POLL_INTERVAL_SECONDS: ${POLL_INTERVAL_SECONDS:-300} + DISCOVERY_BATCH_SIZE: ${DISCOVERY_BATCH_SIZE:-20} INDEXED_ACTOR_REPOLL_SECONDS: ${INDEXED_ACTOR_REPOLL_SECONDS:-21600} DISCOVERY_ERROR_BACKOFF_SECONDS: ${DISCOVERY_ERROR_BACKOFF_SECONDS:-3600} + DISCOVERY_ERROR_BACKOFF_MAX_SECONDS: ${DISCOVERY_ERROR_BACKOFF_MAX_SECONDS:-21600} SRHT_REQUEST_DELAY_SECONDS: ${SRHT_REQUEST_DELAY_SECONDS:-0.5} ACTOR_ALIASES_JSON: ${ACTOR_ALIASES_JSON:-{}} GIT_TRACKED_REPOSITORIES: ${GIT_TRACKED_REPOSITORIES:-[]} diff --git a/src/srht_contrib/config.py b/src/srht_contrib/config.py index 3a7af2c..9221e4c 100644 --- a/src/srht_contrib/config.py +++ b/src/srht_contrib/config.py @@ -34,13 +34,14 @@ class Settings(BaseSettings): alias="DATABASE_URL", ) default_actor: str = Field(default="~unknown", alias="DEFAULT_ACTOR") - poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS") + poll_interval_seconds: int = Field(default=300, 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") sqlite_busy_timeout_seconds: float = Field(default=30.0, alias="SQLITE_BUSY_TIMEOUT_SECONDS") - discovery_batch_size: int = Field(default=5, alias="DISCOVERY_BATCH_SIZE") + discovery_batch_size: int = Field(default=20, 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") + discovery_error_backoff_max_seconds: int = Field(default=21600, alias="DISCOVERY_ERROR_BACKOFF_MAX_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 76d2fb0..3e6772b 100644 --- a/src/srht_contrib/jobs/poller.py +++ b/src/srht_contrib/jobs/poller.py @@ -34,7 +34,7 @@ class PollerService: self.settings = settings self._sync_overlap = timedelta(hours=settings.sync_overlap_hours) - def poll_all(self, db: Session, actor: str) -> int: + def poll_all(self, db: Session, actor: str, *, run_backfill: bool = True) -> int: self.track_actor_request(db, actor, update_last_requested=False) try: inserted = self._poll_actor(db, actor) @@ -45,7 +45,8 @@ class PollerService: raise self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None) - inserted += self._run_backfill_batches(db, actor) + if run_backfill: + inserted += self._run_backfill_batches(db, actor) db.commit() return inserted @@ -56,7 +57,7 @@ class PollerService: results: dict[str, int] = {} actors = db.scalars( - select(TrackedActor.actor) + select(TrackedActor) .where(TrackedActor.is_active.is_(True)) .where( or_( @@ -74,7 +75,12 @@ class PollerService: ) .limit(self.settings.discovery_batch_size) ).all() - for actor in actors: + for tracked_actor in actors: + actor = tracked_actor.actor + run_backfill_only = ( + tracked_actor.last_polled_at is not None + and tracked_actor.recent_backfill_status != "completed" + ) 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) @@ -82,7 +88,12 @@ class PollerService: db.add(claimed_actor) db.commit() try: - results[actor] = self.poll_all(db, actor) + if run_backfill_only: + results[actor] = self.poll_recent_backfill(db, actor) + else: + results[actor] = self.poll_all(db, actor, run_backfill=False) + self._schedule_pending_backfill_or_repoll(db, actor) + db.commit() except SourceHutClientError: logger.exception("Scheduled poll failed for actor=%s", actor) except Exception: @@ -220,15 +231,46 @@ class PollerService: tracked_actor.last_polled_at = now tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) tracked_actor.priority_boosted_at = None + tracked_actor.poll_attempts = 0 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) + backoff_seconds = min( + self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1), + self.settings.discovery_error_backoff_max_seconds, ) + tracked_actor.next_poll_after = now + timedelta(seconds=backoff_seconds) tracked_actor.priority_boosted_at = None db.add(tracked_actor) db.flush() + def poll_recent_backfill(self, db: Session, actor: str) -> int: + try: + inserted = self._run_backfill_batches(db, actor) + except Exception as exc: + db.rollback() + self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc)) + db.commit() + raise + + self._schedule_pending_backfill_or_repoll(db, actor) + db.commit() + return inserted + + def _schedule_pending_backfill_or_repoll(self, db: Session, actor: str) -> None: + tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) + now = datetime.now(tz=UTC) + tracked_actor.discovery_state = "indexed" + tracked_actor.last_poll_status = "indexed" + tracked_actor.last_poll_error = None + tracked_actor.priority_boosted_at = None + tracked_actor.poll_attempts = 0 + if tracked_actor.recent_backfill_status == "completed": + tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) + else: + tracked_actor.next_poll_after = now + db.add(tracked_actor) + db.flush() + def _run_backfill_batches(self, db: Session, actor: str) -> int: tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) if tracked_actor.recent_backfill_status == "completed": diff --git a/src/srht_contrib/scripts/enqueue_actors.py b/src/srht_contrib/scripts/enqueue_actors.py index 03f1f9b..bd838cb 100644 --- a/src/srht_contrib/scripts/enqueue_actors.py +++ b/src/srht_contrib/scripts/enqueue_actors.py @@ -40,7 +40,7 @@ def _iter_usernames(path: Path) -> list[str]: return usernames -def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: datetime | None = None) -> int: +def enqueue_actors(username_file: Path, *, stagger_seconds: int = 60, start_at: datetime | None = None) -> int: settings = Settings() session_factory = make_session_factory(settings) usernames = _iter_usernames(username_file) @@ -95,7 +95,7 @@ def main() -> None: parser.add_argument( "--stagger-seconds", type=int, - default=300, + default=60, help="Seconds to space out each actor's first eligible poll time.", ) args = parser.parse_args() diff --git a/tests/test_ingestion.py b/tests/test_ingestion.py index 87b956b..81c4084 100644 --- a/tests/test_ingestion.py +++ b/tests/test_ingestion.py @@ -9,6 +9,7 @@ from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncSta 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.srht_client import SourceHutClientError from srht_contrib.services.todo import TodoIngestionService, TodoPollResult from srht_contrib.services.types import BackfillBatchResult @@ -112,6 +113,25 @@ class BackfillingTodoService: return self.fetch_backfill_batch(actor, cursor_state) +class FailingTodoService: + service_name = "todo" + + def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult: + raise SourceHutClientError("temporary upstream failure") + + def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: + return BackfillBatchResult(events=[], cursor_state=None, complete=True) + + def fetch_recent_backfill_batch( + self, + actor: str, + cursor_state: dict | None = None, + *, + since: datetime, + ) -> BackfillBatchResult: + return BackfillBatchResult(events=[], cursor_state=None, complete=True) + + class QueueShrinkingTodoService: service_name = "todo" @@ -658,6 +678,50 @@ def test_poll_marks_backfill_complete_and_persists_service_state(db_session) -> assert all(state.status == "completed" for state in service_states) +def test_scheduled_poll_indexes_before_draining_recent_backfill(db_session) -> None: + settings = make_settings(DISCOVERY_BATCH_SIZE=1, INDEXED_ACTOR_REPOLL_SECONDS=3600) + poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService(), settings=settings) + now = datetime.now(tz=UTC) + db_session.add( + TrackedActor( + actor="~ccleberg", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=1), + next_poll_after=now - timedelta(minutes=1), + recent_backfill_status="pending", + ) + ) + db_session.commit() + + first_results = poller.poll_tracked_actors(db_session) + service_states_after_first_poll = db_session.scalars(select(ServiceBackfillState)).all() + tracked_after_first_poll = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg")) + + assert first_results == {"~ccleberg": 0} + assert service_states_after_first_poll == [] + assert tracked_after_first_poll is not None + assert tracked_after_first_poll.discovery_state == "indexed" + assert tracked_after_first_poll.recent_backfill_status == "pending" + assert tracked_after_first_poll.next_poll_after is not None + first_due_at = tracked_after_first_poll.next_poll_after + if first_due_at.tzinfo is None: + first_due_at = first_due_at.replace(tzinfo=UTC) + assert first_due_at <= datetime.now(tz=UTC) + + second_results = poller.poll_tracked_actors(db_session) + + tracked_after_second_poll = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg")) + service_states_after_second_poll = db_session.scalars( + select(ServiceBackfillState).order_by(ServiceBackfillState.scope, ServiceBackfillState.service) + ).all() + + assert second_results == {"~ccleberg": 1} + assert tracked_after_second_poll is not None + assert tracked_after_second_poll.recent_backfill_status == "completed" + assert [f"{state.scope}:{state.service}" for state in service_states_after_second_poll] == ["recent:git", "recent:todo"] + + def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None: settings = make_settings() poller = PollerService( @@ -779,11 +843,61 @@ def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None: 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["~a"].poll_attempts == 0 + assert actors["~b"].poll_attempts == 0 assert actors["~c"].poll_attempts == 0 +def test_scheduled_error_backoff_is_capped_and_success_resets_attempts(db_session) -> None: + settings = make_settings( + DISCOVERY_BATCH_SIZE=1, + DISCOVERY_ERROR_BACKOFF_SECONDS=3600, + DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=7200, + ) + now = datetime.now(tz=UTC) + db_session.add( + TrackedActor( + actor="~flaky", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=1), + next_poll_after=now - timedelta(minutes=1), + poll_attempts=10, + recent_backfill_status="completed", + ) + ) + db_session.commit() + failing_poller = PollerService(todo_service=FailingTodoService(), git_service=EmptyGitService(), settings=settings) + + failing_poller.poll_tracked_actors(db_session) + + failed_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~flaky")) + assert failed_actor is not None + assert failed_actor.discovery_state == "error" + assert failed_actor.poll_attempts == 11 + assert failed_actor.next_poll_after is not None + failed_next_poll_after = failed_actor.next_poll_after + if failed_next_poll_after.tzinfo is None: + failed_next_poll_after = failed_next_poll_after.replace(tzinfo=UTC) + assert failed_next_poll_after <= datetime.now(tz=UTC) + timedelta(seconds=7200, minutes=1) + + failed_actor.next_poll_after = datetime.now(tz=UTC) + db_session.add(failed_actor) + db_session.commit() + successful_poller = PollerService( + todo_service=RecordingTodoService(events_by_call=[[]]), + git_service=EmptyGitService(), + settings=settings, + ) + + successful_poller.poll_tracked_actors(db_session) + + recovered_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~flaky")) + assert recovered_actor is not None + assert recovered_actor.discovery_state == "indexed" + assert recovered_actor.poll_attempts == 0 + + def test_track_actor_request_prioritize_marks_actor_boosted_and_due_now(db_session) -> None: settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) -- cgit v1.2.3