Skip to content

Commit

Permalink
fix: Write UNPROCESSED_STREAM_MESSAGES metric on both failure/succe…
Browse files Browse the repository at this point in the history
…ss (#352)
  • Loading branch information
morgsmccauley authored Nov 1, 2023
1 parent 50a3d6c commit 2096c39
Showing 1 changed file with 9 additions and 9 deletions.
18 changes: 9 additions & 9 deletions runner/src/stream-handler/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,11 @@ void (async function main () {
console.log('Started processing stream: ', streamKey);

let indexerName = '';
const streamType = redisClient.getStreamType(streamKey);

while (true) {
try {
const startTime = performance.now();
const streamType = redisClient.getStreamType(streamKey);

const messages = await redisClient.getNextStreamMessage(streamKey);
const indexerConfig = await redisClient.getStreamStorage(streamKey);
Expand Down Expand Up @@ -52,14 +52,6 @@ void (async function main () {

await redisClient.deleteStreamMessage(streamKey, id);

const unprocessedMessages = await redisClient.getUnprocessedStreamMessages(streamKey);

parentPort?.postMessage({
type: 'UNPROCESSED_STREAM_MESSAGES',
labels: { indexer: indexerName, type: streamType },
value: unprocessedMessages?.length ?? 0,
} satisfies Message);

parentPort?.postMessage({
type: 'EXECUTION_DURATION',
labels: { indexer: indexerName, type: streamType },
Expand All @@ -70,6 +62,14 @@ void (async function main () {
} catch (err) {
await sleep(10000);
console.log(`Failed: ${indexerName}`, err);
} finally {
const unprocessedMessages = await redisClient.getUnprocessedStreamMessages(streamKey);

parentPort?.postMessage({
type: 'UNPROCESSED_STREAM_MESSAGES',
labels: { indexer: indexerName, type: streamType },
value: unprocessedMessages?.length ?? 0,
} satisfies Message);
}
}
})();

0 comments on commit 2096c39

Please sign in to comment.