MCPcopy Create free account
hub / github.com/Effect-TS/effect / notifyWith

Function notifyWith

packages/cluster/src/Runners.ts:130–194  ·  view source on GitHub ↗
(
    message: Message.Outgoing<any>,
    afterPersist: (message: Message.Outgoing<any>, isDuplicate: boolean) => Effect.Effect<void, E>
  )

Source from the content-addressed store, hash-verified

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 }

Callers 2

notifyFunction · 0.85
notifyLocalFunction · 0.85

Calls 5

getMethod · 0.65
dieMessageMethod · 0.65
pipeMethod · 0.65
setMethod · 0.65
withRequestIdMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…