| //#region src/_request.ts | ||
| const StubRequest = /* @__PURE__ */ (() => { | ||
| class StubRequest { | ||
| url; | ||
| _signal; | ||
| _headers; | ||
| _init; | ||
| constructor(url, init = {}) { | ||
| this.url = url; | ||
| this._init = init; | ||
| } | ||
| get headers() { | ||
| if (!this._headers) this._headers = new Headers(this._init?.headers); | ||
| return this._headers; | ||
| } | ||
| clone() { | ||
| return new StubRequest(this.url, this._init); | ||
| } | ||
| get method() { | ||
| return "GET"; | ||
| } | ||
| get signal() { | ||
| return this._signal ??= new AbortSignal(); | ||
| } | ||
| get cache() { | ||
| return "default"; | ||
| } | ||
| get credentials() { | ||
| return "same-origin"; | ||
| } | ||
| get destination() { | ||
| return ""; | ||
| } | ||
| get integrity() { | ||
| return ""; | ||
| } | ||
| get keepalive() { | ||
| return false; | ||
| } | ||
| get redirect() { | ||
| return "follow"; | ||
| } | ||
| get mode() { | ||
| return "cors"; | ||
| } | ||
| get referrer() { | ||
| return "about:client"; | ||
| } | ||
| get referrerPolicy() { | ||
| return ""; | ||
| } | ||
| get body() { | ||
| return null; | ||
| } | ||
| get bodyUsed() { | ||
| return false; | ||
| } | ||
| arrayBuffer() { | ||
| return Promise.resolve(/* @__PURE__ */ new ArrayBuffer(0)); | ||
| } | ||
| blob() { | ||
| return Promise.resolve(new Blob()); | ||
| } | ||
| bytes() { | ||
| return Promise.resolve(new Uint8Array()); | ||
| } | ||
| formData() { | ||
| return Promise.resolve(new FormData()); | ||
| } | ||
| json() { | ||
| return Promise.resolve(JSON.parse("")); | ||
| } | ||
| text() { | ||
| return Promise.resolve(""); | ||
| } | ||
| } | ||
| Object.setPrototypeOf(StubRequest.prototype, globalThis.Request.prototype); | ||
| return StubRequest; | ||
| })(); | ||
| //#endregion | ||
| export { StubRequest as t }; |
| import { a as Hooks } from "./adapter.mjs"; | ||
| import { n as BunOptions } from "./bun.mjs"; | ||
| import { n as CloudflareOptions } from "./cloudflare.mjs"; | ||
| import { n as DenoOptions } from "./deno.mjs"; | ||
| import { n as NodeOptions } from "./node.mjs"; | ||
| import { n as SSEOptions } from "./sse.mjs"; | ||
| import { Server, ServerOptions, ServerPlugin, ServerRequest } from "srvx"; | ||
| //#region src/server/_types.d.ts | ||
| type WSOptions = Partial<Hooks> & { | ||
| resolve?: (req: ServerRequest) => Partial<Hooks> | Promise<Partial<Hooks>>; | ||
| options?: { | ||
| bun?: BunOptions; | ||
| deno?: DenoOptions; | ||
| node?: NodeOptions; | ||
| sse?: SSEOptions; | ||
| cloudflare?: CloudflareOptions; | ||
| }; | ||
| }; | ||
| type ServerWithWSOptions = ServerOptions & { | ||
| websocket?: WSOptions; | ||
| }; | ||
| //#endregion | ||
| export { WSOptions as n, ServerWithWSOptions as t }; |
| import { a as WebSocket } from "./web.mjs"; | ||
| //#region src/error.d.ts | ||
| declare class WSError extends Error { | ||
| constructor(...args: any[]); | ||
| } | ||
| //#endregion | ||
| //#region src/utils.d.ts | ||
| declare const kNodeInspect: unique symbol; | ||
| //#endregion | ||
| //#region src/peer.d.ts | ||
| interface PeerContext extends Record<string, unknown> {} | ||
| interface AdapterInternal { | ||
| ws: unknown; | ||
| request: Request; | ||
| namespace: string; | ||
| peers?: Set<Peer>; | ||
| context?: PeerContext; | ||
| } | ||
| declare abstract class Peer<Internal extends AdapterInternal = AdapterInternal> { | ||
| #private; | ||
| protected _internal: Internal; | ||
| protected _topics: Set<string>; | ||
| protected _id?: string; | ||
| constructor(internal: Internal); | ||
| get context(): PeerContext; | ||
| get namespace(): string; | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the peer. | ||
| */ | ||
| get id(): string; | ||
| /** IP address of the peer */ | ||
| get remoteAddress(): string | undefined; | ||
| /** upgrade request */ | ||
| get request(): Request; | ||
| /** | ||
| * Get the [WebSocket](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket) instance. | ||
| * | ||
| * **Note:** crossws adds polyfill for the following properties if native values are not available: | ||
| * - `protocol`: Extracted from the `sec-websocket-protocol` header. | ||
| * - `extensions`: Extracted from the `sec-websocket-extensions` header. | ||
| * - `url`: Extracted from the request URL (http -> ws). | ||
| * */ | ||
| get websocket(): Partial<WebSocket>; | ||
| /** All connected peers to the server */ | ||
| get peers(): Set<Peer>; | ||
| /** All topics, this peer has been subscribed to. */ | ||
| get topics(): Set<string>; | ||
| abstract close(code?: number, reason?: string): void; | ||
| /** Abruptly close the connection */ | ||
| terminate(): void; | ||
| /** Subscribe to a topic */ | ||
| subscribe(topic: string): void; | ||
| /** Unsubscribe from a topic */ | ||
| unsubscribe(topic: string): void; | ||
| /** Send a message to the peer. */ | ||
| abstract send(data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number | void | undefined; | ||
| /** Send message to subscribes of topic */ | ||
| abstract publish(topic: string, data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): void; | ||
| toString(): string; | ||
| [Symbol.toPrimitive](): string; | ||
| [Symbol.toStringTag](): "WebSocket"; | ||
| [kNodeInspect](): unknown; | ||
| } | ||
| //#endregion | ||
| //#region src/message.d.ts | ||
| declare class Message implements Partial<MessageEvent> { | ||
| #private; | ||
| /** Access to the original [message event](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket/message_event) if available. */ | ||
| readonly event?: MessageEvent; | ||
| /** Access to the Peer that emitted the message. */ | ||
| readonly peer?: Peer; | ||
| /** Raw message data (can be of any type). */ | ||
| readonly rawData: unknown; | ||
| constructor(rawData: unknown, peer: Peer, event?: MessageEvent); | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the message. | ||
| */ | ||
| get id(): string; | ||
| /** | ||
| * Get data as [Uint8Array](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Uint8Array) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| uint8Array(): Uint8Array; | ||
| /** | ||
| * Get data as [ArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/ArrayBuffer) or [SharedArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/SharedArrayBuffer) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| arrayBuffer(): ArrayBuffer | SharedArrayBuffer; | ||
| /** | ||
| * Get data as [Blob](https://developer.mozilla.org/en-US/docs/Web/API/Blob) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. */ | ||
| blob(): Blob; | ||
| /** | ||
| * Get stringified text version of the message. | ||
| * | ||
| * If raw data is in any other format, it will be automatically converted and decoded. | ||
| */ | ||
| text(): string; | ||
| /** | ||
| * Get parsed version of the message text with [`JSON.parse()`](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/JSON/parse). | ||
| */ | ||
| json<T = unknown>(): T; | ||
| /** | ||
| * Message data (value varies based on `peer.websocket.binaryType`). | ||
| */ | ||
| get data(): unknown; | ||
| toString(): string; | ||
| [Symbol.toPrimitive](): string; | ||
| [kNodeInspect](): unknown; | ||
| } | ||
| //#endregion | ||
| //#region src/hooks.d.ts | ||
| declare function defineHooks<T extends Partial<Hooks> = Partial<Hooks>>(hooks: T): T; | ||
| type ResolveHooks = (request: Request & { | ||
| readonly context?: PeerContext; | ||
| }) => Partial<Hooks> | Promise<Partial<Hooks>>; | ||
| type MaybePromise<T> = T | Promise<T>; | ||
| interface Hooks { | ||
| /** | ||
| * Upgrading a request to a WebSocket connection. | ||
| * | ||
| * - You can throw a Response to abort the upgrade. | ||
| * - You can return { headers } to modify the response. | ||
| * - You can return { namespace } to change the pub/sub namespace. | ||
| * - You can return { context } to provide a custom peer context. | ||
| * | ||
| * @param request | ||
| * @throws {Response} | ||
| */ | ||
| upgrade: (request: Request & { | ||
| readonly context?: Record<string, unknown>; | ||
| }) => MaybePromise<{ | ||
| headers?: HeadersInit; | ||
| namespace?: string; | ||
| context?: PeerContext; | ||
| } | Response | void>; | ||
| /** A message is received */ | ||
| message: (peer: Peer, message: Message) => MaybePromise<void>; | ||
| /** A socket is opened */ | ||
| open: (peer: Peer) => MaybePromise<void>; | ||
| /** A socket is closed */ | ||
| close: (peer: Peer, details: { | ||
| code?: number; | ||
| reason?: string; | ||
| }) => MaybePromise<void>; | ||
| /** An error occurs */ | ||
| error: (peer: Peer, error: WSError) => MaybePromise<void>; | ||
| } | ||
| //#endregion | ||
| //#region src/adapter.d.ts | ||
| interface AdapterInstance { | ||
| readonly peers: Map<string, Set<Peer>>; | ||
| readonly publish: (topic: string, data: unknown, options?: { | ||
| compress?: boolean; | ||
| namespace?: string; | ||
| }) => void; | ||
| } | ||
| interface AdapterOptions { | ||
| resolve?: ResolveHooks; | ||
| getNamespace?: (request: Request) => string; | ||
| hooks?: Partial<Hooks>; | ||
| } | ||
| type Adapter<AdapterT extends AdapterInstance = AdapterInstance, Options extends AdapterOptions = AdapterOptions> = (options?: Options) => AdapterT; | ||
| declare function defineWebSocketAdapter<AdapterT extends AdapterInstance = AdapterInstance, Options extends AdapterOptions = AdapterOptions>(factory: Adapter<AdapterT, Options>): Adapter<AdapterT, Options>; | ||
| //#endregion | ||
| export { Hooks as a, Message as c, PeerContext as d, WSError as f, defineWebSocketAdapter as i, AdapterInternal as l, AdapterInstance as n, ResolveHooks as o, AdapterOptions as r, defineHooks as s, Adapter as t, Peer as u }; |
| //#region src/hooks.ts | ||
| var AdapterHookable = class { | ||
| options; | ||
| constructor(options) { | ||
| this.options = options || {}; | ||
| } | ||
| callHook(name, arg1, arg2) { | ||
| const globalHook = this.options.hooks?.[name]; | ||
| const globalPromise = globalHook?.(arg1, arg2); | ||
| const request = arg1.request || arg1; | ||
| const resolveHooksPromise = this.options.resolve?.(request); | ||
| if (!resolveHooksPromise) return globalPromise; | ||
| const resolvePromise = resolveHooksPromise instanceof Promise ? resolveHooksPromise.then((hooks) => hooks?.[name]) : resolveHooksPromise?.[name]; | ||
| return Promise.all([globalPromise, resolvePromise]).then(([globalRes, hook]) => { | ||
| const hookResPromise = hook?.(arg1, arg2); | ||
| return hookResPromise instanceof Promise ? hookResPromise.then((hookRes) => hookRes || globalRes) : hookResPromise || globalRes; | ||
| }); | ||
| } | ||
| async upgrade(request) { | ||
| let namespace = this.options.getNamespace?.(request) ?? new URL(request.url).pathname; | ||
| const context = request.context || {}; | ||
| try { | ||
| const res = await this.callHook("upgrade", request); | ||
| if (!res) return { | ||
| context, | ||
| namespace | ||
| }; | ||
| if (res.namespace) namespace = res.namespace; | ||
| if (res.context) Object.assign(context, res.context); | ||
| if (res instanceof Response) return { | ||
| context, | ||
| namespace, | ||
| endResponse: res | ||
| }; | ||
| if (res.headers) return { | ||
| context, | ||
| namespace, | ||
| upgradeHeaders: res.headers | ||
| }; | ||
| } catch (error) { | ||
| const errResponse = error.response || error; | ||
| if (errResponse instanceof Response) return { | ||
| context, | ||
| namespace, | ||
| endResponse: errResponse | ||
| }; | ||
| throw error; | ||
| } | ||
| return { | ||
| context, | ||
| namespace | ||
| }; | ||
| } | ||
| }; | ||
| function defineHooks(hooks) { | ||
| return hooks; | ||
| } | ||
| //#endregion | ||
| //#region src/adapter.ts | ||
| function adapterUtils(globalPeers) { | ||
| return { | ||
| peers: globalPeers, | ||
| publish(topic, message, options) { | ||
| for (const peers of options?.namespace ? [globalPeers.get(options.namespace) || []] : globalPeers.values()) { | ||
| let firstPeerWithTopic; | ||
| for (const peer of peers) if (peer.topics.has(topic)) { | ||
| firstPeerWithTopic = peer; | ||
| break; | ||
| } | ||
| if (firstPeerWithTopic) { | ||
| firstPeerWithTopic.send(message, options); | ||
| firstPeerWithTopic.publish(topic, message, options); | ||
| } | ||
| } | ||
| } | ||
| }; | ||
| } | ||
| function getPeers(globalPeers, namespace) { | ||
| if (!namespace) throw new Error("Websocket publish namespace missing."); | ||
| let peers = globalPeers.get(namespace); | ||
| if (!peers) { | ||
| peers = /* @__PURE__ */ new Set(); | ||
| globalPeers.set(namespace, peers); | ||
| } | ||
| return peers; | ||
| } | ||
| function defineWebSocketAdapter(factory) { | ||
| return factory; | ||
| } | ||
| //#endregion | ||
| export { defineHooks as a, AdapterHookable as i, defineWebSocketAdapter as n, getPeers as r, adapterUtils as t }; |
| import { d as PeerContext, n as AdapterInstance, r as AdapterOptions, t as Adapter, u as Peer } from "./adapter.mjs"; | ||
| import { Server, ServerWebSocket, WebSocketHandler } from "bun"; | ||
| //#region src/adapters/bun.d.ts | ||
| interface BunAdapter extends AdapterInstance { | ||
| websocket: WebSocketHandler<ContextData>; | ||
| handleUpgrade(req: Request, server: Server<ContextData>): Promise<Response | undefined>; | ||
| } | ||
| interface BunOptions extends AdapterOptions {} | ||
| type ContextData = { | ||
| peer?: BunPeer; | ||
| namespace: string; | ||
| request: Request; | ||
| server?: Server<ContextData>; | ||
| context: PeerContext; | ||
| }; | ||
| declare const bunAdapter: Adapter<BunAdapter, BunOptions>; | ||
| declare class BunPeer extends Peer<{ | ||
| ws: ServerWebSocket<ContextData>; | ||
| namespace: string; | ||
| request: Request; | ||
| peers: Set<BunPeer>; | ||
| }> { | ||
| get remoteAddress(): string; | ||
| get context(): PeerContext; | ||
| send(data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| publish(topic: string, data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| subscribe(topic: string): void; | ||
| unsubscribe(topic: string): void; | ||
| close(code?: number, reason?: string): void; | ||
| terminate(): void; | ||
| } | ||
| //#endregion | ||
| export { BunOptions as n, bunAdapter as r, BunAdapter as t }; |
| import { n as AdapterInstance, r as AdapterOptions, t as Adapter } from "./adapter.mjs"; | ||
| import { a as WebSocket$1 } from "./web.mjs"; | ||
| import { DurableObject } from "cloudflare:workers"; | ||
| import * as CF from "@cloudflare/workers-types"; | ||
| //#region src/adapters/cloudflare.d.ts | ||
| type WSDurableObjectStub = CF.DurableObjectStub & { | ||
| webSocketPublish?: (topic: string, data: unknown, opts: any) => Promise<void>; | ||
| }; | ||
| type ResolveDurableStub = (req: CF.Request | undefined, env: unknown, context: CF.ExecutionContext | undefined) => WSDurableObjectStub | undefined | Promise<WSDurableObjectStub | undefined>; | ||
| interface CloudflareOptions extends AdapterOptions { | ||
| /** | ||
| * Durable Object binding name from environment. | ||
| * | ||
| * **Note:** This option will be ignored if `resolveDurableStub` is provided. | ||
| * | ||
| * @default "$DurableObject" | ||
| */ | ||
| bindingName?: string; | ||
| /** | ||
| * Durable Object instance name. | ||
| * | ||
| * **Note:** This option will be ignored if `resolveDurableStub` is provided. | ||
| * | ||
| * @default "crossws" | ||
| */ | ||
| instanceName?: string; | ||
| /** | ||
| * Custom function that resolves Durable Object binding to handle the WebSocket upgrade. | ||
| * | ||
| * **Note:** This option will override `bindingName` and `instanceName`. | ||
| */ | ||
| resolveDurableStub?: ResolveDurableStub; | ||
| } | ||
| declare const cloudflareAdapter: Adapter<CloudflareDurableAdapter, CloudflareOptions>; | ||
| interface CloudflareDurableAdapter extends AdapterInstance { | ||
| handleUpgrade(req: Request | CF.Request, env: unknown, context: CF.ExecutionContext): Promise<Response>; | ||
| handleDurableInit(obj: DurableObject, state: DurableObjectState, env: unknown): void; | ||
| handleDurableUpgrade(obj: DurableObject, req: Request | CF.Request): Promise<Response>; | ||
| handleDurableMessage(obj: DurableObject, ws: WebSocket | CF.WebSocket | WebSocket$1, message: ArrayBuffer | string): Promise<void>; | ||
| handleDurablePublish: (obj: DurableObject, topic: string, data: unknown, opts: any) => Promise<void>; | ||
| handleDurableClose(obj: DurableObject, ws: WebSocket | CF.WebSocket | WebSocket$1, code: number, reason: string, wasClean: boolean): Promise<void>; | ||
| } | ||
| //#endregion | ||
| export { CloudflareOptions as n, cloudflareAdapter as r, CloudflareDurableAdapter as t }; |
| import { n as AdapterInstance, r as AdapterOptions, t as Adapter } from "./adapter.mjs"; | ||
| //#region src/adapters/deno.d.ts | ||
| interface DenoAdapter extends AdapterInstance { | ||
| handleUpgrade(req: Request, info: ServeHandlerInfo): Promise<Response>; | ||
| } | ||
| interface DenoOptions extends AdapterOptions {} | ||
| type ServeHandlerInfo = { | ||
| remoteAddr?: { | ||
| transport: string; | ||
| hostname: string; | ||
| port: number; | ||
| }; | ||
| }; | ||
| declare const denoAdapter: Adapter<DenoAdapter, DenoOptions>; | ||
| //#endregion | ||
| export { DenoOptions as n, denoAdapter as r, DenoAdapter as t }; |
| //#region src/error.ts | ||
| var WSError = class extends Error { | ||
| constructor(...args) { | ||
| super(...args); | ||
| this.name = "WSError"; | ||
| } | ||
| }; | ||
| //#endregion | ||
| export { WSError as t }; |
Sorry, the diff of this file is too big to display
| import { n as AdapterInstance, r as AdapterOptions, t as Adapter } from "./adapter.mjs"; | ||
| import { EventEmitter } from "events"; | ||
| import { Agent, ClientRequest, ClientRequestArgs, IncomingMessage, OutgoingHttpHeaders, Server } from "node:http"; | ||
| import { Duplex, DuplexOptions } from "node:stream"; | ||
| import { Server as Server$1 } from "node:https"; | ||
| import { SecureContextOptions } from "node:tls"; | ||
| import { URL } from "node:url"; | ||
| import { ZlibOptions } from "node:zlib"; | ||
| //#region types/ws.d.ts | ||
| type BufferLike = string | Buffer | DataView | number | ArrayBufferView | Uint8Array | ArrayBuffer | SharedArrayBuffer | readonly any[] | readonly number[] | { | ||
| valueOf(): ArrayBuffer; | ||
| } | { | ||
| valueOf(): SharedArrayBuffer; | ||
| } | { | ||
| valueOf(): Uint8Array; | ||
| } | { | ||
| valueOf(): readonly number[]; | ||
| } | { | ||
| valueOf(): string; | ||
| } | { | ||
| [Symbol.toPrimitive](hint: string): string; | ||
| }; | ||
| declare class WebSocket extends EventEmitter { | ||
| static readonly createWebSocketStream: typeof createWebSocketStream; | ||
| static readonly WebSocketServer: WebSocketServer; | ||
| static readonly Server: typeof Server$2; | ||
| static readonly WebSocket: typeof WebSocket; | ||
| /** The connection is not yet open. */ | ||
| static readonly CONNECTING: 0; | ||
| /** The connection is open and ready to communicate. */ | ||
| static readonly OPEN: 1; | ||
| /** The connection is in the process of closing. */ | ||
| static readonly CLOSING: 2; | ||
| /** The connection is closed. */ | ||
| static readonly CLOSED: 3; | ||
| binaryType: "nodebuffer" | "arraybuffer" | "fragments"; | ||
| readonly bufferedAmount: number; | ||
| readonly extensions: string; | ||
| /** Indicates whether the websocket is paused */ | ||
| readonly isPaused: boolean; | ||
| readonly protocol: string; | ||
| /** The current state of the connection */ | ||
| readonly readyState: typeof WebSocket.CONNECTING | typeof WebSocket.OPEN | typeof WebSocket.CLOSING | typeof WebSocket.CLOSED; | ||
| readonly url: string; | ||
| /** The connection is not yet open. */ | ||
| readonly CONNECTING: 0; | ||
| /** The connection is open and ready to communicate. */ | ||
| readonly OPEN: 1; | ||
| /** The connection is in the process of closing. */ | ||
| readonly CLOSING: 2; | ||
| /** The connection is closed. */ | ||
| readonly CLOSED: 3; | ||
| onopen: ((event: Event) => void) | null; | ||
| onerror: ((event: ErrorEvent) => void) | null; | ||
| onclose: ((event: CloseEvent) => void) | null; | ||
| onmessage: ((event: MessageEvent) => void) | null; | ||
| constructor(address: null); | ||
| constructor(address: string | URL, options?: ClientOptions | ClientRequestArgs); | ||
| constructor(address: string | URL, protocols?: string | string[], options?: ClientOptions | ClientRequestArgs); | ||
| close(code?: number, data?: string | Buffer): void; | ||
| ping(data?: any, mask?: boolean, cb?: (err: Error) => void): void; | ||
| pong(data?: any, mask?: boolean, cb?: (err: Error) => void): void; | ||
| send(data: BufferLike, cb?: (err?: Error) => void): void; | ||
| send(data: BufferLike, options: { | ||
| mask?: boolean | undefined; | ||
| binary?: boolean | undefined; | ||
| compress?: boolean | undefined; | ||
| fin?: boolean | undefined; | ||
| }, cb?: (err?: Error) => void): void; | ||
| terminate(): void; | ||
| /** | ||
| * Pause the websocket causing it to stop emitting events. Some events can still be | ||
| * emitted after this is called, until all buffered data is consumed. This method | ||
| * is a noop if the ready state is `CONNECTING` or `CLOSED`. | ||
| */ | ||
| pause(): void; | ||
| /** | ||
| * Make a paused socket resume emitting events. This method is a noop if the ready | ||
| * state is `CONNECTING` or `CLOSED`. | ||
| */ | ||
| resume(): void; | ||
| addEventListener(method: "message", cb: (event: MessageEvent) => void, options?: EventListenerOptions): void; | ||
| addEventListener(method: "close", cb: (event: CloseEvent) => void, options?: EventListenerOptions): void; | ||
| addEventListener(method: "error", cb: (event: ErrorEvent) => void, options?: EventListenerOptions): void; | ||
| addEventListener(method: "open", cb: (event: Event) => void, options?: EventListenerOptions): void; | ||
| removeEventListener(method: "message", cb: (event: MessageEvent) => void): void; | ||
| removeEventListener(method: "close", cb: (event: CloseEvent) => void): void; | ||
| removeEventListener(method: "error", cb: (event: ErrorEvent) => void): void; | ||
| removeEventListener(method: "open", cb: (event: Event) => void): void; | ||
| on(event: "close", listener: (this: WebSocket, code: number, reason: Buffer) => void): this; | ||
| on(event: "error", listener: (this: WebSocket, err: Error) => void): this; | ||
| on(event: "upgrade", listener: (this: WebSocket, request: IncomingMessage) => void): this; | ||
| on(event: "message", listener: (this: WebSocket, data: RawData, isBinary: boolean) => void): this; | ||
| on(event: "open", listener: (this: WebSocket) => void): this; | ||
| on(event: "ping" | "pong", listener: (this: WebSocket, data: Buffer) => void): this; | ||
| on(event: "unexpected-response", listener: (this: WebSocket, request: ClientRequest, response: IncomingMessage) => void): this; | ||
| on(event: string | symbol, listener: (this: WebSocket, ...args: any[]) => void): this; | ||
| once(event: "close", listener: (this: WebSocket, code: number, reason: Buffer) => void): this; | ||
| once(event: "error", listener: (this: WebSocket, err: Error) => void): this; | ||
| once(event: "upgrade", listener: (this: WebSocket, request: IncomingMessage) => void): this; | ||
| once(event: "message", listener: (this: WebSocket, data: RawData, isBinary: boolean) => void): this; | ||
| once(event: "open", listener: (this: WebSocket) => void): this; | ||
| once(event: "ping" | "pong", listener: (this: WebSocket, data: Buffer) => void): this; | ||
| once(event: "unexpected-response", listener: (this: WebSocket, request: ClientRequest, response: IncomingMessage) => void): this; | ||
| once(event: string | symbol, listener: (this: WebSocket, ...args: any[]) => void): this; | ||
| off(event: "close", listener: (this: WebSocket, code: number, reason: Buffer) => void): this; | ||
| off(event: "error", listener: (this: WebSocket, err: Error) => void): this; | ||
| off(event: "upgrade", listener: (this: WebSocket, request: IncomingMessage) => void): this; | ||
| off(event: "message", listener: (this: WebSocket, data: RawData, isBinary: boolean) => void): this; | ||
| off(event: "open", listener: (this: WebSocket) => void): this; | ||
| off(event: "ping" | "pong", listener: (this: WebSocket, data: Buffer) => void): this; | ||
| off(event: "unexpected-response", listener: (this: WebSocket, request: ClientRequest, response: IncomingMessage) => void): this; | ||
| off(event: string | symbol, listener: (this: WebSocket, ...args: any[]) => void): this; | ||
| addListener(event: "close", listener: (code: number, reason: Buffer) => void): this; | ||
| addListener(event: "error", listener: (err: Error) => void): this; | ||
| addListener(event: "upgrade", listener: (request: IncomingMessage) => void): this; | ||
| addListener(event: "message", listener: (data: RawData, isBinary: boolean) => void): this; | ||
| addListener(event: "open", listener: () => void): this; | ||
| addListener(event: "ping" | "pong", listener: (data: Buffer) => void): this; | ||
| addListener(event: "unexpected-response", listener: (request: ClientRequest, response: IncomingMessage) => void): this; | ||
| addListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| removeListener(event: "close", listener: (code: number, reason: Buffer) => void): this; | ||
| removeListener(event: "error", listener: (err: Error) => void): this; | ||
| removeListener(event: "upgrade", listener: (request: IncomingMessage) => void): this; | ||
| removeListener(event: "message", listener: (data: RawData, isBinary: boolean) => void): this; | ||
| removeListener(event: "open", listener: () => void): this; | ||
| removeListener(event: "ping" | "pong", listener: (data: Buffer) => void): this; | ||
| removeListener(event: "unexpected-response", listener: (request: ClientRequest, response: IncomingMessage) => void): this; | ||
| removeListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| } | ||
| /** | ||
| * Data represents the raw message payload received over the | ||
| */ | ||
| type RawData = Buffer | ArrayBuffer | Buffer[]; | ||
| /** | ||
| * Data represents the message payload received over the | ||
| */ | ||
| type Data = string | Buffer | ArrayBuffer | Buffer[]; | ||
| /** | ||
| * CertMeta represents the accepted types for certificate & key data. | ||
| */ | ||
| type CertMeta = string | string[] | Buffer | Buffer[]; | ||
| /** | ||
| * VerifyClientCallbackSync is a synchronous callback used to inspect the | ||
| * incoming message. The return value (boolean) of the function determines | ||
| * whether or not to accept the handshake. | ||
| */ | ||
| type VerifyClientCallbackSync<Request extends IncomingMessage = IncomingMessage> = (info: { | ||
| origin: string; | ||
| secure: boolean; | ||
| req: Request; | ||
| }) => boolean; | ||
| /** | ||
| * VerifyClientCallbackAsync is an asynchronous callback used to inspect the | ||
| * incoming message. The return value (boolean) of the function determines | ||
| * whether or not to accept the handshake. | ||
| */ | ||
| type VerifyClientCallbackAsync<Request extends IncomingMessage = IncomingMessage> = (info: { | ||
| origin: string; | ||
| secure: boolean; | ||
| req: Request; | ||
| }, callback: (res: boolean, code?: number, message?: string, headers?: OutgoingHttpHeaders) => void) => void; | ||
| interface ClientOptions extends SecureContextOptions { | ||
| protocol?: string | undefined; | ||
| followRedirects?: boolean | undefined; | ||
| generateMask?(mask: Buffer): void; | ||
| handshakeTimeout?: number | undefined; | ||
| maxRedirects?: number | undefined; | ||
| perMessageDeflate?: boolean | PerMessageDeflateOptions | undefined; | ||
| localAddress?: string | undefined; | ||
| protocolVersion?: number | undefined; | ||
| headers?: { | ||
| [key: string]: string; | ||
| } | undefined; | ||
| origin?: string | undefined; | ||
| agent?: Agent | undefined; | ||
| host?: string | undefined; | ||
| family?: number | undefined; | ||
| checkServerIdentity?(servername: string, cert: CertMeta): boolean; | ||
| rejectUnauthorized?: boolean | undefined; | ||
| maxPayload?: number | undefined; | ||
| skipUTF8Validation?: boolean | undefined; | ||
| } | ||
| interface PerMessageDeflateOptions { | ||
| serverNoContextTakeover?: boolean | undefined; | ||
| clientNoContextTakeover?: boolean | undefined; | ||
| serverMaxWindowBits?: number | undefined; | ||
| clientMaxWindowBits?: number | undefined; | ||
| zlibDeflateOptions?: { | ||
| flush?: number | undefined; | ||
| finishFlush?: number | undefined; | ||
| chunkSize?: number | undefined; | ||
| windowBits?: number | undefined; | ||
| level?: number | undefined; | ||
| memLevel?: number | undefined; | ||
| strategy?: number | undefined; | ||
| dictionary?: Buffer | Buffer[] | DataView | undefined; | ||
| info?: boolean | undefined; | ||
| } | undefined; | ||
| zlibInflateOptions?: ZlibOptions | undefined; | ||
| threshold?: number | undefined; | ||
| concurrencyLimit?: number | undefined; | ||
| } | ||
| interface Event { | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface ErrorEvent { | ||
| error: any; | ||
| message: string; | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface CloseEvent { | ||
| wasClean: boolean; | ||
| code: number; | ||
| reason: string; | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface MessageEvent { | ||
| data: Data; | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface EventListenerOptions { | ||
| once?: boolean | undefined; | ||
| } | ||
| interface ServerOptions<U extends typeof WebSocket = typeof WebSocket, V extends typeof IncomingMessage = typeof IncomingMessage> { | ||
| host?: string | undefined; | ||
| port?: number | undefined; | ||
| backlog?: number | undefined; | ||
| server?: Server<V> | Server$1<V> | undefined; | ||
| verifyClient?: VerifyClientCallbackAsync<InstanceType<V>> | VerifyClientCallbackSync<InstanceType<V>> | undefined; | ||
| handleProtocols?: (protocols: Set<string>, request: InstanceType<V>) => string | false; | ||
| path?: string | undefined; | ||
| noServer?: boolean | undefined; | ||
| clientTracking?: boolean | undefined; | ||
| perMessageDeflate?: boolean | PerMessageDeflateOptions | undefined; | ||
| maxPayload?: number | undefined; | ||
| skipUTF8Validation?: boolean | undefined; | ||
| WebSocket?: U | undefined; | ||
| } | ||
| interface AddressInfo { | ||
| address: string; | ||
| family: string; | ||
| port: number; | ||
| } | ||
| declare class Server$2<T extends typeof WebSocket = typeof WebSocket, U extends typeof IncomingMessage = typeof IncomingMessage> extends EventEmitter { | ||
| options: ServerOptions<T, U>; | ||
| path: string; | ||
| clients: Set<InstanceType<T>>; | ||
| constructor(options?: ServerOptions<T, U>, callback?: () => void); | ||
| address(): AddressInfo | string; | ||
| close(cb?: (err?: Error) => void): void; | ||
| handleUpgrade(request: InstanceType<U>, socket: Duplex, upgradeHead: Buffer, callback: (client: InstanceType<T>, request: InstanceType<U>) => void): void; | ||
| shouldHandle(request: InstanceType<U>): boolean | Promise<boolean>; | ||
| on(event: "connection", cb: (this: Server$2<T>, socket: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| on(event: "error", cb: (this: Server$2<T>, error: Error) => void): this; | ||
| on(event: "headers", cb: (this: Server$2<T>, headers: string[], request: InstanceType<U>) => void): this; | ||
| on(event: "close" | "listening", cb: (this: Server$2<T>) => void): this; | ||
| on(event: string | symbol, listener: (this: Server$2<T>, ...args: any[]) => void): this; | ||
| once(event: "connection", cb: (this: Server$2<T>, socket: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| once(event: "error", cb: (this: Server$2<T>, error: Error) => void): this; | ||
| once(event: "headers", cb: (this: Server$2<T>, headers: string[], request: InstanceType<U>) => void): this; | ||
| once(event: "close" | "listening", cb: (this: Server$2<T>) => void): this; | ||
| once(event: string | symbol, listener: (this: Server$2<T>, ...args: any[]) => void): this; | ||
| off(event: "connection", cb: (this: Server$2<T>, socket: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| off(event: "error", cb: (this: Server$2<T>, error: Error) => void): this; | ||
| off(event: "headers", cb: (this: Server$2<T>, headers: string[], request: InstanceType<U>) => void): this; | ||
| off(event: "close" | "listening", cb: (this: Server$2<T>) => void): this; | ||
| off(event: string | symbol, listener: (this: Server$2<T>, ...args: any[]) => void): this; | ||
| addListener(event: "connection", cb: (client: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| addListener(event: "error", cb: (err: Error) => void): this; | ||
| addListener(event: "headers", cb: (headers: string[], request: InstanceType<U>) => void): this; | ||
| addListener(event: "close" | "listening", cb: () => void): this; | ||
| addListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| removeListener(event: "connection", cb: (client: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| removeListener(event: "error", cb: (err: Error) => void): this; | ||
| removeListener(event: "headers", cb: (headers: string[], request: InstanceType<U>) => void): this; | ||
| removeListener(event: "close" | "listening", cb: () => void): this; | ||
| removeListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| } | ||
| type WebSocketServer = Server$2; | ||
| declare function createWebSocketStream(websocket: WebSocket, options?: DuplexOptions): Duplex; | ||
| //#endregion | ||
| //#region src/adapters/node.d.ts | ||
| interface NodeAdapter extends AdapterInstance { | ||
| handleUpgrade(req: IncomingMessage, socket: Duplex, head: Buffer, webRequest?: Request): Promise<void>; | ||
| closeAll: (code?: number, data?: string | Buffer, force?: boolean) => void; | ||
| } | ||
| interface NodeOptions extends AdapterOptions { | ||
| wss?: WebSocketServer; | ||
| serverOptions?: ServerOptions; | ||
| } | ||
| declare const nodeAdapter: Adapter<NodeAdapter, NodeOptions>; | ||
| //#endregion | ||
| export { NodeOptions as n, nodeAdapter as r, NodeAdapter as t }; |
| //#region src/utils.ts | ||
| const kNodeInspect = /* @__PURE__ */ Symbol.for("nodejs.util.inspect.custom"); | ||
| function toBufferLike(val) { | ||
| if (val === void 0 || val === null) return ""; | ||
| const type = typeof val; | ||
| if (type === "string") return val; | ||
| if (type === "number" || type === "boolean" || type === "bigint") return val.toString(); | ||
| if (type === "function" || type === "symbol") return "{}"; | ||
| if (val instanceof Uint8Array || val instanceof ArrayBuffer) return val; | ||
| if (isPlainObject(val)) return JSON.stringify(val); | ||
| return val; | ||
| } | ||
| function toString(val) { | ||
| if (typeof val === "string") return val; | ||
| const data = toBufferLike(val); | ||
| if (typeof data === "string") return data; | ||
| return `data:application/octet-stream;base64,${btoa(String.fromCharCode(...new Uint8Array(data)))}`; | ||
| } | ||
| function isPlainObject(value) { | ||
| if (value === null || typeof value !== "object") return false; | ||
| const prototype = Object.getPrototypeOf(value); | ||
| if (prototype !== null && prototype !== Object.prototype && Object.getPrototypeOf(prototype) !== null) return false; | ||
| if (Symbol.iterator in value) return false; | ||
| if (Symbol.toStringTag in value) return Object.prototype.toString.call(value) === "[object Module]"; | ||
| return true; | ||
| } | ||
| //#endregion | ||
| //#region src/message.ts | ||
| var Message = class { | ||
| /** Access to the original [message event](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket/message_event) if available. */ | ||
| event; | ||
| /** Access to the Peer that emitted the message. */ | ||
| peer; | ||
| /** Raw message data (can be of any type). */ | ||
| rawData; | ||
| #id; | ||
| #uint8Array; | ||
| #arrayBuffer; | ||
| #blob; | ||
| #text; | ||
| #json; | ||
| constructor(rawData, peer, event) { | ||
| this.rawData = rawData || ""; | ||
| this.peer = peer; | ||
| this.event = event; | ||
| } | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the message. | ||
| */ | ||
| get id() { | ||
| if (!this.#id) this.#id = crypto.randomUUID(); | ||
| return this.#id; | ||
| } | ||
| /** | ||
| * Get data as [Uint8Array](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Uint8Array) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| uint8Array() { | ||
| const _uint8Array = this.#uint8Array; | ||
| if (_uint8Array) return _uint8Array; | ||
| const rawData = this.rawData; | ||
| if (rawData instanceof Uint8Array) return this.#uint8Array = rawData; | ||
| if (rawData instanceof ArrayBuffer || rawData instanceof SharedArrayBuffer) { | ||
| this.#arrayBuffer = rawData; | ||
| return this.#uint8Array = new Uint8Array(rawData); | ||
| } | ||
| if (typeof rawData === "string") { | ||
| this.#text = rawData; | ||
| return this.#uint8Array = new TextEncoder().encode(this.#text); | ||
| } | ||
| if (Symbol.iterator in rawData) return this.#uint8Array = new Uint8Array(rawData); | ||
| if (typeof rawData?.length === "number") return this.#uint8Array = new Uint8Array(rawData); | ||
| if (rawData instanceof DataView) return this.#uint8Array = new Uint8Array(rawData.buffer, rawData.byteOffset, rawData.byteLength); | ||
| throw new TypeError(`Unsupported message type: ${Object.prototype.toString.call(rawData)}`); | ||
| } | ||
| /** | ||
| * Get data as [ArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/ArrayBuffer) or [SharedArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/SharedArrayBuffer) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| arrayBuffer() { | ||
| const _arrayBuffer = this.#arrayBuffer; | ||
| if (_arrayBuffer) return _arrayBuffer; | ||
| const rawData = this.rawData; | ||
| if (rawData instanceof ArrayBuffer || rawData instanceof SharedArrayBuffer) return this.#arrayBuffer = rawData; | ||
| return this.#arrayBuffer = this.uint8Array().buffer; | ||
| } | ||
| /** | ||
| * Get data as [Blob](https://developer.mozilla.org/en-US/docs/Web/API/Blob) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. */ | ||
| blob() { | ||
| const _blob = this.#blob; | ||
| if (_blob) return _blob; | ||
| const rawData = this.rawData; | ||
| if (rawData instanceof Blob) return this.#blob = rawData; | ||
| return this.#blob = new Blob([this.uint8Array()]); | ||
| } | ||
| /** | ||
| * Get stringified text version of the message. | ||
| * | ||
| * If raw data is in any other format, it will be automatically converted and decoded. | ||
| */ | ||
| text() { | ||
| const _text = this.#text; | ||
| if (_text) return _text; | ||
| const rawData = this.rawData; | ||
| if (typeof rawData === "string") return this.#text = rawData; | ||
| return this.#text = new TextDecoder().decode(this.uint8Array()); | ||
| } | ||
| /** | ||
| * Get parsed version of the message text with [`JSON.parse()`](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/JSON/parse). | ||
| */ | ||
| json() { | ||
| const _json = this.#json; | ||
| if (_json) return _json; | ||
| return this.#json = JSON.parse(this.text()); | ||
| } | ||
| /** | ||
| * Message data (value varies based on `peer.websocket.binaryType`). | ||
| */ | ||
| get data() { | ||
| switch (this.peer?.websocket?.binaryType) { | ||
| case "arraybuffer": return this.arrayBuffer(); | ||
| case "blob": return this.blob(); | ||
| case "nodebuffer": return globalThis.Buffer ? Buffer.from(this.uint8Array()) : this.uint8Array(); | ||
| case "uint8array": return this.uint8Array(); | ||
| case "text": return this.text(); | ||
| default: return this.rawData; | ||
| } | ||
| } | ||
| toString() { | ||
| return this.text(); | ||
| } | ||
| [Symbol.toPrimitive]() { | ||
| return this.text(); | ||
| } | ||
| [kNodeInspect]() { | ||
| return { message: { | ||
| id: this.id, | ||
| peer: this.peer, | ||
| text: this.text() | ||
| } }; | ||
| } | ||
| }; | ||
| //#endregion | ||
| //#region src/peer.ts | ||
| var Peer = class { | ||
| _internal; | ||
| _topics; | ||
| _id; | ||
| #ws; | ||
| constructor(internal) { | ||
| this._topics = /* @__PURE__ */ new Set(); | ||
| this._internal = internal; | ||
| } | ||
| get context() { | ||
| return this._internal.context ??= {}; | ||
| } | ||
| get namespace() { | ||
| return this._internal.namespace; | ||
| } | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the peer. | ||
| */ | ||
| get id() { | ||
| if (!this._id) this._id = crypto.randomUUID(); | ||
| return this._id; | ||
| } | ||
| /** IP address of the peer */ | ||
| get remoteAddress() {} | ||
| /** upgrade request */ | ||
| get request() { | ||
| return this._internal.request; | ||
| } | ||
| /** | ||
| * Get the [WebSocket](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket) instance. | ||
| * | ||
| * **Note:** crossws adds polyfill for the following properties if native values are not available: | ||
| * - `protocol`: Extracted from the `sec-websocket-protocol` header. | ||
| * - `extensions`: Extracted from the `sec-websocket-extensions` header. | ||
| * - `url`: Extracted from the request URL (http -> ws). | ||
| * */ | ||
| get websocket() { | ||
| if (!this.#ws) { | ||
| const _ws = this._internal.ws; | ||
| const _request = this._internal.request; | ||
| this.#ws = _request ? createWsProxy(_ws, _request) : _ws; | ||
| } | ||
| return this.#ws; | ||
| } | ||
| /** All connected peers to the server */ | ||
| get peers() { | ||
| return this._internal.peers || /* @__PURE__ */ new Set(); | ||
| } | ||
| /** All topics, this peer has been subscribed to. */ | ||
| get topics() { | ||
| return this._topics; | ||
| } | ||
| /** Abruptly close the connection */ | ||
| terminate() { | ||
| this.close(); | ||
| } | ||
| /** Subscribe to a topic */ | ||
| subscribe(topic) { | ||
| this._topics.add(topic); | ||
| } | ||
| /** Unsubscribe from a topic */ | ||
| unsubscribe(topic) { | ||
| this._topics.delete(topic); | ||
| } | ||
| toString() { | ||
| return this.id; | ||
| } | ||
| [Symbol.toPrimitive]() { | ||
| return this.id; | ||
| } | ||
| [Symbol.toStringTag]() { | ||
| return "WebSocket"; | ||
| } | ||
| [kNodeInspect]() { | ||
| return { peer: { | ||
| id: this.id, | ||
| ip: this.remoteAddress | ||
| } }; | ||
| } | ||
| }; | ||
| function createWsProxy(ws, request) { | ||
| return new Proxy(ws, { get: (target, prop) => { | ||
| const value = Reflect.get(target, prop); | ||
| if (!value) switch (prop) { | ||
| case "protocol": return request?.headers?.get("sec-websocket-protocol") || ""; | ||
| case "extensions": return request?.headers?.get("sec-websocket-extensions") || ""; | ||
| case "url": return request?.url?.replace(/^http/, "ws") || void 0; | ||
| } | ||
| return value; | ||
| } }); | ||
| } | ||
| //#endregion | ||
| export { toString as i, Message as n, toBufferLike as r, Peer as t }; |
| import { createRequire } from "node:module"; | ||
| //#region rolldown:runtime | ||
| var __create = Object.create; | ||
| var __defProp = Object.defineProperty; | ||
| var __getOwnPropDesc = Object.getOwnPropertyDescriptor; | ||
| var __getOwnPropNames = Object.getOwnPropertyNames; | ||
| var __getProtoOf = Object.getPrototypeOf; | ||
| var __hasOwnProp = Object.prototype.hasOwnProperty; | ||
| var __commonJSMin = (cb, mod) => () => (mod || cb((mod = { exports: {} }).exports, mod), mod.exports); | ||
| var __copyProps = (to, from, except, desc) => { | ||
| if (from && typeof from === "object" || typeof from === "function") { | ||
| for (var keys = __getOwnPropNames(from), i = 0, n = keys.length, key; i < n; i++) { | ||
| key = keys[i]; | ||
| if (!__hasOwnProp.call(to, key) && key !== except) { | ||
| __defProp(to, key, { | ||
| get: ((k) => from[k]).bind(null, key), | ||
| enumerable: !(desc = __getOwnPropDesc(from, key)) || desc.enumerable | ||
| }); | ||
| } | ||
| } | ||
| } | ||
| return to; | ||
| }; | ||
| var __toESM = (mod, isNodeMode, target) => (target = mod != null ? __create(__getProtoOf(mod)) : {}, __copyProps(isNodeMode || !mod || !mod.__esModule ? __defProp(target, "default", { | ||
| value: mod, | ||
| enumerable: true | ||
| }) : target, mod)); | ||
| var __require = /* @__PURE__ */ createRequire(import.meta.url); | ||
| //#endregion | ||
| export { __require as n, __toESM as r, __commonJSMin as t }; |
| import { n as AdapterInstance, r as AdapterOptions, t as Adapter } from "./adapter.mjs"; | ||
| //#region src/adapters/sse.d.ts | ||
| interface SSEAdapter extends AdapterInstance { | ||
| fetch(req: Request): Promise<Response>; | ||
| } | ||
| interface SSEOptions extends AdapterOptions { | ||
| bidir?: boolean; | ||
| } | ||
| declare const sseAdapter: Adapter<SSEAdapter, SSEOptions>; | ||
| //#endregion | ||
| export { SSEOptions as n, sseAdapter as r, SSEAdapter as t }; |
| //#region types/web.d.ts | ||
| /** | ||
| * A CloseEvent is sent to clients using WebSockets when the connection is closed. This is delivered to the listener indicated by the WebSocket object's onclose attribute. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent) | ||
| */ | ||
| interface CloseEvent extends Event { | ||
| /** | ||
| * Returns the WebSocket connection close code provided by the server. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent/code) | ||
| */ | ||
| readonly code: number; | ||
| /** | ||
| * Returns the WebSocket connection close reason provided by the server. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent/reason) | ||
| */ | ||
| readonly reason: string; | ||
| /** | ||
| * Returns true if the connection closed cleanly; false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent/wasClean) | ||
| */ | ||
| readonly wasClean: boolean; | ||
| } | ||
| /** | ||
| * An event which takes place in the DOM. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event) | ||
| */ | ||
| interface Event { | ||
| /** | ||
| * Returns true or false depending on how event was initialized. True if event goes through its target's ancestors in reverse tree order, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/bubbles) | ||
| */ | ||
| readonly bubbles: boolean; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/cancelBubble) | ||
| */ | ||
| cancelBubble: boolean; | ||
| /** | ||
| * Returns true or false depending on how event was initialized. Its return value does not always carry meaning, but true can indicate that part of the operation during which event was dispatched, can be canceled by invoking the preventDefault() method. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/cancelable) | ||
| */ | ||
| readonly cancelable: boolean; | ||
| /** | ||
| * Returns true or false depending on how event was initialized. True if event invokes listeners past a ShadowRoot node that is the root of its target, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/composed) | ||
| */ | ||
| readonly composed: boolean; | ||
| /** | ||
| * Returns the object whose event listener's callback is currently being invoked. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/currentTarget) | ||
| */ | ||
| readonly currentTarget: EventTarget | null; | ||
| /** | ||
| * Returns true if preventDefault() was invoked successfully to indicate cancelation, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/defaultPrevented) | ||
| */ | ||
| readonly defaultPrevented: boolean; | ||
| /** | ||
| * Returns the event's phase, which is one of NONE, CAPTURING_PHASE, AT_TARGET, and BUBBLING_PHASE. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/eventPhase) | ||
| */ | ||
| readonly eventPhase: number; | ||
| /** | ||
| * Returns true if event was dispatched by the user agent, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/isTrusted) | ||
| */ | ||
| readonly isTrusted: boolean; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/returnValue) | ||
| */ | ||
| returnValue: boolean; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/srcElement) | ||
| */ | ||
| readonly srcElement: EventTarget | null; | ||
| /** | ||
| * Returns the object to which event is dispatched (its target). | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/target) | ||
| */ | ||
| readonly target: EventTarget | null; | ||
| /** | ||
| * Returns the event's timestamp as the number of milliseconds measured relative to the time origin. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/timeStamp) | ||
| */ | ||
| readonly timeStamp: DOMHighResTimeStamp; | ||
| /** | ||
| * Returns the type of event, e.g. "click", "hashchange", or "submit". | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/type) | ||
| */ | ||
| readonly type: string; | ||
| /** | ||
| * Returns the invocation target objects of event's path (objects on which listeners will be invoked), except for any nodes in shadow trees of which the shadow root's mode is "closed" that are not reachable from event's currentTarget. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/composedPath) | ||
| */ | ||
| composedPath(): EventTarget[]; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/initEvent) | ||
| */ | ||
| initEvent(type: string, bubbles?: boolean, cancelable?: boolean): void; | ||
| /** | ||
| * If invoked when the cancelable attribute value is true, and while executing a listener for the event with passive set to false, signals to the operation that caused event to be dispatched that it needs to be canceled. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/preventDefault) | ||
| */ | ||
| preventDefault(): void; | ||
| /** | ||
| * Invoking this method prevents event from reaching any registered event listeners after the current one finishes running and, when dispatched in a tree, also prevents event from reaching any other objects. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/stopImmediatePropagation) | ||
| */ | ||
| stopImmediatePropagation(): void; | ||
| /** | ||
| * When dispatched in a tree, invoking this method prevents event from reaching any objects other than the current object. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/stopPropagation) | ||
| */ | ||
| stopPropagation(): void; | ||
| readonly NONE: 0; | ||
| readonly CAPTURING_PHASE: 1; | ||
| readonly AT_TARGET: 2; | ||
| readonly BUBBLING_PHASE: 3; | ||
| } | ||
| /** | ||
| * EventTarget is a DOM interface implemented by objects that can receive events and may have listeners for them. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget) | ||
| */ | ||
| interface EventTarget { | ||
| /** | ||
| * Appends an event listener for events whose type attribute value is type. The callback argument sets the callback that will be invoked when the event is dispatched. | ||
| * | ||
| * The options argument sets listener-specific options. For compatibility this can be a boolean, in which case the method behaves exactly as if the value was specified as options's capture. | ||
| * | ||
| * When set to true, options's capture prevents callback from being invoked when the event's eventPhase attribute value is BUBBLING_PHASE. When false (or not present), callback will not be invoked when event's eventPhase attribute value is CAPTURING_PHASE. Either way, callback will be invoked if event's eventPhase attribute value is AT_TARGET. | ||
| * | ||
| * When set to true, options's passive indicates that the callback will not cancel the event by invoking preventDefault(). This is used to enable performance optimizations described in § 2.8 Observing event listeners. | ||
| * | ||
| * When set to true, options's once indicates that the callback will only be invoked once after which the event listener will be removed. | ||
| * | ||
| * If an AbortSignal is passed for options's signal, then the event listener will be removed when signal is aborted. | ||
| * | ||
| * The event listener is appended to target's event listener list and is not appended if it has the same type, callback, and capture. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget/addEventListener) | ||
| */ | ||
| addEventListener(type: string, callback: EventListenerOrEventListenerObject | null, options?: AddEventListenerOptions | boolean): void; | ||
| /** | ||
| * Dispatches a synthetic event event to target and returns true if either event's cancelable attribute value is false or its preventDefault() method was not invoked, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget/dispatchEvent) | ||
| */ | ||
| dispatchEvent(event: Event): boolean; | ||
| /** | ||
| * Removes the event listener in target's event listener list with the same type, callback, and options. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget/removeEventListener) | ||
| */ | ||
| removeEventListener(type: string, callback: EventListenerOrEventListenerObject | null, options?: EventListenerOptions | boolean): void; | ||
| } | ||
| /** | ||
| * A message received by a target object. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent) | ||
| */ | ||
| interface MessageEvent<T = any> extends Event { | ||
| /** | ||
| * Returns the data of the message. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/data) | ||
| */ | ||
| readonly data: T; | ||
| /** | ||
| * Returns the last event ID string, for server-sent events. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/lastEventId) | ||
| */ | ||
| readonly lastEventId: string; | ||
| /** | ||
| * Returns the origin of the message, for server-sent events and cross-document messaging. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/origin) | ||
| */ | ||
| readonly origin: string; | ||
| /** | ||
| * Returns the MessagePort array sent with the message, for cross-document messaging and channel messaging. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/ports) | ||
| */ | ||
| readonly ports: ReadonlyArray<MessagePort>; | ||
| /** | ||
| * Returns the WindowProxy of the source window, for cross-document messaging, and the MessagePort being attached, in the connect event fired at SharedWorkerGlobalScope objects. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/source) | ||
| */ | ||
| readonly source: MessageEventSource | null; | ||
| /** @deprecated */ | ||
| initMessageEvent(type: string, bubbles?: boolean, cancelable?: boolean, data?: any, origin?: string, lastEventId?: string, source?: MessageEventSource | null, ports?: MessagePort[]): void; | ||
| } | ||
| /** | ||
| * Provides the API for creating and managing a WebSocket connection to a server, as well as for sending and receiving data on the connection. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket) | ||
| */ | ||
| interface WebSocket extends EventTarget { | ||
| /** | ||
| * Returns a string that indicates how binary data from the WebSocket object is exposed to scripts: | ||
| * | ||
| * Can be set, to change how binary data is returned. The default is "blob". | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/binaryType) | ||
| */ | ||
| binaryType: BinaryType | (string & {}); | ||
| /** | ||
| * Returns the number of bytes of application data (UTF-8 text and binary data) that have been queued using send() but not yet been transmitted to the network. | ||
| * | ||
| * If the WebSocket connection is closed, this attribute's value will only increase with each call to the send() method. (The number does not reset to zero once the connection closes.) | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/bufferedAmount) | ||
| */ | ||
| readonly bufferedAmount: number; | ||
| /** | ||
| * Returns the extensions selected by the server, if any. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/extensions) | ||
| */ | ||
| readonly extensions: string; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/close_event) */ | ||
| onclose: ((this: WebSocket, ev: CloseEvent) => any) | null; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/error_event) */ | ||
| onerror: ((this: WebSocket, ev: Event) => any) | null; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/message_event) */ | ||
| onmessage: ((this: WebSocket, ev: MessageEvent) => any) | null; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/open_event) */ | ||
| onopen: ((this: WebSocket, ev: Event) => any) | null; | ||
| /** | ||
| * Returns the subprotocol selected by the server, if any. It can be used in conjunction with the array form of the constructor's second argument to perform subprotocol negotiation. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/protocol) | ||
| */ | ||
| readonly protocol: string; | ||
| /** | ||
| * Returns the state of the WebSocket object's connection. It can have the values described below. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/readyState) | ||
| */ | ||
| readonly readyState: number; | ||
| /** | ||
| * Returns the URL that was used to establish the WebSocket connection. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/url) | ||
| */ | ||
| readonly url: string; | ||
| /** | ||
| * Closes the WebSocket connection, optionally using code as the the WebSocket connection close code and reason as the the WebSocket connection close reason. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/close) | ||
| */ | ||
| close(code?: number, reason?: string): void; | ||
| /** | ||
| * Transmits data using the WebSocket connection. data can be a string, a Blob, an ArrayBuffer, or an ArrayBufferView. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/send) | ||
| */ | ||
| send(data: string | ArrayBufferLike | Blob | ArrayBufferView): void; | ||
| readonly CONNECTING: 0; | ||
| readonly OPEN: 1; | ||
| readonly CLOSING: 2; | ||
| readonly CLOSED: 3; | ||
| addEventListener<K extends keyof WebSocketEventMap>(type: K, listener: (this: WebSocket, ev: WebSocketEventMap[K]) => any, options?: boolean | AddEventListenerOptions): void; | ||
| addEventListener(type: string, listener: EventListenerOrEventListenerObject, options?: boolean | AddEventListenerOptions): void; | ||
| removeEventListener<K extends keyof WebSocketEventMap>(type: K, listener: (this: WebSocket, ev: WebSocketEventMap[K]) => any, options?: boolean | EventListenerOptions): void; | ||
| removeEventListener(type: string, listener: EventListenerOrEventListenerObject, options?: boolean | EventListenerOptions): void; | ||
| } | ||
| //#endregion | ||
| export { WebSocket as a, MessageEvent as i, Event as n, EventTarget as r, CloseEvent as t }; |
@@ -1,41 +0,2 @@ | ||
| import { WebSocketHandler, ServerWebSocket, Server } from 'bun'; | ||
| import { Adapter, AdapterInstance, Peer, PeerContext, AdapterOptions } from '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| interface BunAdapter extends AdapterInstance { | ||
| websocket: WebSocketHandler<ContextData>; | ||
| handleUpgrade(req: Request, server: Server<ContextData>): Promise<Response | undefined>; | ||
| } | ||
| interface BunOptions extends AdapterOptions { | ||
| } | ||
| type ContextData = { | ||
| peer?: BunPeer; | ||
| namespace: string; | ||
| request: Request; | ||
| server?: Server<ContextData>; | ||
| context: PeerContext; | ||
| }; | ||
| declare const bunAdapter: Adapter<BunAdapter, BunOptions>; | ||
| declare class BunPeer extends Peer<{ | ||
| ws: ServerWebSocket<ContextData>; | ||
| namespace: string; | ||
| request: Request; | ||
| peers: Set<BunPeer>; | ||
| }> { | ||
| get remoteAddress(): string; | ||
| get context(): PeerContext; | ||
| send(data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| publish(topic: string, data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| subscribe(topic: string): void; | ||
| unsubscribe(topic: string): void; | ||
| close(code?: number, reason?: string): void; | ||
| terminate(): void; | ||
| } | ||
| export { bunAdapter as default }; | ||
| export type { BunAdapter, BunOptions }; | ||
| import { n as BunOptions, r as bunAdapter, t as BunAdapter } from "../_chunks/bun.mjs"; | ||
| export { BunAdapter, BunOptions, bunAdapter as default }; |
+83
-93
@@ -1,99 +0,89 @@ | ||
| import { M as Message, P as Peer, t as toBufferLike } from '../shared/crossws.WpyOHUXc.mjs'; | ||
| import { a as adapterUtils, g as getPeers, A as AdapterHookable } from '../shared/crossws.CPlNx7g8.mjs'; | ||
| import { i as AdapterHookable, r as getPeers, t as adapterUtils } from "../_chunks/adapter.mjs"; | ||
| import { n as Message, r as toBufferLike, t as Peer } from "../_chunks/peer.mjs"; | ||
| //#region src/adapters/bun.ts | ||
| const bunAdapter = (options = {}) => { | ||
| if (typeof Bun === "undefined") { | ||
| throw new Error( | ||
| "[crossws] Using Bun adapter in an incompatible environment." | ||
| ); | ||
| } | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| async handleUpgrade(request, server) { | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) { | ||
| return endResponse; | ||
| } | ||
| const upgradeOK = server.upgrade(request, { | ||
| data: { | ||
| server, | ||
| request, | ||
| context, | ||
| namespace | ||
| }, | ||
| headers: upgradeHeaders | ||
| }); | ||
| if (!upgradeOK) { | ||
| return new Response("Upgrade failed", { status: 500 }); | ||
| } | ||
| }, | ||
| websocket: { | ||
| message: (ws, message) => { | ||
| const peers = getPeers(globalPeers, ws.data.namespace); | ||
| const peer = getPeer(ws, peers); | ||
| hooks.callHook("message", peer, new Message(message, peer)); | ||
| }, | ||
| open: (ws) => { | ||
| const peers = getPeers(globalPeers, ws.data.namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peers.add(peer); | ||
| hooks.callHook("open", peer); | ||
| }, | ||
| close: (ws, code, reason) => { | ||
| const peers = getPeers(globalPeers, ws.data.namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, { code, reason }); | ||
| } | ||
| } | ||
| }; | ||
| if (typeof Bun === "undefined") throw new Error("[crossws] Using Bun adapter in an incompatible environment."); | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| async handleUpgrade(request, server) { | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) return endResponse; | ||
| if (!server.upgrade(request, { | ||
| data: { | ||
| server, | ||
| request, | ||
| context, | ||
| namespace | ||
| }, | ||
| headers: upgradeHeaders | ||
| })) return new Response("Upgrade failed", { status: 500 }); | ||
| }, | ||
| websocket: { | ||
| message: (ws, message) => { | ||
| const peer = getPeer(ws, getPeers(globalPeers, ws.data.namespace)); | ||
| hooks.callHook("message", peer, new Message(message, peer)); | ||
| }, | ||
| open: (ws) => { | ||
| const peers = getPeers(globalPeers, ws.data.namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peers.add(peer); | ||
| hooks.callHook("open", peer); | ||
| }, | ||
| close: (ws, code, reason) => { | ||
| const peers = getPeers(globalPeers, ws.data.namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, { | ||
| code, | ||
| reason | ||
| }); | ||
| } | ||
| } | ||
| }; | ||
| }; | ||
| var bun_default = bunAdapter; | ||
| function getPeer(ws, peers) { | ||
| if (ws.data.peer) { | ||
| return ws.data.peer; | ||
| } | ||
| const peer = new BunPeer({ | ||
| ws, | ||
| request: ws.data.request, | ||
| peers, | ||
| namespace: ws.data.namespace | ||
| }); | ||
| ws.data.peer = peer; | ||
| return peer; | ||
| if (ws.data.peer) return ws.data.peer; | ||
| const peer = new BunPeer({ | ||
| ws, | ||
| request: ws.data.request, | ||
| peers, | ||
| namespace: ws.data.namespace | ||
| }); | ||
| ws.data.peer = peer; | ||
| return peer; | ||
| } | ||
| class BunPeer extends Peer { | ||
| get remoteAddress() { | ||
| return this._internal.ws.remoteAddress; | ||
| } | ||
| get context() { | ||
| return this._internal.ws.data.context; | ||
| } | ||
| send(data, options) { | ||
| return this._internal.ws.send(toBufferLike(data), options?.compress); | ||
| } | ||
| publish(topic, data, options) { | ||
| return this._internal.ws.publish( | ||
| topic, | ||
| toBufferLike(data), | ||
| options?.compress | ||
| ); | ||
| } | ||
| subscribe(topic) { | ||
| this._topics.add(topic); | ||
| this._internal.ws.subscribe(topic); | ||
| } | ||
| unsubscribe(topic) { | ||
| this._topics.delete(topic); | ||
| this._internal.ws.unsubscribe(topic); | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| terminate() { | ||
| this._internal.ws.terminate(); | ||
| } | ||
| } | ||
| var BunPeer = class extends Peer { | ||
| get remoteAddress() { | ||
| return this._internal.ws.remoteAddress; | ||
| } | ||
| get context() { | ||
| return this._internal.ws.data.context; | ||
| } | ||
| send(data, options) { | ||
| return this._internal.ws.send(toBufferLike(data), options?.compress); | ||
| } | ||
| publish(topic, data, options) { | ||
| return this._internal.ws.publish(topic, toBufferLike(data), options?.compress); | ||
| } | ||
| subscribe(topic) { | ||
| this._topics.add(topic); | ||
| this._internal.ws.subscribe(topic); | ||
| } | ||
| unsubscribe(topic) { | ||
| this._topics.delete(topic); | ||
| this._internal.ws.unsubscribe(topic); | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| terminate() { | ||
| this._internal.ws.terminate(); | ||
| } | ||
| }; | ||
| export { bunAdapter as default }; | ||
| //#endregion | ||
| export { bun_default as default }; |
@@ -1,46 +0,2 @@ | ||
| import * as CF from '@cloudflare/workers-types'; | ||
| import { DurableObject } from 'cloudflare:workers'; | ||
| import { Adapter, AdapterInstance, AdapterOptions } from '../index.mjs'; | ||
| import { W as WebSocket$1 } from '../shared/crossws.BQXMA5bH.mjs'; | ||
| type WSDurableObjectStub = CF.DurableObjectStub & { | ||
| webSocketPublish?: (topic: string, data: unknown, opts: any) => Promise<void>; | ||
| }; | ||
| type ResolveDurableStub = (req: CF.Request | undefined, env: unknown, context: CF.ExecutionContext | undefined) => WSDurableObjectStub | undefined | Promise<WSDurableObjectStub | undefined>; | ||
| interface CloudflareOptions extends AdapterOptions { | ||
| /** | ||
| * Durable Object binding name from environment. | ||
| * | ||
| * **Note:** This option will be ignored if `resolveDurableStub` is provided. | ||
| * | ||
| * @default "$DurableObject" | ||
| */ | ||
| bindingName?: string; | ||
| /** | ||
| * Durable Object instance name. | ||
| * | ||
| * **Note:** This option will be ignored if `resolveDurableStub` is provided. | ||
| * | ||
| * @default "crossws" | ||
| */ | ||
| instanceName?: string; | ||
| /** | ||
| * Custom function that resolves Durable Object binding to handle the WebSocket upgrade. | ||
| * | ||
| * **Note:** This option will override `bindingName` and `instanceName`. | ||
| */ | ||
| resolveDurableStub?: ResolveDurableStub; | ||
| } | ||
| declare const cloudflareAdapter: Adapter<CloudflareDurableAdapter, CloudflareOptions>; | ||
| interface CloudflareDurableAdapter extends AdapterInstance { | ||
| handleUpgrade(req: Request | CF.Request, env: unknown, context: CF.ExecutionContext): Promise<Response>; | ||
| handleDurableInit(obj: DurableObject, state: DurableObjectState, env: unknown): void; | ||
| handleDurableUpgrade(obj: DurableObject, req: Request | CF.Request): Promise<Response>; | ||
| handleDurableMessage(obj: DurableObject, ws: WebSocket | CF.WebSocket | WebSocket$1, message: ArrayBuffer | string): Promise<void>; | ||
| handleDurablePublish: (obj: DurableObject, topic: string, data: unknown, opts: any) => Promise<void>; | ||
| handleDurableClose(obj: DurableObject, ws: WebSocket | CF.WebSocket | WebSocket$1, code: number, reason: string, wasClean: boolean): Promise<void>; | ||
| } | ||
| export { cloudflareAdapter as default }; | ||
| export type { CloudflareDurableAdapter, CloudflareOptions }; | ||
| import { n as CloudflareOptions, r as cloudflareAdapter, t as CloudflareDurableAdapter } from "../_chunks/cloudflare.mjs"; | ||
| export { CloudflareDurableAdapter, CloudflareOptions, cloudflareAdapter as default }; |
+173
-219
@@ -1,227 +0,181 @@ | ||
| import { env } from 'cloudflare:workers'; | ||
| import { M as Message, P as Peer, t as toBufferLike } from '../shared/crossws.WpyOHUXc.mjs'; | ||
| import { a as adapterUtils, g as getPeers, A as AdapterHookable } from '../shared/crossws.CPlNx7g8.mjs'; | ||
| import { S as StubRequest } from '../shared/crossws.B31KJMcF.mjs'; | ||
| import { W as WSError } from '../shared/crossws.By9qWDAI.mjs'; | ||
| import { i as AdapterHookable, r as getPeers, t as adapterUtils } from "../_chunks/adapter.mjs"; | ||
| import { n as Message, r as toBufferLike, t as Peer } from "../_chunks/peer.mjs"; | ||
| import { t as StubRequest } from "../_chunks/_request.mjs"; | ||
| import { t as WSError } from "../_chunks/error.mjs"; | ||
| import { env } from "cloudflare:workers"; | ||
| //#region src/adapters/cloudflare.ts | ||
| const cloudflareAdapter = (opts = {}) => { | ||
| const hooks = new AdapterHookable(opts); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| const resolveDurableStub = opts.resolveDurableStub || ((_req, env$1, _context) => { | ||
| const bindingName = opts.bindingName || "$DurableObject"; | ||
| const binding = (env$1 || env)[bindingName]; | ||
| if (binding) { | ||
| const instanceId = binding.idFromName(opts.instanceName || "crossws"); | ||
| return binding.get(instanceId); | ||
| } | ||
| }); | ||
| const { publish: durablePublish, ...utils } = adapterUtils(globalPeers); | ||
| return { | ||
| ...utils, | ||
| handleUpgrade: async (request, cfEnv, cfCtx) => { | ||
| const stub = await resolveDurableStub( | ||
| request, | ||
| cfEnv, | ||
| cfCtx | ||
| ); | ||
| if (stub) { | ||
| return stub.fetch( | ||
| request | ||
| ); | ||
| } | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) { | ||
| return endResponse; | ||
| } | ||
| const peers = getPeers( | ||
| globalPeers, | ||
| namespace | ||
| ); | ||
| const pair = new WebSocketPair(); | ||
| const client = pair[0]; | ||
| const server = pair[1]; | ||
| const peer = new CloudflareFallbackPeer({ | ||
| ws: client, | ||
| peers, | ||
| wsServer: server, | ||
| request, | ||
| cfEnv, | ||
| cfCtx, | ||
| context, | ||
| namespace | ||
| }); | ||
| peers.add(peer); | ||
| server.accept(); | ||
| hooks.callHook("open", peer); | ||
| server.addEventListener("message", (event) => { | ||
| hooks.callHook( | ||
| "message", | ||
| peer, | ||
| new Message(event.data, peer, event) | ||
| ); | ||
| }); | ||
| server.addEventListener("error", (event) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("error", peer, new WSError(event.error)); | ||
| }); | ||
| server.addEventListener("close", (event) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, event); | ||
| server.close(); | ||
| }); | ||
| return new Response(null, { | ||
| status: 101, | ||
| webSocket: client, | ||
| headers: upgradeHeaders | ||
| }); | ||
| }, | ||
| handleDurableInit: async (obj, state, env) => { | ||
| }, | ||
| handleDurableUpgrade: async (obj, request) => { | ||
| const { upgradeHeaders, endResponse, namespace } = await hooks.upgrade( | ||
| request | ||
| ); | ||
| if (endResponse) { | ||
| return endResponse; | ||
| } | ||
| const peers = getPeers(globalPeers, namespace); | ||
| const pair = new WebSocketPair(); | ||
| const client = pair[0]; | ||
| const server = pair[1]; | ||
| const peer = CloudflareDurablePeer._restore( | ||
| obj, | ||
| server, | ||
| request, | ||
| namespace | ||
| ); | ||
| peers.add(peer); | ||
| obj.ctx.acceptWebSocket(server); | ||
| await hooks.callHook("open", peer); | ||
| return new Response(null, { | ||
| status: 101, | ||
| webSocket: client, | ||
| headers: upgradeHeaders | ||
| }); | ||
| }, | ||
| handleDurableMessage: async (obj, ws, message) => { | ||
| const peer = CloudflareDurablePeer._restore(obj, ws); | ||
| await hooks.callHook("message", peer, new Message(message, peer)); | ||
| }, | ||
| handleDurableClose: async (obj, ws, code, reason, wasClean) => { | ||
| const peer = CloudflareDurablePeer._restore(obj, ws); | ||
| const peers = getPeers(globalPeers, peer.namespace); | ||
| peers.delete(peer); | ||
| const details = { code, reason, wasClean }; | ||
| await hooks.callHook("close", peer, details); | ||
| }, | ||
| handleDurablePublish: async (_obj, topic, data, opts2) => { | ||
| return durablePublish(topic, data, opts2); | ||
| }, | ||
| publish: async (topic, data, opts2) => { | ||
| const stub = await resolveDurableStub(void 0, env, void 0); | ||
| if (!stub) { | ||
| throw new Error("[crossws] Durable Object binding cannot be resolved."); | ||
| } | ||
| try { | ||
| return await stub.webSocketPublish(topic, data, opts2); | ||
| } catch (error) { | ||
| console.error(error); | ||
| throw error; | ||
| } | ||
| } | ||
| }; | ||
| const hooks = new AdapterHookable(opts); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| const resolveDurableStub = opts.resolveDurableStub || ((_req, env$1, _context) => { | ||
| const bindingName = opts.bindingName || "$DurableObject"; | ||
| const binding = (env$1 || env)[bindingName]; | ||
| if (binding) { | ||
| const instanceId = binding.idFromName(opts.instanceName || "crossws"); | ||
| return binding.get(instanceId); | ||
| } | ||
| }); | ||
| const { publish: durablePublish, ...utils } = adapterUtils(globalPeers); | ||
| return { | ||
| ...utils, | ||
| handleUpgrade: async (request, cfEnv, cfCtx) => { | ||
| const stub = await resolveDurableStub(request, cfEnv, cfCtx); | ||
| if (stub) return stub.fetch(request); | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) return endResponse; | ||
| const peers = getPeers(globalPeers, namespace); | ||
| const pair = new WebSocketPair(); | ||
| const client = pair[0]; | ||
| const server = pair[1]; | ||
| const peer = new CloudflareFallbackPeer({ | ||
| ws: client, | ||
| peers, | ||
| wsServer: server, | ||
| request, | ||
| cfEnv, | ||
| cfCtx, | ||
| context, | ||
| namespace | ||
| }); | ||
| peers.add(peer); | ||
| server.accept(); | ||
| hooks.callHook("open", peer); | ||
| server.addEventListener("message", (event) => { | ||
| hooks.callHook("message", peer, new Message(event.data, peer, event)); | ||
| }); | ||
| server.addEventListener("error", (event) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("error", peer, new WSError(event.error)); | ||
| }); | ||
| server.addEventListener("close", (event) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, event); | ||
| server.close(); | ||
| }); | ||
| return new Response(null, { | ||
| status: 101, | ||
| webSocket: client, | ||
| headers: upgradeHeaders | ||
| }); | ||
| }, | ||
| handleDurableInit: async (obj, state, env$1) => {}, | ||
| handleDurableUpgrade: async (obj, request) => { | ||
| const { upgradeHeaders, endResponse, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) return endResponse; | ||
| const peers = getPeers(globalPeers, namespace); | ||
| const pair = new WebSocketPair(); | ||
| const client = pair[0]; | ||
| const server = pair[1]; | ||
| const peer = CloudflareDurablePeer._restore(obj, server, request, namespace); | ||
| peers.add(peer); | ||
| obj.ctx.acceptWebSocket(server); | ||
| await hooks.callHook("open", peer); | ||
| return new Response(null, { | ||
| status: 101, | ||
| webSocket: client, | ||
| headers: upgradeHeaders | ||
| }); | ||
| }, | ||
| handleDurableMessage: async (obj, ws, message) => { | ||
| const peer = CloudflareDurablePeer._restore(obj, ws); | ||
| await hooks.callHook("message", peer, new Message(message, peer)); | ||
| }, | ||
| handleDurableClose: async (obj, ws, code, reason, wasClean) => { | ||
| const peer = CloudflareDurablePeer._restore(obj, ws); | ||
| getPeers(globalPeers, peer.namespace).delete(peer); | ||
| const details = { | ||
| code, | ||
| reason, | ||
| wasClean | ||
| }; | ||
| await hooks.callHook("close", peer, details); | ||
| }, | ||
| handleDurablePublish: async (_obj, topic, data, opts$1) => { | ||
| return durablePublish(topic, data, opts$1); | ||
| }, | ||
| publish: async (topic, data, opts$1) => { | ||
| const stub = await resolveDurableStub(void 0, env, void 0); | ||
| if (!stub) throw new Error("[crossws] Durable Object binding cannot be resolved."); | ||
| try { | ||
| return await stub.webSocketPublish(topic, data, opts$1); | ||
| } catch (error) { | ||
| console.error(error); | ||
| throw error; | ||
| } | ||
| } | ||
| }; | ||
| }; | ||
| class CloudflareDurablePeer extends Peer { | ||
| get peers() { | ||
| return new Set( | ||
| this.#getwebsockets().map( | ||
| (ws) => CloudflareDurablePeer._restore(this._internal.durable, ws) | ||
| ) | ||
| ); | ||
| } | ||
| #getwebsockets() { | ||
| return this._internal.durable.ctx.getWebSockets(); | ||
| } | ||
| send(data) { | ||
| return this._internal.ws.send(toBufferLike(data)); | ||
| } | ||
| subscribe(topic) { | ||
| super.subscribe(topic); | ||
| const state = getAttachedState(this._internal.ws); | ||
| if (!state.t) { | ||
| state.t = /* @__PURE__ */ new Set(); | ||
| } | ||
| state.t.add(topic); | ||
| setAttachedState(this._internal.ws, state); | ||
| } | ||
| publish(topic, data) { | ||
| const websockets = this.#getwebsockets(); | ||
| if (websockets.length < 2) { | ||
| return; | ||
| } | ||
| const dataBuff = toBufferLike(data); | ||
| for (const ws of websockets) { | ||
| if (ws === this._internal.ws) { | ||
| continue; | ||
| } | ||
| const state = getAttachedState(ws); | ||
| if (state.t?.has(topic)) { | ||
| ws.send(dataBuff); | ||
| } | ||
| } | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| static _restore(durable, ws, request, namespace) { | ||
| let peer = ws._crosswsPeer; | ||
| if (peer) { | ||
| return peer; | ||
| } | ||
| const state = ws.deserializeAttachment() || {}; | ||
| peer = ws._crosswsPeer = new CloudflareDurablePeer({ | ||
| ws, | ||
| request: request || new StubRequest(state.u || ""), | ||
| namespace: namespace || state.n || "", | ||
| durable | ||
| }); | ||
| if (state.i) { | ||
| peer._id = state.i; | ||
| } | ||
| if (request?.url) { | ||
| state.u = request.url; | ||
| } | ||
| state.i = peer.id; | ||
| setAttachedState(ws, state); | ||
| return peer; | ||
| } | ||
| } | ||
| class CloudflareFallbackPeer extends Peer { | ||
| send(data) { | ||
| this._internal.wsServer.send(toBufferLike(data)); | ||
| return 0; | ||
| } | ||
| publish(_topic, _message) { | ||
| console.warn( | ||
| "[crossws] [cloudflare] pub/sub support requires Durable Objects." | ||
| ); | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| } | ||
| var cloudflare_default = cloudflareAdapter; | ||
| var CloudflareDurablePeer = class CloudflareDurablePeer extends Peer { | ||
| get peers() { | ||
| return new Set(this.#getwebsockets().map((ws) => CloudflareDurablePeer._restore(this._internal.durable, ws))); | ||
| } | ||
| #getwebsockets() { | ||
| return this._internal.durable.ctx.getWebSockets(); | ||
| } | ||
| send(data) { | ||
| return this._internal.ws.send(toBufferLike(data)); | ||
| } | ||
| subscribe(topic) { | ||
| super.subscribe(topic); | ||
| const state = getAttachedState(this._internal.ws); | ||
| if (!state.t) state.t = /* @__PURE__ */ new Set(); | ||
| state.t.add(topic); | ||
| setAttachedState(this._internal.ws, state); | ||
| } | ||
| publish(topic, data) { | ||
| const websockets = this.#getwebsockets(); | ||
| if (websockets.length < 2) return; | ||
| const dataBuff = toBufferLike(data); | ||
| for (const ws of websockets) { | ||
| if (ws === this._internal.ws) continue; | ||
| if (getAttachedState(ws).t?.has(topic)) ws.send(dataBuff); | ||
| } | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| static _restore(durable, ws, request, namespace) { | ||
| let peer = ws._crosswsPeer; | ||
| if (peer) return peer; | ||
| const state = ws.deserializeAttachment() || {}; | ||
| peer = ws._crosswsPeer = new CloudflareDurablePeer({ | ||
| ws, | ||
| request: request || new StubRequest(state.u || ""), | ||
| namespace: namespace || state.n || "", | ||
| durable | ||
| }); | ||
| if (state.i) peer._id = state.i; | ||
| if (request?.url) state.u = request.url; | ||
| state.i = peer.id; | ||
| setAttachedState(ws, state); | ||
| return peer; | ||
| } | ||
| }; | ||
| var CloudflareFallbackPeer = class extends Peer { | ||
| send(data) { | ||
| this._internal.wsServer.send(toBufferLike(data)); | ||
| return 0; | ||
| } | ||
| publish(_topic, _message) { | ||
| console.warn("[crossws] [cloudflare] pub/sub support requires Durable Objects."); | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| }; | ||
| function getAttachedState(ws) { | ||
| let state = ws._crosswsState; | ||
| if (state) { | ||
| return state; | ||
| } | ||
| state = ws.deserializeAttachment() || {}; | ||
| ws._crosswsState = state; | ||
| return state; | ||
| let state = ws._crosswsState; | ||
| if (state) return state; | ||
| state = ws.deserializeAttachment() || {}; | ||
| ws._crosswsState = state; | ||
| return state; | ||
| } | ||
| function setAttachedState(ws, state) { | ||
| ws._crosswsState = state; | ||
| ws.serializeAttachment(state); | ||
| ws._crosswsState = state; | ||
| ws.serializeAttachment(state); | ||
| } | ||
| export { cloudflareAdapter as default }; | ||
| //#endregion | ||
| export { cloudflare_default as default }; |
@@ -1,19 +0,2 @@ | ||
| import { Adapter, AdapterInstance, AdapterOptions } from '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| interface DenoAdapter extends AdapterInstance { | ||
| handleUpgrade(req: Request, info: ServeHandlerInfo): Promise<Response>; | ||
| } | ||
| interface DenoOptions extends AdapterOptions { | ||
| } | ||
| type ServeHandlerInfo = { | ||
| remoteAddr?: { | ||
| transport: string; | ||
| hostname: string; | ||
| port: number; | ||
| }; | ||
| }; | ||
| declare const denoAdapter: Adapter<DenoAdapter, DenoOptions>; | ||
| export { denoAdapter as default }; | ||
| export type { DenoAdapter, DenoOptions }; | ||
| import { n as DenoOptions, r as denoAdapter, t as DenoAdapter } from "../_chunks/deno.mjs"; | ||
| export { DenoAdapter, DenoOptions, denoAdapter as default }; |
+65
-74
@@ -1,78 +0,69 @@ | ||
| import { M as Message, P as Peer, t as toBufferLike } from '../shared/crossws.WpyOHUXc.mjs'; | ||
| import { a as adapterUtils, A as AdapterHookable, g as getPeers } from '../shared/crossws.CPlNx7g8.mjs'; | ||
| import { W as WSError } from '../shared/crossws.By9qWDAI.mjs'; | ||
| import { i as AdapterHookable, r as getPeers, t as adapterUtils } from "../_chunks/adapter.mjs"; | ||
| import { n as Message, r as toBufferLike, t as Peer } from "../_chunks/peer.mjs"; | ||
| import { t as WSError } from "../_chunks/error.mjs"; | ||
| //#region src/adapters/deno.ts | ||
| const denoAdapter = (options = {}) => { | ||
| if (typeof Deno === "undefined") { | ||
| throw new Error( | ||
| "[crossws] Using Deno adapter in an incompatible environment." | ||
| ); | ||
| } | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| handleUpgrade: async (request, info) => { | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) { | ||
| return endResponse; | ||
| } | ||
| const headers = upgradeHeaders instanceof Headers ? upgradeHeaders : new Headers(upgradeHeaders); | ||
| const upgrade = Deno.upgradeWebSocket(request, { | ||
| // @ts-expect-error Setting headers is currently not supported in Deno | ||
| // https://github.com/denoland/deno/issues/19277 | ||
| headers, | ||
| protocol: headers.get("sec-websocket-protocol") ?? "" | ||
| }); | ||
| const peers = getPeers(globalPeers, namespace); | ||
| const peer = new DenoPeer({ | ||
| ws: upgrade.socket, | ||
| request, | ||
| peers, | ||
| denoInfo: info, | ||
| context, | ||
| namespace | ||
| }); | ||
| peers.add(peer); | ||
| upgrade.socket.addEventListener("open", () => { | ||
| hooks.callHook("open", peer); | ||
| }); | ||
| upgrade.socket.addEventListener("message", (event) => { | ||
| hooks.callHook("message", peer, new Message(event.data, peer, event)); | ||
| }); | ||
| upgrade.socket.addEventListener("close", () => { | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, {}); | ||
| }); | ||
| upgrade.socket.addEventListener("error", (error) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("error", peer, new WSError(error)); | ||
| }); | ||
| return upgrade.response; | ||
| } | ||
| }; | ||
| if (typeof Deno === "undefined") throw new Error("[crossws] Using Deno adapter in an incompatible environment."); | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| handleUpgrade: async (request, info) => { | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) return endResponse; | ||
| const headers = upgradeHeaders instanceof Headers ? upgradeHeaders : new Headers(upgradeHeaders); | ||
| const upgrade = Deno.upgradeWebSocket(request, { | ||
| headers, | ||
| protocol: headers.get("sec-websocket-protocol") ?? "" | ||
| }); | ||
| const peers = getPeers(globalPeers, namespace); | ||
| const peer = new DenoPeer({ | ||
| ws: upgrade.socket, | ||
| request, | ||
| peers, | ||
| denoInfo: info, | ||
| context, | ||
| namespace | ||
| }); | ||
| peers.add(peer); | ||
| upgrade.socket.addEventListener("open", () => { | ||
| hooks.callHook("open", peer); | ||
| }); | ||
| upgrade.socket.addEventListener("message", (event) => { | ||
| hooks.callHook("message", peer, new Message(event.data, peer, event)); | ||
| }); | ||
| upgrade.socket.addEventListener("close", () => { | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, {}); | ||
| }); | ||
| upgrade.socket.addEventListener("error", (error) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("error", peer, new WSError(error)); | ||
| }); | ||
| return upgrade.response; | ||
| } | ||
| }; | ||
| }; | ||
| class DenoPeer extends Peer { | ||
| get remoteAddress() { | ||
| return this._internal.denoInfo.remoteAddr?.hostname; | ||
| } | ||
| send(data) { | ||
| return this._internal.ws.send(toBufferLike(data)); | ||
| } | ||
| publish(topic, data) { | ||
| const dataBuff = toBufferLike(data); | ||
| for (const peer of this._internal.peers) { | ||
| if (peer !== this && peer._topics.has(topic)) { | ||
| peer._internal.ws.send(dataBuff); | ||
| } | ||
| } | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| terminate() { | ||
| this._internal.ws.terminate(); | ||
| } | ||
| } | ||
| var deno_default = denoAdapter; | ||
| var DenoPeer = class extends Peer { | ||
| get remoteAddress() { | ||
| return this._internal.denoInfo.remoteAddr?.hostname; | ||
| } | ||
| send(data) { | ||
| return this._internal.ws.send(toBufferLike(data)); | ||
| } | ||
| publish(topic, data) { | ||
| const dataBuff = toBufferLike(data); | ||
| for (const peer of this._internal.peers) if (peer !== this && peer._topics.has(topic)) peer._internal.ws.send(dataBuff); | ||
| } | ||
| close(code, reason) { | ||
| this._internal.ws.close(code, reason); | ||
| } | ||
| terminate() { | ||
| this._internal.ws.terminate(); | ||
| } | ||
| }; | ||
| export { denoAdapter as default }; | ||
| //#endregion | ||
| export { deno_default as default }; |
+2
-299
@@ -1,299 +0,2 @@ | ||
| import { Adapter, AdapterInstance, AdapterOptions } from '../index.mjs'; | ||
| import { Agent, ClientRequestArgs, IncomingMessage, ClientRequest, Server as Server$1, OutgoingHttpHeaders } from 'node:http'; | ||
| import { DuplexOptions, Duplex } from 'node:stream'; | ||
| import { EventEmitter } from 'events'; | ||
| import { Server as Server$2 } from 'node:https'; | ||
| import { SecureContextOptions } from 'node:tls'; | ||
| import { URL } from 'node:url'; | ||
| import { ZlibOptions } from 'node:zlib'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| type BufferLike = string | Buffer | DataView | number | ArrayBufferView | Uint8Array | ArrayBuffer | SharedArrayBuffer | readonly any[] | readonly number[] | { | ||
| valueOf(): ArrayBuffer; | ||
| } | { | ||
| valueOf(): SharedArrayBuffer; | ||
| } | { | ||
| valueOf(): Uint8Array; | ||
| } | { | ||
| valueOf(): readonly number[]; | ||
| } | { | ||
| valueOf(): string; | ||
| } | { | ||
| [Symbol.toPrimitive](hint: string): string; | ||
| }; | ||
| declare class WebSocket extends EventEmitter { | ||
| static readonly createWebSocketStream: typeof createWebSocketStream; | ||
| static readonly WebSocketServer: WebSocketServer; | ||
| static readonly Server: typeof Server; | ||
| static readonly WebSocket: typeof WebSocket; | ||
| /** The connection is not yet open. */ | ||
| static readonly CONNECTING: 0; | ||
| /** The connection is open and ready to communicate. */ | ||
| static readonly OPEN: 1; | ||
| /** The connection is in the process of closing. */ | ||
| static readonly CLOSING: 2; | ||
| /** The connection is closed. */ | ||
| static readonly CLOSED: 3; | ||
| binaryType: "nodebuffer" | "arraybuffer" | "fragments"; | ||
| readonly bufferedAmount: number; | ||
| readonly extensions: string; | ||
| /** Indicates whether the websocket is paused */ | ||
| readonly isPaused: boolean; | ||
| readonly protocol: string; | ||
| /** The current state of the connection */ | ||
| readonly readyState: typeof WebSocket.CONNECTING | typeof WebSocket.OPEN | typeof WebSocket.CLOSING | typeof WebSocket.CLOSED; | ||
| readonly url: string; | ||
| /** The connection is not yet open. */ | ||
| readonly CONNECTING: 0; | ||
| /** The connection is open and ready to communicate. */ | ||
| readonly OPEN: 1; | ||
| /** The connection is in the process of closing. */ | ||
| readonly CLOSING: 2; | ||
| /** The connection is closed. */ | ||
| readonly CLOSED: 3; | ||
| onopen: ((event: Event) => void) | null; | ||
| onerror: ((event: ErrorEvent) => void) | null; | ||
| onclose: ((event: CloseEvent) => void) | null; | ||
| onmessage: ((event: MessageEvent) => void) | null; | ||
| constructor(address: null); | ||
| constructor(address: string | URL, options?: ClientOptions | ClientRequestArgs); | ||
| constructor(address: string | URL, protocols?: string | string[], options?: ClientOptions | ClientRequestArgs); | ||
| close(code?: number, data?: string | Buffer): void; | ||
| ping(data?: any, mask?: boolean, cb?: (err: Error) => void): void; | ||
| pong(data?: any, mask?: boolean, cb?: (err: Error) => void): void; | ||
| send(data: BufferLike, cb?: (err?: Error) => void): void; | ||
| send(data: BufferLike, options: { | ||
| mask?: boolean | undefined; | ||
| binary?: boolean | undefined; | ||
| compress?: boolean | undefined; | ||
| fin?: boolean | undefined; | ||
| }, cb?: (err?: Error) => void): void; | ||
| terminate(): void; | ||
| /** | ||
| * Pause the websocket causing it to stop emitting events. Some events can still be | ||
| * emitted after this is called, until all buffered data is consumed. This method | ||
| * is a noop if the ready state is `CONNECTING` or `CLOSED`. | ||
| */ | ||
| pause(): void; | ||
| /** | ||
| * Make a paused socket resume emitting events. This method is a noop if the ready | ||
| * state is `CONNECTING` or `CLOSED`. | ||
| */ | ||
| resume(): void; | ||
| addEventListener(method: "message", cb: (event: MessageEvent) => void, options?: EventListenerOptions): void; | ||
| addEventListener(method: "close", cb: (event: CloseEvent) => void, options?: EventListenerOptions): void; | ||
| addEventListener(method: "error", cb: (event: ErrorEvent) => void, options?: EventListenerOptions): void; | ||
| addEventListener(method: "open", cb: (event: Event) => void, options?: EventListenerOptions): void; | ||
| removeEventListener(method: "message", cb: (event: MessageEvent) => void): void; | ||
| removeEventListener(method: "close", cb: (event: CloseEvent) => void): void; | ||
| removeEventListener(method: "error", cb: (event: ErrorEvent) => void): void; | ||
| removeEventListener(method: "open", cb: (event: Event) => void): void; | ||
| on(event: "close", listener: (this: WebSocket, code: number, reason: Buffer) => void): this; | ||
| on(event: "error", listener: (this: WebSocket, err: Error) => void): this; | ||
| on(event: "upgrade", listener: (this: WebSocket, request: IncomingMessage) => void): this; | ||
| on(event: "message", listener: (this: WebSocket, data: RawData, isBinary: boolean) => void): this; | ||
| on(event: "open", listener: (this: WebSocket) => void): this; | ||
| on(event: "ping" | "pong", listener: (this: WebSocket, data: Buffer) => void): this; | ||
| on(event: "unexpected-response", listener: (this: WebSocket, request: ClientRequest, response: IncomingMessage) => void): this; | ||
| on(event: string | symbol, listener: (this: WebSocket, ...args: any[]) => void): this; | ||
| once(event: "close", listener: (this: WebSocket, code: number, reason: Buffer) => void): this; | ||
| once(event: "error", listener: (this: WebSocket, err: Error) => void): this; | ||
| once(event: "upgrade", listener: (this: WebSocket, request: IncomingMessage) => void): this; | ||
| once(event: "message", listener: (this: WebSocket, data: RawData, isBinary: boolean) => void): this; | ||
| once(event: "open", listener: (this: WebSocket) => void): this; | ||
| once(event: "ping" | "pong", listener: (this: WebSocket, data: Buffer) => void): this; | ||
| once(event: "unexpected-response", listener: (this: WebSocket, request: ClientRequest, response: IncomingMessage) => void): this; | ||
| once(event: string | symbol, listener: (this: WebSocket, ...args: any[]) => void): this; | ||
| off(event: "close", listener: (this: WebSocket, code: number, reason: Buffer) => void): this; | ||
| off(event: "error", listener: (this: WebSocket, err: Error) => void): this; | ||
| off(event: "upgrade", listener: (this: WebSocket, request: IncomingMessage) => void): this; | ||
| off(event: "message", listener: (this: WebSocket, data: RawData, isBinary: boolean) => void): this; | ||
| off(event: "open", listener: (this: WebSocket) => void): this; | ||
| off(event: "ping" | "pong", listener: (this: WebSocket, data: Buffer) => void): this; | ||
| off(event: "unexpected-response", listener: (this: WebSocket, request: ClientRequest, response: IncomingMessage) => void): this; | ||
| off(event: string | symbol, listener: (this: WebSocket, ...args: any[]) => void): this; | ||
| addListener(event: "close", listener: (code: number, reason: Buffer) => void): this; | ||
| addListener(event: "error", listener: (err: Error) => void): this; | ||
| addListener(event: "upgrade", listener: (request: IncomingMessage) => void): this; | ||
| addListener(event: "message", listener: (data: RawData, isBinary: boolean) => void): this; | ||
| addListener(event: "open", listener: () => void): this; | ||
| addListener(event: "ping" | "pong", listener: (data: Buffer) => void): this; | ||
| addListener(event: "unexpected-response", listener: (request: ClientRequest, response: IncomingMessage) => void): this; | ||
| addListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| removeListener(event: "close", listener: (code: number, reason: Buffer) => void): this; | ||
| removeListener(event: "error", listener: (err: Error) => void): this; | ||
| removeListener(event: "upgrade", listener: (request: IncomingMessage) => void): this; | ||
| removeListener(event: "message", listener: (data: RawData, isBinary: boolean) => void): this; | ||
| removeListener(event: "open", listener: () => void): this; | ||
| removeListener(event: "ping" | "pong", listener: (data: Buffer) => void): this; | ||
| removeListener(event: "unexpected-response", listener: (request: ClientRequest, response: IncomingMessage) => void): this; | ||
| removeListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| } | ||
| /** | ||
| * Data represents the raw message payload received over the | ||
| */ | ||
| type RawData = Buffer | ArrayBuffer | Buffer[]; | ||
| /** | ||
| * Data represents the message payload received over the | ||
| */ | ||
| type Data = string | Buffer | ArrayBuffer | Buffer[]; | ||
| /** | ||
| * CertMeta represents the accepted types for certificate & key data. | ||
| */ | ||
| type CertMeta = string | string[] | Buffer | Buffer[]; | ||
| /** | ||
| * VerifyClientCallbackSync is a synchronous callback used to inspect the | ||
| * incoming message. The return value (boolean) of the function determines | ||
| * whether or not to accept the handshake. | ||
| */ | ||
| type VerifyClientCallbackSync<Request extends IncomingMessage = IncomingMessage> = (info: { | ||
| origin: string; | ||
| secure: boolean; | ||
| req: Request; | ||
| }) => boolean; | ||
| /** | ||
| * VerifyClientCallbackAsync is an asynchronous callback used to inspect the | ||
| * incoming message. The return value (boolean) of the function determines | ||
| * whether or not to accept the handshake. | ||
| */ | ||
| type VerifyClientCallbackAsync<Request extends IncomingMessage = IncomingMessage> = (info: { | ||
| origin: string; | ||
| secure: boolean; | ||
| req: Request; | ||
| }, callback: (res: boolean, code?: number, message?: string, headers?: OutgoingHttpHeaders) => void) => void; | ||
| interface ClientOptions extends SecureContextOptions { | ||
| protocol?: string | undefined; | ||
| followRedirects?: boolean | undefined; | ||
| generateMask?(mask: Buffer): void; | ||
| handshakeTimeout?: number | undefined; | ||
| maxRedirects?: number | undefined; | ||
| perMessageDeflate?: boolean | PerMessageDeflateOptions | undefined; | ||
| localAddress?: string | undefined; | ||
| protocolVersion?: number | undefined; | ||
| headers?: { | ||
| [key: string]: string; | ||
| } | undefined; | ||
| origin?: string | undefined; | ||
| agent?: Agent | undefined; | ||
| host?: string | undefined; | ||
| family?: number | undefined; | ||
| checkServerIdentity?(servername: string, cert: CertMeta): boolean; | ||
| rejectUnauthorized?: boolean | undefined; | ||
| maxPayload?: number | undefined; | ||
| skipUTF8Validation?: boolean | undefined; | ||
| } | ||
| interface PerMessageDeflateOptions { | ||
| serverNoContextTakeover?: boolean | undefined; | ||
| clientNoContextTakeover?: boolean | undefined; | ||
| serverMaxWindowBits?: number | undefined; | ||
| clientMaxWindowBits?: number | undefined; | ||
| zlibDeflateOptions?: { | ||
| flush?: number | undefined; | ||
| finishFlush?: number | undefined; | ||
| chunkSize?: number | undefined; | ||
| windowBits?: number | undefined; | ||
| level?: number | undefined; | ||
| memLevel?: number | undefined; | ||
| strategy?: number | undefined; | ||
| dictionary?: Buffer | Buffer[] | DataView | undefined; | ||
| info?: boolean | undefined; | ||
| } | undefined; | ||
| zlibInflateOptions?: ZlibOptions | undefined; | ||
| threshold?: number | undefined; | ||
| concurrencyLimit?: number | undefined; | ||
| } | ||
| interface Event { | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface ErrorEvent { | ||
| error: any; | ||
| message: string; | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface CloseEvent { | ||
| wasClean: boolean; | ||
| code: number; | ||
| reason: string; | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface MessageEvent { | ||
| data: Data; | ||
| type: string; | ||
| target: WebSocket; | ||
| } | ||
| interface EventListenerOptions { | ||
| once?: boolean | undefined; | ||
| } | ||
| interface ServerOptions<U extends typeof WebSocket = typeof WebSocket, V extends typeof IncomingMessage = typeof IncomingMessage> { | ||
| host?: string | undefined; | ||
| port?: number | undefined; | ||
| backlog?: number | undefined; | ||
| server?: Server$1<V> | Server$2<V> | undefined; | ||
| verifyClient?: VerifyClientCallbackAsync<InstanceType<V>> | VerifyClientCallbackSync<InstanceType<V>> | undefined; | ||
| handleProtocols?: (protocols: Set<string>, request: InstanceType<V>) => string | false; | ||
| path?: string | undefined; | ||
| noServer?: boolean | undefined; | ||
| clientTracking?: boolean | undefined; | ||
| perMessageDeflate?: boolean | PerMessageDeflateOptions | undefined; | ||
| maxPayload?: number | undefined; | ||
| skipUTF8Validation?: boolean | undefined; | ||
| WebSocket?: U | undefined; | ||
| } | ||
| interface AddressInfo { | ||
| address: string; | ||
| family: string; | ||
| port: number; | ||
| } | ||
| declare class Server<T extends typeof WebSocket = typeof WebSocket, U extends typeof IncomingMessage = typeof IncomingMessage> extends EventEmitter { | ||
| options: ServerOptions<T, U>; | ||
| path: string; | ||
| clients: Set<InstanceType<T>>; | ||
| constructor(options?: ServerOptions<T, U>, callback?: () => void); | ||
| address(): AddressInfo | string; | ||
| close(cb?: (err?: Error) => void): void; | ||
| handleUpgrade(request: InstanceType<U>, socket: Duplex, upgradeHead: Buffer, callback: (client: InstanceType<T>, request: InstanceType<U>) => void): void; | ||
| shouldHandle(request: InstanceType<U>): boolean | Promise<boolean>; | ||
| on(event: "connection", cb: (this: Server<T>, socket: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| on(event: "error", cb: (this: Server<T>, error: Error) => void): this; | ||
| on(event: "headers", cb: (this: Server<T>, headers: string[], request: InstanceType<U>) => void): this; | ||
| on(event: "close" | "listening", cb: (this: Server<T>) => void): this; | ||
| on(event: string | symbol, listener: (this: Server<T>, ...args: any[]) => void): this; | ||
| once(event: "connection", cb: (this: Server<T>, socket: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| once(event: "error", cb: (this: Server<T>, error: Error) => void): this; | ||
| once(event: "headers", cb: (this: Server<T>, headers: string[], request: InstanceType<U>) => void): this; | ||
| once(event: "close" | "listening", cb: (this: Server<T>) => void): this; | ||
| once(event: string | symbol, listener: (this: Server<T>, ...args: any[]) => void): this; | ||
| off(event: "connection", cb: (this: Server<T>, socket: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| off(event: "error", cb: (this: Server<T>, error: Error) => void): this; | ||
| off(event: "headers", cb: (this: Server<T>, headers: string[], request: InstanceType<U>) => void): this; | ||
| off(event: "close" | "listening", cb: (this: Server<T>) => void): this; | ||
| off(event: string | symbol, listener: (this: Server<T>, ...args: any[]) => void): this; | ||
| addListener(event: "connection", cb: (client: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| addListener(event: "error", cb: (err: Error) => void): this; | ||
| addListener(event: "headers", cb: (headers: string[], request: InstanceType<U>) => void): this; | ||
| addListener(event: "close" | "listening", cb: () => void): this; | ||
| addListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| removeListener(event: "connection", cb: (client: InstanceType<T>, request: InstanceType<U>) => void): this; | ||
| removeListener(event: "error", cb: (err: Error) => void): this; | ||
| removeListener(event: "headers", cb: (headers: string[], request: InstanceType<U>) => void): this; | ||
| removeListener(event: "close" | "listening", cb: () => void): this; | ||
| removeListener(event: string | symbol, listener: (...args: any[]) => void): this; | ||
| } | ||
| type WebSocketServer = Server; | ||
| declare function createWebSocketStream(websocket: WebSocket, options?: DuplexOptions): Duplex; | ||
| interface NodeAdapter extends AdapterInstance { | ||
| handleUpgrade(req: IncomingMessage, socket: Duplex, head: Buffer, webRequest?: Request): Promise<void>; | ||
| closeAll: (code?: number, data?: string | Buffer, force?: boolean) => void; | ||
| } | ||
| interface NodeOptions extends AdapterOptions { | ||
| wss?: WebSocketServer; | ||
| serverOptions?: ServerOptions; | ||
| } | ||
| declare const nodeAdapter: Adapter<NodeAdapter, NodeOptions>; | ||
| export { nodeAdapter as default }; | ||
| export type { NodeAdapter, NodeOptions }; | ||
| import { n as NodeOptions, r as nodeAdapter, t as NodeAdapter } from "../_chunks/node.mjs"; | ||
| export { NodeAdapter, NodeOptions, nodeAdapter as default }; |
+119
-156
@@ -1,162 +0,125 @@ | ||
| import { M as Message, P as Peer, t as toBufferLike } from '../shared/crossws.WpyOHUXc.mjs'; | ||
| import { g as getPeers, A as AdapterHookable, a as adapterUtils } from '../shared/crossws.CPlNx7g8.mjs'; | ||
| import { W as WSError } from '../shared/crossws.By9qWDAI.mjs'; | ||
| import { _ as _WebSocketServer } from '../shared/crossws.C5pESzqN.mjs'; | ||
| import { S as StubRequest } from '../shared/crossws.B31KJMcF.mjs'; | ||
| import 'stream'; | ||
| import 'events'; | ||
| import 'http'; | ||
| import 'crypto'; | ||
| import 'buffer'; | ||
| import 'zlib'; | ||
| import 'https'; | ||
| import 'net'; | ||
| import 'tls'; | ||
| import 'url'; | ||
| import "../_chunks/rolldown-runtime.mjs"; | ||
| import { i as AdapterHookable, r as getPeers, t as adapterUtils } from "../_chunks/adapter.mjs"; | ||
| import { n as import_websocket_server } from "../_chunks/libs/ws.mjs"; | ||
| import { n as Message, r as toBufferLike, t as Peer } from "../_chunks/peer.mjs"; | ||
| import { t as StubRequest } from "../_chunks/_request.mjs"; | ||
| import { t as WSError } from "../_chunks/error.mjs"; | ||
| //#region src/adapters/node.ts | ||
| const nodeAdapter = (options = {}) => { | ||
| if ("Deno" in globalThis || "Bun" in globalThis) { | ||
| throw new Error( | ||
| "[crossws] Using Node.js adapter in an incompatible environment." | ||
| ); | ||
| } | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| const wss = options.wss || new _WebSocketServer({ | ||
| noServer: true, | ||
| handleProtocols: () => false, | ||
| ...options.serverOptions | ||
| }); | ||
| wss.on("connection", (ws, nodeReq) => { | ||
| const request = new NodeReqProxy(nodeReq); | ||
| const peers = getPeers(globalPeers, nodeReq._namespace); | ||
| const peer = new NodePeer({ | ||
| ws, | ||
| request, | ||
| peers, | ||
| nodeReq, | ||
| namespace: nodeReq._namespace | ||
| }); | ||
| peers.add(peer); | ||
| hooks.callHook("open", peer); | ||
| ws.on("message", (data) => { | ||
| if (Array.isArray(data)) { | ||
| data = Buffer.concat(data); | ||
| } | ||
| hooks.callHook("message", peer, new Message(data, peer)); | ||
| }); | ||
| ws.on("error", (error) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("error", peer, new WSError(error)); | ||
| }); | ||
| ws.on("close", (code, reason) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, { | ||
| code, | ||
| reason: reason?.toString() | ||
| }); | ||
| }); | ||
| }); | ||
| wss.on("headers", (outgoingHeaders, req) => { | ||
| const upgradeHeaders = req._upgradeHeaders; | ||
| if (upgradeHeaders) { | ||
| for (const [key, value] of new Headers(upgradeHeaders)) { | ||
| outgoingHeaders.push(`${key}: ${value}`); | ||
| } | ||
| } | ||
| }); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| handleUpgrade: async (nodeReq, socket, head, webRequest) => { | ||
| const request = webRequest || new NodeReqProxy(nodeReq); | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) { | ||
| return sendResponse(socket, endResponse); | ||
| } | ||
| nodeReq._request = request; | ||
| nodeReq._upgradeHeaders = upgradeHeaders; | ||
| nodeReq._context = context; | ||
| nodeReq._namespace = namespace; | ||
| wss.handleUpgrade(nodeReq, socket, head, (ws) => { | ||
| wss.emit("connection", ws, nodeReq); | ||
| }); | ||
| }, | ||
| closeAll: (code, data, force) => { | ||
| for (const client of wss.clients) { | ||
| if (force) { | ||
| client.terminate(); | ||
| } else { | ||
| client.close(code, data); | ||
| } | ||
| } | ||
| } | ||
| }; | ||
| if ("Deno" in globalThis || "Bun" in globalThis) throw new Error("[crossws] Using Node.js adapter in an incompatible environment."); | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| const wss = options.wss || new import_websocket_server.default({ | ||
| noServer: true, | ||
| handleProtocols: () => false, | ||
| ...options.serverOptions | ||
| }); | ||
| wss.on("connection", (ws, nodeReq) => { | ||
| const request = new NodeReqProxy(nodeReq); | ||
| const peers = getPeers(globalPeers, nodeReq._namespace); | ||
| const peer = new NodePeer({ | ||
| ws, | ||
| request, | ||
| peers, | ||
| nodeReq, | ||
| namespace: nodeReq._namespace | ||
| }); | ||
| peers.add(peer); | ||
| hooks.callHook("open", peer); | ||
| ws.on("message", (data) => { | ||
| if (Array.isArray(data)) data = Buffer.concat(data); | ||
| hooks.callHook("message", peer, new Message(data, peer)); | ||
| }); | ||
| ws.on("error", (error) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("error", peer, new WSError(error)); | ||
| }); | ||
| ws.on("close", (code, reason) => { | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, { | ||
| code, | ||
| reason: reason?.toString() | ||
| }); | ||
| }); | ||
| }); | ||
| wss.on("headers", (outgoingHeaders, req) => { | ||
| const upgradeHeaders = req._upgradeHeaders; | ||
| if (upgradeHeaders) for (const [key, value] of new Headers(upgradeHeaders)) outgoingHeaders.push(`${key}: ${value}`); | ||
| }); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| handleUpgrade: async (nodeReq, socket, head, webRequest) => { | ||
| const request = webRequest || new NodeReqProxy(nodeReq); | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) return sendResponse(socket, endResponse); | ||
| nodeReq._request = request; | ||
| nodeReq._upgradeHeaders = upgradeHeaders; | ||
| nodeReq._context = context; | ||
| nodeReq._namespace = namespace; | ||
| wss.handleUpgrade(nodeReq, socket, head, (ws) => { | ||
| wss.emit("connection", ws, nodeReq); | ||
| }); | ||
| }, | ||
| closeAll: (code, data, force) => { | ||
| for (const client of wss.clients) if (force) client.terminate(); | ||
| else client.close(code, data); | ||
| } | ||
| }; | ||
| }; | ||
| class NodePeer extends Peer { | ||
| get remoteAddress() { | ||
| return this._internal.nodeReq.socket?.remoteAddress; | ||
| } | ||
| get context() { | ||
| return this._internal.nodeReq._context; | ||
| } | ||
| send(data, options) { | ||
| const dataBuff = toBufferLike(data); | ||
| const isBinary = typeof dataBuff !== "string"; | ||
| this._internal.ws.send(dataBuff, { | ||
| compress: options?.compress, | ||
| binary: isBinary, | ||
| ...options | ||
| }); | ||
| return 0; | ||
| } | ||
| publish(topic, data, options) { | ||
| const dataBuff = toBufferLike(data); | ||
| const isBinary = typeof data !== "string"; | ||
| const sendOptions = { | ||
| compress: options?.compress, | ||
| binary: isBinary, | ||
| ...options | ||
| }; | ||
| for (const peer of this._internal.peers) { | ||
| if (peer !== this && peer._topics.has(topic)) { | ||
| peer._internal.ws.send(dataBuff, sendOptions); | ||
| } | ||
| } | ||
| } | ||
| close(code, data) { | ||
| this._internal.ws.close(code, data); | ||
| } | ||
| terminate() { | ||
| this._internal.ws.terminate(); | ||
| } | ||
| } | ||
| class NodeReqProxy extends StubRequest { | ||
| constructor(req) { | ||
| const host = req.headers["host"] || "localhost"; | ||
| const isSecure = req.socket?.encrypted ?? req.headers["x-forwarded-proto"] === "https"; | ||
| const url = `${isSecure ? "https" : "http"}://${host}${req.url}`; | ||
| super(url, { headers: req.headers }); | ||
| } | ||
| } | ||
| var node_default = nodeAdapter; | ||
| var NodePeer = class extends Peer { | ||
| get remoteAddress() { | ||
| return this._internal.nodeReq.socket?.remoteAddress; | ||
| } | ||
| get context() { | ||
| return this._internal.nodeReq._context; | ||
| } | ||
| send(data, options) { | ||
| const dataBuff = toBufferLike(data); | ||
| const isBinary = typeof dataBuff !== "string"; | ||
| this._internal.ws.send(dataBuff, { | ||
| compress: options?.compress, | ||
| binary: isBinary, | ||
| ...options | ||
| }); | ||
| return 0; | ||
| } | ||
| publish(topic, data, options) { | ||
| const dataBuff = toBufferLike(data); | ||
| const isBinary = typeof data !== "string"; | ||
| const sendOptions = { | ||
| compress: options?.compress, | ||
| binary: isBinary, | ||
| ...options | ||
| }; | ||
| for (const peer of this._internal.peers) if (peer !== this && peer._topics.has(topic)) peer._internal.ws.send(dataBuff, sendOptions); | ||
| } | ||
| close(code, data) { | ||
| this._internal.ws.close(code, data); | ||
| } | ||
| terminate() { | ||
| this._internal.ws.terminate(); | ||
| } | ||
| }; | ||
| var NodeReqProxy = class extends StubRequest { | ||
| constructor(req) { | ||
| const host = req.headers["host"] || "localhost"; | ||
| const url = `${req.socket?.encrypted ?? req.headers["x-forwarded-proto"] === "https" ? "https" : "http"}://${host}${req.url}`; | ||
| super(url, { headers: req.headers }); | ||
| } | ||
| }; | ||
| async function sendResponse(socket, res) { | ||
| const head = [ | ||
| `HTTP/1.1 ${res.status || 200} ${res.statusText || ""}`, | ||
| ...[...res.headers.entries()].map( | ||
| ([key, value]) => `${encodeURIComponent(key)}: ${encodeURIComponent(value)}` | ||
| ) | ||
| ]; | ||
| socket.write(head.join("\r\n") + "\r\n\r\n"); | ||
| if (res.body) { | ||
| for await (const chunk of res.body) { | ||
| socket.write(chunk); | ||
| } | ||
| } | ||
| return new Promise((resolve) => { | ||
| socket.end(() => { | ||
| socket.destroy(); | ||
| resolve(); | ||
| }); | ||
| }); | ||
| const head = [`HTTP/1.1 ${res.status || 200} ${res.statusText || ""}`, ...[...res.headers.entries()].map(([key, value]) => `${encodeURIComponent(key)}: ${encodeURIComponent(value)}`)]; | ||
| socket.write(head.join("\r\n") + "\r\n\r\n"); | ||
| if (res.body) for await (const chunk of res.body) socket.write(chunk); | ||
| return new Promise((resolve) => { | ||
| socket.end(() => { | ||
| socket.destroy(); | ||
| resolve(); | ||
| }); | ||
| }); | ||
| } | ||
| export { nodeAdapter as default }; | ||
| //#endregion | ||
| export { node_default as default }; |
@@ -1,13 +0,2 @@ | ||
| import { Adapter, AdapterInstance, AdapterOptions } from '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| interface SSEAdapter extends AdapterInstance { | ||
| fetch(req: Request): Promise<Response>; | ||
| } | ||
| interface SSEOptions extends AdapterOptions { | ||
| bidir?: boolean; | ||
| } | ||
| declare const sseAdapter: Adapter<SSEAdapter, SSEOptions>; | ||
| export { sseAdapter as default }; | ||
| export type { SSEAdapter, SSEOptions }; | ||
| import { n as SSEOptions, r as sseAdapter, t as SSEAdapter } from "../_chunks/sse.mjs"; | ||
| export { SSEAdapter, SSEOptions, sseAdapter as default }; |
+98
-118
@@ -1,122 +0,102 @@ | ||
| import { M as Message, P as Peer, a as toString } from '../shared/crossws.WpyOHUXc.mjs'; | ||
| import { a as adapterUtils, A as AdapterHookable, g as getPeers } from '../shared/crossws.CPlNx7g8.mjs'; | ||
| import { i as AdapterHookable, r as getPeers, t as adapterUtils } from "../_chunks/adapter.mjs"; | ||
| import { i as toString, n as Message, t as Peer } from "../_chunks/peer.mjs"; | ||
| //#region src/adapters/sse.ts | ||
| const sseAdapter = (opts = {}) => { | ||
| const hooks = new AdapterHookable(opts); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| const peersMap = opts.bidir ? /* @__PURE__ */ new Map() : void 0; | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| fetch: async (request) => { | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) { | ||
| return endResponse; | ||
| } | ||
| let peer; | ||
| if (opts.bidir && request.body && request.headers.has("x-crossws-id")) { | ||
| const id = request.headers.get("x-crossws-id"); | ||
| peer = peersMap?.get(id); | ||
| if (!peer) { | ||
| return new Response("invalid peer id", { status: 400 }); | ||
| } | ||
| const stream = request.body.pipeThrough(new TextDecoderStream()); | ||
| try { | ||
| for await (const chunk of stream) { | ||
| hooks.callHook("message", peer, new Message(chunk, peer)); | ||
| } | ||
| } catch { | ||
| await stream.cancel().catch(() => { | ||
| }); | ||
| } | ||
| return new Response(null, {}); | ||
| } else { | ||
| const ws = new SSEWebSocketStub(); | ||
| const peers = getPeers(globalPeers, namespace); | ||
| peer = new SSEPeer({ | ||
| peers, | ||
| peersMap, | ||
| request, | ||
| hooks, | ||
| ws, | ||
| context, | ||
| namespace | ||
| }); | ||
| peers.add(peer); | ||
| if (opts.bidir) { | ||
| peersMap.set(peer.id, peer); | ||
| peer._sendEvent("crossws-id", peer.id); | ||
| } | ||
| } | ||
| let headers = { | ||
| "Content-Type": "text/event-stream", | ||
| "Cache-Control": "no-cache", | ||
| Connection: "keep-alive" | ||
| }; | ||
| if (opts.bidir) { | ||
| headers["x-crossws-id"] = peer.id; | ||
| } | ||
| if (upgradeHeaders) { | ||
| headers = new Headers(headers); | ||
| for (const [key, value] of new Headers(upgradeHeaders)) { | ||
| headers.set(key, value); | ||
| } | ||
| } | ||
| return new Response(peer._sseStream, { headers }); | ||
| } | ||
| }; | ||
| const hooks = new AdapterHookable(opts); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| const peersMap = opts.bidir ? /* @__PURE__ */ new Map() : void 0; | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| fetch: async (request) => { | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(request); | ||
| if (endResponse) return endResponse; | ||
| let peer; | ||
| if (opts.bidir && request.body && request.headers.has("x-crossws-id")) { | ||
| const id = request.headers.get("x-crossws-id"); | ||
| peer = peersMap?.get(id); | ||
| if (!peer) return new Response("invalid peer id", { status: 400 }); | ||
| const stream = request.body.pipeThrough(new TextDecoderStream()); | ||
| try { | ||
| for await (const chunk of stream) hooks.callHook("message", peer, new Message(chunk, peer)); | ||
| } catch { | ||
| await stream.cancel().catch(() => {}); | ||
| } | ||
| return new Response(null, {}); | ||
| } else { | ||
| const ws = new SSEWebSocketStub(); | ||
| const peers = getPeers(globalPeers, namespace); | ||
| peer = new SSEPeer({ | ||
| peers, | ||
| peersMap, | ||
| request, | ||
| hooks, | ||
| ws, | ||
| context, | ||
| namespace | ||
| }); | ||
| peers.add(peer); | ||
| if (opts.bidir) { | ||
| peersMap.set(peer.id, peer); | ||
| peer._sendEvent("crossws-id", peer.id); | ||
| } | ||
| } | ||
| let headers = { | ||
| "Content-Type": "text/event-stream", | ||
| "Cache-Control": "no-cache", | ||
| Connection: "keep-alive" | ||
| }; | ||
| if (opts.bidir) headers["x-crossws-id"] = peer.id; | ||
| if (upgradeHeaders) { | ||
| headers = new Headers(headers); | ||
| for (const [key, value] of new Headers(upgradeHeaders)) headers.set(key, value); | ||
| } | ||
| return new Response(peer._sseStream, { headers }); | ||
| } | ||
| }; | ||
| }; | ||
| class SSEPeer extends Peer { | ||
| _sseStream; | ||
| // server -> client | ||
| _sseStreamController; | ||
| constructor(_internal) { | ||
| super(_internal); | ||
| _internal.ws.readyState = 0; | ||
| this._sseStream = new ReadableStream({ | ||
| start: (controller) => { | ||
| _internal.ws.readyState = 1; | ||
| this._sseStreamController = controller; | ||
| _internal.hooks.callHook("open", this); | ||
| }, | ||
| cancel: () => { | ||
| _internal.ws.readyState = 2; | ||
| _internal.peers.delete(this); | ||
| _internal.peersMap?.delete(this.id); | ||
| Promise.resolve(this._internal.hooks.callHook("close", this)).finally( | ||
| () => { | ||
| _internal.ws.readyState = 3; | ||
| } | ||
| ); | ||
| } | ||
| }).pipeThrough(new TextEncoderStream()); | ||
| } | ||
| _sendEvent(event, data) { | ||
| const lines = data.split("\n"); | ||
| this._sseStreamController?.enqueue( | ||
| `event: ${event} | ||
| ${lines.map((l) => `data: ${l}`)} | ||
| var sse_default = sseAdapter; | ||
| var SSEPeer = class extends Peer { | ||
| _sseStream; | ||
| _sseStreamController; | ||
| constructor(_internal) { | ||
| super(_internal); | ||
| _internal.ws.readyState = 0; | ||
| this._sseStream = new ReadableStream({ | ||
| start: (controller) => { | ||
| _internal.ws.readyState = 1; | ||
| this._sseStreamController = controller; | ||
| _internal.hooks.callHook("open", this); | ||
| }, | ||
| cancel: () => { | ||
| _internal.ws.readyState = 2; | ||
| _internal.peers.delete(this); | ||
| _internal.peersMap?.delete(this.id); | ||
| Promise.resolve(this._internal.hooks.callHook("close", this)).finally(() => { | ||
| _internal.ws.readyState = 3; | ||
| }); | ||
| } | ||
| }).pipeThrough(new TextEncoderStream()); | ||
| } | ||
| _sendEvent(event, data) { | ||
| const lines = data.split("\n"); | ||
| this._sseStreamController?.enqueue(`event: ${event}\n${lines.map((l) => `data: ${l}`)}\n\n`); | ||
| } | ||
| send(data) { | ||
| this._sendEvent("message", toString(data)); | ||
| return 0; | ||
| } | ||
| publish(topic, data) { | ||
| const dataBuff = toString(data); | ||
| for (const peer of this._internal.peers) if (peer !== this && peer._topics.has(topic)) peer._sendEvent("message", dataBuff); | ||
| } | ||
| close() { | ||
| this._sseStreamController?.close(); | ||
| } | ||
| }; | ||
| var SSEWebSocketStub = class { | ||
| readyState; | ||
| }; | ||
| ` | ||
| ); | ||
| } | ||
| send(data) { | ||
| this._sendEvent("message", toString(data)); | ||
| return 0; | ||
| } | ||
| publish(topic, data) { | ||
| const dataBuff = toString(data); | ||
| for (const peer of this._internal.peers) { | ||
| if (peer !== this && peer._topics.has(topic)) { | ||
| peer._sendEvent("message", dataBuff); | ||
| } | ||
| } | ||
| } | ||
| close() { | ||
| this._sseStreamController?.close(); | ||
| } | ||
| } | ||
| class SSEWebSocketStub { | ||
| readyState; | ||
| } | ||
| export { sseAdapter as default }; | ||
| //#endregion | ||
| export { sse_default as default }; |
+44
-44
@@ -1,62 +0,62 @@ | ||
| import { Adapter, AdapterInstance, Peer, PeerContext, AdapterOptions } from '../index.mjs'; | ||
| import { W as WebSocket } from '../shared/crossws.BQXMA5bH.mjs'; | ||
| import uws from 'uWebSockets.js'; | ||
| import { d as PeerContext, n as AdapterInstance, r as AdapterOptions, t as Adapter, u as Peer } from "../_chunks/adapter.mjs"; | ||
| import { a as WebSocket } from "../_chunks/web.mjs"; | ||
| import uws from "uWebSockets.js"; | ||
| //#region src/_request.d.ts | ||
| declare const StubRequest: { | ||
| new (url: string, init?: RequestInit): Request; | ||
| new (url: string, init?: RequestInit): Request; | ||
| }; | ||
| //#endregion | ||
| //#region src/adapters/uws.d.ts | ||
| type UserData = { | ||
| peer?: UWSPeer; | ||
| req: uws.HttpRequest; | ||
| res: uws.HttpResponse; | ||
| webReq: UWSReqProxy; | ||
| protocol: string; | ||
| extensions: string; | ||
| context: PeerContext; | ||
| namespace: string; | ||
| peer?: UWSPeer; | ||
| req: uws.HttpRequest; | ||
| res: uws.HttpResponse; | ||
| webReq: UWSReqProxy; | ||
| protocol: string; | ||
| extensions: string; | ||
| context: PeerContext; | ||
| namespace: string; | ||
| }; | ||
| type WebSocketHandler = uws.WebSocketBehavior<UserData>; | ||
| interface UWSAdapter extends AdapterInstance { | ||
| websocket: WebSocketHandler; | ||
| websocket: WebSocketHandler; | ||
| } | ||
| interface UWSOptions extends AdapterOptions { | ||
| uws?: Exclude<uws.WebSocketBehavior<any>, "close" | "drain" | "message" | "open" | "ping" | "pong" | "subscription" | "upgrade">; | ||
| uws?: Exclude<uws.WebSocketBehavior<any>, "close" | "drain" | "message" | "open" | "ping" | "pong" | "subscription" | "upgrade">; | ||
| } | ||
| declare const uwsAdapter: Adapter<UWSAdapter, UWSOptions>; | ||
| declare class UWSPeer extends Peer<{ | ||
| peers: Set<UWSPeer>; | ||
| request: UWSReqProxy; | ||
| namespace: string; | ||
| uws: uws.WebSocket<UserData>; | ||
| ws: UwsWebSocketProxy; | ||
| uwsData: UserData; | ||
| peers: Set<UWSPeer>; | ||
| request: UWSReqProxy; | ||
| namespace: string; | ||
| uws: uws.WebSocket<UserData>; | ||
| ws: UwsWebSocketProxy; | ||
| uwsData: UserData; | ||
| }> { | ||
| get remoteAddress(): string | undefined; | ||
| get context(): PeerContext; | ||
| send(data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| subscribe(topic: string): void; | ||
| unsubscribe(topic: string): void; | ||
| publish(topic: string, message: string, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| close(code?: number, reason?: uws.RecognizedString): void; | ||
| terminate(): void; | ||
| get remoteAddress(): string | undefined; | ||
| get context(): PeerContext; | ||
| send(data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| subscribe(topic: string): void; | ||
| unsubscribe(topic: string): void; | ||
| publish(topic: string, message: string, options?: { | ||
| compress?: boolean; | ||
| }): number; | ||
| close(code?: number, reason?: uws.RecognizedString): void; | ||
| terminate(): void; | ||
| } | ||
| declare class UWSReqProxy extends StubRequest { | ||
| constructor(req: uws.HttpRequest); | ||
| constructor(req: uws.HttpRequest); | ||
| } | ||
| declare class UwsWebSocketProxy implements Partial<WebSocket> { | ||
| private _uws; | ||
| readyState?: number; | ||
| constructor(_uws: uws.WebSocket<UserData>); | ||
| get bufferedAmount(): number; | ||
| get protocol(): string; | ||
| get extensions(): string; | ||
| private _uws; | ||
| readyState?: number; | ||
| constructor(_uws: uws.WebSocket<UserData>); | ||
| get bufferedAmount(): number; | ||
| get protocol(): string; | ||
| get extensions(): string; | ||
| } | ||
| export { uwsAdapter as default }; | ||
| export type { UWSAdapter, UWSOptions }; | ||
| //#endregion | ||
| export { UWSAdapter, UWSOptions, uwsAdapter as default }; |
+152
-175
@@ -1,181 +0,158 @@ | ||
| import { M as Message, P as Peer, t as toBufferLike } from '../shared/crossws.WpyOHUXc.mjs'; | ||
| import { a as adapterUtils, A as AdapterHookable, g as getPeers } from '../shared/crossws.CPlNx7g8.mjs'; | ||
| import { S as StubRequest } from '../shared/crossws.B31KJMcF.mjs'; | ||
| import { i as AdapterHookable, r as getPeers, t as adapterUtils } from "../_chunks/adapter.mjs"; | ||
| import { n as Message, r as toBufferLike, t as Peer } from "../_chunks/peer.mjs"; | ||
| import { t as StubRequest } from "../_chunks/_request.mjs"; | ||
| //#region src/adapters/uws.ts | ||
| const uwsAdapter = (options = {}) => { | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| websocket: { | ||
| ...options.uws, | ||
| close(ws, code, message) { | ||
| const peers = getPeers(globalPeers, ws.getUserData().namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peer._internal.ws.readyState = 2; | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, { | ||
| code, | ||
| reason: message?.toString() | ||
| }); | ||
| peer._internal.ws.readyState = 3; | ||
| }, | ||
| message(ws, message, isBinary) { | ||
| const peers = getPeers(globalPeers, ws.getUserData().namespace); | ||
| const peer = getPeer(ws, peers); | ||
| hooks.callHook("message", peer, new Message(message, peer)); | ||
| }, | ||
| open(ws) { | ||
| const peers = getPeers(globalPeers, ws.getUserData().namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peers.add(peer); | ||
| hooks.callHook("open", peer); | ||
| }, | ||
| async upgrade(res, req, uwsContext) { | ||
| let aborted = false; | ||
| res.onAborted(() => { | ||
| aborted = true; | ||
| }); | ||
| const webReq = new UWSReqProxy(req); | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(webReq); | ||
| if (endResponse) { | ||
| res.writeStatus(`${endResponse.status} ${endResponse.statusText}`); | ||
| for (const [key, value] of endResponse.headers) { | ||
| res.writeHeader(key, value); | ||
| } | ||
| if (endResponse.body) { | ||
| for await (const chunk of endResponse.body) { | ||
| if (aborted) break; | ||
| res.write(chunk); | ||
| } | ||
| } | ||
| if (!aborted) { | ||
| res.end(); | ||
| } | ||
| return; | ||
| } | ||
| if (aborted) { | ||
| return; | ||
| } | ||
| res.writeStatus("101 Switching Protocols"); | ||
| if (upgradeHeaders) { | ||
| const headers = upgradeHeaders instanceof Headers ? upgradeHeaders : new Headers(upgradeHeaders); | ||
| for (const [key, value] of headers) { | ||
| res.writeHeader(key, value); | ||
| } | ||
| } | ||
| res.cork(() => { | ||
| const key = req.getHeader("sec-websocket-key"); | ||
| const protocol = req.getHeader("sec-websocket-protocol"); | ||
| const extensions = req.getHeader("sec-websocket-extensions"); | ||
| res.upgrade( | ||
| { | ||
| req, | ||
| res, | ||
| webReq, | ||
| protocol, | ||
| extensions, | ||
| context, | ||
| namespace | ||
| }, | ||
| key, | ||
| "", | ||
| extensions, | ||
| uwsContext | ||
| ); | ||
| }); | ||
| } | ||
| } | ||
| }; | ||
| const hooks = new AdapterHookable(options); | ||
| const globalPeers = /* @__PURE__ */ new Map(); | ||
| return { | ||
| ...adapterUtils(globalPeers), | ||
| websocket: { | ||
| ...options.uws, | ||
| close(ws, code, message) { | ||
| const peers = getPeers(globalPeers, ws.getUserData().namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peer._internal.ws.readyState = 2; | ||
| peers.delete(peer); | ||
| hooks.callHook("close", peer, { | ||
| code, | ||
| reason: message?.toString() | ||
| }); | ||
| peer._internal.ws.readyState = 3; | ||
| }, | ||
| message(ws, message, isBinary) { | ||
| const peer = getPeer(ws, getPeers(globalPeers, ws.getUserData().namespace)); | ||
| hooks.callHook("message", peer, new Message(message, peer)); | ||
| }, | ||
| open(ws) { | ||
| const peers = getPeers(globalPeers, ws.getUserData().namespace); | ||
| const peer = getPeer(ws, peers); | ||
| peers.add(peer); | ||
| hooks.callHook("open", peer); | ||
| }, | ||
| async upgrade(res, req, uwsContext) { | ||
| let aborted = false; | ||
| res.onAborted(() => { | ||
| aborted = true; | ||
| }); | ||
| const webReq = new UWSReqProxy(req); | ||
| const { upgradeHeaders, endResponse, context, namespace } = await hooks.upgrade(webReq); | ||
| if (endResponse) { | ||
| res.writeStatus(`${endResponse.status} ${endResponse.statusText}`); | ||
| for (const [key, value] of endResponse.headers) res.writeHeader(key, value); | ||
| if (endResponse.body) for await (const chunk of endResponse.body) { | ||
| if (aborted) break; | ||
| res.write(chunk); | ||
| } | ||
| if (!aborted) res.end(); | ||
| return; | ||
| } | ||
| if (aborted) return; | ||
| res.writeStatus("101 Switching Protocols"); | ||
| if (upgradeHeaders) { | ||
| const headers = upgradeHeaders instanceof Headers ? upgradeHeaders : new Headers(upgradeHeaders); | ||
| for (const [key, value] of headers) res.writeHeader(key, value); | ||
| } | ||
| res.cork(() => { | ||
| const key = req.getHeader("sec-websocket-key"); | ||
| const protocol = req.getHeader("sec-websocket-protocol"); | ||
| const extensions = req.getHeader("sec-websocket-extensions"); | ||
| res.upgrade({ | ||
| req, | ||
| res, | ||
| webReq, | ||
| protocol, | ||
| extensions, | ||
| context, | ||
| namespace | ||
| }, key, "", extensions, uwsContext); | ||
| }); | ||
| } | ||
| } | ||
| }; | ||
| }; | ||
| var uws_default = uwsAdapter; | ||
| function getPeer(uws, peers) { | ||
| const uwsData = uws.getUserData(); | ||
| if (uwsData.peer) { | ||
| return uwsData.peer; | ||
| } | ||
| const peer = new UWSPeer({ | ||
| peers, | ||
| uws, | ||
| ws: new UwsWebSocketProxy(uws), | ||
| request: uwsData.webReq, | ||
| namespace: uwsData.namespace, | ||
| uwsData | ||
| }); | ||
| uwsData.peer = peer; | ||
| return peer; | ||
| const uwsData = uws.getUserData(); | ||
| if (uwsData.peer) return uwsData.peer; | ||
| const peer = new UWSPeer({ | ||
| peers, | ||
| uws, | ||
| ws: new UwsWebSocketProxy(uws), | ||
| request: uwsData.webReq, | ||
| namespace: uwsData.namespace, | ||
| uwsData | ||
| }); | ||
| uwsData.peer = peer; | ||
| return peer; | ||
| } | ||
| class UWSPeer extends Peer { | ||
| get remoteAddress() { | ||
| try { | ||
| const addr = new TextDecoder().decode( | ||
| this._internal.uws.getRemoteAddressAsText() | ||
| ); | ||
| return addr; | ||
| } catch { | ||
| } | ||
| } | ||
| get context() { | ||
| return this._internal.uwsData.context; | ||
| } | ||
| send(data, options) { | ||
| const dataBuff = toBufferLike(data); | ||
| const isBinary = typeof dataBuff !== "string"; | ||
| return this._internal.uws.send(dataBuff, isBinary, options?.compress); | ||
| } | ||
| subscribe(topic) { | ||
| this._topics.add(topic); | ||
| this._internal.uws.subscribe(topic); | ||
| } | ||
| unsubscribe(topic) { | ||
| this._topics.delete(topic); | ||
| this._internal.uws.unsubscribe(topic); | ||
| } | ||
| publish(topic, message, options) { | ||
| const data = toBufferLike(message); | ||
| const isBinary = typeof data !== "string"; | ||
| this._internal.uws.publish(topic, data, isBinary, options?.compress); | ||
| return 0; | ||
| } | ||
| close(code, reason) { | ||
| this._internal.uws.end(code, reason); | ||
| } | ||
| terminate() { | ||
| this._internal.uws.close(); | ||
| } | ||
| } | ||
| class UWSReqProxy extends StubRequest { | ||
| constructor(req) { | ||
| const rawHeaders = []; | ||
| let host = "localhost"; | ||
| let proto = "http"; | ||
| req.forEach((key, value) => { | ||
| if (key === "host") { | ||
| host = value; | ||
| } else if (key === "x-forwarded-proto" && value === "https") { | ||
| proto = "https"; | ||
| } | ||
| rawHeaders.push([key, value]); | ||
| }); | ||
| const query = req.getQuery(); | ||
| const pathname = req.getUrl(); | ||
| const url = `${proto}://${host}${pathname}${query ? `?${query}` : ""}`; | ||
| super(url, { headers: rawHeaders }); | ||
| } | ||
| } | ||
| class UwsWebSocketProxy { | ||
| constructor(_uws) { | ||
| this._uws = _uws; | ||
| } | ||
| readyState = 1; | ||
| get bufferedAmount() { | ||
| return this._uws?.getBufferedAmount(); | ||
| } | ||
| get protocol() { | ||
| return this._uws?.getUserData().protocol; | ||
| } | ||
| get extensions() { | ||
| return this._uws?.getUserData().extensions; | ||
| } | ||
| } | ||
| var UWSPeer = class extends Peer { | ||
| get remoteAddress() { | ||
| try { | ||
| return new TextDecoder().decode(this._internal.uws.getRemoteAddressAsText()); | ||
| } catch {} | ||
| } | ||
| get context() { | ||
| return this._internal.uwsData.context; | ||
| } | ||
| send(data, options) { | ||
| const dataBuff = toBufferLike(data); | ||
| const isBinary = typeof dataBuff !== "string"; | ||
| return this._internal.uws.send(dataBuff, isBinary, options?.compress); | ||
| } | ||
| subscribe(topic) { | ||
| this._topics.add(topic); | ||
| this._internal.uws.subscribe(topic); | ||
| } | ||
| unsubscribe(topic) { | ||
| this._topics.delete(topic); | ||
| this._internal.uws.unsubscribe(topic); | ||
| } | ||
| publish(topic, message, options) { | ||
| const data = toBufferLike(message); | ||
| const isBinary = typeof data !== "string"; | ||
| this._internal.uws.publish(topic, data, isBinary, options?.compress); | ||
| return 0; | ||
| } | ||
| close(code, reason) { | ||
| this._internal.uws.end(code, reason); | ||
| } | ||
| terminate() { | ||
| this._internal.uws.close(); | ||
| } | ||
| }; | ||
| var UWSReqProxy = class extends StubRequest { | ||
| constructor(req) { | ||
| const rawHeaders = []; | ||
| let host = "localhost"; | ||
| let proto = "http"; | ||
| req.forEach((key, value) => { | ||
| if (key === "host") host = value; | ||
| else if (key === "x-forwarded-proto" && value === "https") proto = "https"; | ||
| rawHeaders.push([key, value]); | ||
| }); | ||
| const query = req.getQuery(); | ||
| const pathname = req.getUrl(); | ||
| const url = `${proto}://${host}${pathname}${query ? `?${query}` : ""}`; | ||
| super(url, { headers: rawHeaders }); | ||
| } | ||
| }; | ||
| var UwsWebSocketProxy = class { | ||
| readyState = 1; | ||
| constructor(_uws) { | ||
| this._uws = _uws; | ||
| } | ||
| get bufferedAmount() { | ||
| return this._uws?.getBufferedAmount(); | ||
| } | ||
| get protocol() { | ||
| return this._uws?.getUserData().protocol; | ||
| } | ||
| get extensions() { | ||
| return this._uws?.getUserData().extensions; | ||
| } | ||
| }; | ||
| export { uwsAdapter as default }; | ||
| //#endregion | ||
| export { uws_default as default }; |
+2
-170
@@ -1,170 +0,2 @@ | ||
| import { W as WebSocket } from './shared/crossws.BQXMA5bH.mjs'; | ||
| declare const kNodeInspect: unique symbol; | ||
| interface PeerContext extends Record<string, unknown> { | ||
| } | ||
| interface AdapterInternal { | ||
| ws: unknown; | ||
| request: Request; | ||
| namespace: string; | ||
| peers?: Set<Peer>; | ||
| context?: PeerContext; | ||
| } | ||
| declare abstract class Peer<Internal extends AdapterInternal = AdapterInternal> { | ||
| #private; | ||
| protected _internal: Internal; | ||
| protected _topics: Set<string>; | ||
| protected _id?: string; | ||
| constructor(internal: Internal); | ||
| get context(): PeerContext; | ||
| get namespace(): string; | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the peer. | ||
| */ | ||
| get id(): string; | ||
| /** IP address of the peer */ | ||
| get remoteAddress(): string | undefined; | ||
| /** upgrade request */ | ||
| get request(): Request; | ||
| /** | ||
| * Get the [WebSocket](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket) instance. | ||
| * | ||
| * **Note:** crossws adds polyfill for the following properties if native values are not available: | ||
| * - `protocol`: Extracted from the `sec-websocket-protocol` header. | ||
| * - `extensions`: Extracted from the `sec-websocket-extensions` header. | ||
| * - `url`: Extracted from the request URL (http -> ws). | ||
| * */ | ||
| get websocket(): Partial<WebSocket>; | ||
| /** All connected peers to the server */ | ||
| get peers(): Set<Peer>; | ||
| /** All topics, this peer has been subscribed to. */ | ||
| get topics(): Set<string>; | ||
| abstract close(code?: number, reason?: string): void; | ||
| /** Abruptly close the connection */ | ||
| terminate(): void; | ||
| /** Subscribe to a topic */ | ||
| subscribe(topic: string): void; | ||
| /** Unsubscribe from a topic */ | ||
| unsubscribe(topic: string): void; | ||
| /** Send a message to the peer. */ | ||
| abstract send(data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): number | void | undefined; | ||
| /** Send message to subscribes of topic */ | ||
| abstract publish(topic: string, data: unknown, options?: { | ||
| compress?: boolean; | ||
| }): void; | ||
| toString(): string; | ||
| [Symbol.toPrimitive](): string; | ||
| [Symbol.toStringTag](): "WebSocket"; | ||
| [kNodeInspect](): unknown; | ||
| } | ||
| interface AdapterInstance { | ||
| readonly peers: Map<string, Set<Peer>>; | ||
| readonly publish: (topic: string, data: unknown, options?: { | ||
| compress?: boolean; | ||
| namespace?: string; | ||
| }) => void; | ||
| } | ||
| interface AdapterOptions { | ||
| resolve?: ResolveHooks; | ||
| getNamespace?: (request: Request) => string; | ||
| hooks?: Partial<Hooks>; | ||
| } | ||
| type Adapter<AdapterT extends AdapterInstance = AdapterInstance, Options extends AdapterOptions = AdapterOptions> = (options?: Options) => AdapterT; | ||
| declare function defineWebSocketAdapter<AdapterT extends AdapterInstance = AdapterInstance, Options extends AdapterOptions = AdapterOptions>(factory: Adapter<AdapterT, Options>): Adapter<AdapterT, Options>; | ||
| declare class WSError extends Error { | ||
| constructor(...args: any[]); | ||
| } | ||
| declare class Message implements Partial<MessageEvent> { | ||
| #private; | ||
| /** Access to the original [message event](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket/message_event) if available. */ | ||
| readonly event?: MessageEvent; | ||
| /** Access to the Peer that emitted the message. */ | ||
| readonly peer?: Peer; | ||
| /** Raw message data (can be of any type). */ | ||
| readonly rawData: unknown; | ||
| constructor(rawData: unknown, peer: Peer, event?: MessageEvent); | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the message. | ||
| */ | ||
| get id(): string; | ||
| /** | ||
| * Get data as [Uint8Array](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Uint8Array) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| uint8Array(): Uint8Array; | ||
| /** | ||
| * Get data as [ArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/ArrayBuffer) or [SharedArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/SharedArrayBuffer) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| arrayBuffer(): ArrayBuffer | SharedArrayBuffer; | ||
| /** | ||
| * Get data as [Blob](https://developer.mozilla.org/en-US/docs/Web/API/Blob) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. */ | ||
| blob(): Blob; | ||
| /** | ||
| * Get stringified text version of the message. | ||
| * | ||
| * If raw data is in any other format, it will be automatically converted and decoded. | ||
| */ | ||
| text(): string; | ||
| /** | ||
| * Get parsed version of the message text with [`JSON.parse()`](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/JSON/parse). | ||
| */ | ||
| json<T = unknown>(): T; | ||
| /** | ||
| * Message data (value varies based on `peer.websocket.binaryType`). | ||
| */ | ||
| get data(): unknown; | ||
| toString(): string; | ||
| [Symbol.toPrimitive](): string; | ||
| [kNodeInspect](): unknown; | ||
| } | ||
| declare function defineHooks<T extends Partial<Hooks> = Partial<Hooks>>(hooks: T): T; | ||
| type ResolveHooks = (request: Request & { | ||
| readonly context?: PeerContext; | ||
| }) => Partial<Hooks> | Promise<Partial<Hooks>>; | ||
| type MaybePromise<T> = T | Promise<T>; | ||
| interface Hooks { | ||
| /** | ||
| * Upgrading a request to a WebSocket connection. | ||
| * | ||
| * - You can throw a Response to abort the upgrade. | ||
| * - You can return { headers } to modify the response. | ||
| * - You can return { namespace } to change the pub/sub namespace. | ||
| * - You can return { context } to provide a custom peer context. | ||
| * | ||
| * @param request | ||
| * @throws {Response} | ||
| */ | ||
| upgrade: (request: Request & { | ||
| readonly context?: Record<string, unknown>; | ||
| }) => MaybePromise<{ | ||
| headers?: HeadersInit; | ||
| namespace?: string; | ||
| context?: PeerContext; | ||
| } | Response | void>; | ||
| /** A message is received */ | ||
| message: (peer: Peer, message: Message) => MaybePromise<void>; | ||
| /** A socket is opened */ | ||
| open: (peer: Peer) => MaybePromise<void>; | ||
| /** A socket is closed */ | ||
| close: (peer: Peer, details: { | ||
| code?: number; | ||
| reason?: string; | ||
| }) => MaybePromise<void>; | ||
| /** An error occurs */ | ||
| error: (peer: Peer, error: WSError) => MaybePromise<void>; | ||
| } | ||
| export { Message, Peer, WSError, defineHooks, defineWebSocketAdapter }; | ||
| export type { Adapter, AdapterInstance, AdapterInternal, AdapterOptions, Hooks, PeerContext, ResolveHooks }; | ||
| import { a as Hooks, c as Message, d as PeerContext, f as WSError, i as defineWebSocketAdapter, l as AdapterInternal, n as AdapterInstance, o as ResolveHooks, r as AdapterOptions, s as defineHooks, t as Adapter, u as Peer } from "./_chunks/adapter.mjs"; | ||
| export { type Adapter, type AdapterInstance, type AdapterInternal, type AdapterOptions, type Hooks, type Message, type Peer, type PeerContext, type ResolveHooks, type WSError, defineHooks, defineWebSocketAdapter }; |
+3
-1
@@ -1,1 +0,3 @@ | ||
| export { d as defineHooks, b as defineWebSocketAdapter } from './shared/crossws.CPlNx7g8.mjs'; | ||
| import { a as defineHooks, n as defineWebSocketAdapter } from "./_chunks/adapter.mjs"; | ||
| export { defineHooks, defineWebSocketAdapter }; |
@@ -1,24 +0,11 @@ | ||
| import { ServerPlugin, Server } from 'srvx'; | ||
| import { W as WSOptions, S as ServerWithWSOptions } from '../shared/crossws.CP-89VBK.mjs'; | ||
| import '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| import '../adapters/bun.mjs'; | ||
| import 'bun'; | ||
| import '../adapters/deno.mjs'; | ||
| import '../adapters/node.mjs'; | ||
| import 'node:http'; | ||
| import 'node:stream'; | ||
| import 'events'; | ||
| import 'node:https'; | ||
| import 'node:tls'; | ||
| import 'node:url'; | ||
| import 'node:zlib'; | ||
| import '../adapters/sse.mjs'; | ||
| import '../adapters/cloudflare.mjs'; | ||
| import '@cloudflare/workers-types'; | ||
| import 'cloudflare:workers'; | ||
| import "../_chunks/bun.mjs"; | ||
| import "../_chunks/cloudflare.mjs"; | ||
| import "../_chunks/node.mjs"; | ||
| import { n as WSOptions, t as ServerWithWSOptions } from "../_chunks/_types.mjs"; | ||
| import { Server, ServerPlugin } from "srvx"; | ||
| //#region src/server/bun.d.ts | ||
| declare function plugin(wsOpts: WSOptions): ServerPlugin; | ||
| declare function serve(options: ServerWithWSOptions): Server; | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
+24
-31
@@ -1,37 +0,30 @@ | ||
| import { serve as serve$1 } from 'srvx/bun'; | ||
| import bunAdapter from '../adapters/bun.mjs'; | ||
| import '../shared/crossws.WpyOHUXc.mjs'; | ||
| import '../shared/crossws.CPlNx7g8.mjs'; | ||
| import bun_default from "../adapters/bun.mjs"; | ||
| import { serve as serve$1 } from "srvx/bun"; | ||
| //#region src/server/bun.ts | ||
| function plugin(wsOpts) { | ||
| return (server) => { | ||
| const ws = bunAdapter({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.bun | ||
| }); | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") { | ||
| return ws.handleUpgrade( | ||
| req, | ||
| req.runtime.bun.server | ||
| ); | ||
| } | ||
| return next(); | ||
| }); | ||
| server.options.bun ??= {}; | ||
| if (server.options.bun.websocket) { | ||
| throw new Error("websocket handlers for bun already set!"); | ||
| } | ||
| server.options.bun.websocket = ws.websocket; | ||
| }; | ||
| return (server) => { | ||
| const ws = bun_default({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.bun | ||
| }); | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") return ws.handleUpgrade(req, req.runtime.bun.server); | ||
| return next(); | ||
| }); | ||
| server.options.bun ??= {}; | ||
| if (server.options.bun.websocket) throw new Error("websocket handlers for bun already set!"); | ||
| server.options.bun.websocket = ws.websocket; | ||
| }; | ||
| } | ||
| function serve(options) { | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| } | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
@@ -1,24 +0,11 @@ | ||
| import { ServerPlugin, Server } from 'srvx'; | ||
| import { W as WSOptions, S as ServerWithWSOptions } from '../shared/crossws.CP-89VBK.mjs'; | ||
| import '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| import '../adapters/bun.mjs'; | ||
| import 'bun'; | ||
| import '../adapters/deno.mjs'; | ||
| import '../adapters/node.mjs'; | ||
| import 'node:http'; | ||
| import 'node:stream'; | ||
| import 'events'; | ||
| import 'node:https'; | ||
| import 'node:tls'; | ||
| import 'node:url'; | ||
| import 'node:zlib'; | ||
| import '../adapters/sse.mjs'; | ||
| import '../adapters/cloudflare.mjs'; | ||
| import '@cloudflare/workers-types'; | ||
| import 'cloudflare:workers'; | ||
| import "../_chunks/bun.mjs"; | ||
| import "../_chunks/cloudflare.mjs"; | ||
| import "../_chunks/node.mjs"; | ||
| import { n as WSOptions, t as ServerWithWSOptions } from "../_chunks/_types.mjs"; | ||
| import { Server, ServerPlugin } from "srvx"; | ||
| //#region src/server/cloudflare.d.ts | ||
| declare function plugin(wsOpts: WSOptions): ServerPlugin; | ||
| declare function serve(options: ServerWithWSOptions): Server; | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
@@ -1,36 +0,27 @@ | ||
| import { serve as serve$1 } from 'srvx/cloudflare'; | ||
| import cloudflareAdapter from '../adapters/cloudflare.mjs'; | ||
| import 'cloudflare:workers'; | ||
| import '../shared/crossws.WpyOHUXc.mjs'; | ||
| import '../shared/crossws.CPlNx7g8.mjs'; | ||
| import '../shared/crossws.B31KJMcF.mjs'; | ||
| import '../shared/crossws.By9qWDAI.mjs'; | ||
| import cloudflare_default from "../adapters/cloudflare.mjs"; | ||
| import { serve as serve$1 } from "srvx/cloudflare"; | ||
| //#region src/server/cloudflare.ts | ||
| function plugin(wsOpts) { | ||
| return (server) => { | ||
| const ws = cloudflareAdapter({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.cloudflare | ||
| }); | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") { | ||
| return ws.handleUpgrade( | ||
| req, | ||
| req.runtime.cloudflare.env, | ||
| req.runtime.cloudflare.context | ||
| ); | ||
| } | ||
| return next(); | ||
| }); | ||
| }; | ||
| return (server) => { | ||
| const ws = cloudflare_default({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.cloudflare | ||
| }); | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") return ws.handleUpgrade(req, req.runtime.cloudflare.env, req.runtime.cloudflare.context); | ||
| return next(); | ||
| }); | ||
| }; | ||
| } | ||
| function serve(options) { | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| } | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
@@ -1,24 +0,11 @@ | ||
| import { ServerPlugin, Server } from 'srvx'; | ||
| import { W as WSOptions, S as ServerWithWSOptions } from '../shared/crossws.CP-89VBK.mjs'; | ||
| import '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| import '../adapters/bun.mjs'; | ||
| import 'bun'; | ||
| import '../adapters/deno.mjs'; | ||
| import '../adapters/node.mjs'; | ||
| import 'node:http'; | ||
| import 'node:stream'; | ||
| import 'events'; | ||
| import 'node:https'; | ||
| import 'node:tls'; | ||
| import 'node:url'; | ||
| import 'node:zlib'; | ||
| import '../adapters/sse.mjs'; | ||
| import '../adapters/cloudflare.mjs'; | ||
| import '@cloudflare/workers-types'; | ||
| import 'cloudflare:workers'; | ||
| import "../_chunks/bun.mjs"; | ||
| import "../_chunks/cloudflare.mjs"; | ||
| import "../_chunks/node.mjs"; | ||
| import { n as WSOptions, t as ServerWithWSOptions } from "../_chunks/_types.mjs"; | ||
| import { Server, ServerPlugin } from "srvx"; | ||
| //#region src/server/default.d.ts | ||
| declare function plugin(wsOpts: WSOptions): ServerPlugin; | ||
| declare function serve(options: ServerWithWSOptions): Server; | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
+22
-26
@@ -1,32 +0,28 @@ | ||
| import { serve as serve$1 } from 'srvx'; | ||
| import sseAdapter from '../adapters/sse.mjs'; | ||
| import '../shared/crossws.WpyOHUXc.mjs'; | ||
| import '../shared/crossws.CPlNx7g8.mjs'; | ||
| import sse_default from "../adapters/sse.mjs"; | ||
| import { serve as serve$1 } from "srvx"; | ||
| //#region src/server/default.ts | ||
| function plugin(wsOpts) { | ||
| const ws = sseAdapter({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.sse | ||
| }); | ||
| console.warn( | ||
| "[crossws] Using SSE adapter for WebSocket support. This requires a custom WebSocket client (https://crossws.h3.dev/adapters/sse)." | ||
| ); | ||
| return (server) => { | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") { | ||
| return ws.fetch(req); | ||
| } | ||
| return next(); | ||
| }); | ||
| }; | ||
| const ws = sse_default({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.sse | ||
| }); | ||
| console.warn("[crossws] Using SSE adapter for WebSocket support. This requires a custom WebSocket client (https://crossws.h3.dev/adapters/sse)."); | ||
| return (server) => { | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") return ws.fetch(req); | ||
| return next(); | ||
| }); | ||
| }; | ||
| } | ||
| function serve(options) { | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| } | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
@@ -1,24 +0,11 @@ | ||
| import { ServerPlugin, Server } from 'srvx'; | ||
| import { W as WSOptions, S as ServerWithWSOptions } from '../shared/crossws.CP-89VBK.mjs'; | ||
| import '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| import '../adapters/bun.mjs'; | ||
| import 'bun'; | ||
| import '../adapters/deno.mjs'; | ||
| import '../adapters/node.mjs'; | ||
| import 'node:http'; | ||
| import 'node:stream'; | ||
| import 'events'; | ||
| import 'node:https'; | ||
| import 'node:tls'; | ||
| import 'node:url'; | ||
| import 'node:zlib'; | ||
| import '../adapters/sse.mjs'; | ||
| import '../adapters/cloudflare.mjs'; | ||
| import '@cloudflare/workers-types'; | ||
| import 'cloudflare:workers'; | ||
| import "../_chunks/bun.mjs"; | ||
| import "../_chunks/cloudflare.mjs"; | ||
| import "../_chunks/node.mjs"; | ||
| import { n as WSOptions, t as ServerWithWSOptions } from "../_chunks/_types.mjs"; | ||
| import { Server, ServerPlugin } from "srvx"; | ||
| //#region src/server/deno.d.ts | ||
| declare function plugin(wsOpts: WSOptions): ServerPlugin; | ||
| declare function serve(options: ServerWithWSOptions): Server; | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
+21
-24
@@ -1,30 +0,27 @@ | ||
| import { serve as serve$1 } from 'srvx/deno'; | ||
| import denoAdapter from '../adapters/deno.mjs'; | ||
| import '../shared/crossws.WpyOHUXc.mjs'; | ||
| import '../shared/crossws.CPlNx7g8.mjs'; | ||
| import '../shared/crossws.By9qWDAI.mjs'; | ||
| import deno_default from "../adapters/deno.mjs"; | ||
| import { serve as serve$1 } from "srvx/deno"; | ||
| //#region src/server/deno.ts | ||
| function plugin(wsOpts) { | ||
| return (server) => { | ||
| const ws = denoAdapter({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.deno | ||
| }); | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") { | ||
| return ws.handleUpgrade(req, req.runtime.deno.info); | ||
| } | ||
| return next(); | ||
| }); | ||
| }; | ||
| return (server) => { | ||
| const ws = deno_default({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.deno | ||
| }); | ||
| server.options.middleware.unshift((req, next) => { | ||
| if (req.headers.get("upgrade")?.toLowerCase() === "websocket") return ws.handleUpgrade(req, req.runtime.deno.info); | ||
| return next(); | ||
| }); | ||
| }; | ||
| } | ||
| function serve(options) { | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| } | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
@@ -1,24 +0,11 @@ | ||
| import { ServerPlugin, Server } from 'srvx'; | ||
| import { W as WSOptions, S as ServerWithWSOptions } from '../shared/crossws.CP-89VBK.mjs'; | ||
| import '../index.mjs'; | ||
| import '../shared/crossws.BQXMA5bH.mjs'; | ||
| import '../adapters/bun.mjs'; | ||
| import 'bun'; | ||
| import '../adapters/deno.mjs'; | ||
| import '../adapters/node.mjs'; | ||
| import 'node:http'; | ||
| import 'node:stream'; | ||
| import 'events'; | ||
| import 'node:https'; | ||
| import 'node:tls'; | ||
| import 'node:url'; | ||
| import 'node:zlib'; | ||
| import '../adapters/sse.mjs'; | ||
| import '../adapters/cloudflare.mjs'; | ||
| import '@cloudflare/workers-types'; | ||
| import 'cloudflare:workers'; | ||
| import "../_chunks/bun.mjs"; | ||
| import "../_chunks/cloudflare.mjs"; | ||
| import "../_chunks/node.mjs"; | ||
| import { n as WSOptions, t as ServerWithWSOptions } from "../_chunks/_types.mjs"; | ||
| import { Server, ServerPlugin } from "srvx"; | ||
| //#region src/server/node.d.ts | ||
| declare function plugin(wsOpts: WSOptions): ServerPlugin; | ||
| declare function serve(options: ServerWithWSOptions): Server; | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
+32
-43
@@ -1,49 +0,38 @@ | ||
| import { NodeRequest, serve as serve$1 } from 'srvx/node'; | ||
| import nodeAdapter from '../adapters/node.mjs'; | ||
| import '../shared/crossws.WpyOHUXc.mjs'; | ||
| import '../shared/crossws.CPlNx7g8.mjs'; | ||
| import '../shared/crossws.By9qWDAI.mjs'; | ||
| import '../shared/crossws.C5pESzqN.mjs'; | ||
| import 'stream'; | ||
| import 'events'; | ||
| import 'http'; | ||
| import 'crypto'; | ||
| import 'buffer'; | ||
| import 'zlib'; | ||
| import 'https'; | ||
| import 'net'; | ||
| import 'tls'; | ||
| import 'url'; | ||
| import '../shared/crossws.B31KJMcF.mjs'; | ||
| import "../_chunks/rolldown-runtime.mjs"; | ||
| import "../_chunks/libs/ws.mjs"; | ||
| import node_default from "../adapters/node.mjs"; | ||
| import { NodeRequest, serve as serve$1 } from "srvx/node"; | ||
| //#region src/server/node.ts | ||
| function plugin(wsOpts) { | ||
| return (server) => { | ||
| const ws = nodeAdapter({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.deno | ||
| }); | ||
| const originalServe = server.serve; | ||
| server.serve = () => { | ||
| server.node?.server.on("upgrade", (req, socket, head) => { | ||
| ws.handleUpgrade( | ||
| req, | ||
| socket, | ||
| head, | ||
| // @ts-expect-error (upgrade is not typed) | ||
| new NodeRequest({ req, upgrade: { socket, head } }) | ||
| ); | ||
| }); | ||
| return originalServe.call(server); | ||
| }; | ||
| }; | ||
| return (server) => { | ||
| const ws = node_default({ | ||
| hooks: wsOpts, | ||
| resolve: wsOpts.resolve, | ||
| ...wsOpts.options?.deno | ||
| }); | ||
| const originalServe = server.serve; | ||
| server.serve = () => { | ||
| server.node?.server.on("upgrade", (req, socket, head) => { | ||
| ws.handleUpgrade(req, socket, head, new NodeRequest({ | ||
| req, | ||
| upgrade: { | ||
| socket, | ||
| head | ||
| } | ||
| })); | ||
| }); | ||
| return originalServe.call(server); | ||
| }; | ||
| }; | ||
| } | ||
| function serve(options) { | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| if (options.websocket) { | ||
| options.plugins ||= []; | ||
| options.plugins.push(plugin(options.websocket)); | ||
| } | ||
| return serve$1(options); | ||
| } | ||
| export { plugin, serve }; | ||
| //#endregion | ||
| export { plugin, serve }; |
@@ -0,3 +1,4 @@ | ||
| //#region src/websocket/native.d.ts | ||
| declare const WebSocket: typeof globalThis.WebSocket; | ||
| export { WebSocket as default }; | ||
| //#endregion | ||
| export { WebSocket as default }; |
@@ -0,3 +1,6 @@ | ||
| //#region src/websocket/native.ts | ||
| const WebSocket = globalThis.WebSocket; | ||
| var native_default = WebSocket; | ||
| export { WebSocket as default }; | ||
| //#endregion | ||
| export { native_default as default }; |
@@ -0,3 +1,4 @@ | ||
| //#region src/websocket/node.d.ts | ||
| declare const Websocket: typeof globalThis.WebSocket; | ||
| export { Websocket as default }; | ||
| //#endregion | ||
| export { Websocket as default }; |
@@ -1,15 +0,9 @@ | ||
| import { a as _WebSocket } from '../shared/crossws.C5pESzqN.mjs'; | ||
| import 'stream'; | ||
| import 'events'; | ||
| import 'http'; | ||
| import 'crypto'; | ||
| import 'buffer'; | ||
| import 'zlib'; | ||
| import 'https'; | ||
| import 'net'; | ||
| import 'tls'; | ||
| import 'url'; | ||
| import "../_chunks/rolldown-runtime.mjs"; | ||
| import { t as import_websocket } from "../_chunks/libs/ws.mjs"; | ||
| const Websocket = globalThis.WebSocket || _WebSocket; | ||
| //#region src/websocket/node.ts | ||
| const Websocket = globalThis.WebSocket || import_websocket.default; | ||
| var node_default = Websocket; | ||
| export { Websocket as default }; | ||
| //#endregion | ||
| export { node_default as default }; |
+34
-34
@@ -1,42 +0,42 @@ | ||
| import { E as EventTarget, W as WebSocket, C as CloseEvent, a as Event, M as MessageEvent } from '../shared/crossws.BQXMA5bH.mjs'; | ||
| import { a as WebSocket, i as MessageEvent, n as Event, r as EventTarget, t as CloseEvent } from "../_chunks/web.mjs"; | ||
| //#region src/websocket/sse.d.ts | ||
| type Ctor<T> = { | ||
| prototype: T; | ||
| new (): T; | ||
| prototype: T; | ||
| new (): T; | ||
| }; | ||
| declare const _EventTarget: Ctor<EventTarget>; | ||
| interface WebSocketSSEOptions { | ||
| protocols?: string | string[]; | ||
| /** enabled by default */ | ||
| bidir?: boolean; | ||
| /** enabled by default */ | ||
| stream?: boolean; | ||
| headers?: HeadersInit; | ||
| protocols?: string | string[]; | ||
| /** enabled by default */ | ||
| bidir?: boolean; | ||
| /** enabled by default */ | ||
| stream?: boolean; | ||
| headers?: HeadersInit; | ||
| } | ||
| declare class WebSocketSSE extends _EventTarget implements WebSocket { | ||
| #private; | ||
| static CONNECTING: number; | ||
| static OPEN: number; | ||
| static CLOSING: number; | ||
| static CLOSED: number; | ||
| readonly CONNECTING = 0; | ||
| readonly OPEN = 1; | ||
| readonly CLOSING = 2; | ||
| readonly CLOSED = 3; | ||
| onclose: ((this: WebSocket, ev: CloseEvent) => any) | null; | ||
| onerror: ((this: WebSocket, ev: Event) => any) | null; | ||
| onopen: ((this: WebSocket, ev: Event) => any) | null; | ||
| onmessage: ((this: WebSocket, ev: MessageEvent<any>) => any) | null; | ||
| binaryType: BinaryType; | ||
| readyState: number; | ||
| readonly url: string; | ||
| readonly protocol: string; | ||
| readonly extensions: string; | ||
| readonly bufferedAmount: number; | ||
| constructor(url: string, init?: string | string[] | WebSocketSSEOptions); | ||
| close(_code?: number, _reason?: string): void; | ||
| send(data: any): Promise<void>; | ||
| #private; | ||
| static CONNECTING: number; | ||
| static OPEN: number; | ||
| static CLOSING: number; | ||
| static CLOSED: number; | ||
| readonly CONNECTING = 0; | ||
| readonly OPEN = 1; | ||
| readonly CLOSING = 2; | ||
| readonly CLOSED = 3; | ||
| onclose: ((this: WebSocket, ev: CloseEvent) => any) | null; | ||
| onerror: ((this: WebSocket, ev: Event) => any) | null; | ||
| onopen: ((this: WebSocket, ev: Event) => any) | null; | ||
| onmessage: ((this: WebSocket, ev: MessageEvent<any>) => any) | null; | ||
| binaryType: BinaryType; | ||
| readyState: number; | ||
| readonly url: string; | ||
| readonly protocol: string; | ||
| readonly extensions: string; | ||
| readonly bufferedAmount: number; | ||
| constructor(url: string, init?: string | string[] | WebSocketSSEOptions); | ||
| close(_code?: number, _reason?: string): void; | ||
| send(data: any): Promise<void>; | ||
| } | ||
| export { WebSocketSSE }; | ||
| export type { WebSocketSSEOptions }; | ||
| //#endregion | ||
| export { WebSocketSSE, WebSocketSSEOptions }; |
+112
-121
@@ -0,125 +1,116 @@ | ||
| //#region src/websocket/sse.ts | ||
| const _EventTarget = EventTarget; | ||
| const defaultOptions = Object.freeze({ | ||
| bidir: true, | ||
| stream: true | ||
| bidir: true, | ||
| stream: true | ||
| }); | ||
| class WebSocketSSE extends _EventTarget { | ||
| static CONNECTING = 0; | ||
| static OPEN = 1; | ||
| static CLOSING = 2; | ||
| static CLOSED = 3; | ||
| CONNECTING = 0; | ||
| OPEN = 1; | ||
| CLOSING = 2; | ||
| CLOSED = 3; | ||
| onclose = null; | ||
| onerror = null; | ||
| onopen = null; | ||
| onmessage = null; | ||
| binaryType = "blob"; | ||
| readyState = WebSocketSSE.CONNECTING; | ||
| url; | ||
| protocol = ""; | ||
| extensions = ""; | ||
| bufferedAmount = 0; | ||
| #options = {}; | ||
| #sse; | ||
| #id; | ||
| #sendController; | ||
| #queue = []; | ||
| constructor(url, init) { | ||
| super(); | ||
| this.url = url.replace(/^ws/, "http"); | ||
| if (typeof init === "string") { | ||
| this.#options = { ...defaultOptions, protocols: init }; | ||
| } else if (Array.isArray(init)) { | ||
| this.#options = { ...defaultOptions, protocols: init }; | ||
| } else { | ||
| this.#options = { ...defaultOptions, ...init }; | ||
| } | ||
| this.#sse = new EventSource(this.url); | ||
| this.#sse.addEventListener("open", (_sseEvent) => { | ||
| this.readyState = WebSocketSSE.OPEN; | ||
| const event = new Event("open"); | ||
| this.onopen?.(event); | ||
| this.dispatchEvent(event); | ||
| }); | ||
| this.#sse.addEventListener("message", (sseEvent) => { | ||
| const _event = new MessageEvent("message", { | ||
| data: sseEvent.data | ||
| }); | ||
| this.onmessage?.(_event); | ||
| this.dispatchEvent(_event); | ||
| }); | ||
| if (this.#options.bidir) { | ||
| this.#sse.addEventListener("crossws-id", (sseEvent) => { | ||
| this.#id = sseEvent.data; | ||
| if (this.#options.stream) { | ||
| fetch(this.url, { | ||
| method: "POST", | ||
| // @ts-expect-error | ||
| duplex: "half", | ||
| headers: { | ||
| "content-type": "application/octet-stream", | ||
| "x-crossws-id": this.#id | ||
| }, | ||
| body: new ReadableStream({ | ||
| start: (controller) => { | ||
| this.#sendController = controller; | ||
| }, | ||
| cancel: () => { | ||
| this.#sendController = void 0; | ||
| } | ||
| }).pipeThrough(new TextEncoderStream()) | ||
| }).catch(() => { | ||
| }); | ||
| } | ||
| for (const data of this.#queue) { | ||
| this.send(data); | ||
| } | ||
| this.#queue = []; | ||
| }); | ||
| } | ||
| this.#sse.addEventListener("error", (_sseEvent) => { | ||
| const event = new Event("error"); | ||
| this.onerror?.(event); | ||
| this.dispatchEvent(event); | ||
| }); | ||
| this.#sse.addEventListener("close", (_sseEvent) => { | ||
| this.readyState = WebSocketSSE.CLOSED; | ||
| const event = new Event("close"); | ||
| this.onclose?.(event); | ||
| this.dispatchEvent(event); | ||
| }); | ||
| } | ||
| close(_code, _reason) { | ||
| this.readyState = WebSocketSSE.CLOSING; | ||
| this.#sse.close(); | ||
| this.readyState = WebSocketSSE.CLOSED; | ||
| } | ||
| async send(data) { | ||
| if (!this.#options.bidir) { | ||
| throw new Error("bdir option is not enabled!"); | ||
| } | ||
| if (this.readyState !== WebSocketSSE.OPEN) { | ||
| throw new Error("WebSocket is not open!"); | ||
| } | ||
| if (!this.#id) { | ||
| this.#queue.push(data); | ||
| return; | ||
| } | ||
| if (this.#sendController) { | ||
| this.#sendController.enqueue(data); | ||
| return; | ||
| } | ||
| await fetch(this.url, { | ||
| method: "POST", | ||
| headers: { | ||
| "x-crossws-id": this.#id | ||
| }, | ||
| body: data | ||
| }); | ||
| } | ||
| } | ||
| var WebSocketSSE = class WebSocketSSE extends _EventTarget { | ||
| static CONNECTING = 0; | ||
| static OPEN = 1; | ||
| static CLOSING = 2; | ||
| static CLOSED = 3; | ||
| CONNECTING = 0; | ||
| OPEN = 1; | ||
| CLOSING = 2; | ||
| CLOSED = 3; | ||
| onclose = null; | ||
| onerror = null; | ||
| onopen = null; | ||
| onmessage = null; | ||
| binaryType = "blob"; | ||
| readyState = WebSocketSSE.CONNECTING; | ||
| url; | ||
| protocol = ""; | ||
| extensions = ""; | ||
| bufferedAmount = 0; | ||
| #options = {}; | ||
| #sse; | ||
| #id; | ||
| #sendController; | ||
| #queue = []; | ||
| constructor(url, init) { | ||
| super(); | ||
| this.url = url.replace(/^ws/, "http"); | ||
| if (typeof init === "string") this.#options = { | ||
| ...defaultOptions, | ||
| protocols: init | ||
| }; | ||
| else if (Array.isArray(init)) this.#options = { | ||
| ...defaultOptions, | ||
| protocols: init | ||
| }; | ||
| else this.#options = { | ||
| ...defaultOptions, | ||
| ...init | ||
| }; | ||
| this.#sse = new EventSource(this.url); | ||
| this.#sse.addEventListener("open", (_sseEvent) => { | ||
| this.readyState = WebSocketSSE.OPEN; | ||
| const event = new Event("open"); | ||
| this.onopen?.(event); | ||
| this.dispatchEvent(event); | ||
| }); | ||
| this.#sse.addEventListener("message", (sseEvent) => { | ||
| const _event = new MessageEvent("message", { data: sseEvent.data }); | ||
| this.onmessage?.(_event); | ||
| this.dispatchEvent(_event); | ||
| }); | ||
| if (this.#options.bidir) this.#sse.addEventListener("crossws-id", (sseEvent) => { | ||
| this.#id = sseEvent.data; | ||
| if (this.#options.stream) fetch(this.url, { | ||
| method: "POST", | ||
| duplex: "half", | ||
| headers: { | ||
| "content-type": "application/octet-stream", | ||
| "x-crossws-id": this.#id | ||
| }, | ||
| body: new ReadableStream({ | ||
| start: (controller) => { | ||
| this.#sendController = controller; | ||
| }, | ||
| cancel: () => { | ||
| this.#sendController = void 0; | ||
| } | ||
| }).pipeThrough(new TextEncoderStream()) | ||
| }).catch(() => {}); | ||
| for (const data of this.#queue) this.send(data); | ||
| this.#queue = []; | ||
| }); | ||
| this.#sse.addEventListener("error", (_sseEvent) => { | ||
| const event = new Event("error"); | ||
| this.onerror?.(event); | ||
| this.dispatchEvent(event); | ||
| }); | ||
| this.#sse.addEventListener("close", (_sseEvent) => { | ||
| this.readyState = WebSocketSSE.CLOSED; | ||
| const event = new Event("close"); | ||
| this.onclose?.(event); | ||
| this.dispatchEvent(event); | ||
| }); | ||
| } | ||
| close(_code, _reason) { | ||
| this.readyState = WebSocketSSE.CLOSING; | ||
| this.#sse.close(); | ||
| this.readyState = WebSocketSSE.CLOSED; | ||
| } | ||
| async send(data) { | ||
| if (!this.#options.bidir) throw new Error("bdir option is not enabled!"); | ||
| if (this.readyState !== WebSocketSSE.OPEN) throw new Error("WebSocket is not open!"); | ||
| if (!this.#id) { | ||
| this.#queue.push(data); | ||
| return; | ||
| } | ||
| if (this.#sendController) { | ||
| this.#sendController.enqueue(data); | ||
| return; | ||
| } | ||
| await fetch(this.url, { | ||
| method: "POST", | ||
| headers: { "x-crossws-id": this.#id }, | ||
| body: data | ||
| }); | ||
| } | ||
| }; | ||
| export { WebSocketSSE }; | ||
| //#endregion | ||
| export { WebSocketSSE }; |
+3
-2
| { | ||
| "name": "crossws", | ||
| "version": "0.4.2", | ||
| "version": "0.4.3", | ||
| "description": "Cross-platform WebSocket Servers for Node.js, Deno, Bun and Cloudflare Workers", | ||
@@ -51,3 +51,3 @@ "homepage": "https://crossws.h3.dev", | ||
| "scripts": { | ||
| "build": "unbuild", | ||
| "build": "obuild", | ||
| "dev": "vitest", | ||
@@ -91,2 +91,3 @@ "lint": "eslint --cache . && prettier -c src test", | ||
| "listhen": "^1.9.0", | ||
| "obuild": "^0.4.16", | ||
| "prettier": "^3.8.0", | ||
@@ -93,0 +94,0 @@ "srvx": "^0.10.1", |
| const StubRequest = /* @__PURE__ */ (() => { | ||
| class StubRequest2 { | ||
| url; | ||
| _signal; | ||
| _headers; | ||
| _init; | ||
| constructor(url, init = {}) { | ||
| this.url = url; | ||
| this._init = init; | ||
| } | ||
| get headers() { | ||
| if (!this._headers) { | ||
| this._headers = new Headers(this._init?.headers); | ||
| } | ||
| return this._headers; | ||
| } | ||
| clone() { | ||
| return new StubRequest2(this.url, this._init); | ||
| } | ||
| // --- dummy --- | ||
| get method() { | ||
| return "GET"; | ||
| } | ||
| get signal() { | ||
| return this._signal ??= new AbortSignal(); | ||
| } | ||
| get cache() { | ||
| return "default"; | ||
| } | ||
| get credentials() { | ||
| return "same-origin"; | ||
| } | ||
| get destination() { | ||
| return ""; | ||
| } | ||
| get integrity() { | ||
| return ""; | ||
| } | ||
| get keepalive() { | ||
| return false; | ||
| } | ||
| get redirect() { | ||
| return "follow"; | ||
| } | ||
| get mode() { | ||
| return "cors"; | ||
| } | ||
| get referrer() { | ||
| return "about:client"; | ||
| } | ||
| get referrerPolicy() { | ||
| return ""; | ||
| } | ||
| get body() { | ||
| return null; | ||
| } | ||
| get bodyUsed() { | ||
| return false; | ||
| } | ||
| arrayBuffer() { | ||
| return Promise.resolve(new ArrayBuffer(0)); | ||
| } | ||
| blob() { | ||
| return Promise.resolve(new Blob()); | ||
| } | ||
| bytes() { | ||
| return Promise.resolve(new Uint8Array()); | ||
| } | ||
| formData() { | ||
| return Promise.resolve(new FormData()); | ||
| } | ||
| json() { | ||
| return Promise.resolve(JSON.parse("")); | ||
| } | ||
| text() { | ||
| return Promise.resolve(""); | ||
| } | ||
| } | ||
| Object.setPrototypeOf(StubRequest2.prototype, globalThis.Request.prototype); | ||
| return StubRequest2; | ||
| })(); | ||
| export { StubRequest as S }; |
| /** | ||
| * A CloseEvent is sent to clients using WebSockets when the connection is closed. This is delivered to the listener indicated by the WebSocket object's onclose attribute. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent) | ||
| */ | ||
| interface CloseEvent extends Event { | ||
| /** | ||
| * Returns the WebSocket connection close code provided by the server. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent/code) | ||
| */ | ||
| readonly code: number; | ||
| /** | ||
| * Returns the WebSocket connection close reason provided by the server. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent/reason) | ||
| */ | ||
| readonly reason: string; | ||
| /** | ||
| * Returns true if the connection closed cleanly; false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/CloseEvent/wasClean) | ||
| */ | ||
| readonly wasClean: boolean; | ||
| } | ||
| /** | ||
| * An event which takes place in the DOM. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event) | ||
| */ | ||
| interface Event { | ||
| /** | ||
| * Returns true or false depending on how event was initialized. True if event goes through its target's ancestors in reverse tree order, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/bubbles) | ||
| */ | ||
| readonly bubbles: boolean; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/cancelBubble) | ||
| */ | ||
| cancelBubble: boolean; | ||
| /** | ||
| * Returns true or false depending on how event was initialized. Its return value does not always carry meaning, but true can indicate that part of the operation during which event was dispatched, can be canceled by invoking the preventDefault() method. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/cancelable) | ||
| */ | ||
| readonly cancelable: boolean; | ||
| /** | ||
| * Returns true or false depending on how event was initialized. True if event invokes listeners past a ShadowRoot node that is the root of its target, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/composed) | ||
| */ | ||
| readonly composed: boolean; | ||
| /** | ||
| * Returns the object whose event listener's callback is currently being invoked. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/currentTarget) | ||
| */ | ||
| readonly currentTarget: EventTarget | null; | ||
| /** | ||
| * Returns true if preventDefault() was invoked successfully to indicate cancelation, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/defaultPrevented) | ||
| */ | ||
| readonly defaultPrevented: boolean; | ||
| /** | ||
| * Returns the event's phase, which is one of NONE, CAPTURING_PHASE, AT_TARGET, and BUBBLING_PHASE. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/eventPhase) | ||
| */ | ||
| readonly eventPhase: number; | ||
| /** | ||
| * Returns true if event was dispatched by the user agent, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/isTrusted) | ||
| */ | ||
| readonly isTrusted: boolean; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/returnValue) | ||
| */ | ||
| returnValue: boolean; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/srcElement) | ||
| */ | ||
| readonly srcElement: EventTarget | null; | ||
| /** | ||
| * Returns the object to which event is dispatched (its target). | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/target) | ||
| */ | ||
| readonly target: EventTarget | null; | ||
| /** | ||
| * Returns the event's timestamp as the number of milliseconds measured relative to the time origin. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/timeStamp) | ||
| */ | ||
| readonly timeStamp: DOMHighResTimeStamp; | ||
| /** | ||
| * Returns the type of event, e.g. "click", "hashchange", or "submit". | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/type) | ||
| */ | ||
| readonly type: string; | ||
| /** | ||
| * Returns the invocation target objects of event's path (objects on which listeners will be invoked), except for any nodes in shadow trees of which the shadow root's mode is "closed" that are not reachable from event's currentTarget. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/composedPath) | ||
| */ | ||
| composedPath(): EventTarget[]; | ||
| /** | ||
| * @deprecated | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/initEvent) | ||
| */ | ||
| initEvent(type: string, bubbles?: boolean, cancelable?: boolean): void; | ||
| /** | ||
| * If invoked when the cancelable attribute value is true, and while executing a listener for the event with passive set to false, signals to the operation that caused event to be dispatched that it needs to be canceled. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/preventDefault) | ||
| */ | ||
| preventDefault(): void; | ||
| /** | ||
| * Invoking this method prevents event from reaching any registered event listeners after the current one finishes running and, when dispatched in a tree, also prevents event from reaching any other objects. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/stopImmediatePropagation) | ||
| */ | ||
| stopImmediatePropagation(): void; | ||
| /** | ||
| * When dispatched in a tree, invoking this method prevents event from reaching any objects other than the current object. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Event/stopPropagation) | ||
| */ | ||
| stopPropagation(): void; | ||
| readonly NONE: 0; | ||
| readonly CAPTURING_PHASE: 1; | ||
| readonly AT_TARGET: 2; | ||
| readonly BUBBLING_PHASE: 3; | ||
| } | ||
| /** | ||
| * EventTarget is a DOM interface implemented by objects that can receive events and may have listeners for them. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget) | ||
| */ | ||
| interface EventTarget { | ||
| /** | ||
| * Appends an event listener for events whose type attribute value is type. The callback argument sets the callback that will be invoked when the event is dispatched. | ||
| * | ||
| * The options argument sets listener-specific options. For compatibility this can be a boolean, in which case the method behaves exactly as if the value was specified as options's capture. | ||
| * | ||
| * When set to true, options's capture prevents callback from being invoked when the event's eventPhase attribute value is BUBBLING_PHASE. When false (or not present), callback will not be invoked when event's eventPhase attribute value is CAPTURING_PHASE. Either way, callback will be invoked if event's eventPhase attribute value is AT_TARGET. | ||
| * | ||
| * When set to true, options's passive indicates that the callback will not cancel the event by invoking preventDefault(). This is used to enable performance optimizations described in § 2.8 Observing event listeners. | ||
| * | ||
| * When set to true, options's once indicates that the callback will only be invoked once after which the event listener will be removed. | ||
| * | ||
| * If an AbortSignal is passed for options's signal, then the event listener will be removed when signal is aborted. | ||
| * | ||
| * The event listener is appended to target's event listener list and is not appended if it has the same type, callback, and capture. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget/addEventListener) | ||
| */ | ||
| addEventListener(type: string, callback: EventListenerOrEventListenerObject | null, options?: AddEventListenerOptions | boolean): void; | ||
| /** | ||
| * Dispatches a synthetic event event to target and returns true if either event's cancelable attribute value is false or its preventDefault() method was not invoked, and false otherwise. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget/dispatchEvent) | ||
| */ | ||
| dispatchEvent(event: Event): boolean; | ||
| /** | ||
| * Removes the event listener in target's event listener list with the same type, callback, and options. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/EventTarget/removeEventListener) | ||
| */ | ||
| removeEventListener(type: string, callback: EventListenerOrEventListenerObject | null, options?: EventListenerOptions | boolean): void; | ||
| } | ||
| /** | ||
| * A message received by a target object. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent) | ||
| */ | ||
| interface MessageEvent<T = any> extends Event { | ||
| /** | ||
| * Returns the data of the message. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/data) | ||
| */ | ||
| readonly data: T; | ||
| /** | ||
| * Returns the last event ID string, for server-sent events. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/lastEventId) | ||
| */ | ||
| readonly lastEventId: string; | ||
| /** | ||
| * Returns the origin of the message, for server-sent events and cross-document messaging. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/origin) | ||
| */ | ||
| readonly origin: string; | ||
| /** | ||
| * Returns the MessagePort array sent with the message, for cross-document messaging and channel messaging. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/ports) | ||
| */ | ||
| readonly ports: ReadonlyArray<MessagePort>; | ||
| /** | ||
| * Returns the WindowProxy of the source window, for cross-document messaging, and the MessagePort being attached, in the connect event fired at SharedWorkerGlobalScope objects. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/MessageEvent/source) | ||
| */ | ||
| readonly source: MessageEventSource | null; | ||
| /** @deprecated */ | ||
| initMessageEvent(type: string, bubbles?: boolean, cancelable?: boolean, data?: any, origin?: string, lastEventId?: string, source?: MessageEventSource | null, ports?: MessagePort[]): void; | ||
| } | ||
| /** | ||
| * Provides the API for creating and managing a WebSocket connection to a server, as well as for sending and receiving data on the connection. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket) | ||
| */ | ||
| interface WebSocket extends EventTarget { | ||
| /** | ||
| * Returns a string that indicates how binary data from the WebSocket object is exposed to scripts: | ||
| * | ||
| * Can be set, to change how binary data is returned. The default is "blob". | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/binaryType) | ||
| */ | ||
| binaryType: BinaryType | (string & {}); | ||
| /** | ||
| * Returns the number of bytes of application data (UTF-8 text and binary data) that have been queued using send() but not yet been transmitted to the network. | ||
| * | ||
| * If the WebSocket connection is closed, this attribute's value will only increase with each call to the send() method. (The number does not reset to zero once the connection closes.) | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/bufferedAmount) | ||
| */ | ||
| readonly bufferedAmount: number; | ||
| /** | ||
| * Returns the extensions selected by the server, if any. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/extensions) | ||
| */ | ||
| readonly extensions: string; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/close_event) */ | ||
| onclose: ((this: WebSocket, ev: CloseEvent) => any) | null; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/error_event) */ | ||
| onerror: ((this: WebSocket, ev: Event) => any) | null; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/message_event) */ | ||
| onmessage: ((this: WebSocket, ev: MessageEvent) => any) | null; | ||
| /** [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/open_event) */ | ||
| onopen: ((this: WebSocket, ev: Event) => any) | null; | ||
| /** | ||
| * Returns the subprotocol selected by the server, if any. It can be used in conjunction with the array form of the constructor's second argument to perform subprotocol negotiation. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/protocol) | ||
| */ | ||
| readonly protocol: string; | ||
| /** | ||
| * Returns the state of the WebSocket object's connection. It can have the values described below. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/readyState) | ||
| */ | ||
| readonly readyState: number; | ||
| /** | ||
| * Returns the URL that was used to establish the WebSocket connection. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/url) | ||
| */ | ||
| readonly url: string; | ||
| /** | ||
| * Closes the WebSocket connection, optionally using code as the the WebSocket connection close code and reason as the the WebSocket connection close reason. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/close) | ||
| */ | ||
| close(code?: number, reason?: string): void; | ||
| /** | ||
| * Transmits data using the WebSocket connection. data can be a string, a Blob, an ArrayBuffer, or an ArrayBufferView. | ||
| * | ||
| * [MDN Reference](https://developer.mozilla.org/docs/Web/API/WebSocket/send) | ||
| */ | ||
| send(data: string | ArrayBufferLike | Blob | ArrayBufferView): void; | ||
| readonly CONNECTING: 0; | ||
| readonly OPEN: 1; | ||
| readonly CLOSING: 2; | ||
| readonly CLOSED: 3; | ||
| addEventListener<K extends keyof WebSocketEventMap>(type: K, listener: (this: WebSocket, ev: WebSocketEventMap[K]) => any, options?: boolean | AddEventListenerOptions): void; | ||
| addEventListener(type: string, listener: EventListenerOrEventListenerObject, options?: boolean | AddEventListenerOptions): void; | ||
| removeEventListener<K extends keyof WebSocketEventMap>(type: K, listener: (this: WebSocket, ev: WebSocketEventMap[K]) => any, options?: boolean | EventListenerOptions): void; | ||
| removeEventListener(type: string, listener: EventListenerOrEventListenerObject, options?: boolean | EventListenerOptions): void; | ||
| } | ||
| export type { CloseEvent as C, EventTarget as E, MessageEvent as M, WebSocket as W, Event as a }; |
| class WSError extends Error { | ||
| constructor(...args) { | ||
| super(...args); | ||
| this.name = "WSError"; | ||
| } | ||
| } | ||
| export { WSError as W }; |
Sorry, the diff of this file is too big to display
| import { ServerRequest, ServerOptions } from 'srvx'; | ||
| import { Hooks } from '../index.mjs'; | ||
| import { BunOptions } from '../adapters/bun.mjs'; | ||
| import { DenoOptions } from '../adapters/deno.mjs'; | ||
| import { NodeOptions } from '../adapters/node.mjs'; | ||
| import { SSEOptions } from '../adapters/sse.mjs'; | ||
| import { CloudflareOptions } from '../adapters/cloudflare.mjs'; | ||
| type WSOptions = Partial<Hooks> & { | ||
| resolve?: (req: ServerRequest) => Partial<Hooks> | Promise<Partial<Hooks>>; | ||
| options?: { | ||
| bun?: BunOptions; | ||
| deno?: DenoOptions; | ||
| node?: NodeOptions; | ||
| sse?: SSEOptions; | ||
| cloudflare?: CloudflareOptions; | ||
| }; | ||
| }; | ||
| type ServerWithWSOptions = ServerOptions & { | ||
| websocket?: WSOptions; | ||
| }; | ||
| export type { ServerWithWSOptions as S, WSOptions as W }; |
| class AdapterHookable { | ||
| options; | ||
| constructor(options) { | ||
| this.options = options || {}; | ||
| } | ||
| callHook(name, arg1, arg2) { | ||
| const globalHook = this.options.hooks?.[name]; | ||
| const globalPromise = globalHook?.(arg1, arg2); | ||
| const request = arg1.request || arg1; | ||
| const resolveHooksPromise = this.options.resolve?.(request); | ||
| if (!resolveHooksPromise) { | ||
| return globalPromise; | ||
| } | ||
| const resolvePromise = resolveHooksPromise instanceof Promise ? resolveHooksPromise.then((hooks) => hooks?.[name]) : resolveHooksPromise?.[name]; | ||
| return Promise.all([globalPromise, resolvePromise]).then( | ||
| ([globalRes, hook]) => { | ||
| const hookResPromise = hook?.(arg1, arg2); | ||
| return hookResPromise instanceof Promise ? hookResPromise.then((hookRes) => hookRes || globalRes) : hookResPromise || globalRes; | ||
| } | ||
| ); | ||
| } | ||
| async upgrade(request) { | ||
| let namespace = this.options.getNamespace?.(request) ?? new URL(request.url).pathname; | ||
| const context = request.context || {}; | ||
| try { | ||
| const res = await this.callHook( | ||
| "upgrade", | ||
| request | ||
| ); | ||
| if (!res) { | ||
| return { context, namespace }; | ||
| } | ||
| if (res.namespace) { | ||
| namespace = res.namespace; | ||
| } | ||
| if (res.context) { | ||
| Object.assign( | ||
| context, | ||
| res.context | ||
| ); | ||
| } | ||
| if (res instanceof Response) { | ||
| return { context, namespace, endResponse: res }; | ||
| } | ||
| if (res.headers) { | ||
| return { | ||
| context, | ||
| namespace, | ||
| upgradeHeaders: res.headers | ||
| }; | ||
| } | ||
| } catch (error) { | ||
| const errResponse = error.response || error; | ||
| if (errResponse instanceof Response) { | ||
| return { | ||
| context, | ||
| namespace, | ||
| endResponse: errResponse | ||
| }; | ||
| } | ||
| throw error; | ||
| } | ||
| return { context, namespace }; | ||
| } | ||
| } | ||
| function defineHooks(hooks) { | ||
| return hooks; | ||
| } | ||
| function adapterUtils(globalPeers) { | ||
| return { | ||
| peers: globalPeers, | ||
| publish(topic, message, options) { | ||
| for (const peers of options?.namespace ? [globalPeers.get(options.namespace) || []] : globalPeers.values()) { | ||
| let firstPeerWithTopic; | ||
| for (const peer of peers) { | ||
| if (peer.topics.has(topic)) { | ||
| firstPeerWithTopic = peer; | ||
| break; | ||
| } | ||
| } | ||
| if (firstPeerWithTopic) { | ||
| firstPeerWithTopic.send(message, options); | ||
| firstPeerWithTopic.publish(topic, message, options); | ||
| } | ||
| } | ||
| } | ||
| }; | ||
| } | ||
| function getPeers(globalPeers, namespace) { | ||
| if (!namespace) { | ||
| throw new Error("Websocket publish namespace missing."); | ||
| } | ||
| let peers = globalPeers.get(namespace); | ||
| if (!peers) { | ||
| peers = /* @__PURE__ */ new Set(); | ||
| globalPeers.set(namespace, peers); | ||
| } | ||
| return peers; | ||
| } | ||
| function defineWebSocketAdapter(factory) { | ||
| return factory; | ||
| } | ||
| export { AdapterHookable as A, adapterUtils as a, defineWebSocketAdapter as b, defineHooks as d, getPeers as g }; |
| const kNodeInspect = /* @__PURE__ */ Symbol.for( | ||
| "nodejs.util.inspect.custom" | ||
| ); | ||
| function toBufferLike(val) { | ||
| if (val === void 0 || val === null) { | ||
| return ""; | ||
| } | ||
| const type = typeof val; | ||
| if (type === "string") { | ||
| return val; | ||
| } | ||
| if (type === "number" || type === "boolean" || type === "bigint") { | ||
| return val.toString(); | ||
| } | ||
| if (type === "function" || type === "symbol") { | ||
| return "{}"; | ||
| } | ||
| if (val instanceof Uint8Array || val instanceof ArrayBuffer) { | ||
| return val; | ||
| } | ||
| if (isPlainObject(val)) { | ||
| return JSON.stringify(val); | ||
| } | ||
| return val; | ||
| } | ||
| function toString(val) { | ||
| if (typeof val === "string") { | ||
| return val; | ||
| } | ||
| const data = toBufferLike(val); | ||
| if (typeof data === "string") { | ||
| return data; | ||
| } | ||
| const base64 = btoa(String.fromCharCode(...new Uint8Array(data))); | ||
| return `data:application/octet-stream;base64,${base64}`; | ||
| } | ||
| function isPlainObject(value) { | ||
| if (value === null || typeof value !== "object") { | ||
| return false; | ||
| } | ||
| const prototype = Object.getPrototypeOf(value); | ||
| if (prototype !== null && prototype !== Object.prototype && Object.getPrototypeOf(prototype) !== null) { | ||
| return false; | ||
| } | ||
| if (Symbol.iterator in value) { | ||
| return false; | ||
| } | ||
| if (Symbol.toStringTag in value) { | ||
| return Object.prototype.toString.call(value) === "[object Module]"; | ||
| } | ||
| return true; | ||
| } | ||
| class Message { | ||
| /** Access to the original [message event](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket/message_event) if available. */ | ||
| event; | ||
| /** Access to the Peer that emitted the message. */ | ||
| peer; | ||
| /** Raw message data (can be of any type). */ | ||
| rawData; | ||
| #id; | ||
| #uint8Array; | ||
| #arrayBuffer; | ||
| #blob; | ||
| #text; | ||
| #json; | ||
| constructor(rawData, peer, event) { | ||
| this.rawData = rawData || ""; | ||
| this.peer = peer; | ||
| this.event = event; | ||
| } | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the message. | ||
| */ | ||
| get id() { | ||
| if (!this.#id) { | ||
| this.#id = crypto.randomUUID(); | ||
| } | ||
| return this.#id; | ||
| } | ||
| // --- data views --- | ||
| /** | ||
| * Get data as [Uint8Array](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Uint8Array) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| uint8Array() { | ||
| const _uint8Array = this.#uint8Array; | ||
| if (_uint8Array) { | ||
| return _uint8Array; | ||
| } | ||
| const rawData = this.rawData; | ||
| if (rawData instanceof Uint8Array) { | ||
| return this.#uint8Array = rawData; | ||
| } | ||
| if (rawData instanceof ArrayBuffer || rawData instanceof SharedArrayBuffer) { | ||
| this.#arrayBuffer = rawData; | ||
| return this.#uint8Array = new Uint8Array(rawData); | ||
| } | ||
| if (typeof rawData === "string") { | ||
| this.#text = rawData; | ||
| return this.#uint8Array = new TextEncoder().encode(this.#text); | ||
| } | ||
| if (Symbol.iterator in rawData) { | ||
| return this.#uint8Array = new Uint8Array(rawData); | ||
| } | ||
| if (typeof rawData?.length === "number") { | ||
| return this.#uint8Array = new Uint8Array(rawData); | ||
| } | ||
| if (rawData instanceof DataView) { | ||
| return this.#uint8Array = new Uint8Array( | ||
| rawData.buffer, | ||
| rawData.byteOffset, | ||
| rawData.byteLength | ||
| ); | ||
| } | ||
| throw new TypeError( | ||
| `Unsupported message type: ${Object.prototype.toString.call(rawData)}` | ||
| ); | ||
| } | ||
| /** | ||
| * Get data as [ArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/ArrayBuffer) or [SharedArrayBuffer](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/SharedArrayBuffer) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. | ||
| */ | ||
| arrayBuffer() { | ||
| const _arrayBuffer = this.#arrayBuffer; | ||
| if (_arrayBuffer) { | ||
| return _arrayBuffer; | ||
| } | ||
| const rawData = this.rawData; | ||
| if (rawData instanceof ArrayBuffer || rawData instanceof SharedArrayBuffer) { | ||
| return this.#arrayBuffer = rawData; | ||
| } | ||
| return this.#arrayBuffer = this.uint8Array().buffer; | ||
| } | ||
| /** | ||
| * Get data as [Blob](https://developer.mozilla.org/en-US/docs/Web/API/Blob) value. | ||
| * | ||
| * If raw data is in any other format or string, it will be automatically converted and encoded. */ | ||
| blob() { | ||
| const _blob = this.#blob; | ||
| if (_blob) { | ||
| return _blob; | ||
| } | ||
| const rawData = this.rawData; | ||
| if (rawData instanceof Blob) { | ||
| return this.#blob = rawData; | ||
| } | ||
| return this.#blob = new Blob([this.uint8Array()]); | ||
| } | ||
| /** | ||
| * Get stringified text version of the message. | ||
| * | ||
| * If raw data is in any other format, it will be automatically converted and decoded. | ||
| */ | ||
| text() { | ||
| const _text = this.#text; | ||
| if (_text) { | ||
| return _text; | ||
| } | ||
| const rawData = this.rawData; | ||
| if (typeof rawData === "string") { | ||
| return this.#text = rawData; | ||
| } | ||
| return this.#text = new TextDecoder().decode(this.uint8Array()); | ||
| } | ||
| /** | ||
| * Get parsed version of the message text with [`JSON.parse()`](https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/JSON/parse). | ||
| */ | ||
| json() { | ||
| const _json = this.#json; | ||
| if (_json) { | ||
| return _json; | ||
| } | ||
| return this.#json = JSON.parse(this.text()); | ||
| } | ||
| /** | ||
| * Message data (value varies based on `peer.websocket.binaryType`). | ||
| */ | ||
| get data() { | ||
| switch (this.peer?.websocket?.binaryType) { | ||
| case "arraybuffer": { | ||
| return this.arrayBuffer(); | ||
| } | ||
| case "blob": { | ||
| return this.blob(); | ||
| } | ||
| case "nodebuffer": { | ||
| return globalThis.Buffer ? Buffer.from(this.uint8Array()) : this.uint8Array(); | ||
| } | ||
| case "uint8array": { | ||
| return this.uint8Array(); | ||
| } | ||
| case "text": { | ||
| return this.text(); | ||
| } | ||
| default: { | ||
| return this.rawData; | ||
| } | ||
| } | ||
| } | ||
| // --- inspect --- | ||
| toString() { | ||
| return this.text(); | ||
| } | ||
| [Symbol.toPrimitive]() { | ||
| return this.text(); | ||
| } | ||
| [kNodeInspect]() { | ||
| return { | ||
| message: { | ||
| id: this.id, | ||
| peer: this.peer, | ||
| text: this.text() | ||
| } | ||
| }; | ||
| } | ||
| } | ||
| class Peer { | ||
| _internal; | ||
| _topics; | ||
| _id; | ||
| #ws; | ||
| constructor(internal) { | ||
| this._topics = /* @__PURE__ */ new Set(); | ||
| this._internal = internal; | ||
| } | ||
| get context() { | ||
| return this._internal.context ??= {}; | ||
| } | ||
| get namespace() { | ||
| return this._internal.namespace; | ||
| } | ||
| /** | ||
| * Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the peer. | ||
| */ | ||
| get id() { | ||
| if (!this._id) { | ||
| this._id = crypto.randomUUID(); | ||
| } | ||
| return this._id; | ||
| } | ||
| /** IP address of the peer */ | ||
| get remoteAddress() { | ||
| return void 0; | ||
| } | ||
| /** upgrade request */ | ||
| get request() { | ||
| return this._internal.request; | ||
| } | ||
| /** | ||
| * Get the [WebSocket](https://developer.mozilla.org/en-US/docs/Web/API/WebSocket) instance. | ||
| * | ||
| * **Note:** crossws adds polyfill for the following properties if native values are not available: | ||
| * - `protocol`: Extracted from the `sec-websocket-protocol` header. | ||
| * - `extensions`: Extracted from the `sec-websocket-extensions` header. | ||
| * - `url`: Extracted from the request URL (http -> ws). | ||
| * */ | ||
| get websocket() { | ||
| if (!this.#ws) { | ||
| const _ws = this._internal.ws; | ||
| const _request = this._internal.request; | ||
| this.#ws = _request ? createWsProxy(_ws, _request) : _ws; | ||
| } | ||
| return this.#ws; | ||
| } | ||
| /** All connected peers to the server */ | ||
| get peers() { | ||
| return this._internal.peers || /* @__PURE__ */ new Set(); | ||
| } | ||
| /** All topics, this peer has been subscribed to. */ | ||
| get topics() { | ||
| return this._topics; | ||
| } | ||
| /** Abruptly close the connection */ | ||
| terminate() { | ||
| this.close(); | ||
| } | ||
| /** Subscribe to a topic */ | ||
| subscribe(topic) { | ||
| this._topics.add(topic); | ||
| } | ||
| /** Unsubscribe from a topic */ | ||
| unsubscribe(topic) { | ||
| this._topics.delete(topic); | ||
| } | ||
| // --- inspect --- | ||
| toString() { | ||
| return this.id; | ||
| } | ||
| [Symbol.toPrimitive]() { | ||
| return this.id; | ||
| } | ||
| [Symbol.toStringTag]() { | ||
| return "WebSocket"; | ||
| } | ||
| [kNodeInspect]() { | ||
| return { | ||
| peer: { | ||
| id: this.id, | ||
| ip: this.remoteAddress | ||
| } | ||
| }; | ||
| } | ||
| } | ||
| function createWsProxy(ws, request) { | ||
| return new Proxy(ws, { | ||
| get: (target, prop) => { | ||
| const value = Reflect.get(target, prop); | ||
| if (!value) { | ||
| switch (prop) { | ||
| case "protocol": { | ||
| return request?.headers?.get("sec-websocket-protocol") || ""; | ||
| } | ||
| case "extensions": { | ||
| return request?.headers?.get("sec-websocket-extensions") || ""; | ||
| } | ||
| case "url": { | ||
| return request?.url?.replace(/^http/, "ws") || void 0; | ||
| } | ||
| } | ||
| } | ||
| return value; | ||
| } | ||
| }); | ||
| } | ||
| export { Message as M, Peer as P, toString as a, toBufferLike as t }; |
Major refactor
Supply chain riskPackage has recently undergone a major refactor. It may be unstable or indicate significant internal changes. Use caution when updating to versions that include significant changes.
Network access
Supply chain riskThis module accesses the network.
Found 3 instances
Debug access
Supply chain riskUses debug, reflection and dynamic code execution features.
Major refactor
Supply chain riskPackage has recently undergone a major refactor. It may be unstable or indicate significant internal changes. Use caution when updating to versions that include significant changes.
Network access
Supply chain riskThis module accesses the network.
Found 7 instances
59
13.46%10
-72.97%213919
-11.51%29
3.57%4953
-17.97%5
25%