DEV Community

Cover image for Streaming Bedrock Tokens with AWS Lambda Function URLs: Real Architecture, Real Tradeoffs
Varun Sharma
Varun Sharma

Posted on Originally published at github.com

Streaming Bedrock Tokens with AWS Lambda Function URLs: Real Architecture, Real Tradeoffs

When I first wired up Amazon Bedrock's ConverseStream to a web frontend, I assumed the tricky part would be parsing the EventStream chunks in JavaScript.

It wasn't. The real headache was everything sitting between Bedrock and the browser.

If you deploy a standard Lambda function behind an API Gateway HTTP API, the problem becomes obvious on the first test: nothing reaches the screen until the entire response finishes. For a quick 20-word greeting, nobody notices. But for a multi-paragraph Claude 3.7 or Nova response that takes 15 seconds to generate, your user is left staring at a dead loading spinner.

Bedrock streams tokens natively, but getting those chunks to a client browser over serverless AWS infrastructure without buffering them in transit turns out to have a few sharp edges.


When to Use API Gateway vs. Function URLs

Most tutorials reflexively point you to API Gateway. And to be fair, AWS quietly addressed the biggest historical complaint in late 2025: REST APIs finally support native response streaming (via STREAM transfer mode), and you can raise integration timeouts up to 15 minutes.

So before you bypass API Gateway, look at your existing stack:

If your workload already runs behind an API Gateway REST API—and you rely on its usage plans, request validation, or Cognito authorizers—you probably don't need to reinvent the wheel. Switch your method integration to Transfer Mode: STREAM, configure your client to handle chunked responses, and keep your existing API management in place.

However, there are two practical reasons you might want Lambda Function URLs instead:

First, API Gateway HTTP APIs (the cheaper, lighter variant most teams prefer for greenfield serverless apps) still do not support response streaming. They buffer the entire payload before returning a single byte.

Second, the economics shift at scale. REST APIs charge $3.50 per million calls plus an extra $0.09 per gigabyte of streamed data (after the first 10 GB). If you're building a chat application streaming millions of long completions a month, paying API Gateway's data transfer and request markup on high-volume token streams starts looking like unnecessary overhead.

Lambda Function URLs with InvokeMode: RESPONSE_STREAM give you native HTTP chunked transfer encoding (Transfer-Encoding: chunked) directly from the Lambda microVM to the client. You skip the API Gateway layer and its per-request fees completely.

The trade-off is access control. A Function URL is a naked public HTTPS endpoint unless you lock it down. And doing that behind CloudFront without breaking streaming takes some deliberate design.

flowchart TD
    Client[Web Browser: fetch ReadableStream] -->|SSE Stream| CF[Amazon CloudFront CDN]
    CF -->|Chunked HTTP with X-Origin-Verify| FURL[Lambda Function URL: RESPONSE_STREAM]
    FURL -->|EventStream| Bedrock[Amazon Bedrock ConverseStream]

The CloudFront & Auth Dilemma

AWS documentation frequently suggests locking down Function URLs with AWS_IAM auth and putting CloudFront Origin Access Control (OAC) in front.

Here’s why that falls apart for interactive web apps:

CloudFront OAC cannot sign browser-initiated POST request bodies with SigV4 without a custom edge signing proxy or Lambda@Edge worker. Browser fetch requests cannot generate AWS SigV4 signatures without baking IAM access keys into client-side JavaScript—an obvious security disaster.

If you flip the Function URL to AuthType: NONE so browsers can POST prompts directly, your Lambda endpoint is exposed to anyone on the internet who discovers the URL.

The battle-tested serverless pattern uses a layered compromise:

  1. Origin Restriction: Set the Function URL to AuthType: NONE, but place CloudFront in front of it. Have CloudFront inject a secret custom header (X-Origin-Verify: <secret>) when forwarding requests to the origin. In the Lambda handler, verify this header using constant-time string comparison (crypto.timingSafeEqual in Node, secrets.compare_digest in Python) before touching Bedrock. Any direct hit to the Function URL gets bounced with a 403.
  2. User Authentication: Remember that X-Origin-Verify only proves the request came through CloudFront; it is not user authentication. In a production application, client requests must carry user session tokens (Cognito JWTs, session cookies, or API keys). You can validate these at the CloudFront edge via Lambda@Edge / CloudFront Functions, or inside the Lambda handler.
  3. Edge Buffering Policies: CloudFront's default behavior will buffer chunked streams. You must explicitly attach two AWS-managed policies to the cache behavior:
    • CachePolicyId: 4135ea2d-6df8-44a3-9df3-44ca84e08fad (CachingDisabled): Tells CloudFront edge locations to flush SSE chunks immediately instead of aggregating them.
    • OriginRequestPolicyId: b689b0a8-53d0-40ab-baf2-68738e2966ac (AllViewerExceptHostHeader): Rewrites the viewer's Host header to match the Function URL domain. Without this policy, Lambda rejects the request with a 403 because the Host header does not match the origin.
  4. WAF Protection: Attach AWS WAF to CloudFront to enforce IP rate limits, geo-restrictions, and bot control before requests reach Lambda.

