(
message: Message.Outgoing<any>,
afterPersist: (message: Message.Outgoing<any>, isDuplicate: boolean) => Effect.Effect<void, E>
)
| 128 | const requestIdRewrites = new Map<Snowflake.Snowflake, Snowflake.Snowflake>() |
| 129 | |
| 130 | function notifyWith<E>( |
| 131 | message: Message.Outgoing<any>, |
| 132 | afterPersist: (message: Message.Outgoing<any>, isDuplicate: boolean) => Effect.Effect<void, E> |
| 133 | ): Effect.Effect<void, E | PersistenceError> { |
| 134 | const rpc = message.rpc as any as Rpc.AnyWithProps |
| 135 | const persisted = Context.get(rpc.annotations, Persisted) |
| 136 | if (!persisted) { |
| 137 | return Effect.dieMessage("Runners.notify only supports persisted messages") |
| 138 | } |
| 139 | |
| 140 | if (message._tag === "OutgoingEnvelope") { |
| 141 | const rewriteId = requestIdRewrites.get(message.envelope.requestId) |
| 142 | const requestId = rewriteId ?? message.envelope.requestId |
| 143 | const entry = storageRequests.get(requestId) |
| 144 | if (rewriteId) { |
| 145 | message = new Message.OutgoingEnvelope({ |
| 146 | ...message, |
| 147 | envelope: message.envelope.withRequestId(rewriteId) |
| 148 | }) |
| 149 | } |
| 150 | return storage.saveEnvelope(message).pipe( |
| 151 | Effect.catchTag("MalformedMessage", Effect.die), |
| 152 | Effect.zipRight( |
| 153 | entry ? Effect.zipRight(entry.latch.open, afterPersist(message, false)) : afterPersist(message, false) |
| 154 | ) |
| 155 | ) |
| 156 | } |
| 157 | |
| 158 | // For requests, after persisting the request, we need to check if the |
| 159 | // request is a duplicate. If it is, we need to resume from the last |
| 160 | // received reply. |
| 161 | // |
| 162 | // Otherwise, we notify the remote entity and then reply from storage. |
| 163 | return Effect.flatMap( |
| 164 | Effect.catchTag(storage.saveRequest(message), "MalformedMessage", Effect.die), |
| 165 | MessageStorage.SaveResult.$match({ |
| 166 | Success: () => afterPersist(message, false), |
| 167 | Duplicate: ({ lastReceivedReply, originalId }) => { |
| 168 | // If the last received reply is an exit, we can just return it |
| 169 | // as the response. |
| 170 | if (Option.isSome(lastReceivedReply) && lastReceivedReply.value._tag === "WithExit") { |
| 171 | return message.respond(lastReceivedReply.value.withRequestId(message.envelope.requestId)) |
| 172 | } |
| 173 | requestIdRewrites.set(message.envelope.requestId, originalId) |
| 174 | return afterPersist( |
| 175 | new Message.OutgoingRequest({ |
| 176 | ...message, |
| 177 | lastReceivedReply, |
| 178 | envelope: Envelope.makeRequest({ |
| 179 | ...message.envelope, |
| 180 | requestId: originalId |
| 181 | }), |
| 182 | respond(reply) { |
| 183 | if (reply._tag === "WithExit") { |
| 184 | requestIdRewrites.delete(message.envelope.requestId) |
| 185 | } |
| 186 | return message.respond(reply.withRequestId(message.envelope.requestId)) |
| 187 | } |
no test coverage detected
searching dependent graphs…