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

Function notifyLocal

packages/effect/src/unstable/cluster/Sharding.ts:1034–1088  ·  view source on GitHub ↗
(
    message: M,
    discard: boolean,
    options?: {
      readonly waitUntilRead?: boolean | undefined
    }
  )

Source from the content-addressed store, hash-verified

1032 }
1033 const pendingNotifications = new Map<Snowflake.Snowflake, PendingNotification>()
1034 const notifyLocal = <M extends Message.Outgoing<any> | Message.Incoming<any>>(
1035 message: M,
1036 discard: boolean,
1037 options?: {
1038 readonly waitUntilRead?: boolean | undefined
1039 }
1040 ) =>
1041 Effect.suspend(function loop(): Effect.Effect<
1042 void,
1043 | EntityNotAssignedToRunner
1044 | AlreadyProcessingMessage
1045 | (M extends Message.Incoming<any> ? never : PersistenceError)
1046 > {
1047 const address = message.envelope.address
1048 const state = entityManagers.get(address.entityType)
1049 if (!state) {
1050 return Effect.flatMap(waitForEntityManager(address.entityType), loop)
1051 } else if (state.status === "closed" || !isEntityOnLocalShards(address)) {
1052 return Effect.fail(new EntityNotAssignedToRunner({ address }))
1053 }
1054
1055 const isLocal = isEntityOnLocalShards(address)
1056 const notify = storageEnabled
1057 ? openStorageReadLatch
1058 : () => Effect.die("Sharding.notifyLocal: storage is disabled")
1059
1060 if (message._tag === "IncomingRequest" || message._tag === "IncomingEnvelope") {
1061 if (!isLocal) {
1062 return Effect.fail(new EntityNotAssignedToRunner({ address }))
1063 } else if (
1064 message._tag === "IncomingRequest" && state.manager.isProcessingFor(message, { excludeReplies: true })
1065 ) {
1066 return Effect.fail(new AlreadyProcessingMessage({ address, envelopeId: message.envelope.requestId }))
1067 } else if (message._tag === "IncomingRequest" && options?.waitUntilRead) {
1068 if (!storageEnabled) return notify()
1069 return Effect.callback<void, EntityNotAssignedToRunner>((resume) => {
1070 let entry = pendingNotifications.get(message.envelope.requestId)
1071 if (entry) {
1072 const prevResume = entry.resume
1073 entry.resume = (effect) => {
1074 prevResume(effect)
1075 resume(effect)
1076 }
1077 return
1078 }
1079 entry = { resume, message }
1080 pendingNotifications.set(message.envelope.requestId, entry)
1081 storageReadLatch.openUnsafe()
1082 })
1083 }
1084 return notify()
1085 }
1086
1087 return runnersService.notifyLocal({ message, notify, discard, storageOnly: !isLocal }) as any
1088 })
1089
1090 function sendOutgoing(
1091 message: Message.Outgoing<any>,

Callers 2

sendOutgoingFunction · 0.70
Sharding.tsFile · 0.70

Calls 8

waitForEntityManagerFunction · 0.85
isEntityOnLocalShardsFunction · 0.85
notifyFunction · 0.70
resumeFunction · 0.70
getMethod · 0.65
setMethod · 0.65
openUnsafeMethod · 0.65
failMethod · 0.45

Tested by

no test coverage detected