| @@ -83,6 +83,52 @@ class BackfillingTodoService: |
| 83 | 83 | return BackfillBatchResult(events=[event], cursor_state=None, complete=True) |
| 84 | 84 | |
| 85 | 85 | |
| 86 | class QueueShrinkingTodoService: |
| 87 | service_name = "todo" |
| 88 | |
| 89 | def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult: |
| 90 | return TodoPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat()) |
| 91 | |
| 92 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: |
| 93 | import copy |
| 94 | |
| 95 | state = {"tracker_queue": ["t1", "t2"], "current_tracker": None, "current_ticket": None, "trackers_loaded": True, "trackers_cursor": None} |
| 96 | if cursor_state: |
| 97 | state.update(copy.deepcopy(cursor_state)) |
| 98 | if not state["tracker_queue"]: |
| 99 | return BackfillBatchResult(events=[], cursor_state=None, complete=True) |
| 100 | state["tracker_queue"].pop(0) |
| 101 | return BackfillBatchResult(events=[], cursor_state=state, complete=False) |
| 102 | |
| 103 | |
| 104 | class QueueShrinkingGitService: |
| 105 | service_name = "git" |
| 106 | |
| 107 | def __init__(self) -> None: |
| 108 | self.settings = Settings( |
| 109 | SRHT_TOKEN="x", |
| 110 | DATABASE_URL="sqlite://", |
| 111 | DEFAULT_ACTOR="~ccleberg", |
| 112 | TODO_SRHT_ENDPOINT="https://todo.sr.ht/query", |
| 113 | GIT_SRHT_ENDPOINT="https://git.sr.ht/query", |
| 114 | POLL_INTERVAL_SECONDS=60, |
| 115 | ) |
| 116 | |
| 117 | def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult: |
| 118 | return GitPollResult(events=[], cursor=datetime(2026, 3, 31, tzinfo=UTC).isoformat()) |
| 119 | |
| 120 | def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: |
| 121 | import copy |
| 122 | |
| 123 | state = {"repository_queue": ["r1", "r2"], "current_repository": None, "discovery_complete": True, "discovery_cursor": None} |
| 124 | if cursor_state: |
| 125 | state.update(copy.deepcopy(cursor_state)) |
| 126 | if not state["repository_queue"]: |
| 127 | return BackfillBatchResult(events=[], cursor_state=None, complete=True) |
| 128 | state["repository_queue"].pop(0) |
| 129 | return BackfillBatchResult(events=[], cursor_state=state, complete=False) |
| 130 | |
| 131 | |
| 86 | 132 | def make_settings(**overrides) -> Settings: |
| 87 | 133 | values = { |
| 88 | 134 | "API_KEY": "test-api-key", |
| @@ -470,3 +516,28 @@ def test_poll_marks_backfill_complete_and_persists_service_state(db_session) -> |
| 470 | 516 | assert tracked_actor.backfill_completed_at is not None |
| 471 | 517 | assert [state.service for state in service_states] == ["git", "todo"] |
| 472 | 518 | assert all(state.status == "completed" for state in service_states) |
| 519 | |
| 520 | |
| 521 | def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None: |
| 522 | poller = PollerService(todo_service=QueueShrinkingTodoService(), git_service=QueueShrinkingGitService()) |
| 523 | |
| 524 | poller.poll_all(db_session, "~ccleberg") |
| 525 | first_states = { |
| 526 | state.service: state.cursor_json |
| 527 | for state in db_session.scalars( |
| 528 | select(ServiceBackfillState).where(ServiceBackfillState.actor == "~ccleberg") |
| 529 | ).all() |
| 530 | } |
| 531 | |
| 532 | poller.poll_all(db_session, "~ccleberg") |
| 533 | second_states = { |
| 534 | state.service: state.cursor_json |
| 535 | for state in db_session.scalars( |
| 536 | select(ServiceBackfillState).where(ServiceBackfillState.actor == "~ccleberg") |
| 537 | ).all() |
| 538 | } |
| 539 | |
| 540 | assert first_states["git"]["repository_queue"] == ["r2"] |
| 541 | assert first_states["todo"]["tracker_queue"] == ["t2"] |
| 542 | assert second_states["git"]["repository_queue"] == [] |
| 543 | assert second_states["todo"]["tracker_queue"] == [] |