( payload: EventsQueuePayloadIncomingEvent['payload'], partitionKey: string )
| 184 | }; |
| 185 | |
| 186 | export const produceIncomingEvent = async ( |
| 187 | payload: EventsQueuePayloadIncomingEvent['payload'], |
| 188 | partitionKey: string |
| 189 | ): Promise<void> => { |
| 190 | const p = await getProducer(); |
| 191 | try { |
| 192 | await p.send({ |
| 193 | topic: KAFKA_EVENTS_TOPIC, |
| 194 | timeout: KAFKA_REQUEST_TIMEOUT_MS, |
| 195 | messages: [ |
| 196 | { |
| 197 | key: Buffer.from(partitionKey), |
| 198 | value: Buffer.from(JSON.stringify(payload)), |
| 199 | }, |
| 200 | ], |
| 201 | }); |
| 202 | } catch (err) { |
| 203 | if (isFatalProducerError(err)) { |
| 204 | kafkaLogger.warn( |
| 205 | { err }, |
| 206 | 'kafka producer in fatal state; resetting for next call' |
| 207 | ); |
| 208 | resetProducer(p); |
| 209 | } |
| 210 | throw err; |
| 211 | } |
| 212 | }; |
| 213 | |
| 214 | const consumers = new Set<Consumer>(); |
| 215 |
no test coverage detected