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
5 changes: 4 additions & 1 deletion apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -329,7 +329,10 @@ it.effect.each(
assert.equal(wire.messages[0]?.text, task.prompt);
assert.equal(wire.messages[0]?.scheduledTaskId, task.id);
assert.equal(wire.messages[0]?.createdBy, createdBy);
const turnItem = wire.turnItems.find((item) => item.type === "user_message");
const turnItem = wire.turnItems.find(
(item): item is Extract<typeof item, { type: "user_message" }> =>
item.type === "user_message",
);
assert.equal(turnItem?.text, task.prompt);
assert.equal(turnItem?.scheduledTaskId, task.id);
}).pipe(Effect.provide(Layer.mergeAll(harness.layer, scheduledTasks)));
Expand Down
107 changes: 107 additions & 0 deletions packages/contracts/src/baseSchemas.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
import { describe, expect, it } from "@effect/vitest";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Schema from "effect/Schema";

import {
ForwardCompatibleArray,
ForwardCompatibleUnion,
ForwardCompatibleUnionArray,
isUnknownUnionMember,
UnknownUnionMember,
} from "./baseSchemas.ts";

const at = "2026-10-05T00:00:00.000Z";
const Shape = Schema.Union([
Schema.Struct({ kind: Schema.Literal("circle"), radius: Schema.Number, at: Schema.DateTimeUtc }),
Schema.Struct({ kind: Schema.Literal("square"), side: Schema.Number }),
]);
const members = Shape.members;
/** How clients decode: the JSON wire codec over the runtime schema. */
const fromWire = <S extends Schema.Top>(schema: S) =>
Schema.decodeUnknownSync(Schema.toCodecJson(schema) as never) as (value: unknown) => S["Type"];

describe("ForwardCompatibleArray", () => {
it("decodes elements through their own codec and drops ones it cannot read", () => {
const decoded = fromWire(ForwardCompatibleArray(Shape))([
{ kind: "circle", radius: 1, at },
{ kind: "triangle" },
{ kind: "square", side: "wide" },
]);
expect(decoded).toHaveLength(1);
expect(DateTime.isDateTime((decoded[0] as { at: unknown }).at)).toBe(true);
});

it("applies decoding defaults instead of failing", () => {
const WithDefault = Schema.Struct({
name: Schema.String,
count: Schema.optionalKey(Schema.Number).pipe(
Schema.withDecodingDefaultKey(Effect.succeed(7)),
),
});
expect(Schema.decodeUnknownSync(ForwardCompatibleArray(WithDefault))([{ name: "a" }])).toEqual([
{ name: "a", count: 7 },
]);
});
});

describe("ForwardCompatibleUnion", () => {
const decode = fromWire(ForwardCompatibleUnion(members, "kind"));

it("decodes a member from a newer server as unknown", () => {
const decoded = decode({ kind: "triangle", corners: 3 });
expect(isUnknownUnionMember(decoded)).toBe(true);
expect(decoded).toEqual(new UnknownUnionMember("kind", "triangle"));
});

it("decodes known members in full, transformations included", () => {
const decoded = decode({ kind: "circle", radius: 1, at });
expect(isUnknownUnionMember(decoded)).toBe(false);
expect(DateTime.isDateTime((decoded as { at: unknown }).at)).toBe(true);
});

it("still fails a known member whose payload is broken", () => {
expect(() => decode({ kind: "square", side: "wide" })).toThrow();
});

it("refuses to encode an unknown member", () => {
const encode = Schema.encodeSync(ForwardCompatibleUnion(members, "kind") as never);
expect(() => encode(new UnknownUnionMember("kind", "triangle") as never)).toThrow();
});
});

