// @ts-check
/**
* @typedef {Object} PublishEventOptions
* @property {string} organization - The organization UUID for multi-tenant scoping
* @property {any[] | any} data - The event data payload
* @property {number} [backDate] - Unix timestamp for backdated events
* @property {boolean} [eventLog=true] - Whether to also persist the event to eventLog
*/
/**
* A database connection capable of executing statements (either the transaction
* connection wrapper from @commtool/sql-query or the default pool).
* @typedef {Object} EventLogConnection
* @property {function(string, any[], Object=): Promise<any>} query
*/
import { query } from '@commtool/sql-query';
import redis from 'redis'
import { errorLogger } from './requestLogger.js';
/** @type {redis.RedisClientType | null} */
export let pubClient = null
/** @type {string} */
let redisUrl = 'unknown'
/** @type {boolean} */
let isInitialized = false
/**
* Der laufende Initialisierungsversuch — alle Aufrufer teilen sich **einen**.
*
* Ohne das teilen sie sich nur den Zustand `isInitialized`, und der wird erst
* NACH `await client.connect()` gesetzt. Jeder Aufrufer, der in diesem Fenster
* startet, sieht `false`, erzeugt seinen eigenen Client und schreibt ihn in
* `pubClient`. Die Listener-Zeile liest `pubClient` danach **erneut** und trifft
* damit den ZULETZT erzeugten Client: beim Serverstart waren das 12 verbundene
* Clients (11 davon verwaist) und 11 `error`-Listener auf einem — genau die
* `MaxListenersExceededWarning: 11 error listeners added to [Commander]`
* (`Commander` ist die Client-Klasse von `@redis/client`).
* @type {Promise<boolean>|null}
*/
let initPromise = null
/**
* Einmaliger Verbindungsaufbau. Nur über {@link initializeRedis} aufrufen.
* @async
* @returns {Promise<boolean>} true bei Erfolg, false solange die Secrets fehlen
*/
const connectRedis = async () => {
try {
// Import the merged secrets after they've been loaded
const {mergedSecrets} = await import('../http-server.js');
// Check if secrets are available yet
if (!mergedSecrets || !mergedSecrets.redisServer) {
// Secrets not loaded yet - this is expected during early bootstrap
return false;
}
// the messaging system a seperate redis client is recommended
// Use dedicated event Redis if configured (eventRedisServer in vault/env), otherwise fall back to main Redis
const eventServer = mergedSecrets.eventRedisServer || process.env.EVENT_REDIS_SERVER || mergedSecrets.redisServer;
const eventPort = mergedSecrets.eventRedisPort || process.env.EVENT_REDIS_PORT || mergedSecrets.redisPort || 6379;
const eventPassword = mergedSecrets.eventRedisPassword || process.env.EVENT_REDIS_PASSWORD || mergedSecrets.redisPassword || undefined;
redisUrl=`redis://${eventServer}:${eventPort}`;
const client = redis.createClient({
url: redisUrl,
password: eventPassword,
legacyMode: false,
socket: {
connectTimeout: 10000, // 10 seconds
reconnectStrategy: (retries) => {
if (retries > 5) return new Error('Max retries reached');
return Math.min(retries * 1000, 5000); // Exponential backoff
}
}
});
// Vor connect(): ein Fehler beim Verbindungsaufbau waere sonst ein
// unbehandeltes 'error'-Event am Client.
client.on('error', /** @param {Error} err */ (err) => console.log('client error', err));
await client.connect();
pubClient = client;
isInitialized = true;
console.log(`[pubSubClient] Connected to Redis at ${redisUrl}`);
return true;
} catch (err) {
const errorMessage = `[pubSubClient] Failed to connect to Redis at ${redisUrl || 'unknown'}`;
console.error(errorMessage, err);
errorLogger(err);
throw err;
}
}
/**
* Initialize the Redis connection after secrets are loaded.
*
* Nebenläufigkeitsvertrag: parallele Aufrufe starten **einen** Versuch und
* bekommen dasselbe Ergebnis. Ein Versuch, der `false` liefert (Secrets noch
* nicht geladen) oder wirft, wird **nicht** festgeschrieben — der nächste
* Aufruf versucht es erneut. Genau darauf ist der Startpfad angewiesen:
* `server.js` ruft früh (noch ohne Secrets) auf und bekommt `false`, die erste
* Welle `publishEvent` initialisiert danach parallel.
* @async
* @function initializeRedis
* @returns {Promise<boolean>} Returns true if initialized successfully, false otherwise
*/
export async function initializeRedis() {
if (isInitialized) return true;
if (!initPromise) {
initPromise = connectRedis();
// Fehlschlag nicht festschreiben, sonst gaebe es nie wieder einen Versuch.
// Der zweite Handler ist nur dafuer da, dass diese Ableitung nicht als
// unbehandelte Rejection endet — den Fehler sieht der Aufrufer ueber `initPromise`.
initPromise.then(
(ok) => { if (!ok) initPromise = null; },
() => { initPromise = null; }
);
}
return initPromise;
}
/**
* Build the wire payload for an event. The payload keeps the existing members
* event contract: `{ data, UIDorga, backDate, timestamp }`.
* @param {Object} options
* @param {string|null} [options.organization]
* @param {any} [options.data]
* @param {number|null} [options.backDate]
* @returns {string} JSON payload
*/
export const buildEventPayload = ({ organization = null, data = null, backDate = null }) => {
// we always send the orga as well
const UIDorga = organization ? organization : 'UUID-00000000-0000-0000-0000-000000000000'
return data !== null || backDate !== null ?
JSON.stringify({data: data, UIDorga, backDate: backDate ? backDate : Math.floor(Date.now()/1000), timestamp: Date.now()}) :
JSON.stringify({UIDorga, timestamp: Date.now()});
};
/**
* Publish a prepared payload to Redis (live channel). Never throws.
* @param {string} eventKey
* @param {string} dataJson
* @returns {Promise<void>}
*/
export const publishToRedis = async (eventKey, dataJson) => {
try {
if (!isInitialized) {
const initialized = await initializeRedis();
if (!initialized) {
console.log(`[publishEvent] Skipping event '${eventKey}' - Redis not initialized yet (early bootstrap)`);
return;
}
}
if (!pubClient) {
console.warn(`[publishEvent] No Redis client available, skipping event '${eventKey}'`);
return;
}
pubClient.publish(eventKey, dataJson).catch(/** @param {Error} err */ err => errorLogger(err));
} catch (err) {
console.error('[publishEvent] Redis client error:', err);
errorLogger(err);
}
};
/**
* Persist an event to eventLog. When `connection` is provided (inside a
* transaction) the write joins the same transaction as the business change, so
* no permanent "data change without event" can occur.
*
* The primary key is (Timestamp, EventKey) at microsecond precision; a repeat
* publish within the same microsecond refreshes the payload instead of losing it.
*
* @param {string} eventKey
* @param {string} dataJson
* @param {EventLogConnection} [connection]
* @returns {Promise<void>}
*/
export const writeEventLog = async (eventKey, dataJson, connection = null) => {
const statement = `INSERT INTO eventLog (EventKey,Data) VALUES(?,?)
ON DUPLICATE KEY UPDATE Data=VALUE(Data)`;
if (connection) {
await connection.query(statement, [eventKey, dataJson]);
} else {
await query(statement, [eventKey, dataJson]);
}
};
/**
* Publishes an event to Redis with multi-tenant organization scoping and
* persists it to eventLog for replay.
* @param {string} eventKey - The event key/channel to publish to
* @param {PublishEventOptions} options - Publishing options
* @returns {Promise<void>}
*/
export const publishEvent=async (eventKey, options)=>
{
try {
if (!options || typeof options !== 'object') {
throw new TypeError(`Options must be an object, got ${typeof options}`);
}
const { organization = null, data = null, backDate = null, eventLog = true } = options;
const myData = buildEventPayload({ organization, data, backDate });
await publishToRedis(eventKey, myData);
if (eventLog) {
// fire-and-forget: callers rely on the existing non-blocking behavior
writeEventLog(eventKey, myData);
}
} catch (err) {
console.error('[publishEvent] Event publishing error:', err);
errorLogger(err);
}
}