From 2f455f0140415a59023eb1fefbee2decb8053f11 Mon Sep 17 00:00:00 2001 From: philipp Date: Sat, 6 Jun 2026 16:01:13 +0200 Subject: [PATCH] Bound quote lifecycle rollup catch-up Proof: Quote lifecycle retention now runs without blocking history-writer readiness, prevents concurrent maintenance passes, and bounds historical rollup catch-up before fail-closed pruning; regression tests cover the bounded stale-watermark path. Assumptions: Forty-eight DB-backed rollup granularity windows per pass is sufficient to catch up historical backlog without returning the entire raw quote firehose in one query. Still fake: Initial historical backlog may require multiple scheduled passes before pruning can begin; already-pruned non-success detail cannot be reconstructed. --- src/apps/history-writer.mjs | 63 ++++++++++++++++--------- src/lib/postgres.mjs | 15 ++++-- test/history-writer-static.test.mjs | 5 +- test/quote-lifecycle-retention.test.mjs | 61 ++++++++++++++++++++++++ 4 files changed, 117 insertions(+), 27 deletions(-) diff --git a/src/apps/history-writer.mjs b/src/apps/history-writer.mjs index f7c8764..f59f24a 100644 --- a/src/apps/history-writer.mjs +++ b/src/apps/history-writer.mjs @@ -229,11 +229,15 @@ await refreshQuoteLifecycleRetentionState().catch((error) => { state.retention_error = serializeError(error); }); +let quoteLifecycleRetentionMaintenanceInFlight = null; const quoteLifecycleRetentionMaintenanceTimer = setInterval(() => { void maybeRunQuoteLifecycleRetentionMaintenance(); }, quoteLifecycleRetentionMaintenanceIntervalMs); quoteLifecycleRetentionMaintenanceTimer.unref?.(); -await maybeRunQuoteLifecycleRetentionMaintenance({ force: true }); +const quoteLifecycleRetentionStartupTimer = setTimeout(() => { + void maybeRunQuoteLifecycleRetentionMaintenance({ force: true }); +}, 0); +quoteLifecycleRetentionStartupTimer.unref?.(); for (const historyConsumer of durableConsumers) { await runHistoryConsumer(historyConsumer); @@ -313,6 +317,10 @@ async function runHistoryConsumer(historyConsumer) { } async function maybeRunQuoteLifecycleRetentionMaintenance({ force = false } = {}) { + if (quoteLifecycleRetentionMaintenanceInFlight) { + return quoteLifecycleRetentionMaintenanceInFlight; + } + const nowMs = Date.now(); const lastMaintenanceMs = state.last_retention_maintenance_at ? Date.parse(state.last_retention_maintenance_at) @@ -325,31 +333,39 @@ async function maybeRunQuoteLifecycleRetentionMaintenance({ force = false } = {} return null; } - try { - const result = await runQuoteLifecycleRetentionMaintenance(pool, { - now: new Date(nowMs).toISOString(), - }); - applyRetentionMaintenanceResult(result); - await refreshQuoteLifecycleRetentionState(); - if (result.status === 'success' && Number(result.pruned_counts?.deletedCount || 0) > 0) { - logger.info('quote_lifecycle_retention_maintenance_completed', { - details: result, + quoteLifecycleRetentionMaintenanceInFlight = (async () => { + try { + const result = await runQuoteLifecycleRetentionMaintenance(pool, { + now: new Date(nowMs).toISOString(), }); - } else if (result.status !== 'success') { - logger.warn('quote_lifecycle_retention_maintenance_skipped', { - details: result, + applyRetentionMaintenanceResult(result); + await refreshQuoteLifecycleRetentionState(); + if (result.status === 'success' && Number(result.pruned_counts?.deletedCount || 0) > 0) { + logger.info('quote_lifecycle_retention_maintenance_completed', { + details: result, + }); + } else if (result.status !== 'success') { + logger.warn('quote_lifecycle_retention_maintenance_skipped', { + details: result, + }); + } + return result; + } catch (error) { + state.retention_error = serializeError(error); + state.quote_lifecycle_prune_error = state.retention_error; + logger.error('quote_lifecycle_retention_maintenance_failed', { + details: { + error: state.retention_error, + }, }); + return null; } - return result; - } catch (error) { - state.retention_error = serializeError(error); - state.quote_lifecycle_prune_error = state.retention_error; - logger.error('quote_lifecycle_retention_maintenance_failed', { - details: { - error: state.retention_error, - }, - }); - return null; + })(); + + try { + return await quoteLifecycleRetentionMaintenanceInFlight; + } finally { + quoteLifecycleRetentionMaintenanceInFlight = null; } } @@ -776,6 +792,7 @@ function resumeConsumers() { } async function shutdown() { + clearTimeout(quoteLifecycleRetentionStartupTimer); clearInterval(quoteLifecycleRetentionMaintenanceTimer); await controlApi.close().catch(() => {}); await Promise.allSettled([ diff --git a/src/lib/postgres.mjs b/src/lib/postgres.mjs index 90734cd..da638f3 100644 --- a/src/lib/postgres.mjs +++ b/src/lib/postgres.mjs @@ -53,6 +53,7 @@ const RETENTION_POLICIES_TABLE = 'retention_policies'; const QUOTE_LIFECYCLE_ROLLUPS_TABLE = 'quote_lifecycle_rollups'; const QUOTE_LIFECYCLE_ROLLUP_WATERMARKS_TABLE = 'quote_lifecycle_rollup_watermarks'; const QUOTE_LIFECYCLE_RETENTION_RUNS_TABLE = 'quote_lifecycle_retention_runs'; +const QUOTE_LIFECYCLE_ROLLUP_WINDOWS_PER_PASS = 48; const SUPPORTED_ASSET_IMPORT_RUNS_TABLE = 'supported_asset_import_runs'; const TRADING_ASSETS_TABLE = 'trading_assets'; const TRADING_PAIRS_TABLE = 'trading_pairs'; @@ -2640,9 +2641,15 @@ async function rollupQuoteLifecycleDetail(pool, { }; } + const boundedWindowEnd = new Date(Math.min( + timestampValue(windowEnd), + timestampValue(alignedWindowStart) + ( + policy.rollup_granularity_ms * QUOTE_LIFECYCLE_ROLLUP_WINDOWS_PER_PASS + ), + )).toISOString(); const rows = await loadQuoteLifecycleRollupInputRows(pool, { windowStart: alignedWindowStart, - windowEnd, + windowEnd: boundedWindowEnd, }); const rollups = buildQuoteLifecycleRollups({ rows, @@ -2652,7 +2659,7 @@ async function rollupQuoteLifecycleDetail(pool, { }); await upsertQuoteLifecycleRollups(pool, rollups); await upsertQuoteLifecycleRollupWatermark(pool, { - coveredUntil: windowEnd, + coveredUntil: boundedWindowEnd, computedAt: now, status: 'success', rowCount: rows.length, @@ -2663,7 +2670,9 @@ async function rollupQuoteLifecycleDetail(pool, { row_count: rows.length, rollup_count: rollups.length, window_start: alignedWindowStart, - covered_until: windowEnd, + covered_until: boundedWindowEnd, + target_covered_until: windowEnd, + partial: timestampValue(boundedWindowEnd) < timestampValue(windowEnd), }; } diff --git a/test/history-writer-static.test.mjs b/test/history-writer-static.test.mjs index 7efdbef..c81fab1 100644 --- a/test/history-writer-static.test.mjs +++ b/test/history-writer-static.test.mjs @@ -24,8 +24,11 @@ test('history writer replays durable topics but joins the raw quote firehose liv assert.match(source, /runQuoteLifecycleRetentionMaintenance/); assert.match(source, /loadQuoteLifecycleRetentionSummary/); assert.match(source, /maybeRunQuoteLifecycleRetentionMaintenance/); + assert.match(source, /quoteLifecycleRetentionMaintenanceInFlight/); assert.match(source, /quoteLifecycleRetentionMaintenanceTimer\s*=\s*setInterval/); - assert.match(source, /maybeRunQuoteLifecycleRetentionMaintenance\(\{\s*force:\s*true\s*\}\)/); + assert.match(source, /quoteLifecycleRetentionStartupTimer\s*=\s*setTimeout/); + assert.match(source, /void maybeRunQuoteLifecycleRetentionMaintenance\(\{\s*force:\s*true\s*\}\)/); + assert.match(source, /clearTimeout\(quoteLifecycleRetentionStartupTimer\)/); assert.match(source, /clearInterval\(quoteLifecycleRetentionMaintenanceTimer\)/); assert.match(source, /retention_mode/); assert.match(source, /storage_pressure/); diff --git a/test/quote-lifecycle-retention.test.mjs b/test/quote-lifecycle-retention.test.mjs index 180891a..61ff9d6 100644 --- a/test/quote-lifecycle-retention.test.mjs +++ b/test/quote-lifecycle-retention.test.mjs @@ -297,3 +297,64 @@ test('retention maintenance writes rollups before deleting detail and preserves && query.sql.includes('linked_settlement') ))); }); + +test('retention maintenance bounds historical rollup catch-up and fails pruning closed until covered', async () => { + const queries = []; + let watermarkReads = 0; + const policyRow = { + ...DEFAULT_QUOTE_LIFECYCLE_RETENTION_POLICY, + updated_at: '2026-06-06T09:00:00.000Z', + }; + const pool = { + async query(sql, params = []) { + queries.push({ sql, params }); + if (sql.includes('FROM retention_policies')) return { rows: [policyRow], rowCount: 1 }; + if (sql.includes('pg_database_size')) return { rows: [{ total_bytes: 10 }], rowCount: 1 }; + if (sql.includes('pg_stat_user_tables')) return { rows: [], rowCount: 0 }; + if (sql.includes('COUNT(*)::INT AS count') && sql.includes('quote_outcome_attributions')) { + return { rows: [{ count: 0 }], rowCount: 1 }; + } + if (sql.includes('FROM quote_lifecycle_rollup_watermarks')) { + watermarkReads += 1; + return watermarkReads === 1 + ? { rows: [], rowCount: 0 } + : { + rows: [{ + rollup_key: 'quote_lifecycle', + covered_until: '2026-06-05T04:00:00.000Z', + computed_at: '2026-06-06T11:00:00.000Z', + status: 'success', + row_count: 0, + rollup_count: 0, + }], + rowCount: 1, + }; + } + if (sql.includes('MIN(activity_at)')) { + return { rows: [{ earliest_at: '2026-06-05T00:00:00.000Z' }], rowCount: 1 }; + } + if (sql.includes('WITH subject_events')) return { rows: [], rowCount: 0 }; + if (sql.includes('quote_id IS NULL')) return { rows: [], rowCount: 0 }; + if (sql.includes('INSERT INTO quote_lifecycle_rollup_watermarks')) return { rows: [], rowCount: 1 }; + if (sql.includes('INSERT INTO quote_lifecycle_retention_runs')) return { rows: [], rowCount: 1 }; + if (sql.includes('DELETE FROM')) throw new Error('prune must not run before watermark reaches cutoff'); + return { rows: [], rowCount: 0 }; + }, + }; + + const result = await runQuoteLifecycleRetentionMaintenance(pool, { + now: '2026-06-06T11:00:00.000Z', + }); + const subjectQuery = queries.find((query) => query.sql.includes('WITH subject_events')); + + assert.equal(result.status, 'skipped'); + assert.equal(result.mode, 'blocked_rollup_stale'); + assert.equal(result.rollup_counts.partial, true); + assert.equal(result.rollup_counts.covered_until, '2026-06-05T04:00:00.000Z'); + assert.equal(result.rollup_counts.target_covered_until, '2026-06-06T10:30:00.000Z'); + assert.deepEqual(subjectQuery?.params, [ + '2026-06-05T00:00:00.000Z', + '2026-06-05T04:00:00.000Z', + ]); + assert.equal(queries.some((query) => query.sql.includes('DELETE FROM')), false); +});