The Node.js 22 Gotcha: Bedrock Event Ordering

In Node.js 22, response streaming is supported via awslambda.streamifyResponse(). Wrapping responseStream with awslambda.HttpResponseStream.from(responseStream, headers) lets you emit HTTP status codes and CORS headers before you write SSE (text/event-stream) chunks.

Here is a subtle gotcha that cost me an afternoon:

If you follow naive SDK samples, you'll probably look for token usage inside the chunk.messageStop event. But in Bedrock's actual ConverseStream protocol, chunk.messageStop only carries the stopReason. Token usage metrics (inputTokens, outputTokens) arrive in a subsequent metadata chunk right before the connection closes.

If you attempt to read chunk.metadata?.usage inside the messageStop branch, it evaluates to undefined every single time.

To accurately report token counts, you must collect stopReason and usage across the loop and send your terminal completion event once both have arrived:

// lambda/index.mjs
import crypto from "node:crypto";
import {
  BedrockRuntimeClient,
  ConverseStreamCommand,
} from "@aws-sdk/client-bedrock-runtime";

const bedrock = new BedrockRuntimeClient({
  region: process.env.AWS_REGION || "us-east-1",
});

// Hardcode allowed models — never trust client-supplied model IDs
const ALLOWED_MODELS = new Set([
  "amazon.nova-pro-v1:0",
  "amazon.nova-lite-v1:0",
  "us.anthropic.claude-3-7-sonnet-20250219-v1:0",
  "us.anthropic.claude-3-5-haiku-20241022-v1:0",
]);

const DEFAULT_MODEL_ID = process.env.BEDROCK_MODEL_ID || "amazon.nova-pro-v1:0";
const ALLOWED_ORIGIN = process.env.ALLOWED_ORIGIN || "https://yourdomain.com";
const EXPECTED_ORIGIN_VERIFY = process.env.ORIGIN_VERIFY_SECRET;
const EXPECTED_API_KEY = process.env.APP_API_KEY;

const MAX_TOTAL_PROMPT_CHARS = 4000;
const MAX_SYSTEM_CHARS = 1000;
const MAX_BODY_BYTES = 50 * 1024; // 50 KB

function safeCompare(a, b) {
  if (typeof a !== "string" || typeof b !== "string") return false;
  const bufA = Buffer.from(a);
  const bufB = Buffer.from(b);
  if (bufA.length !== bufB.length) return false;
  return crypto.timingSafeEqual(bufA, bufB);
}

