diff options
| -rw-r--r-- | API.md | 5 | ||||
| -rw-r--r-- | README.md | 2 | ||||
| -rw-r--r-- | alembic/versions/20260412_0007_actor_priority_boost.py | 29 | ||||
| -rw-r--r-- | src/srht_contrib/api/routes_contributions.py | 6 | ||||
| -rw-r--r-- | src/srht_contrib/jobs/poller.py | 16 | ||||
| -rw-r--r-- | src/srht_contrib/models.py | 1 | ||||
| -rw-r--r-- | tests/test_contributions_api.py | 45 | ||||
| -rw-r--r-- | tests/test_ingestion.py | 88 | ||||
| -rw-r--r-- | tests/test_migrations.py | 9 | ||||
| -rw-r--r-- | tests/test_polling_api.py | 4 |
10 files changed, 199 insertions, 6 deletions
@@ -44,6 +44,7 @@ Contribution ranges: - Contribution read endpoints return zero-filled days, so clients do not need to patch missing dates. - Public contribution reads also register the actor for background indexing. A first lookup may therefore return an empty graph while the scheduler catches up. +- Clients may add `prioritize_self=true` on contribution read endpoints to explicitly request temporary indexing priority for the signed-in user's own graph. - Incremental indexing and one-year backfill are separate. An actor can be recently indexed before the retained one-year window is fully filled in. - The service only retains and backfills the most recent 365 days of activity. @@ -101,6 +102,7 @@ Query parameters: - `year` integer, optional - `from` string `YYYY-MM-DD`, optional - `to` string `YYYY-MM-DD`, optional +- `prioritize_self` boolean, optional Rules: @@ -112,6 +114,7 @@ Behavior notes: - This endpoint resolves aliases to a canonical actor before querying data. - This endpoint also registers the actor for background indexing and updates the actor's `last_requested_at` timestamp. +- 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. - The response is always immediate; it does not wait for SourceHut polling to finish. - One-year backfill runs in bounded background batches and may take multiple scheduler passes to complete. @@ -211,10 +214,12 @@ Query parameters: - `year` integer, optional - `from` string `YYYY-MM-DD`, optional - `to` string `YYYY-MM-DD`, optional +- `prioritize_self` boolean, optional Behavior notes: - This endpoint has the same actor-registration and alias-resolution behavior as the calendar endpoint. +- When `prioritize_self=true`, registration also applies the same temporary scheduler boost as the calendar endpoint. - This endpoint returns immediately and does not block on SourceHut polling. - This endpoint also reflects whether the retained one-year history window has been fully backfilled yet. @@ -175,6 +175,8 @@ Scheduled polling only runs when `ENABLE_SCHEDULER=true`. The scheduler seeds `D 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`. +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. + ## Bulk Enqueue Without Immediate Indexing To durably queue a large username list without polling it immediately: diff --git a/alembic/versions/20260412_0007_actor_priority_boost.py b/alembic/versions/20260412_0007_actor_priority_boost.py new file mode 100644 index 0000000..ecceca6 --- /dev/null +++ b/alembic/versions/20260412_0007_actor_priority_boost.py @@ -0,0 +1,29 @@ +"""temporary actor priority boost marker""" + +from __future__ import annotations + +from alembic import op +import sqlalchemy as sa +from sqlalchemy import inspect + + +revision = "20260412_0007" +down_revision = "20260411_0006" +branch_labels = None +depends_on = None + + +def _column_names(table_name: str) -> set[str]: + return {column["name"] for column in inspect(op.get_bind()).get_columns(table_name)} + + +def upgrade() -> None: + columns = _column_names("tracked_actors") + if "priority_boosted_at" not in columns: + op.add_column("tracked_actors", sa.Column("priority_boosted_at", sa.DateTime(timezone=True), nullable=True)) + + +def downgrade() -> None: + columns = _column_names("tracked_actors") + if "priority_boosted_at" in columns: + op.drop_column("tracked_actors", "priority_boosted_at") diff --git a/src/srht_contrib/api/routes_contributions.py b/src/srht_contrib/api/routes_contributions.py index 25bc3c2..82b1948 100644 --- a/src/srht_contrib/api/routes_contributions.py +++ b/src/srht_contrib/api/routes_contributions.py @@ -41,13 +41,14 @@ def get_contributions( year: int | None = Query(default=None, ge=1970, le=3000), from_date: str | None = Query(default=None, alias="from"), to_date: str | None = Query(default=None, alias="to"), + prioritize_self: bool = Query(default=False), poller: PollerService = Depends(get_poller), db: Session = Depends(get_db), actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), ) -> ContributionCalendarResponse: start, end = _resolve_range(year, from_date, to_date) canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) - poller.track_actor_request(db, canonical_actor) + poller.track_actor_request(db, canonical_actor, prioritize=prioritize_self) response = ContributionAggregator().build_calendar(db, canonical_actor, start, end) db.commit() return response @@ -59,13 +60,14 @@ def get_contribution_stats( year: int | None = Query(default=None, ge=1970, le=3000), from_date: str | None = Query(default=None, alias="from"), to_date: str | None = Query(default=None, alias="to"), + prioritize_self: bool = Query(default=False), poller: PollerService = Depends(get_poller), db: Session = Depends(get_db), actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), ) -> ContributionStatsResponse: start, end = _resolve_range(year, from_date, to_date) canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) - poller.track_actor_request(db, canonical_actor) + poller.track_actor_request(db, canonical_actor, prioritize=prioritize_self) response = ContributionAggregator().build_stats(db, canonical_actor, start, end) db.commit() return response diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py index 5e62b55..76d2fb0 100644 --- a/src/srht_contrib/jobs/poller.py +++ b/src/srht_contrib/jobs/poller.py @@ -65,6 +65,8 @@ class PollerService: ) ) .order_by( + TrackedActor.priority_boosted_at.is_not(None).desc(), + TrackedActor.priority_boosted_at.desc(), TrackedActor.next_poll_after.is_(None).desc(), TrackedActor.next_poll_after, TrackedActor.queued_for_discovery_at, @@ -90,7 +92,14 @@ class PollerService: 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: + def track_actor_request( + self, + db: Session, + actor: str, + *, + update_last_requested: bool = True, + prioritize: bool = False, + ) -> TrackedActor: tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) now = datetime.now(tz=UTC) if tracked_actor is None: @@ -111,6 +120,9 @@ class PollerService: tracked_actor.next_poll_after = now if update_last_requested: tracked_actor.last_requested_at = now + if prioritize: + tracked_actor.priority_boosted_at = now + tracked_actor.next_poll_after = now db.flush() return tracked_actor @@ -207,11 +219,13 @@ class PollerService: tracked_actor.discovery_state = "indexed" tracked_actor.last_polled_at = now tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds) + tracked_actor.priority_boosted_at = None elif status == "error": tracked_actor.discovery_state = "error" tracked_actor.next_poll_after = now + timedelta( seconds=self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1) ) + tracked_actor.priority_boosted_at = None db.add(tracked_actor) db.flush() diff --git a/src/srht_contrib/models.py b/src/srht_contrib/models.py index 441b6ec..401e945 100644 --- a/src/srht_contrib/models.py +++ b/src/srht_contrib/models.py @@ -69,6 +69,7 @@ class TrackedActor(Base): is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) discovery_state: Mapped[str] = mapped_column(String(32), nullable=False, default="queued") queued_for_discovery_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + priority_boosted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) next_poll_after: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) last_claimed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) poll_attempts: Mapped[int] = mapped_column(Integer, nullable=False, default=0) diff --git a/tests/test_contributions_api.py b/tests/test_contributions_api.py index e568feb..29b19fb 100644 --- a/tests/test_contributions_api.py +++ b/tests/test_contributions_api.py @@ -7,6 +7,11 @@ from srht_contrib.main import create_app from srht_contrib.models import ContributionEvent, TrackedActor +class _Closable: + def close(self) -> None: + return None + + def test_read_only_contribution_routes_are_public_and_write_routes_require_api_key(settings, db_engine, session_factory) -> None: app = create_app(settings, engine=db_engine, session_factory=session_factory) with TestClient(app) as open_client: @@ -72,6 +77,24 @@ def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_se assert tracked_actor is not None assert tracked_actor.is_active is True assert tracked_actor.last_requested_at is not None + assert tracked_actor.priority_boosted_at is None + + +class RecordingPriorityPoller: + def __init__(self) -> None: + service = type("Service", (), {"client": _Closable()})() + self.todo_service = service + self.git_service = service + self.calls: list[tuple[str, bool]] = [] + + def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True, prioritize: bool = False): + self.calls.append((actor, prioritize)) + + def poll_all(self, db, actor: str) -> int: + return 0 + + def poll_tracked_actors(self, db, default_actor: str | None = None) -> dict[str, int]: + return {} 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 assert response.status_code == 200 assert response.json()["actor"] == "~ccleberg" + + +def test_contribution_route_passes_explicit_priority_signal(settings, db_engine, session_factory) -> None: + poller = RecordingPriorityPoller() + app = create_app(settings, engine=db_engine, session_factory=session_factory, poller=poller) + + with TestClient(app) as client: + response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30&prioritize_self=true") + + assert response.status_code == 200 + assert poller.calls == [("~ccleberg", True)] + + +def test_contribution_stats_route_keeps_non_prioritized_registration_by_default(settings, db_engine, session_factory) -> None: + poller = RecordingPriorityPoller() + app = create_app(settings, engine=db_engine, session_factory=session_factory, poller=poller) + + with TestClient(app) as client: + response = client.get("/api/contributions/~ccleberg/stats?from=2026-03-28&to=2026-03-30") + + assert response.status_code == 200 + assert poller.calls == [("~ccleberg", False)] diff --git a/tests/test_ingestion.py b/tests/test_ingestion.py index 6a1a162..6c06359 100644 --- a/tests/test_ingestion.py +++ b/tests/test_ingestion.py @@ -711,6 +711,94 @@ def test_poll_tracked_actors_limits_to_due_batch_size(db_session) -> None: assert actors["~c"].poll_attempts == 0 +def test_track_actor_request_prioritize_marks_actor_boosted_and_due_now(db_session) -> None: + settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) + poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) + future_due = datetime.now(tz=UTC) + timedelta(hours=2) + db_session.add( + TrackedActor( + actor="~self", + is_active=True, + discovery_state="indexed", + queued_for_discovery_at=datetime.now(tz=UTC) - timedelta(hours=1), + next_poll_after=future_due, + recent_backfill_status="completed", + ) + ) + db_session.commit() + + tracked_actor = poller.track_actor_request(db_session, "~self", prioritize=True) + + assert tracked_actor.priority_boosted_at is not None + assert tracked_actor.next_poll_after is not None + assert tracked_actor.next_poll_after <= tracked_actor.priority_boosted_at + + +def test_poll_tracked_actors_prioritizes_boosted_due_actor_first(db_session) -> None: + settings = make_settings(DISCOVERY_BATCH_SIZE=1, INDEXED_ACTOR_REPOLL_SECONDS=3600) + todo_service = RecordingTodoService(events_by_call=[[]]) + poller = PollerService(todo_service=todo_service, git_service=EmptyGitService(), settings=settings) + + now = datetime.now(tz=UTC) + db_session.add_all( + [ + TrackedActor( + actor="~normal", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=10), + next_poll_after=now - timedelta(minutes=10), + recent_backfill_status="completed", + ), + TrackedActor( + actor="~self", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=1), + next_poll_after=now - timedelta(minutes=1), + priority_boosted_at=now, + recent_backfill_status="completed", + ), + ] + ) + db_session.commit() + + results = poller.poll_tracked_actors(db_session) + + assert list(results) == ["~self"] + remaining = { + actor.actor: actor.discovery_state + for actor in db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() + } + assert remaining["~self"] == "indexed" + assert remaining["~normal"] == "queued" + + +def test_successful_poll_clears_temporary_priority_boost(db_session) -> None: + settings = make_settings(INDEXED_ACTOR_REPOLL_SECONDS=3600) + poller = PollerService(todo_service=RecordingTodoService(events_by_call=[[]]), git_service=EmptyGitService(), settings=settings) + now = datetime.now(tz=UTC) + db_session.add( + TrackedActor( + actor="~self", + is_active=True, + discovery_state="queued", + queued_for_discovery_at=now - timedelta(minutes=1), + next_poll_after=now - timedelta(minutes=1), + priority_boosted_at=now - timedelta(seconds=30), + recent_backfill_status="completed", + ) + ) + db_session.commit() + + poller.poll_all(db_session, "~self") + + tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~self")) + assert tracked_actor is not None + assert tracked_actor.discovery_state == "indexed" + assert tracked_actor.priority_boosted_at is None + + def test_enqueue_actors_staggers_without_polling(tmp_path, monkeypatch) -> None: database_path = tmp_path / "enqueue.db" username_path = tmp_path / "srht_usernames.txt" diff --git a/tests/test_migrations.py b/tests/test_migrations.py index 3ce842b..e646ccf 100644 --- a/tests/test_migrations.py +++ b/tests/test_migrations.py @@ -122,7 +122,14 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None: assert "discovered_repositories" in inspector.get_table_names() assert "tracked_actors" in inspector.get_table_names() assert "service_backfill_states" in inspector.get_table_names() - assert {"discovery_state", "queued_for_discovery_at", "next_poll_after", "last_claimed_at", "poll_attempts"} <= tracked_actor_columns + assert { + "discovery_state", + "queued_for_discovery_at", + "priority_boosted_at", + "next_poll_after", + "last_claimed_at", + "poll_attempts", + } <= tracked_actor_columns def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: diff --git a/tests/test_polling_api.py b/tests/test_polling_api.py index b752524..94188da 100644 --- a/tests/test_polling_api.py +++ b/tests/test_polling_api.py @@ -21,7 +21,7 @@ class InsertingPoller: self.git_service = service self.tracked_poll_calls: list[str] = [] - def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True): + def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True, prioritize: bool = False): tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) if tracked_actor is None: tracked_actor = TrackedActor(actor=actor, is_active=True) @@ -61,7 +61,7 @@ class FailingPoller: self.todo_service = service self.git_service = service - def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True): + def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True, prioritize: bool = False): tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) if tracked_actor is None: tracked_actor = TrackedActor(actor=actor, is_active=True) |
