// @ts-check
/**
* Event replay service - documented service access to eventLog with an opaque,
* strictly ordered cursor. See 080-Workspaces/025-Members-Backend-Implementation-Plan.mdx
* Phase 7 and the event-record.schema.json contract.
*
* Cursor format: `seq-<seq>` where seq is the eventLog AUTO_INCREMENT column
* (monotonic, unique). Consumers treat cursors as opaque.
*
* @import {ExpressRequestAuthorized} from '../../types.js'
*/
import { query, pool } from '@commtool/sql-query';
import { apiError } from '../../utils/apiEnvelope.js';
import { errorLoggerRead } from '../../utils/requestLogger.js';
const DEFAULT_LIMIT = 100;
const MAX_LIMIT = 1000;
const cursorRe = /^seq-(\d+)$/;
/**
* Parse an opaque cursor into its numeric seq.
* @param {string} cursor
* @returns {number}
*/
const parseCursor = (cursor) => {
const m = cursorRe.exec(String(cursor).trim());
if (!m) throw apiError(400, 'INVALID_CURSOR', 'Cursor has an invalid format');
return parseInt(m[1], 10);
};
/**
* Build a cursor from a numeric seq.
* @param {number} seq
* @returns {string}
*/
const buildCursor = (seq) => `seq-${seq}`;
/**
* Replay events for an organization in strict cursor order.
* @param {ExpressRequestAuthorized} req
* @returns {Promise<Object>}
*/
export const replayEvents = async (req) => {
try {
const organization = typeof req.query.organization === 'string' ? req.query.organization : '';
if (!organization.match(/^UUID-[0-9A-Fa-f]{8}-[0-9A-Fa-f]{4}-[0-9A-Fa-f]{4}-[0-9A-Fa-f]{4}-[0-9A-Fa-f]{12}$/)) {
throw apiError(400, 'INVALID_ORGANIZATION', 'organization must be an organization UID');
}
let after = 0;
if (req.query.after !== undefined) {
if (typeof req.query.after !== 'string') throw apiError(400, 'INVALID_CURSOR', 'after must be a single cursor');
after = parseCursor(req.query.after);
}
let limit = DEFAULT_LIMIT;
if (req.query.limit !== undefined) {
limit = parseInt(String(req.query.limit), 10);
if (Number.isNaN(limit) || limit < 1 || limit > MAX_LIMIT) {
throw apiError(400, 'INVALID_LIMIT', `limit must be between 1 and ${MAX_LIMIT}`);
}
}
// Retention check: a cursor below the oldest retained row can never be
// caught up - report 410 instead of silently returning the current state.
const [minRow] = await query(`SELECT MIN(seq) AS min_seq FROM eventLog`);
if (after > 0 && minRow.min_seq !== null && after < minRow.min_seq) {
throw apiError(410, 'EVENT_CURSOR_EXPIRED', 'Cursor is older than the retained event history');
}
const [wmRow] = await query(
`SELECT MAX(seq) AS max_seq FROM eventLog WHERE UIDorga = ?`,
[organization],
);
const highWatermark = buildCursor(wmRow.max_seq ?? 0);
const rows = await query(
`SELECT seq, EventKey, Data
FROM eventLog
WHERE seq > ? AND UIDorga = ?
ORDER BY seq ASC
LIMIT ?`,
[after, organization, limit + 1],
{ cast: ['json'] },
);
const hasMore = rows.length > limit;
const page = hasMore ? rows.slice(0, limit) : rows;
const events = [];
for (const row of page) {
let payload;
try {
payload = typeof row.Data === 'object' && row.Data !== null ? row.Data : JSON.parse(row.Data);
} catch (e) {
payload = { data: [], UIDorga: organization, backDate: 0, timestamp: 0 };
}
events.push({
schema_version: 1,
cursor: buildCursor(row.seq),
event_key: row.EventKey,
payload,
});
}
const nextCursor = events.length ? events[events.length - 1].cursor : buildCursor(after);
return {
success: true,
result: {
events,
next_cursor: nextCursor,
has_more: hasMore,
high_watermark: highWatermark,
},
};
} catch (e) {
errorLoggerRead(e);
throw e;
}
};