export const handler = awslambda.streamifyResponse(
  async (event, responseStream, context) => {
    const headers = {
      "Content-Type": "text/event-stream; charset=utf-8",
      "Cache-Control": "no-cache, no-transform",
      "Connection": "keep-alive",
      "X-Accel-Buffering": "no",
      "Access-Control-Allow-Origin": ALLOWED_ORIGIN,
      "Access-Control-Allow-Methods": "POST, OPTIONS",
      "Access-Control-Allow-Headers": "Content-Type, Authorization, x-api-key, x-origin-verify",
    };

    if (event.requestContext?.http?.method === "OPTIONS") {
      const res = awslambda.HttpResponseStream.from(responseStream, { statusCode: 204, headers });
      res.end();
      return;
    }

    // 1. Verify CloudFront origin secret using constant-time comparison
    if (EXPECTED_ORIGIN_VERIFY) {
      const originHeader = event.headers?.["x-origin-verify"] || event.headers?.["X-Origin-Verify"];
      if (!safeCompare(originHeader, EXPECTED_ORIGIN_VERIFY)) {
        const res = awslambda.HttpResponseStream.from(responseStream, {
          statusCode: 403,
          headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": ALLOWED_ORIGIN },
        });
        res.write(JSON.stringify({ error: "Forbidden: Direct Function URL access blocked." }));
        res.end();
        return;
      }
    }

    // 2. Optional application-layer API key check (useful for service-to-service calls)
    if (EXPECTED_API_KEY) {
      const requestApiKey = event.headers?.["x-api-key"] || event.headers?.["X-Api-Key"];
      if (!safeCompare(requestApiKey, EXPECTED_API_KEY)) {
        const res = awslambda.HttpResponseStream.from(responseStream, {
          statusCode: 401,
          headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": ALLOWED_ORIGIN },
        });
        res.write(JSON.stringify({ error: "Unauthorized." }));
        res.end();
        return;
      }
    }

    // 3. Payload size check
    const rawBody = event.body || "";
    if (Buffer.byteLength(rawBody, "utf8") > MAX_BODY_BYTES) {
      const res = awslambda.HttpResponseStream.from(responseStream, {
        statusCode: 413,
        headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": ALLOWED_ORIGIN },
      });
      res.write(JSON.stringify({ error: "Payload exceeds 50 KB limit." }));
      res.end();
      return;
    }

    const stream = awslambda.HttpResponseStream.from(responseStream, {
      statusCode: 200,
      headers,
    });

    const sendSSE = (data, eventType = "message") => {
      const payload = typeof data === "string" ? data : JSON.stringify(data);
      stream.write(`event: ${eventType}\ndata: ${payload}\n\n`);
    };

    try {
      const body = event.body ? JSON.parse(event.body) : {};
      const prompt = String(body.prompt || "Explain distributed consensus in two sentences.");
      const systemPrompt = String(body.system || "");
      const conversationHistory = Array.isArray(body.messages) ? body.messages : [];
      const modelId = String(body.modelId || DEFAULT_MODEL_ID);

      if (systemPrompt.length > MAX_SYSTEM_CHARS) {
        sendSSE({ error: true, message: `System prompt exceeds ${MAX_SYSTEM_CHARS} char limit.` }, "error");
        stream.end();
        return;
      }

      // Check total characters across prompt string or conversation history
      let totalInputChars = 0;
      if (conversationHistory.length > 0) {
        for (const msg of conversationHistory) {
          if (Array.isArray(msg.content)) {
            for (const part of msg.content) {
              if (part.text) totalInputChars += String(part.text).length;
            }
          }
        }
      } else {
        totalInputChars = prompt.length;
      }

      if (totalInputChars > MAX_TOTAL_PROMPT_CHARS) {
        sendSSE({ error: true, message: `Input text (${totalInputChars} chars) exceeds ${MAX_TOTAL_PROMPT_CHARS} limit.` }, "error");
        stream.end();
        return;
      }

      if (!ALLOWED_MODELS.has(modelId)) {
        sendSSE({ error: true, message: `Model '${modelId}' is not allowed.` }, "error");
        stream.end();
        return;
      }

      sendSSE({ status: "connected", modelId }, "init");

      const messages = conversationHistory.length > 0
        ? conversationHistory
        : [{ role: "user", content: [{ text: prompt }] }];

      const command = new ConverseStreamCommand({
        modelId,
        messages,
        system: systemPrompt ? [{ text: systemPrompt }] : undefined,
        inferenceConfig: { maxTokens: 2048, temperature: 0.7 },
      });

      const bedrockResponse = await bedrock.send(command);

      let stopReason = null;
      let tokenUsage = null;

      for await (const chunk of bedrockResponse.stream) {
        if (chunk.contentBlockDelta?.delta?.text) {
          sendSSE({ text: chunk.contentBlockDelta.delta.text }, "delta");
        }
        if (chunk.messageStop) {
          stopReason = chunk.messageStop.stopReason;
        }
        if (chunk.metadata?.usage) {
          tokenUsage = chunk.metadata.usage;
        }
      }

      // Send completion event once both stopReason and usage have arrived
      sendSSE({ stopReason, usage: tokenUsage }, "done");
    } catch (err) {
      console.error("Bedrock stream error:", err);
      sendSSE({ error: true, message: err.message }, "error");
    } finally {
      stream.end();
    }
  }
);
Enter fullscreen mode Exit fullscreen mode

The Python Concurrency Trap: Event Loop Starvation

