diff options
| author | Christian Cleberg <[email protected]> | 2026-04-11 18:54:01 -0500 |
|---|---|---|
| committer | Christian Cleberg <[email protected]> | 2026-04-11 18:54:01 -0500 |
| commit | 90e61ed3e8e765ff5557bc59a6a7b141124dba01 (patch) | |
| tree | 8d76a215116c7f8ae9c31170c03226af6bccf225 /src | |
| parent | 01c1b50541b0dc1d42cbdaa90052f6b94ceba20c (diff) | |
| download | hutch-stats-90e61ed3e8e765ff5557bc59a6a7b141124dba01.tar.gz hutch-stats-90e61ed3e8e765ff5557bc59a6a7b141124dba01.tar.bz2 hutch-stats-90e61ed3e8e765ff5557bc59a6a7b141124dba01.zip | |
refactor: retain and backfill only the last year of activity
Diffstat (limited to 'src')
| -rw-r--r-- | src/srht_contrib/jobs/poller.py | 52 | ||||
| -rw-r--r-- | src/srht_contrib/schemas.py | 3 | ||||
| -rw-r--r-- | src/srht_contrib/services/aggregator.py | 6 | ||||
| -rw-r--r-- | src/srht_contrib/utils/retention.py | 17 |
4 files changed, 30 insertions, 48 deletions
diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py index c41f3d4..ff0aaae 100644 --- a/src/srht_contrib/jobs/poller.py +++ b/src/srht_contrib/jobs/poller.py @@ -13,14 +13,13 @@ from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.git import GitIngestionService from srht_contrib.services.srht_client import SourceHutClientError from srht_contrib.services.todo import TodoIngestionService +from srht_contrib.utils.retention import RETENTION_DAYS, prune_contribution_events from srht_contrib.utils.repositories import canonicalize_repository_name logger = logging.getLogger(__name__) SYNC_OVERLAP = timedelta(hours=24) -RECENT_BACKFILL_DAYS = 365 RECENT_BACKFILL_BATCHES_PER_SERVICE = 5 -FULL_BACKFILL_BATCHES_PER_SERVICE = 1 class PollerService: @@ -61,6 +60,9 @@ class PollerService: logger.exception("Scheduled poll failed for actor=%s", actor) except Exception: logger.exception("Unexpected scheduled poll failure for actor=%s", actor) + deleted = self.prune_old_events(db) + if deleted: + logger.info("Pruned %s contribution events older than %s days", deleted, RETENTION_DAYS) return results def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: @@ -70,7 +72,6 @@ class PollerService: actor=actor, is_active=True, recent_backfill_status="pending", - backfill_status="pending", ) db.add(tracked_actor) @@ -185,11 +186,11 @@ class PollerService: 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" and tracked_actor.backfill_status == "completed": + if tracked_actor.recent_backfill_status == "completed": return 0 now = datetime.now(tz=UTC) - recent_since = now - timedelta(days=RECENT_BACKFILL_DAYS) + recent_since = now - timedelta(days=RETENTION_DAYS) db.add(tracked_actor) db.flush() total_inserted = 0 @@ -223,34 +224,6 @@ class PollerService: db.add(tracked_actor) db.flush() - if tracked_actor.recent_backfill_status == "completed" and tracked_actor.backfill_status != "completed": - if tracked_actor.backfill_started_at is None: - tracked_actor.backfill_started_at = now - tracked_actor.backfill_status = "in_progress" - tracked_actor.last_backfill_error = None - total_inserted += self._run_backfill_scope( - db, - actor=actor, - scope="full", - services=[ - (self.todo_service.service_name, self.todo_service.fetch_backfill_batch), - (self.git_service.service_name, self.git_service.fetch_backfill_batch), - ], - batches_per_service=FULL_BACKFILL_BATCHES_PER_SERVICE, - since=None, - ) - full_statuses = db.scalars( - select(ServiceBackfillState.status) - .where(ServiceBackfillState.actor == actor) - .where(ServiceBackfillState.scope == "full") - ).all() - if full_statuses and all(status == "completed" for status in full_statuses): - tracked_actor.backfill_status = "completed" - tracked_actor.backfill_completed_at = datetime.now(tz=UTC) - tracked_actor.last_backfill_error = None - db.add(tracked_actor) - db.flush() - return total_inserted def _run_backfill_scope( @@ -328,12 +301,8 @@ class PollerService: state.status = "error" state.last_error = str(exc) state.updated_at = datetime.now(tz=UTC) - if scope == "recent": - tracked_actor.recent_backfill_status = "error" - tracked_actor.last_recent_backfill_error = str(exc) - else: - tracked_actor.backfill_status = "error" - tracked_actor.last_backfill_error = str(exc) + tracked_actor.recent_backfill_status = "error" + tracked_actor.last_recent_backfill_error = str(exc) db.add(state) db.add(tracked_actor) db.flush() @@ -343,3 +312,8 @@ class PollerService: db.flush() return total_inserted + + def prune_old_events(self, db: Session) -> int: + deleted = prune_contribution_events(db) + db.commit() + return deleted diff --git a/src/srht_contrib/schemas.py b/src/srht_contrib/schemas.py index 98d2a74..273cb26 100644 --- a/src/srht_contrib/schemas.py +++ b/src/srht_contrib/schemas.py @@ -31,9 +31,6 @@ class ContributionIndexMetadata(BaseModel): is_recent_window_backfilled: bool = False recent_backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" recent_backfill_completed_at: datetime | None = None - is_backfilled: bool = False - backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" - backfill_completed_at: datetime | None = None class ContributionCalendarResponse(ContributionIndexMetadata): diff --git a/src/srht_contrib/services/aggregator.py b/src/srht_contrib/services/aggregator.py index 4973f16..32baebd 100644 --- a/src/srht_contrib/services/aggregator.py +++ b/src/srht_contrib/services/aggregator.py @@ -65,9 +65,6 @@ class ContributionAggregator: is_recent_window_backfilled=calendar.is_recent_window_backfilled, recent_backfill_state=calendar.recent_backfill_state, recent_backfill_completed_at=calendar.recent_backfill_completed_at, - is_backfilled=calendar.is_backfilled, - backfill_state=calendar.backfill_state, - backfill_completed_at=calendar.backfill_completed_at, ) def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata: @@ -90,9 +87,6 @@ class ContributionAggregator: is_recent_window_backfilled=(tracked_actor.recent_backfill_status == "completed") if tracked_actor is not None else False, recent_backfill_state=(tracked_actor.recent_backfill_status if tracked_actor is not None else "pending"), recent_backfill_completed_at=tracked_actor.recent_backfill_completed_at if tracked_actor is not None else None, - is_backfilled=(tracked_actor.backfill_status == "completed") if tracked_actor is not None else False, - backfill_state=(tracked_actor.backfill_status if tracked_actor is not None else "pending"), - backfill_completed_at=tracked_actor.backfill_completed_at if tracked_actor is not None else None, ) def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: diff --git a/src/srht_contrib/utils/retention.py b/src/srht_contrib/utils/retention.py new file mode 100644 index 0000000..5018b32 --- /dev/null +++ b/src/srht_contrib/utils/retention.py @@ -0,0 +1,17 @@ +from __future__ import annotations + +from datetime import UTC, datetime, timedelta + +from sqlalchemy import delete +from sqlalchemy.orm import Session + +from srht_contrib.models import ContributionEvent + + +RETENTION_DAYS = 365 + + +def prune_contribution_events(db: Session, *, now: datetime | None = None) -> int: + cutoff = (now or datetime.now(tz=UTC)) - timedelta(days=RETENTION_DAYS) + result = db.execute(delete(ContributionEvent).where(ContributionEvent.occurred_at < cutoff)) + return int(result.rowcount or 0) |
