mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
improvement(logs): move per-block progress markers to Redis to cut write amplification (#5248)
* improvement(logs): move per-block progress markers to Redis to cut write amplification
Per-block lastStartedBlock/lastCompletedBlock markers were persisted via a
jsonb_set UPDATE on workflow_execution_logs on every block start and complete
(~2N UPDATEs per run) — the heaviest write query in the DB. These are live
progress breadcrumbs with no DB-polling consumer (live progress comes from the
executor over WebSocket); their only durable value is a breadcrumb folded into
the final record.
Behind the redis-progress-markers flag, markers now live in Redis during the run
and are folded into the single terminal UPDATE at completion, dropping per-run
row UPDATEs from ~2N+1 to 1.
- New progress-markers module: HASH execution:progress:{id}, atomic Lua
monotonic-guard writes preserving the existing <= ordering, reservation-aligned
TTL backstop, graceful no-op when Redis is unavailable
- Deterministic GC: cleared at every terminal/pause boundary; TTL covers crashes
- Flag resolved once per logging session so a run never mixes write paths
- Fold markers into the completion record (Redis wins, falls back to row markers)
- Merge live markers for in-flight detail reads
- Extract shared getExecutionReservationTtlMs so marker and admission-slot TTLs
share one source of truth
* fix(logs): SQL fallback when Redis marker write fails, fold markers on force-fail, validate marker shape
Addresses review feedback on the redis-progress-markers PR:
- persistLast* now falls back to the jsonb_set UPDATE when Redis is unavailable or the write fails (setLast* returns whether it persisted), so a marker is never dropped when the flag is on without a healthy Redis.
- markExecutionAsFailed folds live Redis markers into execution_data before clearing, so the last-started/last-completed breadcrumb survives the force-fail path.
- getProgressMarkers validates marker shape (rebuilds from typed fields), so a stale or wrong-shaped Redis value can never reach API consumers.
* chore(logs): convert inline marker comments to TSDoc
* fix(logs): preserve markers when the completion read fails
getProgressMarkers now returns null on a Redis read error (vs {} for genuinely empty). completeWorkflowExecution and markExecutionAsFailed skip clearProgressMarkers when the read returns null, so a transient read error at completion no longer wipes markers that are still durably in Redis — the TTL reclaims them instead.
* fix(logs): resolve marker store split-brain by latest-timestamp-wins + drain on force-fail
- When a Redis marker write falls back to SQL, Redis and the row can each hold a marker for a different block; reads/folds previously preferred Redis unconditionally and could pick a stale value. Now the completion fold, the in-flight detail read, and the force-fail fold all pick the marker with the later timestamp (pickLatestStartedMarker/pickLatestCompletedMarker; markExecutionAsFailed uses a monotonic SQL guard).
- markAsFailed now drains pending per-block marker writes (not just the completion promise) before folding, so a force-fail racing onBlockStart/onBlockComplete still captures the latest breadcrumb.
* fix(logs): harden Lua marker guard against non-table decoded values
Guard the monotonic-check index with type(decoded) == 'table' so a corrupted Redis field that decodes to a non-table (e.g. a number) can't error the eval; our write path only ever stores JSON objects, so this is defense-in-depth.
* perf(logs): skip completion Redis read/clear when markers went to SQL
completeWorkflowExecution now takes readProgressMarkers (the session's resolved marker mode); when the flag is off it skips the per-completion HGETALL+DEL entirely instead of probing a key that was never written. Sticky to the session so it stays flip-safe (an execution that wrote to Redis always folds+clears Redis). Non-session callers default to true (safe read-and-fold). Also hardened the Lua guard with type(decoded)=='table'.
This commit is contained in:
@@ -5,7 +5,7 @@ import { getPlanTypeForLimits } from '@/lib/billing/plan-helpers'
|
||||
import { isOrgScopedSubscription } from '@/lib/billing/subscriptions/utils'
|
||||
import { isBillingEnabled } from '@/lib/core/config/env-flags'
|
||||
import { getRedisClient } from '@/lib/core/config/redis'
|
||||
import { getMaxExecutionTimeout } from '@/lib/core/execution-limits'
|
||||
import { getExecutionReservationTtlMs } from '@/lib/core/execution-limits'
|
||||
import type { SubscriptionPlan } from '@/lib/core/rate-limiter/types'
|
||||
|
||||
const logger = createLogger('UsageReservation')
|
||||
@@ -41,9 +41,6 @@ const MAX_CONCURRENT_EXECUTIONS: Record<SubscriptionPlan, number> = {
|
||||
*/
|
||||
const SLOT_COST_ESTIMATE = BASE_EXECUTION_CHARGE
|
||||
|
||||
/** Safety buffer added to the reservation TTL beyond the max execution timeout. */
|
||||
const RESERVATION_TTL_BUFFER_MS = 60_000
|
||||
|
||||
const INFLIGHT_KEY_PREFIX = 'usage:inflight:'
|
||||
const POINTER_KEY_PREFIX = 'usage:reservation:'
|
||||
|
||||
@@ -135,7 +132,7 @@ export async function reserveExecutionSlot(
|
||||
const maxConcurrency = getMaxConcurrentExecutions(subscription?.plan)
|
||||
const headroom = Math.max(0, limit - currentUsage)
|
||||
const headroomSlots = Math.floor(headroom / SLOT_COST_ESTIMATE)
|
||||
const ttlMs = getMaxExecutionTimeout() + RESERVATION_TTL_BUFFER_MS
|
||||
const ttlMs = getExecutionReservationTtlMs()
|
||||
const now = Date.now()
|
||||
const expiryScore = now + ttlMs
|
||||
|
||||
|
||||
@@ -77,6 +77,7 @@ export const env = createEnv({
|
||||
TABLE_SNAPSHOT_CACHE: z.boolean().optional(), // Mount tables into sandboxes by reference via a version-keyed CSV snapshot in object storage instead of draining the whole table into web-process heap
|
||||
PII_REDACTION: z.boolean().optional(), // Redact PII from workflow logs via configurable Data Retention rules (Presidio at the logger persist choke point) and expose the Data Retention config UI
|
||||
TRIGGER_EU_REGION: z.boolean().optional(), // Route Trigger.dev runs to eu-central-1 instead of the default us-east-1 (fallback for the trigger-eu-region flag when AppConfig is not the source of truth)
|
||||
REDIS_PROGRESS_MARKERS: z.boolean().optional(), // Write per-block live progress markers to Redis instead of jsonb_set UPDATEs on workflow_execution_logs (fallback for the redis-progress-markers flag when AppConfig is not the source of truth)
|
||||
|
||||
// Table feature limits (per plan). Apply when billing is disabled (free tier defaults) or for billed plans.
|
||||
FREE_TABLES_LIMIT: z.number().optional(), // Max user tables per workspace on free tier (default: 5)
|
||||
|
||||
@@ -97,6 +97,14 @@ const FEATURE_FLAGS = {
|
||||
'resolveTriggerRegion, so the whole deployment switches regions together.',
|
||||
fallback: 'TRIGGER_EU_REGION',
|
||||
},
|
||||
'redis-progress-markers': {
|
||||
description:
|
||||
'Write per-block live progress markers (lastStartedBlock/lastCompletedBlock) to Redis ' +
|
||||
'instead of jsonb_set UPDATEs on workflow_execution_logs, folding them into the single ' +
|
||||
'terminal UPDATE at completion. Eliminates the heaviest write query. Resolved once per ' +
|
||||
'logging session (no user/org context) so an execution never mixes write paths.',
|
||||
fallback: 'REDIS_PROGRESS_MARKERS',
|
||||
},
|
||||
} satisfies Record<string, FeatureFlagDefinition>
|
||||
|
||||
/**
|
||||
|
||||
@@ -75,6 +75,19 @@ export function getMaxExecutionTimeout(): number {
|
||||
return EXECUTION_TIMEOUTS.enterprise.async
|
||||
}
|
||||
|
||||
/** Safety buffer added beyond the max execution timeout for execution-lifetime TTLs. */
|
||||
export const RESERVATION_TTL_BUFFER_MS = 60_000
|
||||
|
||||
/**
|
||||
* TTL (ms) bounding how long a single execution can remain in flight: the max
|
||||
* execution timeout plus a safety buffer. Shared source of truth for the
|
||||
* admission-reservation key and the live progress-marker key so they expire on
|
||||
* the same timeline.
|
||||
*/
|
||||
export function getExecutionReservationTtlMs(): number {
|
||||
return getMaxExecutionTimeout() + RESERVATION_TTL_BUFFER_MS
|
||||
}
|
||||
|
||||
export const DEFAULT_EXECUTION_TIMEOUT_MS = EXECUTION_TIMEOUTS.free.sync
|
||||
|
||||
export function isTimeoutError(error: unknown): boolean {
|
||||
|
||||
@@ -37,6 +37,13 @@ import {
|
||||
replaceLargeValueReferenceKeysWithClient,
|
||||
} from '@/lib/execution/payloads/large-value-metadata'
|
||||
import { type RedactablePayload, redactPIIFromExecution } from '@/lib/logs/execution/pii-redaction'
|
||||
import {
|
||||
clearProgressMarkers,
|
||||
type ExecutionProgressMarkers,
|
||||
getProgressMarkers,
|
||||
pickLatestCompletedMarker,
|
||||
pickLatestStartedMarker,
|
||||
} from '@/lib/logs/execution/progress-markers'
|
||||
import { snapshotService } from '@/lib/logs/execution/snapshot/service'
|
||||
import {
|
||||
externalizeExecutionData,
|
||||
@@ -417,8 +424,15 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
return minimalWithSize.executionData
|
||||
}
|
||||
|
||||
/**
|
||||
* Assemble the final `execution_data` for the terminal UPDATE. Live progress
|
||||
* markers are sourced from `progressMarkers` (Redis, current run) and fall back
|
||||
* to markers already on the row — the legacy SQL path and resumed rows that
|
||||
* folded markers in at their prior pause boundary.
|
||||
*/
|
||||
private buildCompletedExecutionData(params: {
|
||||
existingExecutionData?: WorkflowExecutionLog['executionData']
|
||||
progressMarkers?: ExecutionProgressMarkers
|
||||
traceSpans?: TraceSpan[]
|
||||
finalOutput: BlockOutputData
|
||||
finalizationPath?: ExecutionFinalizationPath
|
||||
@@ -436,6 +450,7 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
}): WorkflowExecutionLog['executionData'] {
|
||||
const {
|
||||
existingExecutionData,
|
||||
progressMarkers,
|
||||
traceSpans,
|
||||
finalOutput,
|
||||
finalizationPath,
|
||||
@@ -446,6 +461,15 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
} = params
|
||||
const traceSpanCount = countTraceSpans(traceSpans)
|
||||
|
||||
const lastStartedBlock = pickLatestStartedMarker(
|
||||
progressMarkers?.lastStartedBlock,
|
||||
existingExecutionData?.lastStartedBlock
|
||||
)
|
||||
const lastCompletedBlock = pickLatestCompletedMarker(
|
||||
progressMarkers?.lastCompletedBlock,
|
||||
existingExecutionData?.lastCompletedBlock
|
||||
)
|
||||
|
||||
return {
|
||||
...(existingExecutionData?.environment
|
||||
? { environment: existingExecutionData.environment }
|
||||
@@ -459,12 +483,8 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
}
|
||||
: {}),
|
||||
...(existingExecutionData?.error ? { error: existingExecutionData.error } : {}),
|
||||
...(existingExecutionData?.lastStartedBlock
|
||||
? { lastStartedBlock: existingExecutionData.lastStartedBlock }
|
||||
: {}),
|
||||
...(existingExecutionData?.lastCompletedBlock
|
||||
? { lastCompletedBlock: existingExecutionData.lastCompletedBlock }
|
||||
: {}),
|
||||
...(lastStartedBlock ? { lastStartedBlock } : {}),
|
||||
...(lastCompletedBlock ? { lastCompletedBlock } : {}),
|
||||
...(completionFailure ? { completionFailure } : {}),
|
||||
...(finalizationPath ? { finalizationPath } : {}),
|
||||
hasTraceSpans: traceSpanCount > 0,
|
||||
@@ -659,6 +679,7 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
isResume?: boolean
|
||||
level?: 'info' | 'error'
|
||||
status?: 'completed' | 'failed' | 'cancelled' | 'pending'
|
||||
readProgressMarkers?: boolean
|
||||
}): Promise<WorkflowExecutionLog> {
|
||||
const {
|
||||
executionId,
|
||||
@@ -674,6 +695,7 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
isResume,
|
||||
level: levelOverride,
|
||||
status: statusOverride,
|
||||
readProgressMarkers = true,
|
||||
} = params
|
||||
|
||||
let execLog = logger.withMetadata({ executionId })
|
||||
@@ -731,8 +753,11 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
models: costSummary.models,
|
||||
}
|
||||
|
||||
const progressMarkers = readProgressMarkers ? await getProgressMarkers(executionId) : null
|
||||
|
||||
const builtExecutionData = this.buildCompletedExecutionData({
|
||||
existingExecutionData,
|
||||
progressMarkers: progressMarkers ?? undefined,
|
||||
traceSpans: mergedTraceSpans,
|
||||
finalOutput,
|
||||
finalizationPath,
|
||||
@@ -893,6 +918,8 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
return log
|
||||
})
|
||||
|
||||
if (progressMarkers !== null) void clearProgressMarkers(executionId)
|
||||
|
||||
try {
|
||||
// Skip workflow lookup if workflow was deleted.
|
||||
const wf = updatedLog.workflowId
|
||||
|
||||
@@ -66,6 +66,31 @@ vi.mock('@/lib/logs/execution/logger', () => ({
|
||||
},
|
||||
}))
|
||||
|
||||
const {
|
||||
isFeatureEnabledMock,
|
||||
setLastStartedBlockMock,
|
||||
setLastCompletedBlockMock,
|
||||
getProgressMarkersMock,
|
||||
clearProgressMarkersMock,
|
||||
} = vi.hoisted(() => ({
|
||||
isFeatureEnabledMock: vi.fn().mockResolvedValue(false),
|
||||
setLastStartedBlockMock: vi.fn().mockResolvedValue(false),
|
||||
setLastCompletedBlockMock: vi.fn().mockResolvedValue(false),
|
||||
getProgressMarkersMock: vi.fn().mockResolvedValue({}),
|
||||
clearProgressMarkersMock: vi.fn().mockResolvedValue(undefined),
|
||||
}))
|
||||
|
||||
vi.mock('@/lib/core/config/feature-flags', () => ({
|
||||
isFeatureEnabled: isFeatureEnabledMock,
|
||||
}))
|
||||
|
||||
vi.mock('@/lib/logs/execution/progress-markers', () => ({
|
||||
setLastStartedBlock: setLastStartedBlockMock,
|
||||
setLastCompletedBlock: setLastCompletedBlockMock,
|
||||
getProgressMarkers: getProgressMarkersMock,
|
||||
clearProgressMarkers: clearProgressMarkersMock,
|
||||
}))
|
||||
|
||||
vi.mock('@/lib/logs/execution/logging-factory', () => ({
|
||||
calculateCostSummary: vi.fn().mockReturnValue({
|
||||
totalCost: 0,
|
||||
@@ -647,4 +672,126 @@ describe('LoggingSession.markExecutionAsFailed workflowId scoping', () => {
|
||||
const combined = String(Array.from(strings)).toLowerCase() + values.join(' ').toLowerCase()
|
||||
expect(combined).toContain('force_failed')
|
||||
})
|
||||
|
||||
it('clears Redis markers when marking failed (terminal boundary outside completeWorkflowExecution)', async () => {
|
||||
await LoggingSession.markExecutionAsFailed('exec-3', 'boom', undefined, 'wf-3')
|
||||
expect(clearProgressMarkersMock).toHaveBeenCalledWith('exec-3')
|
||||
})
|
||||
|
||||
it('folds live Redis markers into the row before clearing on force-fail', async () => {
|
||||
getProgressMarkersMock.mockResolvedValueOnce({
|
||||
lastStartedBlock: { blockId: 'b1', blockName: 'Fetch', blockType: 'api', startedAt: 't1' },
|
||||
lastCompletedBlock: {
|
||||
blockId: 'b1',
|
||||
blockName: 'Fetch',
|
||||
blockType: 'api',
|
||||
endedAt: 't2',
|
||||
success: false,
|
||||
},
|
||||
})
|
||||
|
||||
await LoggingSession.markExecutionAsFailed('exec-9', 'boom', undefined, 'wf-9')
|
||||
|
||||
const folded = dbMocks.sql.mock.calls
|
||||
.map((c) => String(Array.from(c[0] as TemplateStringsArray)))
|
||||
.join(' ')
|
||||
expect(folded).toContain('lastStartedBlock')
|
||||
expect(folded).toContain('lastCompletedBlock')
|
||||
expect(clearProgressMarkersMock).toHaveBeenCalledWith('exec-9')
|
||||
})
|
||||
|
||||
it('does not clear markers when the Redis read fails (avoids wiping the only copy)', async () => {
|
||||
getProgressMarkersMock.mockResolvedValueOnce(null)
|
||||
await LoggingSession.markExecutionAsFailed('exec-readfail', 'boom', undefined, 'wf-x')
|
||||
expect(clearProgressMarkersMock).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('LoggingSession progress-marker write path', () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
startWorkflowExecutionMock.mockResolvedValue({})
|
||||
loadWorkflowStateForExecutionMock.mockResolvedValue({
|
||||
blocks: {},
|
||||
edges: [],
|
||||
loops: {},
|
||||
parallels: {},
|
||||
})
|
||||
dbMocks.execute.mockResolvedValue(undefined)
|
||||
})
|
||||
|
||||
it('writes markers to Redis (not the row) when the flag is on and Redis accepts the write', async () => {
|
||||
isFeatureEnabledMock.mockResolvedValue(true)
|
||||
setLastStartedBlockMock.mockResolvedValue(true)
|
||||
setLastCompletedBlockMock.mockResolvedValue(true)
|
||||
const session = new LoggingSession('wf-1', 'exec-redis', 'manual', 'req-1')
|
||||
await session.start({ workspaceId: 'ws-1' })
|
||||
|
||||
await session.onBlockStart('b1', 'Fetch', 'api', '2026-06-27T10:00:00.000Z')
|
||||
await session.onBlockComplete('b1', 'Fetch', 'api', { endedAt: '2026-06-27T10:00:01.000Z' })
|
||||
|
||||
expect(setLastStartedBlockMock).toHaveBeenCalledWith(
|
||||
'exec-redis',
|
||||
expect.objectContaining({ blockId: 'b1', startedAt: '2026-06-27T10:00:00.000Z' })
|
||||
)
|
||||
expect(setLastCompletedBlockMock).toHaveBeenCalledWith(
|
||||
'exec-redis',
|
||||
expect.objectContaining({ blockId: 'b1', success: true })
|
||||
)
|
||||
expect(dbMocks.execute).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('falls back to the SQL UPDATE when the flag is on but the Redis write fails', async () => {
|
||||
isFeatureEnabledMock.mockResolvedValue(true)
|
||||
setLastStartedBlockMock.mockResolvedValue(false)
|
||||
const session = new LoggingSession('wf-1', 'exec-redis-down', 'manual', 'req-1')
|
||||
await session.start({ workspaceId: 'ws-1' })
|
||||
|
||||
await session.onBlockStart('b1', 'Fetch', 'api', '2026-06-27T10:00:00.000Z')
|
||||
|
||||
expect(setLastStartedBlockMock).toHaveBeenCalled()
|
||||
expect(dbMocks.execute).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('writes markers via jsonb_set UPDATE when the flag is off', async () => {
|
||||
isFeatureEnabledMock.mockResolvedValue(false)
|
||||
const session = new LoggingSession('wf-1', 'exec-sql', 'manual', 'req-1')
|
||||
await session.start({ workspaceId: 'ws-1' })
|
||||
|
||||
await session.onBlockStart('b1', 'Fetch', 'api', '2026-06-27T10:00:00.000Z')
|
||||
|
||||
expect(dbMocks.execute).toHaveBeenCalledTimes(1)
|
||||
expect(setLastStartedBlockMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('falls back to the SQL path when flag resolution throws', async () => {
|
||||
isFeatureEnabledMock.mockRejectedValue(new Error('appconfig unavailable'))
|
||||
const session = new LoggingSession('wf-1', 'exec-fallback', 'manual', 'req-1')
|
||||
await session.start({ workspaceId: 'ws-1' })
|
||||
|
||||
await session.onBlockStart('b1', 'Fetch', 'api', '2026-06-27T10:00:00.000Z')
|
||||
|
||||
expect(dbMocks.execute).toHaveBeenCalledTimes(1)
|
||||
expect(setLastStartedBlockMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('tells completion to read Redis markers only when the flag is on (no wasted ops when off)', async () => {
|
||||
completeWorkflowExecutionMock.mockResolvedValue({})
|
||||
|
||||
isFeatureEnabledMock.mockResolvedValue(true)
|
||||
const onSession = new LoggingSession('wf-1', 'exec-on', 'manual', 'req-1')
|
||||
await onSession.start({ workspaceId: 'ws-1' })
|
||||
await onSession.safeComplete({ finalOutput: { ok: true } })
|
||||
expect(completeWorkflowExecutionMock).toHaveBeenLastCalledWith(
|
||||
expect.objectContaining({ executionId: 'exec-on', readProgressMarkers: true })
|
||||
)
|
||||
|
||||
isFeatureEnabledMock.mockResolvedValue(false)
|
||||
const offSession = new LoggingSession('wf-1', 'exec-off', 'manual', 'req-1')
|
||||
await offSession.start({ workspaceId: 'ws-1' })
|
||||
await offSession.safeComplete({ finalOutput: { ok: true } })
|
||||
expect(completeWorkflowExecutionMock).toHaveBeenLastCalledWith(
|
||||
expect.objectContaining({ executionId: 'exec-off', readProgressMarkers: false })
|
||||
)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -4,6 +4,7 @@ import { createLogger } from '@sim/logger'
|
||||
import { describeError, toError } from '@sim/utils/errors'
|
||||
import { and, eq, sql } from 'drizzle-orm'
|
||||
import { releaseExecutionSlot } from '@/lib/billing/calculations/usage-reservation'
|
||||
import { isFeatureEnabled } from '@/lib/core/config/feature-flags'
|
||||
import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure'
|
||||
import { executionLogger } from '@/lib/logs/execution/logger'
|
||||
import {
|
||||
@@ -13,6 +14,12 @@ import {
|
||||
loadDeployedWorkflowStateForLogging,
|
||||
loadWorkflowStateForExecution,
|
||||
} from '@/lib/logs/execution/logging-factory'
|
||||
import {
|
||||
clearProgressMarkers,
|
||||
getProgressMarkers,
|
||||
setLastCompletedBlock,
|
||||
setLastStartedBlock,
|
||||
} from '@/lib/logs/execution/progress-markers'
|
||||
import type {
|
||||
ExecutionEnvironment,
|
||||
ExecutionFinalizationPath,
|
||||
@@ -129,6 +136,13 @@ export class LoggingSession {
|
||||
private workflowState?: WorkflowState
|
||||
private correlation?: NonNullable<ExecutionTrigger['data']>['correlation']
|
||||
private isResume = false
|
||||
/**
|
||||
* Whether per-block progress markers go to Redis (vs jsonb_set UPDATEs on the
|
||||
* log row). Resolved once in {@link start} and cached so an execution never
|
||||
* mixes write paths across its block callbacks. Defaults to the legacy SQL
|
||||
* path until resolved.
|
||||
*/
|
||||
private useRedisMarkers = false
|
||||
private completed = false
|
||||
/** Synchronous flag to prevent concurrent completion attempts (race condition guard) */
|
||||
private completing = false
|
||||
@@ -167,7 +181,27 @@ export class LoggingSession {
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the per-block marker write path (Redis vs jsonb_set UPDATE) for this
|
||||
* session. Defaults to the legacy SQL path if flag resolution fails.
|
||||
*/
|
||||
private async resolveRedisMarkerMode(): Promise<boolean> {
|
||||
try {
|
||||
return await isFeatureEnabled('redis-progress-markers')
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist the last-started-block marker. Redis is the primary path when the
|
||||
* flag is on; falls back to the durable jsonb_set UPDATE when Redis is
|
||||
* unavailable or the write fails, so a marker is never dropped.
|
||||
*/
|
||||
private async persistLastStartedBlock(marker: ExecutionLastStartedBlock): Promise<void> {
|
||||
if (this.useRedisMarkers && (await setLastStartedBlock(this.executionId, marker))) {
|
||||
return
|
||||
}
|
||||
try {
|
||||
await db.execute(
|
||||
buildStartedMarkerPersistenceQuery({
|
||||
@@ -185,7 +219,15 @@ export class LoggingSession {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist the last-completed-block marker. Redis is the primary path when the
|
||||
* flag is on; falls back to the durable jsonb_set UPDATE when Redis is
|
||||
* unavailable or the write fails, so a marker is never dropped.
|
||||
*/
|
||||
private async persistLastCompletedBlock(marker: ExecutionLastCompletedBlock): Promise<void> {
|
||||
if (this.useRedisMarkers && (await setLastCompletedBlock(this.executionId, marker))) {
|
||||
return
|
||||
}
|
||||
try {
|
||||
await db.execute(
|
||||
buildCompletedMarkerPersistenceQuery({
|
||||
@@ -266,6 +308,7 @@ export class LoggingSession {
|
||||
isResume: this.isResume,
|
||||
level: params.level,
|
||||
status: params.status,
|
||||
readProgressMarkers: this.useRedisMarkers,
|
||||
})
|
||||
|
||||
// Release the admission reservation from preprocessing. Skipped on pause: a
|
||||
@@ -313,6 +356,8 @@ export class LoggingSession {
|
||||
} = params
|
||||
|
||||
try {
|
||||
this.useRedisMarkers = await this.resolveRedisMarkerMode()
|
||||
|
||||
this.trigger = createTriggerObject(this.triggerType, triggerData)
|
||||
this.correlation = triggerData?.correlation
|
||||
this.environment = createEnvironmentObject(
|
||||
@@ -974,8 +1019,14 @@ export class LoggingSession {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Force-fail the execution. Waits for any in-flight completion and drains
|
||||
* pending per-block marker writes first, so a force-fail racing
|
||||
* onBlockStart/onBlockComplete still captures the latest breadcrumb in the fold.
|
||||
*/
|
||||
async markAsFailed(errorMessage?: string): Promise<void> {
|
||||
await this.waitForCompletion()
|
||||
await this.drainPendingProgressWrites()
|
||||
await LoggingSession.markExecutionAsFailed(
|
||||
this.executionId,
|
||||
errorMessage,
|
||||
@@ -984,6 +1035,13 @@ export class LoggingSession {
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Force-fail terminal boundary that bypasses completeWorkflowExecution. Folds
|
||||
* any live Redis progress markers into execution_data before clearing the key,
|
||||
* so a run whose markers only ever lived in Redis still keeps its
|
||||
* last-started/last-completed breadcrumb. Both the fold and clear are no-ops
|
||||
* when the standard completion path already persisted and cleared them.
|
||||
*/
|
||||
static async markExecutionAsFailed(
|
||||
executionId: string,
|
||||
errorMessage: string | undefined,
|
||||
@@ -992,12 +1050,10 @@ export class LoggingSession {
|
||||
): Promise<void> {
|
||||
try {
|
||||
const message = errorMessage || 'Run failed'
|
||||
await db
|
||||
.update(workflowExecutionLogs)
|
||||
.set({
|
||||
level: 'error',
|
||||
status: 'failed',
|
||||
executionData: sql`jsonb_set(
|
||||
|
||||
const markers = await getProgressMarkers(executionId)
|
||||
|
||||
let executionData = sql`jsonb_set(
|
||||
jsonb_set(
|
||||
jsonb_set(
|
||||
COALESCE(execution_data, '{}'::jsonb),
|
||||
@@ -1009,8 +1065,25 @@ export class LoggingSession {
|
||||
),
|
||||
ARRAY['finalizationPath'],
|
||||
to_jsonb('force_failed'::text)
|
||||
)`,
|
||||
})
|
||||
)`
|
||||
if (markers?.lastStartedBlock) {
|
||||
const startedAt = markers.lastStartedBlock.startedAt
|
||||
const startedJson = JSON.stringify(markers.lastStartedBlock)
|
||||
executionData = sql`CASE WHEN COALESCE(jsonb_extract_path_text(execution_data, 'lastStartedBlock', 'startedAt'), '') <= ${startedAt}
|
||||
THEN jsonb_set(${executionData}, ARRAY['lastStartedBlock'], ${startedJson}::jsonb)
|
||||
ELSE ${executionData} END`
|
||||
}
|
||||
if (markers?.lastCompletedBlock) {
|
||||
const endedAt = markers.lastCompletedBlock.endedAt
|
||||
const completedJson = JSON.stringify(markers.lastCompletedBlock)
|
||||
executionData = sql`CASE WHEN COALESCE(jsonb_extract_path_text(execution_data, 'lastCompletedBlock', 'endedAt'), '') <= ${endedAt}
|
||||
THEN jsonb_set(${executionData}, ARRAY['lastCompletedBlock'], ${completedJson}::jsonb)
|
||||
ELSE ${executionData} END`
|
||||
}
|
||||
|
||||
await db
|
||||
.update(workflowExecutionLogs)
|
||||
.set({ level: 'error', status: 'failed', executionData })
|
||||
.where(
|
||||
and(
|
||||
eq(workflowExecutionLogs.executionId, executionId),
|
||||
@@ -1018,6 +1091,8 @@ export class LoggingSession {
|
||||
)
|
||||
)
|
||||
|
||||
if (markers !== null) void clearProgressMarkers(executionId)
|
||||
|
||||
logger.info(`[${requestId || 'unknown'}] Marked execution ${executionId} as failed`)
|
||||
} catch (error) {
|
||||
logger.error(`Failed to mark execution ${executionId} as failed:`, {
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
/**
|
||||
* @vitest-environment node
|
||||
*/
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { ExecutionLastCompletedBlock, ExecutionLastStartedBlock } from '@/lib/logs/types'
|
||||
|
||||
const { mockGetRedisClient, mockRedis } = vi.hoisted(() => {
|
||||
const mockRedis = {
|
||||
eval: vi.fn(),
|
||||
hgetall: vi.fn(),
|
||||
del: vi.fn(),
|
||||
}
|
||||
return { mockGetRedisClient: vi.fn<[], typeof mockRedis | null>(() => mockRedis), mockRedis }
|
||||
})
|
||||
|
||||
vi.mock('@/lib/core/config/redis', () => ({
|
||||
getRedisClient: mockGetRedisClient,
|
||||
}))
|
||||
|
||||
vi.mock('@/lib/core/execution-limits', () => ({
|
||||
getExecutionReservationTtlMs: () => 5_460_000,
|
||||
}))
|
||||
|
||||
import {
|
||||
clearProgressMarkers,
|
||||
getProgressMarkers,
|
||||
pickLatestCompletedMarker,
|
||||
pickLatestStartedMarker,
|
||||
setLastCompletedBlock,
|
||||
setLastStartedBlock,
|
||||
} from '@/lib/logs/execution/progress-markers'
|
||||
|
||||
const EXECUTION_ID = 'exec-1'
|
||||
const KEY = `execution:progress:${EXECUTION_ID}`
|
||||
const EXPECTED_TTL_MS = '5460000' // getExecutionReservationTtlMs() mock value
|
||||
|
||||
const startedMarker: ExecutionLastStartedBlock = {
|
||||
blockId: 'b1',
|
||||
blockName: 'Fetch',
|
||||
blockType: 'api',
|
||||
startedAt: '2026-06-27T10:00:00.000Z',
|
||||
}
|
||||
|
||||
const completedMarker: ExecutionLastCompletedBlock = {
|
||||
blockId: 'b1',
|
||||
blockName: 'Fetch',
|
||||
blockType: 'api',
|
||||
endedAt: '2026-06-27T10:00:01.000Z',
|
||||
success: true,
|
||||
}
|
||||
|
||||
describe('progress-markers', () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
mockGetRedisClient.mockReturnValue(mockRedis)
|
||||
mockRedis.eval.mockResolvedValue(1)
|
||||
mockRedis.hgetall.mockResolvedValue({})
|
||||
mockRedis.del.mockResolvedValue(1)
|
||||
})
|
||||
|
||||
describe('setLastStartedBlock', () => {
|
||||
it('evals the monotonic-guard script with key, started field, timestamp, json, and TTL', async () => {
|
||||
await setLastStartedBlock(EXECUTION_ID, startedMarker)
|
||||
|
||||
expect(mockRedis.eval).toHaveBeenCalledTimes(1)
|
||||
const [, numKeys, key, field, timestampField, timestamp, json, ttl] =
|
||||
mockRedis.eval.mock.calls[0]
|
||||
expect(numKeys).toBe(1)
|
||||
expect(key).toBe(KEY)
|
||||
expect(field).toBe('started')
|
||||
expect(timestampField).toBe('startedAt')
|
||||
expect(timestamp).toBe(startedMarker.startedAt)
|
||||
expect(JSON.parse(json as string)).toEqual(startedMarker)
|
||||
expect(ttl).toBe(EXPECTED_TTL_MS)
|
||||
})
|
||||
|
||||
it('returns true when the Redis write succeeds', async () => {
|
||||
await expect(setLastStartedBlock(EXECUTION_ID, startedMarker)).resolves.toBe(true)
|
||||
})
|
||||
|
||||
it('returns false (caller falls back to SQL) when the eval fails', async () => {
|
||||
mockRedis.eval.mockRejectedValueOnce(new Error('redis down'))
|
||||
await expect(setLastStartedBlock(EXECUTION_ID, startedMarker)).resolves.toBe(false)
|
||||
})
|
||||
|
||||
it('returns false and no-ops when Redis is unavailable', async () => {
|
||||
mockGetRedisClient.mockReturnValue(null)
|
||||
await expect(setLastStartedBlock(EXECUTION_ID, startedMarker)).resolves.toBe(false)
|
||||
expect(mockRedis.eval).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('setLastCompletedBlock', () => {
|
||||
it('evals with the completed field and endedAt timestamp', async () => {
|
||||
await setLastCompletedBlock(EXECUTION_ID, completedMarker)
|
||||
|
||||
const args = mockRedis.eval.mock.calls[0]
|
||||
expect(args[3]).toBe('completed')
|
||||
expect(args[4]).toBe('endedAt')
|
||||
expect(args[5]).toBe(completedMarker.endedAt)
|
||||
expect(JSON.parse(args[6] as string)).toEqual(completedMarker)
|
||||
})
|
||||
})
|
||||
|
||||
describe('getProgressMarkers', () => {
|
||||
it('parses both markers from the hash', async () => {
|
||||
mockRedis.hgetall.mockResolvedValueOnce({
|
||||
started: JSON.stringify(startedMarker),
|
||||
completed: JSON.stringify(completedMarker),
|
||||
})
|
||||
|
||||
const result = await getProgressMarkers(EXECUTION_ID)
|
||||
expect(mockRedis.hgetall).toHaveBeenCalledWith(KEY)
|
||||
expect(result).toEqual({
|
||||
lastStartedBlock: startedMarker,
|
||||
lastCompletedBlock: completedMarker,
|
||||
})
|
||||
})
|
||||
|
||||
it('returns only the present field', async () => {
|
||||
mockRedis.hgetall.mockResolvedValueOnce({ started: JSON.stringify(startedMarker) })
|
||||
const result = await getProgressMarkers(EXECUTION_ID)
|
||||
expect(result).toEqual({ lastStartedBlock: startedMarker })
|
||||
})
|
||||
|
||||
it('returns {} for an empty / missing key', async () => {
|
||||
mockRedis.hgetall.mockResolvedValueOnce({})
|
||||
expect(await getProgressMarkers(EXECUTION_ID)).toEqual({})
|
||||
})
|
||||
|
||||
it('returns {} and does not throw on malformed JSON', async () => {
|
||||
mockRedis.hgetall.mockResolvedValueOnce({ started: '{not json' })
|
||||
expect(await getProgressMarkers(EXECUTION_ID)).toEqual({})
|
||||
})
|
||||
|
||||
it('drops wrong-shaped JSON so malformed markers never reach clients', async () => {
|
||||
mockRedis.hgetall.mockResolvedValueOnce({
|
||||
started: JSON.stringify('just a string'),
|
||||
completed: JSON.stringify({ blockId: 123, blockName: 'x', blockType: 'api', endedAt: 'z' }),
|
||||
})
|
||||
expect(await getProgressMarkers(EXECUTION_ID)).toEqual({})
|
||||
})
|
||||
|
||||
it('strips extra fields, returning only the validated marker shape', async () => {
|
||||
mockRedis.hgetall.mockResolvedValueOnce({
|
||||
started: JSON.stringify({ ...startedMarker, secret: 'leak', extra: 1 }),
|
||||
})
|
||||
expect(await getProgressMarkers(EXECUTION_ID)).toEqual({ lastStartedBlock: startedMarker })
|
||||
})
|
||||
|
||||
it('returns {} when Redis is unavailable', async () => {
|
||||
mockGetRedisClient.mockReturnValue(null)
|
||||
expect(await getProgressMarkers(EXECUTION_ID)).toEqual({})
|
||||
expect(mockRedis.hgetall).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('returns null when the Redis read fails so callers do not clear the only copy', async () => {
|
||||
mockRedis.hgetall.mockRejectedValueOnce(new Error('redis down'))
|
||||
expect(await getProgressMarkers(EXECUTION_ID)).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
describe('clearProgressMarkers', () => {
|
||||
it('deletes the key', async () => {
|
||||
await clearProgressMarkers(EXECUTION_ID)
|
||||
expect(mockRedis.del).toHaveBeenCalledWith(KEY)
|
||||
})
|
||||
|
||||
it('swallows del errors', async () => {
|
||||
mockRedis.del.mockRejectedValueOnce(new Error('redis down'))
|
||||
await expect(clearProgressMarkers(EXECUTION_ID)).resolves.toBeUndefined()
|
||||
})
|
||||
|
||||
it('no-ops when Redis is unavailable', async () => {
|
||||
mockGetRedisClient.mockReturnValue(null)
|
||||
await clearProgressMarkers(EXECUTION_ID)
|
||||
expect(mockRedis.del).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('latest-wins pickers (stale-store safety)', () => {
|
||||
const older = { ...startedMarker, blockId: 'old', startedAt: '2026-06-27T10:00:00.000Z' }
|
||||
const newer = { ...startedMarker, blockId: 'new', startedAt: '2026-06-27T10:00:05.000Z' }
|
||||
|
||||
it('returns the defined side when the other is undefined', () => {
|
||||
expect(pickLatestStartedMarker(older, undefined)).toBe(older)
|
||||
expect(pickLatestStartedMarker(undefined, newer)).toBe(newer)
|
||||
expect(pickLatestStartedMarker(undefined, undefined)).toBeUndefined()
|
||||
})
|
||||
|
||||
it('picks the later startedAt regardless of argument order (row newer than Redis still wins)', () => {
|
||||
expect(pickLatestStartedMarker(older, newer)).toBe(newer)
|
||||
expect(pickLatestStartedMarker(newer, older)).toBe(newer)
|
||||
})
|
||||
|
||||
it('picks the later endedAt for completed markers', () => {
|
||||
const c1 = { ...completedMarker, endedAt: '2026-06-27T10:00:01.000Z' }
|
||||
const c2 = { ...completedMarker, endedAt: '2026-06-27T10:00:09.000Z' }
|
||||
expect(pickLatestCompletedMarker(c1, c2)).toBe(c2)
|
||||
expect(pickLatestCompletedMarker(c2, c1)).toBe(c2)
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,250 @@
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { toError } from '@sim/utils/errors'
|
||||
import { getRedisClient } from '@/lib/core/config/redis'
|
||||
import { getExecutionReservationTtlMs } from '@/lib/core/execution-limits'
|
||||
import type { ExecutionLastCompletedBlock, ExecutionLastStartedBlock } from '@/lib/logs/types'
|
||||
|
||||
const logger = createLogger('ExecutionProgressMarkers')
|
||||
|
||||
/**
|
||||
* Live per-block progress markers (`lastStartedBlock` / `lastCompletedBlock`)
|
||||
* used to be written to `workflow_execution_logs.execution_data` via a
|
||||
* `jsonb_set` UPDATE on every block start and complete — ~2N row UPDATEs per
|
||||
* run and the single heaviest write query in the database. They carry no
|
||||
* value beyond a breadcrumb folded into the final record (no client polls
|
||||
* them; live progress comes from the executor over WebSocket), so they live
|
||||
* in Redis during the run and are folded into the single terminal UPDATE at
|
||||
* completion. See the logs-contention plan for the full rationale.
|
||||
*/
|
||||
|
||||
/** Redis key namespace — matches the `execution:*` family (stream, cancel, budget). */
|
||||
const PROGRESS_KEY_PREFIX = 'execution:progress:'
|
||||
|
||||
const STARTED_FIELD = 'started'
|
||||
const COMPLETED_FIELD = 'completed'
|
||||
|
||||
/**
|
||||
* Single source of the client used for progress markers. Indirection kept so a
|
||||
* future split to a dedicated `LOGS_REDIS_URL` is a one-line change here.
|
||||
*/
|
||||
function getMarkerClient() {
|
||||
return getRedisClient()
|
||||
}
|
||||
|
||||
function markerKey(executionId: string): string {
|
||||
return `${PROGRESS_KEY_PREFIX}${executionId}`
|
||||
}
|
||||
|
||||
/**
|
||||
* Atomic monotonic write: set the field only when the incoming marker's embedded
|
||||
* timestamp is >= the stored one, then refresh the TTL — all in one script so
|
||||
* concurrent block callbacks can't race a read-modify-write. Preserves the exact
|
||||
* `<=` ordering the legacy SQL used (`COALESCE(stored, '') <= incoming`). ISO
|
||||
* UTC timestamps compare correctly lexicographically.
|
||||
*/
|
||||
const SET_MARKER_SCRIPT = `
|
||||
local existing = redis.call('HGET', KEYS[1], ARGV[1])
|
||||
if existing then
|
||||
local ok, decoded = pcall(cjson.decode, existing)
|
||||
if ok and type(decoded) == 'table' and decoded[ARGV[2]] and tostring(decoded[ARGV[2]]) > ARGV[3] then
|
||||
redis.call('PEXPIRE', KEYS[1], ARGV[5])
|
||||
return 0
|
||||
end
|
||||
end
|
||||
redis.call('HSET', KEYS[1], ARGV[1], ARGV[4])
|
||||
redis.call('PEXPIRE', KEYS[1], ARGV[5])
|
||||
return 1
|
||||
`
|
||||
|
||||
/**
|
||||
* Write a marker field under the monotonic guard, refreshing the key TTL. The
|
||||
* TTL is a backstop for executions that die without a terminal/pause boundary
|
||||
* (deterministic cleanup is {@link clearProgressMarkers}); it mirrors the
|
||||
* admission-reservation TTL so a crashed run's marker key and slot expire
|
||||
* together. Returns `true` only when the marker was durably written to Redis,
|
||||
* so callers can fall back to the SQL path on a missing client or a failure.
|
||||
*/
|
||||
async function setMarker(
|
||||
executionId: string,
|
||||
field: string,
|
||||
timestampField: 'startedAt' | 'endedAt',
|
||||
timestamp: string,
|
||||
marker: ExecutionLastStartedBlock | ExecutionLastCompletedBlock
|
||||
): Promise<boolean> {
|
||||
const redis = getMarkerClient()
|
||||
if (!redis) return false
|
||||
|
||||
try {
|
||||
await redis.eval(
|
||||
SET_MARKER_SCRIPT,
|
||||
1,
|
||||
markerKey(executionId),
|
||||
field,
|
||||
timestampField,
|
||||
timestamp,
|
||||
JSON.stringify(marker),
|
||||
getExecutionReservationTtlMs().toString()
|
||||
)
|
||||
return true
|
||||
} catch (error) {
|
||||
logger.error(`Failed to persist progress marker for execution ${executionId}`, {
|
||||
field,
|
||||
error: toError(error).message,
|
||||
})
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist the last-started-block marker. Returns `false` (caller should fall
|
||||
* back to the durable SQL path) when Redis is unavailable or the write fails.
|
||||
*/
|
||||
export async function setLastStartedBlock(
|
||||
executionId: string,
|
||||
marker: ExecutionLastStartedBlock
|
||||
): Promise<boolean> {
|
||||
return setMarker(executionId, STARTED_FIELD, 'startedAt', marker.startedAt, marker)
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist the last-completed-block marker. Returns `false` (caller should fall
|
||||
* back to the durable SQL path) when Redis is unavailable or the write fails.
|
||||
*/
|
||||
export async function setLastCompletedBlock(
|
||||
executionId: string,
|
||||
marker: ExecutionLastCompletedBlock
|
||||
): Promise<boolean> {
|
||||
return setMarker(executionId, COMPLETED_FIELD, 'endedAt', marker.endedAt, marker)
|
||||
}
|
||||
|
||||
export interface ExecutionProgressMarkers {
|
||||
lastStartedBlock?: ExecutionLastStartedBlock
|
||||
lastCompletedBlock?: ExecutionLastCompletedBlock
|
||||
}
|
||||
|
||||
/**
|
||||
* Pick the later of two last-started markers by `startedAt`. Markers can split
|
||||
* across stores — a failed Redis write falls back to the row, so an earlier
|
||||
* successful Redis write may coexist with a newer row marker (or vice versa).
|
||||
* Choosing by timestamp keeps the freshest breadcrumb regardless of which store
|
||||
* holds it. ISO UTC timestamps compare correctly lexicographically.
|
||||
*/
|
||||
export function pickLatestStartedMarker(
|
||||
a: ExecutionLastStartedBlock | undefined,
|
||||
b: ExecutionLastStartedBlock | undefined
|
||||
): ExecutionLastStartedBlock | undefined {
|
||||
if (!a) return b
|
||||
if (!b) return a
|
||||
return a.startedAt >= b.startedAt ? a : b
|
||||
}
|
||||
|
||||
/** Pick the later of two last-completed markers by `endedAt`. See {@link pickLatestStartedMarker}. */
|
||||
export function pickLatestCompletedMarker(
|
||||
a: ExecutionLastCompletedBlock | undefined,
|
||||
b: ExecutionLastCompletedBlock | undefined
|
||||
): ExecutionLastCompletedBlock | undefined {
|
||||
if (!a) return b
|
||||
if (!b) return a
|
||||
return a.endedAt >= b.endedAt ? a : b
|
||||
}
|
||||
|
||||
function safeJsonParse(raw: string | undefined): unknown {
|
||||
if (!raw) return undefined
|
||||
try {
|
||||
return JSON.parse(raw)
|
||||
} catch {
|
||||
return undefined
|
||||
}
|
||||
}
|
||||
|
||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
||||
return typeof value === 'object' && value !== null && !Array.isArray(value)
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a stored last-started marker, rebuilding it from validated fields so a
|
||||
* stale or wrong-shaped Redis value can never reach API consumers.
|
||||
*/
|
||||
function parseStartedMarker(raw: string | undefined): ExecutionLastStartedBlock | undefined {
|
||||
const v = safeJsonParse(raw)
|
||||
if (!isRecord(v)) return undefined
|
||||
const { blockId, blockName, blockType, startedAt } = v
|
||||
if (
|
||||
typeof blockId === 'string' &&
|
||||
typeof blockName === 'string' &&
|
||||
typeof blockType === 'string' &&
|
||||
typeof startedAt === 'string'
|
||||
) {
|
||||
return { blockId, blockName, blockType, startedAt }
|
||||
}
|
||||
return undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a stored last-completed marker, rebuilding it from validated fields so a
|
||||
* stale or wrong-shaped Redis value can never reach API consumers.
|
||||
*/
|
||||
function parseCompletedMarker(raw: string | undefined): ExecutionLastCompletedBlock | undefined {
|
||||
const v = safeJsonParse(raw)
|
||||
if (!isRecord(v)) return undefined
|
||||
const { blockId, blockName, blockType, endedAt, success } = v
|
||||
if (
|
||||
typeof blockId === 'string' &&
|
||||
typeof blockName === 'string' &&
|
||||
typeof blockType === 'string' &&
|
||||
typeof endedAt === 'string' &&
|
||||
typeof success === 'boolean'
|
||||
) {
|
||||
return { blockId, blockName, blockType, endedAt, success }
|
||||
}
|
||||
return undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Read both markers for an execution. Returns an empty object when Redis is
|
||||
* unavailable (markers were never stored here) or the key holds nothing, and
|
||||
* `null` when the Redis read itself failed — callers must treat `null` as
|
||||
* "unknown" and skip {@link clearProgressMarkers}, so a transient read error
|
||||
* never wipes the only copy of markers that are still in Redis.
|
||||
*/
|
||||
export async function getProgressMarkers(
|
||||
executionId: string
|
||||
): Promise<ExecutionProgressMarkers | null> {
|
||||
const redis = getMarkerClient()
|
||||
if (!redis) return {}
|
||||
|
||||
try {
|
||||
const fields = await redis.hgetall(markerKey(executionId))
|
||||
if (!fields || Object.keys(fields).length === 0) return {}
|
||||
|
||||
const result: ExecutionProgressMarkers = {}
|
||||
const started = parseStartedMarker(fields[STARTED_FIELD])
|
||||
if (started) result.lastStartedBlock = started
|
||||
const completed = parseCompletedMarker(fields[COMPLETED_FIELD])
|
||||
if (completed) result.lastCompletedBlock = completed
|
||||
return result
|
||||
} catch (error) {
|
||||
logger.error(`Failed to read progress markers for execution ${executionId}`, {
|
||||
error: toError(error).message,
|
||||
})
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete the markers for an execution. Called at every terminal/pause boundary
|
||||
* after the durable record has been written, so paused executions (which can
|
||||
* live indefinitely) hold no Redis keys. Fire-and-forget; no-op without Redis.
|
||||
*/
|
||||
export async function clearProgressMarkers(executionId: string): Promise<void> {
|
||||
const redis = getMarkerClient()
|
||||
if (!redis) return
|
||||
|
||||
try {
|
||||
await redis.del(markerKey(executionId))
|
||||
} catch (error) {
|
||||
logger.error(`Failed to clear progress markers for execution ${executionId}`, {
|
||||
error: toError(error).message,
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,12 @@ import {
|
||||
} from '@sim/db/schema'
|
||||
import { and, eq, type SQL } from 'drizzle-orm'
|
||||
import type { CostLedger } from '@/lib/api/contracts/logs'
|
||||
import {
|
||||
type ExecutionProgressMarkers,
|
||||
getProgressMarkers,
|
||||
pickLatestCompletedMarker,
|
||||
pickLatestStartedMarker,
|
||||
} from '@/lib/logs/execution/progress-markers'
|
||||
import { materializeExecutionData } from '@/lib/logs/execution/trace-store'
|
||||
import { checkWorkspaceAccess } from '@/lib/workspaces/permissions/utils'
|
||||
|
||||
@@ -77,6 +83,10 @@ interface FetchLogDetailArgs {
|
||||
* Shared loader for the workflow-log detail shape returned by the by-id and
|
||||
* by-execution routes. Returns `null` when no matching row exists in either
|
||||
* the workflow-execution or job-execution tables for this user + workspace.
|
||||
*
|
||||
* For in-flight (running/pending) executions, live progress markers are merged
|
||||
* from Redis, since they are only folded into the row at a terminal/pause
|
||||
* boundary.
|
||||
*/
|
||||
export async function fetchLogDetail({
|
||||
userId,
|
||||
@@ -165,6 +175,20 @@ export async function fetchLogDetail({
|
||||
{ workspaceId, workflowId: log.workflowId, executionId: log.executionId }
|
||||
)
|
||||
|
||||
const liveMarkers =
|
||||
log.status === 'running' || log.status === 'pending'
|
||||
? ((await getProgressMarkers(log.executionId)) ?? {})
|
||||
: {}
|
||||
const rowMarkers = (executionData ?? {}) as ExecutionProgressMarkers
|
||||
const mergedStartedBlock = pickLatestStartedMarker(
|
||||
liveMarkers.lastStartedBlock,
|
||||
rowMarkers.lastStartedBlock
|
||||
)
|
||||
const mergedCompletedBlock = pickLatestCompletedMarker(
|
||||
liveMarkers.lastCompletedBlock,
|
||||
rowMarkers.lastCompletedBlock
|
||||
)
|
||||
|
||||
return {
|
||||
id: log.id,
|
||||
workflowId: log.workflowId,
|
||||
@@ -190,6 +214,8 @@ export async function fetchLogDetail({
|
||||
executionData: {
|
||||
totalDuration: log.totalDurationMs,
|
||||
...executionData,
|
||||
...(mergedStartedBlock ? { lastStartedBlock: mergedStartedBlock } : {}),
|
||||
...(mergedCompletedBlock ? { lastCompletedBlock: mergedCompletedBlock } : {}),
|
||||
enhanced: true as const,
|
||||
},
|
||||
files: log.files ?? null,
|
||||
|
||||
@@ -463,5 +463,12 @@ export interface ExecutionLoggerService {
|
||||
isResume?: boolean
|
||||
level?: 'info' | 'error'
|
||||
status?: 'completed' | 'failed' | 'cancelled' | 'pending'
|
||||
/**
|
||||
* Whether this session wrote live progress markers to Redis. When false, the
|
||||
* completion fold skips the Redis read/clear entirely (markers are already on
|
||||
* the row via the SQL path). Defaults to true so non-session callers keep the
|
||||
* safe read-and-fold behavior.
|
||||
*/
|
||||
readProgressMarkers?: boolean
|
||||
}): Promise<WorkflowExecutionLog>
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user