Files
homelab-frontend/docs/API-sqs.md
T

263 lines
8.8 KiB
Markdown

# API — queue surface (`/sqs/*`)
Fronts the Kafka Management Service (`kmsvc`) in namespace `sqs`. SQS-shaped
message-plane API over Kafka.
Status marks:
**[LIVE]** verified against the running cluster and the committed proto on 2026-08-19.
**[SPEC]** the contract this gateway must implement; not built yet.
Source of truth for shapes:
`~/workplace/kmsvc-proto/proto/kafkamgmt/v1/queue_service.proto`.
---
## The important finding: a REST surface already exists [LIVE]
**Do not build gRPC-to-JSON transcoding.** `kmsvc-manage` already mounts grpc-gateway:
```go
mux := runtime.NewServeMux()
kafkamgmtv1.RegisterQueueServiceHandlerServer(ctx, mux, svc)
```
The upstream serves plain REST/JSON on **:8080** and plain gRPC on **:9090**. Neither
gRPC-Web nor server reflection is enabled.
So `/sqs/*` is a **path-stripping reverse proxy plus authentication**, not a protocol
translator. That makes it dramatically cheaper than the LLM surface.
```
api.riotpiao.com/sqs/v1/queues/{q}/messages
| strip /sqs, authenticate
v
management-service.sqs.svc.cluster.local:8080/v1/queues/{q}/messages
```
Upstream: Deployment `management-service`, 3 replicas, HPA 3-9, Service ClusterIP
`10.98.3.138`, ports `8080` (http) and `9090` (grpc).
---
## Endpoints [LIVE — HTTP annotations from the proto]
Six operations. All unary. No streaming, no subscribe.
| Method | Path (after `/sqs` strip) | RPC |
|---|---|---|
| POST | `/v1/queues/{queue_name}/messages` | `SendMessage` |
| POST | `/v1/queues/{queue_name}/messages:batch` | `SendMessageBatch` |
| GET | `/v1/queues/{queue_name}/messages` | `ReceiveMessage` |
| DELETE | `/v1/queues/{queue_name}/messages/{receipt_handle}` | `DeleteMessage` |
| POST | `/v1/queues/{queue_name}/messages:batchDelete` | `DeleteMessageBatch` |
| PATCH | `/v1/queues/{queue_name}/messages/{receipt_handle}` | `ChangeMessageVisibility` |
---
## Two wire-format traps [LIVE]
Both follow from grpc-gateway defaults, and both will surprise anyone who reads only
the proto.
**1. `bytes` fields are base64 in JSON.** `SendMessageRequest.message_body` and
`Message.body` are proto `bytes`. The JSONPB marshaler encodes them as base64 strings.
Sending raw text will not do what you expect.
**2. Field names are lowerCamelCase.** `cmd/server/main.go` calls bare
`runtime.NewServeMux()` with no marshaler options, so `OrigName` is false. The wire uses
`messageBody`, `receiptHandle`, `maxNumberOfMessages` — not the snake_case names in the
proto.
Document both prominently or every first-time caller loses an hour.
---
## Message shapes [LIVE — from the proto]
### Send
```
POST /sqs/v1/queues/agent-worker-queue/messages
{
"messageBody": "aGVsbG8gd29ybGQ=", // base64 of "hello world"
"messageAttributes": {"values": {"k": "v"}},
"messageGroupId": "", // FIFO only
"messageDeduplicationId": "", // FIFO only
"delaySeconds": 0 // 0-900
}
-> {"messageId": "...", "sequenceNumber": ""} // sequenceNumber FIFO only
```
### Receive — long poll
```
GET /sqs/v1/queues/agent-worker-queue/messages
?maxNumberOfMessages=10 // <= 10
&waitTimeSeconds=20 // 0-20
&visibilityTimeoutSeconds=30 // optional override
-> {"messages": [{
"messageId": "...",
"receiptHandle": "...",
"body": "aGVsbG8gd29ybGQ=",
"attributes": {"values": {}},
"receiveCount": 1,
"messageGroupId": "",
"enqueuedAt": "2026-08-19T16:29:07Z"
}]}
```
### Delete — the ack
```
DELETE /sqs/v1/queues/agent-worker-queue/messages/{receiptHandle}
-> {}
```
### Change visibility
```
PATCH /sqs/v1/queues/agent-worker-queue/messages/{receiptHandle}
{"visibilityTimeoutSeconds": 60} // 0-43200
-> {}
```
### Batch
Both batch calls take `entries[]` with a caller-assigned `id`, and return partial
success:
```json
{"successful": [{"id": "1", "messageId": "..."}],
"failed": [{"id": "2", "error": "..."}]}
```
A batch call can return 200 with entries in `failed`. Callers must inspect the body,
not just the status.
### Limits [LIVE — from the SDK]
`MaxMessageBodyBytes = 262144` (256 KiB), `MaxReceiveMessages = 10`,
`MaxWaitTimeSeconds = 20`.
---
## Semantics
At-least-once, SQS-style. Receive leases a message for the visibility timeout; the
caller must `DeleteMessage` to acknowledge. An un-deleted message reappears after the
timeout and `receiveCount` increments. After `maxReceiveCount` (default 5) it goes to
the DLQ if one is configured.
**Long-polling matters for the gateway.** `waitTimeSeconds` up to 20 means a `GET` can
legitimately hold open for 20 seconds returning nothing. Read timeouts must exceed that
comfortably, and a client disconnect must cancel upstream — the same requirement as the
LLM surface, for the same reason.
---
## Error mapping [SPEC]
The SDK maps gRPC codes to sentinel errors; grpc-gateway maps them to HTTP. Use this as
the gateway's status contract:
| gRPC code | HTTP | SDK sentinel |
|---|---|---|
| `NotFound` | 404 | `ErrQueueNotFound` |
| `AlreadyExists` | 409 | `ErrAlreadyExists` |
| `InvalidArgument` | 400 | `ErrInvalidArgument` |
| `Unauthenticated` | 401 | `ErrUnauthenticated` |
| `ResourceExhausted` | 429 | `ErrMessageTooLarge` |
Upstream errors arrive in the grpc-gateway envelope
`{"code": 5, "message": "Not Found", "details": []}`. Decide deliberately whether
`/sqs/*` passes that through or re-renders it as RFC 9457 to match `/v1/*`.
Recommendation: **pass through**, so the gateway does not become a second, subtly
different error vocabulary for the same upstream.
---
## Queue lifecycle is NOT in this API [LIVE]
There is no `CreateQueue`, `DeleteQueue`, or `ListQueues` RPC. The proto says so
explicitly:
```proto
// Queue lifecycle (create/delete/configure) is managed via the Queue CRD,
// not this service
```
Queues are Kubernetes resources — `queues.kmsvc.io/v1`, namespaced. `kmsvc-cli`'s
`create-queue` and `delete-queue` talk to the Kubernetes API, not to kmsvc.
**This is a hard boundary for the gateway.** Exposing queue creation over `/sqs/*` would
require the gateway to hold Kubernetes write credentials, which violates **G2**. Do not
add it. If declarative queue management ever needs a public surface, it belongs behind a
separate component with its own RBAC — not in the public edge process.
Queue spec fields, for reference when reading a queue's configuration:
`fifoQueue`, `isDLQ`, `deadLetterTargetQueue`, `delaySeconds` (0-900),
`maxReceiveCount` (default 5), `messageRetentionPeriodSeconds` (default 345600),
`visibilityTimeoutSeconds` (default 30), `minShards`, `maxShards` (default 8),
`partitionsPerShard` (default 6), `shardSplitThresholdBytesPerSec`,
`shardSplitCooldownSeconds`.
Kafka topics are named `kmsvc.{queue}.shard-{id}` and are created by `queue-operator`
directly via the Kafka Admin API — there are no `KafkaTopic` CRs.
Currently one queue exists: `agent-worker-queue` in namespace `sqs`, phase `Ready`,
1 shard.
---
## Authentication [SPEC]
`Authorization: Bearer <jwt>`, same as every other gateway surface.
**The upstream enforces nothing.** `kmsvc`'s auth interceptor exists but is never wired,
and the REST surface is mounted with the in-process grpc-gateway variant that bypasses
gRPC interceptors regardless. Both `:8080` and `:9090` are currently open, and
`kmsvc.riotpiao.com` is publicly routed.
The gateway is therefore the only authentication boundary for this surface.
---
## Out of scope
- **Workflow start.** Nothing in kmsvc starts a Temporal workflow — no such RPC exists,
and grep for `ExecuteWorkflow`/`StartWorkflow` across `kmsvc-manage`, `kmsvc-sdk` and
`kmsvc-cli` returns nothing. A caller dials `temporal-frontend.temporal.svc:7233`
with a Temporal SDK directly. A `/workflow/*` surface is net-new code, not a proxy
route — see [task 7.3](../tasks/7.3-workflow-prefix.md).
- **DLQ operations.** `kmsvc-cli`'s `dlq peek` and `dlq redrive` are client-side
compositions of the six RPCs, not server operations. Redrive is a non-atomic
Receive-Send-Delete. If `/sqs/*` should offer redrive, that is new logic with real
failure modes, not a proxied call.
- **Kafka direct access.** No external listener exists; the bootstrap
`kmsvc-kafka-bootstrap.sqs.svc.cluster.local:9092` is cluster-internal only. The
gateway proxies kmsvc, never Kafka.
---
## Examples
```bash
Q=agent-worker-queue
# send (body must be base64)
curl -s -X POST https://api.riotpiao.com/sqs/v1/queues/$Q/messages \
-H 'content-type: application/json' \
-d "{\"messageBody\":\"$(printf 'hello world' | base64)\"}"
# receive, long poll 20s
curl -s "https://api.riotpiao.com/sqs/v1/queues/$Q/messages?maxNumberOfMessages=10&waitTimeSeconds=20"
# acknowledge
curl -s -X DELETE https://api.riotpiao.com/sqs/v1/queues/$Q/messages/$RECEIPT
# extend the lease
curl -s -X PATCH https://api.riotpiao.com/sqs/v1/queues/$Q/messages/$RECEIPT \
-H 'content-type: application/json' -d '{"visibilityTimeoutSeconds":60}'
```