fix: deriver ownership

This commit is contained in:
Rajat Ahuja 2026-01-30 16:54:21 -05:00
parent 87bae0dd86
commit 2c90974ea3
2 changed files with 14 additions and 3 deletions

View File

@ -364,8 +364,11 @@ class QueueManager:
continue
try:
logger.info("POLL_STEP cleanup_stale_start")
await self.cleanup_stale_work_units()
logger.info("POLL_STEP get_and_claim_start")
claimed_work_units = await self.get_and_claim_work_units()
logger.info("POLL_STEP get_and_claim_done claimed=%d", len(claimed_work_units))
if claimed_work_units:
for work_unit_key, aqs_id in claimed_work_units.items():
# Create a new task for processing this work unit
@ -386,13 +389,16 @@ class QueueManager:
settings.DERIVER.POLLING_SLEEP_INTERVAL_SECONDS
)
except Exception as e:
logger.exception("Error in polling loop")
logger.exception("POLL_ERROR in polling loop")
if settings.SENTRY.ENABLED:
sentry_sdk.capture_exception(e)
# Note: rollback is handled by tracked_db dependency
await asyncio.sleep(settings.DERIVER.POLLING_SLEEP_INTERVAL_SECONDS)
except Exception:
logger.exception("POLL_LOOP_CRASHED unexpected exception escaped the loop")
raise
finally:
logger.info("Polling loop stopped")
logger.info("POLL_LOOP_EXIT polling loop stopped")
######################
# Queue Worker Logic #

View File

@ -161,11 +161,16 @@ class ReconcilerScheduler:
try:
enqueued = await self._try_enqueue_task(task)
if enqueued:
logger.debug("Enqueued task: %s", task_name)
logger.info("Enqueued task: %s", task_name)
except Exception as e:
logger.exception("Error enqueueing task %s", task_name)
if settings.SENTRY.ENABLED:
sentry_sdk.capture_exception(e)
logger.info(
"next run for task %s is in %s seconds",
task_name,
task.interval_seconds,
)
# Schedule next run regardless of whether we enqueued
# (if already pending/in-progress, we'll skip next time too)