Commit 6289a28296
Verified · cmc
Layout: unified · split
.env.example +6 −1
| @@ -5,7 +5,12 @@ TODO_SRHT_ENDPOINT=https://todo.sr.ht/query | |||
| 5 | GIT_SRHT_ENDPOINT=https://git.sr.ht/query | 5 | GIT_SRHT_ENDPOINT=https://git.sr.ht/query |
| 6 | DATABASE_URL=sqlite:///./srht_contrib.db | 6 | DATABASE_URL=sqlite:///./srht_contrib.db |
| 7 | DEFAULT_ACTOR=~your-user | 7 | DEFAULT_ACTOR=~your-user |
| 8 | POLL_INTERVAL_SECONDS=900 | 8 | POLL_INTERVAL_SECONDS=300 |
| 9 | DISCOVERY_BATCH_SIZE=20 | ||
| 10 | INDEXED_ACTOR_REPOLL_SECONDS=21600 | ||
| 11 | DISCOVERY_ERROR_BACKOFF_SECONDS=3600 | ||
| 12 | DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=21600 | ||
| 13 | SRHT_REQUEST_DELAY_SECONDS=0.5 | ||
| 9 | # Optional JSON object. Example: | 14 | # Optional JSON object. Example: |
| 10 | # {"~your-user":["you@example.com","Your Name"]} | 15 | # {"~your-user":["you@example.com","Your Name"]} |
| 11 | ACTOR_ALIASES_JSON={} | 16 | ACTOR_ALIASES_JSON={} |
API.md +1
| @@ -54,6 +54,7 @@ Background polling: | |||
| 54 | - The scheduler always seeds `DEFAULT_ACTOR` as a known actor. | 54 | - The scheduler always seeds `DEFAULT_ACTOR` as a known actor. |
| 55 | - Public contribution reads register additional actors for later background polling. | 55 | - Public contribution reads register additional actors for later background polling. |
| 56 | - The scheduler only processes due actors, up to `DISCOVERY_BATCH_SIZE` per pass. | 56 | - The scheduler only processes due actors, up to `DISCOVERY_BATCH_SIZE` per pass. |
| 57 | - Scheduled first indexing skips bounded one-year backfill work, then drains backfill in later scheduled passes so newly requested actors become indexed sooner. | ||
| 57 | - Manual polling remains available through `POST /api/contributions/poll`. | 58 | - Manual polling remains available through `POST /api/contributions/poll`. |
| 58 | 59 | ||
| 59 | Repository names: | 60 | Repository names: |
README.md +7 −5
| @@ -86,6 +86,7 @@ Environment variables: | |||
| 86 | - `DISCOVERY_BATCH_SIZE`: max number of due actors to process per scheduler pass | 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 | 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 | 88 | - `DISCOVERY_ERROR_BACKOFF_SECONDS`: base retry delay after a failed scheduled poll |
| 89 | - `DISCOVERY_ERROR_BACKOFF_MAX_SECONDS`: maximum retry delay after repeated scheduled poll failures | ||
| 89 | - `ACTOR_ALIASES_JSON`: optional JSON object for actor/email/display-name alias mapping | 90 | - `ACTOR_ALIASES_JSON`: optional JSON object for actor/email/display-name alias mapping |
| 90 | - `GIT_TRACKED_REPOSITORIES`: optional JSON array of repository names or `owner/repo` strings to union into git polling | 91 | - `GIT_TRACKED_REPOSITORIES`: optional JSON array of repository names or `owner/repo` strings to union into git polling |
| 91 | 92 | ||
| @@ -99,10 +100,11 @@ TODO_SRHT_ENDPOINT=https://todo.sr.ht/query | |||
| 99 | GIT_SRHT_ENDPOINT=https://git.sr.ht/query | 100 | GIT_SRHT_ENDPOINT=https://git.sr.ht/query |
| 100 | DATABASE_URL=sqlite:///./srht_contrib.db | 101 | DATABASE_URL=sqlite:///./srht_contrib.db |
| 101 | DEFAULT_ACTOR=~your-user | 102 | DEFAULT_ACTOR=~your-user |
| 102 | POLL_INTERVAL_SECONDS=900 | 103 | POLL_INTERVAL_SECONDS=300 |
| 103 | DISCOVERY_BATCH_SIZE=5 | 104 | DISCOVERY_BATCH_SIZE=20 |
| 104 | INDEXED_ACTOR_REPOLL_SECONDS=21600 | 105 | INDEXED_ACTOR_REPOLL_SECONDS=21600 |
| 105 | DISCOVERY_ERROR_BACKOFF_SECONDS=3600 | 106 | DISCOVERY_ERROR_BACKOFF_SECONDS=3600 |
| 107 | DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=21600 | ||
| 106 | ACTOR_ALIASES_JSON={"~your-user":["you@example.com","Your Name"]} | 108 | ACTOR_ALIASES_JSON={"~your-user":["you@example.com","Your Name"]} |
| 107 | GIT_TRACKED_REPOSITORIES=["your-repo","~your-user/your-site"] | 109 | GIT_TRACKED_REPOSITORIES=["your-repo","~your-user/your-site"] |
| 108 | ``` | 110 | ``` |
| @@ -173,7 +175,7 @@ Example response: | |||
| 173 | 175 | ||
| 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. | 176 | 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. |
| 175 | 177 | ||
| 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`. | 178 | 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`. |
| 177 | 179 | ||
| 178 | 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. | 180 | 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. |
| 179 | 181 | ||
| @@ -182,10 +184,10 @@ Clients can explicitly signal that a public contribution read is for the signed- | |||
| 182 | To durably queue a large username list without polling it immediately: | 184 | To durably queue a large username list without polling it immediately: |
| 183 | 185 | ||
| 184 | ```bash | 186 | ```bash |
| 185 | srht-enqueue-actors srht_usernames.txt --stagger-seconds 300 | 187 | srht-enqueue-actors srht_usernames.txt --stagger-seconds 60 |
| 186 | ``` | 188 | ``` |
| 187 | 189 | ||
| 188 | 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. | 190 | 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. |
| 189 | 191 | ||
| 190 | 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: | 192 | 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: |
| 191 | 193 | ||
compose.yml +3 −2
| @@ -13,10 +13,11 @@ services: | |||
| 13 | GIT_SRHT_ENDPOINT: ${GIT_SRHT_ENDPOINT:-https://git.sr.ht/query} | 13 | GIT_SRHT_ENDPOINT: ${GIT_SRHT_ENDPOINT:-https://git.sr.ht/query} |
| 14 | DATABASE_URL: sqlite:////data/srht_contrib.db | 14 | DATABASE_URL: sqlite:////data/srht_contrib.db |
| 15 | DEFAULT_ACTOR: ${DEFAULT_ACTOR} | 15 | DEFAULT_ACTOR: ${DEFAULT_ACTOR} |
| 16 | POLL_INTERVAL_SECONDS: ${POLL_INTERVAL_SECONDS:-900} | 16 | POLL_INTERVAL_SECONDS: ${POLL_INTERVAL_SECONDS:-300} |
| 17 | DISCOVERY_BATCH_SIZE: ${DISCOVERY_BATCH_SIZE:-5} | 17 | DISCOVERY_BATCH_SIZE: ${DISCOVERY_BATCH_SIZE:-20} |
| 18 | INDEXED_ACTOR_REPOLL_SECONDS: ${INDEXED_ACTOR_REPOLL_SECONDS:-21600} | 18 | INDEXED_ACTOR_REPOLL_SECONDS: ${INDEXED_ACTOR_REPOLL_SECONDS:-21600} |
| 19 | DISCOVERY_ERROR_BACKOFF_SECONDS: ${DISCOVERY_ERROR_BACKOFF_SECONDS:-3600} | 19 | DISCOVERY_ERROR_BACKOFF_SECONDS: ${DISCOVERY_ERROR_BACKOFF_SECONDS:-3600} |
| 20 | DISCOVERY_ERROR_BACKOFF_MAX_SECONDS: ${DISCOVERY_ERROR_BACKOFF_MAX_SECONDS:-21600} | ||
| 20 | SRHT_REQUEST_DELAY_SECONDS: ${SRHT_REQUEST_DELAY_SECONDS:-0.5} | 21 | SRHT_REQUEST_DELAY_SECONDS: ${SRHT_REQUEST_DELAY_SECONDS:-0.5} |
| 21 | ACTOR_ALIASES_JSON: ${ACTOR_ALIASES_JSON:-{}} | 22 | ACTOR_ALIASES_JSON: ${ACTOR_ALIASES_JSON:-{}} |
| 22 | GIT_TRACKED_REPOSITORIES: ${GIT_TRACKED_REPOSITORIES:-[]} | 23 | GIT_TRACKED_REPOSITORIES: ${GIT_TRACKED_REPOSITORIES:-[]} |
src/srht_contrib/config.py +3 −2
| @@ -34,13 +34,14 @@ class Settings(BaseSettings): | |||
| 34 | alias="DATABASE_URL", | 34 | alias="DATABASE_URL", |
| 35 | ) | 35 | ) |
| 36 | default_actor: str = Field(default="~unknown", alias="DEFAULT_ACTOR") | 36 | default_actor: str = Field(default="~unknown", alias="DEFAULT_ACTOR") |
| 37 | poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS") | 37 | poll_interval_seconds: int = Field(default=300, 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 | sqlite_busy_timeout_seconds: float = Field(default=30.0, alias="SQLITE_BUSY_TIMEOUT_SECONDS") | 40 | sqlite_busy_timeout_seconds: float = Field(default=30.0, alias="SQLITE_BUSY_TIMEOUT_SECONDS") |
| 41 | discovery_batch_size: int = Field(default=5, alias="DISCOVERY_BATCH_SIZE") | 41 | discovery_batch_size: int = Field(default=20, alias="DISCOVERY_BATCH_SIZE") |
| 42 | indexed_actor_repoll_seconds: int = Field(default=21600, alias="INDEXED_ACTOR_REPOLL_SECONDS") | 42 | indexed_actor_repoll_seconds: int = Field(default=21600, alias="INDEXED_ACTOR_REPOLL_SECONDS") |
| 43 | discovery_error_backoff_seconds: int = Field(default=3600, alias="DISCOVERY_ERROR_BACKOFF_SECONDS") | 43 | discovery_error_backoff_seconds: int = Field(default=3600, alias="DISCOVERY_ERROR_BACKOFF_SECONDS") |
| 44 | discovery_error_backoff_max_seconds: int = Field(default=21600, alias="DISCOVERY_ERROR_BACKOFF_MAX_SECONDS") | ||
| 44 | git_repo_discovery_ttl_seconds: int = Field(default=3600, alias="GIT_REPO_DISCOVERY_TTL_SECONDS") | 45 | git_repo_discovery_ttl_seconds: int = Field(default=3600, alias="GIT_REPO_DISCOVERY_TTL_SECONDS") |
| 45 | actor_aliases_json: dict[str, list[str]] = Field( | 46 | actor_aliases_json: dict[str, list[str]] = Field( |
| 46 | default_factory=dict, | 47 | default_factory=dict, |
src/srht_contrib/jobs/poller.py +49 −7
| @@ -34,7 +34,7 @@ class PollerService: | |||
| 34 | self.settings = settings | 34 | self.settings = settings |
| 35 | self._sync_overlap = timedelta(hours=settings.sync_overlap_hours) | 35 | self._sync_overlap = timedelta(hours=settings.sync_overlap_hours) |
| 36 | 36 | ||
| 37 | def poll_all(self, db: Session, actor: str) -> int: | 37 | def poll_all(self, db: Session, actor: str, *, run_backfill: bool = True) -> int: |
| 38 | self.track_actor_request(db, actor, update_last_requested=False) | 38 | self.track_actor_request(db, actor, update_last_requested=False) |
| 39 | try: | 39 | try: |
| 40 | inserted = self._poll_actor(db, actor) | 40 | inserted = self._poll_actor(db, actor) |
| @@ -45,7 +45,8 @@ class PollerService: | |||
| 45 | raise | 45 | raise |
| 46 | 46 | ||
| 47 | self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None) | 47 | self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None) |
| 48 | inserted += self._run_backfill_batches(db, actor) | 48 | if run_backfill: |
| 49 | inserted += self._run_backfill_batches(db, actor) | ||
| 49 | db.commit() | 50 | db.commit() |
| 50 | return inserted | 51 | return inserted |
| 51 | 52 | ||
| @@ -56,7 +57,7 @@ class PollerService: | |||
| 56 | 57 | ||
| 57 | results: dict[str, int] = {} | 58 | results: dict[str, int] = {} |
| 58 | actors = db.scalars( | 59 | actors = db.scalars( |
| 59 | select(TrackedActor.actor) | 60 | select(TrackedActor) |
| 60 | .where(TrackedActor.is_active.is_(True)) | 61 | .where(TrackedActor.is_active.is_(True)) |
| 61 | .where( | 62 | .where( |
| 62 | or_( | 63 | or_( |
| @@ -74,7 +75,12 @@ class PollerService: | |||
| 74 | ) | 75 | ) |
| 75 | .limit(self.settings.discovery_batch_size) | 76 | .limit(self.settings.discovery_batch_size) |
| 76 | ).all() | 77 | ).all() |
| 77 | for actor in actors: | 78 | for tracked_actor in actors: |
| 79 | actor = tracked_actor.actor | ||
| 80 | run_backfill_only = ( | ||
| 81 | tracked_actor.last_polled_at is not None | ||
| 82 | and tracked_actor.recent_backfill_status != "completed" | ||
| 83 | ) | ||
| 78 | claimed_actor = self.track_actor_request(db, actor, update_last_requested=False) | 84 | claimed_actor = self.track_actor_request(db, actor, update_last_requested=False) |
| 79 | claimed_actor.discovery_state = "in_progress" | 85 | claimed_actor.discovery_state = "in_progress" |
| 80 | claimed_actor.last_claimed_at = datetime.now(tz=UTC) | 86 | claimed_actor.last_claimed_at = datetime.now(tz=UTC) |
| @@ -82,7 +88,12 @@ class PollerService: | |||
| 82 | db.add(claimed_actor) | 88 | db.add(claimed_actor) |
| 83 | db.commit() | 89 | db.commit() |
| 84 | try: | 90 | try: |
| 85 | results[actor] = self.poll_all(db, actor) | 91 | if run_backfill_only: |
| 92 | results[actor] = self.poll_recent_backfill(db, actor) | ||
| 93 | else: | ||
| 94 | results[actor] = self.poll_all(db, actor, run_backfill=False) | ||
| 95 | self._schedule_pending_backfill_or_repoll(db, actor) | ||
| 96 | db.commit() | ||
| 86 | except SourceHutClientError: | 97 | except SourceHutClientError: |
| 87 | logger.exception("Scheduled poll failed for actor=%s", actor) | 98 | logger.exception("Scheduled poll failed for actor=%s", actor) |
| 88 | except Exception: | 99 | except Exception: |
| @@ -220,15 +231,46 @@ class PollerService: | |||
| 220 | tracked_actor.last_polled_at = now | 231 | tracked_actor.last_polled_at = now |
| 221 | tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) | 232 | tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) |
| 222 | tracked_actor.priority_boosted_at = None | 233 | tracked_actor.priority_boosted_at = None |
| 234 | tracked_actor.poll_attempts = 0 | ||
| 223 | elif status == "error": | 235 | elif status == "error": |
| 224 | tracked_actor.discovery_state = "error" | 236 | tracked_actor.discovery_state = "error" |
| 225 | tracked_actor.next_poll_after = now + timedelta( | 237 | backoff_seconds = min( |
| 226 | seconds=self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1) | 238 | self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1), |
| 239 | self.settings.discovery_error_backoff_max_seconds, | ||
| 227 | ) | 240 | ) |
| 241 | tracked_actor.next_poll_after = now + timedelta(seconds=backoff_seconds) | ||
| 228 | tracked_actor.priority_boosted_at = None | 242 | tracked_actor.priority_boosted_at = None |
| 229 | db.add(tracked_actor) | 243 | db.add(tracked_actor) |
| 230 | db.flush() | 244 | db.flush() |
| 231 | 245 | ||
| 246 | def poll_recent_backfill(self, db: Session, actor: str) -> int: | ||
| 247 | try: | ||
| 248 | inserted = self._run_backfill_batches(db, actor) | ||
| 249 | except Exception as exc: | ||
| 250 | db.rollback() | ||
| 251 | self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc)) | ||
| 252 | db.commit() | ||
| 253 | raise | ||
| 254 | |||
| 255 | self._schedule_pending_backfill_or_repoll(db, actor) | ||
| 256 | db.commit() | ||
| 257 | return inserted | ||
| 258 | |||
| 259 | def _schedule_pending_backfill_or_repoll(self, db: Session, actor: str) -> None: | ||
| 260 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) | ||
| 261 | now = datetime.now(tz=UTC) | ||
| 262 | tracked_actor.discovery_state = "indexed" | ||
| 263 | tracked_actor.last_poll_status = "indexed" | ||
| 264 | tracked_actor.last_poll_error = None | ||
| 265 | tracked_actor.priority_boosted_at = None | ||
| 266 | tracked_actor.poll_attempts = 0 | ||
| 267 | if tracked_actor.recent_backfill_status == "completed": | ||
| 268 | tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) | ||
| 269 | else: | ||
| 270 | tracked_actor.next_poll_after = now | ||
| 271 | db.add(tracked_actor) | ||
| 272 | db.flush() | ||
| 273 | |||
| 232 | def _run_backfill_batches(self, db: Session, actor: str) -> int: | 274 | def _run_backfill_batches(self, db: Session, actor: str) -> int: |
| 233 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) | 275 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) |
| 234 | if tracked_actor.recent_backfill_status == "completed": | 276 | if tracked_actor.recent_backfill_status == "completed": |
src/srht_contrib/scripts/enqueue_actors.py +2 −2
| @@ -40,7 +40,7 @@ def _iter_usernames(path: Path) -> list[str]: | |||
| 40 | return usernames | 40 | return usernames |
| 41 | 41 | ||
| 42 | 42 | ||
| 43 | def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: datetime | None = None) -> int: | 43 | def enqueue_actors(username_file: Path, *, stagger_seconds: int = 60, start_at: datetime | None = None) -> int: |
| 44 | settings = Settings() | 44 | settings = Settings() |
| 45 | session_factory = make_session_factory(settings) | 45 | session_factory = make_session_factory(settings) |
| 46 | usernames = _iter_usernames(username_file) | 46 | usernames = _iter_usernames(username_file) |
| @@ -95,7 +95,7 @@ def main() -> None: | |||
| 95 | parser.add_argument( | 95 | parser.add_argument( |
| 96 | "--stagger-seconds", | 96 | "--stagger-seconds", |
| 97 | type=int, | 97 | type=int, |
| 98 | default=300, | 98 | default=60, |
| 99 | help="Seconds to space out each actor's first eligible poll time.", | 99 | help="Seconds to space out each actor's first eligible poll time.", |
| 100 | ) | 100 | ) |
| 101 | args = parser.parse_args() | 101 | args = parser.parse_args() |
tests/test_ingestion.py +116 −2
| @@ -9,6 +9,7 @@ from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncSta | |||
| 9 | from srht_contrib.scripts.enqueue_actors import enqueue_actors | 9 | from srht_contrib.scripts.enqueue_actors import enqueue_actors |
| 10 | from srht_contrib.schemas import NormalizedEvent | 10 | from srht_contrib.schemas import NormalizedEvent |
| 11 | from srht_contrib.services.git import GitIngestionService, GitPollResult | 11 | from srht_contrib.services.git import GitIngestionService, GitPollResult |
| 12 | from srht_contrib.services.srht_client import SourceHutClientError | ||
| 12 | from srht_contrib.services.todo import TodoIngestionService, TodoPollResult | 13 | from srht_contrib.services.todo import TodoIngestionService, TodoPollResult |
| 13 | from srht_contrib.services.types import BackfillBatchResult | 14 | from srht_contrib.services.types import BackfillBatchResult |
| 14 | 15 | ||
| @@ -112,6 +113,25 @@ class BackfillingTodoService: | |||
| 112 | return self.fetch_backfill_batch(actor, cursor_state) | 113 | return self.fetch_backfill_batch(actor, cursor_state) |
| 113 | 114 | ||
| 114 | 115 | ||
| 116 | class FailingTodoService: | ||
| 117 | service_name = "todo" | ||
| 118 | |||
| 119 | def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult: | ||
| 120 | raise SourceHutClientError("temporary upstream failure") | ||
| 121 | |||
| 122 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: | ||
| 123 | return BackfillBatchResult(events=[], cursor_state=None, complete=True) | ||
| 124 | |||
| 125 | def fetch_recent_backfill_batch( | ||
| 126 | self, | ||
| 127 | actor: str, | ||
| 128 | cursor_state: dict | None = None, | ||
| 129 | *, | ||
| 130 | since: datetime, | ||
| 131 | ) -> BackfillBatchResult: | ||
| 132 | return BackfillBatchResult(events=[], cursor_state=None, complete=True) | ||
| 133 | |||
| 134 | |||
| 115 | class QueueShrinkingTodoService: | 135 | class QueueShrinkingTodoService: |
| 116 | service_name = "todo" | 136 | service_name = "todo" |
| 117 | 137 | ||
| @@ -658,6 +678,50 @@ def test_poll_marks_backfill_complete_and_persists_service_state(db_session) -> | |||
| 658 | assert all(state.status == "completed" for state in service_states) | 678 | assert all(state.status == "completed" for state in service_states) |
| 659 | 679 | ||
| 660 | 680 | ||
| 681 | def test_scheduled_poll_indexes_before_draining_recent_backfill(db_session) -> None: | ||
| 682 | settings = make_settings(DISCOVERY_BATCH_SIZE=1, INDEXED_ACTOR_REPOLL_SECONDS=3600) | ||
| 683 | poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService(), settings=settings) | ||
| 684 | now = datetime.now(tz=UTC) | ||
| 685 | db_session.add( | ||
| 686 | TrackedActor( | ||
| 687 | actor="~ccleberg", | ||
| 688 | is_active=True, | ||
| 689 | discovery_state="queued", | ||
| 690 | queued_for_discovery_at=now - timedelta(minutes=1), | ||
| 691 | next_poll_after=now - timedelta(minutes=1), | ||
| 692 | recent_backfill_status="pending", | ||
| 693 | ) | ||
| 694 | ) | ||
| 695 | db_session.commit() | ||
| 696 | |||
| 697 | first_results = poller.poll_tracked_actors(db_session) | ||
| 698 | service_states_after_first_poll = db_session.scalars(select(ServiceBackfillState)).all() | ||
| 699 | tracked_after_first_poll = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg")) | ||
| 700 | |||
| 701 | assert first_results == {"~ccleberg": 0} | ||
| 702 | assert service_states_after_first_poll == [] | ||
| 703 | assert tracked_after_first_poll is not None | ||
| 704 | assert tracked_after_first_poll.discovery_state == "indexed" | ||
| 705 | assert tracked_after_first_poll.recent_backfill_status == "pending" | ||
| 706 | assert tracked_after_first_poll.next_poll_after is not None | ||
| 707 | first_due_at = tracked_after_first_poll.next_poll_after | ||
| 708 | if first_due_at.tzinfo is None: | ||
| 709 | first_due_at = first_due_at.replace(tzinfo=UTC) | ||
| 710 | assert first_due_at <= datetime.now(tz=UTC) | ||
| 711 | |||
| 712 | second_results = poller.poll_tracked_actors(db_session) | ||
| 713 | |||
| 714 | tracked_after_second_poll = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg")) | ||
| 715 | service_states_after_second_poll = db_session.scalars( | ||
| 716 | select(ServiceBackfillState).order_by(ServiceBackfillState.scope, ServiceBackfillState.service) | ||
| 717 | ).all() | ||
| 718 | |||
| 719 | assert second_results == {"~ccleberg": 1} | ||
| 720 | assert tracked_after_second_poll is not None | ||
| 721 | assert tracked_after_second_poll.recent_backfill_status == "completed" | ||
| 722 | assert [f"{state.scope}:{state.service}" for state in service_states_after_second_poll] == ["recent:git", "recent:todo"] | ||
| 723 | |||
| 724 | |||
| 661 | def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None: | 725 | def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None: |
| 662 | settings = make_settings() | 726 | settings = make_settings() |
| 663 | poller = PollerService( | 727 | poller = PollerService( |
| @@ -779,11 +843,61 @@ def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None: | |||
| 779 | assert actors["~a"].discovery_state == "indexed" | 843 | assert actors["~a"].discovery_state == "indexed" |
| 780 | assert actors["~b"].discovery_state == "indexed" | 844 | assert actors["~b"].discovery_state == "indexed" |
| 781 | assert actors["~c"].discovery_state == "queued" | 845 | assert actors["~c"].discovery_state == "queued" |
| 782 | assert actors["~a"].poll_attempts == 1 | 846 | assert actors["~a"].poll_attempts == 0 |
| 783 | assert actors["~b"].poll_attempts == 1 | 847 | assert actors["~b"].poll_attempts == 0 |
| 784 | assert actors["~c"].poll_attempts == 0 | 848 | assert actors["~c"].poll_attempts == 0 |
| 785 | 849 | ||
| 786 | 850 | ||
| 851 | def test_scheduled_error_backoff_is_capped_and_success_resets_attempts(db_session) -> None: | ||
| 852 | settings = make_settings( | ||
| 853 | DISCOVERY_BATCH_SIZE=1, | ||
| 854 | DISCOVERY_ERROR_BACKOFF_SECONDS=3600, | ||
| 855 | DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=7200, | ||
| 856 | ) | ||
| 857 | now = datetime.now(tz=UTC) | ||
| 858 | db_session.add( | ||
| 859 | TrackedActor( | ||
| 860 | actor="~flaky", | ||
| 861 | is_active=True, | ||
| 862 | discovery_state="queued", | ||
| 863 | queued_for_discovery_at=now - timedelta(minutes=1), | ||
| 864 | next_poll_after=now - timedelta(minutes=1), | ||
| 865 | poll_attempts=10, | ||
| 866 | recent_backfill_status="completed", | ||
| 867 | ) | ||
| 868 | ) | ||
| 869 | db_session.commit() | ||
| 870 | failing_poller = PollerService(todo_service=FailingTodoService(), git_service=EmptyGitService(), settings=settings) | ||
| 871 | |||
| 872 | failing_poller.poll_tracked_actors(db_session) | ||
| 873 | |||
| 874 | failed_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~flaky")) | ||
| 875 | assert failed_actor is not None | ||
| 876 | assert failed_actor.discovery_state == "error" | ||
| 877 | assert failed_actor.poll_attempts == 11 | ||
| 878 | assert failed_actor.next_poll_after is not None | ||
| 879 | failed_next_poll_after = failed_actor.next_poll_after | ||
| 880 | if failed_next_poll_after.tzinfo is None: | ||
| 881 | failed_next_poll_after = failed_next_poll_after.replace(tzinfo=UTC) | ||
| 882 | assert failed_next_poll_after <= datetime.now(tz=UTC) + timedelta(seconds=7200, minutes=1) | ||
| 883 | |||
| 884 | failed_actor.next_poll_after = datetime.now(tz=UTC) | ||
| 885 | db_session.add(failed_actor) | ||
| 886 | db_session.commit() | ||
| 887 | successful_poller = PollerService( | ||
| 888 | todo_service=RecordingTodoService(events_by_call=[[]]), | ||
| 889 | git_service=EmptyGitService(), | ||
| 890 | settings=settings, | ||
| 891 | ) | ||
| 892 | |||
| 893 | successful_poller.poll_tracked_actors(db_session) | ||
| 894 | |||
| 895 | recovered_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~flaky")) | ||
| 896 | assert recovered_actor is not None | ||
| 897 | assert recovered_actor.discovery_state == "indexed" | ||
| 898 | assert recovered_actor.poll_attempts == 0 | ||
| 899 | |||
| 900 | |||
| 787 | def test_track_actor_request_prioritize_marks_actor_boosted_and_due_now(db_session) -> None: | 901 | def test_track_actor_request_prioritize_marks_actor_boosted_and_due_now(db_session) -> None: |
| 788 | settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) | 902 | settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) |
| 789 | poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) | 903 | poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) |