Commit 4390684995
Verified · cmc
Layout: unified · split
API.md +5
| @@ -44,6 +44,7 @@ Contribution ranges: | ||
| 44 | 44 | |
| 45 | 45 | - Contribution read endpoints return zero-filled days, so clients do not need to patch missing dates. |
| 46 | 46 | - Public contribution reads also register the actor for background indexing. A first lookup may therefore return an empty graph while the scheduler catches up. |
| 47 | - Clients may add `prioritize_self=true` on contribution read endpoints to explicitly request temporary indexing priority for the signed-in user's own graph. | |
| 47 | 48 | - Incremental indexing and one-year backfill are separate. An actor can be recently indexed before the retained one-year window is fully filled in. |
| 48 | 49 | - The service only retains and backfills the most recent 365 days of activity. |
| 49 | 50 | |
| @@ -101,6 +102,7 @@ Query parameters: | ||
| 101 | 102 | - `year` integer, optional |
| 102 | 103 | - `from` string `YYYY-MM-DD`, optional |
| 103 | 104 | - `to` string `YYYY-MM-DD`, optional |
| 105 | - `prioritize_self` boolean, optional | |
| 104 | 106 | |
| 105 | 107 | Rules: |
| 106 | 108 | |
| @@ -112,6 +114,7 @@ Behavior notes: | ||
| 112 | 114 | |
| 113 | 115 | - This endpoint resolves aliases to a canonical actor before querying data. |
| 114 | 116 | - This endpoint also registers the actor for background indexing and updates the actor's `last_requested_at` timestamp. |
| 117 | - When `prioritize_self=true`, registration also applies a temporary scheduler boost so that due polls for that actor run ahead of the normal due queue. | |
| 115 | 118 | - The response is always immediate; it does not wait for SourceHut polling to finish. |
| 116 | 119 | - One-year backfill runs in bounded background batches and may take multiple scheduler passes to complete. |
| 117 | 120 | |
| @@ -211,10 +214,12 @@ Query parameters: | ||
| 211 | 214 | - `year` integer, optional |
| 212 | 215 | - `from` string `YYYY-MM-DD`, optional |
| 213 | 216 | - `to` string `YYYY-MM-DD`, optional |
| 217 | - `prioritize_self` boolean, optional | |
| 214 | 218 | |
| 215 | 219 | Behavior notes: |
| 216 | 220 | |
| 217 | 221 | - This endpoint has the same actor-registration and alias-resolution behavior as the calendar endpoint. |
| 222 | - When `prioritize_self=true`, registration also applies the same temporary scheduler boost as the calendar endpoint. | |
| 218 | 223 | - This endpoint returns immediately and does not block on SourceHut polling. |
| 219 | 224 | - This endpoint also reflects whether the retained one-year history window has been fully backfilled yet. |
| 220 | 225 | |
README.md +2
| @@ -175,6 +175,8 @@ Scheduled polling only runs when `ENABLE_SCHEDULER=true`. The scheduler seeds `D | ||
| 175 | 175 | |
| 176 | 176 | The scheduler now drains actors gradually instead of polling every tracked actor on every pass. It only claims due actors, up to `DISCOVERY_BATCH_SIZE` per run, then reschedules indexed actors with `INDEXED_ACTOR_REPOLL_SECONDS` and failed actors with backoff based on `DISCOVERY_ERROR_BACKOFF_SECONDS`. |
| 177 | 177 | |
| 178 | Clients can explicitly signal that a public contribution read is for the signed-in user's own graph by sending `prioritize_self=true` on the read request. That temporarily boosts the actor to the front of the due queue for the next indexing pass, then clears the boost after the poll completes. | |
| 179 | ||
| 178 | 180 | ## Bulk Enqueue Without Immediate Indexing |
| 179 | 181 | |
| 180 | 182 | To durably queue a large username list without polling it immediately: |
alembic/versions/20260412_0007_actor_priority_boost.py added +29
| @@ -0,0 +1,29 @@ | ||
| 1 | """temporary actor priority boost marker""" | |
| 2 | ||
| 3 | from __future__ import annotations | |
| 4 | ||
| 5 | from alembic import op | |
| 6 | import sqlalchemy as sa | |
| 7 | from sqlalchemy import inspect | |
| 8 | ||
| 9 | ||
| 10 | revision = "20260412_0007" | |
| 11 | down_revision = "20260411_0006" | |
| 12 | branch_labels = None | |
| 13 | depends_on = None | |
| 14 | ||
| 15 | ||
| 16 | def _column_names(table_name: str) -> set[str]: | |
| 17 | return {column["name"] for column in inspect(op.get_bind()).get_columns(table_name)} | |
| 18 | ||
| 19 | ||
| 20 | def upgrade() -> None: | |
| 21 | columns = _column_names("tracked_actors") | |
| 22 | if "priority_boosted_at" not in columns: | |
| 23 | op.add_column("tracked_actors", sa.Column("priority_boosted_at", sa.DateTime(timezone=True), nullable=True)) | |
| 24 | ||
| 25 | ||
| 26 | def downgrade() -> None: | |
| 27 | columns = _column_names("tracked_actors") | |
| 28 | if "priority_boosted_at" in columns: | |
| 29 | op.drop_column("tracked_actors", "priority_boosted_at") | |
src/srht_contrib/api/routes_contributions.py +4 −2
| @@ -41,13 +41,14 @@ def get_contributions( | ||
| 41 | 41 | year: int | None = Query(default=None, ge=1970, le=3000), |
| 42 | 42 | from_date: str | None = Query(default=None, alias="from"), |
| 43 | 43 | to_date: str | None = Query(default=None, alias="to"), |
| 44 | prioritize_self: bool = Query(default=False), | |
| 44 | 45 | poller: PollerService = Depends(get_poller), |
| 45 | 46 | db: Session = Depends(get_db), |
| 46 | 47 | actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), |
| 47 | 48 | ) -> ContributionCalendarResponse: |
| 48 | 49 | start, end = _resolve_range(year, from_date, to_date) |
| 49 | 50 | canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) |
| 50 | poller.track_actor_request(db, canonical_actor) | |
| 51 | poller.track_actor_request(db, canonical_actor, prioritize=prioritize_self) | |
| 51 | 52 | response = ContributionAggregator().build_calendar(db, canonical_actor, start, end) |
| 52 | 53 | db.commit() |
| 53 | 54 | return response |
| @@ -59,13 +60,14 @@ def get_contribution_stats( | ||
| 59 | 60 | year: int | None = Query(default=None, ge=1970, le=3000), |
| 60 | 61 | from_date: str | None = Query(default=None, alias="from"), |
| 61 | 62 | to_date: str | None = Query(default=None, alias="to"), |
| 63 | prioritize_self: bool = Query(default=False), | |
| 62 | 64 | poller: PollerService = Depends(get_poller), |
| 63 | 65 | db: Session = Depends(get_db), |
| 64 | 66 | actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), |
| 65 | 67 | ) -> ContributionStatsResponse: |
| 66 | 68 | start, end = _resolve_range(year, from_date, to_date) |
| 67 | 69 | canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) |
| 68 | poller.track_actor_request(db, canonical_actor) | |
| 70 | poller.track_actor_request(db, canonical_actor, prioritize=prioritize_self) | |
| 69 | 71 | response = ContributionAggregator().build_stats(db, canonical_actor, start, end) |
| 70 | 72 | db.commit() |
| 71 | 73 | return response |
src/srht_contrib/jobs/poller.py +15 −1
| @@ -65,6 +65,8 @@ class PollerService: | ||
| 65 | 65 | ) |
| 66 | 66 | ) |
| 67 | 67 | .order_by( |
| 68 | TrackedActor.priority_boosted_at.is_not(None).desc(), | |
| 69 | TrackedActor.priority_boosted_at.desc(), | |
| 68 | 70 | TrackedActor.next_poll_after.is_(None).desc(), |
| 69 | 71 | TrackedActor.next_poll_after, |
| 70 | 72 | TrackedActor.queued_for_discovery_at, |
| @@ -90,7 +92,14 @@ class PollerService: | ||
| 90 | 92 | logger.info("Pruned %s contribution events older than %s days", deleted, RETENTION_DAYS) |
| 91 | 93 | return results |
| 92 | 94 | |
| 93 | def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: | |
| 95 | def track_actor_request( | |
| 96 | self, | |
| 97 | db: Session, | |
| 98 | actor: str, | |
| 99 | *, | |
| 100 | update_last_requested: bool = True, | |
| 101 | prioritize: bool = False, | |
| 102 | ) -> TrackedActor: | |
| 94 | 103 | tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) |
| 95 | 104 | now = datetime.now(tz=UTC) |
| 96 | 105 | if tracked_actor is None: |
| @@ -111,6 +120,9 @@ class PollerService: | ||
| 111 | 120 | tracked_actor.next_poll_after = now |
| 112 | 121 | if update_last_requested: |
| 113 | 122 | tracked_actor.last_requested_at = now |
| 123 | if prioritize: | |
| 124 | tracked_actor.priority_boosted_at = now | |
| 125 | tracked_actor.next_poll_after = now | |
| 114 | 126 | db.flush() |
| 115 | 127 | return tracked_actor |
| 116 | 128 | |
| @@ -207,11 +219,13 @@ class PollerService: | ||
| 207 | 219 | tracked_actor.discovery_state = "indexed" |
| 208 | 220 | tracked_actor.last_polled_at = now |
| 209 | 221 | tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) |
| 222 | tracked_actor.priority_boosted_at = None | |
| 210 | 223 | elif status == "error": |
| 211 | 224 | tracked_actor.discovery_state = "error" |
| 212 | 225 | tracked_actor.next_poll_after = now + timedelta( |
| 213 | 226 | seconds=self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1) |
| 214 | 227 | ) |
| 228 | tracked_actor.priority_boosted_at = None | |
| 215 | 229 | db.add(tracked_actor) |
| 216 | 230 | db.flush() |
| 217 | 231 | |
src/srht_contrib/models.py +1
| @@ -69,6 +69,7 @@ class TrackedActor(Base): | ||
| 69 | 69 | is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) |
| 70 | 70 | discovery_state: Mapped[str] = mapped_column(String(32), nullable=False, default="queued") |
| 71 | 71 | queued_for_discovery_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) |
| 72 | priority_boosted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) | |
| 72 | 73 | next_poll_after: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) |
| 73 | 74 | last_claimed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) |
| 74 | 75 | poll_attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0) |
tests/test_contributions_api.py +45
| @@ -7,6 +7,11 @@ from srht_contrib.main import create_app | ||
| 7 | 7 | from srht_contrib.models import ContributionEvent, TrackedActor |
| 8 | 8 | |
| 9 | 9 | |
| 10 | class _Closable: | |
| 11 | def close(self) -> None: | |
| 12 | return None | |
| 13 | ||
| 14 | ||
| 10 | 15 | def test_read_only_contribution_routes_are_public_and_write_routes_require_api_key(settings, db_engine, session_factory) -> None: |
| 11 | 16 | app = create_app(settings, engine=db_engine, session_factory=session_factory) |
| 12 | 17 | with TestClient(app) as open_client: |
| @@ -72,6 +77,24 @@ def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_se | ||
| 72 | 77 | assert tracked_actor is not None |
| 73 | 78 | assert tracked_actor.is_active is True |
| 74 | 79 | assert tracked_actor.last_requested_at is not None |
| 80 | assert tracked_actor.priority_boosted_at is None | |
| 81 | ||
| 82 | ||
| 83 | class RecordingPriorityPoller: | |
| 84 | def __init__(self) -> None: | |
| 85 | service = type("Service", (), {"client": _Closable()})() | |
| 86 | self.todo_service = service | |
| 87 | self.git_service = service | |
| 88 | self.calls: list[tuple[str, bool]] = [] | |
| 89 | ||
| 90 | def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True, prioritize: bool = False): | |
| 91 | self.calls.append((actor, prioritize)) | |
| 92 | ||
| 93 | def poll_all(self, db, actor: str) -> int: | |
| 94 | return 0 | |
| 95 | ||
| 96 | def poll_tracked_actors(self, db, default_actor: str | None = None) -> dict[str, int]: | |
| 97 | return {} | |
| 75 | 98 | |
| 76 | 99 | |
| 77 | 100 | def test_contribution_stats_api(client: TestClient, db_session) -> None: |
| @@ -148,3 +171,25 @@ def test_contribution_routes_use_settings_backed_alias_resolution(settings, db_e | ||
| 148 | 171 | |
| 149 | 172 | assert response.status_code == 200 |
| 150 | 173 | assert response.json()["actor"] == "~ccleberg" |
| 174 | ||
| 175 | ||
| 176 | def test_contribution_route_passes_explicit_priority_signal(settings, db_engine, session_factory) -> None: | |
| 177 | poller = RecordingPriorityPoller() | |
| 178 | app = create_app(settings, engine=db_engine, session_factory=session_factory, poller=poller) | |
| 179 | ||
| 180 | with TestClient(app) as client: | |
| 181 | response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30&prioritize_self=true") | |
| 182 | ||
| 183 | assert response.status_code == 200 | |
| 184 | assert poller.calls == [("~ccleberg", True)] | |
| 185 | ||
| 186 | ||
| 187 | def test_contribution_stats_route_keeps_non_prioritized_registration_by_default(settings, db_engine, session_factory) -> None: | |
| 188 | poller = RecordingPriorityPoller() | |
| 189 | app = create_app(settings, engine=db_engine, session_factory=session_factory, poller=poller) | |
| 190 | ||
| 191 | with TestClient(app) as client: | |
| 192 | response = client.get("/api/contributions/~ccleberg/stats?from=2026-03-28&to=2026-03-30") | |
| 193 | ||
| 194 | assert response.status_code == 200 | |
| 195 | assert poller.calls == [("~ccleberg", False)] | |
tests/test_ingestion.py +88
| @@ -711,6 +711,94 @@ def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None: | ||
| 711 | 711 | assert actors["~c"].poll_attempts == 0 |
| 712 | 712 | |
| 713 | 713 | |
| 714 | def test_track_actor_request_prioritize_marks_actor_boosted_and_due_now(db_session) -> None: | |
| 715 | settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) | |
| 716 | poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) | |
| 717 | future_due = datetime.now(tz=UTC) + timedelta(hours=2) | |
| 718 | db_session.add( | |
| 719 | TrackedActor( | |
| 720 | actor="~self", | |
| 721 | is_active=True, | |
| 722 | discovery_state="indexed", | |
| 723 | queued_for_discovery_at=datetime.now(tz=UTC) - timedelta(hours=1), | |
| 724 | next_poll_after=future_due, | |
| 725 | recent_backfill_status="completed", | |
| 726 | ) | |
| 727 | ) | |
| 728 | db_session.commit() | |
| 729 | ||
| 730 | tracked_actor = poller.track_actor_request(db_session, "~self", prioritize=True) | |
| 731 | ||
| 732 | assert tracked_actor.priority_boosted_at is not None | |
| 733 | assert tracked_actor.next_poll_after is not None | |
| 734 | assert tracked_actor.next_poll_after <= tracked_actor.priority_boosted_at | |
| 735 | ||
| 736 | ||
| 737 | def test_poll_tracked_actors_prioritizes_boosted_due_actor_first(db_session) -> None: | |
| 738 | settings = make_settings(DISCOVERY_BATCH_SIZE=1, INDEXED_ACTOR_REPOLL_SECONDS=3600) | |
| 739 | todo_service = RecordingTodoService(events_by_call=[[]]) | |
| 740 | poller = PollerService(todo_service=todo_service, git_service=EmptyGitService(), settings=settings) | |
| 741 | ||
| 742 | now = datetime.now(tz=UTC) | |
| 743 | db_session.add_all( | |
| 744 | [ | |
| 745 | TrackedActor( | |
| 746 | actor="~normal", | |
| 747 | is_active=True, | |
| 748 | discovery_state="queued", | |
| 749 | queued_for_discovery_at=now - timedelta(minutes=10), | |
| 750 | next_poll_after=now - timedelta(minutes=10), | |
| 751 | recent_backfill_status="completed", | |
| 752 | ), | |
| 753 | TrackedActor( | |
| 754 | actor="~self", | |
| 755 | is_active=True, | |
| 756 | discovery_state="queued", | |
| 757 | queued_for_discovery_at=now - timedelta(minutes=1), | |
| 758 | next_poll_after=now - timedelta(minutes=1), | |
| 759 | priority_boosted_at=now, | |
| 760 | recent_backfill_status="completed", | |
| 761 | ), | |
| 762 | ] | |
| 763 | ) | |
| 764 | db_session.commit() | |
| 765 | ||
| 766 | results = poller.poll_tracked_actors(db_session) | |
| 767 | ||
| 768 | assert list(results) == ["~self"] | |
| 769 | remaining = { | |
| 770 | actor.actor: actor.discovery_state | |
| 771 | for actor in db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() | |
| 772 | } | |
| 773 | assert remaining["~self"] == "indexed" | |
| 774 | assert remaining["~normal"] == "queued" | |
| 775 | ||
| 776 | ||
| 777 | def test_successful_poll_clears_temporary_priority_boost(db_session) -> None: | |
| 778 | settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) | |
| 779 | poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) | |
| 780 | now = datetime.now(tz=UTC) | |
| 781 | db_session.add( | |
| 782 | TrackedActor( | |
| 783 | actor="~self", | |
| 784 | is_active=True, | |
| 785 | discovery_state="queued", | |
| 786 | queued_for_discovery_at=now - timedelta(minutes=1), | |
| 787 | next_poll_after=now - timedelta(minutes=1), | |
| 788 | priority_boosted_at=now - timedelta(seconds=30), | |
| 789 | recent_backfill_status="completed", | |
| 790 | ) | |
| 791 | ) | |
| 792 | db_session.commit() | |
| 793 | ||
| 794 | poller.poll_all(db_session, "~self") | |
| 795 | ||
| 796 | tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~self")) | |
| 797 | assert tracked_actor is not None | |
| 798 | assert tracked_actor.discovery_state == "indexed" | |
| 799 | assert tracked_actor.priority_boosted_at is None | |
| 800 | ||
| 801 | ||
| 714 | 802 | def test_enqueue_actors_staggers_without_polling(tmp_path, monkeypatch) -> None: |
| 715 | 803 | database_path = tmp_path / "enqueue.db" |
| 716 | 804 | username_path = tmp_path / "srht_usernames.txt" |
tests/test_migrations.py +8 −1
| @@ -122,7 +122,14 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None: | ||
| 122 | 122 | assert "discovered_repositories" in inspector.get_table_names() |
| 123 | 123 | assert "tracked_actors" in inspector.get_table_names() |
| 124 | 124 | assert "service_backfill_states" in inspector.get_table_names() |
| 125 | assert {"discovery_state", "queued_for_discovery_at", "next_poll_after", "last_claimed_at", "poll_attempts"} <= tracked_actor_columns | |
| 125 | assert { | |
| 126 | "discovery_state", | |
| 127 | "queued_for_discovery_at", | |
| 128 | "priority_boosted_at", | |
| 129 | "next_poll_after", | |
| 130 | "last_claimed_at", | |
| 131 | "poll_attempts", | |
| 132 | } <= tracked_actor_columns | |
| 126 | 133 | |
| 127 | 134 | |
| 128 | 135 | def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: |
tests/test_polling_api.py +2 −2
| @@ -21,7 +21,7 @@ class InsertingPoller: | ||
| 21 | 21 | self.git_service = service |
| 22 | 22 | self.tracked_poll_calls: list[str] = [] |
| 23 | 23 | |
| 24 | def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True): | |
| 24 | def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True, prioritize: bool = False): | |
| 25 | 25 | tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) |
| 26 | 26 | if tracked_actor is None: |
| 27 | 27 | tracked_actor = TrackedActor(actor=actor, is_active=True) |
| @@ -61,7 +61,7 @@ class FailingPoller: | ||
| 61 | 61 | self.todo_service = service |
| 62 | 62 | self.git_service = service |
| 63 | 63 | |
| 64 | def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True): | |
| 64 | def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True, prioritize: bool = False): | |
| 65 | 65 | tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) |
| 66 | 66 | if tracked_actor is None: |
| 67 | 67 | tracked_actor = TrackedActor(actor=actor, is_active=True) |