mirror of
https://github.com/anomalyco/opencode.git
synced 2026-09-17 14:26:26 +00:00
Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c635a03165 | ||
|
|
99084d8a9b | ||
|
|
c3d278a906 | ||
|
|
e919bcc308 | ||
|
|
0880378abb |
@@ -73,6 +73,11 @@ const driver = (options: Options, body: string): WebSocketChannelDriver => {
|
||||
)
|
||||
if (event.type === "error") {
|
||||
terminal = true
|
||||
yield* OpenResponses.decodeKnownErrorEvent(event).pipe(
|
||||
Effect.mapError((cause) =>
|
||||
ProviderShared.eventError(options.id, `${options.name} returned a malformed error event`, frame, cause),
|
||||
),
|
||||
)
|
||||
return {
|
||||
type: "provider-failure",
|
||||
error: OpenResponses.providerFailure(event, `${options.name} stream error`, frame),
|
||||
|
||||
@@ -108,7 +108,7 @@ const incremental = (
|
||||
return input.slice(baseline.length)
|
||||
}
|
||||
|
||||
const code = (event: OpenResponses.Event) => OpenResponses.errorDetail(event).code
|
||||
const code = (event: OpenResponses.Event) => event.code || event.error?.code || event.response?.error?.code || undefined
|
||||
|
||||
const rejected = (
|
||||
observation: Extract<ChannelObservation, { readonly type: "provider-failure" }>,
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Effect, Option, Schema } from "effect"
|
||||
import { Effect, Option, Schema, SchemaGetter } from "effect"
|
||||
import type { Content } from "@opencode/schema/tool"
|
||||
import { HttpTransport } from "../route/transport/index.js"
|
||||
import { Protocol } from "../route/protocol.js"
|
||||
@@ -333,13 +333,53 @@ export const StreamItem = Schema.StructWithRest(
|
||||
export type StreamItem = Schema.Schema.Type<typeof StreamItem>
|
||||
export type OutputItem = StreamItem & { readonly id: string }
|
||||
|
||||
// Responses-compatible providers put error details at the top level, under `error`, or under
|
||||
// `response.error`, and gateways reshape them freely: strings, numeric codes, extra fields. Those
|
||||
// fields decode as opaque values and `errorDetail` reads them defensively, so an error frame can
|
||||
// only fail on invalid JSON and otherwise always classifies with the raw body as the fallback.
|
||||
// Responses-compatible providers put streaming error details at the top level or
|
||||
// under `error`, and response failures under `response.error`. Accept all three shapes.
|
||||
// https://www.openresponses.org/specification
|
||||
const asText = (value: unknown) =>
|
||||
typeof value === "string" && value.length > 0 ? value : typeof value === "number" ? String(value) : undefined
|
||||
const OpenResponsesErrorObject = Schema.Struct({
|
||||
type: optionalNull(Schema.String),
|
||||
code: optionalNull(Schema.String),
|
||||
message: optionalNull(Schema.String),
|
||||
param: optionalNull(Schema.String),
|
||||
})
|
||||
const OpenResponsesErrorPayload = Schema.Union([Schema.String, OpenResponsesErrorObject]).pipe(
|
||||
Schema.decodeTo(OpenResponsesErrorObject, {
|
||||
decode: SchemaGetter.transform((error) => (typeof error === "string" ? { message: error } : error)),
|
||||
encode: SchemaGetter.passthrough(),
|
||||
}),
|
||||
)
|
||||
type OpenResponsesErrorPayload = Schema.Schema.Type<typeof OpenResponsesErrorPayload>
|
||||
|
||||
const WebSocketErrorHeader = Schema.Union([Schema.String, Schema.Number, Schema.Boolean])
|
||||
export const WebSocketErrorEvent = Schema.StructWithRest(
|
||||
Schema.Struct({
|
||||
type: Schema.tag("error"),
|
||||
status: Schema.optional(Schema.Number),
|
||||
status_code: Schema.optional(Schema.Number),
|
||||
code: optionalNull(Schema.String),
|
||||
message: Schema.optional(Schema.String),
|
||||
param: optionalNull(Schema.String),
|
||||
error: optionalNull(OpenResponsesErrorPayload),
|
||||
headers: Schema.optional(Schema.Record(Schema.String, WebSocketErrorHeader)),
|
||||
}),
|
||||
[Schema.Record(Schema.String, Schema.Unknown)],
|
||||
)
|
||||
const decodeWebSocketErrorEvent = Schema.decodeUnknownEffect(WebSocketErrorEvent)
|
||||
|
||||
export const decodeKnownErrorEvent = (event: Event) =>
|
||||
decodeWebSocketErrorEvent({
|
||||
...event,
|
||||
status: typeof event.status === "number" ? event.status : undefined,
|
||||
status_code: typeof event.status_code === "number" ? event.status_code : undefined,
|
||||
headers: ProviderShared.isRecord(event.headers)
|
||||
? Object.fromEntries(
|
||||
Object.entries(event.headers).filter(
|
||||
(entry): entry is [string, string | number | boolean] =>
|
||||
typeof entry[1] === "string" || typeof entry[1] === "number" || typeof entry[1] === "boolean",
|
||||
),
|
||||
)
|
||||
: undefined,
|
||||
})
|
||||
|
||||
export const Event = Schema.StructWithRest(
|
||||
Schema.Struct({
|
||||
@@ -360,18 +400,31 @@ export const Event = Schema.StructWithRest(
|
||||
incomplete_details: optionalNull(Schema.Struct({ reason: Schema.optional(Schema.String) })),
|
||||
output: Schema.optional(Schema.Array(StreamItem)),
|
||||
usage: optionalNull(OpenResponsesUsage),
|
||||
error: Schema.optional(Schema.Unknown),
|
||||
error: optionalNull(OpenResponsesErrorPayload),
|
||||
}),
|
||||
[Schema.Record(Schema.String, Schema.Unknown)],
|
||||
),
|
||||
),
|
||||
code: Schema.optional(Schema.Unknown),
|
||||
message: Schema.optional(Schema.Unknown),
|
||||
error: Schema.optional(Schema.Unknown),
|
||||
code: optionalNull(Schema.String),
|
||||
message: Schema.optional(Schema.String),
|
||||
param: optionalNull(Schema.String),
|
||||
error: optionalNull(OpenResponsesErrorPayload),
|
||||
status: Schema.optional(Schema.Unknown),
|
||||
status_code: Schema.optional(Schema.Unknown),
|
||||
headers: Schema.optional(Schema.Unknown),
|
||||
}),
|
||||
[Schema.Record(Schema.String, Schema.Unknown)],
|
||||
).pipe(
|
||||
Schema.decode({
|
||||
decode: SchemaGetter.transform((event) => {
|
||||
if (event.type !== "error" || event.error != null) return event
|
||||
const { code, message, param, ...rest } = event
|
||||
if (code === undefined && message === undefined && param === undefined) return event
|
||||
// Flat errors (for example, Meta's) can also arrive through generic Responses endpoints.
|
||||
return { ...rest, error: { code, message, param } }
|
||||
}),
|
||||
encode: SchemaGetter.passthrough(),
|
||||
}),
|
||||
)
|
||||
export type Event = Schema.Schema.Type<typeof Event>
|
||||
export type NormalizedEvent = Event & { readonly item?: OutputItem | null }
|
||||
@@ -380,15 +433,16 @@ const decodeEventValue = Schema.decodeUnknownEffect(Event)
|
||||
const decodeFrame = Schema.decodeUnknownEffect(ProviderShared.Json)
|
||||
|
||||
/**
|
||||
* Decodes one WebSocket frame. Some providers and gateways answer a rejected `response.create` with a bare
|
||||
* `{ "error": ... }` envelope and no event type; that reads as an error event so it classifies instead of
|
||||
* failing decoding.
|
||||
* Decodes one WebSocket frame. xAI answers a rejected `response.create` with `{ "error": { "message", "type" } }` and no
|
||||
* event type; that envelope reads as an error event so the failure classifies instead of failing decoding.
|
||||
*/
|
||||
export const decodeChannelEvent = (frame: string) =>
|
||||
decodeFrame(frame).pipe(
|
||||
Effect.flatMap((value) =>
|
||||
decodeEventValue(
|
||||
ProviderShared.isRecord(value) && value.type === undefined && value.error != null
|
||||
ProviderShared.isRecord(value) &&
|
||||
value.type === undefined &&
|
||||
(typeof value.error === "string" || ProviderShared.isRecord(value.error))
|
||||
? { ...value, type: "error" }
|
||||
: value,
|
||||
),
|
||||
@@ -1368,21 +1422,22 @@ const onResponseFinish = Effect.fn("OpenResponses.onResponseFinish")(function* (
|
||||
return [{ ...current, lifecycle }, events] satisfies StepResult
|
||||
})
|
||||
|
||||
/** Error code and message from wherever the frame put them; top-level fields win over nested ones. */
|
||||
export const errorDetail = (event: Event) => {
|
||||
const raw = event.error ?? event.response?.error
|
||||
const nested = typeof raw === "string" ? { message: raw } : ProviderShared.isRecord(raw) ? raw : undefined
|
||||
return {
|
||||
message: asText(event.message) ?? asText(nested?.message),
|
||||
code: asText(event.code) ?? asText(nested?.code),
|
||||
}
|
||||
// Build the prettiest summary available from whatever the provider supplied.
|
||||
// When both code and message are present, prefix the code so consumers see
|
||||
// the failure mode (e.g. `rate_limit_exceeded: Slow down`) instead of just
|
||||
// the bare message — production rate limits and context-length failures used
|
||||
// to be indistinguishable from generic stream drops. Returns undefined when
|
||||
// the payload carries no usable summary.
|
||||
const providerErrorMessage = (event: Event, nested: OpenResponsesErrorPayload | undefined): string | undefined => {
|
||||
const message = event.message || nested?.message || undefined
|
||||
const code = event.code || nested?.code || undefined
|
||||
if (message && code) return `${code}: ${message}`
|
||||
return message || code
|
||||
}
|
||||
|
||||
// Prefix the code when both are present (`rate_limit_exceeded: Slow down`) so the failure mode is
|
||||
// visible; fall back to the raw frame rather than a generic message when neither decodes.
|
||||
export const providerFailure = (event: Event, fallback: string, body = ProviderShared.encodeJson(event)) => {
|
||||
const detail = errorDetail(event)
|
||||
const summary = detail.message && detail.code ? `${detail.code}: ${detail.message}` : (detail.message ?? detail.code)
|
||||
const nested = event.error ?? event.response?.error ?? undefined
|
||||
const summary = providerErrorMessage(event, nested)
|
||||
const message = summary ?? (body === "{}" ? fallback : body)
|
||||
const status =
|
||||
typeof event.status === "number"
|
||||
@@ -1465,7 +1520,18 @@ export const step = (state: ParserState, event: NormalizedEvent) => {
|
||||
if (event.type === "response.output_item.done") return onOutputItemDone(state, event.item)
|
||||
if (event.type === "response.completed" || event.type === "response.incomplete") return onResponseFinish(state, event)
|
||||
if (event.type === "response.failed") return providerFailure(event, `${state.name} response failed`)
|
||||
if (event.type === "error") return providerFailure(event, `${state.name} stream error`)
|
||||
if (event.type === "error")
|
||||
return decodeKnownErrorEvent(event).pipe(
|
||||
Effect.mapError((cause) =>
|
||||
ProviderShared.eventError(
|
||||
state.id,
|
||||
`${state.name} returned a malformed error event`,
|
||||
ProviderShared.encodeJson(event),
|
||||
cause,
|
||||
),
|
||||
),
|
||||
Effect.flatMap(() => providerFailure(event, `${state.name} stream error`)),
|
||||
)
|
||||
return Effect.succeed<StepResult>([state, NO_EVENTS])
|
||||
}
|
||||
|
||||
|
||||
@@ -11,78 +11,66 @@ import { sseEvents } from "../lib/sse.js"
|
||||
|
||||
const decodeEvent = Schema.decodeUnknownEffect(OpenResponses.protocol.stream.event)
|
||||
|
||||
it.effect("decodes error frames verbatim in shared SSE and WebSocket decoding", () =>
|
||||
it.effect("normalizes flat errors in shared SSE and WebSocket decoding", () =>
|
||||
Effect.gen(function* () {
|
||||
const frame = {
|
||||
type: "error",
|
||||
sequence_number: 4,
|
||||
code: "server_shutting_down",
|
||||
message: "Server is shutting down. Please retry your request.",
|
||||
param: null,
|
||||
}
|
||||
for (const decode of [decodeEvent, OpenResponses.decodeChannelEvent]) {
|
||||
for (const frame of [
|
||||
{ type: "error", sequence_number: 4, code: "server_shutting_down", message: "Shutting down", param: null },
|
||||
const event = yield* decode(JSON.stringify(frame))
|
||||
expect(event).toEqual({
|
||||
type: "error",
|
||||
sequence_number: 4,
|
||||
error: { code: frame.code, message: frame.message, param: null },
|
||||
})
|
||||
|
||||
for (const unchanged of [
|
||||
event,
|
||||
{ type: "error" },
|
||||
{ type: "error", error: "Gateway failed" },
|
||||
{ type: "error", error: { code: 429, message: "slow down" } },
|
||||
{ type: "error", error: 42 },
|
||||
{ type: "error", code: 500, message: ["not", "a", "string"] },
|
||||
{ type: "response.failed", response: { id: "resp_failed", error: "Gateway failed" } },
|
||||
{ type: "response.failed", response: { id: "resp_failed", error: ["weird"] } },
|
||||
{
|
||||
type: "response.failed",
|
||||
response: { id: "resp_failed", error: { code: "server_error", message: "Internal server error" } },
|
||||
},
|
||||
{ type: "response.output_text.delta", item_id: "msg_text", delta: "Hello" },
|
||||
]) {
|
||||
expect(yield* decode(JSON.stringify(frame))).toEqual(frame)
|
||||
expect(yield* decode(JSON.stringify(unchanged))).toEqual(unchanged)
|
||||
}
|
||||
}
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("reads bare WebSocket error envelopes as error events", () =>
|
||||
it.effect("continues to normalize untyped xAI WebSocket errors", () =>
|
||||
Effect.gen(function* () {
|
||||
const message = "gRPC error: Response with id=resp_missing not found"
|
||||
for (const error of [{ type: "api_error", message }, message, 42]) {
|
||||
expect(yield* OpenResponses.decodeChannelEvent(JSON.stringify({ error }))).toEqual({ type: "error", error })
|
||||
}
|
||||
for (const frame of [{ error: null }, { message }]) {
|
||||
expect(yield* OpenResponses.decodeChannelEvent(JSON.stringify(frame)).pipe(Effect.flip)).toBeDefined()
|
||||
for (const error of [{ type: "api_error", message }, message]) {
|
||||
expect(yield* OpenResponses.decodeChannelEvent(JSON.stringify({ error }))).toEqual({
|
||||
type: "error",
|
||||
error: typeof error === "string" ? { message } : error,
|
||||
})
|
||||
}
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("extracts error details from every shape and falls back to the raw frame", () =>
|
||||
it.effect("normalizes string errors in shared SSE and WebSocket decoding", () =>
|
||||
Effect.gen(function* () {
|
||||
const cases: Array<[frame: Record<string, unknown>, message: string, tag: string]> = [
|
||||
[
|
||||
{ type: "error", code: "server_shutting_down", message: "Shutting down" },
|
||||
"server_shutting_down: Shutting down",
|
||||
"UnknownProvider",
|
||||
],
|
||||
[{ type: "error", error: "Gateway failed" }, "Gateway failed", "UnknownProvider"],
|
||||
[{ type: "error", error: { code: 429, message: "slow down" } }, "429: slow down", "UnknownProvider"],
|
||||
[{ type: "error", error: { message: "slow down" }, status: 429 }, "slow down", "RateLimit"],
|
||||
[{ type: "error", code: 500, message: ["not", "a", "string"] }, "500", "UnknownProvider"],
|
||||
[
|
||||
{ type: "response.failed", response: { id: "resp_failed", error: "Gateway failed" } },
|
||||
"Gateway failed",
|
||||
"UnknownProvider",
|
||||
],
|
||||
]
|
||||
for (const [frame, message, tag] of cases) {
|
||||
const event = yield* OpenResponses.decodeChannelEvent(JSON.stringify(frame))
|
||||
const error = OpenResponses.providerFailure(event, "fallback", JSON.stringify(frame))
|
||||
expect(error.message).toBe(message)
|
||||
expect(error.reason._tag).toBe(tag)
|
||||
expect(error.reason.body).toBe(JSON.stringify(frame))
|
||||
for (const decode of [decodeEvent, OpenResponses.decodeChannelEvent]) {
|
||||
expect(yield* decode(JSON.stringify({ type: "error", error: "Gateway failed" }))).toEqual({
|
||||
type: "error",
|
||||
error: { message: "Gateway failed" },
|
||||
})
|
||||
expect(
|
||||
yield* decode(
|
||||
JSON.stringify({ type: "response.failed", response: { id: "resp_failed", error: "Gateway failed" } }),
|
||||
),
|
||||
).toEqual({
|
||||
type: "response.failed",
|
||||
response: { id: "resp_failed", error: { message: "Gateway failed" } },
|
||||
})
|
||||
}
|
||||
for (const frame of [
|
||||
{ type: "error", error: 42 },
|
||||
{ type: "response.failed", response: { id: "resp_failed", error: ["weird"] } },
|
||||
]) {
|
||||
const event = yield* OpenResponses.decodeChannelEvent(JSON.stringify(frame))
|
||||
const error = OpenResponses.providerFailure(event, "fallback", JSON.stringify(frame))
|
||||
expect(error.message).toBe(JSON.stringify(frame))
|
||||
expect(error.reason._tag).toBe("UnknownProvider")
|
||||
}
|
||||
expect(OpenResponses.providerFailure({ type: "error" }, "fallback", "{}").message).toBe("fallback")
|
||||
expect(OpenResponses.providerFailure({ type: "error" }, "fallback", "{}").reason._tag).toBe("ProviderInternal")
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -13,8 +13,8 @@ The idea of code mode was originally introduced by Cloudflare. See
|
||||
|
||||
## How it differs from JavaScript
|
||||
|
||||
- **Only supported APIs are available.** Programs can use the provided tools and supported JavaScript built-ins. APIs
|
||||
such as `fetch`, timers, `process`, filesystem access, imports, and modules are unavailable.
|
||||
- **Only supported APIs are available.** Programs can use the provided tools, supported JavaScript built-ins, and the
|
||||
globals the host adds through extensions. Timers, `process`, filesystem access, imports, and modules are unavailable.
|
||||
- **Unfinished work is interrupted.** Tool calls and async functions start when called. When the program finishes,
|
||||
anything still running is interrupted. Unhandled rejections from un-awaited promises are returned as warnings.
|
||||
- **REPL-style results.** Without an explicit `return`, the final top-level expression becomes the result. `undefined`
|
||||
@@ -94,11 +94,19 @@ receive `{ extension, name, args }`. An `after` hook also receives how the call
|
||||
`failure` with its error, or `interrupted`). A failing `before` hook denies the call, and the program catches the
|
||||
failure as a thrown error.
|
||||
|
||||
### `Values`
|
||||
### `Extension.make`
|
||||
|
||||
`Values` exports the runtime's non-JSON value classes: `Values.URL`, `Values.URLSearchParams`, `Values.Date`,
|
||||
`Values.RegExp`, `Values.Map`, `Values.Set`, and `Values.Promise`. The interpreter recognizes these by class; a
|
||||
program's `new URL(...)` is a `Values.URL` wrapping the host `URL`. `Values.isValue` narrows to the data-like kinds.
|
||||
Extensions are host functions a program calls directly as globals, such as `fetch`. Unlike tools they are not in the
|
||||
catalog, not counted against `maxToolCalls`, and not described to the model; the host decides what they mean.
|
||||
|
||||
```ts
|
||||
const web = Extension.make({ name: "web", globals: { fetch: (url: string) => globalThis.fetch(url) } })
|
||||
const runtime = CodeMode.make({ tools, extensions: [web] })
|
||||
```
|
||||
|
||||
Every value crossing in either direction is converted, never shared: arguments come in as copies, results go out as
|
||||
copies, and a function inside a result is callable the same way. A global that shadows a built-in or another
|
||||
extension throws at `CodeMode.make`.
|
||||
|
||||
### OpenAPI tools
|
||||
|
||||
|
||||
@@ -29,7 +29,7 @@ ultimate source of truth. Upstream test262 files run verbatim from `test/test262
|
||||
Uint8Array is rejected with a hint to encode as text, and own `__proto__` keys are dropped so merging tool
|
||||
inputs or results cannot replace a prototype. In-program `JSON.stringify` keeps JS behavior except for the
|
||||
Error form and a promise, which is a `TypeError` with an await hint rather than a silent `{}`.
|
||||
- [x] Live Date, RegExp, Map, Set, URL, URLSearchParams, and Uint8Array values inside CodeMode.
|
||||
- [x] Live Date, RegExp, Map, Set, URL, URLSearchParams, Headers, and Uint8Array values inside CodeMode.
|
||||
- [x] Tool calls through the host-provided `tools` tree only.
|
||||
- [x] The global `search(...)` built-in: synchronous tool discovery that counts as an admitted tool call and is
|
||||
shadowable by program declarations like other globals.
|
||||
@@ -47,8 +47,8 @@ ultimate source of truth. Upstream test262 files run verbatim from `test/test262
|
||||
## Values and literals
|
||||
|
||||
- [x] `null`, `undefined`, booleans, finite and non-finite numbers, and strings.
|
||||
- [x] Array literals, including holes and spread from arrays, strings, Maps, Sets, URLSearchParams, custom synchronous
|
||||
iterators, and synchronous generators.
|
||||
- [x] Array literals, including holes and spread from arrays, strings, Maps, Sets, URLSearchParams, Headers, custom
|
||||
synchronous iterators, and synchronous generators.
|
||||
- [x] Object literals with shorthand, computed string/number keys, and spread following ToObject: data objects and
|
||||
arrays copy own enumerable keys, strings copy index keys, and other values contribute nothing.
|
||||
- [x] Template literals with interpolation.
|
||||
@@ -95,8 +95,8 @@ ultimate source of truth. Upstream test262 files run verbatim from `test/test262
|
||||
- [x] `if`/`else` and conditional expressions.
|
||||
- [x] `switch`, including default clauses and fallthrough.
|
||||
- [x] `for`, `while`, and `do...while`.
|
||||
- [x] `for...of` over arrays, strings, Maps, Sets, URLSearchParams, custom synchronous iterators, and confined
|
||||
synchronous generators. Abrupt completion invokes the iterator's optional `return()`.
|
||||
- [x] `for...of` over arrays, strings, Maps, Sets, URLSearchParams, Headers, custom synchronous iterators, and
|
||||
confined synchronous generators. Abrupt completion invokes the iterator's optional `return()`.
|
||||
- [x] `for...in` over own keys of plain objects, arrays, strings, and tool references; other values iterate nothing.
|
||||
- [x] Unlabeled `break` and `continue`.
|
||||
- [x] `try`, `catch`, optional catch bindings, and `finally`.
|
||||
@@ -127,7 +127,7 @@ ultimate source of truth. Upstream test262 files run verbatim from `test/test262
|
||||
string). A detached method loses its receiver, as in JS: `values.filter("abc".includes)` is a `TypeError`
|
||||
because `includes` is called without a string `this`.
|
||||
- [x] Constructors work as callbacks with JS call semantics: `Error` types construct (`messages.map(Error)`),
|
||||
and new-requiring constructors (`Map`, `Set`, `URL`, `URLSearchParams`, `Promise`) throw a `TypeError`,
|
||||
and new-requiring constructors (`Map`, `Set`, `URL`, `URLSearchParams`, `Headers`, `Promise`) throw a `TypeError`,
|
||||
like JS.
|
||||
- [x] Tool references and detached `Promise` statics are rejected as callbacks with a hint to wrap them in an
|
||||
arrow function.
|
||||
@@ -179,10 +179,10 @@ ultimate source of truth. Upstream test262 files run verbatim from `test/test262
|
||||
- [x] Sequence expressions (the comma operator).
|
||||
- [x] `await` for CodeMode promises and callable thenables; a plain value passes through unchanged, though every
|
||||
`await` still defers its continuation one reaction turn.
|
||||
- [x] `new` for Array, Object, Error types, Date, RegExp, Map, Set, URL, URLSearchParams, and Promise. `new` on any
|
||||
other value throws a catchable `TypeError` naming the callee: other built-in functions such as `Number` say
|
||||
`new` is unsupported and point at the plain call, user-defined functions report the constructor gap below, and
|
||||
non-callable values are not constructors.
|
||||
- [x] `new` for Array, Object, Error types, Date, RegExp, Map, Set, URL, URLSearchParams, Headers, and Promise. `new`
|
||||
on any other value throws a catchable `TypeError` naming the callee: other built-in functions such as `Number`
|
||||
say `new` is unsupported and point at the plain call, user-defined functions report the constructor gap below,
|
||||
and non-callable values are not constructors.
|
||||
- [x] Arithmetic operators: `+`, `-`, `*`, `/`, `%`, and `**`.
|
||||
- [x] Equality and ordering: `==`, `!=`, `===`, `!==`, `<`, `<=`, `>`, and `>=`.
|
||||
- [x] Bitwise operators: `&`, `|`, `^`, `~`, `<<`, `>>`, and `>>>`.
|
||||
@@ -450,7 +450,13 @@ with a hint to encode as text first (`TextDecoder`, `toBase64`, `toHex`).
|
||||
- [x] `crypto.randomUUID()` and `crypto.getRandomValues(uint8Array)`.
|
||||
- [x] `TextEncoder` and `TextDecoder` for UTF-8 only: any other label is a `RangeError`. `TextDecoder` accepts the
|
||||
`fatal` and `ignoreBOM` options; `decode` takes a Uint8Array or nothing.
|
||||
- [ ] `crypto.subtle`, `Blob`, and `TextDecoder` streaming or non-UTF-8 encodings.
|
||||
- [x] `new Headers()` from records, synchronous iterables of pairs, and Headers, wrapping the host's `Headers`: names
|
||||
fold to lowercase, values are normalized and combined, and invalid names or values throw a `TypeError`.
|
||||
- [x] Headers `append`, `delete`, `get`, `getSetCookie`, `has`, `set`, `forEach`, `keys`, `values`, and `entries`;
|
||||
iteration is live and sorted by name, with `set-cookie` values kept apart.
|
||||
- [x] Headers serialize to a `{ name: value }` object in JSON, in results, and in tool arguments.
|
||||
- [ ] `Request`, `Response`, and `Blob`.
|
||||
- [ ] `crypto.subtle` and `TextDecoder` streaming or non-UTF-8 encodings.
|
||||
|
||||
## Extensions
|
||||
|
||||
@@ -460,8 +466,8 @@ Nothing is exposed unless a host provides it; extension calls are not tool calls
|
||||
- [x] Each global is a function, callable but not constructible, run with `this` undefined. A global that shadows
|
||||
a built-in or another extension throws at `make`.
|
||||
- [x] Every value crossing in either direction is converted, never shared: plain objects and arrays are copied,
|
||||
`Date`, `RegExp`, `URL`, `URLSearchParams`, `Map`, `Set`, and `Uint8Array` become fresh copies with their
|
||||
contents converted (a host `ArrayBuffer` comes in as a `Uint8Array`; other typed arrays cannot come out),
|
||||
`Date`, `RegExp`, `URL`, `URLSearchParams`, `Headers`, `Map`, `Set`, and `Uint8Array` become fresh copies with
|
||||
their contents converted (a host `ArrayBuffer` comes in as a `Uint8Array`; other typed arrays cannot come out),
|
||||
errors cross as errors with their name and message, and a `__proto__` key is dropped. Functions, generators,
|
||||
un-awaited promises, and symbols cannot be passed in; a class instance, a symbol, or a BigInt cannot come out.
|
||||
- [x] A host function inside a result becomes a program function whose calls cross the same way, so a result can
|
||||
|
||||
@@ -17,6 +17,7 @@ import {
|
||||
record,
|
||||
SetObj,
|
||||
URLSearchParamsObj,
|
||||
HeadersObj,
|
||||
} from "./interpreter/objects.js"
|
||||
import { typeofValue } from "./interpreter/references.js"
|
||||
|
||||
@@ -69,6 +70,7 @@ const walk = <R>(
|
||||
)
|
||||
}
|
||||
if (boundary && value instanceof URLSearchParamsObj) return value.params.toString()
|
||||
if (value instanceof HeadersObj) return Object.fromEntries(value.headers)
|
||||
const target = boundary && value instanceof SetObj ? new Arr(ctx.builtins.Array, [...value.set]) : value
|
||||
if (stack.has(target)) throw typeError("Converting circular structure to JSON.")
|
||||
stack.add(target)
|
||||
|
||||
@@ -24,6 +24,7 @@ import {
|
||||
SetObj,
|
||||
URLObj,
|
||||
URLSearchParamsObj,
|
||||
HeadersObj,
|
||||
} from "./objects.js"
|
||||
import { describeValue } from "./references.js"
|
||||
|
||||
@@ -49,6 +50,7 @@ export const extensionGlobals = <R>(
|
||||
if (value instanceof RegExpObj) return new RegExp(value.regex.source, value.regex.flags)
|
||||
if (value instanceof URLObj) return new URL(value.url.href)
|
||||
if (value instanceof URLSearchParamsObj) return new URLSearchParams(value.params)
|
||||
if (value instanceof HeadersObj) return new Headers(value.headers)
|
||||
const next = (item: unknown) => toHost(item, label, depth + 1, seen)
|
||||
if (value instanceof MapObj) return new Map([...value.map].map(([key, item]) => [next(key), next(item)]))
|
||||
if (value instanceof SetObj) return new Set([...value.set].map(next))
|
||||
@@ -96,6 +98,7 @@ export const extensionGlobals = <R>(
|
||||
if (value instanceof URLSearchParams) {
|
||||
return new URLSearchParamsObj(builtins.URLSearchParams, new URLSearchParams(value))
|
||||
}
|
||||
if (value instanceof Headers) return new HeadersObj(builtins.Headers, new Headers(value))
|
||||
const next = (item: unknown, path: string) => fromHost(item, path, depth + 1, seen)
|
||||
if (value instanceof Map) {
|
||||
const wrapped = new MapObj(builtins.Map)
|
||||
|
||||
@@ -11,6 +11,7 @@ import { objectGlobal } from "../stdlib/object.js"
|
||||
import { regexpGlobal } from "../stdlib/regexp.js"
|
||||
import { stringGlobal } from "../stdlib/string.js"
|
||||
import { uriGlobal, urlGlobal, urlSearchParamsGlobal } from "../stdlib/url.js"
|
||||
import { headersGlobal } from "../stdlib/headers.js"
|
||||
import { coercion } from "../stdlib/value.js"
|
||||
import { base64Global, cryptoGlobal } from "../stdlib/web.js"
|
||||
import { ToolReference } from "../tool-runtime.js"
|
||||
@@ -80,6 +81,7 @@ const table: Record<string, Factory> = {
|
||||
Set: (ctx) => setGlobal(ctx),
|
||||
URL: (ctx) => urlGlobal(ctx),
|
||||
URLSearchParams: (ctx) => urlSearchParamsGlobal(ctx),
|
||||
Headers: (ctx) => headersGlobal(ctx),
|
||||
Uint8Array: (ctx) => uint8ArrayGlobal(ctx),
|
||||
TextEncoder: (ctx) => textEncoderGlobal(ctx),
|
||||
TextDecoder: (ctx) => textDecoderGlobal(ctx),
|
||||
|
||||
@@ -84,6 +84,7 @@ import {
|
||||
PromiseObj,
|
||||
SetObj,
|
||||
URLSearchParamsObj,
|
||||
HeadersObj,
|
||||
record,
|
||||
remove,
|
||||
set,
|
||||
@@ -653,7 +654,7 @@ class Frame<R> {
|
||||
const cursor = iterator === undefined ? yield* self.iterate(right, node) : undefined
|
||||
if (iterator === undefined && cursor === undefined) {
|
||||
throw invalidData(
|
||||
`${awaiting ? "for await...of" : "for...of"} requires an array, string, Map, Set, or URLSearchParams, or custom iterator value.`,
|
||||
`${awaiting ? "for await...of" : "for...of"} requires an array, string, Map, Set, URLSearchParams, or Headers, or custom iterator value.`,
|
||||
node,
|
||||
)
|
||||
}
|
||||
@@ -756,9 +757,11 @@ class Frame<R> {
|
||||
? value.set.values()
|
||||
: value instanceof URLSearchParamsObj
|
||||
? value.params.entries()
|
||||
: value instanceof Bytes
|
||||
? value.bytes.values()
|
||||
: undefined
|
||||
: value instanceof HeadersObj
|
||||
? value.headers.entries()
|
||||
: value instanceof Bytes
|
||||
? value.bytes.values()
|
||||
: undefined
|
||||
if (iterator !== undefined) {
|
||||
const proto = this.ctx.builtins.Array
|
||||
return Effect.succeed({
|
||||
@@ -1848,6 +1851,7 @@ class Frame<R> {
|
||||
value instanceof MapObj ||
|
||||
value instanceof SetObj ||
|
||||
value instanceof URLSearchParamsObj ||
|
||||
value instanceof HeadersObj ||
|
||||
value instanceof Bytes
|
||||
) {
|
||||
const cursor = yield* self.iterate(value, node)
|
||||
|
||||
@@ -29,6 +29,7 @@ const builtins = [
|
||||
"Set",
|
||||
"URL",
|
||||
"URLSearchParams",
|
||||
"Headers",
|
||||
"Uint8Array",
|
||||
"TextEncoder",
|
||||
"TextDecoder",
|
||||
@@ -80,6 +81,7 @@ export const createBuiltins = (): Builtins => {
|
||||
Set: plain(),
|
||||
URL: plain(),
|
||||
URLSearchParams: plain(),
|
||||
Headers: plain(),
|
||||
Uint8Array: plain(),
|
||||
TextEncoder: plain(),
|
||||
TextDecoder: plain(),
|
||||
|
||||
@@ -156,6 +156,15 @@ export class URLSearchParamsObj extends Obj {
|
||||
}
|
||||
}
|
||||
|
||||
export class HeadersObj extends Obj {
|
||||
constructor(
|
||||
proto: Obj,
|
||||
readonly headers: Headers,
|
||||
) {
|
||||
super(proto)
|
||||
}
|
||||
}
|
||||
|
||||
export class URLObj extends Obj {
|
||||
readonly searchParams: URLSearchParamsObj
|
||||
constructor(
|
||||
@@ -181,13 +190,14 @@ export class Bytes extends Obj {
|
||||
/** Built-in objects that wrap a host value; data-like, but never plain data. */
|
||||
export const isWrapper = (
|
||||
value: unknown,
|
||||
): value is DateObj | RegExpObj | MapObj | SetObj | URLObj | URLSearchParamsObj | Bytes =>
|
||||
): value is DateObj | RegExpObj | MapObj | SetObj | URLObj | URLSearchParamsObj | HeadersObj | Bytes =>
|
||||
value instanceof DateObj ||
|
||||
value instanceof RegExpObj ||
|
||||
value instanceof MapObj ||
|
||||
value instanceof SetObj ||
|
||||
value instanceof URLObj ||
|
||||
value instanceof URLSearchParamsObj ||
|
||||
value instanceof HeadersObj ||
|
||||
value instanceof Bytes
|
||||
|
||||
const MAX_ARRAY_INDEX = 4_294_967_295
|
||||
|
||||
@@ -16,6 +16,7 @@ import {
|
||||
SetObj,
|
||||
URLObj,
|
||||
URLSearchParamsObj,
|
||||
HeadersObj,
|
||||
} from "./objects.js"
|
||||
|
||||
/** Values that cannot cross the data boundary. */
|
||||
@@ -85,6 +86,7 @@ export const describeValue = (value: unknown): string => {
|
||||
if (value instanceof SetObj) return "a Set"
|
||||
if (value instanceof URLObj) return "a URL"
|
||||
if (value instanceof URLSearchParamsObj) return "a URLSearchParams"
|
||||
if (value instanceof HeadersObj) return "a Headers"
|
||||
if (value instanceof Bytes) return "a Uint8Array"
|
||||
if (value instanceof GeneratorObj) return "a generator"
|
||||
if (isRuntimeReference(value)) return "a function"
|
||||
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
SetObj,
|
||||
URLObj,
|
||||
URLSearchParamsObj,
|
||||
HeadersObj,
|
||||
} from "../interpreter/objects.js"
|
||||
import { containsOpaqueReference, isRuntimeReference } from "../interpreter/references.js"
|
||||
import type { Interpreter } from "../interpreter/interpreter.js"
|
||||
@@ -66,6 +67,7 @@ const formatConsoleValue = (value: unknown, seen: Set<object>, depth: number): s
|
||||
if (value instanceof RegExpObj) return coerceToString(value)
|
||||
if (value instanceof URLObj) return coerceToString(value)
|
||||
if (value instanceof URLSearchParamsObj) return coerceToString(value)
|
||||
if (value instanceof HeadersObj) return `Headers ${JSON.stringify(Object.fromEntries(value.headers))}`
|
||||
if (value instanceof Bytes) return `Uint8Array(${value.bytes.length}) [${value.bytes.join(",")}]`
|
||||
if (depth > MAX_CONSOLE_DEPTH) return "..."
|
||||
if (seen.has(value)) return "[Circular]"
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
import { Effect } from "effect"
|
||||
import { constructor, methods, prototypeFrom, receiver, requiresNew } from "../interpreter/native.js"
|
||||
import { typeError } from "../interpreter/model.js"
|
||||
import { entries, Arr, HeadersObj, Obj } from "../interpreter/objects.js"
|
||||
import { applyCollectionCallback } from "../interpreter/callback.js"
|
||||
import { isRuntimeReference } from "../interpreter/references.js"
|
||||
import type { Interpreter } from "../interpreter/interpreter.js"
|
||||
import { coerceToString } from "./value.js"
|
||||
import { readPairs } from "./url.js"
|
||||
|
||||
// The host validates header names and values and throws its own TypeError; the program gets one of its own.
|
||||
const attempt = <T>(run: () => T): T => {
|
||||
try {
|
||||
return run()
|
||||
} catch (error) {
|
||||
throw typeError(error instanceof Error ? error.message : String(error))
|
||||
}
|
||||
}
|
||||
|
||||
const constructHeaders = <R>(ctx: Interpreter<R>, init: unknown, proto: Obj): Effect.Effect<HeadersObj, unknown, R> => {
|
||||
const wrap = (headers: Headers) => new HeadersObj(proto, headers)
|
||||
if (init === undefined) return Effect.succeed(wrap(new Headers()))
|
||||
return Effect.gen(function* () {
|
||||
const pairs = init instanceof Obj ? yield* readPairs(ctx, init, "new Headers(...)") : undefined
|
||||
if (pairs !== undefined) return wrap(attempt(() => new Headers(pairs)))
|
||||
if (!(init instanceof Obj) || isRuntimeReference(init)) {
|
||||
throw typeError("new Headers(...) expects a record of names to values, iterable [name, value] pairs, or Headers.")
|
||||
}
|
||||
return wrap(
|
||||
attempt(() => new Headers(Object.fromEntries(entries(init).map(([key, value]) => [key, coerceToString(value)])))),
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
export const headersGlobal = <R>(ctx: Interpreter<R>) => {
|
||||
const builtins = ctx.builtins
|
||||
const proto = builtins.Headers
|
||||
const headers = constructor<R>(builtins, proto, {
|
||||
name: "Headers",
|
||||
call: requiresNew("Headers"),
|
||||
construct: (args, newTarget) => constructHeaders(ctx, args[0], prototypeFrom(newTarget, proto)),
|
||||
})
|
||||
const self = (thisValue: unknown, name: string) => receiver(HeadersObj, thisValue, `Headers.prototype.${name}`)
|
||||
const wrap = (items: Array<unknown>) => new Arr(builtins.Array, items)
|
||||
const arg = (args: Array<unknown>, index: number): string => coerceToString(args[index])
|
||||
const requireArgs = (name: string, args: Array<unknown>, count: number): void => {
|
||||
if (args.length < count) throw typeError(`Headers.${name} requires ${count} argument${count === 1 ? "" : "s"}.`)
|
||||
}
|
||||
methods(builtins, proto, [
|
||||
[
|
||||
"append",
|
||||
2,
|
||||
(thisValue, args) => {
|
||||
requireArgs("append", args, 2)
|
||||
const target = self(thisValue, "append").headers
|
||||
return attempt(() => target.append(arg(args, 0), arg(args, 1)))
|
||||
},
|
||||
],
|
||||
[
|
||||
"delete",
|
||||
1,
|
||||
(thisValue, args) => {
|
||||
requireArgs("delete", args, 1)
|
||||
const target = self(thisValue, "delete").headers
|
||||
return attempt(() => target.delete(arg(args, 0)))
|
||||
},
|
||||
],
|
||||
[
|
||||
"get",
|
||||
1,
|
||||
(thisValue, args) => {
|
||||
requireArgs("get", args, 1)
|
||||
const target = self(thisValue, "get").headers
|
||||
return attempt(() => target.get(arg(args, 0)))
|
||||
},
|
||||
],
|
||||
["getSetCookie", 0, (thisValue) => wrap(self(thisValue, "getSetCookie").headers.getSetCookie())],
|
||||
[
|
||||
"has",
|
||||
1,
|
||||
(thisValue, args) => {
|
||||
requireArgs("has", args, 1)
|
||||
const target = self(thisValue, "has").headers
|
||||
return attempt(() => target.has(arg(args, 0)))
|
||||
},
|
||||
],
|
||||
[
|
||||
"set",
|
||||
2,
|
||||
(thisValue, args) => {
|
||||
requireArgs("set", args, 2)
|
||||
const target = self(thisValue, "set").headers
|
||||
return attempt(() => target.set(arg(args, 0), arg(args, 1)))
|
||||
},
|
||||
],
|
||||
["keys", 0, (thisValue) => wrap(Array.from(self(thisValue, "keys").headers.keys()))],
|
||||
["values", 0, (thisValue) => wrap(Array.from(self(thisValue, "values").headers.values()))],
|
||||
[
|
||||
"entries",
|
||||
0,
|
||||
(thisValue) =>
|
||||
wrap(Array.from(self(thisValue, "entries").headers.entries(), ([key, value]) => wrap([key, value]))),
|
||||
],
|
||||
[
|
||||
"forEach",
|
||||
1,
|
||||
(thisValue, args) => {
|
||||
requireArgs("forEach", args, 1)
|
||||
const target = self(thisValue, "forEach")
|
||||
const apply = applyCollectionCallback(ctx, args[0], "Headers.forEach")
|
||||
return Effect.gen(function* () {
|
||||
for (const [key, value] of Array.from(target.headers.entries())) yield* apply([value, key, target])
|
||||
return undefined
|
||||
})
|
||||
},
|
||||
],
|
||||
])
|
||||
return headers
|
||||
}
|
||||
@@ -107,12 +107,10 @@ export const urlGlobal = <R>(ctx: Interpreter<R>) => {
|
||||
return url
|
||||
}
|
||||
|
||||
const readPair = <R>(ctx: Interpreter<R>, value: unknown): Effect.Effect<Array<string>, unknown, R> =>
|
||||
const readPair = <R>(ctx: Interpreter<R>, value: unknown, label: string): Effect.Effect<Array<string>, unknown, R> =>
|
||||
Effect.gen(function* () {
|
||||
const cursor = yield* ctx.iterate(value)
|
||||
if (cursor === undefined) {
|
||||
throw typeError("new URLSearchParams(...) expects iterable [name, value] pairs.")
|
||||
}
|
||||
if (cursor === undefined) throw typeError(`${label} expects iterable [name, value] pairs.`)
|
||||
const items: Array<string> = []
|
||||
while (true) {
|
||||
const step = yield* cursor.next
|
||||
@@ -126,6 +124,29 @@ const readPair = <R>(ctx: Interpreter<R>, value: unknown): Effect.Effect<Array<s
|
||||
}
|
||||
})
|
||||
|
||||
/**
|
||||
* Reads a synchronous iterable of `[name, value]` pairs as strings; `undefined` when `init` is not iterable. As in
|
||||
* WebIDL, the whole sequence is converted before any pair's length is checked.
|
||||
*/
|
||||
export const readPairs = <R>(
|
||||
ctx: Interpreter<R>,
|
||||
init: unknown,
|
||||
label: string,
|
||||
): Effect.Effect<Array<[string, string]> | undefined, unknown, R> =>
|
||||
Effect.gen(function* () {
|
||||
const cursor = yield* ctx.iterate(init)
|
||||
if (cursor === undefined) return undefined
|
||||
const pairs: Array<Array<string>> = []
|
||||
while (true) {
|
||||
const step = yield* cursor.next
|
||||
if (step.done) {
|
||||
if (pairs.some((entry) => entry.length !== 2)) throw typeError(`${label} expects iterable [name, value] pairs.`)
|
||||
return pairs as Array<[string, string]>
|
||||
}
|
||||
pairs.push(yield* preserveConsumerError(cursor, readPair(ctx, step.value, label)))
|
||||
}
|
||||
})
|
||||
|
||||
const constructURLSearchParams = <R>(
|
||||
ctx: Interpreter<R>,
|
||||
init: unknown,
|
||||
@@ -139,20 +160,8 @@ const constructURLSearchParams = <R>(
|
||||
return Effect.succeed(wrap(new URLSearchParams(coerceToString(init))))
|
||||
}
|
||||
return Effect.gen(function* () {
|
||||
const cursor = yield* ctx.iterate(init)
|
||||
if (cursor !== undefined) {
|
||||
const pairs: Array<Array<string>> = []
|
||||
while (true) {
|
||||
const step = yield* cursor.next
|
||||
if (step.done) {
|
||||
if (pairs.some((entry) => entry.length !== 2)) {
|
||||
throw typeError("new URLSearchParams(...) expects iterable [name, value] pairs.")
|
||||
}
|
||||
return wrap(new URLSearchParams(pairs.map((entry): [string, string] => [entry[0] ?? "", entry[1] ?? ""])))
|
||||
}
|
||||
pairs.push(yield* preserveConsumerError(cursor, readPair(ctx, step.value)))
|
||||
}
|
||||
}
|
||||
const pairs = yield* readPairs(ctx, init, "new URLSearchParams(...)")
|
||||
if (pairs !== undefined) return wrap(new URLSearchParams(pairs))
|
||||
if (isRuntimeReference(init)) {
|
||||
throw typeError("new URLSearchParams(...) expects a query string, data object, or synchronous iterable pairs.")
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import {
|
||||
SetObj,
|
||||
URLObj,
|
||||
URLSearchParamsObj,
|
||||
HeadersObj,
|
||||
} from "../interpreter/objects.js"
|
||||
import type { Interpreter } from "../interpreter/interpreter.js"
|
||||
|
||||
@@ -28,6 +29,7 @@ export const coerceToString = (value: unknown): string => {
|
||||
if (value instanceof SetObj) return "[object Set]"
|
||||
if (value instanceof URLObj) return value.url.href
|
||||
if (value instanceof URLSearchParamsObj) return value.params.toString()
|
||||
if (value instanceof HeadersObj) return "[object Headers]"
|
||||
if (value instanceof Bytes) return value.bytes.join(",")
|
||||
if (value instanceof ErrorObj) {
|
||||
// Match Error.prototype.toString: "name: message", or just one when the other is empty.
|
||||
|
||||
@@ -128,6 +128,34 @@ describe("values are converted at the boundary, never shared", () => {
|
||||
expect([...(held[0] as Set<{ z: number }>)][0]).toEqual({ z: 1 })
|
||||
})
|
||||
|
||||
test("Headers cross as copies in both directions", async () => {
|
||||
const stored = new Headers({ "X-A": "1" })
|
||||
const target = CodeMode.make({
|
||||
extensions: [
|
||||
Extension.make({
|
||||
name: "http",
|
||||
globals: {
|
||||
headers: () => stored,
|
||||
keep: (value: Headers) => {
|
||||
held.push(value)
|
||||
return value
|
||||
},
|
||||
},
|
||||
}),
|
||||
],
|
||||
})
|
||||
held.length = 0
|
||||
expect(
|
||||
await value(
|
||||
`const h = headers(); h.set("x-a", "2"); const back = keep(h); back.set("x-a", "3"); return [h instanceof Headers, h.get("x-a"), back === h, back.get("x-a"), [...back]]`,
|
||||
target,
|
||||
),
|
||||
).toEqual([true, "2", false, "3", [["x-a", "3"]]])
|
||||
expect(stored.get("x-a")).toBe("1")
|
||||
expect(held[0]).toBeInstanceOf(Headers)
|
||||
expect((held[0] as Headers).get("x-a")).toBe("2")
|
||||
})
|
||||
|
||||
test("bytes cross as copies in both directions; ArrayBuffer comes in as Uint8Array", async () => {
|
||||
const stored = new Uint8Array([1, 2, 3])
|
||||
const target = CodeMode.make({
|
||||
|
||||
@@ -635,6 +635,154 @@ describe("URL and URI helpers", () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe("Headers", () => {
|
||||
test("constructs from records, pairs, Maps, and Headers; names fold to lowercase and values combine", async () => {
|
||||
expect(
|
||||
await value(`
|
||||
const headers = new Headers({ "Content-Type": "text/plain", "X-Count": 1, "X-Null": null })
|
||||
headers.append("Accept", "text/html")
|
||||
headers.append("accept", "application/json")
|
||||
headers.set("x-count", "2")
|
||||
headers.delete("x-null")
|
||||
const copy = new Headers(headers)
|
||||
copy.set("content-type", "text/html")
|
||||
return {
|
||||
get: headers.get("content-type"),
|
||||
missing: headers.get("x-missing"),
|
||||
combined: headers.get("ACCEPT"),
|
||||
has: [headers.has("Accept"), headers.has("x-null")],
|
||||
count: headers.get("x-count"),
|
||||
copied: [headers.get("content-type"), copy.get("content-type")],
|
||||
pairs: [...new Headers([["b", "2"], ["A", "1"]])],
|
||||
map: [...new Headers(new Map([["k", "v"]]))],
|
||||
keys: headers.keys(),
|
||||
values: headers.values(),
|
||||
entries: headers.entries(),
|
||||
}
|
||||
`),
|
||||
).toEqual({
|
||||
get: "text/plain",
|
||||
missing: null,
|
||||
combined: "text/html, application/json",
|
||||
has: [true, false],
|
||||
count: "2",
|
||||
copied: ["text/plain", "text/html"],
|
||||
pairs: [
|
||||
["a", "1"],
|
||||
["b", "2"],
|
||||
],
|
||||
map: [["k", "v"]],
|
||||
keys: ["accept", "content-type", "x-count"],
|
||||
values: ["text/html, application/json", "text/plain", "2"],
|
||||
entries: [
|
||||
["accept", "text/html, application/json"],
|
||||
["content-type", "text/plain"],
|
||||
["x-count", "2"],
|
||||
],
|
||||
})
|
||||
})
|
||||
|
||||
test("iterates in sorted order everywhere iteration is allowed, and getSetCookie keeps cookies apart", async () => {
|
||||
expect(
|
||||
await value(`
|
||||
const headers = new Headers({ b: "2", a: "1" })
|
||||
headers.append("Set-Cookie", "x=1")
|
||||
headers.append("set-cookie", "y=2")
|
||||
const seen = []
|
||||
headers.forEach((value, name, self) => seen.push(name + "=" + value + ":" + (self === headers)))
|
||||
const [first] = headers
|
||||
function* pairs() { yield* headers }
|
||||
return {
|
||||
seen,
|
||||
first,
|
||||
spread: [...headers],
|
||||
from: Array.from(headers).length,
|
||||
generator: [...pairs()].length,
|
||||
object: Object.fromEntries(headers),
|
||||
cookies: headers.getSetCookie(),
|
||||
}
|
||||
`),
|
||||
).toEqual({
|
||||
seen: ["a=1:true", "b=2:true", "set-cookie=x=1:true", "set-cookie=y=2:true"],
|
||||
first: ["a", "1"],
|
||||
spread: [
|
||||
["a", "1"],
|
||||
["b", "2"],
|
||||
["set-cookie", "x=1"],
|
||||
["set-cookie", "y=2"],
|
||||
],
|
||||
from: 4,
|
||||
generator: 4,
|
||||
object: { a: "1", b: "2", "set-cookie": "y=2" },
|
||||
cookies: ["x=1", "y=2"],
|
||||
})
|
||||
})
|
||||
|
||||
test("serializes as a name-to-value object at the boundary and in JSON; prints for console", async () => {
|
||||
const result = await run(`
|
||||
const headers = new Headers({ "X-A": "1", b: "2" })
|
||||
console.log(headers)
|
||||
return { headers, json: JSON.stringify({ headers }), text: String(headers), type: typeof headers, is: headers instanceof Headers }
|
||||
`)
|
||||
expect(result.ok && result.value).toEqual({
|
||||
headers: { b: "2", "x-a": "1" },
|
||||
json: '{"headers":{"b":"2","x-a":"1"}}',
|
||||
text: "[object Headers]",
|
||||
type: "object",
|
||||
is: true,
|
||||
})
|
||||
expect(result.ok && result.logs?.[0]).toBe('Headers {"b":"2","x-a":"1"}')
|
||||
})
|
||||
|
||||
test("rejects what it cannot build from, and invalid names and values, with TypeErrors the program can catch", async () => {
|
||||
expect(
|
||||
await value(`
|
||||
function message(run) {
|
||||
try { run(); return null } catch (error) { return error instanceof TypeError ? error.message : error }
|
||||
}
|
||||
const headers = new Headers()
|
||||
return [
|
||||
message(() => Headers()),
|
||||
message(() => new Headers(null)),
|
||||
message(() => new Headers(1)),
|
||||
message(() => new Headers("a=1")),
|
||||
message(() => new Headers(new Date())),
|
||||
message(() => new Headers(() => 1)),
|
||||
message(() => new Headers([["name"]])),
|
||||
message(() => new Headers([["a", "b", "c"]])),
|
||||
message(() => new Headers({ "bad name": "x" })),
|
||||
message(() => new Headers({ name: "bad\u0000value" })),
|
||||
message(() => headers.get("invalid\u0100")),
|
||||
message(() => headers.has({})),
|
||||
message(() => headers.set("a", "invalid\u0100")),
|
||||
message(() => headers.append("a")),
|
||||
message(() => headers.forEach()),
|
||||
message(() => headers.forEach(1)),
|
||||
message(() => { const get = headers.get; return get("a") }),
|
||||
]
|
||||
`),
|
||||
).toEqual([
|
||||
"Constructor Headers requires 'new'.",
|
||||
"new Headers(...) expects a record of names to values, iterable [name, value] pairs, or Headers.",
|
||||
"new Headers(...) expects a record of names to values, iterable [name, value] pairs, or Headers.",
|
||||
"new Headers(...) expects a record of names to values, iterable [name, value] pairs, or Headers.",
|
||||
"new Headers(...) expects a record of names to values, iterable [name, value] pairs, or Headers.",
|
||||
"new Headers(...) expects a record of names to values, iterable [name, value] pairs, or Headers.",
|
||||
"new Headers(...) expects iterable [name, value] pairs.",
|
||||
"new Headers(...) expects iterable [name, value] pairs.",
|
||||
expect.stringContaining("bad name"),
|
||||
expect.stringContaining("invalid value"),
|
||||
expect.stringContaining("Invalid header name"),
|
||||
expect.stringContaining("[object Object]"),
|
||||
expect.stringContaining("invalid value"),
|
||||
"Headers.append requires 2 arguments.",
|
||||
"Headers.forEach requires 1 argument.",
|
||||
"Headers.forEach expects a function callback.",
|
||||
"Headers.prototype.get called on incompatible receiver undefined.",
|
||||
])
|
||||
})
|
||||
})
|
||||
|
||||
describe("Map", () => {
|
||||
test("get/set/has/size with chaining", async () => {
|
||||
expect(
|
||||
|
||||
@@ -3,10 +3,13 @@
|
||||
* - html/webappapis/atob/base64.any.js (btoa reference encoder, input list, and atob WebIDL cases)
|
||||
* - fetch/data-urls/resources/base64.json (copied to fixtures/wpt-base64.json)
|
||||
* - WebCryptoAPI/randomUUID.https.any.js
|
||||
* - fetch/api/headers/{headers-basic,headers-errors}.any.js
|
||||
*
|
||||
* Copyright © web-platform-tests contributors. Governed by the 3-Clause BSD license in LICENSE.wpt.
|
||||
*
|
||||
* `assert_throws_dom("InvalidCharacterError", …)` becomes a check for a TypeError: CodeMode has no DOMException.
|
||||
* Headers cases that need `Symbol.iterator`, iterator objects from `keys()`/`values()`/`entries()` (CodeMode returns
|
||||
* arrays), or a custom iterator on a Headers instance are left out.
|
||||
*/
|
||||
import { describe, expect, test } from "bun:test"
|
||||
import { Effect } from "effect"
|
||||
@@ -166,3 +169,225 @@ describe("crypto.randomUUID WPT parity (WebCryptoAPI/randomUUID.https.any.js)",
|
||||
).toEqual([true, true, true, 768])
|
||||
})
|
||||
})
|
||||
|
||||
// Enough of testharness.js to run the Headers files close to verbatim; each `test` records its failure, if any.
|
||||
const testharness = `
|
||||
const failures = []
|
||||
function test(run, name) { try { run() } catch (error) { failures.push(name + ": " + (error && error.message ? error.message : error)) } }
|
||||
function assert_equals(actual, expected, message) { if (actual !== expected) throw new Error((message || "") + " expected " + JSON.stringify(expected) + " got " + JSON.stringify(actual)) }
|
||||
function assert_true(actual, message) { assert_equals(actual, true, message) }
|
||||
function assert_false(actual, message) { assert_equals(actual, false, message) }
|
||||
function assert_array_equals(actual, expected, message) { assert_equals(JSON.stringify(actual), JSON.stringify(expected), message) }
|
||||
function assert_throws_js(type, run) { try { run() } catch (error) { if (error instanceof type) return; throw new Error("threw " + error.name) } throw new Error("did not throw") }
|
||||
function assert_unreached() { throw new Error("unreachable") }
|
||||
`
|
||||
|
||||
describe("Headers WPT parity (fetch/api/headers)", () => {
|
||||
test("headers-basic.any.js", async () => {
|
||||
expect(
|
||||
await value(`
|
||||
${testharness}
|
||||
test(function() { new Headers() }, "Create headers from no parameter")
|
||||
test(function() { new Headers(undefined) }, "Create headers from undefined parameter")
|
||||
test(function() { new Headers({}) }, "Create headers from empty object")
|
||||
var parameters = [null, 1]
|
||||
parameters.forEach(function(parameter) {
|
||||
test(function() { assert_throws_js(TypeError, function() { new Headers(parameter) }) }, "Create headers with " + parameter + " should throw")
|
||||
})
|
||||
var headerDict = {"name1": "value1", "name2": "value2", "name3": "value3", "name4": null, "name5": undefined, "name6": 1, "Content-Type": "value4"}
|
||||
var headerSeq = []
|
||||
for (var name in headerDict) headerSeq.push([name, headerDict[name]])
|
||||
test(function() {
|
||||
var headers = new Headers(headerSeq)
|
||||
for (name in headerDict) assert_equals(headers.get(name), String(headerDict[name]), "name: " + name + " has value: " + headerDict[name])
|
||||
assert_equals(headers.get("length"), null, "init should be treated as a sequence, not as a dictionary")
|
||||
}, "Create headers with sequence")
|
||||
test(function() {
|
||||
var headers = new Headers(headerDict)
|
||||
for (name in headerDict) assert_equals(headers.get(name), String(headerDict[name]), "name: " + name + " has value: " + headerDict[name])
|
||||
}, "Create headers with record")
|
||||
test(function() {
|
||||
var headers = new Headers(headerDict)
|
||||
var headers2 = new Headers(headers)
|
||||
for (name in headerDict) assert_equals(headers2.get(name), String(headerDict[name]), "name: " + name + " has value: " + headerDict[name])
|
||||
}, "Create headers with existing headers")
|
||||
test(function() {
|
||||
var headers = new Headers()
|
||||
for (name in headerDict) {
|
||||
headers.append(name, headerDict[name])
|
||||
assert_equals(headers.get(name), String(headerDict[name]), "name: " + name + " has value: " + headerDict[name])
|
||||
}
|
||||
}, "Check append method")
|
||||
test(function() {
|
||||
var headers = new Headers()
|
||||
for (name in headerDict) {
|
||||
headers.set(name, headerDict[name])
|
||||
assert_equals(headers.get(name), String(headerDict[name]), "name: " + name + " has value: " + headerDict[name])
|
||||
}
|
||||
}, "Check set method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerDict)
|
||||
for (name in headerDict) assert_true(headers.has(name), "headers has name " + name)
|
||||
assert_false(headers.has("nameNotInHeaders"), "headers do not have header: nameNotInHeaders")
|
||||
}, "Check has method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerDict)
|
||||
for (name in headerDict) {
|
||||
assert_true(headers.has(name), "headers have a header: " + name)
|
||||
headers.delete(name)
|
||||
assert_true(!headers.has(name), "headers do not have anymore a header: " + name)
|
||||
}
|
||||
}, "Check delete method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerDict)
|
||||
for (name in headerDict) assert_equals(headers.get(name), String(headerDict[name]), "name: " + name + " has value: " + headerDict[name])
|
||||
assert_equals(headers.get("nameNotInHeaders"), null, "header: nameNotInHeaders has no value")
|
||||
}, "Check get method")
|
||||
var headerEntriesDict = {"name1": "value1", "Name2": "value2", "name": "value3", "content-Type": "value4", "Content-Typ": "value5", "Content-Types": "value6"}
|
||||
var sortedHeaderDict = {}
|
||||
var headerValues = []
|
||||
var sortedHeaderKeys = Object.keys(headerEntriesDict).map(function(value) {
|
||||
sortedHeaderDict[value.toLowerCase()] = headerEntriesDict[value]
|
||||
headerValues.push(headerEntriesDict[value])
|
||||
return value.toLowerCase()
|
||||
}).sort()
|
||||
test(function() {
|
||||
var headers = new Headers(headerEntriesDict)
|
||||
assert_array_equals(headers.keys(), sortedHeaderKeys)
|
||||
for (const key of headers.keys()) assert_true(sortedHeaderKeys.indexOf(key) != -1)
|
||||
}, "Check keys method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerEntriesDict)
|
||||
assert_array_equals(headers.values(), sortedHeaderKeys.map((key) => sortedHeaderDict[key]))
|
||||
for (const value of headers.values()) assert_true(headerValues.indexOf(value) != -1)
|
||||
}, "Check values method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerEntriesDict)
|
||||
assert_array_equals(headers.entries(), sortedHeaderKeys.map((key) => [key, sortedHeaderDict[key]]))
|
||||
for (const entry of headers.entries()) assert_equals(entry[1], sortedHeaderDict[entry[0]])
|
||||
}, "Check entries method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerEntriesDict)
|
||||
assert_array_equals([...headers], sortedHeaderKeys.map((key) => [key, sortedHeaderDict[key]]))
|
||||
}, "Check Symbol.iterator method")
|
||||
test(function() {
|
||||
var headers = new Headers(headerEntriesDict)
|
||||
var index = 0
|
||||
headers.forEach(function(value, key, container) {
|
||||
assert_equals(headers, container)
|
||||
assert_equals(key, sortedHeaderKeys[index])
|
||||
assert_equals(value, sortedHeaderDict[sortedHeaderKeys[index]])
|
||||
index++
|
||||
})
|
||||
assert_equals(index, sortedHeaderKeys.length)
|
||||
}, "Check forEach method")
|
||||
test(() => {
|
||||
const headers = new Headers({"foo": "2", "baz": "1", "BAR": "0"})
|
||||
const actualKeys = []
|
||||
const actualValues = []
|
||||
for (const [header, value] of headers) {
|
||||
actualKeys.push(header)
|
||||
actualValues.push(value)
|
||||
headers.delete("foo")
|
||||
}
|
||||
assert_array_equals(actualKeys, ["bar", "baz"])
|
||||
assert_array_equals(actualValues, ["0", "1"])
|
||||
}, "Iteration skips elements removed while iterating")
|
||||
test(() => {
|
||||
const headers = new Headers({"foo": "2", "baz": "1", "BAR": "0", "quux": "3"})
|
||||
const actualKeys = []
|
||||
const actualValues = []
|
||||
for (const [header, value] of headers) {
|
||||
actualKeys.push(header)
|
||||
actualValues.push(value)
|
||||
if (header === "baz") headers.delete("bar")
|
||||
}
|
||||
assert_array_equals(actualKeys, ["bar", "baz", "quux"])
|
||||
assert_array_equals(actualValues, ["0", "1", "3"])
|
||||
}, "Removing elements already iterated over causes an element to be skipped during iteration")
|
||||
test(() => {
|
||||
const headers = new Headers({"foo": "2", "baz": "1", "BAR": "0", "quux": "3"})
|
||||
const actualKeys = []
|
||||
const actualValues = []
|
||||
for (const [header, value] of headers) {
|
||||
actualKeys.push(header)
|
||||
actualValues.push(value)
|
||||
if (header === "baz") headers.append("X-yZ", "4")
|
||||
}
|
||||
assert_array_equals(actualKeys, ["bar", "baz", "foo", "quux", "x-yz"])
|
||||
assert_array_equals(actualValues, ["0", "1", "2", "3", "4"])
|
||||
}, "Appending a value pair during iteration causes it to be reached during iteration")
|
||||
test(() => {
|
||||
const headers = new Headers({"foo": "2", "baz": "1", "BAR": "0", "quux": "3"})
|
||||
const actualKeys = []
|
||||
const actualValues = []
|
||||
for (const [header, value] of headers) {
|
||||
actualKeys.push(header)
|
||||
actualValues.push(value)
|
||||
if (header === "baz") headers.append("abc", "-1")
|
||||
}
|
||||
assert_array_equals(actualKeys, ["bar", "baz", "baz", "foo", "quux"])
|
||||
assert_array_equals(actualValues, ["0", "1", "1", "2", "3"])
|
||||
}, "Prepending a value pair before the current element position causes it to be skipped during iteration and adds the current element a second time")
|
||||
return failures
|
||||
`),
|
||||
).toEqual([])
|
||||
})
|
||||
|
||||
test("headers-errors.any.js", async () => {
|
||||
expect(
|
||||
await value(`
|
||||
${testharness}
|
||||
test(function() { assert_throws_js(TypeError, function() { new Headers([["name"]]) }) }, "Create headers giving an array having one string as init argument")
|
||||
test(function() { assert_throws_js(TypeError, function() { new Headers([["invalid", "invalidValue1", "invalidValue2"]]) }) }, "Create headers giving an array having three strings as init argument")
|
||||
test(function() { assert_throws_js(TypeError, function() { new Headers([["invalid\u0100", "Value1"]]) }) }, "Create headers giving bad header name as init argument")
|
||||
test(function() { assert_throws_js(TypeError, function() { new Headers([["name", "invalidValue\u0100"]]) }) }, "Create headers giving bad header value as init argument")
|
||||
var badNames = ["invalid\u0100", {}]
|
||||
var badValues = ["invalid\u0100"]
|
||||
badNames.forEach(function(name) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.get(name) }) }, "Check headers get with an invalid name " + name)
|
||||
})
|
||||
badNames.forEach(function(name) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.delete(name) }) }, "Check headers delete with an invalid name " + name)
|
||||
})
|
||||
badNames.forEach(function(name) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.has(name) }) }, "Check headers has with an invalid name " + name)
|
||||
})
|
||||
badNames.forEach(function(name) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.set(name, "Value1") }) }, "Check headers set with an invalid name " + name)
|
||||
})
|
||||
badValues.forEach(function(value) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.set("name", value) }) }, "Check headers set with an invalid value " + value)
|
||||
})
|
||||
badNames.forEach(function(name) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.append("invalid\u0100", "Value1") }) }, "Check headers append with an invalid name " + name)
|
||||
})
|
||||
badValues.forEach(function(value) {
|
||||
test(function() { var headers = new Headers(); assert_throws_js(TypeError, function() { headers.append("name", value) }) }, "Check headers append with an invalid value " + value)
|
||||
})
|
||||
test(function() {
|
||||
var headers = new Headers([["name", "value"]])
|
||||
assert_throws_js(TypeError, function() { headers.forEach() })
|
||||
assert_throws_js(TypeError, function() { headers.forEach(undefined) })
|
||||
assert_throws_js(TypeError, function() { headers.forEach(1) })
|
||||
}, "Headers forEach throws if argument is not callable")
|
||||
test(function() {
|
||||
var headers = new Headers([["name1", "value1"], ["name2", "value2"], ["name3", "value3"]])
|
||||
var counter = 0
|
||||
try {
|
||||
headers.forEach(function(value, name) {
|
||||
counter++
|
||||
if (name == "name2") throw "error"
|
||||
})
|
||||
} catch (e) {
|
||||
assert_equals(counter, 2)
|
||||
assert_equals(e, "error")
|
||||
return
|
||||
}
|
||||
assert_unreached()
|
||||
}, "Headers forEach loop should stop if callback is throwing exception")
|
||||
return failures
|
||||
`),
|
||||
).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
@@ -57,9 +57,6 @@ export const Plugin = define({
|
||||
|
||||
const hook = (event: SessionHooks["context"]) =>
|
||||
Effect.gen(function* () {
|
||||
const session = yield* ctx.session.get({ sessionID: event.sessionID }).pipe(Effect.orDie)
|
||||
if (session.parentID) return
|
||||
|
||||
const active = sessions.get(event.sessionID)
|
||||
const settings = yield* loadSettings()
|
||||
if (!settings) {
|
||||
|
||||
@@ -100,6 +100,7 @@ export const layer = (options?: Options) =>
|
||||
const recoverShell = Effect.fnUntraced(function* (
|
||||
background: Job.Background,
|
||||
recovery: Extract<Job.Recovery, { kind: "shell" }>,
|
||||
suspended: ReadonlySet<SessionSchema.ID>,
|
||||
) {
|
||||
const state = background.status === "running" ? "cancelled" : background.status
|
||||
const text =
|
||||
@@ -123,9 +124,7 @@ export const layer = (options?: Options) =>
|
||||
state,
|
||||
text,
|
||||
}),
|
||||
// Restart notices must not revive idle owners of long-lived shells.
|
||||
// Interrupted executions resume separately after their notices are admitted.
|
||||
resume: false,
|
||||
...(suspended.has(recovery.sessionID) ? { resume: false } : {}),
|
||||
})
|
||||
.pipe(
|
||||
Effect.catchTag("Session.NotFoundError", () => Effect.void),
|
||||
@@ -209,7 +208,7 @@ export const layer = (options?: Options) =>
|
||||
if ((yield* jobs.get(background.id))?.status === "running") return
|
||||
const recovery = background.recovery
|
||||
yield* recovery.kind === "shell"
|
||||
? recoverShell(background, recovery)
|
||||
? recoverShell(background, recovery, suspended)
|
||||
: recoverSubagent(background, recovery, suspended)
|
||||
}),
|
||||
{ discard: true },
|
||||
|
||||
@@ -320,24 +320,15 @@ export const layer = Layer.effect(
|
||||
// which transport actually carries the request, so both hook families are always offered.
|
||||
const webSocket =
|
||||
input.webSocket === "session" && model.transport === "websocket"
|
||||
? transport.bind(session.id, {
|
||||
handshake: (connect) =>
|
||||
hooks
|
||||
.trigger("session", "experimental.ws.handshake", {
|
||||
...scope,
|
||||
url: connect.url,
|
||||
headers: connect.headers,
|
||||
})
|
||||
.pipe(Effect.map((event) => ({ url: event.url, headers: event.headers }))),
|
||||
send: (frame) =>
|
||||
hooks
|
||||
.trigger("session", "experimental.ws.send", { ...scope, frame })
|
||||
.pipe(Effect.map((event) => event.frame)),
|
||||
receive: (frame) =>
|
||||
hooks
|
||||
.trigger("session", "experimental.ws.receive", { ...scope, frame })
|
||||
.pipe(Effect.map((event) => event.frame)),
|
||||
})
|
||||
? transport.bind(session.id, (connect) =>
|
||||
hooks
|
||||
.trigger("session", "experimental.ws.handshake", {
|
||||
...scope,
|
||||
url: connect.url,
|
||||
headers: connect.headers,
|
||||
})
|
||||
.pipe(Effect.map((event) => ({ url: event.url, headers: event.headers }))),
|
||||
)
|
||||
: undefined
|
||||
|
||||
return {
|
||||
|
||||
@@ -59,18 +59,11 @@ export interface Handshake {
|
||||
readonly headers: Record<string, string>
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-exchange taps. `handshake` runs before the connection is selected; `send` sees each outbound
|
||||
* frame after the driver builds it; `receive` sees each inbound frame before the driver observes it.
|
||||
*/
|
||||
export interface Interceptor {
|
||||
readonly handshake?: (connect: Handshake) => Effect.Effect<Handshake>
|
||||
readonly send?: (frame: string) => Effect.Effect<string>
|
||||
readonly receive?: (frame: string) => Effect.Effect<string>
|
||||
}
|
||||
|
||||
export interface Interface {
|
||||
readonly bind: (sessionID: SessionSchema.ID, interceptor?: Interceptor) => WebSocketChannelExecutor
|
||||
readonly bind: (
|
||||
sessionID: SessionSchema.ID,
|
||||
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
|
||||
) => WebSocketChannelExecutor
|
||||
readonly close: (sessionID: SessionSchema.ID) => Effect.Effect<void>
|
||||
readonly closeAll: Effect.Effect<void>
|
||||
}
|
||||
@@ -285,7 +278,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
const start = Effect.fn("SessionModelTransport.start")(function* (
|
||||
owner: State,
|
||||
input: WebSocketChannelExchange,
|
||||
interceptor?: Interceptor,
|
||||
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
|
||||
) {
|
||||
if (owner.closed)
|
||||
return yield* transportError("Session WebSocket owner is closed", {
|
||||
@@ -295,8 +288,8 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
delivery: "not-sent",
|
||||
})
|
||||
if (owner.httpFallback) return fallback(input)
|
||||
const selected = interceptor?.handshake
|
||||
? yield* interceptor.handshake({ url: input.connect.url, headers: { ...input.connect.headers } })
|
||||
const selected = handshake
|
||||
? yield* handshake({ url: input.connect.url, headers: { ...input.connect.headers } })
|
||||
: undefined
|
||||
const exchange: WebSocketChannelExchange = selected
|
||||
? { ...input, connect: { ...input.connect, url: selected.url, headers: Headers.fromInput(selected.headers) } }
|
||||
@@ -361,9 +354,6 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
Effect.onInterrupt(() => closeChannel(owner, channel)),
|
||||
)
|
||||
if (create.mode === "full") channel.checkpoint = undefined
|
||||
const message = interceptor?.send
|
||||
? yield* interceptor.send(create.message).pipe(Effect.onInterrupt(() => closeChannel(owner, channel)))
|
||||
: create.message
|
||||
yield* Effect.logDebug("session websocket sending", {
|
||||
sessionTransport: "websocket",
|
||||
phase: "send",
|
||||
@@ -374,7 +364,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
delivery: "send-attempted",
|
||||
}
|
||||
channel.active = active
|
||||
const sent = yield* channel.connection.sendText(message).pipe(
|
||||
const sent = yield* channel.connection.sendText(create.message).pipe(
|
||||
Effect.withSpan("SessionModelTransport.send"),
|
||||
Effect.onInterrupt(() => closeChannel(owner, channel)),
|
||||
Effect.result,
|
||||
@@ -415,7 +405,6 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
}),
|
||||
),
|
||||
}),
|
||||
Stream.mapEffect((frame) => (interceptor?.receive ? interceptor.receive(frame) : Effect.succeed(frame))),
|
||||
Stream.mapEffect((frame) => exchange.driver.observe(create, frame)),
|
||||
Stream.tap((observation) =>
|
||||
Effect.sync(() => {
|
||||
@@ -493,7 +482,10 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
return { frames, complete, http: channel.connection.http }
|
||||
})
|
||||
|
||||
const bind = (sessionID: SessionSchema.ID, interceptor?: Interceptor): WebSocketChannelExecutor => ({
|
||||
const bind = (
|
||||
sessionID: SessionSchema.ID,
|
||||
handshake?: (connect: Handshake) => Effect.Effect<Handshake>,
|
||||
): WebSocketChannelExecutor => ({
|
||||
execute: (exchange) => {
|
||||
const owner = state(sessionID)
|
||||
let execution: WebSocketChannelExecution | undefined
|
||||
@@ -503,7 +495,7 @@ export const makeLayer = (connector: WebSocketConnector) =>
|
||||
},
|
||||
frames: Stream.unwrap(
|
||||
Effect.acquireRelease(owner.lock.take(1), () => owner.lock.release(1), { interruptible: true }).pipe(
|
||||
Effect.andThen(start(owner, exchange, interceptor)),
|
||||
Effect.andThen(start(owner, exchange, handshake)),
|
||||
Effect.tap((started) =>
|
||||
Effect.sync(() => {
|
||||
execution = started
|
||||
|
||||
@@ -379,7 +379,7 @@ describe("SessionExecution lifecycle", () => {
|
||||
})
|
||||
|
||||
describe("SessionRestart background recovery", () => {
|
||||
it.effect("keeps shell owners idle until a user prompt delivers recovered notices exactly once", () =>
|
||||
it.effect("wakes idle shell owners and delivers recovered notices exactly once", () =>
|
||||
Effect.gen(function* () {
|
||||
const database = yield* Database.Service
|
||||
const store = yield* SessionStore.Service
|
||||
@@ -416,20 +416,6 @@ describe("SessionRestart background recovery", () => {
|
||||
yield* restart.resumeSuspendedSessions
|
||||
yield* Effect.forEach([parent, child], execution.awaitIdle, { discard: true })
|
||||
|
||||
expect(drained).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(1)
|
||||
expect(yield* SessionInbox.list(database.db, child)).toHaveLength(1)
|
||||
expect(yield* restarted.pendingBackground).toEqual([])
|
||||
yield* restart.resumeSuspendedSessions
|
||||
expect(drained).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, parent)).toHaveLength(1)
|
||||
expect(yield* SessionInbox.list(database.db, child)).toHaveLength(1)
|
||||
|
||||
yield* seedInbox(database, parent, ["steer"])
|
||||
yield* seedInbox(database, child, ["steer"])
|
||||
yield* execution.wake(parent)
|
||||
yield* execution.wake(child)
|
||||
yield* Effect.forEach([parent, child], execution.awaitIdle, { discard: true })
|
||||
expect(drained.toSorted()).toEqual([parent, child].toSorted())
|
||||
expect((yield* store.context(parent)).filter((message) => message.type === "synthetic")).toMatchObject([
|
||||
{
|
||||
@@ -523,7 +509,7 @@ describe("SessionRestart background recovery", () => {
|
||||
yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions
|
||||
yield* Context.get(context, SessionExecution.Service).awaitIdle(sessionID)
|
||||
|
||||
expect(drained).toEqual([])
|
||||
expect(drained).toEqual([sessionID])
|
||||
const inbox = yield* SessionInbox.list(database.db, sessionID)
|
||||
expect(inbox).toMatchObject([
|
||||
{
|
||||
@@ -575,6 +561,7 @@ describe("SessionRestart background recovery", () => {
|
||||
expect(yield* restarted.pendingBackground).toEqual([])
|
||||
expect(yield* SessionInbox.list(database.db, sessionID)).toHaveLength(delivered ? 0 : 1)
|
||||
yield* SessionInbox.promote(database.db, bus, sessionID, "steer")
|
||||
// Recovery ends a busy period, so an idle marker follows the notification.
|
||||
const messages = (yield* sessions.messages({ sessionID })).filter((message) => message.type !== "idle")
|
||||
expect(messages).toMatchObject([
|
||||
{
|
||||
|
||||
@@ -80,7 +80,7 @@ describe("SessionModelRequest HTTP hooks", () => {
|
||||
}).pipe(Effect.provideService(SessionModelTransport.Service, transport)),
|
||||
)
|
||||
|
||||
it.effect("offers the WebSocket executor alongside HTTP hooks and routes the WebSocket hooks", () =>
|
||||
it.effect("offers the WebSocket executor alongside HTTP hooks and routes the handshake hook", () =>
|
||||
Effect.gen(function* () {
|
||||
const hooks = yield* PluginHooks.Service
|
||||
const seen: string[] = []
|
||||
@@ -92,31 +92,13 @@ describe("SessionModelRequest HTTP hooks", () => {
|
||||
delete event.headers["api-key"]
|
||||
}),
|
||||
)
|
||||
yield* hooks.register("session", "experimental.ws.send", (event) =>
|
||||
Effect.sync(() => {
|
||||
seen.push(`send:${event.kind}:${event.frame}`)
|
||||
event.frame = `${event.frame}+plugin`
|
||||
}),
|
||||
)
|
||||
yield* hooks.register("session", "experimental.ws.receive", (event) =>
|
||||
Effect.sync(() => {
|
||||
seen.push(`receive:${event.kind}:${event.frame}`)
|
||||
event.frame = event.frame.toUpperCase()
|
||||
}),
|
||||
)
|
||||
const bound: Array<{ url: string; headers: Record<string, string> }> = []
|
||||
const frames: string[] = []
|
||||
const websocketTransport = SessionModelTransport.Service.of({
|
||||
bind: (_sessionID, interceptor) => ({
|
||||
bind: (_sessionID, handshake) => ({
|
||||
execute: () =>
|
||||
Effect.gen(function* () {
|
||||
if (!interceptor?.handshake || !interceptor.send || !interceptor.receive)
|
||||
throw new Error("Expected a full WebSocket interceptor")
|
||||
bound.push(
|
||||
yield* interceptor.handshake({ url: "wss://example.test/v1/responses", headers: { "api-key": "k" } }),
|
||||
)
|
||||
frames.push(yield* interceptor.send("create"))
|
||||
frames.push(yield* interceptor.receive("created"))
|
||||
if (!handshake) throw new Error("Expected a handshake interceptor")
|
||||
bound.push(yield* handshake({ url: "wss://example.test/v1/responses", headers: { "api-key": "k" } }))
|
||||
return { frames: Stream.empty, complete: Effect.void }
|
||||
}),
|
||||
}),
|
||||
@@ -145,12 +127,7 @@ describe("SessionModelRequest HTTP hooks", () => {
|
||||
expect(prepared.options.webSocket).toBeDefined()
|
||||
yield* prepared.options.webSocket!.execute({} as never)
|
||||
expect(bound).toEqual([{ url: "wss://example.test/v1/responses", headers: { authorization: "Bearer minted" } }])
|
||||
expect(frames).toEqual(["create+plugin", "CREATED"])
|
||||
expect(seen).toEqual([
|
||||
"handshake:primary:wss://example.test/v1/responses",
|
||||
"send:primary:create",
|
||||
"receive:primary:created",
|
||||
])
|
||||
expect(seen).toEqual(["handshake:primary:wss://example.test/v1/responses"])
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
||||
@@ -178,13 +178,12 @@ describe("SessionModelTransport", () => {
|
||||
fixture.connector,
|
||||
Effect.gen(function* () {
|
||||
const transport = yield* SessionModelTransport.Service
|
||||
const executor = transport.bind(session, {
|
||||
handshake: (connect) =>
|
||||
Effect.succeed({
|
||||
url: connect.url,
|
||||
headers: { ...connect.headers, authorization: `Bearer ${tokens.shift()}` },
|
||||
}),
|
||||
})
|
||||
const executor = transport.bind(session, (connect) =>
|
||||
Effect.succeed({
|
||||
url: connect.url,
|
||||
headers: { ...connect.headers, authorization: `Bearer ${tokens.shift()}` },
|
||||
}),
|
||||
)
|
||||
yield* collect(executor, exchange("first", { headers: { "api-key": "k" } }))
|
||||
yield* collect(executor, exchange("second", { headers: { "api-key": "k" } }))
|
||||
yield* collect(executor, exchange("third", { headers: { "api-key": "k" } }))
|
||||
@@ -197,36 +196,6 @@ describe("SessionModelTransport", () => {
|
||||
)
|
||||
})
|
||||
|
||||
test("sends the frame the send tap returns and observes the frame the receive tap returns", async () => {
|
||||
const fixture = automatic()
|
||||
const seen: Array<{ tap: "send" | "receive"; frame: string }> = []
|
||||
await run(
|
||||
fixture.connector,
|
||||
Effect.gen(function* () {
|
||||
const transport = yield* SessionModelTransport.Service
|
||||
const executor = transport.bind(session, {
|
||||
send: (frame) => {
|
||||
seen.push({ tap: "send", frame })
|
||||
return Effect.succeed(`${frame}:rewritten`)
|
||||
},
|
||||
receive: (frame) => {
|
||||
seen.push({ tap: "receive", frame })
|
||||
return Effect.succeed(`${frame}:observed`)
|
||||
},
|
||||
})
|
||||
const frames = yield* collect(executor, exchange("first"))
|
||||
|
||||
// The wire carries the rewritten outbound frame; the driver sees the rewritten inbound frame.
|
||||
expect(fixture.connections.map((item) => item.sent)).toEqual([["first:rewritten"]])
|
||||
expect(frames).toEqual(["completed:first:rewritten:observed"])
|
||||
expect(seen).toEqual([
|
||||
{ tap: "send", frame: "first" },
|
||||
{ tap: "receive", frame: "completed:first:rewritten" },
|
||||
])
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
||||
test("does not carry a checkpoint across physical connection rotation", async () => {
|
||||
const fixture = automatic()
|
||||
const checkpoints: Array<unknown> = []
|
||||
|
||||
@@ -99,31 +99,6 @@ export interface SessionWebSocketHandshake {
|
||||
headers: Record<string, string>
|
||||
}
|
||||
|
||||
/**
|
||||
* Outbound frame about to be written to the Session's socket, after the provider driver has built
|
||||
* it. Replacing `frame` sends the replacement verbatim; the driver still tracks state from the
|
||||
* provider's replies, so a rewrite that changes protocol meaning is on the plugin. Experimental.
|
||||
*/
|
||||
export interface SessionWebSocketSend {
|
||||
readonly sessionID: Session.ID
|
||||
readonly agent: Agent.ID
|
||||
readonly model: Model.Ref
|
||||
readonly kind: SessionRequestKind
|
||||
frame: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Inbound frame read from the Session's socket, before the provider driver observes it. Replacing
|
||||
* `frame` hands the replacement to the driver verbatim. Experimental.
|
||||
*/
|
||||
export interface SessionWebSocketReceive {
|
||||
readonly sessionID: Session.ID
|
||||
readonly agent: Agent.ID
|
||||
readonly model: Model.Ref
|
||||
readonly kind: SessionRequestKind
|
||||
frame: string
|
||||
}
|
||||
|
||||
export type SessionRetryDecision = { retry: false } | { retry: true; delay: number }
|
||||
|
||||
export interface SessionRetry {
|
||||
@@ -145,8 +120,6 @@ export interface SessionHooks {
|
||||
readonly "http.request": SessionHttpRequest
|
||||
readonly "http.response": SessionHttpResponse
|
||||
readonly "experimental.ws.handshake": SessionWebSocketHandshake
|
||||
readonly "experimental.ws.send": SessionWebSocketSend
|
||||
readonly "experimental.ws.receive": SessionWebSocketReceive
|
||||
readonly retry: SessionRetry
|
||||
}
|
||||
|
||||
|
||||
@@ -99,31 +99,6 @@ export interface SessionWebSocketHandshake {
|
||||
headers: Record<string, string>
|
||||
}
|
||||
|
||||
/**
|
||||
* Outbound frame about to be written to the Session's socket, after the provider driver has built
|
||||
* it. Replacing `frame` sends the replacement verbatim; the driver still tracks state from the
|
||||
* provider's replies, so a rewrite that changes protocol meaning is on the plugin. Experimental.
|
||||
*/
|
||||
export interface SessionWebSocketSend {
|
||||
readonly sessionID: Session.ID
|
||||
readonly agent: Agent.ID
|
||||
readonly model: Model.Ref
|
||||
readonly kind: SessionRequestKind
|
||||
frame: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Inbound frame read from the Session's socket, before the provider driver observes it. Replacing
|
||||
* `frame` hands the replacement to the driver verbatim. Experimental.
|
||||
*/
|
||||
export interface SessionWebSocketReceive {
|
||||
readonly sessionID: Session.ID
|
||||
readonly agent: Agent.ID
|
||||
readonly model: Model.Ref
|
||||
readonly kind: SessionRequestKind
|
||||
frame: string
|
||||
}
|
||||
|
||||
export type SessionRetryDecision = { retry: false } | { retry: true; delay: number }
|
||||
|
||||
export interface SessionRetry {
|
||||
@@ -145,8 +120,6 @@ export interface SessionHooks {
|
||||
readonly "http.request": SessionHttpRequest
|
||||
readonly "http.response": SessionHttpResponse
|
||||
readonly "experimental.ws.handshake": SessionWebSocketHandshake
|
||||
readonly "experimental.ws.send": SessionWebSocketSend
|
||||
readonly "experimental.ws.receive": SessionWebSocketReceive
|
||||
readonly retry: SessionRetry
|
||||
}
|
||||
|
||||
|
||||
@@ -1277,26 +1277,6 @@ effect: (ctx) =>
|
||||
}),
|
||||
```
|
||||
|
||||
`experimental.ws.send` and `experimental.ws.receive` expose the frames themselves: `send` runs after the provider
|
||||
driver builds an outbound frame, `receive` runs on each inbound frame before the driver observes it. Whatever `frame` holds when the hook returns is what crosses the wire or reaches the driver;
|
||||
OpenCode does not validate it.
|
||||
|
||||
```ts
|
||||
effect: (ctx) =>
|
||||
Effect.gen(function* () {
|
||||
yield* ctx.session.hook(
|
||||
"experimental.ws.send",
|
||||
(event) =>
|
||||
Effect.sync(() => {
|
||||
const body = JSON.parse(event.frame)
|
||||
if (body.type === "response.create") body.metadata = { ...body.metadata, session: event.sessionID }
|
||||
event.frame = JSON.stringify(body)
|
||||
}),
|
||||
{ providerID: "openai" },
|
||||
)
|
||||
}),
|
||||
```
|
||||
|
||||
Override the retry decision for a provider failure or replace its delay in milliseconds. The hook runs after OpenCode
|
||||
classifies the failure and proposes its policy, but before any retry is scheduled. It does not expose how OpenCode
|
||||
internally performs the next attempt.
|
||||
@@ -1337,9 +1317,6 @@ interface SessionHooks {
|
||||
readonly "model.request": SessionModelRequest
|
||||
readonly "http.request": SessionHttpRequest
|
||||
readonly "http.response": SessionHttpResponse
|
||||
readonly "experimental.ws.handshake": SessionWebSocketHandshake
|
||||
readonly "experimental.ws.send": SessionWebSocketSend
|
||||
readonly "experimental.ws.receive": SessionWebSocketReceive
|
||||
readonly retry: SessionRetry
|
||||
}
|
||||
|
||||
|
||||
@@ -1410,31 +1410,7 @@ await ctx.session.hook(
|
||||
)
|
||||
```
|
||||
|
||||
`experimental.ws.send` and `experimental.ws.receive` expose the frames themselves, the WebSocket counterpart of
|
||||
editing an HTTP request or response body. `send` runs after the provider driver builds an outbound frame and before it
|
||||
is written; `receive` runs on each inbound frame before the driver observes it. Both carry the frame as a string and
|
||||
send whatever `frame` holds when the hook returns.
|
||||
|
||||
OpenCode does not validate rewritten frames. The driver tracks state from the provider's replies, so a rewrite that
|
||||
changes protocol meaning is the plugin's responsibility, just as a rewritten HTTP body is.
|
||||
|
||||
```ts
|
||||
await ctx.session.hook(
|
||||
"experimental.ws.send",
|
||||
(event) => {
|
||||
const body = JSON.parse(event.frame)
|
||||
if (body.type === "response.create") body.metadata = { ...body.metadata, session: event.sessionID }
|
||||
event.frame = JSON.stringify(body)
|
||||
},
|
||||
{ providerID: "openai" },
|
||||
)
|
||||
|
||||
await ctx.session.hook("experimental.ws.receive", (event) => {
|
||||
if (event.frame.includes('"type":"error"')) console.error(event.frame)
|
||||
})
|
||||
```
|
||||
|
||||
These hooks are experimental and their names or shapes may change.
|
||||
This hook is experimental and its name or shape may change.
|
||||
|
||||
#### Retry policy
|
||||
|
||||
@@ -1482,8 +1458,6 @@ interface SessionHooks {
|
||||
"http.request": SessionHttpRequestHook
|
||||
"http.response": SessionHttpResponseHook
|
||||
"experimental.ws.handshake": SessionWebSocketHandshakeHook
|
||||
"experimental.ws.send": SessionWebSocketSendHook
|
||||
"experimental.ws.receive": SessionWebSocketReceiveHook
|
||||
retry: SessionRetryHook
|
||||
}
|
||||
|
||||
@@ -1496,22 +1470,6 @@ interface SessionWebSocketHandshakeHook {
|
||||
headers: Record<string, string>
|
||||
}
|
||||
|
||||
interface SessionWebSocketSendHook {
|
||||
readonly sessionID: string
|
||||
readonly agent: string
|
||||
readonly model: { providerID: string; id: string; variant?: string }
|
||||
readonly kind: "primary" | "compaction" | "title" | "generate"
|
||||
frame: string
|
||||
}
|
||||
|
||||
interface SessionWebSocketReceiveHook {
|
||||
readonly sessionID: string
|
||||
readonly agent: string
|
||||
readonly model: { providerID: string; id: string; variant?: string }
|
||||
readonly kind: "primary" | "compaction" | "title" | "generate"
|
||||
frame: string
|
||||
}
|
||||
|
||||
type RetryDecision = { retry: false } | { retry: true; delay: number }
|
||||
|
||||
interface SessionRetryHook {
|
||||
|
||||
Reference in New Issue
Block a user