release: bump version to 0.39.0

This commit is contained in:
rayd1o
2026-04-24 00:48:33 +08:00
parent d5f3784ffb
commit 8b8f7138c0
22 changed files with 1490 additions and 762 deletions

View File

@@ -1,7 +1,6 @@
"""Task Scheduler for running collection jobs."""
import asyncio
import logging
from datetime import UTC, datetime, timedelta
from typing import Any, Dict, Optional
@@ -9,13 +8,14 @@ from apscheduler.schedulers.asyncio import AsyncIOScheduler
from apscheduler.triggers.interval import IntervalTrigger
from sqlalchemy import select
from app.core.logging import get_logger
from app.db.session import async_session_factory
from app.core.time import to_iso8601_utc
from app.models.datasource import DataSource
from app.models.task import CollectionTask
from app.services.collectors.registry import collector_registry
logger = logging.getLogger(__name__)
logger = get_logger(__name__)
scheduler = AsyncIOScheduler()
RUNNING_TASK_GUARD_TIMEOUT_MINUTES = 90
@@ -54,7 +54,11 @@ async def _update_next_run_at(datasource: DataSource, session) -> None:
async def _apply_datasource_schedule(datasource: DataSource, session) -> None:
collector = collector_registry.get(datasource.source)
if not collector:
logger.warning("Collector not found for datasource %s", datasource.source)
logger.warning_event(
"Collector not found for datasource",
event="collector.schedule.collector_missing",
context={"collector_name": datasource.source},
)
return
collector_registry.set_active(datasource.source, datasource.is_active)
@@ -72,13 +76,17 @@ async def _apply_datasource_schedule(datasource: DataSource, session) -> None:
replace_existing=True,
kwargs={"collector_name": datasource.source},
)
logger.info(
"Scheduled collector: %s (every %sm)",
datasource.source,
datasource.frequency_minutes,
logger.info_event(
"Scheduled collector",
event="collector.schedule.updated",
context={"collector_name": datasource.source, "frequency_minutes": datasource.frequency_minutes},
)
else:
logger.info("Collector disabled: %s", datasource.source)
logger.info_event(
"Collector disabled",
event="collector.schedule.disabled",
context={"collector_name": datasource.source},
)
await _update_next_run_at(datasource, session)
@@ -87,18 +95,30 @@ async def run_collector_task(collector_name: str):
"""Run a single collector task."""
collector = collector_registry.get(collector_name)
if not collector:
logger.error("Collector not found: %s", collector_name)
logger.error_event(
"Collector not found",
event="collector.run.collector_missing",
context={"collector_name": collector_name},
)
return
async with async_session_factory() as db:
result = await db.execute(select(DataSource).where(DataSource.source == collector_name))
datasource = result.scalar_one_or_none()
if not datasource:
logger.error("Datasource not found for collector: %s", collector_name)
logger.error_event(
"Datasource not found for collector",
event="collector.run.datasource_missing",
context={"collector_name": collector_name},
)
return
if not datasource.is_active:
logger.info("Skipping disabled collector: %s", collector_name)
logger.info_event(
"Skipping disabled collector",
event="collector.run.skipped_disabled",
context={"collector_name": collector_name},
)
return
running_result = await db.execute(
@@ -122,10 +142,10 @@ async def run_collector_task(collector_name: str):
and (now - started_at) > timedelta(minutes=RUNNING_TASK_GUARD_TIMEOUT_MINUTES)
)
if not is_stale:
logger.warning(
"Skipping collector %s trigger because task %s is already running",
collector_name,
existing_running.id,
logger.warning_event(
"Skipping collector trigger because task is already running",
event="collector.run.skipped_already_running",
context={"collector_name": collector_name, "task_id": existing_running.id},
)
return
@@ -143,31 +163,47 @@ async def run_collector_task(collector_name: str):
else stale_reason
)
await db.commit()
logger.warning(
"Marked stale running task %s as failed before rerun of %s",
existing_running.id,
collector_name,
logger.warning_event(
"Marked stale running task as failed before rerun",
event="collector.run.stale_task_failed",
context={"collector_name": collector_name, "task_id": existing_running.id},
)
try:
collector._datasource_id = datasource.id
logger.info("Running collector: %s (datasource_id=%s)", collector_name, datasource.id)
logger.info_event(
"Running collector",
event="collector.run.started",
context={"collector_name": collector_name, "datasource_id": datasource.id},
)
task_result = await collector.run(db)
datasource.last_run_at = datetime.now(UTC)
datasource.last_status = task_result.get("status")
await _update_next_run_at(datasource, db)
logger.info("Collector %s completed: %s", collector_name, task_result)
logger.info_event(
"Collector completed",
event="collector.run.completed",
context={"collector_name": collector_name, "datasource_id": datasource.id, "result": task_result},
)
except asyncio.CancelledError:
datasource.last_run_at = datetime.now(UTC)
datasource.last_status = "cancelled"
await db.commit()
logger.warning("Collector %s cancelled by operator", collector_name)
logger.warning_event(
"Collector cancelled by operator",
event="collector.run.cancelled",
context={"collector_name": collector_name, "datasource_id": datasource.id},
)
raise
except Exception as exc:
datasource.last_run_at = datetime.now(UTC)
datasource.last_status = "failed"
await db.commit()
logger.exception("Collector %s failed: %s", collector_name, exc)
logger.exception_event(
"Collector failed",
event="collector.run.failed",
context={"collector_name": collector_name, "datasource_id": datasource.id, "error": str(exc)},
)
async def cleanup_stale_running_tasks(max_age_hours: int = 2) -> int:
@@ -194,7 +230,11 @@ async def cleanup_stale_running_tasks(max_age_hours: int = 2) -> int:
if stale_tasks:
await db.commit()
logger.warning("Cleaned up %s stale running collection task(s)", len(stale_tasks))
logger.warning_event(
"Cleaned up stale running collection tasks",
event="collector.cleanup.stale_tasks_cleaned",
context={"count": len(stale_tasks)},
)
return len(stale_tasks)
@@ -203,14 +243,14 @@ def start_scheduler() -> None:
"""Start the scheduler."""
if not scheduler.running:
scheduler.start()
logger.info("Scheduler started")
logger.info_event("Scheduler started", event="scheduler.started")
def stop_scheduler() -> None:
"""Stop the scheduler."""
if scheduler.running:
scheduler.shutdown(wait=False)
logger.info("Scheduler stopped")
logger.info_event("Scheduler stopped", event="scheduler.stopped")
async def sync_scheduler_with_datasources() -> None:
@@ -271,12 +311,20 @@ def run_collector_now(collector_name: str) -> bool:
"""Run a collector immediately (not scheduled)."""
collector = collector_registry.get(collector_name)
if not collector:
logger.error("Collector not found: %s", collector_name)
logger.error_event(
"Collector not found",
event="collector.trigger.collector_missing",
context={"collector_name": collector_name},
)
return False
existing_task = get_running_collector_task(collector_name)
if existing_task is not None and not existing_task.done():
logger.warning("Collector %s is already running in-memory; skipping duplicate trigger", collector_name)
logger.warning_event(
"Collector is already running in-memory; skipping duplicate trigger",
event="collector.trigger.skipped_already_running",
context={"collector_name": collector_name},
)
return False
try:
@@ -289,10 +337,18 @@ def run_collector_now(collector_name: str) -> bool:
RUNNING_COLLECTOR_TASKS.pop(collector_name, None)
task.add_done_callback(_cleanup_task)
logger.info("Triggered collector: %s", collector_name)
logger.info_event(
"Triggered collector",
event="collector.trigger.started",
context={"collector_name": collector_name},
)
return True
except Exception as exc:
logger.error("Failed to trigger collector %s: %s", collector_name, exc)
logger.error_event(
"Failed to trigger collector",
event="collector.trigger.failed",
context={"collector_name": collector_name, "error": str(exc)},
)
return False