| import * as Context from "../../Context.js"; | |
| import * as Data from "../../Data.js"; | |
| import * as Effect from "../../Effect.js"; | |
| import * as Exit from "../../Exit.js"; | |
| import * as Option from "../../Option.js"; | |
| import { hasProperty } from "../../Predicate.js"; | |
| import * as Schema from "../../Schema.js"; | |
| import * as SchemaIssue from "../../SchemaIssue.js"; | |
| import * as SchemaParser from "../../SchemaParser.js"; | |
| import * as SchemaTransformation from "../../SchemaTransformation.js"; | |
| import * as Rpc from "../rpc/Rpc.js"; | |
| import { MalformedMessage } from "./ClusterError.js"; | |
| import { Snowflake, SnowflakeFromBigInt } from "./Snowflake.js"; | |
| const TypeId = "~effect/cluster/Reply"; | |
| /** | |
| * Returns `true` when the supplied value is a runtime cluster reply, based on the | |
| * reply type identifier. | |
| * | |
| * @category guards | |
| * @since 4.0.0 | |
| */ | |
| export const isReply = u => hasProperty(u, TypeId); | |
| /** | |
| * Schema for reply values that are already in encoded form. | |
| * | |
| * **Details** | |
| * | |
| * Per-RPC payload validation is performed by `Reply(rpc)`. | |
| * | |
| * @category schemas | |
| * @since 4.0.0 | |
| */ | |
| export const Encoded = Schema.Any; | |
| /** | |
| * Represents a cluster reply paired with the RPC definition and service context required to | |
| * serialize it for transport. | |
| * | |
| * **When to use** | |
| * | |
| * Use to carry a runtime reply together with the RPC schema and services needed | |
| * to encode it for storage or transport. | |
| * | |
| * @category models | |
| * @since 4.0.0 | |
| */ | |
| export class ReplyWithContext extends /*#__PURE__*/Data.TaggedClass("ReplyWithContext") { | |
| /** | |
| * Creates a terminal reply context that dies with the supplied defect. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static fromDefect(options) { | |
| return new ReplyWithContext({ | |
| reply: new WithExit({ | |
| requestId: options.requestId, | |
| id: options.id, | |
| exit: Exit.die(Schema.encodeSync(Schema.Defect())(options.defect)) | |
| }), | |
| context: Context.empty(), | |
| rpc: neverRpc | |
| }); | |
| } | |
| /** | |
| * Creates a terminal reply context that interrupts the supplied request. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static interrupt(options) { | |
| return new ReplyWithContext({ | |
| reply: new WithExit({ | |
| requestId: options.requestId, | |
| id: options.id, | |
| exit: Exit.interrupt() | |
| }), | |
| context: Context.empty(), | |
| rpc: neverRpc | |
| }); | |
| } | |
| } | |
| const neverRpc = /*#__PURE__*/Rpc.make("Never", { | |
| success: Schema.Never, | |
| error: Schema.Never, | |
| payload: {} | |
| }); | |
| const schemaCache = /*#__PURE__*/new WeakMap(); | |
| /** | |
| * Represents a streaming RPC reply chunk for a request, carrying a non-empty | |
| * batch of success values together with the reply id and sequence number. | |
| * | |
| * @category models | |
| * @since 4.0.0 | |
| */ | |
| export class Chunk extends /*#__PURE__*/Data.TaggedClass("Chunk") { | |
| /** | |
| * Marks this value as a runtime cluster reply. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| [TypeId] = TypeId; | |
| /** | |
| * Creates an empty chunk reply for the supplied request id. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static emptyFrom(requestId) { | |
| return new Chunk({ | |
| requestId, | |
| id: Snowflake(BigInt(0)), | |
| sequence: 0, | |
| values: [undefined] | |
| }); | |
| } | |
| /** | |
| * Schema that accepts any runtime chunk reply without validating payload values. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static Any = /*#__PURE__*/Schema.declare(u => isReply(u) && u._tag === "Chunk"); | |
| /** | |
| * Transformation between encoded chunk records and `Chunk` instances. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static transform = /*#__PURE__*/SchemaTransformation.transform({ | |
| decode: a => new Chunk(a), | |
| encode: a => a | |
| }); | |
| /** | |
| * Builds a chunk schema from the streaming success schema of an RPC. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static schema(rpc) { | |
| const successSchema = rpc.successSchema.success; | |
| if (!successSchema) { | |
| return Schema.Never; | |
| } | |
| return this.schemaFrom(successSchema); | |
| } | |
| /** | |
| * Builds a chunk schema that validates each success value with the supplied schema. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static schemaFrom(success) { | |
| // TODO: extract to a helper function | |
| return Schema.declareConstructor()([success], ([success]) => (input, ast, options) => { | |
| if (!isReply(input) || input._tag !== "Chunk") { | |
| return Effect.fail(new SchemaIssue.InvalidType(ast, Option.some(input))); | |
| } | |
| return Effect.mapBothEager(SchemaParser.decodeEffect(Schema.NonEmptyArray(success))(input.values, options), { | |
| onFailure: issue => new SchemaIssue.Composite(ast, Option.some(input), [new SchemaIssue.Pointer(["values"], issue)]), | |
| onSuccess: values => new Chunk({ | |
| ...input, | |
| values | |
| }) | |
| }); | |
| }, { | |
| expected: "Reply.Chunk", | |
| toCodecJson: ([success]) => Schema.link()(Schema.Struct({ | |
| _tag: Schema.Literal("Chunk"), | |
| requestId: SnowflakeFromBigInt, | |
| id: SnowflakeFromBigInt, | |
| sequence: Schema.Number, | |
| values: Schema.NonEmptyArray(success) | |
| }), SchemaTransformation.transform({ | |
| decode: encoded => new Chunk(encoded), | |
| encode: result => ({ | |
| ...result | |
| }) | |
| })) | |
| }); | |
| } | |
| /** | |
| * Returns a copy of this chunk associated with the supplied request id. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| withRequestId(requestId) { | |
| return new Chunk({ | |
| ...this, | |
| requestId | |
| }); | |
| } | |
| } | |
| /** | |
| * Represents a terminal RPC reply for a request, carrying the final `Exit` for the remote | |
| * call. | |
| * | |
| * **When to use** | |
| * | |
| * Use to represent the final success, typed failure, defect, or interruption | |
| * for a clustered RPC request. | |
| * | |
| * @category models | |
| * @since 4.0.0 | |
| */ | |
| export class WithExit extends /*#__PURE__*/Data.TaggedClass("WithExit") { | |
| /** | |
| * Marks this value as a runtime cluster reply. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| [TypeId] = TypeId; | |
| /** | |
| * Returns `true` when the value is a terminal `WithExit` reply. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static is(u) { | |
| return isReply(u) && u._tag === "WithExit"; | |
| } | |
| /** | |
| * Builds a terminal reply schema from the exit schema of an RPC. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static schema(rpc) { | |
| return this.schemaFrom(Rpc.exitSchema(rpc)); | |
| } | |
| /** | |
| * Builds a terminal reply schema that validates the encoded exit value. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| static schemaFrom(exitSchema) { | |
| // TODO: extract to a helper function | |
| return Schema.declareConstructor()([exitSchema], ([exit]) => (input, ast, options) => { | |
| if (!isReply(input) || input._tag !== "WithExit") { | |
| return Effect.fail(new SchemaIssue.InvalidType(ast, Option.some(input))); | |
| } | |
| return Effect.mapBothEager(SchemaParser.decodeEffect(exit)(input.exit, options), { | |
| onFailure: issue => new SchemaIssue.Composite(ast, Option.some(input), [new SchemaIssue.Pointer(["exit"], issue)]), | |
| onSuccess: exit => new WithExit({ | |
| ...input, | |
| exit: exit | |
| }) | |
| }); | |
| }, { | |
| expected: "Reply.WithExit", | |
| toCodecJson: ([exit]) => Schema.link()(Schema.Struct({ | |
| _tag: Schema.Literal("WithExit"), | |
| requestId: SnowflakeFromBigInt, | |
| id: SnowflakeFromBigInt, | |
| exit | |
| }), SchemaTransformation.transform({ | |
| decode: encoded => new WithExit(encoded), | |
| encode: result => ({ | |
| ...result | |
| }) | |
| })) | |
| }); | |
| } | |
| /** | |
| * Returns a copy of this terminal reply associated with the supplied request id. | |
| * | |
| * @since 4.0.0 | |
| */ | |
| withRequestId(requestId) { | |
| return new WithExit({ | |
| ...this, | |
| requestId | |
| }); | |
| } | |
| } | |
| /** | |
| * Builds the transport codec for replies to the specified RPC, covering terminal | |
| * `WithExit` replies and streaming `Chunk` replies. | |
| * | |
| * @category schemas | |
| * @since 4.0.0 | |
| */ | |
| export const Reply = rpc => { | |
| if (schemaCache.has(rpc)) { | |
| return schemaCache.get(rpc); | |
| } | |
| const schema = Schema.toCodecJson(Schema.Union([WithExit.schema(rpc), Chunk.schema(rpc)])); | |
| schemaCache.set(rpc, schema); | |
| return schema; | |
| }; | |
| /** | |
| * Serializes a `ReplyWithContext` into its encoded wire representation, using the | |
| * reply's RPC schema and context and refailing encoding errors as | |
| * `MalformedMessage`. | |
| * | |
| * @category serialization | |
| * @since 4.0.0 | |
| */ | |
| export const serialize = self => { | |
| const schema = Reply(self.rpc); | |
| return MalformedMessage.refail(Effect.provideContext(Schema.encodeEffect(schema)(self.reply), self.context)); | |
| }; | |
| /** | |
| * Serializes an outgoing request's last received reply when one exists, returning | |
| * `None` when no reply has been received and refailing encoding errors as | |
| * `MalformedMessage`. | |
| * | |
| * @category serialization | |
| * @since 4.0.0 | |
| */ | |
| export const serializeLastReceived = self => { | |
| const lastReceivedReply = self.lastReceivedReply; | |
| if (lastReceivedReply._tag === "None") { | |
| return Effect.succeedNone; | |
| } | |
| const schema = Reply(self.rpc); | |
| return MalformedMessage.refail(Effect.provideContext(Schema.encodeEffect(schema)(lastReceivedReply.value), self.context)).pipe(Effect.map(Option.some)); | |
| }; | |
| //# sourceMappingURL=Reply.js.map |
Xet Storage Details
- Size:
- 8.93 kB
- Xet hash:
- 741d14dd0628f27ba096c7ca01cba9451ac9ba1ad72860eba6f4f957606540d8
·
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.