From 8727a4fa0aa7f92787caea8cbfeb3aed18a07d52 Mon Sep 17 00:00:00 2001 From: Adrian Hope-Bailie Date: Thu, 20 Jun 2019 23:05:22 +0200 Subject: [PATCH] Add rsocket types (#36076) * Add rsocket types * Drop `types` from `tsconfig.json Use triple-slash references * Fix linting errors * Remove private members * Fix linting errors --- types/rsocket-core/RSocketBinaryFraming.d.ts | 231 +++++++++++++++ types/rsocket-core/RSocketBufferUtils.d.ts | 47 +++ types/rsocket-core/RSocketClient.d.ts | 43 +++ types/rsocket-core/RSocketEncoding.d.ts | 40 +++ types/rsocket-core/RSocketFrame.d.ts | 113 +++++++ types/rsocket-core/RSocketMachine.d.ts | 69 +++++ .../RSocketResumableTransport.d.ts | 82 +++++ types/rsocket-core/RSocketSerialization.d.ts | 23 ++ types/rsocket-core/RSocketServer.d.ts | 52 ++++ types/rsocket-core/RSocketVersion.d.ts | 2 + types/rsocket-core/index.d.ts | 69 +++++ types/rsocket-core/tsconfig.json | 33 +++ types/rsocket-core/tslint.json | 1 + types/rsocket-flowable/Flowable.d.ts | 15 + .../rsocket-flowable/FlowableMapOperator.d.ts | 13 + types/rsocket-flowable/FlowableProcessor.d.ts | 12 + .../FlowableRequestOperator.d.ts | 12 + .../FlowableTakeOperator.d.ts | 13 + types/rsocket-flowable/FlowableTimer.d.ts | 13 + types/rsocket-flowable/Single.d.ts | 59 ++++ types/rsocket-flowable/index.d.ts | 14 + types/rsocket-flowable/tsconfig.json | 30 ++ types/rsocket-flowable/tslint.json | 7 + .../rsocket-tcp-client/RSocketTcpClient.d.ts | 45 +++ types/rsocket-tcp-client/index.d.ts | 9 + types/rsocket-tcp-client/tsconfig.json | 24 ++ types/rsocket-tcp-client/tslint.json | 1 + .../rsocket-tcp-server/RSocketTCPServer.d.ts | 23 ++ types/rsocket-tcp-server/index.d.ts | 9 + types/rsocket-tcp-server/tsconfig.json | 24 ++ types/rsocket-tcp-server/tslint.json | 1 + types/rsocket-types/ReactiveSocketTypes.d.ts | 279 ++++++++++++++++++ types/rsocket-types/ReactiveStreamTypes.d.ts | 42 +++ types/rsocket-types/index.d.ts | 8 + types/rsocket-types/tsconfig.json | 25 ++ types/rsocket-types/tslint.json | 6 + .../RSocketWebSocketClient.d.ts | 26 ++ types/rsocket-websocket-client/index.d.ts | 8 + types/rsocket-websocket-client/tsconfig.json | 24 ++ types/rsocket-websocket-client/tslint.json | 1 + .../RSocketWebSocketServer.d.ts | 14 + types/rsocket-websocket-server/index.d.ts | 8 + types/rsocket-websocket-server/tsconfig.json | 24 ++ types/rsocket-websocket-server/tslint.json | 1 + 44 files changed, 1595 insertions(+) create mode 100644 types/rsocket-core/RSocketBinaryFraming.d.ts create mode 100644 types/rsocket-core/RSocketBufferUtils.d.ts create mode 100644 types/rsocket-core/RSocketClient.d.ts create mode 100644 types/rsocket-core/RSocketEncoding.d.ts create mode 100644 types/rsocket-core/RSocketFrame.d.ts create mode 100644 types/rsocket-core/RSocketMachine.d.ts create mode 100644 types/rsocket-core/RSocketResumableTransport.d.ts create mode 100644 types/rsocket-core/RSocketSerialization.d.ts create mode 100644 types/rsocket-core/RSocketServer.d.ts create mode 100644 types/rsocket-core/RSocketVersion.d.ts create mode 100644 types/rsocket-core/index.d.ts create mode 100644 types/rsocket-core/tsconfig.json create mode 100644 types/rsocket-core/tslint.json create mode 100644 types/rsocket-flowable/Flowable.d.ts create mode 100644 types/rsocket-flowable/FlowableMapOperator.d.ts create mode 100644 types/rsocket-flowable/FlowableProcessor.d.ts create mode 100644 types/rsocket-flowable/FlowableRequestOperator.d.ts create mode 100644 types/rsocket-flowable/FlowableTakeOperator.d.ts create mode 100644 types/rsocket-flowable/FlowableTimer.d.ts create mode 100644 types/rsocket-flowable/Single.d.ts create mode 100644 types/rsocket-flowable/index.d.ts create mode 100644 types/rsocket-flowable/tsconfig.json create mode 100644 types/rsocket-flowable/tslint.json create mode 100644 types/rsocket-tcp-client/RSocketTcpClient.d.ts create mode 100644 types/rsocket-tcp-client/index.d.ts create mode 100644 types/rsocket-tcp-client/tsconfig.json create mode 100644 types/rsocket-tcp-client/tslint.json create mode 100644 types/rsocket-tcp-server/RSocketTCPServer.d.ts create mode 100644 types/rsocket-tcp-server/index.d.ts create mode 100644 types/rsocket-tcp-server/tsconfig.json create mode 100644 types/rsocket-tcp-server/tslint.json create mode 100644 types/rsocket-types/ReactiveSocketTypes.d.ts create mode 100644 types/rsocket-types/ReactiveStreamTypes.d.ts create mode 100644 types/rsocket-types/index.d.ts create mode 100644 types/rsocket-types/tsconfig.json create mode 100644 types/rsocket-types/tslint.json create mode 100644 types/rsocket-websocket-client/RSocketWebSocketClient.d.ts create mode 100644 types/rsocket-websocket-client/index.d.ts create mode 100644 types/rsocket-websocket-client/tsconfig.json create mode 100644 types/rsocket-websocket-client/tslint.json create mode 100644 types/rsocket-websocket-server/RSocketWebSocketServer.d.ts create mode 100644 types/rsocket-websocket-server/index.d.ts create mode 100644 types/rsocket-websocket-server/tsconfig.json create mode 100644 types/rsocket-websocket-server/tslint.json diff --git a/types/rsocket-core/RSocketBinaryFraming.d.ts b/types/rsocket-core/RSocketBinaryFraming.d.ts new file mode 100644 index 0000000000..9a371e8f6a --- /dev/null +++ b/types/rsocket-core/RSocketBinaryFraming.d.ts @@ -0,0 +1,231 @@ +/// + + import { + CancelFrame, + ErrorFrame, + Frame, + KeepAliveFrame, + LeaseFrame, + PayloadFrame, + RequestChannelFrame, + RequestFnfFrame, + RequestNFrame, + RequestResponseFrame, + RequestStreamFrame, + ResumeFrame, + ResumeOkFrame, + SetupFrame, +} from 'rsocket-types'; + +import { Encoders } from './RSocketEncoding'; + +export interface FrameWithPayload { + data: any; + flags: number; + metadata: any; +} + +/** + * Frame header is: + * - stream id (uint32 = 4) + * - type + flags (uint 16 = 2) + */ +export const FRAME_HEADER_SIZE = 6; + +/** + * Size of frame length and metadata length fields. + */ +export const UINT24_SIZE = 3; + +/** + * Reads a frame from a buffer that is prefixed with the frame length. + */ +export function deserializeFrameWithLength(buffer: Buffer, encoders?: Encoders): Frame; + +/** + * Given a buffer that may contain zero or more length-prefixed frames followed + * by zero or more bytes of a (partial) subsequent frame, returns an array of + * the frames and a buffer of the leftover bytes. + */ +export function deserializeFrames(buffer: Buffer, encoders?: Encoders): [Frame[], Buffer]; + +/** + * Writes a frame to a buffer with a length prefix. + */ +export function serializeFrameWithLength(frame: Frame, encoders?: Encoders): Buffer; + +/** + * Read a frame from the buffer. + */ +export function deserializeFrame(buffer: Buffer, encoders?: Encoders): Frame; + +/** + * Convert the frame to a (binary) buffer. + */ +export function serializeFrame(frame: Frame, encoders?: Encoders): Buffer; + +/** + * Writes a SETUP frame into a new buffer and returns it. + * + * Prefix size is: + * - version (2x uint16 = 4) + * - keepalive (uint32 = 4) + * - lifetime (uint32 = 4) + * - mime lengths (2x uint8 = 2) + */ +export const SETUP_FIXED_SIZE = 14; +export const RESUME_TOKEN_LENGTH_SIZE = 2; +export function serializeSetupFrame(frame: SetupFrame, encoders: Encoders): Buffer; + +/** + * Reads a SETUP frame from the buffer and returns it. + */ +export function deserializeSetupFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): SetupFrame; + +/** + * Writes an ERROR frame into a new buffer and returns it. + * + * Prefix size is for the error code (uint32 = 4). + */ +export const ERROR_FIXED_SIZE = 4; +export function serializeErrorFrame(frame: ErrorFrame, encoders: Encoders): Buffer; + +/** + * Reads an ERROR frame from the buffer and returns it. + */ +export function deserializeErrorFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): ErrorFrame; + +/** + * Writes a KEEPALIVE frame into a new buffer and returns it. + * + * Prefix size is for the last received position (uint64 = 8). + */ +export const KEEPALIVE_FIXED_SIZE = 8; +export function serializeKeepAliveFrame(frame: KeepAliveFrame, encoders: Encoders): Buffer; + +/** + * Reads a KEEPALIVE frame from the buffer and returns it. + */ +export function deserializeKeepAliveFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): KeepAliveFrame; + +/** + * Writes a LEASE frame into a new buffer and returns it. + * + * Prefix size is for the ttl (uint32) and requestcount (uint32). + */ +export const LEASE_FIXED_SIZE = 8; +export function serializeLeaseFrame(frame: LeaseFrame, encoders: Encoders): Buffer; + +/** + * Reads a LEASE frame from the buffer and returns it. + */ +export function deserializeLeaseFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): LeaseFrame; +/** + * Writes a REQUEST_FNF or REQUEST_RESPONSE frame to a new buffer and returns + * it. + * + * Note that these frames have the same shape and only differ in their type. + */ +export function serializeRequestFrame(frame: RequestFnfFrame | RequestResponseFrame, encoders: Encoders): Buffer; + +export function deserializeRequestFnfFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): RequestFnfFrame; + +export function deserializeRequestResponseFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): RequestResponseFrame; + +/** + * Writes a REQUEST_STREAM or REQUEST_CHANNEL frame to a new buffer and returns + * it. + * + * Note that these frames have the same shape and only differ in their type. + * + * Prefix size is for requestN (uint32 = 4). + */ +export const REQUEST_MANY_HEADER = 4; +export function serializeRequestManyFrame(frame: RequestStreamFrame | RequestChannelFrame, encoders: Encoders): Buffer; + +export function deserializeRequestStreamFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): RequestStreamFrame; + +export function deserializeRequestChannelFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): RequestChannelFrame; + +/** + * Writes a REQUEST_N frame to a new buffer and returns it. + * + * Prefix size is for requestN (uint32 = 4). + */ +export const REQUEST_N_HEADER = 4; +export function serializeRequestNFrame(frame: RequestNFrame, encoders: Encoders): Buffer; + +export function deserializeRequestNFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): RequestNFrame; + +/** + * Writes a CANCEL frame to a new buffer and returns it. + */ +export function serializeCancelFrame(frame: CancelFrame, encoders: Encoders): Buffer; + +export function deserializeCancelFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): CancelFrame; + +/** + * Writes a PAYLOAD frame to a new buffer and returns it. + */ +export function serializePayloadFrame(frame: PayloadFrame, encoders: Encoders): Buffer; + +export function deserializePayloadFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): PayloadFrame; + +/** + * Writes a RESUME frame into a new buffer and returns it. + * + * Fixed size is: + * - major version (uint16 = 2) + * - minor version (uint16 = 2) + * - token length (uint16 = 2) + * - client position (uint64 = 8) + * - server position (uint64 = 8) + */ +export const RESUME_FIXED_SIZE = 22; +export function serializeResumeFrame(frame: ResumeFrame, encoders: Encoders): Buffer; + +export function deserializeResumeFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): ResumeFrame; + +/** + * Writes a RESUME_OK frame into a new buffer and returns it. + * + * Fixed size is: + * - client position (uint64 = 8) + */ +export const RESUME_OK_FIXED_SIZE = 8; +export function serializeResumeOkFrame(frame: ResumeOkFrame, encoders: Encoders): Buffer; + +export function deserializeResumeOkFrame(buffer: Buffer, streamId: number, flags: number, encoders: Encoders): ResumeOkFrame; + +/** + * Write the header of the frame into the buffer. + */ +export function writeHeader(frame: Frame, buffer: Buffer): number; + +/** + * Determine the length of the payload section of a frame. Only applies to + * frame types that MAY have both metadata and data. + */ +export function getPayloadLength(frame: FrameWithPayload, encoders: Encoders): number; + +/** + * Write the payload of a frame into the given buffer. Only applies to frame + * types that MAY have both metadata and data. + */ +export function writePayload( + frame: FrameWithPayload, + buffer: Buffer, + encoders: Encoders, + offset: number +): void; + +/** + * Read the payload from a buffer and write it into the frame. Only applies to + * frame types that MAY have both metadata and data. + */ +export function readPayload( + buffer: Buffer, + frame: FrameWithPayload, + encoders: Encoders, + offset: number +): void; diff --git a/types/rsocket-core/RSocketBufferUtils.d.ts b/types/rsocket-core/RSocketBufferUtils.d.ts new file mode 100644 index 0000000000..f809da74ce --- /dev/null +++ b/types/rsocket-core/RSocketBufferUtils.d.ts @@ -0,0 +1,47 @@ +/// + + export type Encoding = "ascii" | "base64" | "hex" | "utf8"; + +/** + * Mimimum value that would overflow bitwise operators (2^32). + */ +export const BITWISE_OVERFLOW = 0x100000000; + +/** + * Read a uint24 from a buffer starting at the given offset. + */ +export function readUInt24BE(buffer: Buffer, offset: number): number; + +/** + * Writes a uint24 to a buffer starting at the given offset, returning the + * offset of the next byte. + */ +export function writeUInt24BE(buffer: Buffer, value: number, offset: number): number; + +/** + * Read a uint64 (technically supports up to 53 bits per JS number + * representation). + */ +export function readUInt64BE(buffer: Buffer, offset: number): number; + +/** + * Write a uint64 (technically supports up to 53 bits per JS number + * representation). + */ +export function writeUInt64BE(buffer: Buffer, value: number, offset: number): number; + +/** + * Determine the number of bytes it would take to encode the given data with the + * given encoding. + */ +export function byteLength(data: any, encoding: Encoding): number; + +/** + * Attempts to construct a buffer from the input, throws if invalid. + */ +export function toBuffer(data: any): Buffer; + +/** + * Function to create a buffer of a given sized filled with zeros. + */ +export function createBuffer(...args: any[]): Buffer; diff --git a/types/rsocket-core/RSocketClient.d.ts b/types/rsocket-core/RSocketClient.d.ts new file mode 100644 index 0000000000..0fd2517b9f --- /dev/null +++ b/types/rsocket-core/RSocketClient.d.ts @@ -0,0 +1,43 @@ +/// + +import { ConnectionStatus, DuplexConnection, Payload, ReactiveSocket, SetupFrame, Responder } from 'rsocket-types'; +import { PayloadSerializers } from './RSocketSerialization'; +import { Flowable, Single } from 'rsocket-flowable'; +export interface ClientConfig { + serializers?: PayloadSerializers; + setup: { + dataMimeType: string; + keepAlive: number; + lifetime: number; + metadataMimeType: string; + }; + transport: DuplexConnection; + responder?: Responder; +} + +/** + * RSocketClient: A client in an RSocket connection that will communicates with + * the peer via the given transport client. Provides methods for establishing a + * connection and initiating the RSocket interactions: + * - fireAndForget() + * - requestResponse() + * - requestStream() + * - requestChannel() + * - metadataPush() + */ +export default class RSocketClient { + constructor(config: ClientConfig); + close(): undefined; + connect(): Single>; +} + +export class RSocketClientSocket implements ReactiveSocket { + constructor(config: ClientConfig, connection: DuplexConnection); + fireAndForget(payload: Payload): undefined; + requestResponse(payload: Payload): Single>; + requestStream(payload: Payload): Flowable>; + requestChannel(payloads: Flowable>): Flowable>; + metadataPush(payload: Payload): Single; + close(): undefined; + connectionStatus(): Flowable; +} diff --git a/types/rsocket-core/RSocketEncoding.d.ts b/types/rsocket-core/RSocketEncoding.d.ts new file mode 100644 index 0000000000..5381f8337f --- /dev/null +++ b/types/rsocket-core/RSocketEncoding.d.ts @@ -0,0 +1,40 @@ +/// + +import { Encodable } from 'rsocket-types'; +import { byteLength } from './RSocketBufferUtils'; + +/** + * Commonly used subset of the allowed Node Buffer Encoder types. + */ +export interface Encoder { + byteLength: (value: Encodable) => number; + encode: (value: Encodable, buffer: Buffer, start: number, end: number) => number; + decode: (buffer: Buffer, start: number, end: number) => T; +} + +/** + * The Encoders object specifies how values should be serialized/deserialized + * to/from binary. + */ +export interface Encoders { + data: Encoder; + dataMimeType: Encoder; + message: Encoder; + metadata: Encoder; + metadataMimeType: Encoder; + resumeToken: Encoder; +} + +export const UTF8Encoder: Encoder; + +export const BufferEncoder: Encoder; + +/** + * Encode all values as UTF8 strings. + */ +export const Utf8Encoders: Encoders; + +/** + * Encode all values as buffers. + */ +export const BufferEncoders: Encoders; diff --git a/types/rsocket-core/RSocketFrame.d.ts b/types/rsocket-core/RSocketFrame.d.ts new file mode 100644 index 0000000000..88331f1772 --- /dev/null +++ b/types/rsocket-core/RSocketFrame.d.ts @@ -0,0 +1,113 @@ +import { ErrorFrame, Frame } from 'rsocket-types'; +export const CONNECTION_STREAM_ID = 0; +export const FRAME_TYPES: { + CANCEL: number; + ERROR: number; + EXT: number; + KEEPALIVE: number; + LEASE: number; + METADATA_PUSH: number; + PAYLOAD: number; + REQUEST_CHANNEL: number; + REQUEST_FNF: number; + REQUEST_N: number; + REQUEST_RESPONSE: number; + REQUEST_STREAM: number; + RESERVED: number; + RESUME: number; + RESUME_OK: number; + SETUP: number; +}; +export const FRAME_TYPE_NAMES: {}; +export const FLAGS: { + COMPLETE: number; + FOLLOWS: number; + IGNORE: number; + LEASE: number; + METADATA: number; + NEXT: number; + RESPOND: number; + RESUME_ENABLE: number; +}; +export const ERROR_CODES: { + APPLICATION_ERROR: number; + CANCELED: number; + CONNECTION_CLOSE: number; + CONNECTION_ERROR: number; + INVALID: number; + INVALID_SETUP: number; + REJECTED: number; + REJECTED_RESUME: number; + REJECTED_SETUP: number; + RESERVED: number; + RESERVED_EXTENSION: number; + UNSUPPORTED_SETUP: number; +}; +export const ERROR_EXPLANATIONS: {}; +export const FLAGS_MASK = 1023; +export const FRAME_TYPE_OFFFSET = 10; +export const MAX_CODE = 2147483647; +export const MAX_KEEPALIVE = 2147483647; +export const MAX_LIFETIME = 2147483647; +export const MAX_METADATA_LENGTH = 16777215; +export const MAX_MIME_LENGTH = 255; +export const MAX_REQUEST_COUNT = 2147483647; +export const MAX_REQUEST_N = 2147483647; +export const MAX_RESUME_LENGTH = 65535; +export const MAX_STREAM_ID = 2147483647; +export const MAX_TTL = 2147483647; +export const MAX_VERSION = 65535; +/** + * Returns true iff the flags have the IGNORE bit set. + */ +export function isIgnore(flags: number): boolean; +/** + * Returns true iff the flags have the METADATA bit set. + */ +export function isMetadata(flags: number): boolean; +/** + * Returns true iff the flags have the COMPLETE bit set. + */ +export function isComplete(flags: number): boolean; +/** + * Returns true iff the flags have the NEXT bit set. + */ +export function isNext(flags: number): boolean; +/** + * Returns true iff the flags have the RESPOND bit set. + */ +export function isRespond(flags: number): boolean; +/** + * Returns true iff the flags have the RESUME_ENABLE bit set. + */ +export function isResumeEnable(flags: number): boolean; +/** + * Returns true iff the flags have the LEASE bit set. + */ +export function isLease(flags: number): boolean; +/** + * Returns true iff the frame type is counted toward the implied + * client/server position used for the resumption protocol. + */ +export function isResumePositionFrameType(type: number): boolean; +export function getFrameTypeName(type: number): string; +/** + * Constructs an Error object given the contents of an error frame. The + * `source` property contains metadata about the error for use in introspecting + * the error at runtime: + * - `error.source.code: number`: the error code returned by the server. + * - `error.source.explanation: string`: human-readable explanation of the code + * (this value is not standardized and may change). + * - `error.source.message: string`: the error string returned by the server. + */ +export function createErrorFromFrame(frame: ErrorFrame): Error; +/** + * Given a RSocket error code, returns a human-readable explanation of that + * code, following the names used in the protocol specification. + */ +export function getErrorCodeExplanation(code: number): string; +/** + * Pretty-prints the frame for debugging purposes, with types, flags, and + * error codes annotated with descriptive names. + */ +export function printFrame(frame: Frame): string; diff --git a/types/rsocket-core/RSocketMachine.d.ts b/types/rsocket-core/RSocketMachine.d.ts new file mode 100644 index 0000000000..4e45e09cad --- /dev/null +++ b/types/rsocket-core/RSocketMachine.d.ts @@ -0,0 +1,69 @@ +import { + CancelFrame, + ConnectionStatus, + DuplexConnection, + Frame, + FrameWithData, + Payload, + Responder, + ReactiveSocket, + RequestFnfFrame, + RequestNFrame, + RequestResponseFrame, + RequestStreamFrame, + RequestChannelFrame, + ISubject, + ISubscription, + ISubscriber +} from 'rsocket-types'; + +import { Flowable, FlowableProcessor, Single } from 'rsocket-flowable'; +import { + createErrorFromFrame, + getFrameTypeName, + isComplete, + isNext, + isRespond, + CONNECTION_STREAM_ID, + ERROR_CODES, + FLAGS, + FRAME_TYPES, + MAX_REQUEST_N, + MAX_STREAM_ID, +} from './RSocketFrame'; +import { PayloadSerializers, IdentitySerializers } from './RSocketSerialization'; + +export type Role = 'CLIENT' | 'SERVER'; + +export class ResponderWrapper implements Responder { + constructor(responder: Partial>) + setResponder(responder: Partial>): void; + fireAndForget(payload: Payload): void; + requestResponse(payload: Payload): Single>; + requestStream(payload: Payload): Flowable>; + requestChannel(payloads: Flowable>): Flowable>; + metadataPush(payload: Payload): Single; +} + +export interface RSocketMachine extends ReactiveSocket { + setRequestHandler(requestHandler?: Partial>): void; +} + +export function createServerMachine( + connection: DuplexConnection, + connectionPublisher: (partialSubscriber: Partial>) => void, + serializers?: PayloadSerializers, + requestHandler?: Partial>, +): RSocketMachine; + +export function createClientMachine( + connection: DuplexConnection, + connectionPublisher: (partialSubscriber: Partial>) => void, + serializers?: PayloadSerializers, + requestHandler?: Partial>, +): RSocketMachine; + +export function deserializePayload( + serializers: PayloadSerializers, + frame: FrameWithData, +): Payload; diff --git a/types/rsocket-core/RSocketResumableTransport.d.ts b/types/rsocket-core/RSocketResumableTransport.d.ts new file mode 100644 index 0000000000..005f9d38ef --- /dev/null +++ b/types/rsocket-core/RSocketResumableTransport.d.ts @@ -0,0 +1,82 @@ +import { ConnectionStatus, DuplexConnection, Frame, SetupFrame, ISubject, ISubscription, CONNECTION_STATUS } from 'rsocket-types'; + +import { Flowable } from 'rsocket-flowable'; +import { + createErrorFromFrame, + isResumePositionFrameType, + CONNECTION_STREAM_ID, + FLAGS, + FRAME_TYPES, +} from './RSocketFrame'; + +export interface Options { + bufferSize: number; + resumeToken: string; +} + +/** + * NOTE: This implementation conforms to an upcoming version of the RSocket protocol + * and will not work with version 1.0 servers. + * + * An implementation of the DuplexConnection interface that supports automatic + * resumption per the RSocket protocol. + * + * # Example + * + * Create a client instance: + * ``` + * const client = new RSocketClient({ + * ..., + * transport: new RSocketResumableTransport( + * () => new RSocketWebSocketClient(...), // provider for low-level transport instances + * { + * bufferSize: 10, // max number of sent & pending frames to buffer before failing + * resumeToken: 'abc123', // string to uniquely identify the session across connections + * } + * ), + * }) + * + * Open the connection. After this if the connection dies it will be auto-resumed: + * ``` + * client.connect().subscribe(...); + * ``` + * + * Optionally, subscribe to the status of the connection: + * ``` + * client.connectionStatus().subscribe(...); + * ``` + * + * # Implementation Notes + * + * This transport maintains: + * - _currentConnection: a current low-level transport, which is null when not + * connected + * - _sentFrames: a buffer of frames written to a low-level transport (which + * may or may not have been received by the server) + * - _pendingFrames: a buffer of frames not yet written to the low-level + * connection, because they were sent while not connected. + * + * The initial connection is simple: connect using the low-level transport and + * flush any _pendingFrames (write them and add them to _sentFrames). + * + * Thereafter if the low-level transport drops, this transport attempts resumption. + * It obtains a fresh low-level transport from the given transport `source` + * and attempts to connect. Once connected, it sends a RESUME frame and waits. + * If RESUME_OK is received, _sentFrames and _pendingFrames are adjusted such + * that: + * - any frames the server has received are removed from _sentFrames + * - the remaining frames are merged (in correct order) into _pendingFrames + * + * Then the connection proceeds as above, where all pending frames are flushed. + * If anything other than RESUME_OK is received, resumption is considered to + * have failed and the connection is set to the ERROR status. + */ +export default class RSocketResumableTransport implements DuplexConnection { + constructor(source: () => DuplexConnection, options: Options) + close(): void; + connect(): void; + connectionStatus(): Flowable; + receive(): Flowable; + sendOne(frame: Frame): void; + send(frames: Flowable): void; +} diff --git a/types/rsocket-core/RSocketSerialization.d.ts b/types/rsocket-core/RSocketSerialization.d.ts new file mode 100644 index 0000000000..08cedbc754 --- /dev/null +++ b/types/rsocket-core/RSocketSerialization.d.ts @@ -0,0 +1,23 @@ +import { Encodable } from 'rsocket-types'; +/** + * A Serializer transforms data between the application encoding used in + * Payloads and the Encodable type accepted by the transport client. + */ +export interface Serializer { + deserialize: (data?: Encodable) => T | undefined; + serialize: (data?: T) => Encodable | undefined; +} +export interface PayloadSerializers { + data: Serializer; + metadata: Serializer; +} +export const JsonSerializer: Serializer; +export const JsonSerializers: { + data: Serializer; + metadata: Serializer; +}; +export const IdentitySerializer: Serializer; +export const IdentitySerializers: { + data: Serializer; + metadata: Serializer; +}; diff --git a/types/rsocket-core/RSocketServer.d.ts b/types/rsocket-core/RSocketServer.d.ts new file mode 100644 index 0000000000..88c53a8d86 --- /dev/null +++ b/types/rsocket-core/RSocketServer.d.ts @@ -0,0 +1,52 @@ +import { + DuplexConnection, + Frame, + FrameWithData, + Payload, + Responder, + ReactiveSocket, + ISubscription, + ISubscriber +} from 'rsocket-types'; +import { IdentitySerializers, PayloadSerializers } from './RSocketSerialization'; + +import { Flowable } from 'rsocket-flowable'; +import { + getFrameTypeName, + CONNECTION_STREAM_ID, + ERROR_CODES, + FRAME_TYPES, +} from './RSocketFrame'; +import { createServerMachine } from './RSocketMachine'; + +export interface TransportServer { + start: () => Flowable; + stop: () => void; +} + +export interface ServerConfig { + getRequestHandler: (socket: ReactiveSocket, payload: Payload) => Partial>; + serializers?: PayloadSerializers; + transport: TransportServer; +} + +/** + * RSocketServer: A server in an RSocket connection that accepts connections + * from peers via the given transport server. + */ +export default class RSocketServer { + constructor(config: ServerConfig); + start(): void; + stop(): void; +} + +export class SubscriberSwapper implements ISubscriber { + constructor(target?: Partial>); + swap(next: Partial>): ISubscriber; + onComplete(): void; + onError(error: Error): void; + onNext(value: T): void; + onSubscribe(subscription: ISubscription): void; +} + +export function deserializePayload(serializers: PayloadSerializers, frame: FrameWithData): Payload; diff --git a/types/rsocket-core/RSocketVersion.d.ts b/types/rsocket-core/RSocketVersion.d.ts new file mode 100644 index 0000000000..01500d3fe0 --- /dev/null +++ b/types/rsocket-core/RSocketVersion.d.ts @@ -0,0 +1,2 @@ +export const MAJOR_VERSION = 1; +export const MINOR_VERSION = 0; diff --git a/types/rsocket-core/index.d.ts b/types/rsocket-core/index.d.ts new file mode 100644 index 0000000000..420538cdfd --- /dev/null +++ b/types/rsocket-core/index.d.ts @@ -0,0 +1,69 @@ +// Type definitions for rsocket-core 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +export { ClientConfig } from "./RSocketClient"; +export { ServerConfig, TransportServer } from "./RSocketServer"; +export { Encodable } from "rsocket-types"; +export { Encoder, Encoders } from "./RSocketEncoding"; +export { Serializer, PayloadSerializers } from "./RSocketSerialization"; +import RSocketClient from "./RSocketClient"; +export { RSocketClient }; +import RSocketServer from "./RSocketServer"; +export { RSocketServer }; +import RSocketResumableTransport from "./RSocketResumableTransport"; +export { RSocketResumableTransport }; +export { + CONNECTION_STREAM_ID, + ERROR_CODES, + ERROR_EXPLANATIONS, + FLAGS_MASK, + FLAGS, + FRAME_TYPE_OFFFSET, + FRAME_TYPES, + MAX_CODE, + MAX_KEEPALIVE, + MAX_LIFETIME, + MAX_MIME_LENGTH, + MAX_RESUME_LENGTH, + MAX_STREAM_ID, + MAX_VERSION, + createErrorFromFrame, + getErrorCodeExplanation, + isComplete, + isIgnore, + isLease, + isMetadata, + isNext, + isRespond, + isResumeEnable, + printFrame +} from "./RSocketFrame"; +export { + deserializeFrame, + deserializeFrameWithLength, + deserializeFrames, + serializeFrame, + serializeFrameWithLength +} from "./RSocketBinaryFraming"; +export { + byteLength, + createBuffer, + readUInt24BE, + toBuffer, + writeUInt24BE +} from "./RSocketBufferUtils"; +export { + BufferEncoders, + BufferEncoder, + Utf8Encoders, + UTF8Encoder +} from "./RSocketEncoding"; +export { + IdentitySerializer, + IdentitySerializers, + JsonSerializer, + JsonSerializers +} from "./RSocketSerialization"; diff --git a/types/rsocket-core/tsconfig.json b/types/rsocket-core/tsconfig.json new file mode 100644 index 0000000000..5ab85bd408 --- /dev/null +++ b/types/rsocket-core/tsconfig.json @@ -0,0 +1,33 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "RSocketBinaryFraming.d.ts", + "RSocketBufferUtils.d.ts", + "RSocketClient.d.ts", + "RSocketEncoding.d.ts", + "RSocketFrame.d.ts", + "RSocketMachine.d.ts", + "RSocketResumableTransport.d.ts", + "RSocketSerialization.d.ts", + "RSocketServer.d.ts", + "RSocketVersion.d.ts" + ] +} diff --git a/types/rsocket-core/tslint.json b/types/rsocket-core/tslint.json new file mode 100644 index 0000000000..3db14f85ea --- /dev/null +++ b/types/rsocket-core/tslint.json @@ -0,0 +1 @@ +{ "extends": "dtslint/dt.json" } diff --git a/types/rsocket-flowable/Flowable.d.ts b/types/rsocket-flowable/Flowable.d.ts new file mode 100644 index 0000000000..07a7af17fd --- /dev/null +++ b/types/rsocket-flowable/Flowable.d.ts @@ -0,0 +1,15 @@ +import { IPublisher, ISubscriber } from 'rsocket-types'; +export type Source = (subscriber: ISubscriber) => void; +/** + * Implements the ReactiveStream `Publisher` interface with Rx-style operators. + */ +export default class Flowable implements IPublisher { + static just(...values: U[]): Flowable; + static error(error: Error): Flowable<{}>; + static never(): Flowable<{}>; + constructor(source: Source, max?: number); + subscribe(subscriberOrCallback?: Partial> | ((a: T) => void)): void; + lift(onSubscribeLift: (subscriber: ISubscriber) => ISubscriber): Flowable; + map(fn: (data: T) => R): Flowable; + take(toTake: number): Flowable; +} diff --git a/types/rsocket-flowable/FlowableMapOperator.d.ts b/types/rsocket-flowable/FlowableMapOperator.d.ts new file mode 100644 index 0000000000..9e2213e8d5 --- /dev/null +++ b/types/rsocket-flowable/FlowableMapOperator.d.ts @@ -0,0 +1,13 @@ +import { ISubscriber, ISubscription } from 'rsocket-types'; +/** + * An operator that acts like Array.map, applying a given function to + * all values provided by its `Subscription` and passing the result to its + * `Subscriber`. + */ +export default class FlowableMapOperator implements ISubscriber { + constructor(subscriber: ISubscriber, fn: (t: T) => R); + onComplete(): void; + onError(error: Error): void; + onNext(t: T): void; + onSubscribe(subscription: ISubscription): void; +} diff --git a/types/rsocket-flowable/FlowableProcessor.d.ts b/types/rsocket-flowable/FlowableProcessor.d.ts new file mode 100644 index 0000000000..dae604305c --- /dev/null +++ b/types/rsocket-flowable/FlowableProcessor.d.ts @@ -0,0 +1,12 @@ +import { IPublisher, ISubscription, ISubscriber } from 'rsocket-types'; +export default class FlowableProcessor implements IPublisher, ISubscriber, ISubscription { + constructor(source: IPublisher, fn?: (a: T) => R); + onSubscribe(subscription: ISubscription): void; + onNext(t: T): void; + onError(error: Error): void; + onComplete(): void; + subscribe(subscriber?: Partial>): void; + map(fn: (a: R) => S): IPublisher; + request(n: number): void; + cancel(): void; +} diff --git a/types/rsocket-flowable/FlowableRequestOperator.d.ts b/types/rsocket-flowable/FlowableRequestOperator.d.ts new file mode 100644 index 0000000000..9c41e72414 --- /dev/null +++ b/types/rsocket-flowable/FlowableRequestOperator.d.ts @@ -0,0 +1,12 @@ +import { ISubscriber, ISubscription } from 'rsocket-types'; +/** + * An operator that `request()`s the given number of items immediately upon + * being subscribed. + */ +export default class FlowableRequestOperator implements ISubscriber { + constructor(subscriber: ISubscriber, toRequest: number); + onComplete(): void; + onError(error: Error): void; + onNext(t: T): void; + onSubscribe(subscription: ISubscription): void; +} diff --git a/types/rsocket-flowable/FlowableTakeOperator.d.ts b/types/rsocket-flowable/FlowableTakeOperator.d.ts new file mode 100644 index 0000000000..86771f855f --- /dev/null +++ b/types/rsocket-flowable/FlowableTakeOperator.d.ts @@ -0,0 +1,13 @@ +import { ISubscriber, ISubscription } from 'rsocket-types'; +/** + * An operator that requests a fixed number of values from its source + * `Subscription` and forwards them to its `Subscriber`, cancelling the + * subscription when the requested number of items has been reached. + */ +export default class FlowableTakeOperator implements ISubscriber { + constructor(subscriber: ISubscriber, toTake: number); + onComplete(): void; + onError(error: Error): void; + onNext(t: T): void; + onSubscribe(subscription: ISubscription): void; +} diff --git a/types/rsocket-flowable/FlowableTimer.d.ts b/types/rsocket-flowable/FlowableTimer.d.ts new file mode 100644 index 0000000000..c9a41f962d --- /dev/null +++ b/types/rsocket-flowable/FlowableTimer.d.ts @@ -0,0 +1,13 @@ +import Flowable from './Flowable'; +/** + * Returns a Publisher that provides the current time (Date.now()) every `ms` + * milliseconds. + * + * The timer is established on the first call to `request`: on each + * interval a value is published if there are outstanding requests, + * otherwise nothing occurs for that interval. This approach ensures + * that the interval between `onNext` calls is as regular as possible + * and means that overlapping `request` calls (ie calling again before + * the previous values have been vended) behaves consistently. + */ +export function every(ms: number): Flowable; diff --git a/types/rsocket-flowable/Single.d.ts b/types/rsocket-flowable/Single.d.ts new file mode 100644 index 0000000000..f5679c5acc --- /dev/null +++ b/types/rsocket-flowable/Single.d.ts @@ -0,0 +1,59 @@ +export type Source = (subject: IFutureSubject) => undefined; +export type CancelCallback = () => undefined; + +export interface IFutureSubscriber { + onComplete: (value: T) => undefined; + onError: (error: Error) => undefined; + onSubscribe: (cancel: CancelCallback) => undefined; +} +export interface IFutureSubject { + onComplete: (value: T) => undefined; + onError: (error: Error) => undefined; + onSubscribe: (cancel: CancelCallback | null | undefined) => undefined; +} +/** + * Represents a lazy computation that will either produce a value of type T + * or fail with an error. Calling `subscribe()` starts the + * computation and returns a subscription object, which has an `unsubscribe()` + * method that can be called to prevent completion/error callbacks from being + * invoked and, where supported, to also cancel the computation. + * Implementations may optionally implement cancellation; if they do not + * `cancel()` is a no-op. + * + * Note: Unlike Promise, callbacks (onComplete/onError) may be invoked + * synchronously. + * + * Example: + * + * ``` + * const value = new Single(subscriber => { + * const id = setTimeout( + * () => subscriber.onComplete('Hello!'), + * 250 + * ); + * // Optional: Call `onSubscribe` with a cancellation callback + * subscriber.onSubscribe(() => clearTimeout(id)); + * }); + * + * // Start the computation. onComplete will be called after the timeout + * // with 'hello' unless `cancel()` is called first. + * value.subscribe({ + * onComplete: value => console.log(value), + * onError: error => console.error(error), + * onSubscribe: cancel => ... + * }); + * ``` + */ +export default class Single { + static of(value: U): Single; + static error(error: Error): Single<{}>; + constructor(source: Source); + subscribe(partialSubscriber?: Partial>): void; + flatMap(fn: (data: T) => Single): Single; + /** + * Return a new Single that resolves to the value of this Single applied to + * the given mapping function. + */ + map(fn: (data: T) => R): Single; + then(successFn?: (data: T) => void, errorFn?: (error: Error) => void): void; +} diff --git a/types/rsocket-flowable/index.d.ts b/types/rsocket-flowable/index.d.ts new file mode 100644 index 0000000000..0e200d6686 --- /dev/null +++ b/types/rsocket-flowable/index.d.ts @@ -0,0 +1,14 @@ +// Type definitions for rsocket-flowable 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +import Flowable from './Flowable'; +import Single from './Single'; +import FlowableProcessor from './FlowableProcessor'; +import { every } from './FlowableTimer'; +/** + * The public API of the `flowable` package. + */ +export { Flowable, FlowableProcessor, Single, every }; diff --git a/types/rsocket-flowable/tsconfig.json b/types/rsocket-flowable/tsconfig.json new file mode 100644 index 0000000000..738fd880b5 --- /dev/null +++ b/types/rsocket-flowable/tsconfig.json @@ -0,0 +1,30 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "Flowable.d.ts", + "FlowableMapOperator.d.ts", + "FlowableProcessor.d.ts", + "FlowableRequestOperator.d.ts", + "FlowableTakeOperator.d.ts", + "FlowableTimer.d.ts", + "Single.d.ts" + ] +} diff --git a/types/rsocket-flowable/tslint.json b/types/rsocket-flowable/tslint.json new file mode 100644 index 0000000000..e7ef30a8db --- /dev/null +++ b/types/rsocket-flowable/tslint.json @@ -0,0 +1,7 @@ +{ + "extends": "dtslint/dt.json", + "rules": { + "interface-name": false + } + +} diff --git a/types/rsocket-tcp-client/RSocketTcpClient.d.ts b/types/rsocket-tcp-client/RSocketTcpClient.d.ts new file mode 100644 index 0000000000..5b65c58145 --- /dev/null +++ b/types/rsocket-tcp-client/RSocketTcpClient.d.ts @@ -0,0 +1,45 @@ +/// + +import { ConnectionStatus, DuplexConnection, Frame, ISubject, ISubscriber, ISubscription, CONNECTION_STATUS } from 'rsocket-types'; + +import * as net from 'net'; +import * as tls from 'tls'; +import { Flowable } from 'rsocket-flowable'; +import { + Encoders, + createBuffer, + deserializeFrames, + serializeFrameWithLength, +} from 'rsocket-core'; + +/** + * A TCP transport client for use in node environments. + */ +export class RSocketTcpConnection implements DuplexConnection { + constructor(socket?: net.Socket, encoders?: Encoders); + close(): void; + connect(): void; + setupSocket(socket: net.Socket): void; + connectionStatus(): Flowable; + receive(): Flowable; + sendOne(frame: Frame): void; + send(frames: Flowable): void; + getConnectionState(): ConnectionStatus; + setConnectionStatus(status: ConnectionStatus): void; +} + +/** + * A TCP transport client for use in node environments. + */ +export class RSocketTcpClient extends RSocketTcpConnection { + constructor(options: net.TcpSocketConnectOpts, encoders?: Encoders); + connect(): void; +} + +/** + * A TLS transport client for use in node environments. + */ +export class RSocketTlsClient extends RSocketTcpConnection { + constructor(options: tls.ConnectionOptions, encoders?: Encoders); + connect(): void; +} diff --git a/types/rsocket-tcp-client/index.d.ts b/types/rsocket-tcp-client/index.d.ts new file mode 100644 index 0000000000..b150ce6398 --- /dev/null +++ b/types/rsocket-tcp-client/index.d.ts @@ -0,0 +1,9 @@ +// Type definitions for rsocket-tcp-client 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +import { RSocketTcpClient, RSocketTcpConnection } from './RSocketTcpClient'; +export default RSocketTcpClient; +export { RSocketTcpConnection }; diff --git a/types/rsocket-tcp-client/tsconfig.json b/types/rsocket-tcp-client/tsconfig.json new file mode 100644 index 0000000000..3085677ce1 --- /dev/null +++ b/types/rsocket-tcp-client/tsconfig.json @@ -0,0 +1,24 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "RSocketTcpClient.d.ts" + ] +} diff --git a/types/rsocket-tcp-client/tslint.json b/types/rsocket-tcp-client/tslint.json new file mode 100644 index 0000000000..3db14f85ea --- /dev/null +++ b/types/rsocket-tcp-client/tslint.json @@ -0,0 +1 @@ +{ "extends": "dtslint/dt.json" } diff --git a/types/rsocket-tcp-server/RSocketTCPServer.d.ts b/types/rsocket-tcp-server/RSocketTCPServer.d.ts new file mode 100644 index 0000000000..ba09573376 --- /dev/null +++ b/types/rsocket-tcp-server/RSocketTCPServer.d.ts @@ -0,0 +1,23 @@ +/// + +import { DuplexConnection } from 'rsocket-types'; +import { Encoders, TransportServer } from 'rsocket-core'; +import * as EventEmitter from 'events'; +import * as net from 'net'; +import { Flowable } from 'rsocket-flowable'; +export interface ServerOptions { + host?: string; + port: number; + serverFactory?: (onConnect: (socket: net.Socket) => undefined) => net.Server; +} + +/** + * A TCP transport server. + * + * //FIXME: Inconsistent casing between TCPServer and TcpClient matches library + */ +export default class RSocketTCPServer implements TransportServer { + constructor(options: ServerOptions, encoders?: Encoders); + start(): Flowable; + stop(): undefined; +} diff --git a/types/rsocket-tcp-server/index.d.ts b/types/rsocket-tcp-server/index.d.ts new file mode 100644 index 0000000000..0fb5b43610 --- /dev/null +++ b/types/rsocket-tcp-server/index.d.ts @@ -0,0 +1,9 @@ +// Type definitions for rsocket-tcp-server 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +import RSocketTCPServer, { ServerOptions } from './RSocketTCPServer'; +export default RSocketTCPServer; +export { ServerOptions }; diff --git a/types/rsocket-tcp-server/tsconfig.json b/types/rsocket-tcp-server/tsconfig.json new file mode 100644 index 0000000000..e3331da499 --- /dev/null +++ b/types/rsocket-tcp-server/tsconfig.json @@ -0,0 +1,24 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "RSocketTCPServer.d.ts" + ] +} diff --git a/types/rsocket-tcp-server/tslint.json b/types/rsocket-tcp-server/tslint.json new file mode 100644 index 0000000000..3db14f85ea --- /dev/null +++ b/types/rsocket-tcp-server/tslint.json @@ -0,0 +1 @@ +{ "extends": "dtslint/dt.json" } diff --git a/types/rsocket-types/ReactiveSocketTypes.d.ts b/types/rsocket-types/ReactiveSocketTypes.d.ts new file mode 100644 index 0000000000..66c6ce1a8d --- /dev/null +++ b/types/rsocket-types/ReactiveSocketTypes.d.ts @@ -0,0 +1,279 @@ +/// + +import { Flowable, Single } from 'rsocket-flowable'; + +export interface Responder { + /** + * Fire and Forget interaction model of `ReactiveSocket`. The returned + * Publisher resolves when the passed `payload` is successfully handled. + */ + fireAndForget(payload: Payload): void; + + /** + * Request-Response interaction model of `ReactiveSocket`. The returned + * Publisher resolves with the response. + */ + requestResponse(payload: Payload): Single>; + + /** + * Request-Stream interaction model of `ReactiveSocket`. The returned + * Publisher returns values representing the response(s). + */ + requestStream(payload: Payload): Flowable>; + + /** + * Request-Channel interaction model of `ReactiveSocket`. The returned + * Publisher returns values representing the response(s). + */ + requestChannel(payloads: Flowable>): Flowable>; + + /** + * Metadata-Push interaction model of `ReactiveSocket`. The returned Publisher + * resolves when the passed `payload` is successfully handled. + */ + metadataPush(payload: Payload): Single; +} + +/** + * A contract providing different interaction models per the [ReactiveSocket protocol] + * (https://github.com/ReactiveSocket/reactivesocket/blob/master/Protocol.md). + */ +export interface ReactiveSocket extends Responder { + /** + * Close this `ReactiveSocket` and the underlying transport connection. + */ + close(): void; + + /** + * Returns a Flowable that immediately publishes the current connection + * status and thereafter updates as it changes. Once a connection is in + * the CLOSED or ERROR state, it may not be connected again. + * Implementations must publish values per the comments on ConnectionStatus. + */ + connectionStatus(): Flowable; +} + +/** + * Represents a network connection with input/output used by a ReactiveSocket to + * send/receive data. + */ +export interface DuplexConnection { + /** + * Send a single frame on the connection. + */ + sendOne(frame: Frame): void; + + /** + * Send all the `input` frames on this connection. + * + * Notes: + * - Implementations must not cancel the subscription. + * - Implementations must signal any errors by calling `onError` on the + * `receive()` Publisher. + */ + send(input: Flowable): void; + + /** + * Returns a stream of all `Frame`s received on this connection. + * + * Notes: + * - Implementations must call `onComplete` if the underlying connection is + * closed by the peer or by calling `close()`. + * - Implementations must call `onError` if there are any errors + * sending/receiving frames. + * - Implemenations may optionally support multi-cast receivers. Those that do + * not should throw if `receive` is called more than once. + */ + receive(): Flowable; + + /** + * Close the underlying connection, emitting `onComplete` on the receive() + * Publisher. + */ + close(): void; + + /** + * Open the underlying connection. Throws if the connection is already in + * the CLOSED or ERROR state. + */ + connect(): void; + + /** + * Returns a Flowable that immediately publishes the current connection + * status and thereafter updates as it changes. Once a connection is in + * the CLOSED or ERROR state, it may not be connected again. + * Implementations must publish values per the comments on ConnectionStatus. + */ + connectionStatus(): Flowable; +} + +/** + * Describes the connection status of a ReactiveSocket/DuplexConnection. + * - NOT_CONNECTED: no connection established or pending. + * - CONNECTING: when `connect()` has been called but a connection is not yet + * established. + * - CONNECTED: when a connection is established. + * - CLOSED: when the connection has been explicitly closed via `close()`. + * - ERROR: when the connection has been closed for any other reason. + */ +export type ConnectionStatus = + {kind: 'NOT_CONNECTED'} | + {kind: 'CONNECTING'} | + {kind: 'CONNECTED'} | + {kind: 'CLOSED'} | + {kind: 'ERROR', error: Error}; + +export const CONNECTION_STATUS: ConnectionStatus; + +/** + * A type that can be written to a buffer. + */ +export type Encodable = string | Buffer | Uint8Array; + +/** + * A single unit of data exchanged between the peers of a `ReactiveSocket`. + */ +export interface Payload { + data?: D; + metadata?: M; +} + +export type Frame = + CancelFrame | + ErrorFrame | + KeepAliveFrame | + LeaseFrame | + PayloadFrame | + RequestChannelFrame | + RequestFnfFrame | + RequestNFrame | + RequestResponseFrame | + RequestStreamFrame | + ResumeFrame | + ResumeOkFrame | + SetupFrame | + UnsupportedFrame; + +export interface FrameWithData { + data?: Encodable; + metadata?: Encodable; +} + +export interface CancelFrame { + type: 0x09; + flags: number; + streamId: number; +} + +export interface ErrorFrame { + type: 0x0B; + flags: number; + code: number; + message: string; + streamId: number; +} + +export interface KeepAliveFrame { + type: 0x03; + flags: number; + data?: Encodable; + lastReceivedPosition: number; + streamId: 0; +} + +export interface LeaseFrame { + type: 0x02; + flags: number; + ttl: number; + requestCount: number; + metadata?: Encodable; + streamId: 0; +} + +export interface PayloadFrame { + type: 0x0A; + flags: number; + data?: Encodable; + metadata?: Encodable; + streamId: number; +} + +export interface RequestChannelFrame { + type: 0x07; + data?: Encodable; + metadata?: Encodable; + flags: number; + requestN: number; + streamId: number; +} + +export interface RequestFnfFrame { + type: 0x05; + data?: Encodable; + metadata?: Encodable; + flags: number; + streamId: number; +} + +export interface RequestNFrame { + type: 0x08; + flags: number; + requestN: number; + streamId: number; +} + +export interface RequestResponseFrame { + type: 0x04; + data?: Encodable; + metadata?: Encodable; + flags: number; + streamId: number; +} + +export interface RequestStreamFrame { + type: 0x06; + data: Encodable; + metadata: Encodable; + flags: number; + requestN: number; + streamId: number; +} + +export interface ResumeFrame { + type: 0x0d; + clientPosition: number; + flags: number; + majorVersion: number; + minorVersion: number; + resumeToken: Encodable; + serverPosition: number; + streamId: 0; +} + +export interface ResumeOkFrame { + type: 0x0e; + clientPosition: number; + flags: number; + streamId: 0; +} + +export interface SetupFrame { + type: 0x01; + data?: Encodable; + dataMimeType: string; + flags: number; + keepAlive: number; + lifetime: number; + metadata?: Encodable; + metadataMimeType: string; + resumeToken?: Encodable; + streamId: 0; + majorVersion: number; + minorVersion: number; +} + +export interface UnsupportedFrame { + type: 0x3f | 0x0c | 0x00; + streamId: 0; + flags: number; +} diff --git a/types/rsocket-types/ReactiveStreamTypes.d.ts b/types/rsocket-types/ReactiveStreamTypes.d.ts new file mode 100644 index 0000000000..2eb5ff1cac --- /dev/null +++ b/types/rsocket-types/ReactiveStreamTypes.d.ts @@ -0,0 +1,42 @@ +/** + * Core types per the [ReactiveStreams specification](http://www.reactive-streams.org/) + */ +/** + * Represents an asynchronous, pull-based stream of values. Calling + * `subscribe()` causes the subscriber's `onSubscribe()` method to be be invoked + * with a Subscription object that has two methods: + * - `cancel()`: stops the publisher from publishing any more values. + * - `request(n)`: requests `n` additional values. + * + * The subscriber can use `request(n)` to pull additional values from the + * stream. + */ + +export interface IPublisher { + subscribe: (subscriber?: Partial>) => void; + map: (fn: (data: T) => R) => IPublisher; +} +/** + * An underlying source of values for a Publisher. + */ +export interface ISubscription { + cancel: () => void; + request: (n: number) => void; +} +/** + * A handler for values provided by a Publisher. + */ +export interface ISubscriber { + onComplete: () => void; + onError: (error: Error) => void; + onNext: (value: T) => void; + onSubscribe: (subscription: ISubscription) => void; +} +/** + * Similar to Subscriber, but without onSubscribe. + */ +export interface ISubject { + onComplete: () => void; + onError: (error: Error) => void; + onNext: (value: T) => void; +} diff --git a/types/rsocket-types/index.d.ts b/types/rsocket-types/index.d.ts new file mode 100644 index 0000000000..d718ff4c97 --- /dev/null +++ b/types/rsocket-types/index.d.ts @@ -0,0 +1,8 @@ +// Type definitions for rsocket-types 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +export * from './ReactiveSocketTypes'; +export * from './ReactiveStreamTypes'; diff --git a/types/rsocket-types/tsconfig.json b/types/rsocket-types/tsconfig.json new file mode 100644 index 0000000000..b4487b7404 --- /dev/null +++ b/types/rsocket-types/tsconfig.json @@ -0,0 +1,25 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "ReactiveSocketTypes.d.ts", + "ReactiveStreamTypes.d.ts" + ] +} diff --git a/types/rsocket-types/tslint.json b/types/rsocket-types/tslint.json new file mode 100644 index 0000000000..2c7c1bed53 --- /dev/null +++ b/types/rsocket-types/tslint.json @@ -0,0 +1,6 @@ +{ + "extends": "dtslint/dt.json", + "rules": { + "interface-name": false + } +} diff --git a/types/rsocket-websocket-client/RSocketWebSocketClient.d.ts b/types/rsocket-websocket-client/RSocketWebSocketClient.d.ts new file mode 100644 index 0000000000..65fd4db112 --- /dev/null +++ b/types/rsocket-websocket-client/RSocketWebSocketClient.d.ts @@ -0,0 +1,26 @@ +import { ConnectionStatus, DuplexConnection, Frame, ISubject, ISubscriber, ISubscription, CONNECTION_STATUS } from 'rsocket-types'; + +import { Flowable } from 'rsocket-flowable'; +import { + deserializeFrame, + deserializeFrameWithLength, + Encoders, + printFrame, + serializeFrame, + serializeFrameWithLength, + toBuffer, +} from 'rsocket-core'; +import * as ws from 'ws'; + +/** + * A WebSocket transport client for use in browser environments. + */ +export default class RSocketWebSocketClient implements DuplexConnection { + constructor(options: ws.ClientOptions, encoders?: Encoders) + close(): void; + connect(): void; + connectionStatus(): Flowable; + receive(): Flowable; + sendOne(frame: Frame): void; + send(frames: Flowable): void; +} diff --git a/types/rsocket-websocket-client/index.d.ts b/types/rsocket-websocket-client/index.d.ts new file mode 100644 index 0000000000..306fbef218 --- /dev/null +++ b/types/rsocket-websocket-client/index.d.ts @@ -0,0 +1,8 @@ +// Type definitions for rsocket-websocket-client 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +import RSocketWebSocketClient from './RSocketWebSocketClient'; +export default RSocketWebSocketClient; diff --git a/types/rsocket-websocket-client/tsconfig.json b/types/rsocket-websocket-client/tsconfig.json new file mode 100644 index 0000000000..023c8cfd18 --- /dev/null +++ b/types/rsocket-websocket-client/tsconfig.json @@ -0,0 +1,24 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "RSocketWebSocketClient.d.ts" + ] +} diff --git a/types/rsocket-websocket-client/tslint.json b/types/rsocket-websocket-client/tslint.json new file mode 100644 index 0000000000..3db14f85ea --- /dev/null +++ b/types/rsocket-websocket-client/tslint.json @@ -0,0 +1 @@ +{ "extends": "dtslint/dt.json" } diff --git a/types/rsocket-websocket-server/RSocketWebSocketServer.d.ts b/types/rsocket-websocket-server/RSocketWebSocketServer.d.ts new file mode 100644 index 0000000000..50b55b89a9 --- /dev/null +++ b/types/rsocket-websocket-server/RSocketWebSocketServer.d.ts @@ -0,0 +1,14 @@ +import { DuplexConnection } from 'rsocket-types'; +import { Encoders, TransportServer } from 'rsocket-core'; +import { EventEmitter } from 'events'; +import { Flowable } from 'rsocket-flowable'; +import * as ws from 'ws'; + +/** + * A WebSocket transport server. + */ +export class RSocketWebSocketServer implements TransportServer { + constructor(options: ws.ServerOptions, encoders?: Encoders); + start(): Flowable; + stop(): void; +} diff --git a/types/rsocket-websocket-server/index.d.ts b/types/rsocket-websocket-server/index.d.ts new file mode 100644 index 0000000000..fe62edb45e --- /dev/null +++ b/types/rsocket-websocket-server/index.d.ts @@ -0,0 +1,8 @@ +// Type definitions for rsocket-websocket-server 0.0 +// Project: https://github.com/rsocket/rsocket-js/ +// Definitions by: Adrian Hope-Bailie +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped +// TypeScript Version: 2.1 + +import { RSocketWebSocketServer } from './RSocketWebSocketServer'; +export default RSocketWebSocketServer; diff --git a/types/rsocket-websocket-server/tsconfig.json b/types/rsocket-websocket-server/tsconfig.json new file mode 100644 index 0000000000..762bbc6457 --- /dev/null +++ b/types/rsocket-websocket-server/tsconfig.json @@ -0,0 +1,24 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "dom", + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": true, + "strictFunctionTypes": true, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "RSocketWebSocketServer.d.ts" + ] +} diff --git a/types/rsocket-websocket-server/tslint.json b/types/rsocket-websocket-server/tslint.json new file mode 100644 index 0000000000..3db14f85ea --- /dev/null +++ b/types/rsocket-websocket-server/tslint.json @@ -0,0 +1 @@ +{ "extends": "dtslint/dt.json" }