(
message: M,
discard: boolean,
options?: {
readonly waitUntilRead?: boolean | undefined
}
)
| 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>, |
no test coverage detected