| 475 | * @since 4.0.0 |
| 476 | */ |
| 477 | export const make = ( |
| 478 | storage: Omit< |
| 479 | MessageStorage["Service"], |
| 480 | "registerReplyHandler" | "unregisterReplyHandler" | "unregisterShardReplyHandlers" |
| 481 | > |
| 482 | ): Effect.Effect<MessageStorage["Service"]> => |
| 483 | Effect.sync(() => { |
| 484 | type ReplyHandler = { |
| 485 | readonly message: Message.OutgoingRequest<any> | Message.IncomingRequest<any> |
| 486 | readonly shardSet: Set<ReplyHandler> |
| 487 | readonly respond: (reply: Reply.ReplyWithContext<any>) => Effect.Effect<void, PersistenceError | MalformedMessage> |
| 488 | readonly resume: (effect: Effect.Effect<void, EntityNotAssignedToRunner>) => void |
| 489 | } |
| 490 | const replyHandlers = new Map<Snowflake.Snowflake, Array<ReplyHandler>>() |
| 491 | const replyHandlersShard = new Map<string, Set<ReplyHandler>>() |
| 492 | return MessageStorage.of({ |
| 493 | ...storage, |
| 494 | registerReplyHandler: (message) => { |
| 495 | const requestId = message.envelope.requestId |
| 496 | return Effect.callback<void, EntityNotAssignedToRunner>((resume) => { |
| 497 | const shardId = message.envelope.address.shardId.toString() |
| 498 | let handlers = replyHandlers.get(requestId) |
| 499 | if (handlers === undefined) { |
| 500 | handlers = [] |
| 501 | replyHandlers.set(requestId, handlers) |
| 502 | } |
| 503 | let shardSet = replyHandlersShard.get(shardId) |
| 504 | if (!shardSet) { |
| 505 | shardSet = new Set() |
| 506 | replyHandlersShard.set(shardId, shardSet) |
| 507 | } |
| 508 | const entry: ReplyHandler = { |
| 509 | message, |
| 510 | shardSet, |
| 511 | respond: message._tag === "IncomingRequest" ? message.respond : (reply) => message.respond(reply.reply), |
| 512 | resume |
| 513 | } |
| 514 | handlers.push(entry) |
| 515 | shardSet.add(entry) |
| 516 | return Effect.sync(() => { |
| 517 | const index = handlers.indexOf(entry) |
| 518 | handlers.splice(index, 1) |
| 519 | shardSet.delete(entry) |
| 520 | }) |
| 521 | }) |
| 522 | }, |
| 523 | unregisterReplyHandler: (requestId) => |
| 524 | Effect.sync(() => { |
| 525 | const handlers = replyHandlers.get(requestId) |
| 526 | if (!handlers) return |
| 527 | replyHandlers.delete(requestId) |
| 528 | for (let i = 0; i < handlers.length; i++) { |
| 529 | const handler = handlers[i] |
| 530 | handler.shardSet.delete(handler) |
| 531 | handler.resume(Effect.fail( |
| 532 | new EntityNotAssignedToRunner({ |
| 533 | address: handler.message.envelope.address |
| 534 | }) |