diff --git a/.changeset/tidy-lambdas-wait.md b/.changeset/tidy-lambdas-wait.md new file mode 100644 index 000000000..850bb275c --- /dev/null +++ b/.changeset/tidy-lambdas-wait.md @@ -0,0 +1,5 @@ +--- +"@reflag/node-sdk": patch +--- + +Wait for already in-flight bulk deliveries when flushing, including batches started by timers or reaching the batch size limit. Document AWS Lambda invocation lifecycle handling. diff --git a/packages/node-sdk/README.md b/packages/node-sdk/README.md index a2f756e9f..73f8593de 100644 --- a/packages/node-sdk/README.md +++ b/packages/node-sdk/README.md @@ -468,6 +468,44 @@ The SDK caches flag definitions in memory for fast performance. The first reques `EdgeClient` uses `flagsSyncMode: "in-request"`. Refresh fetch starts are throttled to at most once per second, and Cloudflare Workers cannot rely on delayed timer callbacks to run follow-up refreshes later. That means `refreshFlags()` calls made during the throttle window only mark a refresh as pending, so the call itself may resolve before the fetch runs. The queued refresh runs on the next request/access or `refreshFlags()` call after the throttle window expires. +## AWS Lambda + +Like Workers, Lambda cannot reliably finish background work after an invocation. +Unlike Workers' `ctx.waitUntil()`, an async Lambda handler must **await** flushing +before returning. Keep the client outside the handler for warm cache reuse, disable +periodic work, and flush in `finally`, including when application code throws: + +```typescript +import { ReflagClient } from "@reflag/node-sdk"; + +const reflag = new ReflagClient({ + flagsSyncMode: "in-request", + batchOptions: { intervalMs: 0, flushOnExit: false }, +}); + +export const handler = async () => { + try { + await reflag.initialize(); + const { isEnabled } = reflag.getFlag( + { user: { id: "user-id" }, company: { id: "company-id" } }, + "my-flag", + ); + return { statusCode: 200, body: JSON.stringify({ isEnabled }) }; + } finally { + await reflag.flush(); + } +}; +``` + +Process-exit hooks do not run when Lambda freezes an execution environment. +Setting `callbackWaitsForEmptyEventLoop` is not a substitute for awaiting `flush()`. +Allow enough invocation time for application work, initialization/refresh and +flushing (bulk HTTP requests have a 10-second timeout). Flushing waits for delivery +attempts; failed batches are logged and discarded, not retried. + +See [examples/aws-lambda](https://github.com/reflagcom/javascript/tree/main/packages/node-sdk/examples/aws-lambda) +for a deployable lifecycle probe and warm/cold invocation testing instructions. + ## Error Handling The SDK is designed to fail gracefully and never throw exceptions to the caller. Instead, it logs errors and provides diff --git a/packages/node-sdk/examples/aws-lambda/.gitignore b/packages/node-sdk/examples/aws-lambda/.gitignore new file mode 100644 index 000000000..c3b1df2cb --- /dev/null +++ b/packages/node-sdk/examples/aws-lambda/.gitignore @@ -0,0 +1,3 @@ +build/ +.aws-sam/ +samconfig.toml diff --git a/packages/node-sdk/examples/aws-lambda/README.md b/packages/node-sdk/examples/aws-lambda/README.md new file mode 100644 index 000000000..8c4217a67 --- /dev/null +++ b/packages/node-sdk/examples/aws-lambda/README.md @@ -0,0 +1,80 @@ +# AWS Lambda lifecycle probe + +Tests the **local Node SDK source**, not the published package. Requires Node, the +repository dependencies, AWS CLI v2, AWS SAM CLI, and a sandbox AWS profile with +CloudFormation, Lambda, IAM role creation/passing, S3 deployment, and Logs access. +Deployment creates billable AWS resources. No public function URL is created. + +## Build and deploy + +From the repository root: + +```sh +./node_modules/.bin/esbuild packages/node-sdk/examples/aws-lambda/handler.ts --bundle --platform=node --target=node22 --outfile=packages/node-sdk/examples/aws-lambda/build/handler.js +cd packages/node-sdk/examples/aws-lambda +export AWS_PROFILE=your-sandbox-profile AWS_REGION=eu-west-1 +aws sts get-caller-identity +sam deploy --template-file template.yaml --stack-name reflag-lambda-probe --resolve-s3 --capabilities CAPABILITY_IAM +FUNCTION=$(aws cloudformation describe-stacks --stack-name reflag-lambda-probe --query 'Stacks[0].Outputs[?OutputKey==`FunctionName`].OutputValue' --output text) +``` + +Use environment-based temporary credentials instead of a profile if preferred. +Never commit credentials or include them in invocation payloads. + +## Reproduce and compare + +The default endpoint is a delayed local HTTP server inside Lambda, requiring no +Reflag secret. It deterministically exposes the in-flight flush race, not external +network timeout behavior. The local server freezes with the function. + +```sh +for mode in timer full; do + for gap in 0 1 15 60; do + sleep "$gap" + aws lambda invoke --function-name "$FUNCTION" --cli-binary-format raw-in-base64-out \ + --payload "{\"mode\":\"$mode\",\"flush\":true}" /tmp/reflag-lambda-result.json + python3 -m json.tool /tmp/reflag-lambda-result.json + done +done +aws logs tail "/aws/lambda/$FUNCTION" --since 30m +``` + +- `timer`: lets an automatic timer start a request before calling `flush()`. +- `full`: fills a batch before calling `flush()`, like unawaited evaluation events. +- Both must report `pendingAtReturn: 0` with `flush: true` on the patched SDK. +- The original SDK returns with `pendingAtReturn > 0` against the delayed endpoint. +- Set `flush: false` as a negative control; pending work can cross invocations. +- Match `environmentId` to confirm warm reuse; Lambda does not guarantee reuse. +- Counters are cumulative per environment. Check failures, not just pending count: + SDK flush resolves even when delivery fails and events are discarded. + +For an A/B comparison, build/deploy this harness against the original +`src/batch-buffer.ts` in a separate worktree/stack, then repeat identical payloads. + +For **real network diagnosis**, set `REFLAG_API_BASE_URL=https://front.reflag.com` +and `REFLAG_SECRET_KEY` on the test Lambda using your approved secret-management +workflow. Use a disposable Reflag environment: this sends real tracking events. +Repeat the warm/cold/gap matrix and compare request durations, HTTP failures, error +names, and pending work. The probe intentionally only tests bulk delivery; it does +not initialize/fetch flag definitions. Also test the customer's runtime, memory, +VPC/NAT configuration, and handler duration before attributing their timeouts to +freeze/thaw. Do not increase timeouts as a substitute for lifecycle coordination. + +## Production pattern + +Keep a client outside the handler, configure `flagsSyncMode: "in-request"` and +`batchOptions: { intervalMs: 0, flushOnExit: false }`, initialize inside the handler, +and `await client.flush()` in `finally` after application work. Unlike Workers' +`ctx.waitUntil`, Lambda has no equivalent that keeps unawaited SDK work alive after +an async handler returns. Budget Lambda execution time for initialization, flag +refresh, application work and flushing; bulk requests currently have a 10s timeout. +`callbackWaitsForEmptyEventLoop` is not a replacement for awaiting SDK work. + +## Cleanup + +```sh +sam delete --stack-name reflag-lambda-probe +``` + +Verify deletion of the stack and log group. SAM's shared managed deployment bucket +may remain; do not remove it if other deployments use it. diff --git a/packages/node-sdk/examples/aws-lambda/handler.ts b/packages/node-sdk/examples/aws-lambda/handler.ts new file mode 100644 index 000000000..a022c3620 --- /dev/null +++ b/packages/node-sdk/examples/aws-lambda/handler.ts @@ -0,0 +1,131 @@ +import { randomUUID } from "node:crypto"; +import { createServer } from "node:http"; + +import { ReflagClient } from "../../src"; +import fetchClient from "../../src/fetch-http-client"; + +// A new ID means a cold start. Repeat invocations with gaps to observe thawing. +const environmentId = randomUUID(); +let invocation = 0; +let pending = 0; +let started = 0; +let succeeded = 0; +let failed = 0; +let setup: Promise | undefined; + +async function createClient() { + let apiBaseUrl = process.env.REFLAG_API_BASE_URL; + if (!apiBaseUrl) { + // Deterministic race probe: real HTTP, but no Reflag credentials/traffic. + // This server freezes too; use the real endpoint for network diagnosis. + const server = createServer((request, response) => { + request.resume(); + request.on("end", () => { + setTimeout(() => { + response.writeHead(200, { "content-type": "application/json" }); + response.end(JSON.stringify({ success: true })); + }, 250); + }); + }); + await new Promise((resolve) => + server.listen(0, "127.0.0.1", resolve), + ); + server.unref(); + const address = server.address(); + if (!address || typeof address === "string") + throw new Error("No server address"); + apiBaseUrl = `http://127.0.0.1:${address.port}`; + } else if (!process.env.REFLAG_SECRET_KEY) { + throw new Error("Real endpoint mode requires a test REFLAG_SECRET_KEY"); + } + + return new ReflagClient({ + secretKey: + process.env.REFLAG_SECRET_KEY || "lambda-local-probe-not-a-real-secret", + apiBaseUrl, + flagsSyncMode: "in-request", + batchOptions: { intervalMs: 100, maxSize: 100, flushOnExit: false }, + httpClient: { + ...fetchClient, + async post( + url: string, + headers: Record, + body: TBody, + ) { + pending++; + started++; + const start = Date.now(); + try { + const response = await fetchClient.post( + url, + headers, + body, + ); + if (response.ok) succeeded++; + else failed++; + return response; + } catch (error) { + failed++; + console.log( + JSON.stringify({ + environmentId, + error: error instanceof Error ? error.name : "unknown", + }), + ); + throw error; + } finally { + pending--; + console.log( + JSON.stringify({ + environmentId, + elapsedMs: Date.now() - start, + pending, + }), + ); + } + }, + }, + }); +} + +export async function handler( + event: { mode?: "timer" | "full"; flush?: boolean } = {}, + context: { + callbackWaitsForEmptyEventLoop: boolean; + getRemainingTimeInMillis(): number; + }, +) { + context.callbackWaitsForEmptyEventLoop = false; + const pendingAtEntry = pending; + const client = await (setup ??= createClient()); + invocation++; + const start = Date.now(); + const tracking: Promise[] = []; + // Full batches start I/O synchronously, as fire-and-forget flag checks can. + for (let i = 0; i < (event.mode === "full" ? 100 : 1); i++) { + tracking.push( + client.track(`lambda-probe-${environmentId}`, "lambda-lifecycle-probe"), + ); + } + if (event.mode !== "full") { + await Promise.all(tracking); + await new Promise((resolve) => setTimeout(resolve, 150)); + } + if (event.flush !== false) await client.flush(); + // Deliberately don't await full-batch tracking separately: test flush's contract. + const result = { + environmentId, + invocation, + mode: event.mode ?? "timer", + flush: event.flush !== false, + pendingAtEntry, + pendingAtReturn: pending, + started, + succeeded, + failed, + elapsedMs: Date.now() - start, + remainingMs: context.getRemainingTimeInMillis(), + }; + console.log(JSON.stringify(result)); + return result; +} diff --git a/packages/node-sdk/examples/aws-lambda/template.yaml b/packages/node-sdk/examples/aws-lambda/template.yaml new file mode 100644 index 000000000..28a278501 --- /dev/null +++ b/packages/node-sdk/examples/aws-lambda/template.yaml @@ -0,0 +1,23 @@ +AWSTemplateFormatVersion: "2010-09-09" +Transform: AWS::Serverless-2016-10-31 +Description: Temporary Node SDK Lambda lifecycle probe (no public endpoint) +Resources: + Probe: + Type: AWS::Serverless::Function + Properties: + CodeUri: build/ + Handler: handler.handler + Runtime: nodejs22.x + Timeout: 30 + MemorySize: 256 + ReservedConcurrentExecutions: 1 + Policies: + - AWSLambdaBasicExecutionRole + ProbeLogs: + Type: AWS::Logs::LogGroup + Properties: + LogGroupName: !Sub /aws/lambda/${Probe} + RetentionInDays: 1 +Outputs: + FunctionName: + Value: !Ref Probe diff --git a/packages/node-sdk/src/batch-buffer.ts b/packages/node-sdk/src/batch-buffer.ts index ec9f68584..edd9a2eea 100644 --- a/packages/node-sdk/src/batch-buffer.ts +++ b/packages/node-sdk/src/batch-buffer.ts @@ -8,6 +8,7 @@ import { isObject, ok } from "./utils"; */ export default class BatchBuffer { private buffer: T[] = []; + private inFlight = new Set>(); private flushHandler: (items: T[]) => Promise; private logger?: Logger; private maxSize: number; @@ -66,12 +67,24 @@ export default class BatchBuffer { if (this.buffer.length === 0) { this.logger?.debug("buffer is empty. nothing to flush"); - return; + } else { + const flushingBuffer = this.buffer; + this.buffer = []; + const request = this.send(flushingBuffer); + this.inFlight.add(request); + // Keep completed requests out of subsequent flushes, including failures. + void request.then( + () => this.inFlight.delete(request), + () => this.inFlight.delete(request), + ); } - const flushingBuffer = this.buffer; - this.buffer = []; + // Automatic flushes remove events from the buffer before delivery finishes. + // Explicit flushes must also await those requests before a runtime can freeze. + await Promise.all(this.inFlight); + } + private async send(flushingBuffer: T[]): Promise { try { await this.flushHandler(flushingBuffer); diff --git a/packages/node-sdk/test/batch-buffer.test.ts b/packages/node-sdk/test/batch-buffer.test.ts index ebafe2b7c..c802836f0 100644 --- a/packages/node-sdk/test/batch-buffer.test.ts +++ b/packages/node-sdk/test/batch-buffer.test.ts @@ -70,6 +70,7 @@ describe("BatchBuffer", () => { expect(buffer).toEqual({ buffer: [], + inFlight: new Set(), flushHandler: mockFlushHandler, timer: null, intervalMs: 33, @@ -82,6 +83,7 @@ describe("BatchBuffer", () => { const buffer = new BatchBuffer({ flushHandler: mockFlushHandler }); expect(buffer).toEqual({ buffer: [], + inFlight: new Set(), flushHandler: mockFlushHandler, intervalMs: BATCH_INTERVAL_MS, maxSize: BATCH_MAX_SIZE, @@ -155,6 +157,77 @@ describe("BatchBuffer", () => { expect(itemsFlushed).toBe(1); }); + it.each(["full batch", "timer"])( + "waits for an in-flight %s even when the buffer is empty", + async (trigger) => { + vi.useFakeTimers(); + try { + let release!: () => void; + const delivery = new Promise((resolve) => { + release = resolve; + }); + const send = vi.fn(() => delivery); + const buffer = new BatchBuffer({ + flushHandler: send, + maxSize: trigger === "full batch" ? 1 : 10, + intervalMs: 100, + }); + const adding = buffer.add("event"); + if (trigger === "timer") vi.advanceTimersByTime(100); + expect(send).toHaveBeenCalledOnce(); + let completed = false; + const flushing = buffer.flush().then(() => { + completed = true; + }); + await Promise.resolve(); + await Promise.resolve(); + expect(completed).toBe(false); + release(); + await Promise.all([adding, flushing]); + expect(completed).toBe(true); + expect(send).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } + }, + ); + + it("awaits all overlapping batches, including failed delivery", async () => { + let release!: () => void; + let reject!: (error: Error) => void; + const first = new Promise((resolve) => { + release = resolve; + }); + const second = new Promise((_, fail) => { + reject = fail; + }); + const send = vi + .fn() + .mockReturnValueOnce(first) + .mockReturnValueOnce(second); + const buffer = new BatchBuffer({ + flushHandler: send, + maxSize: 1, + logger: mockLogger, + }); + const addingFirst = buffer.add("first"); + const addingSecond = buffer.add("second"); + let completed = false; + const flushing = buffer.flush().then(() => { + completed = true; + }); + reject(new Error("timeout")); + await Promise.resolve(); + await Promise.resolve(); + expect(completed).toBe(false); + release(); + await Promise.all([addingFirst, addingSecond, flushing]); + expect(completed).toBe(true); + expect(mockLogger.warn).toHaveBeenCalledOnce(); + await buffer.flush(); + expect(send).toHaveBeenCalledTimes(2); + }); + it("should flush buffer", async () => { const buffer = new BatchBuffer({ flushHandler: mockFlushHandler,