Skip to content

Commit 05ea678

Browse files
[web] Add actions for waking up workflow from sleep and re-enqueue runs for debugging (#582)
1 parent b56aae3 commit 05ea678

13 files changed

Lines changed: 1257 additions & 239 deletions

File tree

‎.changeset/fine-moles-sit.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
"@workflow/web-shared": patch
3+
"@workflow/web": patch
4+
---
5+
6+
Add buttons to wake up workflow from sleep or scheduling issues

‎packages/web-shared/src/api/workflow-api-client.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,10 @@ import {
2323
fetchStreams,
2424
readStreamServerAction,
2525
recreateRun as recreateRunServerAction,
26+
wakeUpRun as wakeUpRunServerAction,
27+
type StopSleepOptions,
28+
type StopSleepResult,
29+
reenqueueRun as reenqueueRunServerAction,
2630
} from './workflow-server-actions';
2731

2832
const MAX_ITEMS = 1000;
@@ -1101,6 +1105,35 @@ export async function recreateRun(env: EnvMap, runId: string): Promise<string> {
11011105
return resultData;
11021106
}
11031107

1108+
/**
1109+
* Wake up a workflow run by re-enqueuing it
1110+
*/
1111+
export async function reenqueueRun(env: EnvMap, runId: string): Promise<void> {
1112+
const { error } = await unwrapServerActionResult(
1113+
reenqueueRunServerAction(env, runId)
1114+
);
1115+
if (error) {
1116+
throw error;
1117+
}
1118+
}
1119+
1120+
/**
1121+
* Wake up a workflow run by interrupting any pending sleep() calls
1122+
*/
1123+
export async function wakeUpRun(
1124+
env: EnvMap,
1125+
runId: string,
1126+
options?: StopSleepOptions
1127+
): Promise<StopSleepResult> {
1128+
const { error, result: resultData } = await unwrapServerActionResult(
1129+
wakeUpRunServerAction(env, runId, options)
1130+
);
1131+
if (error) {
1132+
throw error;
1133+
}
1134+
return resultData;
1135+
}
1136+
11041137
function isServerActionError(value: unknown): value is ServerActionError {
11051138
return (
11061139
typeof value === 'object' &&

‎packages/web-shared/src/api/workflow-server-actions.ts‎

Lines changed: 126 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -507,6 +507,132 @@ export async function recreateRun(
507507
}
508508
}
509509

510+
/**
511+
* Re-enqueue a workflow run.
512+
*
513+
* This re-enqueues the workflow orchestration layer. It's a no-op unless the workflow
514+
* got stuck due to an implementation issue in the World. Useful for debugging custom Worlds.
515+
*/
516+
export async function reenqueueRun(
517+
worldEnv: EnvMap,
518+
runId: string
519+
): Promise<ServerActionResult<void>> {
520+
try {
521+
const world = getWorldFromEnv({ ...worldEnv });
522+
const run = await world.runs.get(runId);
523+
const deploymentId = run.deploymentId;
524+
525+
await world.queue(
526+
`__wkf_workflow_${run.workflowName}`,
527+
{
528+
runId,
529+
},
530+
{
531+
deploymentId,
532+
}
533+
);
534+
535+
return createResponse(undefined);
536+
} catch (error) {
537+
return createServerActionError<void>(error, 'reenqueueRun', { runId });
538+
}
539+
}
540+
541+
export interface StopSleepResult {
542+
/** Number of pending sleeps that were stopped */
543+
stoppedCount: number;
544+
}
545+
546+
export interface StopSleepOptions {
547+
/**
548+
* Optional list of specific correlation IDs to target.
549+
* If provided, only these sleep calls will be interrupted.
550+
* If not provided, all pending sleep calls will be interrupted.
551+
*/
552+
correlationIds?: string[];
553+
}
554+
555+
/**
556+
* Wake up a workflow run by interrupting pending sleep() calls.
557+
*
558+
* This finds wait_created events without matching wait_completed events,
559+
* creates wait_completed events for them, and then re-enqueues the run.
560+
*
561+
* @param worldEnv - Environment configuration for the World
562+
* @param runId - The run ID to wake up
563+
* @param options - Optional settings to narrow down targeting (specific correlation IDs)
564+
*/
565+
export async function wakeUpRun(
566+
worldEnv: EnvMap,
567+
runId: string,
568+
options?: StopSleepOptions
569+
): Promise<ServerActionResult<StopSleepResult>> {
570+
try {
571+
const world = getWorldFromEnv({ ...worldEnv });
572+
const run = await world.runs.get(runId);
573+
const deploymentId = run.deploymentId;
574+
575+
// Fetch all events for the run
576+
const eventsResult = await world.events.list({
577+
runId,
578+
pagination: { limit: 1000 },
579+
resolveData: 'none',
580+
});
581+
582+
// Find wait_created events without matching wait_completed events
583+
const waitCreatedEvents = eventsResult.data.filter(
584+
(e) => e.eventType === 'wait_created'
585+
);
586+
const waitCompletedCorrelationIds = new Set(
587+
eventsResult.data
588+
.filter((e) => e.eventType === 'wait_completed')
589+
.map((e) => e.correlationId)
590+
);
591+
592+
let pendingWaits = waitCreatedEvents.filter(
593+
(e) => !waitCompletedCorrelationIds.has(e.correlationId)
594+
);
595+
596+
// If specific correlation IDs are provided, filter to only those
597+
if (options?.correlationIds && options.correlationIds.length > 0) {
598+
const targetCorrelationIds = new Set(options.correlationIds);
599+
pendingWaits = pendingWaits.filter(
600+
(e) => e.correlationId && targetCorrelationIds.has(e.correlationId)
601+
);
602+
}
603+
604+
// Create wait_completed events for each pending wait
605+
for (const waitEvent of pendingWaits) {
606+
if (waitEvent.correlationId) {
607+
await world.events.create(runId, {
608+
eventType: 'wait_completed',
609+
correlationId: waitEvent.correlationId,
610+
});
611+
}
612+
}
613+
614+
// Re-enqueue the run to wake it up
615+
if (pendingWaits.length > 0) {
616+
await world.queue(
617+
`__wkf_workflow_${run.workflowName}`,
618+
{
619+
runId,
620+
},
621+
{
622+
deploymentId,
623+
}
624+
);
625+
}
626+
627+
return createResponse({ stoppedCount: pendingWaits.length });
628+
} catch (error) {
629+
return createServerActionError<StopSleepResult>(error, 'wakeUpRun', {
630+
runId,
631+
correlationIds: options?.correlationIds,
632+
});
633+
}
634+
}
635+
510636
export async function readStreamServerAction(
511637
env: EnvMap,
512638
streamId: string,

‎packages/web-shared/src/index.ts‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,17 +2,28 @@ export {
22
parseStepName,
33
parseWorkflowName,
44
} from '@workflow/core/parse-name';
5-
65
export type { Event, Hook, Step, WorkflowRun } from '@workflow/world';
6+
77
export * from './api/workflow-api-client';
88
export type { EnvMap } from './api/workflow-server-actions';
9+
10+
export type { EventAnalysis } from './lib/event-analysis';
11+
export {
12+
analyzeEvents,
13+
hasPendingHooksFromEvents,
14+
hasPendingSleepsFromEvents,
15+
hasPendingStepsFromEvents,
16+
isTerminalStatus,
17+
shouldShowReenqueueButton,
18+
} from './lib/event-analysis';
919
export type { StreamStep } from './lib/utils';
1020
export {
1121
extractConversation,
1222
formatDuration,
1323
identifyStreamSteps,
1424
isDoStreamStep,
1525
} from './lib/utils';
26+
1627
export { RunTraceView } from './run-trace-view';
1728
export { ConversationView } from './sidebar/conversation-view';
1829
export { StreamViewer } from './stream-viewer';

0 commit comments

Comments
 (0)