krz/hutch-stats
Server-side utility for calculating contributions for sourcehut users.
clone: git clone https://gitbay.org/krz/hutch-stats.git
main: tests/test_ingestion.py · raw
1from datetime import UTC, datetime, timedelta
2from pathlib import Path
3
4from sqlalchemy import select
5
6from srht_contrib.config import Settings
7from srht_contrib.jobs.poller import PollerService
8from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository
9from srht_contrib.scripts.enqueue_actors import enqueue_actors
10from srht_contrib.schemas import NormalizedEvent
11from srht_contrib.services.git import GitIngestionService, GitPollResult
12from srht_contrib.services.srht_client import SourceHutClientError
13from srht_contrib.services.todo import TodoIngestionService, TodoPollResult
14from srht_contrib.services.types import BackfillBatchResult
15
16
17class StubClient:
18 def __init__(self, payload: dict | None = None, payloads_by_query: dict[str, dict] | None = None) -> None:
19 self.payload = payload or {}
20 self.payloads_by_query = payloads_by_query or {}
21 self.calls: list[tuple[str, dict | None]] = []
22
23 def execute(self, query: str, variables: dict | None = None) -> dict:
24 self.calls.append((query, variables))
25 for marker, payload in self.payloads_by_query.items():
26 if marker in query:
27 return payload
28 return self.payload
29
30
31class RecordingTodoService:
32 service_name = "todo"
33
34 def __init__(self, events_by_call: list[list[NormalizedEvent]]) -> None:
35 self.events_by_call = events_by_call
36 self.calls: list[datetime] = []
37
38 def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult:
39 assert since is not None
40 self.calls.append(since)
41 events = self.events_by_call.pop(0)
42 return TodoPollResult(events=events, cursor="2026-03-31T00:00:00+00:00")
43
44 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
45 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
46
47 def fetch_recent_backfill_batch(
48 self,
49 actor: str,
50 cursor_state: dict | None = None,
51 *,
52 since: datetime,
53 ) -> BackfillBatchResult:
54 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
55
56
57class EmptyGitService:
58 service_name = "git"
59
60 def __init__(self) -> None:
61 self.settings = Settings(
62 SRHT_TOKEN="x",
63 DATABASE_URL="sqlite://",
64 DEFAULT_ACTOR="~ccleberg",
65 TODO_SRHT_ENDPOINT="https://todo.sr.ht/query",
66 GIT_SRHT_ENDPOINT="https://git.sr.ht/query",
67 POLL_INTERVAL_SECONDS=60,
68 )
69
70 def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None, db=None) -> GitPollResult:
71 return GitPollResult(events=[], cursor="2026-03-31T00:00:00+00:00")
72
73 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
74 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
75
76 def fetch_recent_backfill_batch(
77 self,
78 actor: str,
79 cursor_state: dict | None = None,
80 *,
81 since: datetime,
82 ) -> BackfillBatchResult:
83 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
84
85
86class BackfillingTodoService:
87 service_name = "todo"
88
89 def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult:
90 return TodoPollResult(events=[], cursor="2026-03-31T00:00:00+00:00")
91
92 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
93 event = NormalizedEvent(
94 service="todo",
95 event_type="ticket_created",
96 actor=actor,
97 repo_name="todo",
98 resource_id="backfill-ticket",
99 external_uid=f"todo:backfill:{actor}",
100 occurred_at=datetime(2026, 1, 1, 12, 0, tzinfo=UTC),
101 weight=1.0,
102 raw_payload_json=None,
103 )
104 return BackfillBatchResult(events=[event], cursor_state=None, complete=True)
105
106 def fetch_recent_backfill_batch(
107 self,
108 actor: str,
109 cursor_state: dict | None = None,
110 *,
111 since: datetime,
112 ) -> BackfillBatchResult:
113 return self.fetch_backfill_batch(actor, cursor_state)
114
115
116class 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
135class QueueShrinkingTodoService:
136 service_name = "todo"
137
138 def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult:
139 return TodoPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat())
140
141 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
142 import copy
143
144 state = {
145 "tracker_queue": ["t1", "t2", "t3", "t4", "t5", "t6"],
146 "current_tracker": None,
147 "current_ticket": None,
148 "trackers_loaded": True,
149 "trackers_cursor": None,
150 }
151 if cursor_state:
152 state.update(copy.deepcopy(cursor_state))
153 if not state["tracker_queue"]:
154 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
155 state["tracker_queue"].pop(0)
156 return BackfillBatchResult(events=[], cursor_state=state, complete=False)
157
158 def fetch_recent_backfill_batch(
159 self,
160 actor: str,
161 cursor_state: dict | None = None,
162 *,
163 since: datetime,
164 ) -> BackfillBatchResult:
165 return self.fetch_backfill_batch(actor, cursor_state)
166
167
168class QueueShrinkingGitService:
169 service_name = "git"
170
171 def __init__(self) -> None:
172 self.settings = Settings(
173 SRHT_TOKEN="x",
174 DATABASE_URL="sqlite://",
175 DEFAULT_ACTOR="~ccleberg",
176 TODO_SRHT_ENDPOINT="https://todo.sr.ht/query",
177 GIT_SRHT_ENDPOINT="https://git.sr.ht/query",
178 POLL_INTERVAL_SECONDS=60,
179 )
180
181 def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None, db=None) -> GitPollResult:
182 return GitPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat())
183
184 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
185 import copy
186
187 state = {
188 "repository_queue": ["r1", "r2", "r3", "r4", "r5", "r6"],
189 "current_repository": None,
190 "discovery_complete": True,
191 "discovery_cursor": None,
192 }
193 if cursor_state:
194 state.update(copy.deepcopy(cursor_state))
195 if not state["repository_queue"]:
196 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
197 state["repository_queue"].pop(0)
198 return BackfillBatchResult(events=[], cursor_state=state, complete=False)
199
200 def fetch_recent_backfill_batch(
201 self,
202 actor: str,
203 cursor_state: dict | None = None,
204 *,
205 since: datetime,
206 ) -> BackfillBatchResult:
207 return self.fetch_backfill_batch(actor, cursor_state)
208
209
210def make_settings(**overrides) -> Settings:
211 values = {
212 "API_KEY": "test-api-key",
213 "ENABLE_SCHEDULER": False,
214 "SRHT_TOKEN": "x",
215 "DATABASE_URL": "sqlite://",
216 "DEFAULT_ACTOR": "~ccleberg",
217 "TODO_SRHT_ENDPOINT": "https://todo.sr.ht/query",
218 "GIT_SRHT_ENDPOINT": "https://git.sr.ht/query",
219 "POLL_INTERVAL_SECONDS": 60,
220 }
221 values.update(overrides)
222 return Settings(**values)
223
224
225def branch_payload(*branches: str) -> dict:
226 return {
227 "user": {
228 "repository": {
229 "references": {
230 "results": [{"name": branch, "target": "abc123"} for branch in branches],
231 "cursor": None,
232 }
233 }
234 }
235 }
236
237
238def test_todo_ingestion_is_idempotent(db_session) -> None:
239 settings = make_settings()
240 # The poller polls the last 30 days, so keep fixture events inside that window.
241 now = datetime.now(tz=UTC)
242 payload = {
243 "me": {"canonicalName": "~ccleberg"},
244 "events": {
245 "results": [
246 {
247 "id": "1001",
248 "created": (now - timedelta(days=3)).isoformat(),
249 "ticket": {
250 "id": "123",
251 "ref": "~ccleberg/todo/123",
252 "status": "RESOLVED",
253 "resolution": "CLOSED",
254 "tracker": {"name": "todo"},
255 },
256 "changes": [
257 {
258 "__typename": "Created",
259 "eventType": "CREATED",
260 "ticket": {"id": "123"},
261 "author": {"canonicalName": "~ccleberg"},
262 }
263 ],
264 },
265 {
266 "id": "1002",
267 "created": (now - timedelta(days=2, hours=1)).isoformat(),
268 "ticket": {
269 "id": "123",
270 "ref": "~ccleberg/todo/123",
271 "status": "RESOLVED",
272 "resolution": "CLOSED",
273 "tracker": {"name": "todo"},
274 },
275 "changes": [
276 {
277 "__typename": "Comment",
278 "eventType": "COMMENT",
279 "ticket": {"id": "123"},
280 "author": {"canonicalName": "~ccleberg"},
281 }
282 ],
283 },
284 {
285 "id": "1003",
286 "created": (now - timedelta(days=2)).isoformat(),
287 "ticket": {
288 "id": "123",
289 "ref": "~ccleberg/todo/123",
290 "status": "RESOLVED",
291 "resolution": "CLOSED",
292 "tracker": {"name": "todo"},
293 },
294 "changes": [
295 {
296 "__typename": "StatusChange",
297 "eventType": "STATUS_CHANGE",
298 "ticket": {"id": "123"},
299 "editor": {"canonicalName": "~ccleberg"},
300 "oldStatus": "IN_PROGRESS",
301 "newStatus": "RESOLVED",
302 "oldResolution": "UNRESOLVED",
303 "newResolution": "CLOSED",
304 }
305 ],
306 },
307 ],
308 "cursor": None,
309 },
310 }
311
312 todo_service = TodoIngestionService(StubClient(payload), settings)
313 git_service = GitIngestionService(StubClient(payload={}), settings)
314 poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings)
315
316 first_inserted = poller.poll_all(db_session, "~ccleberg")
317 second_inserted = poller.poll_all(db_session, "~ccleberg")
318
319 assert first_inserted == 3
320 assert second_inserted == 0
321
322
323def test_todo_ingestion_falls_back_to_tracker_crawl(db_session) -> None:
324 settings = make_settings()
325 client = StubClient(
326 payloads_by_query={
327 "query TodoActivity": {
328 "me": {"canonicalName": "~ccleberg"},
329 "events": {"results": [], "cursor": None},
330 },
331 "query TodoTrackers": {
332 "me": {
333 "canonicalName": "~ccleberg",
334 "trackers": {"results": [{"id": "1", "rid": "tracker-rid", "name": "todo"}], "cursor": None},
335 }
336 },
337 "query TodoTrackerTickets": {
338 "tracker": {
339 "id": "1",
340 "name": "todo",
341 "tickets": {
342 "results": [
343 {
344 "id": 123,
345 "ref": "~ccleberg/todo/123",
346 "created": "2026-03-29T09:00:00Z",
347 "updated": "2026-03-30T09:00:00Z",
348 "status": "RESOLVED",
349 "resolution": "CLOSED",
350 "submitter": {"canonicalName": "~ccleberg"},
351 }
352 ],
353 "cursor": None,
354 },
355 }
356 },
357 "query TodoTicketEvents": {
358 "tracker": {
359 "ticket": {
360 "id": 123,
361 "ref": "~ccleberg/todo/123",
362 "status": "RESOLVED",
363 "resolution": "CLOSED",
364 "events": {
365 "results": [
366 {
367 "id": "evt-1",
368 "created": "2026-03-30T09:00:00Z",
369 "changes": [
370 {
371 "__typename": "Comment",
372 "eventType": "COMMENT",
373 "ticket": {"id": "123"},
374 "author": {"canonicalName": "~ccleberg"},
375 }
376 ],
377 }
378 ],
379 "cursor": None,
380 },
381 }
382 }
383 },
384 }
385 )
386 todo_service = TodoIngestionService(client, settings)
387 git_service = GitIngestionService(StubClient(payload={}), settings)
388 poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings)
389
390 inserted = poller.poll_all(db_session, "~ccleberg")
391
392 assert inserted == 1
393 assert any("query TodoTrackers" in call[0] for call in client.calls)
394
395
396def test_unsupported_todo_changes_are_ignored(db_session) -> None:
397 settings = make_settings()
398 payload = {
399 "me": {"canonicalName": "~ccleberg"},
400 "events": {
401 "results": [
402 {
403 "id": "1001",
404 "created": "2026-03-29T10:00:00Z",
405 "ticket": {
406 "id": "123",
407 "ref": "~ccleberg/todo/123",
408 "status": "OPEN",
409 "resolution": "UNRESOLVED",
410 "tracker": {"name": "todo"},
411 },
412 "changes": [
413 {"__typename": "LabelUpdate", "eventType": "LABEL_UPDATE", "ticket": {"id": "123"}},
414 {"__typename": "TicketMention", "eventType": "TICKET_MENTION", "ticket": {"id": "123"}},
415 ],
416 }
417 ],
418 "cursor": None,
419 },
420 }
421
422 todo_service = TodoIngestionService(StubClient(payload), settings)
423 git_service = GitIngestionService(StubClient(payload={}), settings)
424 poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings)
425
426 inserted = poller.poll_all(db_session, "~ccleberg")
427
428 assert inserted == 0
429
430
431def test_git_ingestion_normalizes_commit_aliases_and_repository_names(db_session) -> None:
432 settings = make_settings(
433 ACTOR_ALIASES_JSON={"~ccleberg": ["cmc@example.com", "Chris Cleberg"]},
434 GIT_TRACKED_REPOSITORIES=["Hutch"],
435 )
436 git_payload = {
437 "user": {
438 "repository": {
439 "name": "Hutch",
440 "owner": {"canonicalName": "~ccleberg"},
441 "log": {
442 "results": [
443 {
444 "id": "abc123",
445 "shortId": "abc123",
446 "author": {
447 "name": "Chris Cleberg",
448 "email": "cmc@example.com",
449 "time": "2026-03-30T12:00:00Z",
450 },
451 "committer": {
452 "name": "Chris Cleberg",
453 "email": "cmc@example.com",
454 "time": "2026-03-30T12:00:00Z",
455 },
456 "message": "Add contribution calendar",
457 }
458 ],
459 "cursor": None,
460 },
461 }
462 }
463 }
464
465 todo_service = TodoIngestionService(
466 StubClient(payload={"me": {"canonicalName": "~ccleberg"}, "events": {"results": [], "cursor": None}}),
467 settings,
468 )
469 git_service = GitIngestionService(
470 StubClient(
471 payloads_by_query={
472 "query RepositoryBranches": branch_payload("refs/heads/main"),
473 "query RepositoryLog": git_payload,
474 }
475 ),
476 settings,
477 )
478 poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings)
479
480 inserted = poller.poll_all(db_session, "~ccleberg")
481
482 assert inserted == 1
483
484 tracked_repositories = db_session.scalars(select(TrackedRepository.repo_name)).all()
485 assert tracked_repositories == ["~ccleberg/Hutch"]
486
487
488def test_git_ingestion_auto_discovers_owned_repositories(db_session) -> None:
489 settings = make_settings(
490 ACTOR_ALIASES_JSON={"~ccleberg": ["cmc@example.com", "Chris Cleberg"]},
491 GIT_TRACKED_REPOSITORIES=[],
492 )
493 client = StubClient(
494 payloads_by_query={
495 "query UserRepositories": {
496 "user": {
497 "repositories": {
498 "results": [
499 {
500 "name": "Hutch",
501 "visibility": "PUBLIC",
502 "owner": {"canonicalName": "~ccleberg"},
503 },
504 ],
505 "cursor": None,
506 }
507 }
508 },
509 "query RepositoryBranches": branch_payload("refs/heads/main"),
510 "query RepositoryLog": {
511 "user": {
512 "repository": {
513 "name": "Hutch",
514 "owner": {"canonicalName": "~ccleberg"},
515 "log": {
516 "results": [
517 {
518 "id": "abc123",
519 "shortId": "abc123",
520 "author": {
521 "name": "Chris Cleberg",
522 "email": "cmc@example.com",
523 "time": "2026-03-30T12:00:00Z",
524 },
525 "committer": {
526 "name": "Chris Cleberg",
527 "email": "cmc@example.com",
528 "time": "2026-03-30T12:00:00Z",
529 },
530 "message": "Auto-discovered repo commit",
531 }
532 ],
533 "cursor": None,
534 },
535 }
536 }
537 },
538 }
539 )
540
541 todo_service = TodoIngestionService(
542 StubClient(payload={"me": {"canonicalName": "~ccleberg"}, "events": {"results": [], "cursor": None}}),
543 settings,
544 )
545 git_service = GitIngestionService(client, settings)
546 poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings)
547
548 inserted = poller.poll_all(db_session, "~ccleberg")
549
550 assert inserted == 1
551 assert any("query UserRepositories" in call[0] for call in client.calls)
552
553
554def test_git_ingestion_reads_repository_logs_from_all_branches(db_session) -> None:
555 settings = make_settings(
556 ACTOR_ALIASES_JSON={"~ccleberg": ["cmc@example.com", "Chris Cleberg"]},
557 GIT_TRACKED_REPOSITORIES=["Hutch"],
558 )
559 client = StubClient(
560 payloads_by_query={
561 "query RepositoryBranches": branch_payload(
562 "refs/heads/main",
563 "refs/heads/trunk",
564 "refs/tags/v1.0.0",
565 ),
566 "query RepositoryLog": {
567 "user": {
568 "repository": {
569 "name": "Hutch",
570 "owner": {"canonicalName": "~ccleberg"},
571 "log": {
572 "results": [
573 {
574 "id": "abc123",
575 "shortId": "abc123",
576 "author": {
577 "name": "Chris Cleberg",
578 "email": "cmc@example.com",
579 "time": "2026-03-30T12:00:00Z",
580 },
581 "committer": {
582 "name": "Chris Cleberg",
583 "email": "cmc@example.com",
584 "time": "2026-03-30T12:00:00Z",
585 },
586 "message": "Commit reachable from more than one branch",
587 }
588 ],
589 "cursor": None,
590 },
591 }
592 }
593 },
594 }
595 )
596 git_service = GitIngestionService(client, settings)
597
598 result = git_service.fetch_recent_events("~ccleberg", since=datetime(2026, 3, 1, tzinfo=UTC))
599
600 repository_log_calls = [call for call in client.calls if "query RepositoryLog" in call[0]]
601 assert len(result.events) == 1
602 assert [call[1]["from"] for call in repository_log_calls] == ["refs/heads/main", "refs/heads/trunk"]
603
604
605def test_sync_overlap_reuses_cursor_window_and_suppresses_duplicates(db_session) -> None:
606 settings = make_settings()
607 event = NormalizedEvent(
608 service="todo",
609 event_type="ticket_created",
610 actor="~ccleberg",
611 repo_name="todo",
612 resource_id="123",
613 external_uid="todo:event:123:created:123",
614 occurred_at=datetime(2026, 3, 30, 10, 0, tzinfo=UTC),
615 weight=1.0,
616 raw_payload_json=None,
617 )
618 todo_service = RecordingTodoService(events_by_call=[[event], [event]])
619 poller = PollerService(todo_service=todo_service, git_service=EmptyGitService(), settings=settings)
620
621 first_inserted = poller.poll_all(db_session, "~ccleberg")
622 second_inserted = poller.poll_all(db_session, "~ccleberg")
623
624 state = db_session.scalar(select(SyncState).where(SyncState.service == "todo").where(SyncState.actor == "~ccleberg"))
625
626 assert first_inserted == 1
627 assert second_inserted == 0
628 assert state is not None
629 assert len(todo_service.calls) == 2
630 assert todo_service.calls[1].isoformat() == "2026-03-30T23:00:00+00:00"
631
632
633def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) -> None:
634 settings = make_settings()
635 event = NormalizedEvent(
636 service="todo",
637 event_type="ticket_created",
638 actor="~known",
639 repo_name="todo",
640 resource_id="123",
641 external_uid="todo:event:known:created:123",
642 occurred_at=datetime(2026, 3, 30, 10, 0, tzinfo=UTC),
643 weight=1.0,
644 raw_payload_json=None,
645 )
646 todo_service = RecordingTodoService(events_by_call=[[], [event]])
647 poller = PollerService(todo_service=todo_service, git_service=EmptyGitService(), settings=settings)
648
649 db_session.add(TrackedActor(actor="~known", is_active=True))
650 db_session.commit()
651
652 results = poller.poll_tracked_actors(db_session, default_actor="~default")
653
654 tracked_actors = db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all()
655
656 assert results == {"~default": 1, "~known": 0}
657 assert [actor.actor for actor in tracked_actors] == ["~default", "~known"]
658 assert all(actor.last_poll_status == "indexed" for actor in tracked_actors)
659 assert all(actor.last_polled_at is not None for actor in tracked_actors)
660
661
662def test_poll_marks_backfill_complete_and_persists_service_state(db_session) -> None:
663 settings = make_settings()
664 poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService(), settings=settings)
665
666 inserted = poller.poll_all(db_session, "~ccleberg")
667
668 tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg"))
669 service_states = db_session.scalars(
670 select(ServiceBackfillState)
671 .where(ServiceBackfillState.actor == "~ccleberg")
672 .order_by(ServiceBackfillState.scope, ServiceBackfillState.service)
673 ).all()
674
675 assert inserted == 1
676 assert tracked_actor is not None
677 assert tracked_actor.recent_backfill_status == "completed"
678 assert tracked_actor.recent_backfill_completed_at is not None
679 assert [f"{state.scope}:{state.service}" for state in service_states] == ["recent:git", "recent:todo"]
680 assert all(state.status == "completed" for state in service_states)
681
682
683def test_scheduled_poll_indexes_before_draining_recent_backfill(db_session) -> None:
684 settings = make_settings(DISCOVERY_BATCH_SIZE=1, INDEXED_ACTOR_REPOLL_SECONDS=3600)
685 poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService(), settings=settings)
686 now = datetime.now(tz=UTC)
687 db_session.add(
688 TrackedActor(
689 actor="~ccleberg",
690 is_active=True,
691 discovery_state="queued",
692 queued_for_discovery_at=now - timedelta(minutes=1),
693 next_poll_after=now - timedelta(minutes=1),
694 recent_backfill_status="pending",
695 )
696 )
697 db_session.commit()
698
699 first_results = poller.poll_tracked_actors(db_session)
700 service_states_after_first_poll = db_session.scalars(select(ServiceBackfillState)).all()
701 tracked_after_first_poll = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg"))
702
703 assert first_results == {"~ccleberg": 0}
704 assert service_states_after_first_poll == []
705 assert tracked_after_first_poll is not None
706 assert tracked_after_first_poll.discovery_state == "indexed"
707 assert tracked_after_first_poll.recent_backfill_status == "pending"
708 assert tracked_after_first_poll.next_poll_after is not None
709 first_due_at = tracked_after_first_poll.next_poll_after
710 if first_due_at.tzinfo is None:
711 first_due_at = first_due_at.replace(tzinfo=UTC)
712 assert first_due_at <= datetime.now(tz=UTC)
713
714 second_results = poller.poll_tracked_actors(db_session)
715
716 tracked_after_second_poll = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg"))
717 service_states_after_second_poll = db_session.scalars(
718 select(ServiceBackfillState).order_by(ServiceBackfillState.scope, ServiceBackfillState.service)
719 ).all()
720
721 assert second_results == {"~ccleberg": 1}
722 assert tracked_after_second_poll is not None
723 assert tracked_after_second_poll.recent_backfill_status == "completed"
724 assert [f"{state.scope}:{state.service}" for state in service_states_after_second_poll] == ["recent:git", "recent:todo"]
725
726
727def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None:
728 settings = make_settings()
729 poller = PollerService(
730 todo_service=QueueShrinkingTodoService(),
731 git_service=QueueShrinkingGitService(),
732 settings=settings,
733 )
734
735 poller.poll_all(db_session, "~ccleberg")
736 first_states = {
737 state.service: state.cursor_json
738 for state in db_session.scalars(
739 select(ServiceBackfillState)
740 .where(ServiceBackfillState.actor == "~ccleberg")
741 .where(ServiceBackfillState.scope == "recent")
742 ).all()
743 }
744
745 poller.poll_all(db_session, "~ccleberg")
746 second_states = {
747 state.service: state.cursor_json
748 for state in db_session.scalars(
749 select(ServiceBackfillState)
750 .where(ServiceBackfillState.actor == "~ccleberg")
751 .where(ServiceBackfillState.scope == "recent")
752 ).all()
753 }
754
755 assert first_states["git"]["repository_queue"] == ["r6"]
756 assert first_states["todo"]["tracker_queue"] == ["t6"]
757 assert second_states["git"] is None
758 assert second_states["todo"] is None
759
760
761def test_prune_old_events_removes_data_older_than_one_year(db_session) -> None:
762 settings = make_settings()
763 poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService(), settings=settings)
764 db_session.add_all(
765 [
766 ContributionEvent(
767 service="todo",
768 event_type="ticket_created",
769 actor="~ccleberg",
770 repo_name="todo",
771 resource_id="old",
772 external_uid="todo:old",
773 occurred_at=datetime(2025, 1, 1, 12, 0, tzinfo=UTC),
774 weight=1.0,
775 raw_payload_json=None,
776 ),
777 ContributionEvent(
778 service="todo",
779 event_type="ticket_created",
780 actor="~ccleberg",
781 repo_name="todo",
782 resource_id="recent",
783 external_uid="todo:recent",
784 occurred_at=datetime(2026, 4, 1, 12, 0, tzinfo=UTC),
785 weight=1.0,
786 raw_payload_json=None,
787 ),
788 ]
789 )
790 db_session.commit()
791
792 deleted = poller.prune_old_events(db_session)
793 remaining = db_session.scalars(
794 select(ContributionEvent.external_uid).order_by(ContributionEvent.external_uid)
795 ).all()
796
797 assert deleted == 1
798 assert remaining == ["todo:recent"]
799
800
801def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None:
802 settings = make_settings(DISCOVERY_BATCH_SIZE=2, INDEXED_ACTOR_REPOLL_SECONDS=3600)
803 todo_service = RecordingTodoService(events_by_call=[[], []])
804 git_service = EmptyGitService()
805 poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings)
806
807 now = datetime.now(tz=UTC)
808 db_session.add_all(
809 [
810 TrackedActor(
811 actor="~a",
812 is_active=True,
813 discovery_state="queued",
814 queued_for_discovery_at=now - timedelta(minutes=3),
815 next_poll_after=now - timedelta(minutes=3),
816 recent_backfill_status="completed",
817 ),
818 TrackedActor(
819 actor="~b",
820 is_active=True,
821 discovery_state="queued",
822 queued_for_discovery_at=now - timedelta(minutes=2),
823 next_poll_after=now - timedelta(minutes=2),
824 recent_backfill_status="completed",
825 ),
826 TrackedActor(
827 actor="~c",
828 is_active=True,
829 discovery_state="queued",
830 queued_for_discovery_at=now - timedelta(minutes=1),
831 next_poll_after=now - timedelta(minutes=1),
832 recent_backfill_status="completed",
833 ),
834 ]
835 )
836 db_session.commit()
837
838 results = poller.poll_tracked_actors(db_session)
839
840 assert set(results) == {"~a", "~b"}
841 actors = {
842 actor.actor: actor
843 for actor in db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all()
844 }
845 assert actors["~a"].discovery_state == "indexed"
846 assert actors["~b"].discovery_state == "indexed"
847 assert actors["~c"].discovery_state == "queued"
848 assert actors["~a"].poll_attempts == 0
849 assert actors["~b"].poll_attempts == 0
850 assert actors["~c"].poll_attempts == 0
851
852
853def test_scheduled_error_backoff_is_capped_and_success_resets_attempts(db_session) -> None:
854 settings = make_settings(
855 DISCOVERY_BATCH_SIZE=1,
856 DISCOVERY_ERROR_BACKOFF_SECONDS=3600,
857 DISCOVERY_ERROR_BACKOFF_MAX_SECONDS=7200,
858 )
859 now = datetime.now(tz=UTC)
860 db_session.add(
861 TrackedActor(
862 actor="~flaky",
863 is_active=True,
864 discovery_state="queued",
865 queued_for_discovery_at=now - timedelta(minutes=1),
866 next_poll_after=now - timedelta(minutes=1),
867 poll_attempts=10,
868 recent_backfill_status="completed",
869 )
870 )
871 db_session.commit()
872 failing_poller = PollerService(todo_service=FailingTodoService(), git_service=EmptyGitService(), settings=settings)
873
874 failing_poller.poll_tracked_actors(db_session)
875
876 failed_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~flaky"))
877 assert failed_actor is not None
878 assert failed_actor.discovery_state == "error"
879 assert failed_actor.poll_attempts == 11
880 assert failed_actor.next_poll_after is not None
881 failed_next_poll_after = failed_actor.next_poll_after
882 if failed_next_poll_after.tzinfo is None:
883 failed_next_poll_after = failed_next_poll_after.replace(tzinfo=UTC)
884 assert failed_next_poll_after <= datetime.now(tz=UTC) + timedelta(seconds=7200, minutes=1)
885
886 failed_actor.next_poll_after = datetime.now(tz=UTC)
887 db_session.add(failed_actor)
888 db_session.commit()
889 successful_poller = PollerService(
890 todo_service=RecordingTodoService(events_by_call=[[]]),
891 git_service=EmptyGitService(),
892 settings=settings,
893 )
894
895 successful_poller.poll_tracked_actors(db_session)
896
897 recovered_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~flaky"))
898 assert recovered_actor is not None
899 assert recovered_actor.discovery_state == "indexed"
900 assert recovered_actor.poll_attempts == 0
901
902
903def test_track_actor_request_prioritize_marks_actor_boosted_and_due_now(db_session) -> None:
904 settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600)
905 poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings)
906 future_due = datetime.now(tz=UTC) + timedelta(hours=2)
907 db_session.add(
908 TrackedActor(
909 actor="~self",
910 is_active=True,
911 discovery_state="indexed",
912 queued_for_discovery_at=datetime.now(tz=UTC) - timedelta(hours=1),
913 next_poll_after=future_due,
914 recent_backfill_status="completed",
915 )
916 )
917 db_session.commit()
918
919 tracked_actor = poller.track_actor_request(db_session, "~self", prioritize=True)
920
921 assert tracked_actor.priority_boosted_at is not None
922 assert tracked_actor.next_poll_after is not None
923 assert tracked_actor.next_poll_after <= tracked_actor.priority_boosted_at
924
925
926def test_poll_tracked_actors_prioritizes_boosted_due_actor_first(db_session) -> None:
927 settings = make_settings(DISCOVERY_BATCH_SIZE=1, INDEXED_ACTOR_REPOLL_SECONDS=3600)
928 todo_service = RecordingTodoService(events_by_call=[[]])
929 poller = PollerService(todo_service=todo_service, git_service=EmptyGitService(), settings=settings)
930
931 now = datetime.now(tz=UTC)
932 db_session.add_all(
933 [
934 TrackedActor(
935 actor="~normal",
936 is_active=True,
937 discovery_state="queued",
938 queued_for_discovery_at=now - timedelta(minutes=10),
939 next_poll_after=now - timedelta(minutes=10),
940 recent_backfill_status="completed",
941 ),
942 TrackedActor(
943 actor="~self",
944 is_active=True,
945 discovery_state="queued",
946 queued_for_discovery_at=now - timedelta(minutes=1),
947 next_poll_after=now - timedelta(minutes=1),
948 priority_boosted_at=now,
949 recent_backfill_status="completed",
950 ),
951 ]
952 )
953 db_session.commit()
954
955 results = poller.poll_tracked_actors(db_session)
956
957 assert list(results) == ["~self"]
958 remaining = {
959 actor.actor: actor.discovery_state
960 for actor in db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all()
961 }
962 assert remaining["~self"] == "indexed"
963 assert remaining["~normal"] == "queued"
964
965
966def test_successful_poll_clears_temporary_priority_boost(db_session) -> None:
967 settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600)
968 poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings)
969 now = datetime.now(tz=UTC)
970 db_session.add(
971 TrackedActor(
972 actor="~self",
973 is_active=True,
974 discovery_state="queued",
975 queued_for_discovery_at=now - timedelta(minutes=1),
976 next_poll_after=now - timedelta(minutes=1),
977 priority_boosted_at=now - timedelta(seconds=30),
978 recent_backfill_status="completed",
979 )
980 )
981 db_session.commit()
982
983 poller.poll_all(db_session, "~self")
984
985 tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~self"))
986 assert tracked_actor is not None
987 assert tracked_actor.discovery_state == "indexed"
988 assert tracked_actor.priority_boosted_at is None
989
990
991def test_enqueue_actors_staggers_without_polling(tmp_path, monkeypatch) -> None:
992 database_path = tmp_path / "enqueue.db"
993 username_path = tmp_path / "srht_usernames.txt"
994 username_path.write_text("alice\nbob\nalice\n~carol\n", encoding="utf-8")
995 monkeypatch.setenv("DATABASE_URL", f"sqlite:///{database_path}")
996 monkeypatch.setenv("SRHT_TOKEN", "test-token")
997 monkeypatch.setenv("DEFAULT_ACTOR", "~ccleberg")
998
999 from srht_contrib.db import Base, make_engine, make_session_factory
1000
1001 settings = Settings()
1002 engine = make_engine(settings)
1003 Base.metadata.create_all(bind=engine)
1004 session_factory = make_session_factory(settings)
1005 queued_at = datetime(2026, 4, 11, 12, 0, tzinfo=UTC)
1006
1007 inserted = enqueue_actors(Path(username_path), stagger_seconds=60, start_at=queued_at)
1008
1009 with session_factory() as db:
1010 actors = db.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all()
1011
1012 assert inserted == 3
1013 assert [actor.actor for actor in actors] == ["~alice", "~bob", "~carol"]
1014 assert all(actor.discovery_state == "queued" for actor in actors)
1015 assert actors[0].next_poll_after == queued_at.replace(tzinfo=None)
1016 assert actors[1].next_poll_after == (queued_at + timedelta(seconds=60)).replace(tzinfo=None)
1017 assert actors[2].next_poll_after == (queued_at + timedelta(seconds=120)).replace(tzinfo=None)
1018
1019
1020def test_enqueue_actors_skips_invalid_usernames(tmp_path, monkeypatch) -> None:
1021 database_path = tmp_path / "enqueue-invalid.db"
1022 username_path = tmp_path / "srht_usernames.txt"
1023 username_path.write_text("-0\n.\n~bad-\nvalid_user\nok.ok\n", encoding="utf-8")
1024 monkeypatch.setenv("DATABASE_URL", f"sqlite:///{database_path}")
1025 monkeypatch.setenv("SRHT_TOKEN", "test-token")
1026 monkeypatch.setenv("DEFAULT_ACTOR", "~ccleberg")
1027
1028 from srht_contrib.db import Base, make_engine, make_session_factory
1029
1030 settings = Settings()
1031 engine = make_engine(settings)
1032 Base.metadata.create_all(bind=engine)
1033 session_factory = make_session_factory(settings)
1034
1035 inserted = enqueue_actors(Path(username_path), stagger_seconds=60)
1036
1037 with session_factory() as db:
1038 actors = db.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all()
1039
1040 assert inserted == 2
1041 assert [actor.actor for actor in actors] == ["~ok.ok", "~valid_user"]