Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
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
13 changes: 13 additions & 0 deletions .changeset/zod-45-compiled-schemas.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
---
'@workflow/ai': patch
'@workflow/cli': patch
'@workflow/core': patch
'@workflow/world': patch
'@workflow/world-local': patch
'@workflow/world-postgres': patch
'@workflow/world-testing': patch
'@workflow/world-vercel': patch
'workflow': patch
---

Upgrade runtime validation to Zod 4.5 and enable compilation on SDK-owned Zod schemas.
13 changes: 8 additions & 5 deletions packages/world-local/src/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ const MAX_SAFE_TIMEOUT_MS = 2147483647;
// The local workers share the same Node.js process and event loop,
// so we need to limit concurrency to avoid overwhelming the system.
const DEFAULT_CONCURRENCY_LIMIT = 1000;

const WORKFLOW_LOCAL_QUEUE_CONCURRENCY =
parseInt(process.env.WORKFLOW_LOCAL_QUEUE_CONCURRENCY ?? '0', 10) ||
DEFAULT_CONCURRENCY_LIMIT;
Expand Down Expand Up @@ -397,11 +398,13 @@ export function createQueue(config: Partial<Config>): LocalQueue {
return { messageId };
};

const HeaderParser = z.object({
'x-vqs-queue-name': ValidQueueName,
'x-vqs-message-id': MessageId,
'x-vqs-message-attempt': z.coerce.number(),
});
const HeaderParser = z.compile(
z.object({
'x-vqs-queue-name': ValidQueueName,
'x-vqs-message-id': MessageId,
'x-vqs-message-attempt': z.coerce.number(),
})
);

