Source: Router/events/service.js

// @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;
    }
};