feat: real plugin system with loadable instances + OpenBao secrets (v1.17.0)
Generalize the half-built discovery plugins into a real plugin system: plugin TYPES (the plugins/<category>/<type>.js modules with manifests) and loadable, configurable, multi-copy plugin INSTANCES (PluginInstance ORM model) managed from a dedicated /plugins page and /api/plugins API, with per-instance secrets in OpenBao at secret/plugins/<id>/conf. - plugin_registry.js: getTypes/getModule/splitConfig/mask + required-field helpers - PluginInstance model (Sequelize): id/pluginType/category/name/slug(unique)/ enabled/cron/config(json, non-secret)/lastRun*; registered in models/index.js - plugin_secrets.js: read/write/remove/mergeForRun over @simpleworkjs/bao-conf - scheduler.js: schedules from the DB registry; per-instance stable BullMQ JobScheduler ids (plugin:<id>) for load/unload; legacy migration from conf.discovery.plugins on first boot (idempotent, empty-table-guarded) - api_plugins.js (replaces routes/plugins.js): types/list/get/create/update/ secrets/test/load/unload/run/delete/runs; admin-gated; secrets always masked - /plugins page (plugins.ejs) + nav; Agents & Scheduler tab removed from /directory; /docs/agents aliased to /docs/plugins - proxmox/unifi/nmap gained manifests (configSchema/validate/run alias) - tests/plugins.test.js: registry unit + plugin_secrets (mocked bao-conf) + PluginInstance model round-trip/unique-slug - docs (plugins.md, vault.md, _config.yml, API.md) + 1.16.1 -> 1.17.0 Requires theta-suite >= v1.30.1 for the sso-broker secret/plugins/* grant; fails-soft with a clear error if absent. Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
+181
-59
@@ -1,5 +1,27 @@
|
||||
'use strict';
|
||||
|
||||
// Discovery / plugin scheduler.
|
||||
//
|
||||
// Generalized from the one-shot discovery-plugin loader: plugin *types* live
|
||||
// under nodejs/plugins/<category>/<type>.js (see services/plugin_registry.js),
|
||||
// and configured, loadable/unloadable *instances* live in the PluginInstance
|
||||
// table (models/plugin_instance.js). This module schedules enabled instances
|
||||
// on cron via BullMQ JobSchedulers and runs them in a Worker.
|
||||
//
|
||||
// Each instance owns a stable JobScheduler id (`plugin:<instanceId>`) so load/
|
||||
// unload can add/remove a single schedule without disturbing the others —
|
||||
// `upsertJobScheduler`/`removeJobScheduler` (BullMQ v6) take that id directly.
|
||||
//
|
||||
// Per-instance secrets are merged in from OpenBao (utils/plugin_secrets.js) at
|
||||
// run time; the plugin's run()/discover() receives the combined non-secret
|
||||
// config + secret values as a single `config` object, exactly as the legacy
|
||||
// static-config path did.
|
||||
|
||||
const { Queue, Worker } = require('bullmq');
|
||||
const { DiscoveryReconciler } = require('./discovery_reconciler');
|
||||
const pluginRegistry = require('./plugin_registry');
|
||||
const pluginSecrets = require('../utils/plugin_secrets');
|
||||
const { PluginInstance, STATUS } = require('../models/plugin_instance');
|
||||
const Redis = require('ioredis');
|
||||
|
||||
// Ensure Redis connection works for BullMQ
|
||||
@@ -8,80 +30,180 @@ const connection = new Redis(process.env.REDIS_URL || 'redis://127.0.0.1:6379',
|
||||
|
||||
const discoveryQueue = new Queue('discovery', { connection });
|
||||
|
||||
// Load plugins
|
||||
const fs = require('fs');
|
||||
const path = require('path');
|
||||
const pluginsDir = path.join(__dirname, '../plugins/discovery');
|
||||
|
||||
let plugins = {};
|
||||
|
||||
if (fs.existsSync(pluginsDir)) {
|
||||
fs.readdirSync(pluginsDir).forEach(file => {
|
||||
if (file.endsWith('.js')) {
|
||||
const name = path.basename(file, '.js');
|
||||
plugins[name] = require(path.join(pluginsDir, file));
|
||||
}
|
||||
});
|
||||
}
|
||||
const RUN = 'run_plugin';
|
||||
const GC = 'garbage_collect';
|
||||
function pluginSchedulerId(id) { return `plugin:${id}`; }
|
||||
|
||||
const worker = new Worker('discovery', async job => {
|
||||
if (job.name === 'run_plugin') {
|
||||
const { pluginName, config } = job.data;
|
||||
if (plugins[pluginName]) {
|
||||
console.log(`[Scheduler] Running plugin: ${pluginName}`);
|
||||
try {
|
||||
const payload = await plugins[pluginName].discover(config);
|
||||
await DiscoveryReconciler.reconcile(pluginName, payload);
|
||||
} catch (err) {
|
||||
console.error(`[Scheduler] Plugin ${pluginName} failed:`, err);
|
||||
}
|
||||
}
|
||||
} else if (job.name === 'garbage_collect') {
|
||||
console.log(`[Scheduler] Running garbage collection`);
|
||||
if (job.name === RUN) {
|
||||
await runPluginJob(job.data && job.data.instanceId);
|
||||
} else if (job.name === GC) {
|
||||
console.log('[Scheduler] Running garbage collection');
|
||||
await DiscoveryReconciler.garbageCollect();
|
||||
}
|
||||
}, { connection });
|
||||
|
||||
// Function to start scheduling
|
||||
async function initScheduler(discoveryConfig) {
|
||||
// Clear old repeatable jobs (BullMQ v6 uses JobSchedulers)
|
||||
try {
|
||||
const schedulers = await discoveryQueue.getJobSchedulers();
|
||||
for (const job of schedulers) {
|
||||
await discoveryQueue.removeJobScheduler(job.id);
|
||||
}
|
||||
} catch (e) {
|
||||
console.log('[Scheduler] Could not clear old job schedulers (may not be supported or none exist)');
|
||||
// Run one plugin instance. Loads the row (skip silently if it was deleted or
|
||||
// disabled after the job was enqueued), merges its OpenBao secrets into its
|
||||
// config, calls the plugin's run()/discover(), and — for discovery plugins —
|
||||
// reconciles the result into the resource graph under the instance's slug.
|
||||
// Bookkeeping (lastRunAt/lastStatus/lastError) is stamped on the row so the UI
|
||||
// can show run state without querying BullMQ.
|
||||
async function runPluginJob(instanceId) {
|
||||
if (!instanceId) { console.warn('[Scheduler] run_plugin job with no instanceId'); return; }
|
||||
const instance = await PluginInstance.get(instanceId);
|
||||
if (!instance) { console.warn(`[Scheduler] instance ${instanceId} gone — skipping`); return; }
|
||||
if (!instance.enabled) { console.warn(`[Scheduler] instance ${instance.slug} (${instanceId}) disabled — skipping`); return; }
|
||||
|
||||
let mod;
|
||||
try { mod = pluginRegistry.getModule(instance.pluginType); }
|
||||
catch (err) {
|
||||
console.error(`[Scheduler] instance ${instance.slug}: type ${instance.pluginType} unavailable:`, err.message);
|
||||
await instance.update({ lastRunAt: Date.now(), lastStatus: STATUS.ERROR, lastError: `plugin type unavailable: ${instance.pluginType}` });
|
||||
return;
|
||||
}
|
||||
|
||||
// Schedule Garbage Collection
|
||||
await discoveryQueue.add('garbage_collect', {}, { repeat: { pattern: '0 0 * * *' } }); // Daily
|
||||
const runFn = mod.run || mod.discover;
|
||||
if (typeof runFn !== 'function') {
|
||||
console.error(`[Scheduler] instance ${instance.slug}: type ${instance.pluginType} has no run()/discover()`);
|
||||
await instance.update({ lastRunAt: Date.now(), lastStatus: STATUS.ERROR, lastError: 'plugin type has no run()/discover()' });
|
||||
return;
|
||||
}
|
||||
|
||||
// Load plugin overrides from Redis
|
||||
let overrides = {};
|
||||
console.log(`[Scheduler] Running plugin: ${instance.slug} (${instance.pluginType})`);
|
||||
await instance.update({ lastRunAt: Date.now(), lastStatus: STATUS.RUNNING, lastError: null });
|
||||
try {
|
||||
const data = await connection.hgetall('discovery_plugins');
|
||||
for (const [k, v] of Object.entries(data)) {
|
||||
overrides[k] = JSON.parse(v);
|
||||
const cfg = await pluginSecrets.mergeForRun(instance);
|
||||
const payload = await runFn(cfg);
|
||||
if (instance.category === 'discovery') {
|
||||
await DiscoveryReconciler.reconcile(instance.slug, payload);
|
||||
}
|
||||
await instance.update({ lastStatus: STATUS.OK, lastError: null });
|
||||
} catch (err) {
|
||||
console.error('[Scheduler] Failed to load plugin overrides from Redis', err);
|
||||
console.error(`[Scheduler] Plugin ${instance.slug} failed:`, err.message);
|
||||
await instance.update({ lastStatus: STATUS.ERROR, lastError: String(err.message || err) });
|
||||
}
|
||||
}
|
||||
|
||||
// Schedule Plugins based on config + overrides
|
||||
if (discoveryConfig && discoveryConfig.plugins) {
|
||||
for (const [name, config] of Object.entries(discoveryConfig.plugins)) {
|
||||
const mergedConfig = { ...config, ...(overrides[name] || {}) };
|
||||
if (mergedConfig.enabled && plugins[name]) {
|
||||
const cron = mergedConfig.cron || '0 * * * *'; // Default hourly
|
||||
await discoveryQueue.add('run_plugin', { pluginName: name, config: mergedConfig }, { repeat: { pattern: cron } });
|
||||
console.log(`[Scheduler] Scheduled plugin ${name} with cron ${cron}`);
|
||||
|
||||
// Also run once immediately
|
||||
await discoveryQueue.add('run_plugin', { pluginName: name, config: mergedConfig });
|
||||
}
|
||||
// Schedule one instance: upsert a repeatable JobScheduler keyed by its id. Does
|
||||
// NOT trigger an immediate run — call runInstanceNow(id) separately for that
|
||||
// (used on boot and on "load"). Safe to call repeatedly (upsert is idempotent
|
||||
// and will update the cron if it changed).
|
||||
async function scheduleInstance(instance) {
|
||||
if (!instance || !instance.id) return;
|
||||
if (!instance.enabled) { await unscheduleInstance(instance.id); return; }
|
||||
const cron = instance.cron || '0 * * * *';
|
||||
await discoveryQueue.upsertJobScheduler(pluginSchedulerId(instance.id), { pattern: cron }, {
|
||||
name: RUN,
|
||||
data: { instanceId: instance.id }
|
||||
});
|
||||
console.log(`[Scheduler] Scheduled instance ${instance.slug} with cron ${cron}`);
|
||||
}
|
||||
|
||||
// Remove an instance's repeatable schedule. No-op if it had none.
|
||||
async function unscheduleInstance(id) {
|
||||
if (!id) return;
|
||||
try { await discoveryQueue.removeJobScheduler(pluginSchedulerId(id)); }
|
||||
catch (err) { /* missing scheduler is fine */ }
|
||||
}
|
||||
|
||||
// Enqueue a single immediate run for an instance (the "Run now" button / boot
|
||||
// kick). Runs once regardless of enabled, on top of any schedule.
|
||||
async function runInstanceNow(id) {
|
||||
if (!id) return;
|
||||
await discoveryQueue.add(RUN, { instanceId: id });
|
||||
}
|
||||
|
||||
// One-time legacy migration: if the PluginInstance table is empty AND
|
||||
// conf.discovery.plugins has entries (the old static-config shape), seed one
|
||||
// instance per configured type and copy its secret fields into OpenBao. After
|
||||
// the first boot, the table is non-empty and the static config is ignored.
|
||||
// Idempotent (guarded by the empty-table check).
|
||||
async function migrateLegacyPlugins(discoveryConfig) {
|
||||
const existing = await PluginInstance.list();
|
||||
if (existing && existing.length) return;
|
||||
|
||||
const legacy = discoveryConfig && discoveryConfig.plugins;
|
||||
if (!legacy || typeof legacy !== 'object') return;
|
||||
const names = Object.keys(legacy);
|
||||
if (!names.length) return;
|
||||
|
||||
console.log(`[Scheduler] Migrating ${names.length} legacy discovery plugin(s) to instances…`);
|
||||
for (const name of names) {
|
||||
const entry = legacy[name] || {};
|
||||
const manifest = pluginRegistry.getManifest(name);
|
||||
if (!manifest) {
|
||||
console.warn(`[Scheduler] legacy plugin '${name}' has no registered type — skipping`);
|
||||
continue;
|
||||
}
|
||||
// splitConfig keeps only declared configSchema fields and separates secret
|
||||
// from non-secret. Legacy `enabled`/`cron` are not in configSchema, so they
|
||||
// are dropped here and read from the entry directly below.
|
||||
const { config, secrets } = pluginRegistry.splitConfig(name, entry);
|
||||
const instance = await PluginInstance.create({
|
||||
pluginType: name,
|
||||
category: manifest.category,
|
||||
name: manifest.name,
|
||||
slug: name,
|
||||
enabled: entry.enabled !== false,
|
||||
cron: entry.cron || '0 * * * *',
|
||||
config,
|
||||
created_by: 'legacy-migration'
|
||||
});
|
||||
try {
|
||||
await pluginSecrets.write(instance.id, secrets);
|
||||
console.log(`[Scheduler] migrated '${name}' -> instance ${instance.id} (slug ${instance.slug})`);
|
||||
} catch (err) {
|
||||
// The instance row exists; if we can't write secrets (e.g. the sso-broker
|
||||
// policy predates theta-suite v1.30.1) the operator gets a clear error
|
||||
// from the API on edit, and the instance still runs with its non-secret
|
||||
// config. Don't delete the row — the operator just needs to re-run
|
||||
// setup.sh and edit/save the secrets.
|
||||
console.error(`[Scheduler] migrated '${name}' row but FAILED to write secrets:`, err.message);
|
||||
await instance.update({ lastStatus: STATUS.ERROR, lastError: `secret migration failed: ${err.message}` });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = { initScheduler, discoveryQueue, connection };
|
||||
// Boot-time initialization: clear stale schedulers, schedule garbage collection,
|
||||
// migrate any legacy static-config plugins, then schedule every enabled
|
||||
// instance and kick one immediate run for each.
|
||||
async function initScheduler(discoveryConfig) {
|
||||
// Clear stale plugin/gc schedulers from a previous boot. Other-named
|
||||
// schedulers (none in this app) are left alone.
|
||||
try {
|
||||
const schedulers = await discoveryQueue.getJobSchedulers();
|
||||
for (const s of schedulers) {
|
||||
if (s.name === RUN || s.name === GC) {
|
||||
await discoveryQueue.removeJobScheduler(s.key || s.id);
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
console.log('[Scheduler] Could not clear old job schedulers:', e.message);
|
||||
}
|
||||
|
||||
// Daily garbage collection of stale discovery resources.
|
||||
await discoveryQueue.upsertJobScheduler(GC, { pattern: '0 0 * * *' }, { name: GC, data: {} });
|
||||
|
||||
try {
|
||||
await migrateLegacyPlugins(discoveryConfig);
|
||||
} catch (err) {
|
||||
console.error('[Scheduler] legacy migration failed:', err.message);
|
||||
}
|
||||
|
||||
const enabled = await PluginInstance.listEnabled();
|
||||
for (const instance of enabled) {
|
||||
await scheduleInstance(instance);
|
||||
await runInstanceNow(instance.id); // boot kick
|
||||
}
|
||||
console.log(`[Scheduler] initialized — ${enabled.length} instance(s) scheduled`);
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
initScheduler,
|
||||
scheduleInstance,
|
||||
unscheduleInstance,
|
||||
runInstanceNow,
|
||||
discoveryQueue,
|
||||
connection
|
||||
};
|
||||
Reference in New Issue
Block a user