Files

327 lines
11 KiB
Python

# retoor <retoor@molodetz.nl>
import asyncio
import json
import logging
from collections import deque
from datetime import datetime, timedelta, timezone
from molodetz.database import (
ensure_service_state,
get_service_state,
get_setting,
now_iso,
parse_iso,
set_setting,
write_service_state,
)
logger = logging.getLogger(__name__)
HEARTBEAT_SECONDS = 8
LIVENESS_SECONDS = 15
METRICS_SECONDS = 15
FIELD_TYPES = ("int", "str", "url", "text", "password", "bool", "select")
class ConfigField:
def __init__(self, key, type="str", default="", label="", help="", minimum=None, maximum=None, options=None, secret=False, group="General"):
if type not in FIELD_TYPES:
raise ValueError(f"unknown field type {type}")
self.key = key
self.type = type
self.default = default
self.label = label or key
self.help = help
self.minimum = minimum
self.maximum = maximum
self.options = options or []
self.secret = secret or type == "password"
self.group = group
def coerce(self, raw):
value = "" if raw is None else str(raw).strip()
if self.type == "int":
number = int(value)
if self.minimum is not None and number < self.minimum:
raise ValueError(f"{self.label} must be at least {self.minimum}")
if self.maximum is not None and number > self.maximum:
raise ValueError(f"{self.label} must be at most {self.maximum}")
return str(number)
if self.type == "bool":
if value not in ("0", "1"):
raise ValueError(f"{self.label} must be 0 or 1")
return value
if self.type == "select":
if value not in [str(option) for option in self.options]:
raise ValueError(f"{self.label} must be one of {self.options}")
return value
if self.type == "url" and value and not value.startswith(("http://", "https://")):
raise ValueError(f"{self.label} must be an http(s) URL")
return value
def read(self, service_name):
raw = get_setting(f"service_{service_name}_{self.key}", str(self.default))
try:
return self.coerce(raw)
except (TypeError, ValueError):
return str(self.default)
def spec(self, service_name):
value = self.read(service_name)
return {
"key": self.key,
"type": self.type,
"label": self.label,
"help": self.help,
"group": self.group,
"options": self.options,
"minimum": self.minimum,
"maximum": self.maximum,
"secret": self.secret,
"value": "" if self.secret and value else value,
"has_value": bool(value),
}
class BaseService:
name = "base"
title = "Service"
description = ""
default_enabled = True
default_interval = 60
min_interval = 1
config_fields: list[ConfigField] = []
metrics_interval = METRICS_SECONDS
def __init__(self):
self.logs = deque(maxlen=self.log_size())
self._task = None
self._row_id = None
self._last_heartbeat = 0.0
self._last_metrics = 0.0
self._last_command = None
self._was_enabled = None
self._next_run = None
@property
def interval_key(self):
return f"service_{self.name}_interval"
def fields(self):
framework = [
ConfigField("enabled", "bool", "1" if self.default_enabled else "0", "Enabled", group="General"),
ConfigField("interval", "int", self.default_interval, "Interval (s)", minimum=self.min_interval, group="General"),
]
tail = [ConfigField("log_size", "int", 200, "Log lines", minimum=10, maximum=5000, group="Logs")]
return framework + list(self.config_fields) + tail
def field(self, key):
for field in self.fields():
if field.key == key:
return field
raise KeyError(key)
def get_config(self):
return {field.key: field.read(self.name) for field in self.fields()}
def is_enabled(self):
return self.field("enabled").read(self.name) == "1"
def current_interval(self):
return max(self.min_interval, int(self.field("interval").read(self.name)))
def log_size(self):
try:
return max(10, int(get_setting(f"service_{self.name}_log_size", "200")))
except ValueError:
return 200
def log(self, message):
line = f"{datetime.now(timezone.utc).strftime('%d/%m/%Y %H:%M:%S')} {message}"
self.logs.append(line)
logger.info("[%s] %s", self.name, message)
async def run_once(self):
raise NotImplementedError
async def on_enable(self):
return None
async def on_disable(self):
return None
def collect_metrics(self):
return {"cards": [], "table": None}
def _command(self):
return get_setting(f"service_{self.name}_command", "")
def _persist(self, status, force_metrics=False):
loop_time = asyncio.get_running_loop().time()
fields = {"status": status, "heartbeat": now_iso(), "logs": list(self.logs)}
if self._next_run:
fields["next_run"] = self._next_run.isoformat()
if force_metrics or loop_time - self._last_metrics >= self.metrics_interval:
fields["metrics"] = self.collect_metrics()
self._last_metrics = loop_time
write_service_state(self._row_id, **fields)
self._last_heartbeat = loop_time
async def _tick(self):
enabled = self.is_enabled()
if enabled != self._was_enabled:
if enabled:
write_service_state(self._row_id, started_at=now_iso())
await self.on_enable()
self.log("ingeschakeld")
elif self._was_enabled is not None:
await self.on_disable()
self.log("uitgeschakeld")
self._was_enabled = enabled
self._persist("running" if enabled else "stopped", force_metrics=True)
command = self._command()
run_now = False
if command and command != self._last_command:
if self._last_command is not None:
verb = command.split(":", 1)[0]
if verb == "run":
run_now = True
elif verb == "clear":
self.logs.clear()
self._last_command = command
if self.logs.maxlen != self.log_size():
self.logs = deque(self.logs, maxlen=self.log_size())
now = datetime.now(timezone.utc)
if enabled and (run_now or self._next_run is None or now >= self._next_run):
try:
await self.run_once()
except Exception as exc:
self.log(f"error: {exc}")
logger.exception("service %s failed", self.name)
self._next_run = datetime.now(timezone.utc) + timedelta(seconds=self.current_interval())
write_service_state(self._row_id, last_run=now_iso())
self._persist("running", force_metrics=True)
elif asyncio.get_running_loop().time() - self._last_heartbeat >= HEARTBEAT_SECONDS:
self._persist("running" if enabled else "stopped")
async def loop(self):
self._row_id = ensure_service_state(self.name)
self._last_command = None
while True:
try:
await self._tick()
except asyncio.CancelledError:
raise
except Exception:
logger.exception("service loop %s crashed a tick", self.name)
await asyncio.sleep(1)
def start(self):
if self._task is None or self._task.done():
self._task = asyncio.get_running_loop().create_task(self.loop())
async def stop(self):
if self._task is not None:
self._task.cancel()
try:
await self._task
except asyncio.CancelledError:
pass
self._task = None
if self._row_id is not None:
write_service_state(self._row_id, status="stopped")
def derived_status(service, row):
if not service.is_enabled():
return "stopped"
if row is None:
return "stalled"
heartbeat = parse_iso(row.get("heartbeat"))
if heartbeat and (datetime.now(timezone.utc) - heartbeat).total_seconds() <= LIVENESS_SECONDS:
return "running"
return "stalled"
class ServiceManager:
def __init__(self):
self.services = {}
self.supervising = False
def register(self, service):
self.services[service.name] = service
return service
def get(self, name):
return self.services.get(name)
def describe(self, service):
row = get_service_state(service.name)
metrics = {}
logs = []
if row:
try:
metrics = json.loads(row.get("metrics") or "{}")
except ValueError:
metrics = {}
try:
logs = json.loads(row.get("logs") or "[]")
except ValueError:
logs = []
return {
"name": service.name,
"title": service.title,
"description": service.description,
"status": derived_status(service, row),
"enabled": service.is_enabled(),
"last_run": row.get("last_run") if row else None,
"next_run": row.get("next_run") if row else None,
"heartbeat": row.get("heartbeat") if row else None,
"metrics": metrics,
"logs": logs,
"fields": [field.spec(service.name) for field in service.fields()],
}
def describe_all(self):
return [self.describe(service) for service in self.services.values()]
def set_enabled(self, name, enabled):
set_setting(f"service_{name}_enabled", "1" if enabled else "0")
def send_command(self, name, verb):
current = get_setting(f"service_{name}_command", "")
counter = 0
if ":" in current:
try:
counter = int(current.split(":", 1)[1])
except ValueError:
counter = 0
set_setting(f"service_{name}_command", f"{verb}:{counter + 1}")
def save_config(self, name, values):
service = self.services[name]
updates = {}
for field in service.fields():
if field.key not in values:
continue
raw = values[field.key]
if field.secret and (raw is None or str(raw).strip() == ""):
continue
updates[field.key] = field.coerce(raw)
for key, value in updates.items():
set_setting(f"service_{name}_{key}", value)
return updates
def supervise(self):
self.supervising = True
for service in self.services.values():
service.start()
async def shutdown_all(self):
for service in self.services.values():
await service.stop()
self.supervising = False
manager = ServiceManager()