--- name: ak-cloud-deploy description: > Deploy an Agent Kernel project to AWS, Azure, or GCP using Terraform modules, or to any Kubernetes cluster (on-prem, baremetal, EKS) using the official Helm chart. Supports serverless and containerized modes for all three clouds. AWS supports execution modes (rest_sync, rest_async, async, stream), queue-based scalable processing, custom API Gateway authorizers, EventBridge Scheduler for scheduled tasks, and SSM Parameter Store secret resolution (`ssm_enabled`). GCP supports Cloud Run serverless (scale-to-zero) and containerized (always-on) with Redis or Firestore session backends. Kubernetes runs the queue pipeline over NATS JetStream, Kafka, or SQS with KEDA autoscaling and an optional WebSocket gateway tier. license: Apache-2.0 metadata: author: yaalalabs version: "0.9.3" category: user --- # Deploy to Cloud Use this skill to deploy your Agent Kernel project to AWS, Azure, or GCP. ## Instructions for the Agent When the user wants cloud deployment, follow this workflow. ### Step 1: Identify the Project Check for: - `pyproject.toml` with `agentkernel` dependency - agent entry file (`app.py`, `lambda.py`, or similar) - `config.yaml` If missing, suggest `ak-init` first. ### Step 2: Ask Deployment Questions 1. Target platform: AWS, Azure, GCP, or Kubernetes (on-prem / baremetal / EKS via the Helm chart; see the On-Prem / Kubernetes section)? 2. Runtime mode: - Serverless - Containerized 3. Execution pattern (AWS only): - Synchronous HTTP (`rest_sync`, supports standard or queue/scalable mode; AWS serverless or containerized) - Asynchronous REST (`rest_async`, queue/scalable mode; AWS serverless or containerized) - WebSocket full-response (`async`, works with or without queue mode — `queue_mode = false` runs the agent inline, `queue_mode = true` enqueues to a separately-scalable Agent Runner) — AWS serverless or AWS containerized/ECS - WebSocket token streaming (`stream`, works with or without queue mode, same as `async` above) — AWS serverless or AWS containerized/ECS; also available via SSE (`POST /api/v1/chat` with `execution.mode: stream`, no Terraform changes required) wherever the built-in FastAPI REST server runs: local/self-hosted, AWS ECS single-container REST, Azure Container Apps, GCP Cloud Run — not on AWS Lambda or Azure Functions, which use WebSocket instead 4. Scalability (AWS serverless only): standard or queue/scalable mode? 5. Session store: Redis, Valkey (AWS only), DynamoDB (AWS), Cosmos DB (Azure), Firestore (GCP)? 6. Security: custom authorizer required (AWS serverless only)? 7. Resource naming: a single `prefix` applied to every resource name (e.g. `myapp-dev-chat`). ### Step 3: Choose the Correct Terraform Module Use official modules: - AWS serverless: `yaalalabs/ak-serverless/aws` - AWS containerized: `yaalalabs/ak-containerized/aws` - Azure serverless: `yaalalabs/ak-serverless/azurerm` - Azure containerized: `yaalalabs/ak-containerized/azurerm` - GCP serverless: `yaalalabs/ak-serverless/google` - GCP containerized: `yaalalabs/ak-containerized/google` Use current module version (`0.9.3`) unless user requests another. Kubernetes does not use Terraform: the Helm chart lives at `ak-deployment/ak-k8s/chart` in the Agent Kernel repository and is published as an OCI artifact (`oci://ghcr.io/yaalalabs/charts/agent-kernel`, versioned with each release). Install the published chart, as in the Kubernetes section below, not a repository checkout. All modules are provider-agnostic: they declare `required_providers` but do not configure them internally. Configure each provider (`aws`/`docker`, `azurerm`, or `google`/`google-beta`/`docker`) in the root module and pass it explicitly via the module's `providers = { ... }` argument, as shown in the examples below. Azure's containerized module builds and pushes its image via a nested submodule with its own internal `docker` provider, so no `docker` provider needs to be configured or passed by the caller there. AWS-only features in this skill: - `execution_mode` - `queue_mode` - `authorizer` Azure examples use the provider module's standard top-level inputs (`function_name`, `package_path`, `container_port`, `publisher_email`) rather than AWS serverless execution modes. GCP modules use Cloud Run for both serverless (scale-to-zero, `min_instance_count=0`) and containerized (always-on, `min_instance_count≥1`) modes. Both use `agentkernel.gcp.CloudRun` as the entry point. ### Step 4: Align App Dependencies and Session Config When the user selects a session store, always update both app dependencies and `config.yaml`. #### Redis sessions - Dependency extras: include `redis` - Typical dependency example: ```toml dependencies = [ "agentkernel[openai,api,redis]>=0.9.3" ] ``` - `config.yaml` session block: ```yaml session: type: redis namespace: chat redis: host: ${REDIS_HOST} port: 6379 db: 0 ``` #### Valkey sessions (AWS) [Valkey](https://valkey.io/) is the open-source, Linux Foundation-governed fork of Redis. It is wire-compatible with Redis and available on AWS ElastiCache at a lower price point than the Redis OSS engine. Agent Kernel treats it as a first-class session and response store backend on AWS. - Dependency extras: include `valkey` - Typical dependency example: ```toml dependencies = [ "agentkernel[openai,api,aws,valkey]>=0.9.3" ] ``` - `config.yaml` session block: ```yaml session: type: valkey valkey: url: ${AK_SESSION__VALKEY__URL} # valkeys:// for SSL ttl: 604800 prefix: "ak:sessions:" ``` - Terraform provisions the cluster with `create_valkey_cluster = true` (serverless and containerized) and injects `AK_SESSION__VALKEY__URL` — `session.type: valkey` must still be set in `config.yaml` since the type itself is never injected. - For queue-mode response storage, set `create_valkey_response_store = true` (serverless only) and `execution.response_store.type: valkey` in `config.yaml`. #### DynamoDB sessions (AWS) - Dependency extras: include `aws` - Typical dependency example: ```toml dependencies = [ "agentkernel[openai,api,aws]>=0.9.3" ] ``` - `config.yaml` session block: ```yaml session: type: dynamodb dynamodb: table_name: ${DYNAMODB_MEMORY_TABLE} region: ${AWS_REGION} ``` #### Cosmos DB sessions (Azure) - Dependency extras: include `azure` - Typical dependency example: ```toml dependencies = [ "agentkernel[openai,api,azure]>=0.9.3" ] ``` - `config.yaml` session block: ```yaml session: type: cosmosdb cosmosdb: endpoint: ${AZURE_COSMOS_ENDPOINT} database_name: ${AZURE_COSMOS_DATABASE} container_name: ${AZURE_COSMOS_CONTAINER} ``` #### Firestore sessions (GCP) - Dependency extras: include `gcp` - Typical dependency example: ```toml dependencies = [ "agentkernel[openai,api,gcp]>=0.9.3" ] ``` - `config.yaml` session block: ```yaml session: type: firestore firestore: collection_name: ${AK_SESSION__FIRESTORE__COLLECTION_NAME} project_id: ${PROJECT_ID} ttl: 604800 ``` - The Terraform module injects `AK_SESSION__TYPE=firestore` and `AK_SESSION__FIRESTORE__COLLECTION_NAME` automatically when Firestore is enabled. - A TTL policy must be set on the Firestore collection pointing to the `expiry_time` field for automatic document expiry. If the Terraform module creates the backing store (for example `create_redis_cluster = true`, `create_valkey_cluster = true`, or `create_dynamodb_memory_table = true`), make sure output values are wired to the app environment variables used by `config.yaml`. ### Deploying Conversation Thread Storage Conversation threads (`ak-add-capabilities`) deploy the same way sessions do: the app mounts `AgentThreadRequestHandler` (which is what enables the feature; the self-hosted REST API is the only surface that records threads — queue-mode runners and the other deployment adapters do not) and declares `thread.type` in `config.yaml`, and Terraform provisions the backend, injecting only the connection detail. Terraform never sets `AK_THREAD__TYPE` itself. - **AWS (serverless + containerized)**: `create_dynamodb_thread_table = true` provisions a DynamoDB table (partition `session_id`, sort `sk`, TTL on `expiry_time`) and injects `AK_THREAD__DYNAMODB__TABLE_NAME`. - **GCP (serverless + containerized)**: `create_firestore_thread_collection = true` requires `create_firestore_database = true`, reuses that database, and injects `AK_THREAD__FIRESTORE__COLLECTION_NAME`, `__PROJECT_ID`, and `__DATABASE_ID`. No new resource is provisioned — the collection is created on first write. - **Redis / Valkey**: no dedicated Terraform flag — reuse whatever `create_redis_cluster` / `create_valkey_cluster` already provisions and declare `thread: {type: redis}` (or `valkey`) with the cluster URL in `config.yaml`. Setting the Terraform flag *without* declaring `thread.type` in `config.yaml` silently leaves threads on the non-durable in-memory backend (any `AK_THREAD__*` var materialises `AKConfig.thread`, but `type` still defaults to `in_memory`) — always pair the flag with the matching `thread.type`. ### Deploying Scheduled Tasks Scheduled tasks (`ak-add-capabilities`) deploy on the same two-step split: the app declares `schedule.provider.type` and `schedule.store.type` in `config.yaml`, and Terraform provisions the backends, injecting only the coordinates. Terraform never sets `AK_SCHEDULE__PROVIDER__TYPE` or `AK_SCHEDULE__STORE__TYPE`. **Scheduling requires queue mode** — occurrences are delivered into the input queue, so `queue_mode = true` is mandatory on AWS. `enable_scheduling` without it is rejected. - **AWS (serverless + containerized)**: `enable_scheduling = true` provisions an EventBridge Scheduler schedule group and the execution role Scheduler assumes to deliver triggers to the input queue, grants both the request/REST role and the agent-runner role `scheduler:CreateSchedule|UpdateSchedule|DeleteSchedule|GetSchedule` plus `iam:PassRole` on that role, and injects `AK_SCHEDULE__PROVIDER__EVENTBRIDGE__GROUP_NAME`, `__ROLE_ARN` and `__QUEUE_ARN`. It also flips the input queue to content-based deduplication, because Scheduler cannot set a `MessageDeduplicationId` (application senders are unaffected — their explicit id takes precedence). - **AWS (serverless + containerized)**: `create_dynamodb_schedule_table = true` provisions a DynamoDB table (partition `task_id`, no sort key, no GSI, TTL on `expiry_time`) and injects `AK_SCHEDULE__STORE__DYNAMODB__TABLE_NAME`. - **Redis / Valkey task store**: no dedicated Terraform flag — reuse whatever `create_redis_cluster` / `create_valkey_cluster` already provisions and declare `schedule: {store: {type: redis}}` (or `valkey`) with the cluster URL in `config.yaml`. - **Azure / GCP**: no scheduling provider ships for these clouds yet. The `local` provider only works in a single always-on process with the `in_memory` transport and store, which no cloud deployment mode gives you. ```hcl queue_mode = true enable_scheduling = true create_dynamodb_schedule_table = true ``` ```yaml schedule: provider: type: eventbridge store: type: dynamodb ``` Setting the Terraform flags *without* declaring both types in `config.yaml` silently leaves scheduling on the in-process `local` provider and `in_memory` store (any `AK_SCHEDULE__*` var materialises `AKConfig.schedule`, but the types still default), while the provisioned group and table sit unused — always pair the flags with the matching types. The reverse is safe: declaring the types without the flags fails at startup with an `AKConfigError` on the missing `group_name` / `role_arn` / `queue_arn`. The blocks above enable deferring and the agent tools. The `/api/v1/schedules` management routes are **not** mounted from config — on ECS queue mode, pass the handler to the IO container's entrypoint, where it is served alongside the queue-producing chat route: ```python # app_rest_service.py from agentkernel.aws import ECSIOHandler from agentkernel.schedule import ScheduleRESTRequestHandler if __name__ == "__main__": # add authoriser=... to scope listings to the caller; without one the routes are open ECSIOHandler.run(handlers=[ScheduleRESTRequestHandler()]) ``` Remember to list the schedule paths in `gateway_endpoints` — the gateway only proxies paths it is told about. New outputs on both AWS stacks (null unless the matching flag is set): `schedule_group_name`, `schedule_group_arn`, `scheduler_execution_role_arn`, `schedule_table_name`, `schedule_table_arn`. ### Deploying Secret Resolution (AWS SSM) Secret resolution (`ak-add-capabilities`) follows the same split: the app declares `secret.provider.type: aws_ssm` in `config.yaml` and resolves keys with `SecretManager.current().get("OPENAI_API_KEY")`; Terraform grants read access and injects only the deployment scope. Terraform never sets `AK_SECRET__PROVIDER__TYPE`. - **AWS (serverless + containerized)**: `ssm_enabled = true` grants only the tier that runs the agent (the agent runner when `queue_mode = true`; otherwise the serverless request handler / containerized REST service) `ssm:GetParameter` on `arn:aws:ssm:::parameter/ak//*` and injects `AK_SECRET__PREFIX = ` (the module's own `prefix`). The response handler and WebSocket connection handler never get SSM access, so resolve secrets in agent-runner code only. - **Terraform does not create the parameters** — create each one yourself as a `SecureString` (AWS-managed `alias/aws/ssm` key) at `/ak//`, e.g. `aws ssm put-parameter --name /ak//openai_api_key --type SecureString --value sk-...`. - Drop the key from `environment_variables`: a set, non-empty environment variable always wins over SSM, so leaving `OPENAI_API_KEY` there bypasses the store. - The app needs the `aws` extra (`agentkernel[aws]`) for the `aws_ssm` provider. ```hcl ssm_enabled = true ``` ```yaml secret: provider: type: aws_ssm # prefix is injected by Terraform as AK_SECRET__PREFIX ``` ```python from agentkernel.secret import SecretManager from agents import set_default_openai_key set_default_openai_key(SecretManager.current().get("OPENAI_API_KEY")) ``` Declaring `aws_ssm` without `ssm_enabled` (and without setting `secret.prefix` yourself) fails at startup with an `AKConfigError` on the missing prefix; with a prefix but no IAM grant, lookups raise `SecretError` (AccessDenied). Setting `ssm_enabled` without declaring the provider leaves resolution on the `env` provider. See `examples/aws-serverless/openai` and `examples/aws-containerized/openai-dynamodb-scalable`. ## AWS Serverless (Lambda + API Gateway) ### Agent Code Pattern Use Lambda handler entrypoint: ```python from agentkernel.aws import Lambda from agentkernel.openai import OpenAIModule # Register agents with your framework module OpenAIModule([...]) handler = Lambda.handler ``` ### A) Basic Synchronous REST (`rest_sync`) Use this for straightforward request/response APIs. This is the single-Lambda pattern: use `request_handler` plus any `gateway_endpoints`. Do not add `agent_runner` or `response_handler` unless you are using queue mode. ```hcl module "serverless_agents" { source = "yaalalabs/ak-serverless/aws" version = "0.9.3" prefix = var.prefix product_display_name = "AK Serverless" region = var.region execution_mode = "rest_sync" request_handler = { function_name = "chat-handler" function_description = "AK request handler" handler_path = "lambda.handler" package_type = "Image" package_path = "../dist" timeout = 45 memory_size = 256 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } gateway_endpoints = [ { path = "app", method = "GET" }, { path = "app_info", method = "POST" } ] } ``` ### B) Scalable Queue Mode (`rest_sync` or `rest_async`) Use this for high throughput and long-running requests. This is the multi-artifact pattern used by the scalable example: a request handler zip, an agent-runner image, and a response-handler zip. Each Lambda can use one of three `package_type` values: | `package_type` | Artifact source | Required field(s) | What Terraform does | | -------------- | ---------------------------------- | ---------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | `LocalZip` | Local ZIP file or source directory | `package_path` | Uses the Lambda module's local packaging support. No shared source bucket is managed by this module. | | `S3Zip` | ZIP artifact stored in S3 | **Either** `package_path` **or** `lambda_package_s3` | If `package_path` is provided, this module creates/uses a shared source bucket, uploads the ZIP, and deploys from S3. If `lambda_package_s3` (`{ bucket, key, version_id? }`) is provided, Terraform uses the existing S3 object directly; set `version_id` (from a versioned bucket) so re-uploading changed code redeploys the function. | | `Image` | Container image in ECR | **Either** `ecr_image_uri` **or** `package_path` | If `ecr_image_uri` is provided, Terraform uses the existing image. If `package_path` is provided, Terraform builds and pushes an image to ECR and deploys it. | `package_path` and `lambda_package_s3`/`ecr_image_uri` are mutually exclusive — set only one. **Development / local build** (build artifacts locally before running `terraform apply`): ```hcl module "serverless_agents" { source = "yaalalabs/ak-serverless/aws" version = "0.9.3" prefix = var.prefix region = var.region product_display_name = "AK Scalable REST" queue_mode = true execution_mode = "rest_sync" # can also be "rest_async" create_dynamodb_memory_table = true create_dynamodb_response_store = true # Uncomment with `schedule: {provider: {type: eventbridge}, store: {type: dynamodb}}` in config.yaml # enable_scheduling = true # create_dynamodb_schedule_table = true request_handler = { function_name = "request-handler" function_description = "Receives REST requests" handler_path = "lambda_request_handler.handler" package_type = "LocalZip" package_path = "../dist_request_handler.zip" timeout = 45 memory_size = 256 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } agent_runner = { function_name = "agent-runner" function_description = "Processes queued requests" handler_path = "lambda_agent_runner.handler" package_type = "Image" package_path = "../dist_agent_runner" timeout = 45 memory_size = 512 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } response_handler = { function_name = "response-handler" function_description = "Handles async response completion" handler_path = "lambda_response_handler.handler" package_type = "LocalZip" package_path = "../dist_response_handler.zip" timeout = 45 memory_size = 256 } queue_config = { input_queue_visibility_timeout = 60 input_queue_max_receive_count = 3 input_queue_create_dlq = false input_queue_message_retention_seconds = 300 output_queue_visibility_timeout = 60 output_queue_max_receive_count = 3 output_queue_create_dlq = false output_queue_message_retention_seconds = 300 batch_size = 10 maximum_batching_window_in_seconds = 0 } } ``` **Production: external artifact sources (S3 / ECR)** Build and push artifacts in CI/CD, then point Terraform at them so `terraform apply` contains no local build step: ```hcl request_handler = { function_name = "request-handler" handler_path = "lambda_request_handler.handler" package_type = "S3Zip" lambda_package_s3 = { bucket = "my-lambda-packages-bucket" key = "dist_request_handler.zip" version_id = "" # from your versioned bucket → redeploys work (#548) } timeout = 45 memory_size = 256 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } agent_runner = { function_name = "agent-runner" handler_path = "lambda_agent_runner.handler" package_type = "Image" ecr_image_uri = "123456789012.dkr.ecr.us-west-2.amazonaws.com/agent-runner:latest" timeout = 45 memory_size = 512 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } response_handler = { function_name = "response-handler" handler_path = "lambda_response_handler.handler" package_type = "S3Zip" lambda_package_s3 = { bucket = "my-lambda-packages-bucket" key = "dist_response_handler.zip" version_id = "" # from your versioned bucket → redeploys work (#548) } timeout = 45 memory_size = 256 } ``` See [examples/aws-serverless/scalable-openai](https://github.com/yaalalabs/agent-kernel/tree/develop/examples/aws-serverless/scalable-openai) for a complete working example of this pattern. Terraform only replaces a Lambda's code when `s3_bucket`, `s3_key`, `s3_object_version`, or `source_code_hash` changes. If you re-upload to the **same** S3 key in an **unversioned** bucket, none of those change and the function keeps running the old code. To make `S3Zip` redeploys reliable when pointing at your own artifact (`lambda_package_s3`): enable versioning on the bucket holding the ZIPs, then pass the uploaded object's version through `lambda_package_s3.version_id`. **Queue mode `config.yaml`** (bundled into every Lambda package — `execution.mode`, queue URLs, table names, and `max_receive_count` are all injected automatically by Terraform as environment variables; only set values that are NOT injected): ```yaml # For rest_sync or rest_async execution: response_store: type: dynamodb # or redis / valkey — not injected, must be set here retry_count: 5 delay: 5 session: type: dynamodb # or redis / valkey — not injected, must be set here ``` - `rest_sync`: request handler sends to queue, polls the response store until the result is available, then returns it synchronously. - `rest_async`: request handler returns `{status: "ACCEPTED", request_id}` immediately; the client polls `GET /api/{version}/{endpoint}` with `{"request_id": ""}` in the JSON body. - SQS visibility timeout for each queue **must be >= the corresponding Lambda timeout**. **Required `pyproject.toml` extras for queue mode**: ```toml dependencies = [ "agentkernel[openai,api,aws]>=0.9.3" # include 'redis' if using Redis, or 'valkey' if using Valkey session/response store ] ``` ### C) WebSocket Async (`async`) Use this for realtime bidirectional interactions. This follows the current websocket example shape: the request handler stays on the module's top-level Lambda inputs, then queue workers and WebSocket handlers are added via nested blocks. ```hcl module "serverless_agents" { source = "yaalalabs/ak-serverless/aws" version = "0.9.3" prefix = var.prefix region = var.region product_display_name = "AK WebSocket Serverless Example" queue_mode = true execution_mode = "async" create_redis_cluster = true create_dynamodb_response_store = true request_handler = { function_name = "request-handler" function_description = "Receives WebSocket requests" handler_path = "lambda_request_handler.handler" package_type = "LocalZip" package_path = "../dist_request_handler.zip" timeout = 45 memory_size = 256 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } agent_runner = { function_name = "agent-runner" function_description = "Processes chat tasks" handler_path = "lambda_agent_runner.handler" package_type = "Image" package_path = "../dist_agent_runner" timeout = 45 memory_size = 512 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } response_handler = { function_name = "response-handler" function_description = "Sends responses to connections" handler_path = "lambda_response_handler.handler" package_type = "LocalZip" package_path = "../dist_response_handler.zip" timeout = 45 memory_size = 256 } ws_connection_handler = { function_name = "ws-connection-handler" function_description = "Handles $connect/$disconnect" handler_path = "lambda_ws_connection_handler.handler" package_path = "../dist_ws_connection_handler.zip" timeout = 45 memory_size = 256 } ws_routes = [ { route = "app" }, { route = "app_info" } ] } ``` **WebSocket `lambda_ws_connection_handler.py`** — implement `AuthValidator` to validate the bearer token on `$connect`: ```python import jwt import os from agentkernel.auth import AuthValidator, ValidationResult from agentkernel.aws import WebsocketConnectionHandler class MyAuthValidator(AuthValidator): def validate(self, token: str) -> ValidationResult: try: payload = jwt.decode( token, os.environ["JWT_SECRET"], algorithms=["HS256"], issuer=os.environ["JWT_ISSUER"], audience=os.environ["JWT_AUDIENCE"], ) user_id = payload.get("userId", "") if user_id: return ValidationResult(is_valid=True, claims={"userId": user_id}) return ValidationResult(is_valid=False, error_msg="userId claim missing") except Exception as e: return ValidationResult(is_valid=False, error_msg=str(e)) handler = WebsocketConnectionHandler.set_auth_validator(MyAuthValidator()).handler ``` **WebSocket `lambda_request_handler.py`** — routes are registered by **route name** (no leading slash, no HTTP method) because the WebSocket API uses `$request.body.route` for dispatch: ```python import json from agentkernel.aws import Lambda @Lambda.register("app") # WebSocket route name, no slash/method def app_handler(event, context): return {"response": "Hello from app route"} @Lambda.register("app_info") def app_info_handler(event, context): payload = json.loads(event.get("body") or "{}") return {"response": "Hello from app_info route"} handler = Lambda.handler ``` Contrast with REST mode where routes use a leading slash and HTTP method: ```python @Lambda.register("/app", method="GET") # REST mode: path + method ``` **WebSocket `config.yaml`** (bundled into every Lambda package — `execution.mode`, `websocket_api.chat_route`, queue URLs, and connection table name are all injected automatically by Terraform; only set values that are NOT injected): ```yaml session: type: redis # or valkey / dynamodb — not injected, must be set here redis: prefix: "ak:myapp:" ``` - Authentication is **mandatory** for WebSocket mode. `AuthValidator.validate()` must return `claims["userId"]`; connections without a valid token are rejected at `$connect`. - The `authorizer` Terraform block is not supported in WebSocket mode — leave it unset. - Only `LocalZip` is supported for `ws_connection_handler.package_path`. - `ws_routes` custom routes must have a matching `@Lambda.register("route_name")` entry. `ws_chat_route` is the built-in AI chat route — it is handled internally by the framework and does **not** need a `@Lambda.register` entry. **Required `pyproject.toml` extras for WebSocket mode**: ```toml dependencies = [ "agentkernel[openai,api,aws,redis,auth]>=0.9.3" ] ``` ### D) WebSocket Token Streaming (`stream`) Use this when the client should receive each stream event as soon as it is produced, instead of waiting for the full response. Same Terraform shape as WebSocket Async (`request_handler`, `agent_runner`, `response_handler`, `ws_connection_handler`, `ws_routes`) — only `execution_mode` changes: ```hcl module "serverless_agents" { source = "yaalalabs/ak-serverless/aws" version = "0.9.3" prefix = var.prefix region = var.region product_display_name = "AK Streaming WebSocket Example" queue_mode = true execution_mode = "stream" create_redis_cluster = true # create_redis_response_store, create_valkey_response_store, and create_dynamodb_response_store must stay false/unset: # WebSocket modes (async/stream) push responses over the connection and Terraform # validation fails if a response store is enabled for them. request_handler = { ... } # same shape as async mode agent_runner = { ... } # runs ServerlessStreamAgentRunner, streams chunks to output queue response_handler = { ... } # broadcasts each chunk as a STREAM_CHUNK message ws_connection_handler = { ... } ws_routes = [ { route = "app" }, { route = "app_info" } ] } ``` **`config.yaml`** — the only required setting beyond WebSocket async mode is the execution mode itself: ```yaml execution: mode: stream ``` This makes `ServerlessAgentRunner.handle()` dispatch to `ServerlessStreamAgentRunner` (queue mode) — no code change needed in the agent runner Lambda beyond the standard `Lambda.handler` entrypoint. **Message format** — clients receive a sequence of `STREAM_CHUNK` messages instead of one `CHAT_RESPONSE`: ```json {"type": "STREAM_CHUNK", "delta": "Hello", "done": false, "session_id": "user-1"} {"type": "STREAM_CHUNK", "delta": " world", "done": false, "session_id": "user-1"} {"type": "STREAM_CHUNK", "delta": "!", "done": true, "session_id": "user-1"} ``` On an unrecoverable error, the final chunk carries `error` instead of `delta`, with `done: true`. - Queue disabled (`queue_mode = false`): the request handler Lambda streams tokens directly to the WebSocket client without SQS. - `create_redis_response_store` / `create_valkey_response_store` / `create_dynamodb_response_store` must be `false` for `stream` (same constraint as `async`) — Terraform validation enforces this. - See [examples/aws-serverless/streaming-openai](https://github.com/yaalalabs/agent-kernel/tree/develop/examples/aws-serverless/streaming-openai) for a complete working example. **Containerized / direct streaming (no Terraform WebSocket setup)**: wherever the built-in FastAPI REST server runs — local/self-hosted, AWS ECS single-container REST, Azure Container Apps, GCP Cloud Run — SSE event streaming can be enabled by setting `execution.mode: stream` in `config.yaml`. `POST /api/v1/chat` and `/api/v1/chat-multipart` then return `text/event-stream` responses instead of JSON — no queue or WebSocket infrastructure is required for this mode. Not available on AWS Lambda or Azure Functions (serverless), which use WebSocket for streaming instead. ### E) API Gateway Custom Authorizer (AWS) If the user needs token verification, include `authorizer` block: ```hcl authorizer = { description = "API Gateway Lambda Authorizer" function_name = "gateway-authorizer" handler_path = "lambda_auth.handler" package_path = "../dist_auth.zip" package_type = "LocalZip" result_ttl_in_seconds = 0 environment_variables = { SOME_OTHER_KEY = "Some Other Value" } } ``` ## AWS Containerized (ECS/Fargate) ### A) Basic Mode (single container, direct execution) ```python from agentkernel.api import RESTAPI from agentkernel.openai import OpenAIModule OpenAIModule([...]) if __name__ == "__main__": RESTAPI.run() ``` ```hcl module "containerized_agents" { source = "yaalalabs/ak-containerized/aws" version = "0.9.3" prefix = var.prefix region = var.region product_display_name = "AK ECS Deployment" rest_service = { package_path = "../dist" container_port = 8000 desired_count = 2 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } create_dynamodb_memory_table = true } ``` ### B) Scalable Queue Mode (two containers) Use this for high-throughput or long-running agents. Two separate ECS services share SQS queues: - **IO container** (`ECSIOHandler`) — Thread 1: FastAPI REST API; Thread 2: output queue consumer → DynamoDB / WebSocket. Thread 2 is actually `queues.output.no_of_consumers` (default 5) parallel polling threads. - **Agent Runner container** (`ECSAgentRunner`) — runs `queues.input.no_of_consumers` (default 5) parallel threads, each polling the input queue, executing the agent, and putting the result on the output queue. **`app_rest_service.py`** (IO container entrypoint — NO agent definitions here): ```python from agentkernel.aws import ECSIOHandler runner = ECSIOHandler.run if __name__ == "__main__": runner() ``` Pass `auth_validator=MyAuthValidator()` to secure the REST routes: `ECSIOHandler.run` binds it onto every REST route via `AWSRestAPI.add_auth_handlers` (the same Bearer-token mechanism as `RESTAPI.add_auth_handlers` in Basic Mode); omit it to leave the REST routes unauthenticated. **`app_agent_runner.py`** (Agent Runner container entrypoint): ```python from agentkernel.aws import ECSAgentRunner from agentkernel.openai import OpenAIModule OpenAIModule([...]) # register agents here only handler = ECSAgentRunner.run if __name__ == "__main__": handler() ``` **`config.yaml`** (same file included in both images — queue URLs and table names are injected by Terraform; `batch_size` is Terraform-only, never set here): ```yaml execution: queues: type: sqs # mandatory in a declared queues block: it is what selects the transport input: no_of_consumers: 5 # parallel input-poll threads in the Agent Runner container (default 5) output: no_of_consumers: 5 # parallel output-poll threads in the IO container (default 5) response_store: type: dynamodb retry_count: 30 delay: 2 session: type: dynamodb # or redis / valkey (create_redis_cluster / create_valkey_cluster) ``` **Terraform:** ```hcl module "containerized_agents" { source = "yaalalabs/ak-containerized/aws" version = "0.9.3" prefix = var.prefix region = var.region rest_service = { package_path = "../dist-rest-service" cpu = 512 memory = 1024 desired_count = 2 command = ["python", "app_rest_service.py"] environment_variables = { OPENAI_API_KEY = var.openai_api_key } } queue_mode = true execution_mode = "rest_sync" # or "rest_async" # Uncomment with `schedule: {provider: {type: eventbridge}, store: {type: dynamodb}}` in config.yaml # enable_scheduling = true # create_dynamodb_schedule_table = true queue_config = { input_queue_visibility_timeout = 120 output_queue_visibility_timeout = 60 input_queue_create_dlq = true output_queue_create_dlq = true } # Agent Runner container — separate image with agent definitions agent_runner = { cpu = 1024 memory = 2048 desired_count = 1 package_path = "../dist-agent-runner" command = ["python", "app_agent_runner.py"] environment_variables = { OPENAI_API_KEY = var.openai_api_key } } # Optional: auto-scale Agent Runner based on Input Queue depth scaling_config = { enabled = true min_count = 1 max_count = 10 backlog_target = 5 scale_in_cooldown = 180 scale_out_cooldown = 60 } create_dynamodb_memory_table = true } ``` **Required `pyproject.toml` extras:** ```toml dependencies = [ "agentkernel[openai,api,aws]>=0.9.3" ] ``` **Key rules:** - Agent definitions (`OpenAIModule([...])`) go in `app_agent_runner.py` only — never in `app_rest_service.py`. - `ECSIOHandler` starts two threads via `ThreadRunner`; if the output-consumer pool crashes, sibling consumer threads finish their in-flight message first (graceful drain via a shared `shutdown_event`), then the container exits (`os._exit(1)`) so ECS can restart it. The REST API thread doesn't participate in the drain — it's just terminated at that point. - `ECSAgentRunner` and `ECSOutputConsumer` both extend `ECSSQSConsumer` (itself a `RawQueueConsumer`, the same base `LambdaSQSConsumer` extends); extend either class to customise message processing. ### C) WebSocket Mode (`async` / `stream`) A WebSocket API Gateway proxies frames to the ECS REST service via a VPC Link V1 + internal NLB in front of the existing ALB. Supports both **direct** (`queue_mode = false`, one ECS service, agent runs inline) and **queue** (`queue_mode = true`, same two-container split as section B, with the REST/IO service enqueueing chat frames and its output-queue consumer pushing replies back over the socket) variants. `execution_mode = "async"` delivers the full reply as one `CHAT_RESPONSE` push; `execution_mode = "stream"` delivers it token-by-token as a sequence of `STREAM_CHUNK` pushes (terminated by a chunk with `"done": true`) — in queue mode `ECSAgentRunner.run()` dispatches to `ECSStreamAgentRunner`, which fans out one Output Queue message per chunk instead of one for the full reply; in direct mode the chat route runs `ChatService.process_stream_chat_async()` and broadcasts each chunk inline. See `examples/aws-containerized/openai-stream` (direct mode) or `examples/aws-containerized/openai-stream-queue-mode` (queue mode) for full streaming examples. **`app.py`** (direct/single-container variant; register custom routes before calling `run()`): ```python from agentkernel.aws import AWSWebsocketAPI from agentkernel.auth import AuthValidator, ValidationContext, ValidationResult from agentkernel.openai import OpenAIModule from typing import Optional OpenAIModule([...]) class CustomAuthValidator(AuthValidator): def validate(self, token: str, context: Optional[ValidationContext] = None) -> ValidationResult: # userId claim is mandatory — it's how a connection maps to a user for reply routing ... return ValidationResult(is_valid=True, claims={"userId": "userId-claim-value"}) @AWSWebsocketAPI.register("status") # bare route name, no leading slash, no HTTP verb async def status(ctx: dict) -> dict: # ctx = {"message": , "user_id": ...} — no BaseRequest schema imposed return {"status": "OK", "user_id": ctx["user_id"]} # dict return -> broadcast as SYSTEM_RESPONSE; None -> no broadcast if __name__ == "__main__": # Direct mode calls AWSWebsocketAPI directly — NOT ECSIOHandler, which always starts a # second thread polling an output queue that doesn't exist in direct (non-queue) mode. AWSWebsocketAPI.set_auth_handler(auth_validator=CustomAuthValidator()).run() ``` Queue mode is different: it splits into `app_rest_service.py` (`ECSIOHandler.run(auth_validator=...)` — correct here, since `ECSIOHandler` also starts the output-queue consumer thread — custom routes registered here, no agent definitions) and `app_agent_runner.py` (`ECSAgentRunner.run`, `OpenAIModule([...])` here only) — identical shape to section B. **Terraform:** ```hcl module "containerized_agents" { source = "yaalalabs/ak-containerized/aws" version = "0.9.3" providers = { aws = aws, docker = docker } prefix = var.prefix region = var.region vpc_id = var.vpc_id private_subnet_ids = var.private_subnet_ids rest_service = { package_path = "../dist" container_port = 8000 environment_variables = { OPENAI_API_KEY = var.openai_api_key } } queue_mode = false # or true for the two-container queue variant (see section B) execution_mode = "async" # "stream" delivers the reply as STREAM_CHUNK messages instead ws_chat_route = "chat" # optional, defaults to "chat" ws_routes = [ # every @AWSWebsocketAPI.register(...) route must be listed here too { route = "status" }, ] create_dynamodb_memory_table = true } ``` **Key rules:** - Auth is **mandatory**: `AuthValidator.validate()` must resolve a `userId` claim, or the framework raises at `$connect` construction time. There is no API Gateway authorizer for WebSocket mode. - Every custom route needs both a Python `@AWSWebsocketAPI.register("name")` decorator **and** a matching Terraform `ws_routes = [{ route = "name" }]` entry — Python cannot create the API Gateway route/integration, so both sides must declare it. `ws_chat_route` is Terraform-only — it just names the API Gateway route key that forwards to the container's hardcoded `/ws/chat` endpoint; the container never reads it from config, and renaming it needs no Python change. - The DynamoDB response store is simply never created in WebSocket mode (no `response_store` input to set) — replies are always pushed over the connection instead. `gateway_endpoints` specifically has a Terraform validation rule rejecting it in `async`/`stream` modes. - Client wire format: `{"route": "chat", "body": {...}}` — the WebSocket API's `route_selection_expression` is `$request.body.route`. ## Azure Serverless (Functions + APIM) ### Agent Code Pattern ```python from agentkernel.azure import AzureFunctions from agentkernel.openai import OpenAIModule OpenAIModule([...]) handler = AzureFunctions.handler ``` ### Terraform Example ```hcl module "serverless_agents" { source = "yaalalabs/ak-serverless/azurerm" version = "0.9.3" prefix = var.prefix region = var.region resource_group_name = var.resource_group_name publisher_email = var.publisher_email product_display_name = "AK Azure Functions Example" function_name = "openai-agents" function_description = "Agent Kernel OpenAI Sample Azure Function" module_type = "python" package_path = "../dist.zip" create_redis_cluster = true create_cosmosdb_cluster = false environment_variables = { OPENAI_API_KEY = var.openai_api_key } gateway_endpoints = [ { function_name = "AgentFunction" path = "/chat" method = "POST" }, { function_name = "CustomFunction" path = "/custom" method = "POST" } ] } ``` ## Azure Containerized (Container Apps + APIM) ### Terraform Example ```hcl module "containerized_agents" { source = "yaalalabs/ak-containerized/azurerm" version = "0.9.3" prefix = var.prefix region = var.region resource_group_name = var.resource_group_name publisher_email = var.publisher_email package_path = "../dist" product_display_name = "AK Azure Container Apps" container_port = 8000 create_cosmosdb_cluster = true environment_variables = { OPENAI_API_KEY = var.openai_api_key } } ``` ## GCP Serverless (Cloud Run — scale-to-zero) ### Agent Code Pattern ```python from agentkernel.gcp import CloudRun from agentkernel.openai import OpenAIModule OpenAIModule([...]) @CloudRun.register("/app", method="GET") def app_handler() -> dict: return {"status": "ok"} def main() -> None: CloudRun.run() ``` ### Terraform Example (with Redis sessions) ```hcl module "serverless_agent" { source = "yaalalabs/ak-serverless/google" version = "0.9.3" providers = { google = google, google-beta = google-beta, docker = docker } project_id = var.project_id region = var.region prefix = var.prefix product_display_name = "AK GCP Serverless" package_path = "${path.module}/../dist" create_redis_cluster = true environment_variables = { OPENAI_API_KEY = var.openai_api_key } gateway_endpoints = [ { path = "app", method = "GET", overwrite_path = "/app" }, { path = "app_info", method = "POST", overwrite_path = "/app_info" } ] } ``` ### Terraform Example (with Firestore sessions) ```hcl module "serverless_agent" { source = "yaalalabs/ak-serverless/google" version = "0.9.3" providers = { google = google, google-beta = google-beta, docker = docker } project_id = var.project_id region = var.region prefix = var.prefix product_display_name = "AK GCP Serverless Firestore" package_path = "${path.module}/../dist" create_firestore_db = true environment_variables = { OPENAI_API_KEY = var.openai_api_key } } ``` The module injects `AK_SESSION__TYPE=firestore` and `AK_SESSION__FIRESTORE__COLLECTION_NAME` automatically when `create_firestore_db = true`. **Required `pyproject.toml` extras:** ```toml dependencies = [ "agentkernel[openai,api,gcp]>=0.9.3" # for Firestore sessions # or: "agentkernel[openai,api,redis]>=0.9.3" # for Redis sessions ] ``` ## GCP Containerized (Cloud Run — always-on) ### Agent Code Pattern Same as GCP Serverless — use `agentkernel.gcp.CloudRun`: ```python from agentkernel.gcp import CloudRun from agentkernel.openai import OpenAIModule OpenAIModule([...]) def main() -> None: CloudRun.run() ``` ### Terraform Example ```hcl module "containerized_agent" { source = "yaalalabs/ak-containerized/google" version = "0.9.3" providers = { google = google, google-beta = google-beta, docker = docker } project_id = var.project_id region = var.region prefix = var.prefix product_display_name = "AK GCP Containerized" package_path = "${path.module}/../dist" min_instance_count = 1 # always-on — no cold starts container_port = 8000 create_redis_cluster = true environment_variables = { OPENAI_API_KEY = var.openai_api_key } } ``` **Key difference from GCP Serverless:** `min_instance_count = 1` keeps at least one instance running at all times, eliminating cold starts. ## On-Prem / Kubernetes (Helm Chart) Deploys the queue pipeline to any Kubernetes cluster: an `io-handler` Deployment (REST API + Response Handler), an `agent-runner` Deployment (the consumers executing the agents), and an optional `ws-gateway` Deployment for WebSocket modes. Backing services (Valkey, NATS) are condition-gated dependencies of the chart; flavors (dev, baremetal, EKS) are values files. Reference: the chart at `ak-deployment/ak-k8s/` and the end-to-end example at `examples/k8s/openai-queue-mode/` in the Agent Kernel repository. **Entry files** (one image per component; the chart's `command` defaults match these names): ```python # app_io_handler.py: the io-handler image from agentkernel.pipeline import IOHandler def main(): IOHandler.run() if __name__ == "__main__": main() ``` ```python # app_agent_runner.py: the agent-runner image (registers the agent modules) from agentkernel.openai import OpenAIModule from agentkernel.pipeline import AgentRunner from agents import Agent # ... define agents ... OpenAIModule([...]) def main(): AgentRunner.run() if __name__ == "__main__": main() ``` ```python # app_ws_gateway.py: only for async/stream modes from agentkernel.pipeline import WebSocketGateway def main(): WebSocketGateway.run(auth_validator=YourAuthValidator()) # claims must include 'userId' if __name__ == "__main__": main() ``` **Config split:** the image's `config.yaml` declares WHAT runs (`execution.mode`, agents, logging); the chart injects WHERE it runs (broker, response store, session store) as `AK_*` env vars that override matching `config.yaml` fields. Point `execution.queues.type` at `nats` (recommended on-prem), `kafka`, or `sqs`, and use a shared response/session store (`valkey`) on any multi-pod topology. **Images:** `python:3.12-slim` base, dependencies staged with pip's `--target` layout, one Dockerfile per component (mirror `examples/k8s/openai-queue-mode/deploy/package.sh`, which also cross-installs Linux wheels so builds work from macOS). **Install:** ```bash helm pull oci://ghcr.io/yaalalabs/charts/agent-kernel --version 0.9.3 --untar # unpacks the flavor values files helm install ak oci://ghcr.io/yaalalabs/charts/agent-kernel --version 0.9.3 \ -f agent-kernel/values-dev.yaml \ --set ioHandler.image.repository= \ --set agentRunner.image.repository= --set image.tag= \ --set 'extraEnv[0].name=OPENAI_API_KEY' \ --set 'extraEnv[0].valueFrom.secretKeyRef.name=openai' \ --set 'extraEnv[0].valueFrom.secretKeyRef.key=api-key' ``` The chart version tracks the Agent Kernel release (the pin above is the current one); the flavor values files ship inside the chart, so no repository checkout is needed. `ak-deployment/ak-k8s/chart` in the repository is for developing the chart itself. - `values-dev.yaml`: micro-clusters (k3d/kind/microk8s/k3s), single replicas, auto-provisioned JetStream, port-forward entry; also documents the single-process profile (one pod, `in_memory` transport, zero backing services). - `values-baremetal.yaml`: Envoy Gateway + MetalLB + cert-manager, NACK-managed JetStream objects (`autoProvision: false`), OpenEBS hostpath storage. - `values-eks.yaml`: AWS Load Balancer Controller gateway classes, EBS gp3, Pod Identity; `sqs`, `kafka`, and `nats` all valid. - WebSocket modes: `--set execution.mode=stream --set wsGateway.enabled=true` plus a `wsGateway.auth.token`; the gateway needs a shared session store for its connection table. - Autoscaling: `keda.enabled=true` scales the agent-runner on queue depth (KEDA prerequisite). **Verify:** ```bash kubectl port-forward service/ak-agent-kernel-io 8000:80 curl -s -X POST http://localhost:8000/api/v1/chat \ -H 'Content-Type: application/json' \ -d '{"prompt": "Hello", "session_id": "s1", "agent": ""}' ``` **Teardown:** `helm uninstall ak` (cluster prerequisites like Gateway API CRDs, KEDA, NACK, and Strimzi are installed per cluster, never by the chart, so they remain). ## Config and Packaging Notes ### Packaging For AWS async/scalable deployments, build and package separately: - **request handler** — `LocalZip`, `S3Zip`, or `Image` - **agent runner** — `LocalZip`, `S3Zip`, or `Image` (typically `Image` for heavier dependencies) - **response handler** — `LocalZip`, `S3Zip`, or `Image` - **ws-connection-handler** — `LocalZip` only (WebSocket mode; no `package_type` field, only `package_path`) Azure serverless uses a single zip `package_path` for the Functions deployment. Azure containerized uses `package_path` as the Docker build context for Container Apps. ### config.yaml per mode Each Lambda package must include a `config.yaml`. Most runtime values (`execution.mode`, queue URLs, table names, `max_receive_count`, `websocket_api.chat_route`) are **injected automatically by Terraform as environment variables** and do not need to be set in `config.yaml`. Only the values in the table below must be configured manually: | Config key | Required in | Notes | |------------|------------|-------| | `execution.response_store.type` | request-handler, response-handler (queue modes) | `redis`, `valkey`, or `dynamodb` — not injected | | `execution.response_store.retry_count` | request-handler (`rest_sync`) | how many times to poll for response | | `execution.response_store.delay` | request-handler (`rest_sync`) | seconds between polls | | `session.type` | agent-runner | `redis`, `valkey`, or `dynamodb` — not injected | | `session.redis.prefix` | agent-runner (Redis sessions) | key namespace prefix | | `session.valkey.prefix` | agent-runner (Valkey sessions) | key namespace prefix | ### Timeouts and visibility - SQS `input_queue_visibility_timeout` must be **>= agent runner Lambda timeout**. - SQS `output_queue_visibility_timeout` must be **>= response handler Lambda timeout**. - FIFO queues are used by default; `message_group_id` maps to session ID. ## Prerequisites Checklist ### AWS - AWS CLI configured (`aws configure`) - Terraform >= 1.9.5 - S3 bucket for remote state (optional but recommended) - DynamoDB table for state lock (optional but recommended) - Container registry access (if `package_type = "Image"`) ### Azure - Azure CLI logged in (`az login`) - Terraform >= 1.9.5 - Resource group and subscription permissions - Storage for remote Terraform state (optional but recommended) ### GCP - GCP CLI configured (`gcloud auth application-default login`) - Terraform >= 1.9.5 - Docker installed (for local image building and push to Artifact Registry) - An existing GCP project with appropriate permissions - Required APIs enabled: Cloud Run, API Gateway, Artifact Registry, VPC Access ## Deploy ```bash cd deploy terraform init terraform plan terraform apply ``` Kubernetes deployments use `helm install` / `helm upgrade` instead; see the On-Prem / Kubernetes section above. ## Teardown ```bash cd deploy terraform destroy ``` Kubernetes: `helm uninstall `. ## What to Do Next - Add capabilities (`ak-add-capabilities`) before production rollout. - Add tests (`ak-test`) against deployed endpoints. - Add integrations (`ak-add-integration`) for Slack/WhatsApp/Telegram/Teams. - Iterate on tools and agents (`ak-build`) and re-run `terraform apply` (or `helm upgrade`).