const createQueueHandler: Queue['createQueueHandler'] = (prefix, handler) => {
return async (req) => {
Expand Down
46 changes: 25 additions & 21 deletions packages/world-local/src/storage/events-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@ import type {
} from '@workflow/world';
import {
applyAttributeChanges,
EventSchema,
eventIdToSlot,
FIRST_EVENT_SLOT,
getMaxEventsPerRun,
Expand Down Expand Up @@ -107,6 +106,7 @@ import { handleLegacyEvent } from './legacy.js';
import {
purgeRunEntityData,
purgesUserDataOnFinish,
ReadEventSchema,
withRunPayloadsPurged,
} from './run-retention.js';
import { signalRunTerminal } from './run-status-signal.js';
Expand Down Expand Up @@ -169,9 +169,11 @@ function getHookRetentionLimitMs(): number {
* lifetimes can never share one marker (see
* `hookRecoveryMarkerPath`).
*/
const HookRecoveryMarkerSchema = z.object({
eventId: z.string(),
});
const HookRecoveryMarkerSchema = z.compile(
z.object({
eventId: z.string(),
})
);

/**
* Durable `(runId, resumeId)` claim for a lazy hook resume. Written via
Expand All @@ -182,13 +184,15 @@ const HookRecoveryMarkerSchema = z.object({
* records the content hash so a reused `resumeId` carrying a different payload
* can be rejected as a conflict, matching the server's constraint.
*/
const HookResumeClaimSchema = z.object({
runId: z.string(),
resumeId: z.string(),
hookId: z.string(),
eventId: z.string(),
payloadDigest: z.string().optional(),
});
const HookResumeClaimSchema = z.compile(
z.object({
runId: z.string(),
resumeId: z.string(),
hookId: z.string(),
eventId: z.string(),
payloadDigest: z.string().optional(),
})
);

/**
* Whether `event` is the `hook_received` a resume claim stands for.
Expand Down Expand Up @@ -237,7 +241,7 @@ async function findCommittedResumeEvent(
basedir,
'events',
`${runId}-${eventId}`,
EventSchema,
ReadEventSchema,
tag
);
if (
Expand Down Expand Up @@ -361,7 +365,7 @@ async function findExistingHookCreatedEventId(
): Promise<string | null> {
const result = await paginatedFileSystemQuery({
directory: path.join(basedir, 'events'),
schema: EventSchema,
schema: ReadEventSchema,
filePrefix: `${runId}-`,
filter: (event) =>
event.eventType === 'hook_created' &&
Expand Down Expand Up @@ -405,7 +409,7 @@ async function repairHookEntityFromPersistedEvent(
basedir,
'events',
compositeKey,
EventSchema,
ReadEventSchema,
tag
);
if (
Expand Down Expand Up @@ -851,7 +855,7 @@ export function createEventsStorage(
return;
}

const cachedEvent = EventSchema.safeParse(
const cachedEvent = ReadEventSchema.safeParse(
JSON.parse(serializedEvent, jsonReviver)
);
if (cachedEvent.success) {
Expand Down Expand Up @@ -905,7 +909,7 @@ export function createEventsStorage(
const queryRunEvents = (runId: string, pagination: PaginationOptions) =>
paginatedFileSystemQuery({
directory: path.join(basedir, 'events'),
schema: EventSchema,
schema: ReadEventSchema,
cachedItems: eventCache,
filePrefix: `${runId}-`,
sortOrder: pagination.sortOrder ?? 'asc',
Expand Down Expand Up @@ -1389,7 +1393,7 @@ export function createEventsStorage(
basedir,
'events',
`${effectiveRunId}-${committedClaim.eventId}`,
EventSchema,
ReadEventSchema,
tag
);
const committedEvent =
Expand Down Expand Up @@ -1482,7 +1486,7 @@ export function createEventsStorage(
basedir,
'events',
`${effectiveRunId}-${claim.eventId}`,
EventSchema,
ReadEventSchema,
tag
);
if (atClaimedId && isResumeEvent(atClaimedId, claim)) {
Expand Down Expand Up @@ -2895,7 +2899,7 @@ export function createEventsStorage(
basedir,
'events',
`${effectiveRunId}-${eventId}`,
EventSchema,
ReadEventSchema,
tag
);
if (
Expand Down Expand Up @@ -3108,7 +3112,7 @@ export function createEventsStorage(
basedir,
'events',
compositeKey,
EventSchema,
ReadEventSchema,
tag
);
if (!event) {
Expand Down Expand Up @@ -3147,7 +3151,7 @@ export function createEventsStorage(
const resolveData = params.resolveData ?? DEFAULT_RESOLVE_DATA_OPTION;
const result = await paginatedFileSystemQuery({
directory: path.join(basedir, 'events'),
schema: EventSchema,
schema: ReadEventSchema,
cachedItems: eventCache,
// Scoped to the run's own event files, since a correlation id
// identifies a step or wait only within its run: a slot-numbered
Expand Down
22 changes: 12 additions & 10 deletions packages/world-local/src/storage/helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -348,16 +348,18 @@ export function hookTokenClaimPath(basedir: string, token: string): string {
return path.join(basedir, 'hooks', 'tokens', `${hashToken(token)}.json`);
}

export const HookTokenClaimSchema = z.object({
// Legacy claims omitted hookId. Keeping it optional preserves their
// existing cross-hook conflict behavior (see #2283).
hookId: z.string().optional(),
runId: z.string(),
// Legacy claims also omitted eventId. Their recovery marker pins the
// canonical hook_created event before concurrent retries publish it.
eventId: z.string().optional(),
tokenRetentionUntil: z.coerce.date().optional(),
});
export const HookTokenClaimSchema = z.compile(
z.object({
// Legacy claims omitted hookId. Keeping it optional preserves their
// existing cross-hook conflict behavior (see #2283).
hookId: z.string().optional(),
runId: z.string(),
// Legacy claims also omitted eventId. Their recovery marker pins the
// canonical hook_created event before concurrent retries publish it.
eventId: z.string().optional(),
tokenRetentionUntil: z.coerce.date().optional(),
})
);

export type HookTokenClaim = z.infer<typeof HookTokenClaimSchema>;

Expand Down
18 changes: 11 additions & 7 deletions packages/world-local/src/storage/hook-index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,14 +39,18 @@ import { hashToken } from './helpers.js';
* (`ensureHookIndexes`) guarded by a completion marker.
*/

const IndexEntrySchema = z.object({
runId: z.string(),
});
const IndexEntrySchema = z.compile(
z.object({
runId: z.string(),
})
);

const ByRunMarkerSchema = z.object({
hookId: z.string(),
tag: z.string().optional(),
});
const ByRunMarkerSchema = z.compile(
z.object({
hookId: z.string(),
tag: z.string().optional(),
})
);

// No `.json` extension so entity listings never pick it up.
const INDEX_COMPLETE_MARKER = '.hook-index-complete';
Expand Down
41 changes: 40 additions & 1 deletion packages/world-local/src/storage/run-retention.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import path from 'node:path';
import type { WorkflowRun } from '@workflow/world';
import type { Event, WorkflowRun } from '@workflow/world';
import {
EventSchema,
getEventDataRefFields,
RETENTION_ATTRIBUTE,
readRunRetention,
Expand Down Expand Up @@ -110,6 +111,44 @@ export async function purgeRunEntityData(
*/
const RawEntitySchema = z.record(z.string(), z.any());

/**
* `EventSchema`, tolerant of a zero-retention purge having deleted an
* event's ref field (`eventData.input`/`result`/`error`/`payload`, see
* {@link getEventDataRefFields}) outright.
*
* `scrubEntityFiles` above deletes the key rather than writing it back as an
* explicit `undefined`, because JSON has no way to represent "present but
* undefined" and a schema round-trip would drop it either way. A *missing*
* key and a key *present with value `undefined`* are different things to
* Zod, though: a required field whose own schema tolerates `undefined`
* (`SerializedDataSchema`'s trailing `z.any()`) only accepts the latter.
* Every read of a stored event goes through this schema instead of the bare
* `EventSchema` so a purged run's log stays parseable, while `EventSchema`
* itself stays strict for the create-time contract (`CreateEventSchema`
* shares the same per-type schemas, and a run genuinely being created must
* still supply `input`).
*/
export const ReadEventSchema: z.ZodType<Event> = z.preprocess((raw) => {
if (
raw &&
typeof raw === 'object' &&
'eventType' in raw &&
'eventData' in raw
) {
const eventData = (raw as { eventData: unknown }).eventData;
if (eventData && typeof eventData === 'object') {
for (const field of getEventDataRefFields(
String((raw as { eventType: unknown }).eventType)
)) {
if (!(field in eventData)) {
(eventData as Record<string, unknown>)[field] = undefined;
}
}
}
}
return raw;
}, EventSchema);

/** Rewrite every file in `directory` belonging to `runId` with `scrub` applied. */
async function scrubEntityFiles(
directory: string,
Expand Down
8 changes: 5 additions & 3 deletions packages/world-local/src/streamer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,11 @@ const chunkIds = globalSingleton(
const monotonicUlid = (seedTime?: number): string => chunkIds.next(seedTime);

// Schema for the run-to-streams mapping file
const RunStreamsSchema = z.object({
streams: z.array(z.string()),
});
const RunStreamsSchema = z.compile(
z.object({
streams: z.array(z.string()),
})
);

/**
* A chunk consists of a boolean `eof` indicating if it's the last chunk,
Expand Down
26 changes: 14 additions & 12 deletions packages/world-postgres/src/message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,18 @@ import { Base64Buffer } from './zod.js';
* the body to ensure binary safety
* maybe later we can have a `blobs` table for larger payloads
*/
export const MessageData = z.object({
attempt: z.number().describe('The attempt number of the message'),
messageId: MessageId.describe('The unique ID of the message'),
idempotencyKey: z.string().optional(),
headers: z.record(z.string(), z.string()).optional(),
id: z
.string()
.describe(
"The ID of the sub-queue. For workflows, it's the workflow name. For steps, it's the step name."
),
data: Base64Buffer.describe('The message that was sent'),
});
export const MessageData = z.compile(
z.object({
attempt: z.number().describe('The attempt number of the message'),
messageId: MessageId.describe('The unique ID of the message'),
idempotencyKey: z.string().optional(),
headers: z.record(z.string(), z.string()).optional(),
id: z
.string()
.describe(
"The ID of the sub-queue. For workflows, it's the workflow name. For steps, it's the step name."
),
data: Base64Buffer.describe('The message that was sent'),
})
);
export type MessageData = z.infer<typeof MessageData>;
14 changes: 8 additions & 6 deletions packages/world-postgres/src/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,12 +55,14 @@ const graphileLogger = createGraphileLogger();
const COMPLETED_IDEMPOTENCY_CACHE_LIMIT = 10_000;
// Core records MAX_DELIVERIES_EXCEEDED on delivery 49.
const MAX_GRAPHILE_JOB_ATTEMPTS = 49;
const GraphileHelpers = z.object({
abortSignal: z.instanceof(AbortSignal).optional(),
job: z.object({
attempts: z.number().int().positive(),
}),
});
const GraphileHelpers = z.compile(
z.object({
abortSignal: z.instanceof(AbortSignal).optional(),
job: z.object({
attempts: z.number().int().positive(),
}),
})
);

type HttpExecutionResult =
| { type: 'completed' }
Expand Down
10 changes: 6 additions & 4 deletions packages/world-postgres/src/streamer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,12 @@ import * as z from 'zod';
import { type Drizzle, Schema } from './drizzle/index.js';
import { Mutex } from './util.js';

const StreamPublishMessage = z.object({
streamId: z.string(),
chunkId: z.templateLiteral(['chnk_', z.string()]),
});
const StreamPublishMessage = z.compile(
z.object({
streamId: z.string(),
chunkId: z.templateLiteral(['chnk_', z.string()]),
})
);

interface StreamChunkEvent {
id: `chnk_${string}`;
Expand Down
18 changes: 10 additions & 8 deletions packages/world-postgres/src/zod.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
import { z } from 'zod/v4';

export const Base64Buffer = z.codec(z.base64(), z.instanceof(Buffer), {
decode(b64) {
return Buffer.from(b64, 'base64');
},
encode(buf) {
return buf.toString('base64');
},
});
export const Base64Buffer = z.compile(
z.codec(z.base64(), z.instanceof(Buffer), {
decode(b64) {
return Buffer.from(b64, 'base64');
},
encode(buf) {
return buf.toString('base64');
},
})
);
Loading
Loading