Commit 90e61ed3e8
Verified · cmc
Layout: unified · split
API.md +13 −21
| @@ -140,9 +140,6 @@ Response `200 OK`: | |||
| 140 | "is_recent_window_backfilled": false, | 140 | "is_recent_window_backfilled": false, |
| 141 | "recent_backfill_state": "in_progress", | 141 | "recent_backfill_state": "in_progress", |
| 142 | "recent_backfill_completed_at": null, | 142 | "recent_backfill_completed_at": null, |
| 143 | "is_backfilled": false, | ||
| 144 | "backfill_state": "in_progress", | ||
| 145 | "backfill_completed_at": null, | ||
| 146 | "days": [ | 143 | "days": [ |
| 147 | { "date": "2026-03-01", "count": 0, "score": 0.0 }, | 144 | { "date": "2026-03-01", "count": 0, "score": 0.0 }, |
| 148 | { "date": "2026-03-02", "count": 3, "score": 2.5 } | 145 | { "date": "2026-03-02", "count": 3, "score": 2.5 } |
| @@ -158,12 +155,9 @@ Response fields: | |||
| 158 | - `is_indexed` boolean: whether the service has already completed at least one successful recent/incremental poll for this actor | 155 | - `is_indexed` boolean: whether the service has already completed at least one successful recent/incremental poll for this actor |
| 159 | - `last_polled_at` string or `null`: most recent successful poll time, if any | 156 | - `last_polled_at` string or `null`: most recent successful poll time, if any |
| 160 | - `indexing_state` string: one of `pending`, `indexed`, or `error` | 157 | - `indexing_state` string: one of `pending`, `indexed`, or `error` |
| 161 | - `is_recent_window_backfilled` boolean: whether the prioritized recent history window has completed backfill | 158 | - `is_recent_window_backfilled` boolean: whether the service has finished filling the retained one-year history window |
| 162 | - `recent_backfill_state` string: one of `pending`, `in_progress`, `completed`, or `error` | 159 | - `recent_backfill_state` string: one of `pending`, `in_progress`, `completed`, or `error` |
| 163 | - `recent_backfill_completed_at` string or `null`: when recent-window backfill completed, if it has | 160 | - `recent_backfill_completed_at` string or `null`: when one-year backfill completed, if it has |
| 164 | - `is_backfilled` boolean: whether historical backfill has completed for this actor | ||
| 165 | - `backfill_state` string: one of `pending`, `in_progress`, `completed`, or `error` | ||
| 166 | - `backfill_completed_at` string or `null`: when full historical backfill completed, if it has | ||
| 167 | - `days` array: | 161 | - `days` array: |
| 168 | - `date` string `YYYY-MM-DD` | 162 | - `date` string `YYYY-MM-DD` |
| 169 | - `count` integer contribution count for the day | 163 | - `count` integer contribution count for the day |
| @@ -175,15 +169,19 @@ Indexing state semantics: | |||
| 175 | - `indexed`: at least one successful poll has completed for the actor | 169 | - `indexed`: at least one successful poll has completed for the actor |
| 176 | - `error`: the most recent poll attempt for the actor failed | 170 | - `error`: the most recent poll attempt for the actor failed |
| 177 | 171 | ||
| 178 | Backfill state semantics: | 172 | Recent backfill semantics: |
| 179 | 173 | ||
| 180 | - recent-window fields: | 174 | - the service only retains and backfills the most recent 365 days of activity |
| 181 | - represent the prioritized visible-history window for client UX | 175 | - `pending`: the actor has not started one-year backfill yet |
| 182 | - `pending`: the actor has not started historical backfill yet | 176 | - `in_progress`: one-year backfill is actively progressing in bounded background batches |
| 183 | - `in_progress`: historical backfill is actively progressing in bounded background batches | 177 | - `completed`: the retained one-year window is fully backfilled |
| 184 | - `completed`: historical backfill has completed for all supported services | ||
| 185 | - `error`: the most recent backfill attempt failed | 178 | - `error`: the most recent backfill attempt failed |
| 186 | 179 | ||
| 180 | Retention notes: | ||
| 181 | |||
| 182 | - activity older than 365 days is not retained | ||
| 183 | - scheduled polling periodically prunes contribution rows older than the retained window | ||
| 184 | |||
| 187 | Possible errors: | 185 | Possible errors: |
| 188 | 186 | ||
| 189 | - `400 Bad Request` for invalid or conflicting date input | 187 | - `400 Bad Request` for invalid or conflicting date input |
| @@ -218,7 +216,7 @@ Behavior notes: | |||
| 218 | 216 | ||
| 219 | - This endpoint has the same actor-registration and alias-resolution behavior as the calendar endpoint. | 217 | - This endpoint has the same actor-registration and alias-resolution behavior as the calendar endpoint. |
| 220 | - This endpoint returns immediately and does not block on SourceHut polling. | 218 | - This endpoint returns immediately and does not block on SourceHut polling. |
| 221 | - This endpoint also reflects historical backfill state so clients can distinguish recent indexing from complete history. | 219 | - This endpoint also reflects whether the retained one-year history window has been fully backfilled yet. |
| 222 | 220 | ||
| 223 | Example: | 221 | Example: |
| 224 | 222 | ||
| @@ -239,9 +237,6 @@ Response `200 OK`: | |||
| 239 | "is_recent_window_backfilled": true, | 237 | "is_recent_window_backfilled": true, |
| 240 | "recent_backfill_state": "completed", | 238 | "recent_backfill_state": "completed", |
| 241 | "recent_backfill_completed_at": "2026-04-11T18:02:00Z", | 239 | "recent_backfill_completed_at": "2026-04-11T18:02:00Z", |
| 242 | "is_backfilled": false, | ||
| 243 | "backfill_state": "in_progress", | ||
| 244 | "backfill_completed_at": null, | ||
| 245 | "total_events": 126, | 240 | "total_events": 126, |
| 246 | "total_score": 116.75, | 241 | "total_score": 116.75, |
| 247 | "active_days": 14, | 242 | "active_days": 14, |
| @@ -261,9 +256,6 @@ Response fields: | |||
| 261 | - `is_recent_window_backfilled` boolean | 256 | - `is_recent_window_backfilled` boolean |
| 262 | - `recent_backfill_state` string | 257 | - `recent_backfill_state` string |
| 263 | - `recent_backfill_completed_at` string or `null` | 258 | - `recent_backfill_completed_at` string or `null` |
| 264 | - `is_backfilled` boolean | ||
| 265 | - `backfill_state` string | ||
| 266 | - `backfill_completed_at` string or `null` | ||
| 267 | - `total_events` integer | 259 | - `total_events` integer |
| 268 | - `total_score` float | 260 | - `total_score` float |
| 269 | - `active_days` integer | 261 | - `active_days` integer |
README.md +3 −8
| @@ -211,9 +211,6 @@ Example response: | |||
| 211 | "is_recent_window_backfilled": true, | 211 | "is_recent_window_backfilled": true, |
| 212 | "recent_backfill_state": "completed", | 212 | "recent_backfill_state": "completed", |
| 213 | "recent_backfill_completed_at": "2026-04-11T18:02:00Z", | 213 | "recent_backfill_completed_at": "2026-04-11T18:02:00Z", |
| 214 | "is_backfilled": false, | ||
| 215 | "backfill_state": "in_progress", | ||
| 216 | "backfill_completed_at": null, | ||
| 217 | "days": [ | 214 | "days": [ |
| 218 | {"date": "2026-03-28", "count": 3, "score": 3.5}, | 215 | {"date": "2026-03-28", "count": 3, "score": 3.5}, |
| 219 | {"date": "2026-03-29", "count": 0, "score": 0.0}, | 216 | {"date": "2026-03-29", "count": 0, "score": 0.0}, |
| @@ -241,9 +238,6 @@ Example response: | |||
| 241 | "is_recent_window_backfilled": true, | 238 | "is_recent_window_backfilled": true, |
| 242 | "recent_backfill_state": "completed", | 239 | "recent_backfill_state": "completed", |
| 243 | "recent_backfill_completed_at": "2026-04-11T18:02:00Z", | 240 | "recent_backfill_completed_at": "2026-04-11T18:02:00Z", |
| 244 | "is_backfilled": false, | ||
| 245 | "backfill_state": "in_progress", | ||
| 246 | "backfill_completed_at": null, | ||
| 247 | "total_events": 42, | 241 | "total_events": 42, |
| 248 | "total_score": 37.5, | 242 | "total_score": 37.5, |
| 249 | "active_days": 18, | 243 | "active_days": 18, |
| @@ -326,14 +320,15 @@ The SourceHut-specific assumptions are isolated to the service modules: | |||
| 326 | 320 | ||
| 327 | - `src/srht_contrib/services/todo.py` uses the authenticated `events(cursor)` feed first, then falls back to tracker/ticket event traversal for reliable contribution discovery. | 321 | - `src/srht_contrib/services/todo.py` uses the authenticated `events(cursor)` feed first, then falls back to tracker/ticket event traversal for reliable contribution discovery. |
| 328 | - `src/srht_contrib/services/git.py` discovers owned repositories for an actor, polls each repository `log(cursor)`, and attributes commits through the configured alias map. | 322 | - `src/srht_contrib/services/git.py` discovers owned repositories for an actor, polls each repository `log(cursor)`, and attributes commits through the configured alias map. |
| 323 | - the service only retains the most recent 365 days of contribution history and periodically prunes older rows | ||
| 329 | 324 | ||
| 330 | ## Known Limitations | 325 | ## Known Limitations |
| 331 | 326 | ||
| 332 | - `git.sr.ht` polling assumes the actor's repositories are discoverable through the SourceHut GraphQL API | 327 | - `git.sr.ht` polling assumes the actor's repositories are discoverable through the SourceHut GraphQL API |
| 333 | - scheduled polling runs in-process, so it is not a distributed scheduler | 328 | - scheduled polling runs in-process, so it is not a distributed scheduler |
| 334 | - newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs | 329 | - newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs |
| 335 | - the recent visible-history window is prioritized first, but deep-history backfill can still take many scheduler passes for active users | 330 | - the service is intentionally limited to a rolling one-year history window; older activity is not retained |
| 336 | - full historical backfill can take many scheduler passes for active users because it runs in bounded batches | 331 | - one-year backfill still runs in bounded batches, so a newly requested actor may take multiple scheduler passes before their visible graph is fully filled in |
| 337 | - alias management is config-driven; there is no alias CRUD API yet | 332 | - alias management is config-driven; there is no alias CRUD API yet |
| 338 | - current deployment model is trusted-operator V1, not a public multi-tenant service | 333 | - current deployment model is trusted-operator V1, not a public multi-tenant service |
| 339 | 334 | ||
src/srht_contrib/jobs/poller.py +13 −39
| @@ -13,14 +13,13 @@ from srht_contrib.schemas import NormalizedEvent | |||
| 13 | from srht_contrib.services.git import GitIngestionService | 13 | from srht_contrib.services.git import GitIngestionService |
| 14 | from srht_contrib.services.srht_client import SourceHutClientError | 14 | from srht_contrib.services.srht_client import SourceHutClientError |
| 15 | from srht_contrib.services.todo import TodoIngestionService | 15 | from srht_contrib.services.todo import TodoIngestionService |
| 16 | from srht_contrib.utils.retention import RETENTION_DAYS, prune_contribution_events | ||
| 16 | from srht_contrib.utils.repositories import canonicalize_repository_name | 17 | from srht_contrib.utils.repositories import canonicalize_repository_name |
| 17 | 18 | ||
| 18 | 19 | ||
| 19 | logger = logging.getLogger(__name__) | 20 | logger = logging.getLogger(__name__) |
| 20 | SYNC_OVERLAP = timedelta(hours=24) | 21 | SYNC_OVERLAP = timedelta(hours=24) |
| 21 | RECENT_BACKFILL_DAYS = 365 | ||
| 22 | RECENT_BACKFILL_BATCHES_PER_SERVICE = 5 | 22 | RECENT_BACKFILL_BATCHES_PER_SERVICE = 5 |
| 23 | FULL_BACKFILL_BATCHES_PER_SERVICE = 1 | ||
| 24 | 23 | ||
| 25 | 24 | ||
| 26 | class PollerService: | 25 | class PollerService: |
| @@ -61,6 +60,9 @@ class PollerService: | |||
| 61 | logger.exception("Scheduled poll failed for actor=%s", actor) | 60 | logger.exception("Scheduled poll failed for actor=%s", actor) |
| 62 | except Exception: | 61 | except Exception: |
| 63 | logger.exception("Unexpected scheduled poll failure for actor=%s", actor) | 62 | logger.exception("Unexpected scheduled poll failure for actor=%s", actor) |
| 63 | deleted = self.prune_old_events(db) | ||
| 64 | if deleted: | ||
| 65 | logger.info("Pruned %s contribution events older than %s days", deleted, RETENTION_DAYS) | ||
| 64 | return results | 66 | return results |
| 65 | 67 | ||
| 66 | def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: | 68 | def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: |
| @@ -70,7 +72,6 @@ class PollerService: | |||
| 70 | actor=actor, | 72 | actor=actor, |
| 71 | is_active=True, | 73 | is_active=True, |
| 72 | recent_backfill_status="pending", | 74 | recent_backfill_status="pending", |
| 73 | backfill_status="pending", | ||
| 74 | ) | 75 | ) |
| 75 | db.add(tracked_actor) | 76 | db.add(tracked_actor) |
| 76 | 77 | ||
| @@ -185,11 +186,11 @@ class PollerService: | |||
| 185 | 186 | ||
| 186 | def _run_backfill_batches(self, db: Session, actor: str) -> int: | 187 | def _run_backfill_batches(self, db: Session, actor: str) -> int: |
| 187 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) | 188 | tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) |
| 188 | if tracked_actor.recent_backfill_status == "completed" and tracked_actor.backfill_status == "completed": | 189 | if tracked_actor.recent_backfill_status == "completed": |
| 189 | return 0 | 190 | return 0 |
| 190 | 191 | ||
| 191 | now = datetime.now(tz=UTC) | 192 | now = datetime.now(tz=UTC) |
| 192 | recent_since = now - timedelta(days=RECENT_BACKFILL_DAYS) | 193 | recent_since = now - timedelta(days=RETENTION_DAYS) |
| 193 | db.add(tracked_actor) | 194 | db.add(tracked_actor) |
| 194 | db.flush() | 195 | db.flush() |
| 195 | total_inserted = 0 | 196 | total_inserted = 0 |
| @@ -223,34 +224,6 @@ class PollerService: | |||
| 223 | db.add(tracked_actor) | 224 | db.add(tracked_actor) |
| 224 | db.flush() | 225 | db.flush() |
| 225 | 226 | ||
| 226 | if tracked_actor.recent_backfill_status == "completed" and tracked_actor.backfill_status != "completed": | ||
| 227 | if tracked_actor.backfill_started_at is None: | ||
| 228 | tracked_actor.backfill_started_at = now | ||
| 229 | tracked_actor.backfill_status = "in_progress" | ||
| 230 | tracked_actor.last_backfill_error = None | ||
| 231 | total_inserted += self._run_backfill_scope( | ||
| 232 | db, | ||
| 233 | actor=actor, | ||
| 234 | scope="full", | ||
| 235 | services=[ | ||
| 236 | (self.todo_service.service_name, self.todo_service.fetch_backfill_batch), | ||
| 237 | (self.git_service.service_name, self.git_service.fetch_backfill_batch), | ||
| 238 | ], | ||
| 239 | batches_per_service=FULL_BACKFILL_BATCHES_PER_SERVICE, | ||
| 240 | since=None, | ||
| 241 | ) | ||
| 242 | full_statuses = db.scalars( | ||
| 243 | select(ServiceBackfillState.status) | ||
| 244 | .where(ServiceBackfillState.actor == actor) | ||
| 245 | .where(ServiceBackfillState.scope == "full") | ||
| 246 | ).all() | ||
| 247 | if full_statuses and all(status == "completed" for status in full_statuses): | ||
| 248 | tracked_actor.backfill_status = "completed" | ||
| 249 | tracked_actor.backfill_completed_at = datetime.now(tz=UTC) | ||
| 250 | tracked_actor.last_backfill_error = None | ||
| 251 | db.add(tracked_actor) | ||
| 252 | db.flush() | ||
| 253 | |||
| 254 | return total_inserted | 227 | return total_inserted |
| 255 | 228 | ||
| 256 | def _run_backfill_scope( | 229 | def _run_backfill_scope( |
| @@ -328,12 +301,8 @@ class PollerService: | |||
| 328 | state.status = "error" | 301 | state.status = "error" |
| 329 | state.last_error = str(exc) | 302 | state.last_error = str(exc) |
| 330 | state.updated_at = datetime.now(tz=UTC) | 303 | state.updated_at = datetime.now(tz=UTC) |
| 331 | if scope == "recent": | 304 | tracked_actor.recent_backfill_status = "error" |
| 332 | tracked_actor.recent_backfill_status = "error" | 305 | tracked_actor.last_recent_backfill_error = str(exc) |
| 333 | tracked_actor.last_recent_backfill_error = str(exc) | ||
| 334 | else: | ||
| 335 | tracked_actor.backfill_status = "error" | ||
| 336 | tracked_actor.last_backfill_error = str(exc) | ||
| 337 | db.add(state) | 306 | db.add(state) |
| 338 | db.add(tracked_actor) | 307 | db.add(tracked_actor) |
| 339 | db.flush() | 308 | db.flush() |
| @@ -343,3 +312,8 @@ class PollerService: | |||
| 343 | db.flush() | 312 | db.flush() |
| 344 | 313 | ||
| 345 | return total_inserted | 314 | return total_inserted |
| 315 | |||
| 316 | def prune_old_events(self, db: Session) -> int: | ||
| 317 | deleted = prune_contribution_events(db) | ||
| 318 | db.commit() | ||
| 319 | return deleted | ||
src/srht_contrib/schemas.py −3
| @@ -31,9 +31,6 @@ class ContributionIndexMetadata(BaseModel): | |||
| 31 | is_recent_window_backfilled: bool = False | 31 | is_recent_window_backfilled: bool = False |
| 32 | recent_backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" | 32 | recent_backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" |
| 33 | recent_backfill_completed_at: datetime | None = None | 33 | recent_backfill_completed_at: datetime | None = None |
| 34 | is_backfilled: bool = False | ||
| 35 | backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" | ||
| 36 | backfill_completed_at: datetime | None = None | ||
| 37 | 34 | ||
| 38 | 35 | ||
| 39 | class ContributionCalendarResponse(ContributionIndexMetadata): | 36 | class ContributionCalendarResponse(ContributionIndexMetadata): |
src/srht_contrib/services/aggregator.py −6
| @@ -65,9 +65,6 @@ class ContributionAggregator: | |||
| 65 | is_recent_window_backfilled=calendar.is_recent_window_backfilled, | 65 | is_recent_window_backfilled=calendar.is_recent_window_backfilled, |
| 66 | recent_backfill_state=calendar.recent_backfill_state, | 66 | recent_backfill_state=calendar.recent_backfill_state, |
| 67 | recent_backfill_completed_at=calendar.recent_backfill_completed_at, | 67 | recent_backfill_completed_at=calendar.recent_backfill_completed_at, |
| 68 | is_backfilled=calendar.is_backfilled, | ||
| 69 | backfill_state=calendar.backfill_state, | ||
| 70 | backfill_completed_at=calendar.backfill_completed_at, | ||
| 71 | ) | 68 | ) |
| 72 | 69 | ||
| 73 | def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata: | 70 | def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata: |
| @@ -90,9 +87,6 @@ class ContributionAggregator: | |||
| 90 | is_recent_window_backfilled=(tracked_actor.recent_backfill_status == "completed") if tracked_actor is not None else False, | 87 | is_recent_window_backfilled=(tracked_actor.recent_backfill_status == "completed") if tracked_actor is not None else False, |
| 91 | recent_backfill_state=(tracked_actor.recent_backfill_status if tracked_actor is not None else "pending"), | 88 | recent_backfill_state=(tracked_actor.recent_backfill_status if tracked_actor is not None else "pending"), |
| 92 | recent_backfill_completed_at=tracked_actor.recent_backfill_completed_at if tracked_actor is not None else None, | 89 | recent_backfill_completed_at=tracked_actor.recent_backfill_completed_at if tracked_actor is not None else None, |
| 93 | is_backfilled=(tracked_actor.backfill_status == "completed") if tracked_actor is not None else False, | ||
| 94 | backfill_state=(tracked_actor.backfill_status if tracked_actor is not None else "pending"), | ||
| 95 | backfill_completed_at=tracked_actor.backfill_completed_at if tracked_actor is not None else None, | ||
| 96 | ) | 90 | ) |
| 97 | 91 | ||
| 98 | def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: | 92 | def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: |
src/srht_contrib/utils/retention.py added +17
| @@ -0,0 +1,17 @@ | |||
| 1 | from __future__ import annotations | ||
| 2 | |||
| 3 | from datetime import UTC, datetime, timedelta | ||
| 4 | |||
| 5 | from sqlalchemy import delete | ||
| 6 | from sqlalchemy.orm import Session | ||
| 7 | |||
| 8 | from srht_contrib.models import ContributionEvent | ||
| 9 | |||
| 10 | |||
| 11 | RETENTION_DAYS = 365 | ||
| 12 | |||
| 13 | |||
| 14 | def prune_contribution_events(db: Session, *, now: datetime | None = None) -> int: | ||
| 15 | cutoff = (now or datetime.now(tz=UTC)) - timedelta(days=RETENTION_DAYS) | ||
| 16 | result = db.execute(delete(ContributionEvent).where(ContributionEvent.occurred_at < cutoff)) | ||
| 17 | return int(result.rowcount or 0) | ||
tests/test_contributions_api.py −4
| @@ -68,9 +68,6 @@ def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_se | |||
| 68 | assert response.json()["is_recent_window_backfilled"] is False | 68 | assert response.json()["is_recent_window_backfilled"] is False |
| 69 | assert response.json()["recent_backfill_state"] == "pending" | 69 | assert response.json()["recent_backfill_state"] == "pending" |
| 70 | assert response.json()["recent_backfill_completed_at"] is None | 70 | assert response.json()["recent_backfill_completed_at"] is None |
| 71 | assert response.json()["is_backfilled"] is False | ||
| 72 | assert response.json()["backfill_state"] == "pending" | ||
| 73 | assert response.json()["backfill_completed_at"] is None | ||
| 74 | assert response.json()["last_polled_at"] is None | 71 | assert response.json()["last_polled_at"] is None |
| 75 | assert tracked_actor is not None | 72 | assert tracked_actor is not None |
| 76 | assert tracked_actor.is_active is True | 73 | assert tracked_actor.is_active is True |
| @@ -117,7 +114,6 @@ def test_contribution_stats_api(client: TestClient, db_session) -> None: | |||
| 117 | assert response.json()["indexing_state"] == "indexed" | 114 | assert response.json()["indexing_state"] == "indexed" |
| 118 | assert response.json()["is_recent_window_backfilled"] is False | 115 | assert response.json()["is_recent_window_backfilled"] is False |
| 119 | assert response.json()["recent_backfill_state"] == "pending" | 116 | assert response.json()["recent_backfill_state"] == "pending" |
| 120 | assert response.json()["is_backfilled"] is False | ||
| 121 | 117 | ||
| 122 | 118 | ||
| 123 | def test_invalid_date_input_returns_400(client: TestClient) -> None: | 119 | def test_invalid_date_input_returns_400(client: TestClient) -> None: |
tests/test_ingestion.py +61 −13
| @@ -4,7 +4,7 @@ from sqlalchemy import select | |||
| 4 | 4 | ||
| 5 | from srht_contrib.config import Settings | 5 | from srht_contrib.config import Settings |
| 6 | from srht_contrib.jobs.poller import PollerService | 6 | from srht_contrib.jobs.poller import PollerService |
| 7 | from srht_contrib.models import ServiceBackfillState, SyncState, TrackedActor, TrackedRepository | 7 | from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository |
| 8 | from srht_contrib.schemas import NormalizedEvent | 8 | from srht_contrib.schemas import NormalizedEvent |
| 9 | from srht_contrib.services.git import GitIngestionService, GitPollResult | 9 | from srht_contrib.services.git import GitIngestionService, GitPollResult |
| 10 | from srht_contrib.services.todo import TodoIngestionService, TodoPollResult | 10 | from srht_contrib.services.todo import TodoIngestionService, TodoPollResult |
| @@ -94,7 +94,7 @@ class BackfillingTodoService: | |||
| 94 | repo_name="todo", | 94 | repo_name="todo", |
| 95 | resource_id="backfill-ticket", | 95 | resource_id="backfill-ticket", |
| 96 | external_uid=f"todo:backfill:{actor}", | 96 | external_uid=f"todo:backfill:{actor}", |
| 97 | occurred_at=datetime(2024, 1, 1, 12, 0, tzinfo=UTC), | 97 | occurred_at=datetime(2026, 1, 1, 12, 0, tzinfo=UTC), |
| 98 | weight=1.0, | 98 | weight=1.0, |
| 99 | raw_payload_json=None, | 99 | raw_payload_json=None, |
| 100 | ) | 100 | ) |
| @@ -119,7 +119,13 @@ class QueueShrinkingTodoService: | |||
| 119 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: | 119 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: |
| 120 | import copy | 120 | import copy |
| 121 | 121 | ||
| 122 | state = {"tracker_queue": ["t1", "t2"], "current_tracker": None, "current_ticket": None, "trackers_loaded": True, "trackers_cursor": None} | 122 | state = { |
| 123 | "tracker_queue": ["t1", "t2", "t3", "t4", "t5", "t6"], | ||
| 124 | "current_tracker": None, | ||
| 125 | "current_ticket": None, | ||
| 126 | "trackers_loaded": True, | ||
| 127 | "trackers_cursor": None, | ||
| 128 | } | ||
| 123 | if cursor_state: | 129 | if cursor_state: |
| 124 | state.update(copy.deepcopy(cursor_state)) | 130 | state.update(copy.deepcopy(cursor_state)) |
| 125 | if not state["tracker_queue"]: | 131 | if not state["tracker_queue"]: |
| @@ -156,7 +162,12 @@ class QueueShrinkingGitService: | |||
| 156 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: | 162 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: |
| 157 | import copy | 163 | import copy |
| 158 | 164 | ||
| 159 | state = {"repository_queue": ["r1", "r2"], "current_repository": None, "discovery_complete": True, "discovery_cursor": None} | 165 | state = { |
| 166 | "repository_queue": ["r1", "r2", "r3", "r4", "r5", "r6"], | ||
| 167 | "current_repository": None, | ||
| 168 | "discovery_complete": True, | ||
| 169 | "discovery_cursor": None, | ||
| 170 | } | ||
| 160 | if cursor_state: | 171 | if cursor_state: |
| 161 | state.update(copy.deepcopy(cursor_state)) | 172 | state.update(copy.deepcopy(cursor_state)) |
| 162 | if not state["repository_queue"]: | 173 | if not state["repository_queue"]: |
| @@ -561,9 +572,7 @@ def test_poll_marks_backfill_complete_and_persists_service_state(db_session) -> | |||
| 561 | assert tracked_actor is not None | 572 | assert tracked_actor is not None |
| 562 | assert tracked_actor.recent_backfill_status == "completed" | 573 | assert tracked_actor.recent_backfill_status == "completed" |
| 563 | assert tracked_actor.recent_backfill_completed_at is not None | 574 | assert tracked_actor.recent_backfill_completed_at is not None |
| 564 | assert tracked_actor.backfill_status == "completed" | 575 | assert [f"{state.scope}:{state.service}" for state in service_states] == ["recent:git", "recent:todo"] |
| 565 | assert tracked_actor.backfill_completed_at is not None | ||
| 566 | assert [f"{state.scope}:{state.service}" for state in service_states] == ["full:git", "full:todo", "recent:git", "recent:todo"] | ||
| 567 | assert all(state.status == "completed" for state in service_states) | 576 | assert all(state.status == "completed" for state in service_states) |
| 568 | 577 | ||
| 569 | 578 | ||
| @@ -576,7 +585,7 @@ def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None | |||
| 576 | for state in db_session.scalars( | 585 | for state in db_session.scalars( |
| 577 | select(ServiceBackfillState) | 586 | select(ServiceBackfillState) |
| 578 | .where(ServiceBackfillState.actor == "~ccleberg") | 587 | .where(ServiceBackfillState.actor == "~ccleberg") |
| 579 | .where(ServiceBackfillState.scope == "full") | 588 | .where(ServiceBackfillState.scope == "recent") |
| 580 | ).all() | 589 | ).all() |
| 581 | } | 590 | } |
| 582 | 591 | ||
| @@ -586,11 +595,50 @@ def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None | |||
| 586 | for state in db_session.scalars( | 595 | for state in db_session.scalars( |
| 587 | select(ServiceBackfillState) | 596 | select(ServiceBackfillState) |
| 588 | .where(ServiceBackfillState.actor == "~ccleberg") | 597 | .where(ServiceBackfillState.actor == "~ccleberg") |
| 589 | .where(ServiceBackfillState.scope == "full") | 598 | .where(ServiceBackfillState.scope == "recent") |
| 590 | ).all() | 599 | ).all() |
| 591 | } | 600 | } |
| 592 | 601 | ||
| 593 | assert first_states["git"]["repository_queue"] == ["r2"] | 602 | assert first_states["git"]["repository_queue"] == ["r6"] |
| 594 | assert first_states["todo"]["tracker_queue"] == ["t2"] | 603 | assert first_states["todo"]["tracker_queue"] == ["t6"] |
| 595 | assert second_states["git"]["repository_queue"] == [] | 604 | assert second_states["git"] is None |
| 596 | assert second_states["todo"]["tracker_queue"] == [] | 605 | assert second_states["todo"] is None |
| 606 | |||
| 607 | |||
| 608 | def test_prune_old_events_removes_data_older_than_one_year(db_session) -> None: | ||
| 609 | poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService()) | ||
| 610 | db_session.add_all( | ||
| 611 | [ | ||
| 612 | ContributionEvent( | ||
| 613 | service="todo", | ||
| 614 | event_type="ticket_created", | ||
| 615 | actor="~ccleberg", | ||
| 616 | repo_name="todo", | ||
| 617 | resource_id="old", | ||
| 618 | external_uid="todo:old", | ||
| 619 | occurred_at=datetime(2025, 1, 1, 12, 0, tzinfo=UTC), | ||
| 620 | weight=1.0, | ||
| 621 | raw_payload_json=None, | ||
| 622 | ), | ||
| 623 | ContributionEvent( | ||
| 624 | service="todo", | ||
| 625 | event_type="ticket_created", | ||
| 626 | actor="~ccleberg", | ||
| 627 | repo_name="todo", | ||
| 628 | resource_id="recent", | ||
| 629 | external_uid="todo:recent", | ||
| 630 | occurred_at=datetime(2026, 4, 1, 12, 0, tzinfo=UTC), | ||
| 631 | weight=1.0, | ||
| 632 | raw_payload_json=None, | ||
| 633 | ), | ||
| 634 | ] | ||
| 635 | ) | ||
| 636 | db_session.commit() | ||
| 637 | |||
| 638 | deleted = poller.prune_old_events(db_session) | ||
| 639 | remaining = db_session.scalars( | ||
| 640 | select(ContributionEvent.external_uid).order_by(ContributionEvent.external_uid) | ||
| 641 | ).all() | ||
| 642 | |||
| 643 | assert deleted == 1 | ||
| 644 | assert remaining == ["todo:recent"] | ||