(message: FromServer<Rpcs>)
| 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 | } |
no test coverage detected
searching dependent graphs…