summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorChristian Cleberg <[email protected]>2026-04-12 11:43:27 -0500
committerChristian Cleberg <[email protected]>2026-04-12 11:43:27 -0500
commit4390684995c9933f213faec2cb20f39af6d506d2 (patch)
treefac5252756c038e9fdb3444a1a4d68863e1683f3 /src
parentdaaeba8be8df4d0876312e540bdd420e16ff724c (diff)
downloadhutch-stats-4390684995c9933f213faec2cb20f39af6d506d2.tar.gz
hutch-stats-4390684995c9933f213faec2cb20f39af6d506d2.tar.bz2
hutch-stats-4390684995c9933f213faec2cb20f39af6d506d2.zip
add temporary self-graph indexing priority
Diffstat (limited to 'src')
-rw-r--r--src/srht_contrib/api/routes_contributions.py6
-rw-r--r--src/srht_contrib/jobs/poller.py16
-rw-r--r--src/srht_contrib/models.py1
3 files changed, 20 insertions, 3 deletions
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)