feat(agent): scheduler, healthcheck, and container entrypoint
This commit is contained in:
@@ -0,0 +1,37 @@
|
||||
"""Agent entrypoint: bootstrap the database, then run scheduled jobs forever."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from agent import __version__
|
||||
from agent.config import load_config
|
||||
from agent.db import ensure_database, run_migrations
|
||||
from agent.jobs.heartbeat import run_heartbeat
|
||||
from agent.scheduler import build_scheduler
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||||
)
|
||||
logger = logging.getLogger("agent")
|
||||
|
||||
|
||||
def main() -> None:
|
||||
logger.info("firefly-agent %s starting", __version__)
|
||||
cfg = load_config()
|
||||
ensure_database(cfg)
|
||||
applied = run_migrations(cfg)
|
||||
if applied:
|
||||
logger.info("applied migrations: %s", ", ".join(applied))
|
||||
status = run_heartbeat(cfg)
|
||||
logger.info("startup heartbeat: %s", status)
|
||||
logger.info(
|
||||
"scheduler starting; heartbeat every %d minutes",
|
||||
cfg.heartbeat_interval_minutes,
|
||||
)
|
||||
build_scheduler(cfg).start()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,23 @@
|
||||
"""Docker HEALTHCHECK entry: healthy = agentdb answers SELECT 1."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sys
|
||||
|
||||
from agent.config import load_config
|
||||
from agent.db import connect
|
||||
|
||||
|
||||
def main() -> int:
|
||||
try:
|
||||
cfg = load_config()
|
||||
with connect(cfg) as conn:
|
||||
conn.execute("SELECT 1")
|
||||
return 0
|
||||
except Exception as exc:
|
||||
print(f"unhealthy: {exc}", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -0,0 +1,23 @@
|
||||
"""APScheduler assembly: one place where jobs get registered."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from apscheduler.schedulers.base import BaseScheduler
|
||||
from apscheduler.schedulers.blocking import BlockingScheduler
|
||||
|
||||
from agent.config import AgentConfig
|
||||
from agent.jobs.heartbeat import run_heartbeat
|
||||
|
||||
|
||||
def build_scheduler(
|
||||
cfg: AgentConfig, *, scheduler_class: type = BlockingScheduler
|
||||
) -> BaseScheduler:
|
||||
scheduler = scheduler_class(timezone="UTC")
|
||||
scheduler.add_job(
|
||||
run_heartbeat,
|
||||
"interval",
|
||||
args=[cfg],
|
||||
minutes=cfg.heartbeat_interval_minutes,
|
||||
id="heartbeat",
|
||||
)
|
||||
return scheduler
|
||||
@@ -0,0 +1,30 @@
|
||||
from datetime import timedelta
|
||||
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
|
||||
from agent.config import AgentConfig
|
||||
from agent.scheduler import build_scheduler
|
||||
|
||||
|
||||
def _cfg(minutes: int) -> AgentConfig:
|
||||
return AgentConfig(
|
||||
db_host="db",
|
||||
db_port=5432,
|
||||
db_user="u",
|
||||
db_password="p",
|
||||
agent_db_name="agentdb",
|
||||
firefly_url="http://app:8080",
|
||||
firefly_token="tok",
|
||||
heartbeat_interval_minutes=minutes,
|
||||
)
|
||||
|
||||
|
||||
def test_build_scheduler_registers_heartbeat_at_configured_interval():
|
||||
scheduler = build_scheduler(_cfg(minutes=15), scheduler_class=BackgroundScheduler)
|
||||
scheduler.start(paused=True)
|
||||
try:
|
||||
job = scheduler.get_job("heartbeat")
|
||||
assert job is not None
|
||||
assert job.trigger.interval == timedelta(minutes=15)
|
||||
finally:
|
||||
scheduler.shutdown(wait=False)
|
||||
Reference in New Issue
Block a user