| @@ -8,6 +8,7 @@ from sqlalchemy import select |
| 8 | from sqlalchemy.exc import IntegrityError |
8 | from sqlalchemy.exc import IntegrityError |
| 9 | from sqlalchemy.orm import Session |
9 | from sqlalchemy.orm import Session |
| 10 | |
10 | |
| |
11 | from srht_contrib.config import Settings |
| 11 | from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository |
12 | from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository |
| 12 | from srht_contrib.schemas import NormalizedEvent |
13 | from srht_contrib.schemas import NormalizedEvent |
| 13 | from srht_contrib.services.git import GitIngestionService |
14 | from srht_contrib.services.git import GitIngestionService |
| @@ -18,14 +19,19 @@ from srht_contrib.utils.repositories import canonicalize_repository_name |
| 18 | |
19 | |
| 19 | |
20 | |
| 20 | logger = logging.getLogger(__name__) |
21 | logger = logging.getLogger(__name__) |
| 21 | SYNC_OVERLAP = timedelta(hours=24) |
| |
| 22 | RECENT_BACKFILL_BATCHES_PER_SERVICE = 5 |
22 | RECENT_BACKFILL_BATCHES_PER_SERVICE = 5 |
| 23 | |
23 | |
| 24 | |
24 | |
| 25 | class PollerService: |
25 | class PollerService: |
| 26 | def __init__(self, todo_service: TodoIngestionService, git_service: GitIngestionService) -> None: |
26 | def __init__( |
| |
27 | self, |
| |
28 | todo_service: TodoIngestionService, |
| |
29 | git_service: GitIngestionService, |
| |
30 | settings: Settings, |
| |
31 | ) -> None: |
| 27 | self.todo_service = todo_service |
32 | self.todo_service = todo_service |
| 28 | self.git_service = git_service |
33 | self.git_service = git_service |
| |
34 | self._sync_overlap = timedelta(hours=settings.sync_overlap_hours) |
| 29 | |
35 | |
| 30 | def poll_all(self, db: Session, actor: str) -> int: |
36 | def poll_all(self, db: Session, actor: str) -> int: |
| 31 | self.track_actor_request(db, actor, update_last_requested=False) |
37 | self.track_actor_request(db, actor, update_last_requested=False) |
| @@ -104,7 +110,7 @@ class PollerService: |
| 104 | ) |
110 | ) |
| 105 | since = datetime.now(tz=UTC) - timedelta(days=30) |
111 | since = datetime.now(tz=UTC) - timedelta(days=30) |
| 106 | if state and state.cursor_value: |
112 | if state and state.cursor_value: |
| 107 | since = datetime.fromisoformat(state.cursor_value.replace("Z", "+00:00")).astimezone(UTC) - SYNC_OVERLAP |
113 | since = datetime.fromisoformat(state.cursor_value.replace("Z", "+00:00")).astimezone(UTC) - self._sync_overlap |
| 108 | logger.info( |
114 | logger.info( |
| 109 | "Using sync cursor for %s actor=%s with overlap; since=%s", |
115 | "Using sync cursor for %s actor=%s with overlap; since=%s", |
| 110 | service_name, |
116 | service_name, |