If you build this in Python with FastAPI and run it on Lambda via AWS Lambda Web Adapter (AWS_LWA_INVOKE_MODE=response_stream), there is a deceptive concurrency trap waiting for you.

The reason this catches so many developers off guard is that FastAPI looks and feels fully asynchronous. You declare async def chat_stream(...), you return a StreamingResponse, and you assume FastAPI handles concurrent users effortlessly.

The illusion breaks because Boto3's converse_stream isn't an asyncio library. Under the hood, botocore's EventStream wraps a synchronous blocking network socket. If you write:

# WARNING: This destroys concurrency across all users
for event in response.get("stream"):  # BLOCKS the main event loop thread!
    yield f"data: {event}\n\n"
Enter fullscreen mode Exit fullscreen mode

When you write for event in stream:, Python performs synchronous socket reads directly on the main thread running the asyncio event loop. Every time the worker waits 30–80ms for Bedrock to push the next token chunk, the entire event loop halts. Under concurrent load, all incoming requests freeze waiting on each other's network packets.

The clean solution is a classic producer-consumer pattern:

Offload the blocking Boto3 socket iterator to a background daemon thread, push events into an asyncio.Queue using loop.call_soon_threadsafe(), and let your FastAPI route await queue.get().

And make sure to handle client disconnects: if a user closes their tab mid-generation, your generator loop exits, but the background thread will keep happily burning Bedrock tokens until the model finishes unless you signal a threading.Event to abort it:

# lambda/python_adapter/main.py
import os
import json
import asyncio
import secrets
import threading
from typing import AsyncGenerator, Optional
import boto3
from fastapi import FastAPI, HTTPException, Request, Header
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field

app = FastAPI()
bedrock = boto3.client("bedrock-runtime", region_name="us-east-1")
_STREAM_END = object()

ALLOWED_MODELS = {
    "amazon.nova-pro-v1:0",
    "amazon.nova-lite-v1:0",
    "us.anthropic.claude-3-7-sonnet-20250219-v1:0",
    "us.anthropic.claude-3-5-haiku-20241022-v1:0",
}

EXPECTED_ORIGIN_VERIFY = os.getenv("ORIGIN_VERIFY_SECRET")

def safe_compare(val: Optional[str], expected: Optional[str]) -> bool:
    if not val or not expected:
        return False
    return secrets.compare_digest(val.strip(), expected.strip())

class ChatRequest(BaseModel):
    prompt: Optional[str] = Field(None, max_length=4000)
    system: Optional[str] = Field("You are a concise, accurate AI assistant.", max_length=1000)
    modelId: Optional[str] = Field("amazon.nova-pro-v1:0")

async def stream_bedrock_events(request: ChatRequest, http_req: Request) -> AsyncGenerator[str, None]:
    if request.modelId not in ALLOWED_MODELS:
        yield f"event: error\ndata: {json.dumps({'error': f'Model {request.modelId} not allowed.'})}\n\n"
        return

    queue = asyncio.Queue()
    loop = asyncio.get_running_loop()
    stop_event = threading.Event()

    def worker():
        try:
            response = bedrock.converse_stream(
                modelId=request.modelId,
                messages=[{"role": "user", "content": [{"text": request.prompt or "Hello"}]}],
            )
            for event in response.get("stream"):
                if stop_event.is_set():
                    break
                loop.call_soon_threadsafe(queue.put_nowait, event)
            loop.call_soon_threadsafe(queue.put_nowait, _STREAM_END)
        except Exception as e:
            loop.call_soon_threadsafe(queue.put_nowait, e)

    thread = threading.Thread(target=worker, daemon=True)
    thread.start()

    yield f"event: init\ndata: {json.dumps({'status': 'connected', 'modelId': request.modelId})}\n\n"

    stop_reason = None
    token_usage = None

    try:
        while True:
            # Check client disconnection to avoid paying for unread tokens
            if await http_req.is_disconnected():
                stop_event.set()
                break

            try:
                item = await asyncio.wait_for(queue.get(), timeout=0.5)
            except asyncio.TimeoutError:
                continue

            if item is _STREAM_END:
                break
            if isinstance(item, Exception):
                yield f"event: error\ndata: {json.dumps({'error': str(item)})}\n\n"
                break

            if "contentBlockDelta" in item:
                text = item["contentBlockDelta"]["delta"]["text"]
                yield f"event: delta\ndata: {json.dumps({'text': text})}\n\n"
            elif "messageStop" in item:
                stop_reason = item["messageStop"].get("stopReason")
            elif "metadata" in item:
                token_usage = item["metadata"].get("usage")

        if not stop_event.is_set():
            yield f"event: done\ndata: {json.dumps({'stopReason': stop_reason, 'usage': token_usage})}\n\n"
    finally:
        stop_event.set()

