summaryrefslogtreecommitdiff
path: root/src/srht_contrib
diff options
context:
space:
mode:
Diffstat (limited to 'src/srht_contrib')
-rw-r--r--src/srht_contrib/config.py5
-rw-r--r--src/srht_contrib/jobs/poller.py56
-rw-r--r--src/srht_contrib/scripts/enqueue_actors.py4
3 files changed, 54 insertions, 11 deletions
diff --git a/src/srht_contrib/config.py b/src/srht_contrib/config.py
index 3a7af2c..9221e4c 100644
--- a/src/srht_contrib/config.py
+++ b/src/srht_contrib/config.py
@@ -34,13 +34,14 @@ class Settings(BaseSettings):
alias="DATABASE_URL",
)
default_actor: str = Field(default="~unknown", alias="DEFAULT_ACTOR")
- poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS")
+ poll_interval_seconds: int = Field(default=300, alias="POLL_INTERVAL_SECONDS")
sync_overlap_hours: int = Field(default=1, alias="SYNC_OVERLAP_HOURS")
srht_request_delay_seconds: float = Field(default=0.5, alias="SRHT_REQUEST_DELAY_SECONDS")
sqlite_busy_timeout_seconds: float = Field(default=30.0, alias="SQLITE_BUSY_TIMEOUT_SECONDS")
- discovery_batch_size: int = Field(default=5, alias="DISCOVERY_BATCH_SIZE")
+ discovery_batch_size: int = Field(default=20, alias="DISCOVERY_BATCH_SIZE")
indexed_actor_repoll_seconds: int = Field(default=21600, alias="INDEXED_ACTOR_REPOLL_SECONDS")
discovery_error_backoff_seconds: int = Field(default=3600, alias="DISCOVERY_ERROR_BACKOFF_SECONDS")
+ discovery_error_backoff_max_seconds: int = Field(default=21600, alias="DISCOVERY_ERROR_BACKOFF_MAX_SECONDS")
git_repo_discovery_ttl_seconds: int = Field(default=3600, alias="GIT_REPO_DISCOVERY_TTL_SECONDS")
actor_aliases_json: dict[str, list[str]] = Field(
default_factory=dict,
diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py
index 76d2fb0..3e6772b 100644
--- a/src/srht_contrib/jobs/poller.py
+++ b/src/srht_contrib/jobs/poller.py
@@ -34,7 +34,7 @@ class PollerService:
self.settings = settings
self._sync_overlap = timedelta(hours=settings.sync_overlap_hours)
- def poll_all(self, db: Session, actor: str) -> int:
+ def poll_all(self, db: Session, actor: str, *, run_backfill: bool = True) -> int:
self.track_actor_request(db, actor, update_last_requested=False)
try:
inserted = self._poll_actor(db, actor)
@@ -45,7 +45,8 @@ class PollerService:
raise
self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None)
- inserted += self._run_backfill_batches(db, actor)
+ if run_backfill:
+ inserted += self._run_backfill_batches(db, actor)
db.commit()
return inserted
@@ -56,7 +57,7 @@ class PollerService:
results: dict[str, int] = {}
actors = db.scalars(
- select(TrackedActor.actor)
+ select(TrackedActor)
.where(TrackedActor.is_active.is_(True))
.where(
or_(
@@ -74,7 +75,12 @@ class PollerService:
)
.limit(self.settings.discovery_batch_size)
).all()
- for actor in actors:
+ for tracked_actor in actors:
+ actor = tracked_actor.actor
+ run_backfill_only = (
+ tracked_actor.last_polled_at is not None
+ and tracked_actor.recent_backfill_status != "completed"
+ )
claimed_actor = self.track_actor_request(db, actor, update_last_requested=False)
claimed_actor.discovery_state = "in_progress"
claimed_actor.last_claimed_at = datetime.now(tz=UTC)
@@ -82,7 +88,12 @@ class PollerService:
db.add(claimed_actor)
db.commit()
try:
- results[actor] = self.poll_all(db, actor)
+ if run_backfill_only:
+ results[actor] = self.poll_recent_backfill(db, actor)
+ else:
+ results[actor] = self.poll_all(db, actor, run_backfill=False)
+ self._schedule_pending_backfill_or_repoll(db, actor)
+ db.commit()
except SourceHutClientError:
logger.exception("Scheduled poll failed for actor=%s", actor)
except Exception:
@@ -220,15 +231,46 @@ class PollerService:
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
+ tracked_actor.poll_attempts = 0
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)
+ backoff_seconds = min(
+ self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1),
+ self.settings.discovery_error_backoff_max_seconds,
)
+ tracked_actor.next_poll_after = now + timedelta(seconds=backoff_seconds)
tracked_actor.priority_boosted_at = None
db.add(tracked_actor)
db.flush()
+ def poll_recent_backfill(self, db: Session, actor: str) -> int:
+ try:
+ inserted = self._run_backfill_batches(db, actor)
+ except Exception as exc:
+ db.rollback()
+ self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc))
+ db.commit()
+ raise
+
+ self._schedule_pending_backfill_or_repoll(db, actor)
+ db.commit()
+ return inserted
+
+ def _schedule_pending_backfill_or_repoll(self, db: Session, actor: str) -> None:
+ tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
+ now = datetime.now(tz=UTC)
+ tracked_actor.discovery_state = "indexed"
+ tracked_actor.last_poll_status = "indexed"
+ tracked_actor.last_poll_error = None
+ tracked_actor.priority_boosted_at = None
+ tracked_actor.poll_attempts = 0
+ if tracked_actor.recent_backfill_status == "completed":
+ tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds)
+ else:
+ tracked_actor.next_poll_after = now
+ db.add(tracked_actor)
+ db.flush()
+
def _run_backfill_batches(self, db: Session, actor: str) -> int:
tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
if tracked_actor.recent_backfill_status == "completed":
diff --git a/src/srht_contrib/scripts/enqueue_actors.py b/src/srht_contrib/scripts/enqueue_actors.py
index 03f1f9b..bd838cb 100644
--- a/src/srht_contrib/scripts/enqueue_actors.py
+++ b/src/srht_contrib/scripts/enqueue_actors.py
@@ -40,7 +40,7 @@ def _iter_usernames(path: Path) -> list[str]:
return usernames
-def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: datetime | None = None) -> int:
+def enqueue_actors(username_file: Path, *, stagger_seconds: int = 60, start_at: datetime | None = None) -> int:
settings = Settings()
session_factory = make_session_factory(settings)
usernames = _iter_usernames(username_file)
@@ -95,7 +95,7 @@ def main() -> None:
parser.add_argument(
"--stagger-seconds",
type=int,
- default=300,
+ default=60,
help="Seconds to space out each actor's first eligible poll time.",
)
args = parser.parse_args()