MCPcopy Create free account
hub / github.com/effect-app/libs / write

Function write

repos/effect/packages/effect/src/unstable/rpc/RpcClient.ts:556–595  ·  view source on GitHub ↗
(message: FromServer<Rpcs>)

Source from the content-addressed store, hash-verified

554 )
555 fiber.addObserver(() => {
556 resume(Effect.void)
557 })
558 })
559
560 const write = (message: FromServer<Rpcs>): Effect.Effect<void> => {
561 switch (message._tag) {
562 case "Chunk": {
563 const requestId = message.requestId
564 const entry = entries.get(requestId)
565 if (!entry || entry._tag !== "Queue") return Effect.void
566 return Queue.offerAll(entry.queue, message.values).pipe(
567 supportsAck
568 ? Effect.flatMap(() =>
569 options.onFromClient({
570 message: { _tag: "Ack", requestId: message.requestId },
571 context: entry.context,
572 discard: false
573 })
574 )
575 : identity,
576 Effect.catchCause((cause) => Queue.failCause(entry.queue, cause))
577 )
578 }
579 case "Exit": {
580 const requestId = message.requestId
581 const entry = entries.get(requestId)
582 if (!entry) return Effect.void
583 entries.delete(requestId)
584 if (entry._tag === "Effect") {
585 entry.resume(message.exit)
586 return Effect.void
587 }
588 return message.exit._tag === "Success"
589 ? Queue.end(entry.queue)
590 : Queue.failCause(entry.queue, message.exit.cause)
591 }
592 case "Defect": {
593 return clearEntries(Exit.die(message.defect))
594 }
595 case "ClientEnd": {
596 return Effect.void
597 }
598 }

Callers 3

RpcClient.tsFile · 0.70
makeProtocolSocketFunction · 0.70
sendFunction · 0.70

Calls 4

offerAllMethod · 0.80
getMethod · 0.65
pipeMethod · 0.65
endMethod · 0.65

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…