Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .changeset/zstd-step-error-display.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
'@workflow/world-vercel': patch
'@workflow/web-shared': patch
'@workflow/web': patch
---

Decompress gzip- and zstd-prefixed serialized data returned from Vercel Workflow storage, and route OSS web hydration through the async WASM-capable path for compressed payloads.
12 changes: 8 additions & 4 deletions packages/web-shared/src/components/sidebar/events-list.tsx
Original file line number Diff line number Diff line change
@@ -1,8 +1,12 @@
'use client';

import { EVENT_DATA_REF_FIELDS, type Event } from '@workflow/world';
import type { Event } from '@workflow/world';
import { useCallback, useLayoutEffect, useMemo, useRef, useState } from 'react';
import { hasEncryptedFields, isExpiredMarker } from '../../lib/hydration';
import {
getEventDataRefFields,
hasEncryptedFields,
isExpiredMarker,
} from '../../lib/hydration';
import {
Collapsible,
CollapsibleContent,
Expand Down Expand Up @@ -230,15 +234,15 @@ function EventItem({

/**
* Check if an eventData object has only expired marker values in ref/payload
* fields for this event type (see {@link EVENT_DATA_REF_FIELDS}). Other keys
* fields for this event type (see {@link getEventDataRefFields}). Other keys
* (e.g. `resumeAt`, `stepName`) are ignored.
*/
function hasOnlyExpiredFields(data: unknown, eventType: string): boolean {
if (data === null || typeof data !== 'object' || Array.isArray(data)) {
return false;
}
const record = data as Record<string, unknown>;
const refKeys = EVENT_DATA_REF_FIELDS[eventType] ?? [];
const refKeys = getEventDataRefFields(eventType);
const presentKeys = refKeys.filter((k) => k in record);
return (
presentKeys.length > 0 &&
Expand Down
15 changes: 8 additions & 7 deletions packages/web-shared/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,13 +14,6 @@ export {
waitEventsToWaitEntity,
} from './components/workflow-traces/trace-span-construction';
export type { EventAnalysis } from './lib/event-analysis';
export {
parseExactWorkflowSearchId,
looksLikeWorkflowIdSearchInput,
type ExactWorkflowSearchId,
type ExactWorkflowSearchIdKind,
type ExactIdSearchResult,
} from './lib/exact-event-search-id';
export {
analyzeEvents,
hasPendingHooksFromEvents,
Expand All @@ -40,6 +33,13 @@ export {
materializeSteps,
materializeWaits,
} from './lib/event-materialization';
export {
type ExactIdSearchResult,
type ExactWorkflowSearchId,
type ExactWorkflowSearchIdKind,
looksLikeWorkflowIdSearchInput,
parseExactWorkflowSearchId,
} from './lib/exact-event-search-id';
export type { Revivers, StreamRef } from './lib/hydration';
export {
CLASS_INSTANCE_REF_TYPE,
Expand All @@ -49,6 +49,7 @@ export {
getWebRevivers,
hasEncryptedFields,
hydrateResourceIO,
hydrateResourceIOAsync,
hydrateResourceIOWithKey,
isClassInstanceRef,
isEncryptedMarker,
Expand Down
73 changes: 48 additions & 25 deletions packages/web-shared/src/lib/hydration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,20 @@ import {
} from '@workflow/core/serialization-format';
import { EVENT_DATA_REF_FIELDS } from '@workflow/world';

const V4_EXTRA_EVENT_DATA_REF_FIELDS: Record<string, string[]> = {
run_started: ['input'],
step_started: ['input'],
};

export function getEventDataRefFields(eventType: string): string[] {
return [
...new Set([
...(EVENT_DATA_REF_FIELDS[eventType] ?? []),
...(V4_EXTRA_EVENT_DATA_REF_FIELDS[eventType] ?? []),
]),
];
}

// Re-export types and utilities that consumers need
export {
CLASS_INSTANCE_REF_TYPE,
Expand Down Expand Up @@ -446,7 +460,7 @@ function replaceEncryptedAndExpiredWithMarkers<T>(resource: T): T {
if (result.eventData && typeof result.eventData === 'object') {
const eventType =
typeof result.eventType === 'string' ? result.eventType : '';
const refKeys = EVENT_DATA_REF_FIELDS[eventType] ?? [];
const refKeys = getEventDataRefFields(eventType);
const ed = { ...(result.eventData as Record<string, unknown>) };
for (const key of refKeys) {
if (key in ed) {
Expand All @@ -465,43 +479,56 @@ function replaceEncryptedAndExpiredWithMarkers<T>(resource: T): T {
* When a key is provided, encrypted fields are decrypted before hydration.
* This is the async version used when the user clicks "Decrypt" in the web UI.
*
* Handles both top-level fields (input, output, metadata) and nested
* eventData subfields per `EVENT_DATA_REF_FIELDS` from `@workflow/world` for that event type.
* Handles both top-level fields (input, output, metadata) and nested eventData
* payload fields for that event type.
*/
export async function hydrateResourceIOWithKey<T>(
resource: T,
key: Uint8Array
): Promise<T> {
return hydrateResourceIOAsync(resource, key);
}

/**
* Async hydration for web display.
*
* This follows the same resource-field mapping as {@link hydrateResourceIO},
* but can also inflate compressed browser payloads through the registered
* zstd WASM decoder. When a key is provided, encrypted fields are decrypted
* first and then inflated/hydrated.
*/
export async function hydrateResourceIOAsync<T>(
resource: T,
key?: Uint8Array
): Promise<T> {
const { hydrateDataWithKey } = await import(
'@workflow/core/serialization-format'
);
const { importKey } = await import('@workflow/core/encryption');
// Payloads may be zstd-compressed (the Web DecompressionStream has no zstd);
// register the WASM-backed browser decoder before hydrating. Idempotent and
// lazy — the WASM is only compiled when a zstd payload is actually decoded.
const { ensureZstdDecoderRegistered } = await import(
'./zstd-browser-decoder.js'
);
ensureZstdDecoderRegistered();
const cryptoKey = await importKey(key);
const cryptoKey = key
? await import('@workflow/core/encryption').then(({ importKey }) =>
importKey(key)
)
: undefined;
const revivers = getRevivers();

/** Extract original encrypted bytes from a marker or raw Uint8Array, then decrypt + hydrate */
async function decryptField(
value: unknown,
rev: Revivers,
k: Awaited<ReturnType<typeof importKey>>
): Promise<unknown> {
async function hydrateField(value: unknown): Promise<unknown> {
// Already-hydrated: encrypted marker with stored bytes
if (isEncryptedMarker(value)) {
const raw = (value as any).__encryptedData as Uint8Array;
return hydrateDataWithKey(raw, rev, k);
return cryptoKey ? hydrateDataWithKey(raw, revivers, cryptoKey) : value;
}
// Raw encrypted Uint8Array (not yet hydrated)
// Raw Uint8Array: may be encrypted, compressed, or plain devalue.
if (value instanceof Uint8Array) {
return hydrateDataWithKey(value, rev, k);
return hydrateDataWithKey(value, revivers, cryptoKey);
}
// Not encrypted — return as-is
// Not serialized — return as-is.
return value;
}

Expand All @@ -511,29 +538,25 @@ export async function hydrateResourceIOWithKey<T>(
// Decrypt + hydrate top-level serialized fields (runs, steps, hooks)
for (const field of ['input', 'output', 'metadata', 'error']) {
if (field in result) {
result[field] = await decryptField(result[field], revivers, cryptoKey);
result[field] = await hydrateField(result[field]);
}
}

// Decrypt + hydrate eventData subfields (events)
// Hydrate eventData subfields (events)
if (result.eventData && typeof result.eventData === 'object') {
const eventType =
typeof result.eventType === 'string' ? result.eventType : '';
const refKeys = EVENT_DATA_REF_FIELDS[eventType] ?? [];
const refKeys = getEventDataRefFields(eventType);
const eventData = { ...(result.eventData as Record<string, unknown>) };
for (const field of refKeys) {
if (field in eventData) {
eventData[field] = await decryptField(
eventData[field],
revivers,
cryptoKey
);
eventData[field] = await hydrateField(eventData[field]);
}
}
result.eventData = eventData;
}

return result as T;
return replaceEncryptedAndExpiredWithMarkers(result as T);
}

/**
Expand All @@ -552,7 +575,7 @@ export function hasEncryptedFields(resource: unknown): boolean {

if (r.eventData && typeof r.eventData === 'object') {
const eventType = typeof r.eventType === 'string' ? r.eventType : '';
const refKeys = EVENT_DATA_REF_FIELDS[eventType] ?? [];
const refKeys = getEventDataRefFields(eventType);
const ed = r.eventData as Record<string, unknown>;
for (const key of refKeys) {
if (key in ed && isEncryptedMarker(ed[key])) return true;
Expand Down
111 changes: 109 additions & 2 deletions packages/web-shared/test/hydration.test.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,19 @@
import { dehydrateStepError } from '@workflow/core/serialization';
import { importKey } from '@workflow/core/encryption';
import {
dehydrateStepError,
dehydrateStepReturnValue,
} from '@workflow/core/serialization';
import { hydrateData } from '@workflow/core/serialization-format';
import { FatalError, RetryableError } from '@workflow/errors';
import { describe, expect, it } from 'vitest';
import { getWebRevivers } from '../src/lib/hydration.js';
import {
getWebRevivers,
hasEncryptedFields,
hydrateResourceIO,
hydrateResourceIOAsync,
hydrateResourceIOWithKey,
isEncryptedMarker,
} from '../src/lib/hydration.js';

/**
* The web reviver set must mirror every key in `SerializableSpecial` (see
Expand All @@ -20,6 +31,7 @@ import { getWebRevivers } from '../src/lib/hydration.js';
*/

const REVIVERS = getWebRevivers();
const textDecoder = new TextDecoder();

/** Run a real value through the production wire path with no encryption. */
async function roundTrip<T>(value: unknown): Promise<T> {
Expand Down Expand Up @@ -141,3 +153,98 @@ describe('getWebRevivers — error family', () => {
expect(revived.retryAfter).toBeUndefined();
});
});

describe('front hydration — encrypted compressed payloads', () => {
const runId = 'wrun_test';
const rawKey = new Uint8Array(32).fill(7);

function formatPrefix(value: unknown): string {
expect(value).toBeInstanceOf(Uint8Array);
return textDecoder.decode((value as Uint8Array).subarray(0, 4));
}

it('keeps encrypted compressed step errors as markers until decrypting them with the run key', async () => {
const cryptoKey = await importKey(rawKey);
const original = new Error(
`boom ${'front encrypted payload '.repeat(400)}`
);
const wire = await dehydrateStepError(
original,
runId,
cryptoKey,
[],
globalThis,
true
);

expect(formatPrefix(wire)).toBe('encr');

const hydrated = hydrateResourceIO({
stepId: 'step_test',
error: wire,
});

expect(isEncryptedMarker(hydrated.error)).toBe(true);
expect(hasEncryptedFields(hydrated)).toBe(true);

const decrypted = await hydrateResourceIOWithKey(hydrated, rawKey);
expect(decrypted.error).toBeInstanceOf(Error);
expect((decrypted.error as Error).message).toBe(original.message);
});

it('hydrates unencrypted compressed step errors through the async web path', async () => {
const original = new Error(
`boom ${'oss web compressed payload '.repeat(400)}`
);
const wire = await dehydrateStepError(
original,
runId,
undefined,
[],
globalThis,
true
);

expect(['gzip', 'zstd']).toContain(formatPrefix(wire));

const hydrated = await hydrateResourceIOAsync({
stepId: 'step_test',
error: wire,
});

expect(hydrated.error).toBeInstanceOf(Error);
expect((hydrated.error as Error).message).toBe(original.message);
});

it('decrypts encrypted compressed v4 step_started input payloads', async () => {
const cryptoKey = await importKey(rawKey);
const input = ['probe', { message: 'encrypted front payload' }];
const wire = await dehydrateStepReturnValue(
input,
runId,
cryptoKey,
[],
globalThis,
false,
false,
true
);

expect(formatPrefix(wire)).toBe('encr');

const hydrated = hydrateResourceIO({
eventId: 'evnt_test',
eventType: 'step_started',
eventData: {
stepName: 'probe',
input: wire,
},
});

expect(isEncryptedMarker(hydrated.eventData.input)).toBe(true);
expect(hasEncryptedFields(hydrated)).toBe(true);

const decrypted = await hydrateResourceIOWithKey(hydrated, rawKey);
expect(decrypted.eventData.input).toEqual(input);
});
});
17 changes: 9 additions & 8 deletions packages/web/app/components/run-detail-view.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,7 @@ import {
DecryptButton,
ErrorBoundary,
EventListView,
hydrateResourceIO,
hydrateResourceIOWithKey,
hydrateResourceIOAsync,
NewTraceViewer,
type SidebarDataContextValue,
StreamViewer,
Expand Down Expand Up @@ -301,9 +300,10 @@ export function RunDetailView({
if (error) {
throw error;
}
const fullEvent = encryptionKeyRef.current
? await hydrateResourceIOWithKey(result, encryptionKeyRef.current)
: hydrateResourceIO(result);
const fullEvent = await hydrateResourceIOAsync(
result,
encryptionKeyRef.current ?? undefined
);
if ('eventData' in fullEvent) {
return fullEvent.eventData;
}
Expand All @@ -321,9 +321,10 @@ export function RunDetailView({
if (error) {
throw error;
}
const fullEvent = encryptionKeyRef.current
? await hydrateResourceIOWithKey(result, encryptionKeyRef.current)
: hydrateResourceIO(result);
const fullEvent = await hydrateResourceIOAsync(
result,
encryptionKeyRef.current ?? undefined
);
if ('eventData' in fullEvent) {
return fullEvent.eventData;
}
Expand Down
Loading
Loading