describe("ForwardCompatibleUnionArray", () => {
it("keeps a known member however its fields are named", () => {
const Flagged = Schema.Struct({
kind: Schema.Literal("flag"),
_unknown: Schema.Literal(true),
tag: Schema.String,
value: Schema.String,
});
const flag = { kind: "flag", _unknown: true, tag: "kind", value: "x" };
expect(fromWire(ForwardCompatibleUnionArray([Flagged], "kind"))([flag])).toEqual([flag]);
});

const decode = fromWire(ForwardCompatibleUnionArray(members, "kind"));

it("drops members a newer server added and keeps the rest", () => {
const decoded = decode([
{ kind: "triangle" },
{ kind: "circle", radius: 1, at },
{ kind: "square", side: 2 },
]);
expect(decoded.map((shape) => shape.kind)).toEqual(["circle", "square"]);
});

it("fails when a known member is broken, so real bugs stay visible", () => {
expect(() => decode([{ kind: "square", side: "wide" }])).toThrow();
});

it("encodes as a plain array", () => {
const encode = Schema.encodeSync(
Schema.toCodecJson(ForwardCompatibleUnionArray(members, "kind")),
);
const square = { kind: "square" as const, side: 2 };
expect(encode([square])).toEqual([square]);
});
});
175 changes: 160 additions & 15 deletions packages/contracts/src/baseSchemas.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
import * as Schema from "effect/Schema";
import * as SchemaGetter from "effect/SchemaGetter";
import * as SchemaTransformation from "effect/SchemaTransformation";

