diff --git a/src/metrics/dogstatsd.spec.ts b/src/metrics/dogstatsd.spec.ts index 498be7f7..c93e4c1e 100644 --- a/src/metrics/dogstatsd.spec.ts +++ b/src/metrics/dogstatsd.spec.ts @@ -1,13 +1,35 @@ import * as dgram from "node:dgram"; + +import { logDebug } from "../utils"; import { LambdaDogStatsD } from "./dogstatsd"; jest.mock("node:dgram", () => ({ createSocket: jest.fn(), })); +jest.mock("../utils", () => ({ + logDebug: jest.fn(), +})); describe("LambdaDogStatsD", () => { - let mockSend: jest.Mock; - let mockUnref: jest.Mock; + let mockSend: jest.Mock void]>; + let mockUnref: jest.Mock; + + function useControlledSocket(): Array<(error?: Error) => void> { + const callbacks: Array<(error?: Error) => void> = []; + mockSend.mockImplementation((_message, _port, _host, callback) => { + callbacks.push(callback); + }); + return callbacks; + } + + /** + * @param {Promise} promise + */ + function observeCompletion(promise: Promise): jest.Mock { + const completed = jest.fn(); + void promise.then(completed); + return completed; + } beforeEach(() => { // A send() that immediately calls its callback @@ -23,6 +45,7 @@ describe("LambdaDogStatsD", () => { }); afterEach(() => { + jest.useRealTimers(); jest.clearAllMocks(); }); @@ -62,8 +85,10 @@ describe("LambdaDogStatsD", () => { }); it("flush() resolves immediately when there are no sends", async () => { + jest.useFakeTimers(); const client = new LambdaDogStatsD(); await expect(client.flush()).resolves.toBeUndefined(); + expect(jest.getTimerCount()).toBe(0); }); it("unrefs socket on client creation", () => { @@ -72,26 +97,145 @@ describe("LambdaDogStatsD", () => { expect(mockUnref).toHaveBeenCalledTimes(1); }); - it("flush() times out if a send never invokes its callback", async () => { - // replace socket.send with a never‐calling callback - (dgram.createSocket as jest.Mock).mockReturnValue({ - send: jest.fn(), // never calls callback - unref: jest.fn(), - getSendBufferSize: jest.fn(), - setSendBufferSize: jest.fn(), - bind: jest.fn(), + it("flush() waits for one pending send and clears its timeout", async () => { + jest.useFakeTimers(); + const callbacks = useControlledSocket(); + const client = new LambdaDogStatsD(); + client.distribution("pending", 1); + + const flush = client.flush(); + const completed = observeCompletion(flush); + await Promise.resolve(); + expect(completed).not.toHaveBeenCalled(); + + callbacks[0](); + await expect(flush).resolves.toBeUndefined(); + expect(jest.getTimerCount()).toBe(0); + }); + + it("flush() waits for every pending send", async () => { + const callbacks = useControlledSocket(); + const client = new LambdaDogStatsD(); + client.distribution("first", 1); + client.distribution("second", 2); + client.distribution("third", 3); + + const flush = client.flush(); + const completed = observeCompletion(flush); + callbacks[0](); + callbacks[1](); + await Promise.resolve(); + expect(completed).not.toHaveBeenCalled(); + + callbacks[2](); + await expect(flush).resolves.toBeUndefined(); + }); + + it("flush() waits for sends started while it is pending", async () => { + const callbacks = useControlledSocket(); + const client = new LambdaDogStatsD(); + client.distribution("first", 1); + + const flush = client.flush(); + const completed = observeCompletion(flush); + client.distribution("second", 2); + callbacks[0](); + await Promise.resolve(); + expect(completed).not.toHaveBeenCalled(); + + callbacks[1](); + await expect(flush).resolves.toBeUndefined(); + }); + + it("flush() settles concurrent callers after the pending send", async () => { + const callbacks = useControlledSocket(); + const client = new LambdaDogStatsD(); + client.distribution("pending", 1); + + const firstFlush = client.flush(); + const secondFlush = client.flush(); + expect(secondFlush).toBe(firstFlush); + callbacks[0](); + + await Promise.all([firstFlush, secondFlush]); + }); + + it("logs callback errors and completes the send", async () => { + jest.useFakeTimers(); + const error = new Error("callback failure"); + const callbacks = useControlledSocket(); + const client = new LambdaDogStatsD(); + + client.distribution("metric", 1); + + const flush = client.flush(); + callbacks[0](error); + expect(jest.getTimerCount()).toBe(0); + await expect(flush).resolves.toBeUndefined(); + expect(logDebug).toHaveBeenCalledWith("Unable to send metric packet: callback failure"); + }); + + it("logs synchronous socket errors without throwing", async () => { + jest.useFakeTimers(); + mockSend.mockImplementation(() => { + throw new Error("synchronous failure"); }); + const client = new LambdaDogStatsD(); + client.distribution("metric", 1); + + const flush = client.flush(); + expect(jest.getTimerCount()).toBe(0); + await expect(flush).resolves.toBeUndefined(); + expect(logDebug).toHaveBeenCalledWith("Unable to send metric packet: synchronous failure"); + }); + + it("normalizes non-Error values thrown synchronously by the socket", async () => { + jest.useFakeTimers(); + mockSend.mockImplementation(() => { + throw "non-error failure"; + }); const client = new LambdaDogStatsD(); - client.distribution("will", 9); + client.distribution("metric", 1); + + const flush = client.flush(); + expect(jest.getTimerCount()).toBe(0); + await expect(flush).resolves.toBeUndefined(); + expect(logDebug).toHaveBeenCalledWith("Unable to send metric packet: Unknown socket send failure"); + }); + + it("flush() times out if a send never invokes its callback", async () => { jest.useFakeTimers(); - const p = client.flush(); - // advance past the 1000ms MAX_FLUSH_TIMEOUT - jest.advanceTimersByTime(1100); + useControlledSocket(); - // expect the Promise returned by flush() to resolve successfully - await expect(p).resolves.toBeUndefined(); - jest.useRealTimers(); + const client = new LambdaDogStatsD(); + client.distribution("will", 9); + + const flush = client.flush(); + jest.advanceTimersByTime(1000); + + await expect(flush).resolves.toBeUndefined(); + expect(logDebug).toHaveBeenCalledWith("Timed out before sending all metric payloads"); + }); + + it("ignores late callbacks from a timed-out flush when the client is reused", async () => { + jest.useFakeTimers(); + const callbacks = useControlledSocket(); + const client = new LambdaDogStatsD(); + client.distribution("old", 1); + const oldFlush = client.flush(); + jest.advanceTimersByTime(1000); + await oldFlush; + + client.distribution("new", 2); + const newFlush = client.flush(); + const completed = observeCompletion(newFlush); + callbacks[0](); + await Promise.resolve(); + expect(completed).not.toHaveBeenCalled(); + + callbacks[1](); + await expect(newFlush).resolves.toBeUndefined(); }); }); diff --git a/src/metrics/dogstatsd.ts b/src/metrics/dogstatsd.ts index 1c559bf7..47641a7e 100644 --- a/src/metrics/dogstatsd.ts +++ b/src/metrics/dogstatsd.ts @@ -12,9 +12,14 @@ export class LambdaDogStatsD { private static readonly TAG_SUB = "_"; // The maximum amount to wait while flushing pending sends, so we don't block forever. private static readonly MAX_FLUSH_TIMEOUT = 1000; + static readonly #resolvedFlush = Promise.resolve(); private readonly socket: dgram.Socket; - private readonly pendingSends = new Set>(); + #pendingSends = 0; + #sendGeneration = 0; + #flushPromise: Promise | undefined; + #resolveFlush: (() => void) | undefined; + #flushTimeout: NodeJS.Timeout | undefined; constructor() { this.socket = dgram.createSocket(LambdaDogStatsD.SOCKET_TYPE); @@ -48,34 +53,78 @@ export class LambdaDogStatsD { this.send(payload); } - private send(packet: string) { - const msg = Buffer.from(packet, LambdaDogStatsD.ENCODING); - const promise = new Promise((resolve) => { - this.socket.send(msg, LambdaDogStatsD.PORT, LambdaDogStatsD.HOST, (err) => { - if (err) { - logDebug(`Unable to send metric packet: ${err.message}`); - } + /** + * @param {string} packet + */ + private send(packet: string): void { + const message = Buffer.from(packet, LambdaDogStatsD.ENCODING); + const generation = this.#sendGeneration; + this.#pendingSends++; - resolve(); + try { + this.socket.send(message, LambdaDogStatsD.PORT, LambdaDogStatsD.HOST, (error) => { + this.#completeSend(generation, error); }); - }); + } catch (error) { + const sendError = error instanceof Error ? error : new Error("Unknown socket send failure"); + this.#completeSend(generation, sendError); + } + } + + /** + * @param {number} generation + * @param {Error | null} [error] + */ + #completeSend(generation: number, error?: Error | null): void { + if (error) { + logDebug(`Unable to send metric packet: ${error.message}`); + } + + if (generation !== this.#sendGeneration) { + return; + } - this.pendingSends.add(promise); - void promise.finally(() => this.pendingSends.delete(promise)); + this.#pendingSends--; + if (this.#pendingSends === 0 && this.#resolveFlush !== undefined) { + this.#finishFlush(false); + } } - /** Block until all in-flight sends have settled */ - public async flush(): Promise { - const allSettled = Promise.allSettled(this.pendingSends); - const maxTimeout = new Promise<"timeout">((resolve) => { - setTimeout(() => resolve("timeout"), LambdaDogStatsD.MAX_FLUSH_TIMEOUT); - }); + /** + * @param {boolean} timedOut + */ + #finishFlush(timedOut: boolean): void { + const resolve = this.#resolveFlush!; + + if (this.#flushTimeout !== undefined) { + clearTimeout(this.#flushTimeout); + } + this.#flushPromise = undefined; + this.#resolveFlush = undefined; + this.#flushTimeout = undefined; - const winner = await Promise.race([allSettled, maxTimeout]); - if (winner === "timeout") { + if (timedOut) { + this.#pendingSends = 0; + this.#sendGeneration++; logDebug("Timed out before sending all metric payloads"); } - this.pendingSends.clear(); + resolve(); + } + + /** Block until all in-flight sends have settled or the flush timeout expires. */ + public flush(): Promise { + if (this.#pendingSends === 0) { + return LambdaDogStatsD.#resolvedFlush; + } + if (this.#flushPromise !== undefined) { + return this.#flushPromise; + } + + this.#flushPromise = new Promise((resolve) => { + this.#resolveFlush = resolve; + }); + this.#flushTimeout = setTimeout(() => this.#finishFlush(true), LambdaDogStatsD.MAX_FLUSH_TIMEOUT); + return this.#flushPromise; } }