feat: auto-restart policies, SSL monitoring, DNS propagation, dependency tracking, config drift detection
CI / Test & Lint (push) Has been cancelled
CI / Security audit (push) Has been cancelled

This commit is contained in:
Hermes
2026-06-10 14:43:46 -07:00
parent afcccf811e
commit 954be9e868
15 changed files with 3048 additions and 6 deletions
+1 -1
View File
@@ -1 +1 @@
1.8.0
1.9.0
+503
View File
@@ -0,0 +1,503 @@
/**
* Auto-Restart Manager - Per-container restart policies with retry tracking
*
* When a container goes down, attempts automatic restart up to N times
* (configurable per-service). Sends notifications on each attempt and
* when max retries are exceeded. Integrates with HealthChecker events.
*
* @module auto-restart-manager
*/
const EventEmitter = require('events');
const path = require('path');
const { readJsonFile, writeJsonFile } = require('./fs-helpers');
/**
* Default policy values applied when a new policy is created.
* @readonly
*/
const DEFAULT_POLICY = {
enabled: true,
maxRetries: 3,
retryIntervalMs: 5000,
windowMinutes: 10,
currentRetries: 0,
lastRestartAt: null,
cooldownUntil: null,
};
/**
* Manages automatic container restart policies and execution.
*
* @extends EventEmitter
*
* @fires AutoRestartManager#auto-restart-attempt
* @fires AutoRestartManager#auto-restart-success
* @fires AutoRestartManager#auto-restart-failed
* @fires AutoRestartManager#auto-restart-max-reached
*/
class AutoRestartManager extends EventEmitter {
/**
* @param {Object} ctx - Shared application context
* @param {Object} ctx.docker - Docker client wrapper ({ client: Dockerode })
* @param {Object} ctx.healthChecker - HealthChecker singleton
* @param {Object} ctx.notification - NotificationManager instance
* @param {Object} ctx.log - Logger instance
* @param {Function} ctx.logError - Error logging function
* @param {string} ctx.SERVICES_FILE - Path to services.json (used to derive data dir)
*/
constructor(ctx) {
super();
this.ctx = ctx;
this.log = ctx.log || console;
this.logError = ctx.logError || ((_ctx, err) => console.error(err));
this.docker = ctx.docker;
this.healthChecker = ctx.healthChecker;
this.notification = ctx.notification;
/** @type {Map<string, Object>} serviceId -> policy */
this.policies = new Map();
/** Path to the JSON file that persists policies */
this.policiesFile = path.join(path.dirname(ctx.SERVICES_FILE), 'auto-restart-policies.json');
/** Track previous health status per service for transition detection */
this._previousHealth = new Map();
/** Bound handlers so we can remove them on stop() */
this._onStatusCheck = this._handleStatusCheck.bind(this);
this._started = false;
}
// ─── Lifecycle ────────────────────────────────────────────────────────
/**
* Load persisted policies, then wire into HealthChecker events.
* @returns {Promise<void>}
*/
async start() {
if (this._started) return;
// Load persisted policies from disk
try {
const data = await readJsonFile(this.policiesFile, {});
for (const [serviceId, policy] of Object.entries(data)) {
this.policies.set(serviceId, { ...DEFAULT_POLICY, ...policy });
}
this.log.info('auto-restart', 'Policies loaded', { count: this.policies.size });
} catch (err) {
this.log.error('auto-restart', 'Failed to load policies', { error: err.message });
}
// Listen to health checker status transitions
if (this.healthChecker) {
this.healthChecker.on('status-check', this._onStatusCheck);
}
this._started = true;
this.log.info('auto-restart', 'Manager started');
}
/**
* Remove event listeners and stop processing health events.
*/
stop() {
if (!this._started) return;
if (this.healthChecker) {
this.healthChecker.removeListener('status-check', this._onStatusCheck);
}
this._started = false;
this.log.info('auto-restart', 'Manager stopped');
}
// ─── Policy CRUD ─────────────────────────────────────────────────────
/**
* Create or update a restart policy for a service.
*
* @param {string} serviceId - Unique service identifier
* @param {Object} policy - Partial policy fields to merge
* @param {boolean} [policy.enabled=true]
* @param {number} [policy.maxRetries=3]
* @param {number} [policy.retryIntervalMs=5000]
* @param {number} [policy.windowMinutes=10]
* @returns {Promise<Object>} The resulting policy
* @throws {Error} If serviceId is invalid
*/
async setPolicy(serviceId, policy) {
if (!serviceId || typeof serviceId !== 'string') {
throw new Error('serviceId is required');
}
const existing = this.policies.get(serviceId) || { ...DEFAULT_POLICY, serviceId };
const merged = {
...existing,
...policy,
serviceId,
// Never allow caller to override runtime counters directly
currentRetries: existing.currentRetries || 0,
lastRestartAt: existing.lastRestartAt,
cooldownUntil: existing.cooldownUntil,
};
this.policies.set(serviceId, merged);
await this._savePolicies();
this.log.info('auto-restart', 'Policy set', { serviceId, enabled: merged.enabled });
return { ...merged };
}
/**
* Retrieve the policy for a service.
*
* @param {string} serviceId
* @returns {Object|null} Policy object or null if none exists
*/
getPolicy(serviceId) {
const policy = this.policies.get(serviceId);
return policy ? { ...policy } : null;
}
/**
* Return all policies as an array.
* @returns {Object[]}
*/
listPolicies() {
return Array.from(this.policies.values()).map(p => ({ ...p }));
}
/**
* Remove a service's restart policy.
*
* @param {string} serviceId
* @returns {Promise<boolean>} true if a policy was removed
*/
async removePolicy(serviceId) {
if (!this.policies.has(serviceId)) return false;
this.policies.delete(serviceId);
await this._savePolicies();
this.log.info('auto-restart', 'Policy removed', { serviceId });
return true;
}
// ─── Core Restart Logic ──────────────────────────────────────────────
/**
* Called when a container is detected as down.
*
* Checks policy, cooldown, and retry count, then either attempts a
* Docker restart or notifies that max retries were exceeded.
*
* @param {string} serviceId - Service identifier
* @param {string} containerId - Docker container ID to restart
* @returns {Promise<Object>} Result of the operation
*/
async handleContainerDown(serviceId, containerId) {
const policy = this.policies.get(serviceId);
if (!policy) {
return { action: 'ignored', reason: 'no-policy' };
}
if (!policy.enabled) {
return { action: 'ignored', reason: 'disabled' };
}
// Check cooldown window
const now = Date.now();
if (policy.cooldownUntil && now < policy.cooldownUntil) {
this.log.info('auto-restart', 'Skipping — cooldown active', {
serviceId,
cooldownUntil: new Date(policy.cooldownUntil).toISOString(),
});
return { action: 'skipped', reason: 'cooldown' };
}
// Max retries exceeded — notify and enter cooldown
if (policy.currentRetries >= policy.maxRetries) {
const cooldownMs = policy.windowMinutes * 60 * 1000;
policy.cooldownUntil = now + cooldownMs;
policy.currentRetries = 0; // Reset so next window can try again
await this._savePolicies();
const eventData = {
serviceId,
containerId,
maxRetries: policy.maxRetries,
cooldownUntil: policy.cooldownUntil,
timestamp: new Date().toISOString(),
};
/**
* @event AutoRestartManager#auto-restart-max-reached
* @type {Object}
*/
this.emit('auto-restart-max-reached', eventData);
// Send notification
try {
await this._notify('auto-restart', {
containerName: serviceId,
message: `⛔ Max auto-restart retries (${policy.maxRetries}) exceeded for "${serviceId}". Cooldown until ${new Date(policy.cooldownUntil).toISOString()}.`,
...eventData,
});
} catch (notifErr) {
this.log.error('auto-restart', 'Notification failed', { error: notifErr.message });
}
return { action: 'max-reached', ...eventData };
}
// Wait for the configured retry interval before attempting
if (policy.retryIntervalMs > 0 && policy.lastRestartAt) {
const elapsed = now - new Date(policy.lastRestartAt).getTime();
if (elapsed < policy.retryIntervalMs) {
const waitMs = policy.retryIntervalMs - elapsed;
this.log.info('auto-restart', 'Waiting for retry interval', { serviceId, waitMs });
await new Promise(resolve => setTimeout(resolve, waitMs));
}
}
// Attempt restart
policy.currentRetries += 1;
const attemptNum = policy.currentRetries;
const maxRetries = policy.maxRetries;
/**
* @event AutoRestartManager#auto-restart-attempt
* @type {Object}
*/
this.emit('auto-restart-attempt', {
serviceId,
containerId,
attempt: attemptNum,
maxRetries,
timestamp: new Date().toISOString(),
});
try {
if (!this.docker?.client) {
throw new Error('Docker client not available');
}
const container = this.docker.client.getContainer(containerId);
await container.start();
policy.lastRestartAt = new Date().toISOString();
await this._savePolicies();
const successData = {
serviceId,
containerId,
attempt: attemptNum,
maxRetries,
timestamp: new Date().toISOString(),
};
/**
* @event AutoRestartManager#auto-restart-success
* @type {Object}
*/
this.emit('auto-restart-success', successData);
// Notify
try {
await this._notify('auto-restart', {
containerName: serviceId,
message: `🔄 Auto-restart attempt ${attemptNum}/${maxRetries} succeeded for "${serviceId}".`,
...successData,
});
} catch (notifErr) {
this.log.error('auto-restart', 'Notification failed', { error: notifErr.message });
}
this.log.info('auto-restart', 'Container restarted', {
serviceId,
attempt: attemptNum,
maxRetries,
});
return { action: 'restarted', ...successData };
} catch (restartErr) {
policy.lastRestartAt = new Date().toISOString();
await this._savePolicies();
const failData = {
serviceId,
containerId,
attempt: attemptNum,
maxRetries,
error: restartErr.message,
timestamp: new Date().toISOString(),
};
/**
* @event AutoRestartManager#auto-restart-failed
* @type {Object}
*/
this.emit('auto-restart-failed', failData);
// Notify
try {
await this._notify('auto-restart', {
containerName: serviceId,
message: `❌ Auto-restart attempt ${attemptNum}/${maxRetries} failed for "${serviceId}": ${restartErr.message}`,
...failData,
});
} catch (notifErr) {
this.log.error('auto-restart', 'Notification failed', { error: notifErr.message });
}
this.log.error('auto-restart', 'Restart failed', {
serviceId,
attempt: attemptNum,
error: restartErr.message,
});
return { action: 'failed', ...failData };
}
}
/**
* Called when a container recovers to healthy state.
* Resets the retry counter for the associated service.
*
* @param {string} serviceId
* @returns {Promise<void>}
*/
async handleContainerUp(serviceId) {
const policy = this.policies.get(serviceId);
if (!policy) return;
if (policy.currentRetries > 0) {
policy.currentRetries = 0;
policy.cooldownUntil = null;
await this._savePolicies();
this.log.info('auto-restart', 'Retries reset after recovery', { serviceId });
}
}
// ─── Health Event Bridge ─────────────────────────────────────────────
/**
* Internal handler for HealthChecker `status-check` events.
* Detects healthy→unhealthy and unhealthy→healthy transitions for tracked services.
*
* @param {Object} status - HealthChecker status object
* @param {string} status.serviceId
* @param {string} status.status - "up" or "down"
* @private
*/
async _handleStatusCheck(status) {
const { serviceId, status: currentStatus } = status;
if (!serviceId) return;
// Only process services that have a restart policy
if (!this.policies.has(serviceId)) return;
const previousStatus = this._previousHealth.get(serviceId);
this._previousHealth.set(serviceId, currentStatus);
// Transition: healthy → unhealthy
if (previousStatus === 'up' && currentStatus === 'down') {
// Find the containerId from the health checker config or status details
const containerId = this._resolveContainerId(serviceId, status);
if (containerId) {
try {
await this.handleContainerDown(serviceId, containerId);
} catch (err) {
this.logError('auto-restart-health-bridge', err);
}
}
}
// Transition: unhealthy → healthy (recovery)
if (previousStatus === 'down' && currentStatus === 'up') {
try {
await this.handleContainerUp(serviceId);
} catch (err) {
this.logError('auto-restart-health-bridge', err);
}
}
}
/**
* Attempt to find the containerId for a service from various sources.
*
* @param {string} serviceId
* @param {Object} status - The status-check event data
* @returns {string|null}
* @private
*/
_resolveContainerId(serviceId, status) {
// Check if it's in the status details (some health checks embed it)
if (status.details?.containerId) return status.details.containerId;
// Look in the health checker config
const hcService = this.healthChecker?.config?.services?.[serviceId];
if (hcService?.containerId) return hcService.containerId;
// Try to look it up from the services state manager
try {
const servicesStateManager = this.ctx.servicesStateManager;
if (servicesStateManager) {
const readResult = servicesStateManager.read();
if (readResult && typeof readResult.then === 'function') {
// It returns a promise — fire-and-forget lookup
readResult.then(list => {
const found = (list || []).find(s => s.id === serviceId);
return found?.containerId || null;
}).catch(() => null);
} else {
const found = (readResult || []).find(s => s.id === serviceId);
if (found?.containerId) return found.containerId;
}
}
} catch (_) { /* best effort */ }
return null;
}
// ─── Persistence ─────────────────────────────────────────────────────
/**
* Persist current policies to disk.
* @returns {Promise<void>}
* @private
*/
async _savePolicies() {
try {
const obj = {};
for (const [serviceId, policy] of this.policies.entries()) {
obj[serviceId] = { ...policy };
}
await writeJsonFile(this.policiesFile, obj);
} catch (err) {
this.log.error('auto-restart', 'Failed to save policies', { error: err.message });
}
}
// ─── Helpers ─────────────────────────────────────────────────────────
/**
* Send a notification via the notification manager.
*
* @param {string} event - Event type (e.g. 'auto-restart')
* @param {Object} data - Notification payload
* @returns {Promise<Object>}
* @private
*/
async _notify(event, data) {
if (this.notification?.send) {
return this.notification.send(event, data);
}
return { success: false, reason: 'no-notification-manager' };
}
}
module.exports = { AutoRestartManager, DEFAULT_POLICY };
+376
View File
@@ -0,0 +1,376 @@
/**
* Config Drift Detector - Compares services.json with live Docker state
*
* Detects discrepancies between the configured service list and what is
* actually running in Docker, including missing containers, unknown
* containers, port mismatches, state mismatches, and stale records.
*
* @module config-drift-detector
*/
const EventEmitter = require('events');
/**
* @typedef {Object} DriftReport
* @property {string} checkedAt - ISO timestamp of the check
* @property {Object[]} missingContainers - Services with containerId but container absent in Docker
* @property {Object[]} unknownContainers - Running Docker containers with sami.managed label but not in services.json
* @property {Object[]} portMismatch - Service port != container mapped port
* @property {Object[]} stateMismatch - Service expected up but container stopped/absent
* @property {Object[]} staleRecords - Services with containerId pointing to removed containers
* @property {boolean} hasDrift - Whether any drift category is non-empty
*/
/**
* Detects and reports configuration drift between services.json and Docker.
*
* @extends EventEmitter
*
* @fires ConfigDriftDetector#drift-detected
*/
class ConfigDriftDetector extends EventEmitter {
/**
* @param {Object} ctx - Shared application context
* @param {Object} ctx.docker - Docker client wrapper ({ client: Dockerode })
* @param {Object} ctx.servicesStateManager - StateManager for services.json
* @param {Object} ctx.notification - NotificationManager instance
* @param {Object} ctx.log - Logger instance
* @param {Function} ctx.logError - Error logging function
*/
constructor(ctx) {
super();
this.ctx = ctx;
this.log = ctx.log || console;
this.logError = ctx.logError || ((_c, err) => console.error(err));
this.docker = ctx.docker;
this.servicesStateManager = ctx.servicesStateManager;
this.notification = ctx.notification;
/** @type {DriftReport|null} Cached report from last detection */
this.lastReport = null;
/** @type {NodeJS.Timeout|null} Polling timer reference */
this._pollTimer = null;
/** Whether polling is currently active */
this._polling = false;
}
// ─── Detection ───────────────────────────────────────────────────────
/**
* Run a full drift detection and return the report.
*
* Reads services from servicesStateManager and live containers from Docker,
* then compares them across five drift categories.
*
* @returns {Promise<DriftReport>}
*/
async detect() {
const checkedAt = new Date().toISOString();
// Gather configured services
let services = [];
try {
const data = await this.servicesStateManager.read();
services = Array.isArray(data) ? data : (data.services || []);
} catch (err) {
this.log.error('drift', 'Failed to read services', { error: err.message });
}
// Gather live Docker containers
let containers = [];
try {
containers = await this.docker.client.listContainers({ all: true });
} catch (err) {
this.log.error('drift', 'Failed to list containers', { error: err.message });
}
// Build lookup maps
const containerById = new Map(); // containerId (short or long) → container info
const containerByName = new Map(); // container name → container info
for (const c of containers) {
// Store by full ID
containerById.set(c.Id, c);
// Store by short ID (first 12 chars)
if (c.Id && c.Id.length >= 12) {
containerById.set(c.Id.substring(0, 12), c);
}
// Store by name (strip leading /)
for (const name of (c.Names || [])) {
containerByName.set(name.replace(/^\//, ''), c);
}
}
// Build set of service containerIds for reverse lookup
const serviceContainerIds = new Set();
const serviceByContainerId = new Map();
for (const svc of services) {
if (svc.containerId) {
serviceContainerIds.add(svc.containerId);
// Index by both full and short ID
serviceByContainerId.set(svc.containerId, svc);
if (svc.containerId.length >= 12) {
serviceByContainerId.set(svc.containerId.substring(0, 12), svc);
}
}
}
const missingContainers = [];
const portMismatch = [];
const stateMismatch = [];
const staleRecords = [];
for (const svc of services) {
if (!svc.containerId) continue;
// Look up the container
const container = containerById.get(svc.containerId)
|| containerById.get(svc.containerId.substring(0, 12));
if (!container) {
// Container ID referenced but not found in Docker at all
staleRecords.push({
serviceId: svc.id,
name: svc.name,
containerId: svc.containerId,
reason: 'Container not found in Docker',
});
continue;
}
// Missing container — service expects it but it's not running
if (container.State !== 'running') {
missingContainers.push({
serviceId: svc.id,
name: svc.name,
containerId: svc.containerId,
containerState: container.State,
containerStatus: container.Status,
});
// Also a state mismatch if the service is expected to be up
stateMismatch.push({
serviceId: svc.id,
name: svc.name,
expectedState: 'running',
actualState: container.State,
containerId: svc.containerId,
});
}
// Port mismatch detection
if (svc.port && container.State === 'running') {
const actualPorts = this._extractContainerPorts(container);
if (actualPorts.length > 0 && !actualPorts.includes(svc.port)) {
portMismatch.push({
serviceId: svc.id,
name: svc.name,
configuredPort: svc.port,
actualPorts,
containerId: svc.containerId,
});
}
}
}
// Unknown managed containers: Docker containers with sami.managed label
// that are NOT in services.json
const unknownContainers = [];
for (const c of containers) {
const isManaged = c.Labels && c.Labels['sami.managed'] === 'true';
if (!isManaged) continue;
const isInServices = serviceByContainerId.has(c.Id)
|| serviceByContainerId.has(c.Id.substring(0, 12));
if (!isInServices) {
unknownContainers.push({
containerId: c.Id,
name: (c.Names && c.Names[0] || '').replace(/^\//, ''),
image: c.Image,
state: c.State,
status: c.Status,
app: c.Labels?.['sami.app'] || null,
subdomain: c.Labels?.['sami.subdomain'] || null,
});
}
}
const report = {
checkedAt,
missingContainers,
unknownContainers,
portMismatch,
stateMismatch,
staleRecords,
hasDrift: missingContainers.length > 0
|| unknownContainers.length > 0
|| portMismatch.length > 0
|| stateMismatch.length > 0
|| staleRecords.length > 0,
};
// Cache for quick API access
this.lastReport = report;
// Emit and notify if drift detected
if (report.hasDrift) {
/**
* @event ConfigDriftDetector#drift-detected
* @type {DriftReport}
*/
this.emit('drift-detected', report);
try {
await this._sendDriftNotification(report);
} catch (notifErr) {
this.log.error('drift', 'Failed to send drift notification', {
error: notifErr.message,
});
}
}
this.log.info('drift', 'Detection complete', {
hasDrift: report.hasDrift,
missing: report.missingContainers.length,
unknown: report.unknownContainers.length,
portMismatch: report.portMismatch.length,
stateMismatch: report.stateMismatch.length,
stale: report.staleRecords.length,
});
return report;
}
// ─── Auto-fix ────────────────────────────────────────────────────────
/**
* Attempt to auto-fix drift:
* - Remove stale records (services referencing removed containers)
* - Flag unknown containers for review
*
* @returns {Promise<{ staleRemoved: number, unknownFlagged: number }>}
*/
async autoFix() {
const report = await this.detect();
let staleRemoved = 0;
// Remove stale records from services.json
if (report.staleRecords.length > 0) {
const staleIds = new Set(report.staleRecords.map(r => r.serviceId));
await this.servicesStateManager.update(services => {
const before = services.length;
const cleaned = services.filter(s => !staleIds.has(s.id));
staleRemoved = before - cleaned.length;
return cleaned;
});
}
const unknownFlagged = report.unknownContainers.length;
this.log.info('drift', 'Auto-fix applied', { staleRemoved, unknownFlagged });
return { staleRemoved, unknownFlagged };
}
// ─── Polling ─────────────────────────────────────────────────────────
/**
* Start periodic drift detection.
*
* @param {number} [intervalMs=300000] - Polling interval in milliseconds (default 5 min)
*/
startPolling(intervalMs = 300000) {
this.stopPolling();
this._polling = true;
this._pollTimer = setInterval(async () => {
try {
await this.detect();
} catch (err) {
this.logError('drift-poll', err);
}
}, intervalMs);
this.log.info('drift', 'Polling started', { intervalMs });
}
/**
* Stop periodic drift detection.
*/
stopPolling() {
if (this._pollTimer) {
clearInterval(this._pollTimer);
this._pollTimer = null;
}
this._polling = false;
this.log.info('drift', 'Polling stopped');
}
/**
* Whether polling is currently active.
* @returns {boolean}
*/
isPolling() {
return this._polling;
}
// ─── Helpers ─────────────────────────────────────────────────────────
/**
* Extract mapped host ports from a Docker container info object.
*
* @param {Object} container - Dockerode container info
* @returns {number[]} Array of host port numbers
* @private
*/
_extractContainerPorts(container) {
const ports = [];
if (!container.Ports) return ports;
for (const p of container.Ports) {
if (p.PublicPort) {
ports.push(p.PublicPort);
}
}
return ports;
}
/**
* Send a notification about detected drift.
*
* @param {DriftReport} report
* @returns {Promise<Object>}
* @private
*/
async _sendDriftNotification(report) {
if (!this.notification?.send) {
return { success: false, reason: 'no-notification-manager' };
}
const parts = [];
if (report.missingContainers.length > 0) {
parts.push(`Missing containers: ${report.missingContainers.map(c => c.name).join(', ')}`);
}
if (report.unknownContainers.length > 0) {
parts.push(`Unknown managed containers: ${report.unknownContainers.map(c => c.name).join(', ')}`);
}
if (report.portMismatch.length > 0) {
parts.push(`Port mismatches: ${report.portMismatch.map(c => c.name).join(', ')}`);
}
if (report.staleRecords.length > 0) {
parts.push(`Stale records: ${report.staleRecords.map(c => c.name).join(', ')}`);
}
return this.notification.send('drift-detected', {
text: `⚠️ Configuration drift detected:\n${parts.join('\n')}`,
report,
});
}
}
module.exports = { ConfigDriftDetector };
+605
View File
@@ -0,0 +1,605 @@
/**
* Dependency Manager - Service dependency tracking with ordered restart chains
*
* Manages directed acyclic graph (DAG) of service dependencies. Services can
* declare which other services they depend on, and this manager provides:
* - Full dependency graph inspection
* - Topological ordering for safe restart chains
* - Circular dependency detection
* - Health-aware restart with per-service polling
*
* Dependencies are stored directly on service objects in services.json:
* { id, name, ..., dependsOn: ['service-id-1', 'service-id-2'] }
*
* @module dependency-manager
*/
const EventEmitter = require('events');
/** Maximum seconds to wait for a single container to become healthy after restart */
const HEALTH_CHECK_TIMEOUT_MS = 30_000;
/** Interval between container health polls */
const HEALTH_CHECK_INTERVAL_MS = 1_000;
/**
* @typedef {Object} ServiceNode
* @property {string} serviceId
* @property {string} name
* @property {string|null} containerId
*/
/**
* @typedef {Object} DependencyEdge
* @property {string} from - The service that depends
* @property {string} to - The service being depended upon
*/
/**
* @typedef {Object} DependencyGraph
* @property {ServiceNode[]} nodes
* @property {DependencyEdge[]} edges
*/
/**
* @typedef {Object} DependencyStatusEntry
* @property {string} serviceId
* @property {string} name
* @property {boolean} isUp
* @property {string} [error]
*/
/**
* DependencyManager — tracks service dependencies and orchestrates ordered restarts.
*
* Events emitted:
* - `dependency-restart-start` ({ serviceId, chain: string[] })
* - `dependency-restart-progress` ({ serviceId, currentServiceId, index, total })
* - `dependency-restart-complete` ({ serviceId, chain: string[], results: Array })
* - `dependency-restart-failed` ({ serviceId, failedServiceId, error, chain: string[] })
*
* @extends EventEmitter
*/
class DependencyManager extends EventEmitter {
/**
* @param {Object} ctx - Application context
* @param {Object} ctx.servicesStateManager - StateManager for services.json
* @param {Object} ctx.docker - Docker context ({ client: Dockerode })
* @param {Object} ctx.notification - NotificationManager instance
* @param {Object} ctx.log - Logger instance
*/
constructor(ctx) {
super();
/** @private */
this.ctx = ctx;
/** @private */
this._servicesStateManager = ctx.servicesStateManager;
/** @private */
this._docker = ctx.docker;
/** @private */
this._notification = ctx.notification;
/** @private */
this._log = ctx.log || console;
}
// ---------------------------------------------------------------------------
// Core helpers
// ---------------------------------------------------------------------------
/**
* Load all services from the state manager.
* @private
* @returns {Promise<Object[]>}
*/
async _loadServices() {
const data = await this._servicesStateManager.read();
return Array.isArray(data) ? data : (data.services || []);
}
/**
* Find a single service by ID.
* @private
* @param {string} serviceId
* @returns {Promise<Object|null>}
*/
async _findService(serviceId) {
const services = await this._loadServices();
return services.find(s => s.id === serviceId) || null;
}
// ---------------------------------------------------------------------------
// Graph queries
// ---------------------------------------------------------------------------
/**
* Return the full dependency graph for visualisation.
*
* @returns {Promise<DependencyGraph>}
*/
async getDependencyGraph() {
const services = await this._loadServices();
const nodes = services.map(s => ({
serviceId: s.id,
name: s.name,
containerId: s.containerId || null,
}));
const edges = [];
for (const service of services) {
const deps = service.dependsOn || [];
for (const depId of deps) {
edges.push({ from: service.id, to: depId });
}
}
return { nodes, edges };
}
/**
* Return the services that depend on the given service (reverse deps).
*
* @param {string} serviceId
* @returns {Promise<Object[]>} Services whose `dependsOn` includes `serviceId`.
*/
async getDependents(serviceId) {
const services = await this._loadServices();
return services.filter(s => (s.dependsOn || []).includes(serviceId));
}
/**
* Return the direct dependencies for a service.
*
* @param {string} serviceId
* @returns {Promise<Object[]>} Services that `serviceId` depends on.
*/
async getDependencies(serviceId) {
const services = await this._loadServices();
const service = services.find(s => s.id === serviceId);
if (!service) return [];
const depIds = service.dependsOn || [];
return services.filter(s => depIds.includes(s.id));
}
// ---------------------------------------------------------------------------
// Topological sort
// ---------------------------------------------------------------------------
/**
* Build an adjacency list for the current dependency graph.
* Edge direction: service → its dependencies (i.e. what it depends on).
*
* @private
* @param {Object[]} services
* @returns {Map<string, string[]>}
*/
_buildAdjacencyList(services) {
const adj = new Map();
for (const service of services) {
adj.set(service.id, (service.dependsOn || []).slice());
}
return adj;
}
/**
* DFS-based topological sort with cycle detection (white/gray/black coloring).
*
* Returns services in restart order: dependencies first, dependents last.
* The target service is included at the end.
*
* @private
* @param {string} serviceId - Target service (will be last in the result).
* @param {Object[]} services - All services.
* @param {Map<string, string[]>} adj - Adjacency list (service → deps).
* @returns {string[]} Ordered service IDs for restart.
* @throws {Error} If a circular dependency is detected.
*/
_topologicalSort(serviceId, services, adj) {
// Collect only the reachable sub-graph from serviceId
const visited = new Set();
const reachable = new Set();
const collectReachable = (id) => {
if (reachable.has(id)) return;
reachable.add(id);
for (const dep of (adj.get(id) || [])) {
collectReachable(dep);
}
};
collectReachable(serviceId);
// DFS topological sort on the reachable sub-graph
const WHITE = 0, GRAY = 1, BLACK = 2;
const color = new Map();
for (const id of reachable) color.set(id, WHITE);
const result = [];
const dfs = (id) => {
if (color.get(id) === BLACK) return;
if (color.get(id) === GRAY) {
throw new Error(`Circular dependency detected involving service "${id}"`);
}
color.set(id, GRAY);
for (const dep of (adj.get(id) || [])) {
dfs(dep);
}
color.set(id, BLACK);
result.push(id);
};
// Visit the target last so it ends up at the end of the result
// Actually, we want deps *first* then the target.
// The DFS naturally puts deps before dependents, so starting from
// serviceId will place it last (which is correct for restart order).
dfs(serviceId);
return result;
}
/**
* Get the topologically ordered restart chain for a service.
*
* The returned array lists all services that must be restarted,
* starting with leaf dependencies and ending with the target service.
*
* @param {string} serviceId - The service to build the chain for.
* @returns {Promise<string[]>} Ordered service IDs.
* @throws {Error} If `serviceId` doesn't exist or a circular dependency is found.
*/
async getOrderedRestartChain(serviceId) {
const services = await this._loadServices();
const service = services.find(s => s.id === serviceId);
if (!service) {
throw new Error(`Service "${serviceId}" not found`);
}
const adj = this._buildAdjacencyList(services);
return this._topologicalSort(serviceId, services, adj);
}
// ---------------------------------------------------------------------------
// Validation
// ---------------------------------------------------------------------------
/**
* Validate a proposed set of dependencies for a service.
*
* Checks:
* - All referenced service IDs exist.
* - Adding these dependencies would not create a circular dependency.
* - A service cannot depend on itself.
*
* @param {string} serviceId - The service to set dependencies on.
* @param {string[]} dependsOn - Proposed dependency IDs.
* @returns {Promise<{ valid: boolean, errors: string[] }>}
*/
async validateDependencies(serviceId, dependsOn) {
const errors = [];
if (!Array.isArray(dependsOn)) {
return { valid: false, errors: ['dependsOn must be an array'] };
}
const services = await this._loadServices();
const allIds = new Set(services.map(s => s.id));
// Service must exist
if (!allIds.has(serviceId)) {
return { valid: false, errors: [`Service "${serviceId}" not found`] };
}
// Self-dependency
if (dependsOn.includes(serviceId)) {
errors.push(`Service "${serviceId}" cannot depend on itself`);
}
// Existence check
for (const depId of dependsOn) {
if (!allIds.has(depId)) {
errors.push(`Dependency service "${depId}" does not exist`);
}
}
if (errors.length > 0) {
return { valid: false, errors };
}
// Circular dependency check: temporarily set the proposed dependsOn
// and attempt a topological sort.
const tempServices = services.map(s => {
if (s.id === serviceId) {
return { ...s, dependsOn: dependsOn.slice() };
}
return { ...s };
});
const adj = this._buildAdjacencyList(tempServices);
// Check every node for cycles with the new edges
try {
const WHITE = 0, GRAY = 1, BLACK = 2;
const color = new Map();
for (const s of tempServices) color.set(s.id, WHITE);
const dfs = (id) => {
if (color.get(id) === BLACK) return;
if (color.get(id) === GRAY) {
throw new Error(`Circular dependency detected involving service "${id}"`);
}
color.set(id, GRAY);
for (const dep of (adj.get(id) || [])) {
dfs(dep);
}
color.set(id, BLACK);
};
for (const s of tempServices) {
if (color.get(s.id) === WHITE) {
dfs(s.id);
}
}
} catch (err) {
errors.push(err.message);
}
return { valid: errors.length === 0, errors };
}
// ---------------------------------------------------------------------------
// Health status
// ---------------------------------------------------------------------------
/**
* Get the current container status for a service and all its transitive dependencies.
*
* @param {string} serviceId
* @returns {Promise<DependencyStatusEntry[]>}
* @throws {Error} If `serviceId` doesn't exist.
*/
async getDependencyStatus(serviceId) {
const services = await this._loadServices();
const service = services.find(s => s.id === serviceId);
if (!service) {
throw new Error(`Service "${serviceId}" not found`);
}
// Collect all transitive dependencies via BFS
const serviceMap = new Map(services.map(s => [s.id, s]));
const visited = new Set();
const queue = [serviceId];
const allRelated = [];
while (queue.length > 0) {
const currentId = queue.shift();
if (visited.has(currentId)) continue;
visited.add(currentId);
const svc = serviceMap.get(currentId);
if (!svc) continue;
allRelated.push(svc);
for (const depId of (svc.dependsOn || [])) {
if (!visited.has(depId)) {
queue.push(depId);
}
}
}
// Query container status for each
const results = [];
for (const svc of allRelated) {
const entry = {
serviceId: svc.id,
name: svc.name,
isUp: false,
};
if (!svc.containerId) {
entry.error = 'No container associated with this service';
results.push(entry);
continue;
}
try {
const container = this._docker.client.getContainer(svc.containerId);
const info = await container.inspect();
entry.isUp = info.State?.Running === true;
} catch (err) {
entry.error = err.message || 'Unable to inspect container';
}
results.push(entry);
}
return results;
}
// ---------------------------------------------------------------------------
// Restart with dependencies
// ---------------------------------------------------------------------------
/**
* Wait for a container to report as running after a restart.
*
* @private
* @param {string} containerId
* @param {number} [timeoutMs=30000]
* @returns {Promise<boolean>} `true` if healthy, `false` if timed out.
*/
async _waitForContainerHealthy(containerId, timeoutMs = HEALTH_CHECK_TIMEOUT_MS) {
const start = Date.now();
while (Date.now() - start < timeoutMs) {
try {
const container = this._docker.client.getContainer(containerId);
const info = await container.inspect();
if (info.State?.Running === true) {
return true;
}
} catch {
// Container might not be inspectable during restart — keep polling
}
await new Promise(r => setTimeout(r, HEALTH_CHECK_INTERVAL_MS));
}
return false;
}
/**
* Restart a service and all its dependencies in topological order.
*
* Emits progress events and sends a notification on completion/failure.
* This method is designed to be called from the route handler and
* **does not throw** — errors are reported via events and notifications.
*
* @param {string} serviceId - Target service to restart (with deps).
* @returns {Promise<{ success: boolean, chain: string[], results: Array }>}
*/
async restartWithDependencies(serviceId) {
const service = await this._findService(serviceId);
if (!service) {
const err = new Error(`Service "${serviceId}" not found`);
this.emit('dependency-restart-failed', {
serviceId,
failedServiceId: serviceId,
error: err.message,
chain: [],
});
throw err;
}
let chain;
try {
chain = await this.getOrderedRestartChain(serviceId);
} catch (err) {
this.emit('dependency-restart-failed', {
serviceId,
failedServiceId: serviceId,
error: err.message,
chain: [],
});
throw err;
}
const services = await this._loadServices();
const serviceMap = new Map(services.map(s => [s.id, s]));
this._log.info('dependency', 'Starting dependency restart chain', {
serviceId,
chain,
});
this.emit('dependency-restart-start', { serviceId, chain });
const results = [];
const total = chain.length;
for (let i = 0; i < total; i++) {
const currentId = chain[i];
const svc = serviceMap.get(currentId);
this.emit('dependency-restart-progress', {
serviceId,
currentServiceId: currentId,
index: i,
total,
});
if (!svc || !svc.containerId) {
const msg = !svc
? `Service "${currentId}" not found in state`
: `Service "${currentId}" has no container — skipping restart`;
this._log.warn('dependency', msg);
results.push({ serviceId: currentId, restarted: false, skipped: true, reason: msg });
continue;
}
try {
const container = this._docker.client.getContainer(svc.containerId);
this._log.info('dependency', `Restarting container for service "${currentId}"`, {
containerId: svc.containerId,
});
await container.restart();
// Wait for it to come back up
const healthy = await this._waitForContainerHealthy(svc.containerId);
if (!healthy) {
const msg = `Container for service "${currentId}" did not become healthy within ${HEALTH_CHECK_TIMEOUT_MS / 1000}s`;
this._log.warn('dependency', msg);
results.push({ serviceId: currentId, restarted: true, healthy: false, error: msg });
// Abort chain — dependency didn't come back
this.emit('dependency-restart-failed', {
serviceId,
failedServiceId: currentId,
error: msg,
chain,
});
await this._notifyRestartResult(serviceId, false, chain, results, currentId);
return { success: false, chain, results };
}
this._log.info('dependency', `Service "${currentId}" is healthy after restart`);
results.push({ serviceId: currentId, restarted: true, healthy: true });
} catch (err) {
const msg = err.message || 'Unknown error during restart';
this._log.error('dependency', `Failed to restart service "${currentId}"`, {
error: msg,
});
results.push({ serviceId: currentId, restarted: false, error: msg });
this.emit('dependency-restart-failed', {
serviceId,
failedServiceId: currentId,
error: msg,
chain,
});
await this._notifyRestartResult(serviceId, false, chain, results, currentId);
return { success: false, chain, results };
}
}
this.emit('dependency-restart-complete', { serviceId, chain, results });
await this._notifyRestartResult(serviceId, true, chain, results);
return { success: true, chain, results };
}
/**
* Send a notification about the restart result.
*
* @private
* @param {string} serviceId
* @param {boolean} success
* @param {string[]} chain
* @param {Array} results
* @param {string} [failedServiceId]
*/
async _notifyRestartResult(serviceId, success, chain, results, failedServiceId) {
if (!this._notification) return;
try {
if (success) {
await this._notification.send('dependency-restart-complete', {
text: `✅ Dependency restart chain completed for "${serviceId}". Restarted: ${chain.join(' → ')}`,
serviceId,
chain,
results,
});
} else {
await this._notification.send('dependency-restart-failed', {
text: `❌ Dependency restart chain failed for "${serviceId}" at "${failedServiceId}". Chain: ${chain.join(' → ')}`,
serviceId,
failedServiceId,
chain,
results,
});
}
} catch (err) {
this._log.error('dependency', 'Failed to send restart notification', {
error: err.message,
});
}
}
}
module.exports = DependencyManager;
+273
View File
@@ -0,0 +1,273 @@
/**
* DNS Propagation Checker
* Verifies DNS record propagation by querying multiple resolvers.
* Runs as background jobs with configurable timeout and interval.
*
* @module dns-propagation
*/
const dns = require('dns').promises;
const EventEmitter = require('events');
/** Default verification options */
const DEFAULT_OPTIONS = {
timeout: 300000, // 5 minutes
interval: 10000, // 10 seconds
resolvers: ['1.1.1.1', '8.8.8.8', '9.9.9.9']
};
/** Maximum age for stored verification results (1 hour) */
const MAX_RESULT_AGE_MS = 3600000;
class DNSPropagationChecker extends EventEmitter {
/**
* Create a DNSPropagationChecker instance.
* @param {Object} ctx - Shared application context
* @param {Object} ctx.notification - NotificationManager instance
* @param {Object} ctx.log - Logger instance
*/
constructor(ctx) {
super();
this.ctx = ctx;
this.log = ctx.log || console;
/** @type {Map<string, Object>} domain → verification status */
this.verifications = new Map();
}
/**
* Verify that a DNS record has propagated by querying multiple resolvers.
* Retries every `interval` ms until `timeout` is reached.
*
* @param {string} domain - The domain to check (e.g., 'test.sami')
* @param {string} expectedIp - The expected IP address
* @param {Object} [options={}] - Verification options
* @param {number} [options.timeout=300000] - Maximum time to wait (ms)
* @param {number} [options.interval=10000] - Time between retries (ms)
* @param {string[]} [options.resolvers] - DNS resolvers to query
* @returns {Promise<Object>} Verification result
*/
async verifyRecord(domain, expectedIp, options = {}) {
const startTime = Date.now();
const {
timeout = DEFAULT_OPTIONS.timeout,
interval = DEFAULT_OPTIONS.interval,
resolvers = DEFAULT_OPTIONS.resolvers
} = options;
const allResults = [];
let propagated = false;
while (Date.now() - startTime < timeout) {
const roundResults = [];
for (const resolver of resolvers) {
const checkStart = Date.now();
try {
// Use dns.resolve4 with a custom resolver
const resolverInstance = new dns.Resolver();
resolverInstance.setServers([resolver]);
resolverInstance.setTimeout(5000);
const addresses = await resolverInstance.resolve4(domain);
const matched = addresses.includes(expectedIp);
const result = {
resolver,
ips: addresses,
matched,
checkedAt: new Date().toISOString(),
responseTime: Date.now() - checkStart
};
roundResults.push(result);
if (matched) {
propagated = true;
}
} catch (err) {
roundResults.push({
resolver,
ips: [],
matched: false,
checkedAt: new Date().toISOString(),
error: err.code || err.message,
responseTime: Date.now() - checkStart
});
}
}
allResults.push(...roundResults);
// Emit progress event
this.emit('propagation-check', {
domain,
expectedIp,
roundResults,
elapsed: Date.now() - startTime,
propagated
});
if (propagated) {
break;
}
// Wait before next attempt
await new Promise(resolve => setTimeout(resolve, interval));
}
const totalTime = Date.now() - startTime;
return {
domain,
expectedIp,
propagated,
results: allResults,
totalTime,
checkedAt: new Date().toISOString()
};
}
/**
* Start a background DNS propagation verification.
* Does not block — returns immediately with the job reference.
*
* @param {string} domain - The domain to verify
* @param {string} expectedIp - The expected IP address
* @param {Object} [options={}] - Verification options
* @returns {Object} Job status object
*/
startVerification(domain, expectedIp, options = {}) {
// If there's already a running verification for this domain, return it
const existing = this.verifications.get(domain);
if (existing && existing.status === 'running') {
return existing;
}
const job = {
domain,
expectedIp,
status: 'running',
startedAt: new Date().toISOString(),
progress: [],
result: null
};
this.verifications.set(domain, job);
// Run verification in background (non-blocking)
this.verifyRecord(domain, expectedIp, options)
.then(result => {
job.status = 'completed';
job.result = result;
job.completedAt = new Date().toISOString();
if (result.propagated) {
this.emit('propagation-complete', result);
if (this.ctx.notification) {
this.ctx.notification.send('dns-propagation', {
text: `✅ DNS record for ${domain} propagated successfully to ${expectedIp}`,
domain,
expectedIp,
totalTime: result.totalTime
}, 'success').catch(err => {
this.log.error('dns-propagation', 'Failed to send propagation notification', {
error: err.message
});
});
}
} else {
this.emit('propagation-timeout', result);
if (this.ctx.notification) {
this.ctx.notification.send('dns-propagation', {
text: `⏱️ DNS propagation timeout for ${domain} — expected ${expectedIp} not found after ${Math.round(result.totalTime / 1000)}s`,
domain,
expectedIp,
totalTime: result.totalTime
}, 'warning').catch(err => {
this.log.error('dns-propagation', 'Failed to send timeout notification', {
error: err.message
});
});
}
}
})
.catch(err => {
job.status = 'error';
job.error = err.message;
job.completedAt = new Date().toISOString();
this.log.error('dns-propagation', `Verification failed for ${domain}`, {
error: err.message
});
});
return job;
}
/**
* Get the current verification status for a domain.
*
* @param {string} domain - The domain to look up
* @returns {Object|null} Verification status or null if not found
*/
getVerificationStatus(domain) {
const job = this.verifications.get(domain);
if (!job) return null;
return {
domain: job.domain,
expectedIp: job.expectedIp,
status: job.status,
startedAt: job.startedAt,
completedAt: job.completedAt || null,
result: job.result || null,
error: job.error || null
};
}
/**
* Get all recent verifications.
*
* @returns {Object[]} Array of verification statuses
*/
getAllVerifications() {
const results = [];
for (const [domain, job] of this.verifications.entries()) {
results.push({
domain,
expectedIp: job.expectedIp,
status: job.status,
startedAt: job.startedAt,
completedAt: job.completedAt || null,
propagated: job.result?.propagated || null,
totalTime: job.result?.totalTime || null,
error: job.error || null
});
}
return results;
}
/**
* Remove verifications older than 1 hour.
*/
cleanup() {
const now = Date.now();
for (const [domain, job] of this.verifications.entries()) {
const completedAt = job.completedAt ? new Date(job.completedAt).getTime() : null;
const startedAt = new Date(job.startedAt).getTime();
// Clean up completed/error jobs older than 1 hour
// Also clean up stale running jobs that started over 2 hours ago
const age = completedAt ? (now - completedAt) : (now - startedAt);
const maxAge = job.status === 'running' ? MAX_RESULT_AGE_MS * 2 : MAX_RESULT_AGE_MS;
if (age > maxAge) {
this.verifications.delete(domain);
}
}
}
}
module.exports = DNSPropagationChecker;
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "dashcaddy-api",
"version": "1.8.0",
"version": "1.9.0",
"description": "DashCaddy API server - Dashboard backend for Docker, Caddy & DNS management",
"main": "server.js",
"scripts": {
+164
View File
@@ -0,0 +1,164 @@
/**
* Auto-Restart Policy Routes
*
* CRUD endpoints for per-container auto-restart policies.
* Also provides a dry-run test endpoint.
*
* @module routes/auto-restart
*/
const express = require('express');
const { success } = require('../response-helpers');
const { ValidationError, NotFoundError } = require('../errors');
/**
* Auto-restart route factory
*
* @param {Object} deps - Explicit dependencies
* @param {Object} deps.autoRestartManager - AutoRestartManager instance
* @param {Function} deps.asyncHandler - Async route handler wrapper
* @param {Function} deps.logError - Error logging function
* @returns {express.Router}
*/
module.exports = function ({ autoRestartManager, asyncHandler, logError }) {
const router = express.Router();
/**
* GET /auto-restart/policies
* List all configured auto-restart policies.
*/
router.get('/policies', asyncHandler(async (_req, res) => {
const policies = autoRestartManager.listPolicies();
success(res, { policies });
}, 'auto-restart-list'));
/**
* GET /auto-restart/policies/:serviceId
* Get the restart policy for a single service.
*/
router.get('/policies/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
if (!serviceId || !/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,100}$/.test(serviceId)) {
throw new ValidationError('Invalid service ID format');
}
const policy = autoRestartManager.getPolicy(serviceId);
if (!policy) {
throw new NotFoundError(`Auto-restart policy for "${serviceId}"`);
}
success(res, { policy });
}, 'auto-restart-get'));
/**
* POST /auto-restart/policies/:serviceId
* Create or update a restart policy.
*
* Body: { enabled, maxRetries, retryIntervalMs, windowMinutes }
*/
router.post('/policies/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
if (!serviceId || !/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,100}$/.test(serviceId)) {
throw new ValidationError('Invalid service ID format');
}
const { enabled, maxRetries, retryIntervalMs, windowMinutes } = req.body;
// Validate inputs
if (enabled !== undefined && typeof enabled !== 'boolean') {
throw new ValidationError('enabled must be a boolean');
}
if (maxRetries !== undefined) {
if (!Number.isInteger(maxRetries) || maxRetries < 0 || maxRetries > 100) {
throw new ValidationError('maxRetries must be an integer between 0 and 100');
}
}
if (retryIntervalMs !== undefined) {
if (!Number.isInteger(retryIntervalMs) || retryIntervalMs < 0 || retryIntervalMs > 3600000) {
throw new ValidationError('retryIntervalMs must be an integer between 0 and 3600000');
}
}
if (windowMinutes !== undefined) {
if (!Number.isInteger(windowMinutes) || windowMinutes < 0 || windowMinutes > 1440) {
throw new ValidationError('windowMinutes must be an integer between 0 and 1440');
}
}
const policy = await autoRestartManager.setPolicy(serviceId, {
...(enabled !== undefined && { enabled }),
...(maxRetries !== undefined && { maxRetries }),
...(retryIntervalMs !== undefined && { retryIntervalMs }),
...(windowMinutes !== undefined && { windowMinutes }),
});
success(res, { policy, message: `Policy ${serviceId} saved` });
}, 'auto-restart-set'));
/**
* DELETE /auto-restart/policies/:serviceId
* Remove a restart policy.
*/
router.delete('/policies/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
if (!serviceId || !/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,100}$/.test(serviceId)) {
throw new ValidationError('Invalid service ID format');
}
const removed = await autoRestartManager.removePolicy(serviceId);
if (!removed) {
throw new NotFoundError(`Auto-restart policy for "${serviceId}"`);
}
success(res, { message: `Policy for "${serviceId}" removed` });
}, 'auto-restart-delete'));
/**
* POST /auto-restart/policies/:serviceId/test
* Dry-run: simulate a restart attempt without actually restarting.
* Returns what *would* happen given the current policy state.
*/
router.post('/policies/:serviceId/test', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
if (!serviceId || !/^[a-zA-Z0-9][a-zA-Z0-9_.-]{0,100}$/.test(serviceId)) {
throw new ValidationError('Invalid service ID format');
}
const policy = autoRestartManager.getPolicy(serviceId);
if (!policy) {
throw new NotFoundError(`Auto-restart policy for "${serviceId}"`);
}
const now = Date.now();
const inCooldown = policy.cooldownUntil && now < policy.cooldownUntil;
const wouldRetry = !inCooldown && policy.currentRetries < policy.maxRetries;
const nextAttempt = policy.currentRetries + 1;
success(res, {
dryRun: true,
serviceId,
policy: {
enabled: policy.enabled,
currentRetries: policy.currentRetries,
maxRetries: policy.maxRetries,
cooldownUntil: policy.cooldownUntil,
inCooldown,
},
wouldRestart: policy.enabled && wouldRetry,
wouldMaxOut: !wouldRetry && !inCooldown,
nextAttempt: wouldRetry ? nextAttempt : null,
message: !policy.enabled
? 'Policy is disabled — no restart would occur'
: inCooldown
? `In cooldown until ${new Date(policy.cooldownUntil).toISOString()} — would skip`
: wouldRetry
? `Would attempt restart ${nextAttempt}/${policy.maxRetries}`
: `Max retries (${policy.maxRetries}) already reached — would enter cooldown`,
});
}, 'auto-restart-test'));
return router;
};
+92
View File
@@ -0,0 +1,92 @@
/**
* Config Drift Detection Routes
*
* API endpoints for running drift detection, reading cached reports,
* auto-fixing drift, and controlling periodic polling.
*
* @module routes/config-drift
*/
const express = require('express');
const { success } = require('../response-helpers');
const { ValidationError, NotFoundError } = require('../errors');
/**
* Config-drift route factory
*
* @param {Object} deps - Explicit dependencies
* @param {Object} deps.driftDetector - ConfigDriftDetector instance
* @param {Function} deps.asyncHandler - Async route handler wrapper
* @param {Function} deps.logError - Error logging function
* @returns {express.Router}
*/
module.exports = function ({ driftDetector, asyncHandler, logError }) {
const router = express.Router();
/**
* GET /config-drift/report
* Run a fresh drift detection and return the full report.
*/
router.get('/report', asyncHandler(async (_req, res) => {
const report = await driftDetector.detect();
success(res, { report });
}, 'drift-report'));
/**
* GET /config-drift/last
* Return the last cached drift report (no re-detection).
*/
router.get('/last', asyncHandler(async (_req, res) => {
if (!driftDetector.lastReport) {
throw new NotFoundError('No cached drift report — run detection first');
}
success(res, { report: driftDetector.lastReport });
}, 'drift-last'));
/**
* POST /config-drift/fix
* Auto-fix detected drift: remove stale records, flag unknown containers.
*/
router.post('/fix', asyncHandler(async (_req, res) => {
const result = await driftDetector.autoFix();
success(res, {
message: 'Auto-fix applied',
staleRemoved: result.staleRemoved,
unknownFlagged: result.unknownFlagged,
});
}, 'drift-fix'));
/**
* POST /config-drift/polling
* Enable or disable periodic drift detection polling.
*
* Body: { enabled: boolean, intervalMs?: number }
*/
router.post('/polling', asyncHandler(async (req, res) => {
const { enabled, intervalMs } = req.body;
if (typeof enabled !== 'boolean') {
throw new ValidationError('enabled must be a boolean');
}
if (intervalMs !== undefined) {
if (!Number.isInteger(intervalMs) || intervalMs < 10000 || intervalMs > 86400000) {
throw new ValidationError('intervalMs must be an integer between 10000 and 86400000 (10s 24h)');
}
}
if (enabled) {
driftDetector.startPolling(intervalMs || 300000);
success(res, {
message: 'Drift polling enabled',
intervalMs: intervalMs || 300000,
});
} else {
driftDetector.stopPolling();
success(res, { message: 'Drift polling disabled' });
}
}, 'drift-polling'));
return router;
};
+235
View File
@@ -0,0 +1,235 @@
/**
* Dependencies Route — REST API for service dependency tracking
*
* Endpoints:
* GET /dependencies/graph Full dependency graph
* GET /dependencies/validate Validate a proposed dep chain
* GET /dependencies/:serviceId Direct deps for one service
* GET /dependencies/:serviceId/chain Ordered restart chain
* GET /dependencies/:serviceId/status Dependency health status
* POST /dependencies/:serviceId Set dependencies
* DELETE /dependencies/:serviceId Remove all dependencies
* POST /dependencies/:serviceId/restart Restart with dependency chain
*
* @module routes/dependencies
*/
const express = require('express');
const { success, error: errorResponse } = require('../response-helpers');
const { NotFoundError, ValidationError } = require('../errors');
/**
* Dependencies route factory
*
* @param {Object} deps - Explicit dependencies
* @param {Object} deps.dependencyManager - DependencyManager instance
* @param {Object} deps.servicesStateManager - State manager for services.json
* @param {Object} deps.docker - Docker client wrapper
* @param {Function} deps.asyncHandler - Async route handler wrapper
* @param {Function} deps.logError - Error logging function
* @param {Function} deps.resyncHealthChecker - Health checker resync function
* @param {Object} deps.log - Logger instance
* @returns {express.Router}
*/
module.exports = function({
dependencyManager,
servicesStateManager,
docker,
asyncHandler,
logError,
resyncHealthChecker,
log,
}) {
const router = express.Router();
// -------------------------------------------------------------------------
// GET /dependencies/graph — Full dependency graph
// -------------------------------------------------------------------------
router.get('/graph', asyncHandler(async (req, res) => {
const graph = await dependencyManager.getDependencyGraph();
success(res, { graph });
}, 'dep-graph'));
// -------------------------------------------------------------------------
// GET /dependencies/validate — Validate a proposed dep chain (query params)
// -------------------------------------------------------------------------
router.get('/validate', asyncHandler(async (req, res) => {
const { serviceId, dependsOn } = req.query;
if (!serviceId) {
throw new ValidationError('serviceId query parameter is required');
}
// dependsOn may be a comma-separated string or already an array
let parsed;
if (Array.isArray(dependsOn)) {
parsed = dependsOn;
} else if (typeof dependsOn === 'string' && dependsOn.length > 0) {
parsed = dependsOn.split(',').map(s => s.trim()).filter(Boolean);
} else {
parsed = [];
}
const result = await dependencyManager.validateDependencies(serviceId, parsed);
success(res, result);
}, 'dep-validate'));
// -------------------------------------------------------------------------
// GET /dependencies/:serviceId — Direct deps for one service
// -------------------------------------------------------------------------
router.get('/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
const dependencies = await dependencyManager.getDependencies(serviceId);
const dependents = await dependencyManager.getDependents(serviceId);
// Read the service's current dependsOn array
const services = await servicesStateManager.read();
const allServices = Array.isArray(services) ? services : (services.services || []);
const service = allServices.find(s => s.id === serviceId);
if (!service) {
throw new NotFoundError(`Service "${serviceId}"`);
}
success(res, {
serviceId,
dependsOn: service.dependsOn || [],
dependencies,
dependents: dependents.map(d => ({ id: d.id, name: d.name })),
});
}, 'dep-get'));
// -------------------------------------------------------------------------
// GET /dependencies/:serviceId/chain — Ordered restart chain
// -------------------------------------------------------------------------
router.get('/:serviceId/chain', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
const chain = await dependencyManager.getOrderedRestartChain(serviceId);
success(res, { serviceId, chain });
}, 'dep-chain'));
// -------------------------------------------------------------------------
// GET /dependencies/:serviceId/status — Dependency health status
// -------------------------------------------------------------------------
router.get('/:serviceId/status', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
const statuses = await dependencyManager.getDependencyStatus(serviceId);
success(res, { serviceId, statuses });
}, 'dep-status'));
// -------------------------------------------------------------------------
// POST /dependencies/:serviceId — Set dependencies
// -------------------------------------------------------------------------
router.post('/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
const { dependsOn } = req.body;
if (!Array.isArray(dependsOn)) {
throw new ValidationError('Request body must include dependsOn as an array of service IDs');
}
// Validate first
const validation = await dependencyManager.validateDependencies(serviceId, dependsOn);
if (!validation.valid) {
return errorResponse(res, validation.errors.join('; '), 400);
}
// Update the service
let found = false;
await servicesStateManager.update(services => {
const arr = Array.isArray(services) ? services : [];
return arr.map(s => {
if (s.id === serviceId) {
found = true;
return { ...s, dependsOn: dependsOn.slice() };
}
return s;
});
});
if (!found) {
throw new NotFoundError(`Service "${serviceId}"`);
}
log.info('dependency', 'Dependencies updated', { serviceId, dependsOn });
success(res, {
message: `Dependencies updated for "${serviceId}"`,
serviceId,
dependsOn,
});
}, 'dep-set'));
// -------------------------------------------------------------------------
// DELETE /dependencies/:serviceId — Remove all dependencies for a service
// -------------------------------------------------------------------------
router.delete('/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
let found = false;
await servicesStateManager.update(services => {
const arr = Array.isArray(services) ? services : [];
return arr.map(s => {
if (s.id === serviceId) {
found = true;
const updated = { ...s };
delete updated.dependsOn;
return updated;
}
return s;
});
});
if (!found) {
throw new NotFoundError(`Service "${serviceId}"`);
}
log.info('dependency', 'Dependencies removed', { serviceId });
success(res, {
message: `All dependencies removed for "${serviceId}"`,
serviceId,
});
}, 'dep-delete'));
// -------------------------------------------------------------------------
// POST /dependencies/:serviceId/restart — Restart with dependency chain
// -------------------------------------------------------------------------
router.post('/:serviceId/restart', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
// Verify the service exists
const services = await servicesStateManager.read();
const allServices = Array.isArray(services) ? services : (services.services || []);
if (!allServices.find(s => s.id === serviceId)) {
throw new NotFoundError(`Service "${serviceId}"`);
}
// Get the chain first for the response (before async restart begins)
let chain;
try {
chain = await dependencyManager.getOrderedRestartChain(serviceId);
} catch (err) {
return errorResponse(res, err.message, 400);
}
// Respond immediately with the chain order
success(res, {
message: `Dependency restart initiated for "${serviceId}"`,
serviceId,
chain,
});
// Run the restart chain asynchronously so the client doesn't block
dependencyManager.restartWithDependencies(serviceId).catch(err => {
if (log) {
log.error('dependency', 'Async dependency restart failed', {
serviceId,
error: err.message,
});
}
});
}, 'dep-restart'));
return router;
};
+73 -1
View File
@@ -26,7 +26,8 @@ module.exports = function({
log,
safeErrorMessage,
fetchT,
credentialManager
credentialManager,
dnsPropagationChecker
}) {
const router = express.Router();
@@ -139,6 +140,14 @@ module.exports = function({
});
if (result.status === 'ok') {
// Start DNS propagation verification in background
if (dnsPropagationChecker && ip) {
const fullDomain = domain;
dnsPropagationChecker.startVerification(fullDomain, ip).catch(err => {
log('DNS propagation check start failed:', err.message);
});
}
success(res, { message: `DNS record ${domain} -> ${ip} created` });
} else {
// Error handled by middleware
@@ -641,5 +650,68 @@ module.exports = function({
}
}, 'dns-update'));
// ===== DNS PROPAGATION =====
// GET /propagation — Get all recent DNS propagation checks
router.get('/propagation', asyncHandler(async (req, res) => {
if (!dnsPropagationChecker) {
return success(res, { verifications: [], message: 'DNS propagation checker not available' });
}
// Cleanup old entries
dnsPropagationChecker.cleanup();
const verifications = dnsPropagationChecker.getAllVerifications();
success(res, { verifications });
}, 'dns-propagation-all'));
// POST /propagation/verify — Manually trigger DNS propagation verification
router.post('/propagation/verify', asyncHandler(async (req, res) => {
if (!dnsPropagationChecker) {
return errorResponse(res, 'DNS propagation checker not available', 503);
}
const { domain, expectedIp } = req.body;
if (!domain || !expectedIp) {
throw new ValidationError('domain and expectedIp are required');
}
// Validate domain format
if (!REGEX.DOMAIN.test(domain)) {
throw new ValidationError('[DC-301] Invalid domain format');
}
// Validate IP address
const validatorLib = require('validator');
if (!validatorLib.isIP(expectedIp)) {
throw new ValidationError('[DC-210] Invalid IP address');
}
const job = dnsPropagationChecker.startVerification(domain, expectedIp);
success(res, {
message: 'DNS propagation verification started',
domain,
expectedIp,
status: job.status
});
}, 'dns-propagation-verify'));
// GET /propagation/:domain — Get propagation status for a specific domain
router.get('/propagation/:domain', asyncHandler(async (req, res) => {
if (!dnsPropagationChecker) {
return success(res, { verification: null, message: 'DNS propagation checker not available' });
}
const { domain } = req.params;
const status = dnsPropagationChecker.getVerificationStatus(domain);
if (!status) {
throw new NotFoundError(`No propagation check found for domain: ${domain}`);
}
success(res, { verification: status });
}, 'dns-propagation-domain'));
return router;
};
+44 -1
View File
@@ -8,9 +8,10 @@ const express = require('express');
* @param {Object} deps.healthChecker - Health checker
* @param {Object} deps.updateManager - Update manager
* @param {Function} deps.logError - Error logging function
* @param {Object} deps.dependencyManager - Dependency manager for restart chain events
* @returns {express.Router}
*/
module.exports = function({ resourceMonitor, healthChecker, updateManager, logError }) {
module.exports = function({ resourceMonitor, healthChecker, updateManager, logError, dependencyManager, autoRestartManager, driftDetector, sslMonitor, dnsPropagationChecker }) {
const router = express.Router();
const clients = new Set();
@@ -74,6 +75,48 @@ module.exports = function({ resourceMonitor, healthChecker, updateManager, logEr
});
}
// Dependency manager events
if (dependencyManager) {
dependencyManager.on('dependency-restart-start', (data) => {
broadcast('dependency-restart-start', data);
});
dependencyManager.on('dependency-restart-progress', (data) => {
broadcast('dependency-restart-progress', data);
});
dependencyManager.on('dependency-restart-complete', (data) => {
broadcast('dependency-restart-complete', data);
});
dependencyManager.on('dependency-restart-failed', (data) => {
broadcast('dependency-restart-failed', data);
});
}
// Auto-restart manager events
if (autoRestartManager) {
autoRestartManager.on('auto-restart-attempt', (data) => broadcast('auto-restart-attempt', data));
autoRestartManager.on('auto-restart-success', (data) => broadcast('auto-restart-success', data));
autoRestartManager.on('auto-restart-failed', (data) => broadcast('auto-restart-failed', data));
autoRestartManager.on('auto-restart-max-reached', (data) => broadcast('auto-restart-max-reached', data));
}
// Config drift detector events
if (driftDetector) {
driftDetector.on('drift-detected', (data) => broadcast('drift-detected', data));
}
// SSL monitor events
if (sslMonitor) {
sslMonitor.on('cert-expiring', (data) => broadcast('cert-expiring', data));
sslMonitor.on('cert-critical', (data) => broadcast('cert-critical', data));
}
// DNS propagation checker events
if (dnsPropagationChecker) {
dnsPropagationChecker.on('propagation-check', (data) => broadcast('dns-propagation-check', data));
dnsPropagationChecker.on('propagation-complete', (data) => broadcast('dns-propagation-complete', data));
dnsPropagationChecker.on('propagation-timeout', (data) => broadcast('dns-propagation-timeout', data));
}
// SSE endpoint
router.get('/stream', (req, res) => {
res.writeHead(200, {
+113
View File
@@ -0,0 +1,113 @@
/**
* SSL Monitor Routes
* REST API endpoints for SSL certificate monitoring.
*
* @module routes/ssl-monitor
*/
const express = require('express');
const { success, error: errorResponse, notFound } = require('../response-helpers');
/**
* SSL Monitor route factory
* @param {Object} deps - Explicit dependencies
* @param {Object} deps.sslMonitor - SSLMonitor instance
* @param {Function} deps.asyncHandler - Async route handler wrapper
* @param {Function} deps.logError - Error logging function
* @returns {express.Router}
*/
module.exports = function({ sslMonitor, asyncHandler, logError }) {
const router = express.Router();
/**
* GET /ssl/certificates
* Get all SSL certificate statuses
*/
router.get('/certificates', asyncHandler(async (req, res) => {
const status = sslMonitor.getStatus();
success(res, { certificates: status });
}, 'ssl-certificates'));
/**
* GET /ssl/certificates/:serviceId
* Get SSL certificate status for a specific service
*/
router.get('/certificates/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
const certStatus = sslMonitor.getServiceCertStatus(serviceId);
if (!certStatus) {
return notFound(res, `No SSL certificate status found for service: ${serviceId}`);
}
success(res, { certificate: certStatus });
}, 'ssl-certificate-service'));
/**
* POST /ssl/check
* Trigger an on-demand check of all SSL certificates
*/
router.post('/check', asyncHandler(async (req, res) => {
const results = await sslMonitor.checkAll();
success(res, { certificates: results, message: 'SSL check completed' });
}, 'ssl-check-all'));
/**
* POST /ssl/check/:serviceId
* Check the SSL certificate for a specific service
*/
router.post('/check/:serviceId', asyncHandler(async (req, res) => {
const { serviceId } = req.params;
// Look up the existing cert status to find the hostname
const existingCert = sslMonitor.getServiceCertStatus(serviceId);
if (!existingCert) {
return notFound(res, `No HTTPS URL found for service: ${serviceId}`);
}
try {
const result = await sslMonitor.checkCert(existingCert.hostname, existingCert.port);
success(res, { certificate: { ...result, serviceId } });
} catch (err) {
errorResponse(res, `Failed to check SSL certificate: ${err.message}`, 500);
}
}, 'ssl-check-service'));
/**
* GET /ssl/config
* Get current SSL monitoring configuration
*/
router.get('/config', asyncHandler(async (req, res) => {
const config = sslMonitor.getConfig();
success(res, { config });
}, 'ssl-config-get'));
/**
* POST /ssl/config
* Update SSL monitoring configuration
* Body: { enabled: boolean, intervalMs: number }
*/
router.post('/config', asyncHandler(async (req, res) => {
const { enabled, intervalMs } = req.body;
// Validate inputs
if (enabled !== undefined && typeof enabled !== 'boolean') {
return errorResponse(res, 'enabled must be a boolean', 400);
}
if (intervalMs !== undefined) {
if (typeof intervalMs !== 'number' || intervalMs < 60000) {
return errorResponse(res, 'intervalMs must be a number >= 60000 (1 minute)', 400);
}
}
const updates = {};
if (enabled !== undefined) updates.enabled = enabled;
if (intervalMs !== undefined) updates.intervalMs = intervalMs;
sslMonitor.updateConfig(updates);
const config = sslMonitor.getConfig();
success(res, { config, message: 'SSL monitoring config updated' });
}, 'ssl-config-update'));
return router;
};
+74 -2
View File
@@ -77,6 +77,15 @@ const themesRoutes = require('../routes/themes');
const dockerResourcesRoutes = require('../routes/docker-resources');
const eventsRoutes = require('../routes/events');
const workflowsRoutes = require('../routes/workflows');
const dependenciesRoutes = require('../routes/dependencies');
const DependencyManager = require('../dependency-manager');
const autoRestartRoutes = require('../routes/auto-restart');
const configDriftRoutes = require('../routes/config-drift');
const sslMonitorRoutes = require('../routes/ssl-monitor');
const { AutoRestartManager } = require('../auto-restart-manager');
const { ConfigDriftDetector } = require('../config-drift-detector');
const { SSLMonitor } = require('../ssl-monitor');
const { DNSPropagationChecker } = require('../dns-propagation');
// Constants
const { APP } = require('../constants');
@@ -335,6 +344,39 @@ async function createApp() {
}
}
// Initialize dependency manager
const dependencyManager = new DependencyManager({
servicesStateManager,
docker: ctx.docker,
notification: ctx.notification,
log,
});
ctx.dependencyManager = dependencyManager;
log.info('app', 'Dependency manager initialized');
// Initialize auto-restart manager
const autoRestartManager = new AutoRestartManager(ctx);
ctx.autoRestartManager = autoRestartManager;
autoRestartManager.start();
log.info('app', 'Auto-restart manager initialized');
// Initialize config drift detector
const driftDetector = new ConfigDriftDetector(ctx);
ctx.driftDetector = driftDetector;
driftDetector.startPolling(300000); // 5 min
log.info('app', 'Config drift detector initialized');
// Initialize SSL monitor
const sslMonitor = new SSLMonitor(ctx);
ctx.sslMonitor = sslMonitor;
sslMonitor.start(3600000); // 1 hour
log.info('app', 'SSL monitor initialized');
// Initialize DNS propagation checker
const dnsPropagationChecker = new DNSPropagationChecker(ctx);
ctx.dnsPropagationChecker = dnsPropagationChecker;
log.info('app', 'DNS propagation checker initialized');
// Build versioned API router
const apiRouter = express.Router();
@@ -375,7 +417,8 @@ async function createApp() {
log: ctx.log,
safeErrorMessage: ctx.safeErrorMessage,
fetchT: ctx.fetchT,
credentialManager: ctx.credentialManager
credentialManager: ctx.credentialManager,
dnsPropagationChecker: ctx.dnsPropagationChecker
}));
apiRouter.use('/notifications', notificationRoutes({
notification: ctx.notification,
@@ -489,13 +532,42 @@ async function createApp() {
resourceMonitor: ctx.resourceMonitor,
healthChecker: ctx.healthChecker,
updateManager: ctx.updateManager,
logError: ctx.logError
logError: ctx.logError,
dependencyManager: ctx.dependencyManager,
autoRestartManager: ctx.autoRestartManager,
driftDetector: ctx.driftDetector,
sslMonitor: ctx.sslMonitor,
dnsPropagationChecker: ctx.dnsPropagationChecker
}));
apiRouter.use(workflowsRoutes({
workflowEngine: ctx.workflowEngine,
licenseManager: ctx.licenseManager,
asyncHandler: ctx.asyncHandler
}));
apiRouter.use('/dependencies', dependenciesRoutes({
dependencyManager: ctx.dependencyManager,
servicesStateManager: ctx.servicesStateManager,
docker: ctx.docker,
asyncHandler: ctx.asyncHandler,
logError: ctx.logError,
resyncHealthChecker: ctx.resyncHealthChecker,
log: ctx.log,
}));
apiRouter.use(autoRestartRoutes({
autoRestartManager: ctx.autoRestartManager,
asyncHandler: ctx.asyncHandler,
logError: ctx.logError,
}));
apiRouter.use(configDriftRoutes({
driftDetector: ctx.driftDetector,
asyncHandler: ctx.asyncHandler,
logError: ctx.logError,
}));
apiRouter.use(sslMonitorRoutes({
sslMonitor: ctx.sslMonitor,
asyncHandler: ctx.asyncHandler,
logError: ctx.logError,
}));
// Inline API routes
apiRouter.get('/health', (req, res) => {
+411
View File
@@ -0,0 +1,411 @@
/**
* SSL Certificate Monitor
* Periodically checks SSL certificates on services with HTTPS URLs.
* Alerts at 30, 14, and 7 days before expiry.
*
* @module ssl-monitor
*/
const tls = require('tls');
const EventEmitter = require('events');
const path = require('path');
const { readJsonFile, writeJsonFile } = require('./fs-helpers');
const { resolveServiceUrl } = require('./url-resolver');
/** Default check interval: 1 hour */
const DEFAULT_INTERVAL_MS = 3600000;
/** Alert thresholds in days */
const THRESHOLDS = {
WARNING: 30,
URGENT: 14,
CRITICAL: 7
};
/** TLS connection timeout in milliseconds */
const TLS_TIMEOUT_MS = 10000;
class SSLMonitor extends EventEmitter {
/**
* Create an SSLMonitor instance.
* @param {Object} ctx - Shared application context
* @param {Object} ctx.servicesStateManager - State manager for reading services
* @param {Function} ctx.buildServiceUrl - URL builder helper
* @param {Object} ctx.siteConfig - Site configuration
* @param {Object} ctx.notification - NotificationManager instance
* @param {Object} ctx.log - Logger instance
* @param {string} [ctx.SSL_CACHE_FILE] - Path to persist SSL cache
*/
constructor(ctx) {
super();
this.ctx = ctx;
this.log = ctx.log || console;
/** @type {Map<string, Object>} hostname → last cert check result */
this.certStatus = new Map();
/** @type {Map<string, number>} hostname → last notified threshold level */
this.notifiedThresholds = new Map();
/** @type {Map<string, string>} hostname → service ID mapping */
this.hostnameToServiceId = new Map();
/** @type {NodeJS.Timeout|null} */
this.intervalHandle = null;
/** Current config */
this.config = {
enabled: true,
intervalMs: DEFAULT_INTERVAL_MS
};
/** Cache file path */
this.cacheFile = ctx.SSL_CACHE_FILE ||
path.join(path.dirname(ctx.SERVICES_FILE || './data'), 'ssl-cache.json');
}
/**
* Check the SSL certificate for a given hostname and port.
* Connects via TLS with rejectUnauthorized: false to retrieve certificate info.
*
* @param {string} hostname - The hostname to check
* @param {number} [port=443] - The port to connect to
* @returns {Promise<Object>} Certificate information
*/
async checkCert(hostname, port = 443) {
return new Promise((resolve, reject) => {
const socket = tls.connect({
host: hostname,
port,
rejectUnauthorized: false,
servername: hostname,
timeout: TLS_TIMEOUT_MS
}, () => {
try {
const cert = socket.getPeerCertificate();
if (!cert || Object.keys(cert).length === 0) {
socket.destroy();
return reject(new Error(`No certificate returned for ${hostname}:${port}`));
}
const validFrom = new Date(cert.valid_from);
const validTo = new Date(cert.valid_to);
const now = new Date();
const msRemaining = validTo.getTime() - now.getTime();
const daysRemaining = Math.ceil(msRemaining / (1000 * 60 * 60 * 24));
const result = {
hostname,
port,
subject: cert.subject?.CN || cert.subject?.O || 'Unknown',
issuer: cert.issuer?.CN || cert.issuer?.O || 'Unknown',
validFrom: cert.valid_from,
validTo: cert.valid_to,
daysRemaining,
fingerprint: cert.fingerprint || null,
isExpiring: daysRemaining <= THRESHOLDS.WARNING,
checkedAt: new Date().toISOString()
};
socket.destroy();
resolve(result);
} catch (err) {
socket.destroy();
reject(err);
}
});
socket.on('error', (err) => {
reject(new Error(`TLS connect error for ${hostname}:${port}: ${err.message}`));
});
socket.setTimeout(TLS_TIMEOUT_MS, () => {
socket.destroy(new Error(`TLS connection timeout for ${hostname}:${port}`));
reject(new Error(`TLS connection timeout for ${hostname}:${port}`));
});
});
}
/**
* Check SSL certificates for all services that have HTTPS URLs.
* Reads services from ctx.servicesStateManager, resolves URLs, and checks each HTTPS cert.
*
* @returns {Promise<Object>} Map of hostname → cert status
*/
async checkAll() {
if (!this.config.enabled) {
this.log.info('ssl-monitor', 'SSL monitoring is disabled, skipping check');
return this.getStatus();
}
let servicesData;
try {
servicesData = await this.ctx.servicesStateManager.read();
} catch (err) {
this.log.error('ssl-monitor', 'Failed to read services', { error: err.message });
return this.getStatus();
}
const services = Array.isArray(servicesData) ? servicesData : (servicesData.services || []);
for (const service of services) {
const serviceId = service.id || service.name?.toLowerCase();
if (!serviceId) continue;
try {
const url = resolveServiceUrl(serviceId, service, this.ctx.siteConfig, this.ctx.buildServiceUrl);
if (!url) continue;
const parsed = new URL(url);
if (parsed.protocol !== 'https:') continue;
const hostname = parsed.hostname;
const port = parseInt(parsed.port) || 443;
// Map hostname back to service ID
this.hostnameToServiceId.set(hostname, serviceId);
const result = await this.checkCert(hostname, port);
// Store result
this.certStatus.set(hostname, result);
// Emit check event
this.emit('cert-check', { serviceId, hostname, result });
// Check alert thresholds
await this._checkAndNotify(hostname, result, serviceId);
} catch (err) {
this.log.warn('ssl-monitor', `Failed to check cert for service ${serviceId}`, {
error: err.message
});
}
}
// Persist results
await this._saveCache();
return this.getStatus();
}
/**
* Start periodic SSL certificate checking.
*
* @param {number} [intervalMs=3600000] - Check interval in milliseconds
*/
start(intervalMs) {
if (intervalMs !== undefined) {
this.config.intervalMs = intervalMs;
}
if (this.intervalHandle) {
this.log.warn('ssl-monitor', 'SSL monitor is already running');
return;
}
this.config.enabled = true;
// Load cached data
this._loadCache().catch(err => {
this.log.warn('ssl-monitor', 'Failed to load SSL cache', { error: err.message });
});
// Initial check (non-blocking)
this.checkAll().catch(err => {
this.log.error('ssl-monitor', 'Initial SSL check failed', { error: err.message });
});
// Schedule periodic checks
this.intervalHandle = setInterval(() => {
this.checkAll().catch(err => {
this.log.error('ssl-monitor', 'Periodic SSL check failed', { error: err.message });
});
}, this.config.intervalMs);
this.log.info('ssl-monitor', 'SSL monitoring started', {
intervalMs: this.config.intervalMs
});
}
/**
* Stop periodic SSL certificate checking.
*/
stop() {
if (this.intervalHandle) {
clearInterval(this.intervalHandle);
this.intervalHandle = null;
}
this.config.enabled = false;
this.log.info('ssl-monitor', 'SSL monitoring stopped');
}
/**
* Get the current SSL certificate status for all checked hostnames.
*
* @returns {Object} Map of hostname → cert status
*/
getStatus() {
const status = {};
for (const [hostname, cert] of this.certStatus.entries()) {
status[hostname] = { ...cert };
}
return status;
}
/**
* Get the SSL certificate status for a specific service.
*
* @param {string} serviceId - The service ID to look up
* @returns {Object|null} Certificate status or null if not found
*/
getServiceCertStatus(serviceId) {
// Find hostname mapped to this service
for (const [hostname, id] of this.hostnameToServiceId.entries()) {
if (id === serviceId) {
const cert = this.certStatus.get(hostname);
return cert ? { ...cert, serviceId } : null;
}
}
return null;
}
/**
* Get current monitoring configuration.
*
* @returns {Object} Config with interval and enabled state
*/
getConfig() {
return { ...this.config };
}
/**
* Update monitoring configuration.
*
* @param {Object} updates - Config updates
* @param {boolean} [updates.enabled] - Enable/disable monitoring
* @param {number} [updates.intervalMs] - Check interval in milliseconds
*/
updateConfig(updates) {
if (typeof updates.enabled === 'boolean') {
this.config.enabled = updates.enabled;
if (!updates.enabled && this.intervalHandle) {
this.stop();
}
}
if (typeof updates.intervalMs === 'number' && updates.intervalMs >= 60000) {
this.config.intervalMs = updates.intervalMs;
// Restart interval if running
if (this.intervalHandle) {
clearInterval(this.intervalHandle);
this.intervalHandle = setInterval(() => {
this.checkAll().catch(err => {
this.log.error('ssl-monitor', 'Periodic SSL check failed', { error: err.message });
});
}, this.config.intervalMs);
}
}
}
// ===== Private Methods =====
/**
* Check alert thresholds and send notifications if thresholds are crossed.
* Only sends one notification per threshold per hostname.
*
* @param {string} hostname
* @param {Object} certResult
* @param {string} serviceId
*/
async _checkAndNotify(hostname, certResult, serviceId) {
const { daysRemaining } = certResult;
const key = hostname;
const lastNotified = this.notifiedThresholds.get(key) || Infinity;
let level = null;
let eventType = null;
let message = null;
if (daysRemaining <= THRESHOLDS.CRITICAL) {
level = THRESHOLDS.CRITICAL;
eventType = 'cert-critical';
message = `🔒 CRITICAL: SSL certificate for ${hostname} expires in ${daysRemaining} days!`;
} else if (daysRemaining <= THRESHOLDS.URGENT) {
level = THRESHOLDS.URGENT;
eventType = 'cert-expiring';
message = `⚠️ URGENT: SSL certificate for ${hostname} expires in ${daysRemaining} days`;
} else if (daysRemaining <= THRESHOLDS.WARNING) {
level = THRESHOLDS.WARNING;
eventType = 'cert-expiring';
message = `⚠️ SSL certificate for ${hostname} expires in ${daysRemaining} days`;
}
if (level !== null && level < lastNotified) {
// New threshold crossed — send notification
this.notifiedThresholds.set(key, level);
this.emit(eventType, { hostname, serviceId, daysRemaining, level });
if (this.ctx.notification) {
try {
await this.ctx.notification.send('ssl-cert-expiry', {
text: message,
hostname,
serviceId,
daysRemaining,
level,
validTo: certResult.validTo
}, level <= THRESHOLDS.CRITICAL ? 'error' : 'warning');
} catch (err) {
this.log.error('ssl-monitor', 'Failed to send SSL notification', { error: err.message });
}
}
} else if (level === null) {
// Cert is healthy — reset notification tracking
this.notifiedThresholds.delete(key);
}
}
/**
* Persist cert status cache to disk.
*/
async _saveCache() {
try {
const data = {
lastChecked: new Date().toISOString(),
certs: {},
hostnameToServiceId: Object.fromEntries(this.hostnameToServiceId)
};
for (const [hostname, cert] of this.certStatus.entries()) {
data.certs[hostname] = cert;
}
await writeJsonFile(this.cacheFile, data);
} catch (err) {
this.log.warn('ssl-monitor', 'Failed to save SSL cache', { error: err.message });
}
}
/**
* Load cert status cache from disk.
*/
async _loadCache() {
try {
const data = await readJsonFile(this.cacheFile, null);
if (data && data.certs) {
for (const [hostname, cert] of Object.entries(data.certs)) {
this.certStatus.set(hostname, cert);
}
if (data.hostnameToServiceId) {
for (const [hostname, serviceId] of Object.entries(data.hostnameToServiceId)) {
this.hostnameToServiceId.set(hostname, serviceId);
}
}
this.log.info('ssl-monitor', 'Loaded SSL cache', {
certCount: this.certStatus.size
});
}
} catch (err) {
this.log.warn('ssl-monitor', 'Failed to load SSL cache', { error: err.message });
}
}
}
module.exports = SSLMonitor;