@app.post("/stream")
async def chat_endpoint(
    req: ChatRequest,
    http_req: Request,
    x_origin_verify: Optional[str] = Header(None, alias="x-origin-verify"),
):
    if EXPECTED_ORIGIN_VERIFY and not safe_compare(x_origin_verify, EXPECTED_ORIGIN_VERIFY):
        raise HTTPException(status_code=403, detail="Forbidden: Direct Function URL access blocked.")

    return StreamingResponse(
        stream_bedrock_events(req, http_req),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
    )
Enter fullscreen mode Exit fullscreen mode

The Practical Threat Model

When you put an LLM behind an open HTTP endpoint, the threat model is less about data breaches and more about financial denial-of-service. An attacker doesn't need to crack your encryption; they just need to flood your endpoint with requests asking for 4,000 output tokens each.

Here are the guardrails you need in place before deploying:

  • Scope your IAM execution role: Avoid Resource: "*". In fact, avoid foundation-model/* as well. Pin your policy strictly to the foundation model and inference profile ARNs you intend to support:
  - Effect: Allow
    Action:
      - bedrock:InvokeModelWithResponseStream
      - bedrock:ConverseStream
    Resource:
      - !Sub "arn:${AWS::Partition}:bedrock:${AWS::Region}::foundation-model/amazon.nova-pro-v1:0"
      - !Sub "arn:${AWS::Partition}:bedrock:${AWS::Region}::foundation-model/amazon.nova-lite-v1:0"
      - !Sub "arn:${AWS::Partition}:bedrock:${AWS::Region}::foundation-model/anthropic.claude-3-5-haiku-20241022-v1:0"
      - !Sub "arn:${AWS::Partition}:bedrock:${AWS::Region}::foundation-model/anthropic.claude-3-7-sonnet-20250219-v1:0"
      - !Sub "arn:${AWS::Partition}:bedrock:${AWS::Region}:${AWS::AccountId}:inference-profile/*"
      - !Sub "arn:${AWS::Partition}:bedrock:*:*:inference-profile/us.anthropic.claude-3-7-sonnet-20250219-v1:0"
      - !Sub "arn:${AWS::Partition}:bedrock:*:*:inference-profile/us.anthropic.claude-3-5-haiku-20241022-v1:0"
Enter fullscreen mode Exit fullscreen mode
  • Enforce an in-code allowlist: Never pass req.body.modelId straight to the SDK. Validate every request against an explicit set of approved models.
  • Cap input text & payloads: Tally character counts across prompts, conversation histories (4,000 chars), and system prompts (1,000 chars), and reject raw bodies over 50 KB before calling Bedrock.
  • Halt on client disconnect: As shown in the Python implementation, abort the stream worker if the client disconnects so you aren't paying for tokens no viewer will ever see.
  • WAF rate limiting: Attach AWS WAF to CloudFront with an IP-based rate limit (e.g. 100 requests per 5 minutes) to absorb automated scraping.

Bandwidth & Throughput Limits

Lambda response streaming delivers an initial 6 MB unthrottled burst, after which subsequent throughput is capped at 2 MB/s up to a maximum 200 MB response. For text-based LLM token streaming, that ceiling is orders of magnitude higher than the generation speed of current foundation models.


Wrap-Up

There is no single "right" serverless streaming pattern on AWS:

If your application already lives behind API Gateway REST APIs and depends on its usage plans, API keys, or Cognito integration, flip the integration transfer mode to STREAM and keep your existing infrastructure.

But if you want a direct, low-latency streaming pipeline for a web or mobile client without paying API Gateway invocation and data processing markup, Lambda Function URLs with CloudFront origin verification deliver a clean, fast, serverless streaming architecture.

All of the code—including working SAM templates, CDK v2 stack, Node and Python handlers, and local test scripts—is on GitHub:

🔗 GitHub Repository: sharma-the-karma/serverless-bedrock-token-streaming

Top comments (0)