import json
import logging
import os
import secrets
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from threading import Event, Thread
from typing import Any
from urllib.parse import parse_qs, urlparse
from .alerts import build_alert_summary
from .remote import parse_agent_registration, parse_remote_event
from .reports import _row
from .timeutil import utc_now
from .watcher import Watcher, build_collector
LOGGER = logging.getLogger(__name__)
WATCHER_HTML = """<!doctype html>
<html lang="de">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>Worker-Überwachung · Agent Solutions</title>
<style>
:root { color-scheme: dark; --bg:#081416; --panel:#0d2022; --line:#244044; --text:#f4e4d2; --muted:#9bb0ae; --green:#54d98a; --yellow:#f4c95d; --red:#ff7777; --blue:#8dc8ff; }
* { box-sizing:border-box; } body { margin:0; background:var(--bg); color:var(--text); font:15px/1.45 system-ui,-apple-system,sans-serif; }
header { border-bottom:1px solid var(--line); }
.header-inner { width:min(1280px,90vw); margin:0 auto; padding:24px 0 20px; display:flex; justify-content:space-between; gap:20px; align-items:center; }
.header-title { flex:1 1 auto; min-width:0; }
.header-actions { display:flex; align-items:center; gap:16px; flex:0 1 auto; min-width:0; }
.refresh-meta { color:var(--muted); font-size:12px; line-height:1.45; text-align:right; }
.refresh-meta div + div { margin-top:2px; }
.refresh-message { color:inherit; margin-left:6px; }
.refresh-message.error { color:var(--red); }
#refresh { min-width:144px; display:inline-flex; align-items:center; justify-content:center; white-space:nowrap; }
.spinner { display:inline-block; flex:0 0 12px; width:12px; height:12px; margin-right:6px; border:2px solid currentColor; border-right-color:transparent; border-radius:50%; vertical-align:-2px; animation:spin .7s linear infinite; }
.spinner[hidden] { display:inline-block; visibility:hidden; animation:none; }
@keyframes spin { to { transform:rotate(360deg); } }
h1 { margin:0; font-size:28px; } h2 { margin:0 0 16px; font-size:19px; } .sub { color:var(--muted); margin-top:4px; }
main { width:min(1280px,90vw); margin:28px auto 50px; } .grid { display:grid; grid-template-columns:repeat(auto-fit,minmax(140px,1fr)); gap:12px; margin-bottom:24px; }
.card,.section { background:var(--panel); border:1px solid var(--line); border-radius:10px; } .card { padding:16px; }
.label { color:var(--muted); font-size:12px; text-transform:uppercase; letter-spacing:.08em; } .value { font-size:28px; margin-top:4px; }
.section { padding:20px; margin-bottom:20px; overflow:auto; } table { border-collapse:collapse; width:100%; min-width:700px; } .job-table { min-width:1100px; } th,td { text-align:left; padding:10px 8px; border-bottom:1px solid var(--line); vertical-align:top; } th { color:var(--muted); font-size:12px; text-transform:uppercase; letter-spacing:.06em; white-space:nowrap; }
.sort-button { color:inherit; background:transparent; border:0; border-radius:4px; padding:2px 0; font:inherit; letter-spacing:inherit; text-transform:inherit; text-align:left; } .sort-button:hover,.sort-button:focus-visible { color:var(--text); outline:2px solid var(--blue); outline-offset:2px; } .sort-arrow { color:var(--blue); margin-left:4px; }
.pill { display:inline-block; padding:3px 9px; border-radius:999px; border:1px solid var(--line); font-size:12px; } .healthy,.online,.running,.completed { color:var(--green); border-color:#246b48; } .failed,.stale,.stopped { color:var(--red); border-color:#713b3b; } .degraded,.unknown,.unavailable,.waiting,.paused,.scheduled { color:var(--yellow); border-color:#725d2c; } .fresh { color:var(--green); } .muted { color:var(--muted); }
button { color:var(--blue); background:transparent; border:1px solid var(--line); border-radius:6px; padding:8px 12px; cursor:pointer; } button:hover { border-color:var(--blue); } button:disabled { opacity:.7; cursor:wait; }
.empty { color:var(--muted); padding:10px 0; } @media (max-width:850px) { .grid { grid-template-columns:repeat(2,1fr); } .header-inner { align-items:flex-start; flex-direction:column; } .header-actions { width:100%; justify-content:space-between; } } @media (max-width:500px) { .grid { grid-template-columns:1fr; } .header-actions { align-items:flex-start; flex-direction:column; } .refresh-meta { text-align:left; } }
</style>
</head>
<body>
<header><div class="header-inner"><div class="header-title"><h1>Worker-Überwachung</h1><div class="sub">Zentrale Nur-Lese-Überwachung von KI-Workern und Automationen</div></div><div class="header-actions"><div class="refresh-meta"><div>Zuletzt aktualisiert: <time id="last-updated">Unbekannt</time></div><div>Letzte Watcher-Prüfung: <time id="last-watcher-cycle">Unbekannt</time><span id="refresh-message" class="refresh-message" aria-live="polite"></span></div></div><button id="refresh" type="button" aria-busy="false"><span id="refresh-spinner" class="spinner" aria-hidden="true" hidden></span><span id="refresh-label">Aktualisieren</span></button></div></div></header>
<main>
<div class="grid" id="summary"></div>
<section class="section"><h2>Zentrale Infrastruktur</h2><div id="jobs-core" class="empty">Lade Daten …</div></section>
<section class="section"><h2>Projekt-Worker</h2><div id="jobs-projects" class="empty">Lade Daten …</div></section>
<section class="section"><h2>Automationen und Hintergrund-Worker</h2><div id="jobs-automation" class="empty">Lade Daten …</div></section>
<section class="section"><h2>Weitere überwachte Aufgaben</h2><div id="jobs-other" class="empty">Lade Daten …</div></section>
<section class="section"><h2>Ausführungsorte und Remote-Agenten</h2><div id="agents" class="empty">Lade Daten …</div></section>
<section class="section"><h2>Letzte Ereignisse</h2><div id="events" class="empty">Lade Daten …</div></section>
</main>
<script>
const esc = value => String(value ?? '—').replace(/[&<>"']/g, c => ({'&':'&','<':'<','>':'>','"':'"',"'":'''}[c]));
const statusTranslations = {
healthy:'Fehlerfrei', degraded:'Gestört', failed:'Fehlgeschlagen', unknown:'Unbekannt',
fresh:'Aktuell', unavailable:'Nicht verfügbar', stale:'Veraltet', running:'Läuft', stopped:'Gestoppt',
waiting:'Wartet', paused:'Pausiert', scheduled:'Geplant', completed:'Abgeschlossen',
critical:'Kritisch', warning:'Warnung', info:'Information', online:'Online', offline:'Offline',
pending:'Ausstehend', activated:'Aktiviert', inactive:'Inaktiv', active:'Aktiv', dead:'Beendet',
exited:'Beendet', activating:'Wird gestartet', deactivating:'Wird beendet', restarting:'Wird neu gestartet',
created:'Erstellt', starting:'Wird gestartet', unhealthy:'Fehlerhaft', none:'Nicht geprüft',
success:'Erfolgreich', failure:'Fehlerhaft', timeout:'Zeitüberschreitung', signal:'Signal', error:'Fehler'
};
const eventTranslations = {
new_instance:'Neue Instanz', status_changed:'Zustand geändert', error_raised:'Fehler aufgetreten',
error_resolved:'Fehler behoben', collector_failed:'Collector fehlgeschlagen', collector_succeeded:'Collector erfolgreich'
};
const statusLabel = value => statusTranslations[String(value ?? '').toLowerCase()] || value || '—';
const eventLabel = value => eventTranslations[String(value ?? '').toLowerCase()] || statusLabel(value);
const pill = value => `<span class="pill ${esc(value)}">${esc(statusLabel(value))}</span>`;
const get = path => fetch(path, {cache:'no-store'}).then(r => { if (!r.ok) throw new Error(`HTTP ${r.status}`); return r.json(); });
const tableStates = {};
const collator = new Intl.Collator('de-DE', {sensitivity:'base', numeric:true});
const healthRank = {failed:0, degraded:1, unknown:2, stale:3, healthy:4};
const observationRank = {unavailable:0, stale:1, fresh:2};
const operationalRank = {stopped:0, paused:1, waiting:2, scheduled:3, running:4, completed:5, unknown:6};
const severityRank = {critical:0, warning:1, info:2};
const agentRank = {offline:0, pending:1, online:2};
const emptyValue = value => value === null || value === undefined || value === '' || Number.isNaN(value);
const compareValues = (left, right, column) => {
const leftEmpty = emptyValue(left.value);
const rightEmpty = emptyValue(right.value);
if (leftEmpty || rightEmpty) return leftEmpty === rightEmpty ? 0 : (leftEmpty ? 1 : -1);
let result = 0;
if (column.type === 'number' || column.type === 'date' || column.type === 'duration') {
result = Number(left.value) - Number(right.value);
} else if (column.type === 'status') {
const ranks = column.rank || {};
result = (ranks[String(left.value).toLowerCase()] ?? 999) - (ranks[String(right.value).toLowerCase()] ?? 999);
} else {
result = collator.compare(String(left.value), String(right.value));
}
return result;
};
const defaultJobCompare = (left, right) => {
for (const [key, column] of [['health', {type:'status', rank:healthRank}], ['activity', {type:'date'}], ['task', {type:'text'}]]) {
const result = compareValues(left.cells[key], right.cells[key], column);
if (result) return result;
}
return 0;
};
const defaultEventCompare = (left, right) => compareValues(right.cells.time, left.cells.time, {type:'date'});
const defaultAgentCompare = (left, right) => {
const status = compareValues(left.cells.status, right.cells.status, {type:'status', rank:agentRank});
return status || compareValues(right.cells.heartbeat, left.cells.heartbeat, {type:'date'});
};
const sortRows = (rows, columns, state, defaultCompare) => {
const ordered = [...rows];
if (!state.key) return ordered.sort(defaultCompare);
const column = columns.find(item => item.key === state.key);
if (!column) return ordered.sort(defaultCompare);
return ordered.sort((left, right) => compareValues(left.cells[state.key], right.cells[state.key], column) * state.direction);
};
const sortableHeader = (target, column, state) => {
const active = state.key === column.key && state.direction !== 0;
const arrow = active ? (state.direction === 1 ? ' ▲' : ' ▼') : '';
const direction = active ? (state.direction === 1 ? 'absteigend' : 'aufsteigend') : 'aufsteigend';
return `<th scope="col" aria-sort="${active ? (state.direction === 1 ? 'ascending' : 'descending') : 'none'}"><button class="sort-button" type="button" data-sort-target="${esc(target)}" data-sort-key="${esc(column.key)}" aria-label="Nach ${esc(column.label)} ${direction} sortieren" title="Nach ${esc(column.label)} sortieren">${esc(column.label)}<span class="sort-arrow" aria-hidden="true">${arrow}</span></button></th>`;
};
function table(target, columns, rows, defaultCompare) {
if (!rows.length) { document.getElementById(target).innerHTML = '<div class="empty">Noch keine Daten vorhanden.</div>'; return; }
const available = new Set(columns.map(column => column.key));
const state = tableStates[target] || {key:null, direction:0};
if (state.key && !available.has(state.key)) { state.key = null; state.direction = 0; }
tableStates[target] = state;
const ordered = sortRows(rows, columns, state, defaultCompare);
const markup = ordered.map(row => `<tr>${columns.map(column => `<td>${row.cells[column.key].html}</td>`).join('')}</tr>`).join('');
const container = document.getElementById(target);
container.innerHTML = `<table class="${target.startsWith('jobs-') ? 'job-table' : ''}"><thead><tr>${columns.map(column => sortableHeader(target, column, state)).join('')}</tr></thead><tbody>${markup}</tbody></table>`;
container.querySelectorAll('.sort-button').forEach(button => button.addEventListener('click', () => {
const key = button.dataset.sortKey;
const current = tableStates[target];
if (current.key !== key) { current.key = key; current.direction = 1; }
else if (current.direction === 1) current.direction = -1;
else { current.key = null; current.direction = 0; }
refresh();
}));
}
const projectNames = {lead:'Lead Engine',social:'Social Media',telefon:'Telefon Agent',vento:'Vento',youtube:'YouTube Research',content:'Content Extraction',solar:'Solar',mirofish:'MiroFish'};
const projectName = job => projectNames[job.job_key.split('.')[1]] || 'Projekt';
const cell = (value, html) => ({value, html: html ?? esc(value)});
const statusCell = value => cell(value, pill(value));
const formatDate = value => {
if (!value) return '—';
const date = new Date(value);
if (Number.isNaN(date.getTime())) return value;
return `${new Intl.DateTimeFormat('de-DE', {dateStyle:'medium', timeStyle:'medium', timeZone:'UTC'}).format(date)} UTC`;
};
const parseTime = value => {
if (!value) return null;
const timestamp = new Date(value).getTime();
return Number.isNaN(timestamp) ? null : timestamp;
};
const relativeDuration = milliseconds => {
const seconds = Math.max(0, Math.floor(milliseconds / 1000));
if (seconds < 60) return `vor ${seconds} ${seconds === 1 ? 'Sekunde' : 'Sekunden'}`;
const minutes = Math.floor(seconds / 60);
if (minutes < 60) return `vor ${minutes} ${minutes === 1 ? 'Minute' : 'Minuten'}`;
const hours = Math.floor(minutes / 60);
if (hours < 24) return `vor ${hours} ${hours === 1 ? 'Stunde' : 'Stunden'}`;
const days = Math.floor(hours / 24);
return `vor ${days} ${days === 1 ? 'Tag' : 'Tagen'}`;
};
const evidenceValue = (job, key) => (job.evidence_json || []).find(item => item.key === key)?.value;
// Ersatzreihenfolge: fachliche Aktivität, Laufende, Erfolg, Lebenszeichen, Systemd-Aktivierung.
const activityData = job => {
const raw = job.last_activity_at || job.last_finished_at || job.last_started_at || job.last_success_at || job.heartbeat_at || evidenceValue(job, 'systemd_ActiveEnterTimestamp');
const timestamp = parseTime(raw);
return {timestamp, label: timestamp === null ? 'Unbekannt' : formatDate(timestamp)};
};
const nextRunData = job => {
const timestamp = parseTime(job.expected_next_run_at);
if (timestamp !== null) {
return timestamp <= Date.now()
? cell(timestamp, `Überfällig seit ${relativeDuration(Date.now() - timestamp)}`)
: cell(timestamp, formatDate(timestamp));
}
const continuous = job.operational_state === 'running' && ['service', 'worker'].includes(job.job_type);
return cell(null, continuous ? 'Dauerbetrieb' : 'Kein Zeitplan bekannt');
};
const sourceStateLabel = value => statusLabel(value);
const reasonText = value => {
const text = String(value ?? '');
const translations = {
'reported source status':'Gemeldeter Quellenstatus',
'source is healthy':'Quelle ist fehlerfrei',
'run failed with exit code 15':'Ausführung mit Exitcode 15 fehlgeschlagen',
'collector could not read the configured source':'Collector konnte die konfigurierte Quelle nicht lesen',
'last known job status retained; current observation unavailable':'Letzter bekannter Aufgabenstatus beibehalten; aktuelle Beobachtung nicht verfügbar',
'source was readable but did not provide enough evidence':'Quelle war lesbar, lieferte aber nicht genügend Nachweise',
'remote agent heartbeat overdue':'Lebenszeichen des Remote-Agenten überfällig'
};
const missingSource = text.match(/^error: no such object: (.+)$/i);
return translations[text] || (missingSource ? `Fehler: Quelle nicht gefunden: ${missingSource[1]}` : text || '—');
};
const reason = job => {
if (job.observation_status === 'unavailable') return `Quelle nicht verfügbar: ${job.source_reference || job.worker_key}`;
if (job.status_reason === 'reported source status') {
const source = evidenceValue(job, 'source');
const active = evidenceValue(job, 'systemd_ActiveState');
const sub = evidenceValue(job, 'systemd_SubState');
const result = evidenceValue(job, 'systemd_Result');
const state = evidenceValue(job, 'reported_state');
const container = evidenceValue(job, 'container_name');
const cstate = evidenceValue(job, 'container_state');
const chealth = evidenceValue(job, 'container_health');
if (active || sub) return `Systemd ${sourceStateLabel(active)}/${sourceStateLabel(sub)}${result ? `, Ergebnis ${sourceStateLabel(result)}` : ''}`;
if (container || cstate) return `Docker ${container || '—'}: ${sourceStateLabel(cstate)}${chealth ? ` (${sourceStateLabel(chealth)})` : ''}`;
if (source === 'procfs') return evidenceValue(job, 'process_exists') ? 'Prozess vorhanden' : 'Prozess nicht vorhanden';
return state ? `Gemeldeter Zustand: ${sourceStateLabel(state)}` : 'Quelle meldet normalen Zustand';
}
return reasonText(job.status_reason);
};
const jobColumns = [
{key:'task', label:'Aufgabe / Instanz', type:'text'},
{key:'worker', label:'Worker / Dienst', type:'text'},
{key:'location', label:'Ort', type:'text'},
{key:'activity', label:'Letzte Aktivität', type:'date'},
{key:'since', label:'Seit', type:'duration'},
{key:'operational', label:'Betriebszustand', type:'status', rank:operationalRank},
{key:'health', label:'Gesundheit', type:'status', rank:healthRank},
{key:'observation', label:'Beobachtung', type:'status', rank:observationRank},
{key:'next', label:'Nächster Lauf / Überfällig', type:'date'},
{key:'reason', label:'Grund', type:'text'}
];
const projectColumns = [{key:'project', label:'Projekt', type:'text'}, ...jobColumns];
const jobCells = job => {
const activity = activityData(job);
const since = activity.timestamp === null ? cell(null, 'Unbekannt') : cell(Date.now() - activity.timestamp, relativeDuration(Date.now() - activity.timestamp));
const taskValue = `${job.name} ${job.job_key}`;
const workerValue = job.worker_key || job.source_reference || '';
const reasonValue = reason(job);
return {
task: cell(taskValue, `${esc(job.name)}<br><span class="muted">${esc(job.job_key)}</span>`),
worker: cell(workerValue, `${esc(workerValue)}${job.source_reference && job.source_reference !== workerValue ? `<br><span class="muted">${esc(job.source_reference)}</span>` : ''}`),
location: cell(job.location_name || job.host || null),
activity: cell(activity.timestamp, esc(activity.label)),
since,
operational: statusCell(job.operational_state),
health: statusCell(job.health_status),
observation: statusCell(job.observation_status),
next: nextRunData(job),
reason: cell(reasonValue, esc(reasonValue))
};
};
const jobRow = job => ({cells: jobCells(job)});
const projectJobRow = job => ({cells: {project: cell(projectName(job)), ...jobCells(job)}});
const agentColumns = [
{key:'agent', label:'Agent', type:'text'}, {key:'environment', label:'Umgebung', type:'text'},
{key:'host', label:'Host', type:'text'}, {key:'runtime', label:'Laufzeit', type:'text'},
{key:'status', label:'Zustand', type:'status', rank:agentRank}, {key:'heartbeat', label:'Letzte Rückmeldung', type:'date'}
];
const agentRow = agent => {
const heartbeat = parseTime(agent.last_heartbeat_at);
return {cells:{
agent:cell(agent.agent_name, `${esc(agent.agent_name)}<br><span class="muted">${esc(agent.agent_key)}</span>`),
environment:cell(agent.environment), host:cell(agent.host), runtime:cell(agent.runtime),
status:statusCell(agent.status), heartbeat:cell(heartbeat, esc(heartbeat === null ? 'Unbekannt' : formatDate(heartbeat)))
}};
};
const eventColumns = [
{key:'time', label:'Zeit', type:'date'}, {key:'type', label:'Typ', type:'text'},
{key:'severity', label:'Schweregrad', type:'status', rank:severityRank}, {key:'component', label:'Komponente', type:'text'},
{key:'reason', label:'Grund', type:'text'}
];
const eventRow = event => {
const time = parseTime(event.occurred_at);
const reasonValue = reasonText(event.reason || event.message);
return {cells:{
time:cell(time, esc(time === null ? 'Unbekannt' : formatDate(time))), type:cell(eventLabel(event.event_type)),
severity:statusCell(event.new_severity || event.severity), component:cell(event.component || event.instance_id),
reason:cell(reasonValue, esc(reasonValue))
}};
};
const refreshButton = document.getElementById('refresh');
const refreshLabel = document.getElementById('refresh-label');
const refreshSpinner = document.getElementById('refresh-spinner');
const refreshMessage = document.getElementById('refresh-message');
const lastUpdated = document.getElementById('last-updated');
const lastWatcherCycle = document.getElementById('last-watcher-cycle');
let refreshInFlight = null;
let hasRenderedData = false;
let lastSuccessfulRefreshAt = null;
const formatLocalTimestamp = value => {
if (!value) return 'Unbekannt';
const timestamp = new Date(value).getTime();
if (Number.isNaN(timestamp)) return 'Unbekannt';
return new Intl.DateTimeFormat('de-DE', {
day:'2-digit', month:'2-digit', year:'numeric', hour:'2-digit', minute:'2-digit', second:'2-digit',
hourCycle:'h23', timeZone:'Europe/Madrid'
}).format(new Date(timestamp));
};
const updateRefreshMeta = () => {
lastUpdated.textContent = lastSuccessfulRefreshAt
? `${formatLocalTimestamp(lastSuccessfulRefreshAt)} · ${relativeRefreshAge((Date.now() - lastSuccessfulRefreshAt) / 1000)}`
: 'Unbekannt';
};
const relativeRefreshAge = seconds => {
const value = Math.max(0, Math.floor(Number(seconds)));
if (value < 60) return 'gerade eben';
if (value < 3600) return `vor ${Math.floor(value / 60)} Min.`;
if (value < 86400) return `vor ${Math.floor(value / 3600)} Std.`;
return `vor ${Math.floor(value / 86400)} T.`;
};
const setRefreshUi = loading => {
refreshButton.disabled = loading;
refreshButton.setAttribute('aria-busy', String(loading));
refreshSpinner.hidden = !loading;
refreshLabel.textContent = 'Aktualisieren';
};
let refreshMessageTimer = null;
const setRefreshMessage = (message, error = false, duration = 0) => {
if (refreshMessageTimer) window.clearTimeout(refreshMessageTimer);
refreshMessage.textContent = message;
refreshMessage.classList.toggle('error', error);
if (duration > 0) {
refreshMessageTimer = window.setTimeout(() => {
refreshMessage.textContent = '';
refreshMessage.classList.remove('error');
refreshMessageTimer = null;
}, duration);
}
};
function refresh(manual = false) {
if (refreshInFlight) return refreshInFlight;
setRefreshUi(true);
if (manual) setRefreshMessage('Wird aktualisiert …');
refreshInFlight = Promise.all([get('/api/v1/status'), get('/api/v1/jobs'), get('/api/v1/agents'), get('/api/v1/events?limit=25')])
.then(([status, jobs, agents, events]) => {
const s = status;
document.getElementById('summary').innerHTML = [['Watcher',statusLabel(s.watcher_status)],['Aufgaben',s.jobs_total],['Fehlerfrei',s.jobs_healthy],['Gestört',s.jobs_degraded],['Fehlgeschlagen',s.jobs_failed],['Unbekannt',s.jobs_unknown],['Agenten',agents.items.length]].map(x => `<div class="card"><div class="label">${esc(x[0])}</div><div class="value">${esc(x[1])}</div></div>`).join('');
const groups = {core:[], projects:[], automation:[], other:[]};
jobs.items.forEach(j => { const key = j.job_key; if (key.startsWith('project.')) groups.projects.push(j); else if (key.startsWith('automation.')) groups.automation.push(j); else if (key.startsWith('agent.') || key.startsWith('dashboard.') || key.startsWith('portal.') || key.startsWith('admin.') || key.startsWith('commander.') || key.startsWith('editor.') || key.startsWith('watcher.') || key.startsWith('process.')) groups.core.push(j); else groups.other.push(j); });
table('jobs-core',jobColumns,groups.core.map(jobRow),defaultJobCompare);
table('jobs-projects',projectColumns,groups.projects.map(projectJobRow),defaultJobCompare);
table('jobs-automation',jobColumns,groups.automation.map(jobRow),defaultJobCompare);
table('jobs-other',jobColumns,groups.other.map(jobRow),defaultJobCompare);
table('agents',agentColumns,agents.items.map(agentRow),defaultAgentCompare);
table('events',eventColumns,events.items.map(eventRow),defaultEventCompare);
hasRenderedData = true;
lastSuccessfulRefreshAt = Date.now();
lastUpdated.dateTime = new Date(lastSuccessfulRefreshAt).toISOString();
updateRefreshMeta();
lastWatcherCycle.textContent = formatLocalTimestamp(s.last_cycle_finished_at);
lastWatcherCycle.dateTime = s.last_cycle_finished_at || '';
if (manual) setRefreshMessage('Aktualisiert', false, 2000);
else setRefreshMessage('');
})
.catch(error => {
if (!hasRenderedData) document.getElementById('summary').innerHTML = `<div class="card"><div class="label">Verbindung</div><div class="value">Fehler</div><div class="muted">${esc(error instanceof Error ? error.message : error)}</div></div>`;
setRefreshMessage('Aktualisierung fehlgeschlagen', true, 2000);
})
.finally(() => {
setRefreshUi(false);
refreshInFlight = null;
});
return refreshInFlight;
}
refreshButton.addEventListener('click', () => refresh(true));
refresh();
setInterval(refresh, 15000);
setInterval(updateRefreshMeta, 1000);
</script>
</body>
</html>"""
def collector_config_rows(watcher: Watcher) -> list[dict[str, Any]]:
rows: list[dict[str, Any]] = []
for config in watcher.config.collectors:
latest = watcher.connection.execute(
"SELECT * FROM collector_runs WHERE collector_name=? ORDER BY id DESC LIMIT 1", (config.name,)
).fetchone()
try:
version = build_collector(config.type).version
except ValueError:
version = "unknown"
item: dict[str, Any] = {
"name": config.name,
"type": config.type,
"enabled": config.enabled,
"timeout_seconds": config.timeout_seconds,
"definitions_count": len(config.definitions),
"version": version,
}
if latest is not None:
item.update({
"last_run_at": latest["started_at"],
"last_finished_at": latest["finished_at"],
"last_success": bool(latest["success"]),
"last_error_type": latest["error_type"],
"last_error_message": latest["error_message"],
"jobs_seen": latest["jobs_seen"],
"observations_created": latest["observations_created"],
})
else:
item.update({"last_run_at": None, "last_finished_at": None, "last_success": None,
"last_error_type": None, "last_error_message": None, "jobs_seen": 0,
"observations_created": 0})
rows.append(item)
return rows
def alert_summary(watcher: Watcher) -> dict[str, Any]:
health = watcher.health()
jobs = [_row(row) for row in watcher.repository.status_rows()]
events = [_row(row) for row in watcher.connection.execute("SELECT * FROM status_events ORDER BY occurred_at ASC")]
return build_alert_summary(health, jobs, events, utc_now(), watcher.config.watcher.cycle_interval_seconds)
def create_handler(watcher: Watcher) -> type[BaseHTTPRequestHandler]:
class Handler(BaseHTTPRequestHandler):
def log_message(self, format: str, *args: Any) -> None:
return
def send_json(self, status: int, payload: Any) -> None:
body = json.dumps(payload, ensure_ascii=False, default=str).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def send_html(self, status: int, body: str) -> None:
encoded = body.encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "text/html; charset=utf-8")
self.send_header("Content-Length", str(len(encoded)))
self.end_headers()
self.wfile.write(encoded)
def read_json(self) -> dict[str, Any] | None:
try:
length = int(self.headers.get("Content-Length", "0"))
if length <= 0 or length > watcher.config.watcher.max_remote_event_bytes:
return None
payload = json.loads(self.rfile.read(length))
return payload if isinstance(payload, dict) else None
except (ValueError, json.JSONDecodeError):
return None
def authorized_agent_request(self) -> bool:
expected = os.environ.get(watcher.config.watcher.agent_auth_token_env)
if not expected:
return True
provided = self.headers.get("Authorization", "")
return secrets.compare_digest(provided, f"Bearer {expected}")
def do_GET(self) -> None:
parsed = urlparse(self.path)
path = parsed.path.rstrip("/") or "/"
if path == "/":
self.send_html(200, WATCHER_HTML)
elif path == "/health/live":
self.send_json(200, {"status": "live"})
elif path == "/health/ready":
try:
watcher.connection.execute("SELECT 1").fetchone()
self.send_json(200, {"status": "ready", "schema_version": watcher.schema_version})
except Exception as exc: # noqa: BLE001 - readiness must never expose internals
self.send_json(503, {"status": "not_ready", "error": type(exc).__name__})
elif path == "/health" or path == "/api/v1/status":
self.send_json(200, watcher.health())
elif path == "/api/v1/alerts/summary":
self.send_json(200, alert_summary(watcher))
elif path == "/api/v1/jobs":
query = parse_qs(parsed.query)
rows = watcher.repository.status_rows()
status = query.get("status", [None])[0]
if status:
rows = [row for row in rows if row["health_status"] == status or row["operational_state"] == status]
self.send_json(200, {"items": [_row(row) for row in rows]})
elif path == "/api/v1/workers":
self.send_json(200, {"items": [_row(row) for row in watcher.repository.worker_rows()]})
elif path == "/api/v1/agents":
self.send_json(200, {"items": [_row(row) for row in watcher.repository.agent_rows()]})
elif path == "/api/v1/collectors":
self.send_json(200, {"items": collector_config_rows(watcher)})
elif path == "/api/v1/collectors/runs":
limit = int(parse_qs(parsed.query).get("limit", [100])[0])
self.send_json(200, {"items": [_row(row) for row in watcher.repository.collector_rows(limit)]})
elif path.startswith("/api/v1/collectors/"):
collector_name = path.rsplit("/", 1)[1]
item = next((row for row in collector_config_rows(watcher) if row["name"] == collector_name), None)
if item is None:
self.send_json(404, {"error": "collector not found"})
else:
self.send_json(200, item)
elif path == "/api/v1/remote-events":
limit = int(parse_qs(parsed.query).get("limit", [100])[0])
self.send_json(200, {"items": [_row(row) for row in watcher.repository.remote_event_rows(limit)]})
elif path.startswith("/api/v1/jobs/"):
try:
instance_id = int(path.rsplit("/", 1)[1])
except ValueError:
self.send_json(404, {"error": "instance not found"})
return
row = watcher.connection.execute("SELECT * FROM job_status_current WHERE instance_id=?", (instance_id,)).fetchone()
self.send_json(200, _row(row) if row else {"error": "instance not found"}) if row else self.send_json(404, {"error": "instance not found"})
elif path == "/api/v1/events":
limit = int(parse_qs(parsed.query).get("limit", [50])[0])
self.send_json(200, {"items": [_row(row) for row in watcher.repository.events(limit)]})
else:
self.send_json(404, {"error": "not found"})
def do_POST(self) -> None:
path = urlparse(self.path).path.rstrip("/") or "/"
if path not in {"/api/v1/agents/register", "/api/v1/agents/events"}:
self.send_json(404, {"error": "not found"})
return
if not self.authorized_agent_request():
self.send_json(401, {"error": "agent authentication required"})
return
if not watcher.config.watcher.remote_ingest_enabled:
self.send_json(404, {"error": "remote ingest disabled"})
return
payload = self.read_json()
if payload is None:
self.send_json(400, {"error": "JSON object required and limited to 512 KiB"})
return
try:
if path.endswith("/register"):
registration = parse_agent_registration(payload)
agent_id = watcher.register_remote_agent(registration)
self.send_json(201, {"accepted": True, "agent_id": agent_id, "agent_key": registration.agent_key})
else:
event = parse_remote_event(payload)
self.send_json(200, watcher.ingest_remote_event(event))
except (KeyError, TypeError, ValueError) as exc:
self.send_json(422, {"error": str(exc)})
return Handler
def serve(watcher: Watcher, host: str, port: int) -> None:
stop = Event()
def scan_loop() -> None:
while not stop.is_set():
try:
watcher.run_once()
except Exception:
LOGGER.exception("watcher cycle failed")
stop.wait(watcher.config.watcher.cycle_interval_seconds)
scan_thread = Thread(target=scan_loop, name="watcher-cycle", daemon=True)
scan_thread.start()
server = ThreadingHTTPServer((host, port), create_handler(watcher))
try:
server.serve_forever()
finally:
stop.set()
scan_thread.join(timeout=2)
server.server_close()