export const TrimmedString = Schema.String.pipe(
Expand Down Expand Up @@ -35,14 +36,6 @@ export type DpopFailureReason = typeof DpopFailureReason.Type;
export const IsoDateTime = Schema.String;
export type IsoDateTime = typeof IsoDateTime.Type;

/**
* Wire codec for server→client arrays whose element unions grow over time
* (new literal members, new struct variants). Decoding drops elements the
* current build cannot decode instead of failing the whole payload — a client
* has to keep decoding configs sent by servers newer than itself, and
* rejecting the payload would take down the connection over data the client
* couldn't act on anyway. Encoding is the plain array encoding.
*/
/**
* Same idea for one optional value whose literal set grows over time: a
* member this build does not know decodes as absent rather than failing the
Expand Down Expand Up @@ -107,20 +100,172 @@ export const OmittedWhenNull = <Value extends Schema.Top>(value: Value) => {
);
};

export const ForwardCompatibleArray = <Element extends Schema.Top>(element: Element) => {
const decodeElement = Schema.decodeUnknownOption(element as never);
return Schema.Array(Schema.Unknown).pipe(
/**
* Wire codec for a server→client value whose shape grows over time: a union
* a newer server may add members to, or an array of such values. Clients keep
* decoding payloads from servers newer than themselves instead of failing the
* connection over data they could not act on anyway.
*
* Decoding runs each value through its own schema, so transformations (dates,
* trimming, decoding defaults) apply as usual. Only values that schema
* rejects are dropped; encoding is the plain encoding.
*
* For a tagged union, prefer {@link ForwardCompatibleUnion}: it drops only
* values whose tag this build does not know, so a known member with a broken
* payload still fails loudly.
*/
export const ForwardCompatibleArray = <Element extends Schema.Top>(element: Element) =>
Schema.Array(
Schema.UndefinedOr(element).pipe(
// An element this build cannot read becomes a hole, filtered out below.
Schema.catchDecoding(() => Effect.succeedSome(undefined)),
),
).pipe(
Schema.decodeTo(
Schema.Array(Schema.toType(element)),
SchemaTransformation.transform<
ReadonlyArray<Element["Type"]>,
ReadonlyArray<Element["Type"] | undefined>
>({
decode: (values) => values.filter((value) => value !== undefined),
encode: (values) => values,
}),
),
) as unknown as ForwardCompatibleArray<Element>;
export type ForwardCompatibleArray<Element extends Schema.Top> = Schema.Codec<
ReadonlyArray<Element["Type"]>,
ReadonlyArray<Element["Encoded"]>,
Element["DecodingServices"],
Element["EncodingServices"]
>;

/**
* A member of a {@link ForwardCompatibleUnion} whose tag this build does not
* know. A class, so only the codec can make one: no decoded payload, however
* it is shaped, is ever mistaken for it.
*/
export class UnknownUnionMember<Tag extends string = string> {
readonly tag: Tag;
readonly value: string;
constructor(tag: Tag, value: string) {
this.tag = tag;
this.value = value;
}
}

/**
* Wire codec for a server→client tagged union a newer server may add members
* to. A value whose `tag` field holds a value this build does not know decodes
* to {@link UnknownUnionMember} instead of failing; a known member decodes
* through its own schema, so a broken known payload still fails. Callers
* decide what an unknown member means: skip it, or show a fallback.
*
* Unknown members are decode-only; encoding one fails.
*/
export const ForwardCompatibleUnion = <
const Members extends ReadonlyArray<Schema.Top & { readonly fields: object }>,
const Tag extends string,
>(
members: Members,
tag: Tag,
) => {
const known = knownTags(members, tag);
const unknownMember = Schema.Struct({
[tag]: Schema.String.check(
Schema.makeFilter(
(value: string) => !known.has(value) || `A known ${tag} must decode in full.`,
),
),
} as Record<Tag, Schema.String>).pipe(
Schema.decodeTo(Schema.instanceOf(UnknownUnionMember<Tag>), {
decode: SchemaGetter.transform(
(raw: Record<string, string>) => new UnknownUnionMember(tag, raw[tag]!),
),
encode: SchemaGetter.forbidden(() => `Unknown ${tag} values are never sent.`),
}),
);
return Schema.Union([...members, unknownMember]) as unknown as ForwardCompatibleUnion<
Members,
Tag
>;
};
export type ForwardCompatibleUnion<
Members extends ReadonlyArray<Schema.Top>,
Tag extends string,
> = Schema.Codec<
Members[number]["Type"] | UnknownUnionMember<Tag>,
Members[number]["Encoded"],
Members[number]["DecodingServices"],
Members[number]["EncodingServices"]
>;

/**
* Whether a raw, undecoded value carries a `tag` this build does not know
* among `members`. For checks outside a {@link ForwardCompatibleUnion}, such
* as an envelope that must skip a payload of an unknown kind.
*/
export const hasUnknownUnionTag = (
members: ReadonlyArray<Schema.Top & { readonly fields: object }>,
tag: string,
): ((value: unknown) => boolean) => {
const known = knownTags(members, tag);
return (value) =>
typeof value === "object" &&
value !== null &&
tag in value &&
!known.has((value as Record<string, unknown>)[tag] as string);
};

/** Whether a decoded {@link ForwardCompatibleUnion} value is a member this build does not know. */
export const isUnknownUnionMember = <A>(
value: A,
): value is Extract<A, UnknownUnionMember<string>> => value instanceof UnknownUnionMember;

/**
* An array of a growing tagged union that drops members this build does not
* know. Known members decode through their own schema and still fail when
* broken. Encoding is the plain array encoding.
*/
export const ForwardCompatibleUnionArray = <
const Members extends ReadonlyArray<Schema.Top & { readonly fields: object }>,
const Tag extends string,
>(
members: Members,
tag: Tag,
) =>
Schema.Array(ForwardCompatibleUnion(members, tag)).pipe(
Schema.decodeTo(
Schema.Array(element),
SchemaTransformation.transform<ReadonlyArray<Element["Encoded"]>, ReadonlyArray<unknown>>({
Schema.Array(Schema.toType(Schema.Union(members))),
SchemaTransformation.transform<
ReadonlyArray<Members[number]["Type"]>,
ReadonlyArray<Members[number]["Type"] | UnknownUnionMember<Tag>>
>({
decode: (values) =>
values.filter((value) => Option.isSome(decodeElement(value))) as ReadonlyArray<
Element["Encoded"]
values.filter((value) => !isUnknownUnionMember(value)) as ReadonlyArray<
Members[number]["Type"]
>,
encode: (values) => values,
}),
),
);

const knownTags = (
members: ReadonlyArray<Schema.Top & { readonly fields: object }>,
tag: string,
): ReadonlySet<string> => {
const tags = new Set<string>();
for (const member of members) {
const field = (member.fields as Record<string, unknown>)[tag] as
| { readonly literal?: unknown; readonly literals?: ReadonlyArray<unknown> }
| undefined;
for (const literal of field?.literals ?? [field?.literal]) {
if (typeof literal !== "string") {
throw new Error(`Every ForwardCompatibleUnion member needs a string literal ${tag}.`);
}
tags.add(literal);
}
}
return tags;
};

/**
Expand Down
Loading
Loading