| 661 | } |
| 662 | const entries = new Map<RequestId, ClientEntry>() |
| 663 | |
| 664 | const { client, write } = yield* makeNoSerialization(group, { |
| 665 | ...options, |
| 666 | supportsAck, |
| 667 | onFromClient({ message }) { |
| 668 | switch (message._tag) { |
| 669 | case "Request": { |
| 670 | const rpc = group.requests.get(message.tag)! as any as Rpc.AnyWithProps |
| 671 | const collector = supportsTransferables ? Transferable.makeCollectorUnsafe() : undefined |
| 672 | |
| 673 | const fiber = Fiber.getCurrent()! |
| 674 | |
| 675 | const entry: ClientEntry = { |
| 676 | rpc, |
| 677 | context: collector ? Context.add(fiber.context, Transferable.Collector, collector) : fiber.context, |
| 678 | schemas: rpcSchemas(rpc) |
| 679 | } |
| 680 | entries.set(message.id, entry) |
| 681 | |
| 682 | return entry.schemas.encodePayload(message.payload).pipe( |
| 683 | Effect.provideContext(entry.context), |
| 684 | Effect.orDie, |
| 685 | Effect.flatMap((payload) => |
| 686 | send(clientId, { |
| 687 | ...message, |
| 688 | id: message.id, |
| 689 | payload, |
| 690 | headers: Object.entries(message.headers) |
| 691 | }, collector && collector.readUnsafe()) |
| 692 | ) |
| 693 | ) as Effect.Effect<void, RpcClientError> |
| 694 | } |
| 695 | case "Ack": { |
| 696 | const entry = entries.get(message.requestId) |
| 697 | if (!entry) return Effect.void |
| 698 | return send(clientId, { |
| 699 | _tag: "Ack", |
| 700 | requestId: message.requestId |
| 701 | }) as Effect.Effect<void, RpcClientError> |
| 702 | } |
| 703 | case "Interrupt": { |
| 704 | const entry = entries.get(message.requestId) |
| 705 | if (!entry) return Effect.void |
| 706 | entries.delete(message.requestId) |
| 707 | return send(clientId, { |
| 708 | _tag: "Interrupt", |
| 709 | requestId: message.requestId |
| 710 | }) as Effect.Effect<void, RpcClientError> |
| 711 | } |
| 712 | case "Eof": { |
| 713 | return Effect.void |
| 714 | } |
| 715 | } |