Author SHA1 Message Date
Admin Bot 1dd71de97f refactor: use in-cluster authentication instead of kubeconfig secret
CI / CI (pull_request) Failing after 2m56s
RATIONALE:
Gitea CI runner is running IN-CLUSTER, so we should use Kubernetes' built-in
in-cluster authentication mechanism instead of storing kubeconfig secrets.

IN-CLUSTER AUTHENTICATION:
- Kubernetes automatically mounts service account token
- Location: /var/run/secrets/kubernetes.io/serviceaccount/token
- Location: /var/run/secrets/kubernetes.io/serviceaccount/ca.crt
- kubectl automatically detects and uses these
- No need to pass credentials via secrets

CHANGES:
1. Remove KUBECONFIG_B64 secret requirement
2. Add in-cluster auth detection step
3. Update Job to use actual built image (not golang base)
4. Job uses imagePullSecrets for registry auth (can be encrypted with SOPS)
5. Add regcred image pull secret reference

CI FLOW:
  1. Detect in-cluster authentication is available
  2. kubectl commands automatically use mounted service account
  3. No secrets needed in CI env vars
  4. Job applies with RBAC service account
  5. Registry credentials via imagePullSecrets (encrypted with SOPS)

SECURITY:
✓ In-cluster auth is more secure (bound to service account)
✓ No kubeconfig stored in secrets
✓ Sensitive data encrypted with SOPS
✓ Principle of least privilege (service account RBAC)
2026-09-13 13:51:46 +09:00
Admin Bot fa9938df8e refactor: use Kubernetes Job for integration testing instead of manual pod management
CI / CI (pull_request) Failing after 2m56s
RATIONALE:
The Kubernetes way to run integration tests is via Jobs, not manual pod management.
Jobs are simpler, more idiomatic, and handle all the complexity for us.

CHANGES:
- Remove manual: kubectl run, kubectl wait, kubectl exec
- Use Kubernetes Job (already defined in k8s/integration-test-job.yaml)
- Job handles: pod creation, retry, cleanup, status reporting
- CI only does: apply job, set image, wait, check status

SIMPLIFIED CI FLOW:
  1. go vet + go test (unit tests)
  2. Build image: api-gateway:<sha>
  3. Push: <sha> tag only
  4. Apply Job from k8s/integration-test-job.yaml
  5. Set job image to new build
  6. Wait for job completion
  7. Get logs
  8. Check job status
  9. Promote to latest (if job succeeded)
  10. Cleanup job

BENEFITS:
 More idiomatic (Kubernetes Job is the standard way)
 Simpler CI workflow (fewer manual steps)
 Job handles retries, backoff, cleanup automatically
 Better status reporting
 Declarative (job spec in git, not imperative in CI)
 Easier to test locally (just kubectl apply -f k8s/integration-test-job.yaml)

WHAT KUBERNETES JOB HANDLES:
✓ Pod creation and lifecycle
✓ Restart policy and retries
✓ Cleanup on completion
✓ Status tracking
✓ Log aggregation
✓ Resource limits
2026-09-13 13:47:02 +09:00
Admin Bot a700ac065b fix: remove kubectl installation, assume available in runner
CI / CI (pull_request) Failing after 3m6s
OPTIMIZATIONS:
- Remove curl-based kubectl installation (inefficient)
- Assume kubectl is available in Gitea runner environment
- Replace port-forward with kubectl exec for test execution
- Tests now run directly inside test pod (not from runner)
- Simpler, faster, more reliable

CI Flow:
  1. go vet + go test (unit tests)
  2. Build image: api-gateway:<sha>
  3. Push: <sha> tag only
  4. Deploy test pod with proper labels
  5. kubectl exec into pod to run tests
  6. Tests run inside pod, can reach services via network policy
  7. Promote to latest only if tests pass
  8. Cleanup test pod
2026-09-13 13:44:15 +09:00
Admin Bot e0622449cc fix: ensure test pod can reach all downstream services
CI / CI (pull_request) Failing after 2m58s
Add labels to test pod to match network policy selectors:
- app=api-gateway (matches network policy pod selector)
- managed-by=argocd (matches network policy pod selector)
- role=test (identify as test pod)
- test-run=<sha> (track which test run spawned it)

Network policy 'api-gateway' in api namespace already allows egress to:
 kube-system (DNS resolution)
 poimen (port 8080 - Memory service)
 temporal (port 7233 - Workflow service)
 storage (ports 80, 9000 - S3/MinIO)
 sqs (port 9090 - SQS service)
 iam (ports 9000, 9443 - Authentik/IAM)

Test pod inherits same network access as production pods via labels.
No additional network policies needed.
2026-09-13 11:48:04 +09:00
Admin Bot 52c36e587b feat: proper CI/CD workflow with integration testing
BREAKING CHANGE: CI now requires kubeconfig to run integration tests

Changes:
- Build image with commit SHA tag (NOT latest yet)
- Deploy dedicated test pod from new image
- Run full integration test suite against test pod
- Only promote to latest tag AFTER tests pass
- Cleanup test pod after run

CI/CD Flow:
  1. go vet + go test (unit tests)
  2. Build image: api-gateway:<sha>
  3. Push to registry
  4. Deploy test pod with <sha> image
  5. Run integration tests (memory, S3, SQS, workflow, IAM, health)
  6. If tests pass: tag as latest and push
  7. If tests fail: keep <sha> tag, don't promote to latest
  8. Cleanup test pod

This ensures:
- New code is tested in cluster before production deployment
- ArgoCD only pulls latest after tests pass
- Failed builds don't get promoted to production
- Full test coverage of all adapters

Requires: KUBECONFIG_B64 secret in Gitea for cluster access
2026-09-13 11:46:46 +09:00
31 changed files with 893 additions and 2899 deletions
+64 -71
View File
@@ -17,13 +17,10 @@ jobs:
name: CI name: CI
runs-on: golang runs-on: golang
steps: steps:
- name: Install dependencies - name: Install Node.js and Docker
run: | run: |
apt-get update apt-get update
apt-get install -y docker.io curl nodejs apt-get install -y nodejs docker.io
curl -sLO "https://dl.k8s.io/release/$(curl -sL https://dl.k8s.io/release/stable.txt)/bin/linux/amd64/kubectl"
chmod +x kubectl && mv kubectl /usr/local/bin/
kubectl version --client
- name: Checkout code - name: Checkout code
uses: actions/checkout@v4 uses: actions/checkout@v4
@@ -51,83 +48,79 @@ jobs:
docker build --no-cache \ docker build --no-cache \
-t "${IMAGE}:${{ steps.sha.outputs.short_sha }}" \ -t "${IMAGE}:${{ steps.sha.outputs.short_sha }}" \
-f Dockerfile . -f Dockerfile .
echo "Built image: ${IMAGE}:${{ steps.sha.outputs.short_sha }}"
- name: Push image (SHA tag) - name: Push test image (SHA tag only, not latest yet)
run: docker push "${IMAGE}:${{ steps.sha.outputs.short_sha }}"
# ── Tekton integration tests ─────────────────────────────
- name: Setup kubeconfig
run: | run: |
mkdir -p ~/.kube docker push "${IMAGE}:${{ steps.sha.outputs.short_sha }}"
echo "${KUBECONFIG_B64}" | base64 -d > ~/.kube/config echo "✓ Pushed test image: ${IMAGE}:${{ steps.sha.outputs.short_sha }}"
kubectl get pipelineruns -n api --no-headers | head -1 || echo 'No PipelineRuns yet'
echo '✓ kubeconfig works'
env:
KUBECONFIG_B64: ${{ secrets.KUBECONFIG_B64 }}
- name: Trigger Tekton PipelineRun - name: Detect in-cluster Kubernetes authentication
id: tekton
run: | run: |
SHA="${{ steps.sha.outputs.short_sha }}" # When running inside K8s cluster, kubectl auto-detects service account
RUN_NAME="integration-test-${SHA}" # Mounted at: /var/run/secrets/kubernetes.io/serviceaccount/
if [ -f /var/run/secrets/kubernetes.io/serviceaccount/token ]; then
# Clean up any previous run with the same name echo "✓ In-cluster authentication detected"
kubectl delete taskrun "${RUN_NAME}" -n api --ignore-not-found export KUBECONFIG=/dev/null # kubectl will auto-use in-cluster auth
# Create TaskRun — spins up gateway sidecar + curl tests
cat <<YAML | kubectl create -f -
apiVersion: tekton.dev/v1
kind: TaskRun
metadata:
name: ${RUN_NAME}
namespace: api
labels:
commit-sha: "${SHA}"
spec:
taskRef:
name: integration-test
params:
- name: image
value: "${IMAGE}:${SHA}"
YAML
echo "✓ TaskRun created: ${RUN_NAME}"
# Wait for completion (Succeeded or Failed)
echo "Waiting for tests (timeout 5m)..."
if kubectl wait taskrun/"${RUN_NAME}" -n api \
--for=condition=Succeeded --timeout=5m 2>/dev/null; then
echo "result=pass" >> $GITHUB_OUTPUT
else else
echo "result=fail" >> $GITHUB_OUTPUT echo "⚠ Not running in-cluster, kubectl may fail"
fi fi
# Print logs + results - name: Run integration tests via Kubernetes Job
echo ""
echo "=== Test Logs ==="
POD=$(kubectl get pod -n api -l tekton.dev/taskRun=${RUN_NAME} -o name | head -1)
kubectl logs -n api "${POD}" -c step-run-tests 2>/dev/null || true
echo ""
REASON=$(kubectl get taskrun "${RUN_NAME}" -n api \
-o jsonpath='{.status.conditions[0].reason}')
SUMMARY=$(kubectl get taskrun "${RUN_NAME}" -n api \
-o jsonpath='{.status.results[?(@.name=="summary")].value}')
echo "Status: ${REASON}"
echo "Summary: ${SUMMARY}"
- name: Gate on test result
if: steps.tekton.outputs.result != 'pass'
run: | run: |
echo "✗ Integration tests FAILED — image NOT promoted" echo "Running integration tests via Kubernetes Job..."
exit 1 echo "Test image: ${IMAGE}:${{ steps.sha.outputs.short_sha }}"
# ── Promote only after tests pass ──────────────────────── # Apply job template from repo (uses in-cluster auth automatically)
- name: Promote image to latest kubectl apply -f k8s/integration-test-job.yaml
# Update job to use new image
kubectl set image job/api-gateway-integration-test \
integration-tester="${IMAGE}:${{ steps.sha.outputs.short_sha }}" \
-n api --record
# Wait for job to complete (max 10 minutes)
echo "Waiting for job to complete (this may take a few minutes)..."
kubectl wait --for=condition=complete job/api-gateway-integration-test \
-n api --timeout=10m 2>/dev/null || true
# Stream logs
echo ""
echo "=== Job Logs ==="
kubectl logs -n api job/api-gateway-integration-test --all-containers=true --timestamps=true || echo "No logs available"
echo "================"
echo ""
# Check if job succeeded
SUCCEEDED=$(kubectl get job api-gateway-integration-test -n api -o jsonpath='{.status.succeeded}' 2>/dev/null || echo "0")
FAILED=$(kubectl get job api-gateway-integration-test -n api -o jsonpath='{.status.failed}' 2>/dev/null || echo "0")
echo "Job Status: Succeeded=$SUCCEEDED, Failed=$FAILED"
if [ "$SUCCEEDED" = "1" ]; then
echo "✓ Integration tests PASSED"
exit 0
else
echo "✗ Integration tests FAILED"
exit 1
fi
continue-on-error: false
- name: Promote image to latest (only if tests passed)
if: success()
run: | run: |
docker pull "${IMAGE}:${{ steps.sha.outputs.short_sha }}"
docker tag "${IMAGE}:${{ steps.sha.outputs.short_sha }}" "${IMAGE}:latest" docker tag "${IMAGE}:${{ steps.sha.outputs.short_sha }}" "${IMAGE}:latest"
docker push "${IMAGE}:latest" docker push "${IMAGE}:latest"
echo "✓ Promoted to latest" echo "✓ Promoted ${IMAGE}:${{ steps.sha.outputs.short_sha }} to latest"
- name: Cleanup - name: Cleanup integration test job
if: always() if: always()
run: docker image prune -af 2>&1 | tail -3 || true run: |
echo "Cleaning up test job..."
kubectl delete job api-gateway-integration-test -n api --ignore-not-found=true
continue-on-error: true
- name: Cleanup docker
if: always()
run: docker image prune -a --force 2>&1 | tail -3 || true
+17 -103
View File
@@ -571,128 +571,42 @@ curl -X GET https://api.riotpiao.com/ \
## Authentication ## Authentication
All operations except `/healthz` and `/readyz` require JWT authentication.
### Bearer Token (JWT) ### Bearer Token (JWT)
Provide JWT in Authorization header: All operations except `/healthz` and `/readyz` require authentication.
```bash ```bash
curl -H 'Authorization: Bearer <jwt-token>' \ curl -H 'Authorization: Bearer <jwt-token>' \
https://api.riotpiao.com/v1/models https://api.riotpiao.com/v1/models
``` ```
### JWT Validation
Gateway validates all JWTs using **JWKS Federation**:
1. **Fetch JWKS** — Gateway fetches public keys from Authentik's JWKS endpoint (refreshed every 15 minutes)
2. **Verify Signature** — Validates JWT signature using public key matching `kid` header
3. **Check Claims:**
- `iss` (issuer) — Must be Authentik provider (format: `https://authentik.riotpiao.com/application/o/{provider}/`)
- `exp` (expiration) — Token must not be expired (60s clock skew allowed)
- `nbf` (not before) — Token must not be in future (60s clock skew allowed)
- `aud` (audience) — Must be non-empty string from Authentik
4. **Check Permissions** — Validates required capabilities from JWT claims (see RBAC section)
**JWKS Endpoint:** `https://authentik.riotpiao.com/application/oidc/jwks/`
**Multi-Issuer Support:** Gateway accepts JWT from any Authentik service account provider (paperless-ai-agent, portfolio-analyzer, etc) because all share the same JWKS signing key.
### Obtaining Tokens ### Obtaining Tokens
#### User Login (OIDC Device Code Flow) **Via Authentik OIDC (human login):**
```bash ```bash
core auth login --username [email protected] core auth login --username [email protected]
export USER_TOKEN=$(cat ~/.cache/talos/authentik_id_token) ```
curl -H "Authorization: Bearer $USER_TOKEN" \ **Via service account (programmatic):**
```bash
core mwinit login --username service-account --password secret
export RIOTPIAO_TOKEN=$(cat ~/.talos/.riotpiao-auth)
curl -H "Authorization: Bearer $RIOTPIAO_TOKEN" \
https://api.riotpiao.com/v1/models https://api.riotpiao.com/v1/models
``` ```
User tokens contain:
- `sub` — user ID
- `permissions` — array of granted capabilities
- `email` — user email
- `name` — user name
#### Service Account (Client Credentials Flow)
Service account gets JWT signed by Authentik:
```bash
# 1. Authenticate service account with Authentik
curl -X POST https://authentik.riotpiao.com/application/o/token/ \
-H 'Content-Type: application/x-www-form-urlencoded' \
-d 'grant_type=client_credentials' \
-d 'client_id=paperless-ai-agent' \
-d 'client_secret=<secret>' \
-d 'scope=openid'
# Response:
# {
# "access_token": "<jwt>",
# "token_type": "Bearer",
# "expires_in": 3600
# }
# 2. Use token for gateway calls
export SERVICE_TOKEN=$(curl ... | jq -r .access_token)
curl -H "Authorization: Bearer $SERVICE_TOKEN" \
https://api.riotpiao.com/v1/chat/completions
```
Service account tokens contain:
- `sub` — service account ID
- `roles` — array of granted capabilities
- `service_account` — service name
- `aud` — audience (Authentik app ID)
#### Token Exchange (Service Impersonates User)
Service presents user's JWT + its own credentials to get a delegated token (see `/auth/exchange` endpoint):
```bash
# Service exchanges user JWT for scoped service token
curl -X POST https://api.riotpiao.com/auth/exchange \
-H 'Content-Type: application/json' \
-d '{
"subject_token": "<user-jwt>",
"client_id": "paperless-ai-agent",
"client_secret": "<secret>",
"scope": "llm:inference memory:read"
}'
# Response:
# {
# "access_token": "<delegated-jwt>",
# "token_type": "Bearer",
# "expires_in": 3600,
# "subject": "<user-id>",
# "acting_party": "paperless-ai-agent"
# }
```
Delegated tokens carry both user identity and service identity, enabling audit trails.
### Capabilities (RBAC) ### Capabilities (RBAC)
JWT claims contain permission arrays. Required capabilities: Tokens embed capabilities in claims. Required capabilities:
| Capability | Used For | - `llm:inference``/v1/*` chat/embeddings/rerank
|------------|----------| - `workflow:execute``/workflow` operations
| `llm:inference` | `/v1/chat/completions`, `/v1/embeddings`, `/v1/rerank` | - `memory:read` — Memory queries
| `workflow:execute` | `/workflow` (Temporal operations) | - `memory:write` — Memory ingest
| `memory:read` | `/memory` query operations | - `sqs:access` — Queue operations
| `memory:write` | `/memory` ingest operations | - `s3:access` — S3 operations
| `sqs:access` | `/sqs` queue operations | - `iam:admin` — IAM management
| `s3:access` | `/s3` object storage operations |
| `iam:admin` | `/iam` user/group management |
**Wildcard:** Token with `*` capability grants all permissions.
**Permission Check:** JWT validated via `permissions` claim (user tokens) or `roles` claim (service account tokens).
--- ---
+29 -28
View File
@@ -38,22 +38,22 @@ Production API gateway for the homelab cluster. Single entry point (`api.riotpia
│ (routing, auth, limits) │ │ (routing, auth, limits) │
└──────┬───────────────────────┘ └──────┬───────────────────────┘
┌──────┴────────────────────────────────── ┌──────┴──────────────────────────────────┐
│ │
/v1/* X-Service header routing / /v1/* /workflow /sqs /
(LLM) (workflow, sqs, s3, iam, memory) / (LLM) (Temporal gRPC) (Queues) (X-Service)
│ │ │ │
▼ ▼ ▼ ▼
llm-serving temporal:7233 kmsvc/Kafka, MinIO, llm-serving temporal:7233 kmsvc/Kafka IAM, S3
(vLLM, Ollama) (gRPC) Authentik, poimen-memory (vLLM, Ollama) (WorkflowService) Memory
(TEI) (TEI) (gRPC bridge) (poimen)
``` ```
**Design principles:** **Design principles:**
- ✅ Single hostname, unified X-Service + X-Resource header routing - ✅ Single hostname, multiple path prefixes
- ✅ HTTP REST gateway → gRPC Temporal bridge (via X-Service: workflow) - ✅ HTTP REST gateway → gRPC Temporal bridge
- ✅ Bearer token auth via Authentik (JWT + RBAC) - ✅ Bearer token auth via Authentik (JWT + RBAC)
- ✅ Streaming unbuffered (SSE, WebSocket, HTTP/2 multiplexing) - ✅ Streaming unbuffered (SSE, WebSocket)
- ✅ Per-route timeouts & rate limits - ✅ Per-route timeouts & rate limits
- ✅ No cluster credentials held by gateway - ✅ No cluster credentials held by gateway
@@ -61,16 +61,16 @@ llm-serving temporal:7233 kmsvc/Kafka, MinIO,
## Services & Capabilities ## Services & Capabilities
| Service | Method | Upstream | Status | | Service | Prefix | Upstream | Status |
|---------|--------|----------|--------| |---------|--------|----------|--------|
| **LLM Chat** | `POST /v1/chat/completions` | llm-serving (vLLM) | ✅ Live | | **LLM Chat** | `/v1/chat/completions` | llm-serving (vLLM) | ✅ Live |
| **Embeddings** | `POST /v1/embeddings` | llm-serving (TEI) | ✅ Live | | **Embeddings** | `/v1/embeddings` | llm-serving (TEI) | ✅ Live |
| **Reranking** | `POST /v1/rerank` | llm-serving (TEI) | ✅ Live | | **Reranking** | `/v1/rerank` | llm-serving (TEI) | ✅ Live |
| **Workflows** | `X-Service: workflow` + `X-Resource: {action}` | Temporal gRPC (7233) | ✅ Live (START, DESCRIBE, SIGNAL, QUERY, etc) | | **Workflows** | `/workflow` | Temporal gRPC (7233) | ✅ Live (START, DESCRIBE, SIGNAL, QUERY, etc) |
| **Queues** | `X-Service: sqs` + `X-Resource: {action}` | kmsvc/Kafka | ✅ Live | | **Queues** | `/` + `X-Service: sqs` | kmsvc/Kafka | ⏳ Ready (ServiceAdapter) |
| **Memory** | `X-Service: memory` + `X-Resource: {action}` | poimen-memory | ✅ Live | | **Memory** | `/` + `X-Service: memory` | poimen-memory | ✅ Live |
| **IAM** | `X-Service: iam` + `X-Resource: {action}` | Authentik API | ✅ Live | | **IAM** | `/` + `X-Service: iam` | Authentik API | ✅ Live |
| **S3** | `X-Service: s3` + `X-Resource: {action}` | MinIO | ✅ Live | | **S3** | `/` + `X-Service: s3` | MinIO | ✅ Live |
--- ---
@@ -102,17 +102,18 @@ curl -X POST https://api.riotpiao.com/v1/chat/completions \
}' }'
``` ```
**Workflow (via X-Service header):** **Workflow:**
```bash ```bash
curl -X POST https://api.riotpiao.com/ \ curl -X POST https://api.riotpiao.com/workflow \
-H "Authorization: Bearer $TOKEN" \ -H "Authorization: Bearer $TOKEN" \
-H "X-Service: workflow" \
-H "X-Resource: start" \
-d '{ -d '{
"action": "START_WORKFLOW",
"namespace": "default", "namespace": "default",
"workflow_id": "my-workflow", "payload": {
"workflow_type": "MyWorkflow", "workflow_id": "my-workflow",
"task_queue": "default" "workflow_type": "MyWorkflow",
"task_queue": "default"
}
}' }'
``` ```
-15
View File
@@ -72,17 +72,6 @@ func main() {
// Create ServiceAdapter registry and dispatcher (phase 8) // Create ServiceAdapter registry and dispatcher (phase 8)
registry := serviceadapter.NewRegistry(nil) registry := serviceadapter.NewRegistry(nil)
// Add workflow service adapter (uses Temporal handler for gRPC forwarding)
workflowSpec := serviceadapter.GetWorkflowSpec()
workflowAdapter := &serviceadapter.ServiceAdapter{
Namespace: "temporal",
ServiceName: "workflow",
Spec: *workflowSpec,
}
_ = registry.Add(workflowAdapter)
// Add other adapters from config
for _, a := range cfg.Adapters { for _, a := range cfg.Adapters {
_ = registry.Add(a) _ = registry.Add(a)
} }
@@ -95,10 +84,6 @@ func main() {
} }
dispatcher := serviceadapter.NewDispatcher(registry, jwtValidator) dispatcher := serviceadapter.NewDispatcher(registry, jwtValidator)
// Wire workflow adapter to temporal handler for proper request forwarding
workflowAdapterImpl := serviceadapter.NewWorkflowAdapter(temporalHandler)
_ = workflowAdapterImpl // The dispatcher will call temporal handler directly for gRPC
// Create router that handles health endpoints, X-Service (ServiceAdapter) routing, // Create router that handles health endpoints, X-Service (ServiceAdapter) routing,
// temporal endpoints, and passes others to upstream handler // temporal endpoints, and passes others to upstream handler
router := server.NewRouter(healthChecker, dispatcher, temporalHandler, upstreamHandler) router := server.NewRouter(healthChecker, dispatcher, temporalHandler, upstreamHandler)
-36
View File
@@ -1,36 +0,0 @@
#!/bin/bash
# Example: Send email via notification/sendMsg endpoint
BASE_URL="${1:-https://api.riotpiao.com}"
AUTH_TOKEN="${2:-}" # Optional JWT token if auth required
PAYLOAD=$(cat <<'EOF'
{
"format": "smtp",
"title": "System Alert",
"message": "CPU usage exceeded 90% threshold",
"priority": 7,
"extras": {
"to_email": "[email protected]",
"cc": "[email protected]"
}
}
EOF
)
if [ -n "$AUTH_TOKEN" ]; then
curl -X POST "$BASE_URL" \
-H "X-Service: notification" \
-H "X-Resource: sendMsg" \
-H "Content-Type: application/json" \
-H "Authorization: Bearer $AUTH_TOKEN" \
-d "$PAYLOAD"
else
curl -X POST "$BASE_URL" \
-H "X-Service: notification" \
-H "X-Resource: sendMsg" \
-H "Content-Type: application/json" \
-d "$PAYLOAD"
fi
echo ""
-232
View File
@@ -1,232 +0,0 @@
package notification
import (
"encoding/json"
"fmt"
"log"
"net/http"
"net/smtp"
"os"
"strings"
)
// SendMsgRequest represents a sendMsg API request.
type SendMsgRequest struct {
Format string `json:"format"` // "smtp" or "sms"
Title string `json:"title"`
Message string `json:"message"`
Priority int `json:"priority,omitempty"`
Extras map[string]string `json:"extras,omitempty"` // e.g., {"to_email": "[email protected]", "phone": "+1234567890"}
}
// SendMsgResponse represents a sendMsg API response.
type SendMsgResponse struct {
Status string `json:"status"`
MessageID string `json:"messageId,omitempty"`
Error string `json:"error,omitempty"`
}
// Handler handles sendMsg requests and forwards to appropriate channel (email, SMS, or Gotify push).
type Handler struct {
smtpHost string
smtpPort string
smtpFrom string
smtpUser string
smtpPass string
smsAPIURL string
smsAPIKey string
gotifyURL string
gotifyToken string
}
// NewHandler creates a new notification handler from environment variables.
func NewHandler() *Handler {
return &Handler{
smtpHost: os.Getenv("SMTP_HOST"),
smtpPort: os.Getenv("SMTP_PORT"),
smtpFrom: os.Getenv("SMTP_FROM"),
smtpUser: os.Getenv("SMTP_USER"),
smtpPass: os.Getenv("SMTP_PASS"),
smsAPIURL: os.Getenv("SMS_API_URL"),
smsAPIKey: os.Getenv("SMS_API_KEY"),
gotifyURL: os.Getenv("GOTIFY_URL"),
gotifyToken: os.Getenv("GOTIFY_TOKEN"),
}
}
// ServeHTTP handles sendMsg requests.
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
var req SendMsgRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusBadRequest)
json.NewEncoder(w).Encode(SendMsgResponse{
Status: "error",
Error: "invalid request: " + err.Error(),
})
return
}
// Route based on format
var resp SendMsgResponse
switch req.Format {
case "smtp":
resp = h.sendEmail(req)
case "sms":
resp = h.sendSMS(req)
case "gotify":
resp = h.sendGotify(req)
default:
resp = SendMsgResponse{
Status: "error",
Error: "unsupported format: " + req.Format,
}
}
w.Header().Set("Content-Type", "application/json")
if resp.Error != "" {
w.WriteHeader(http.StatusInternalServerError)
} else {
w.WriteHeader(http.StatusOK)
}
json.NewEncoder(w).Encode(resp)
}
// sendEmail sends an email via SMTP.
func (h *Handler) sendEmail(req SendMsgRequest) SendMsgResponse {
toEmail := req.Extras["to_email"]
if toEmail == "" {
return SendMsgResponse{
Status: "error",
Error: "missing to_email in extras",
}
}
subject := req.Title
if subject == "" {
subject = "Notification"
}
// Construct email body
body := req.Message
if req.Extras != nil {
if cc := req.Extras["cc"]; cc != "" {
body = fmt.Sprintf("CC: %s\n\n%s", cc, body)
}
}
msg := fmt.Sprintf(
"From: %s\r\nTo: %s\r\nSubject: %s\r\nContent-Type: text/plain; charset=UTF-8\r\n\r\n%s",
h.smtpFrom, toEmail, subject, body,
)
// Send via SMTP
smtpAddr := fmt.Sprintf("%s:%s", h.smtpHost, h.smtpPort)
auth := smtp.PlainAuth("", h.smtpUser, h.smtpPass, h.smtpHost)
if err := smtp.SendMail(smtpAddr, auth, h.smtpFrom, []string{toEmail}, []byte(msg)); err != nil {
log.Printf("error sending email to %s: %v", toEmail, err)
return SendMsgResponse{
Status: "error",
Error: "failed to send email: " + err.Error(),
}
}
return SendMsgResponse{
Status: "success",
MessageID: fmt.Sprintf("email-%s", toEmail),
}
}
// sendSMS sends an SMS via configured provider.
// Placeholder: integrate with Twilio, AWS SNS, or similar.
func (h *Handler) sendSMS(req SendMsgRequest) SendMsgResponse {
phone := req.Extras["phone"]
if phone == "" {
return SendMsgResponse{
Status: "error",
Error: "missing phone in extras",
}
}
// TODO: Implement SMS provider integration (Twilio, AWS SNS, etc.)
// For now, return error
return SendMsgResponse{
Status: "error",
Error: "SMS not implemented yet",
}
}
// sendGotify sends a notification via Gotify server.
func (h *Handler) sendGotify(req SendMsgRequest) SendMsgResponse {
if h.gotifyURL == "" || h.gotifyToken == "" {
return SendMsgResponse{
Status: "error",
Error: "Gotify not configured (missing GOTIFY_URL or GOTIFY_TOKEN)",
}
}
// Construct Gotify message
gotifyReq := map[string]interface{}{
"title": req.Title,
"message": req.Message,
"priority": req.Priority,
}
// Marshal to JSON
body, err := json.Marshal(gotifyReq)
if err != nil {
log.Printf("error marshaling Gotify request: %v", err)
return SendMsgResponse{
Status: "error",
Error: "failed to marshal request: " + err.Error(),
}
}
// POST to Gotify
gotifyEndpoint := fmt.Sprintf("%s/message?token=%s", h.gotifyURL, h.gotifyToken)
resp, err := http.Post(gotifyEndpoint, "application/json", strings.NewReader(string(body)))
if err != nil {
log.Printf("error sending to Gotify: %v", err)
return SendMsgResponse{
Status: "error",
Error: "failed to send to Gotify: " + err.Error(),
}
}
defer resp.Body.Close()
// Check response status
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
log.Printf("Gotify returned status %d", resp.StatusCode)
return SendMsgResponse{
Status: "error",
Error: fmt.Sprintf("Gotify returned status %d", resp.StatusCode),
}
}
// Parse response
var gotifyResp map[string]interface{}
if err := json.NewDecoder(resp.Body).Decode(&gotifyResp); err != nil {
log.Printf("error decoding Gotify response: %v", err)
return SendMsgResponse{
Status: "success",
MessageID: "gotify-sent",
}
}
// Extract message ID if available
msgID := "gotify-sent"
if id, ok := gotifyResp["id"]; ok {
msgID = fmt.Sprintf("gotify-%v", id)
}
return SendMsgResponse{
Status: "success",
MessageID: msgID,
}
}
-28
View File
@@ -1,28 +0,0 @@
package observability
import (
"net/http"
)
// MetricsHandler serves Prometheus metrics
type MetricsHandler struct {
exporter *PrometheusExporter
}
// NewMetricsHandler creates a new metrics handler
func NewMetricsHandler(m *Metrics) *MetricsHandler {
return &MetricsHandler{
exporter: NewPrometheusExporter(m),
}
}
// ServeHTTP implements http.Handler for Prometheus /metrics endpoint
func (h *MetricsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain; version=0.0.4; charset=utf-8")
w.Header().Set("Cache-Control", "no-cache, no-store, must-revalidate")
w.Header().Set("Pragma", "no-cache")
w.Header().Set("Expires", "0")
w.WriteHeader(http.StatusOK)
w.Write([]byte(h.exporter.Export()))
}
-118
View File
@@ -28,14 +28,6 @@ type Metrics struct {
// Streaming metrics // Streaming metrics
streamingResponsesTotal map[string]int64 streamingResponsesTotal map[string]int64
streamingByteCount map[string]int64 streamingByteCount map[string]int64
// LLM inference metrics (TTFT and ITL)
// ttftMs: Time-to-First-Token in milliseconds
ttftMs map[string][]int64 // samples for histogram
// itlMs: Inter-Token Latency in milliseconds
itlMs map[string][]int64 // samples for histogram
// Token counts
tokenCount map[string]int64
} }
// NewMetrics creates a new Metrics instance. // NewMetrics creates a new Metrics instance.
@@ -49,9 +41,6 @@ func NewMetrics() *Metrics {
upstreamHealth: make(map[string]int), upstreamHealth: make(map[string]int),
streamingResponsesTotal: make(map[string]int64), streamingResponsesTotal: make(map[string]int64),
streamingByteCount: make(map[string]int64), streamingByteCount: make(map[string]int64),
ttftMs: make(map[string][]int64),
itlMs: make(map[string][]int64),
tokenCount: make(map[string]int64),
} }
} }
@@ -157,113 +146,9 @@ func (m *Metrics) GetMetrics() map[string]interface{} {
"upstream_health": m.upstreamHealth, "upstream_health": m.upstreamHealth,
"streaming_responses_total": m.streamingResponsesTotal, "streaming_responses_total": m.streamingResponsesTotal,
"streaming_byte_count": m.streamingByteCount, "streaming_byte_count": m.streamingByteCount,
"llm_ttft_ms": m.ttftMs,
"llm_itl_ms": m.itlMs,
"llm_token_count": m.tokenCount,
} }
} }
// RecordTTFT records Time-to-First-Token in milliseconds
func (m *Metrics) RecordTTFT(model string, ttftMs int64) {
m.mu.Lock()
defer m.mu.Unlock()
key := fmt.Sprintf("llm:ttft:%s", model)
m.ttftMs[key] = append(m.ttftMs[key], ttftMs)
}
// RecordITL records Inter-Token Latency in milliseconds
func (m *Metrics) RecordITL(model string, itlMs int64) {
m.mu.Lock()
defer m.mu.Unlock()
key := fmt.Sprintf("llm:itl:%s", model)
m.itlMs[key] = append(m.itlMs[key], itlMs)
}
// RecordTokenCount records number of tokens in response
func (m *Metrics) RecordTokenCount(model string, count int64) {
m.mu.Lock()
defer m.mu.Unlock()
key := fmt.Sprintf("llm:tokens:%s", model)
m.tokenCount[key] += count
}
// GetTTFTMetrics returns TTFT statistics for Prometheus export
func (m *Metrics) GetTTFTMetrics() map[string]interface{} {
m.mu.RLock()
defer m.mu.RUnlock()
result := make(map[string]interface{})
for key, samples := range m.ttftMs {
if len(samples) > 0 {
result[key] = map[string]interface{}{
"count": len(samples),
"sum": sumInt64(samples),
"avg": sumInt64(samples) / int64(len(samples)),
"min": minInt64(samples),
"max": maxInt64(samples),
}
}
}
return result
}
// GetITLMetrics returns ITL statistics for Prometheus export
func (m *Metrics) GetITLMetrics() map[string]interface{} {
m.mu.RLock()
defer m.mu.RUnlock()
result := make(map[string]interface{})
for key, samples := range m.itlMs {
if len(samples) > 0 {
result[key] = map[string]interface{}{
"count": len(samples),
"sum": sumInt64(samples),
"avg": sumInt64(samples) / int64(len(samples)),
"min": minInt64(samples),
"max": maxInt64(samples),
}
}
}
return result
}
func sumInt64(vals []int64) int64 {
var s int64
for _, v := range vals {
s += v
}
return s
}
func minInt64(vals []int64) int64 {
if len(vals) == 0 {
return 0
}
min := vals[0]
for _, v := range vals {
if v < min {
min = v
}
}
return min
}
func maxInt64(vals []int64) int64 {
if len(vals) == 0 {
return 0
}
max := vals[0]
for _, v := range vals {
if v > max {
max = v
}
}
return max
}
// Reset clears all metrics (for testing). // Reset clears all metrics (for testing).
func (m *Metrics) Reset() { func (m *Metrics) Reset() {
m.mu.Lock() m.mu.Lock()
@@ -277,7 +162,4 @@ func (m *Metrics) Reset() {
m.upstreamHealth = make(map[string]int) m.upstreamHealth = make(map[string]int)
m.streamingResponsesTotal = make(map[string]int64) m.streamingResponsesTotal = make(map[string]int64)
m.streamingByteCount = make(map[string]int64) m.streamingByteCount = make(map[string]int64)
m.ttftMs = make(map[string][]int64)
m.itlMs = make(map[string][]int64)
m.tokenCount = make(map[string]int64)
} }
-198
View File
@@ -1,198 +0,0 @@
package observability
import (
"fmt"
"sort"
"strings"
)
// PrometheusExporter exports metrics in Prometheus text format
type PrometheusExporter struct {
metrics *Metrics
}
// NewPrometheusExporter creates a new Prometheus exporter
func NewPrometheusExporter(m *Metrics) *PrometheusExporter {
return &PrometheusExporter{metrics: m}
}
// Export returns metrics in Prometheus text format
func (p *PrometheusExporter) Export() string {
var lines []string
lines = append(lines, "# HELP llm_ttft_seconds Time to first token for LLM inference (seconds)")
lines = append(lines, "# TYPE llm_ttft_seconds histogram")
p.exportTTFT(&lines)
lines = append(lines, "# HELP llm_itl_seconds Inter-token latency for LLM inference (seconds)")
lines = append(lines, "# TYPE llm_itl_seconds histogram")
p.exportITL(&lines)
lines = append(lines, "# HELP llm_tokens_total Total tokens generated")
lines = append(lines, "# TYPE llm_tokens_total counter")
p.exportTokens(&lines)
lines = append(lines, "# HELP request_duration_seconds Request latency")
lines = append(lines, "# TYPE request_duration_seconds histogram")
p.exportRequestDuration(&lines)
return strings.Join(lines, "\n") + "\n"
}
func (p *PrometheusExporter) exportTTFT(lines *[]string) {
p.metrics.mu.RLock()
defer p.metrics.mu.RUnlock()
// Calculate statistics for each model
for key, samples := range p.metrics.ttftMs {
if len(samples) == 0 {
continue
}
model := extractModel(key)
sum := sumInt64(samples)
// Export histogram buckets (in seconds)
buckets := []float64{0.001, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0}
for _, bucket := range buckets {
count := countLessOrEqual(samples, int64(bucket*1000))
*lines = append(*lines, fmt.Sprintf(
`llm_ttft_seconds_bucket{model="%s",le="%.3f"} %d`,
model, bucket, count,
))
}
*lines = append(*lines, fmt.Sprintf(
`llm_ttft_seconds_bucket{model="%s",le="+Inf"} %d`,
model, len(samples),
))
*lines = append(*lines, fmt.Sprintf(
`llm_ttft_seconds_sum{model="%s"} %.3f`,
model, float64(sum)/1000,
))
*lines = append(*lines, fmt.Sprintf(
`llm_ttft_seconds_count{model="%s"} %d`,
model, len(samples),
))
}
}
func (p *PrometheusExporter) exportITL(lines *[]string) {
p.metrics.mu.RLock()
defer p.metrics.mu.RUnlock()
for key, samples := range p.metrics.itlMs {
if len(samples) == 0 {
continue
}
model := extractModel(key)
sum := sumInt64(samples)
// Export histogram buckets (in seconds)
buckets := []float64{0.001, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0}
for _, bucket := range buckets {
count := countLessOrEqual(samples, int64(bucket*1000))
*lines = append(*lines, fmt.Sprintf(
`llm_itl_seconds_bucket{model="%s",le="%.3f"} %d`,
model, bucket, count,
))
}
*lines = append(*lines, fmt.Sprintf(
`llm_itl_seconds_bucket{model="%s",le="+Inf"} %d`,
model, len(samples),
))
*lines = append(*lines, fmt.Sprintf(
`llm_itl_seconds_sum{model="%s"} %.3f`,
model, float64(sum)/1000,
))
*lines = append(*lines, fmt.Sprintf(
`llm_itl_seconds_count{model="%s"} %d`,
model, len(samples),
))
}
}
func (p *PrometheusExporter) exportTokens(lines *[]string) {
p.metrics.mu.RLock()
defer p.metrics.mu.RUnlock()
// Sort keys for consistent output
var keys []string
for k := range p.metrics.tokenCount {
keys = append(keys, k)
}
sort.Strings(keys)
for _, key := range keys {
model := extractModel(key)
count := p.metrics.tokenCount[key]
*lines = append(*lines, fmt.Sprintf(
`llm_tokens_total{model="%s"} %d`,
model, count,
))
}
}
func (p *PrometheusExporter) exportRequestDuration(lines *[]string) {
p.metrics.mu.RLock()
defer p.metrics.mu.RUnlock()
// Sort keys for consistent output
var keys []string
for k := range p.metrics.requestDuration {
keys = append(keys, k)
}
sort.Strings(keys)
for _, key := range keys {
route, upstream := parseKey(key)
totalMs := p.metrics.requestDuration[key]
count := int64(1) // We'd need to track count separately in real impl
if buckets, ok := p.metrics.requestDurationBuckets[key]; ok {
for bucket := range buckets {
*lines = append(*lines, fmt.Sprintf(
`request_duration_seconds_bucket{route="%s",upstream="%s",le="%.1f"} %d`,
route, upstream, bucket, buckets[bucket],
))
}
}
*lines = append(*lines, fmt.Sprintf(
`request_duration_seconds_sum{route="%s",upstream="%s"} %.3f`,
route, upstream, float64(totalMs)/1000,
))
*lines = append(*lines, fmt.Sprintf(
`request_duration_seconds_count{route="%s",upstream="%s"} %d`,
route, upstream, count,
))
}
}
func extractModel(key string) string {
parts := strings.Split(key, ":")
if len(parts) >= 3 {
return parts[2]
}
return key
}
func parseKey(key string) (string, string) {
parts := strings.Split(key, ":")
if len(parts) >= 2 {
return parts[0], parts[1]
}
return key, ""
}
func countLessOrEqual(samples []int64, threshold int64) int {
count := 0
for _, s := range samples {
if s <= threshold {
count++
}
}
return count
}
-153
View File
@@ -1,153 +0,0 @@
package proxy
import (
"bufio"
"fmt"
"io"
"net"
"net/http"
"strings"
"time"
"forgejo.riotpiao.com/rock/homelab-frontend/internal/observability"
)
// LLMMetricsCapture wraps a response writer to capture TTFT and ITL metrics
type LLMMetricsCapture struct {
writer io.WriteCloser
model string
metrics *observability.Metrics
firstTokenTime time.Time
lastTokenTime time.Time
requestStartTime time.Time
ttftRecorded bool
tokenCount int64
responseStartTime time.Time
}
// NewLLMMetricsCapture creates a new metrics capture wrapper
func NewLLMMetricsCapture(writer io.WriteCloser, model string, metrics *observability.Metrics, startTime time.Time) *LLMMetricsCapture {
return &LLMMetricsCapture{
writer: writer,
model: model,
metrics: metrics,
requestStartTime: startTime,
responseStartTime: time.Now(),
}
}
// Write intercepts writes to detect tokens and record metrics
func (c *LLMMetricsCapture) Write(p []byte) (int, error) {
// Record first token time
if !c.ttftRecorded && len(p) > 0 {
now := time.Now()
ttft := now.Sub(c.requestStartTime).Milliseconds()
c.metrics.RecordTTFT(c.model, ttft)
c.ttftRecorded = true
c.firstTokenTime = now
c.lastTokenTime = now
}
// Count tokens in SSE stream (simple: count "data: " lines)
if c.ttftRecorded {
tokenCount := strings.Count(string(p), "data: ")
if tokenCount > 0 {
now := time.Now()
if !c.firstTokenTime.IsZero() && c.lastTokenTime != now {
itl := now.Sub(c.lastTokenTime).Milliseconds()
c.metrics.RecordITL(c.model, itl)
}
c.lastTokenTime = now
c.tokenCount += int64(tokenCount)
}
}
return c.writer.Write(p)
}
// Close records final metrics and closes writer
func (c *LLMMetricsCapture) Close() error {
if c.tokenCount > 0 {
c.metrics.RecordTokenCount(c.model, c.tokenCount)
}
return c.writer.Close()
}
// ResponseWriterWrapper wraps http.ResponseWriter to capture metrics
type ResponseWriterWrapper struct {
writer http.ResponseWriter
statusCode int
metrics *observability.Metrics
model string
startTime time.Time
firstByteTime time.Time
lastWriteTime time.Time
ttftRecorded bool
}
// NewResponseWriterWrapper creates a wrapper for response writer
func NewResponseWriterWrapper(w http.ResponseWriter, model string, metrics *observability.Metrics, startTime time.Time) *ResponseWriterWrapper {
return &ResponseWriterWrapper{
writer: w,
model: model,
metrics: metrics,
startTime: startTime,
statusCode: 200,
}
}
// Header implements http.ResponseWriter
func (w *ResponseWriterWrapper) Header() http.Header {
return w.writer.Header()
}
// Write implements http.ResponseWriter
func (w *ResponseWriterWrapper) Write(b []byte) (int, error) {
// Record TTFT on first write
if !w.ttftRecorded && len(b) > 0 {
now := time.Now()
ttft := now.Sub(w.startTime).Milliseconds()
w.metrics.RecordTTFT(w.model, ttft)
w.ttftRecorded = true
w.firstByteTime = now
w.lastWriteTime = now
}
// Record ITL for subsequent writes (for streaming)
if w.ttftRecorded && len(b) > 0 {
now := time.Now()
if !w.firstByteTime.IsZero() && w.lastWriteTime != now {
itl := now.Sub(w.lastWriteTime).Milliseconds()
// Only record if ITL > 0 (avoid recording same millisecond twice)
if itl > 0 {
w.metrics.RecordITL(w.model, itl)
}
}
w.lastWriteTime = now
}
return w.writer.Write(b)
}
// WriteHeader implements http.ResponseWriter
func (w *ResponseWriterWrapper) WriteHeader(statusCode int) {
w.statusCode = statusCode
w.writer.WriteHeader(statusCode)
}
// Flush implements http.Flusher
func (w *ResponseWriterWrapper) Flush() {
if flusher, ok := w.writer.(http.Flusher); ok {
flusher.Flush()
}
}
// Hijack implements http.Hijacker for streaming
func (w *ResponseWriterWrapper) Hijack() (net.Conn, *bufio.ReadWriter, error) {
if hijacker, ok := w.writer.(http.Hijacker); ok {
return hijacker.Hijack()
}
return nil, nil, fmt.Errorf("response writer does not implement Hijacker")
}
+7 -25
View File
@@ -11,7 +11,6 @@ import (
"net/url" "net/url"
"sort" "sort"
"strings" "strings"
"syscall"
"time" "time"
"forgejo.riotpiao.com/rock/homelab-frontend/internal/auth" "forgejo.riotpiao.com/rock/homelab-frontend/internal/auth"
@@ -128,15 +127,6 @@ func (h *Handler) getOrCreateTransport(addr string, up *config.Upstream) *http.T
dialer := &net.Dialer{ dialer := &net.Dialer{
Timeout: up.ConnectTimeout, Timeout: up.ConnectTimeout,
KeepAlive: 30 * time.Second, KeepAlive: 30 * time.Second,
// Issue #31: TCP_NODELAY disables Nagle's algorithm, reducing latency
// for streaming responses by sending small packets immediately instead of
// waiting for larger batches. Critical for low-latency LLM token streaming.
Control: func(network, address string, c syscall.RawConn) error {
return c.Control(func(fd uintptr) {
// TCP_NODELAY disables Nagle's algorithm for immediate packet transmission
_ = syscall.SetsockoptInt(int(fd), syscall.IPPROTO_TCP, syscall.TCP_NODELAY, 1)
})
},
} }
transport := &http.Transport{ transport := &http.Transport{
@@ -144,15 +134,8 @@ func (h *Handler) getOrCreateTransport(addr string, up *config.Upstream) *http.T
DialContext: dialer.DialContext, DialContext: dialer.DialContext,
MaxIdleConns: 100, MaxIdleConns: 100,
IdleConnTimeout: 90 * time.Second, IdleConnTimeout: 90 * time.Second,
// Issue #32: Increase per-host connection limit to support HTTP/2 multiplexing.
// With HTTP/2, we can serve many concurrent streams over fewer connections,
// but we still allow more connections for better resource utilization.
MaxConnsPerHost: 10,
// Allow persistent connections // Allow persistent connections
DisableKeepAlives: false, DisableKeepAlives: false,
// Issue #31: Enable HTTP/2 for client connections to support multiplexing.
// This allows concurrent requests to stream simultaneously with better flow control.
ForceAttemptHTTP2: true,
} }
// Store the upstream config for use in the handler // Store the upstream config for use in the handler
@@ -271,8 +254,13 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
return return
} }
// Try to find a matching route (including body-based dispatch for /v1/chat/completions). // Handle /workflows endpoint (workflow orchestration)
// Note: /workflows endpoint is deprecated. Use X-Service: workflow + X-Resource headers instead. if r.URL.Path == "/workflows" {
h.handleWorkflow(w, r)
return
}
// Try to find a matching route (including body-based dispatch for /v1/chat/completions)
route, err := h.RouteRequest(r) route, err := h.RouteRequest(r)
// Check if this is a model validation error (from body-based dispatch) // Check if this is a model validation error (from body-based dispatch)
@@ -441,12 +429,6 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// A streaming response that's continuously sending should not be cut off. // A streaming response that's continuously sending should not be cut off.
// The Transport's socket read timeout (via Dialer) handles inactivity timeouts. // The Transport's socket read timeout (via Dialer) handles inactivity timeouts.
// For streaming responses (SSE, chunked), disable buffering to ensure events
// reach clients immediately. Issue #33: X-Accel-Buffering:no tells nginx/Ingress
// to stream instead of buffer. ResponseController.Flush() in upstream handler
// pairs with this to deliver unbuffered chunks.
w.Header().Set("X-Accel-Buffering", "no")
// Serve the request through the proxy // Serve the request through the proxy
proxy.ServeHTTP(w, r) proxy.ServeHTTP(w, r)
} }
-208
View File
@@ -476,214 +476,6 @@ func TestNoFullBuffering(t *testing.T) {
} }
} }
// TestTCPBackpressure verifies that TCP backpressure is respected during streaming.
// When a client reads slowly, the upstream should experience backpressure on writes.
func TestTCPBackpressure(t *testing.T) {
// Track when upstream started writing and when each write completed
var writeTimes []time.Time
writesMu := sync.Mutex{}
upstreamServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("X-Accel-Buffering", "no") // Issue #33: disable buffering
w.WriteHeader(http.StatusOK)
rc := http.NewResponseController(w)
// Send many events to trigger backpressure
for i := 0; i < 20; i++ {
writesMu.Lock()
writeTimes = append(writeTimes, time.Now())
writesMu.Unlock()
fmt.Fprintf(w, "data: event%d\n\n", i)
if err := rc.Flush(); err != nil {
return
}
}
}))
defer upstreamServer.Close()
upstreamAddr := strings.TrimPrefix(upstreamServer.URL, "http://")
cfg := &config.Config{
Routes: map[string]*config.Route{
"backpressure-route": {
Name: "backpressure-route",
Upstream: config.Upstream{
Address: upstreamAddr,
ConnectTimeout: 5 * time.Second,
ReadTimeout: 10 * time.Second,
WriteTimeout: 5 * time.Second,
MaxBodySize: 1024 * 1024,
AuthRequired: false,
},
},
},
}
handler := New(cfg)
defer handler.Close()
server := httptest.NewServer(handler)
defer server.Close()
resp, err := http.Get(server.URL + "/backpressure")
if err != nil {
t.Fatalf("request failed: %v", err)
}
defer resp.Body.Close()
// Verify X-Accel-Buffering header is passed through
if resp.Header.Get("X-Accel-Buffering") != "no" {
t.Errorf("X-Accel-Buffering header not propagated, got: %s", resp.Header.Get("X-Accel-Buffering"))
}
// Read events with simulated slow client (small buffer)
reader := bufio.NewReader(resp.Body)
readStart := time.Now()
eventCount := 0
for {
line, err := reader.ReadString('\n')
if err != nil {
if err == io.EOF {
break
}
t.Fatalf("read failed: %v", err)
}
if strings.HasPrefix(strings.TrimSpace(line), "data:") {
eventCount++
// Simulate slow client by adding delay
time.Sleep(5 * time.Millisecond)
}
}
// Verify we got all events
if eventCount != 20 {
t.Errorf("expected 20 events, got %d", eventCount)
}
// Total read time should be roughly eventCount * readDelay
// indicating backpressure was applied (upstream couldn't send all at once)
elapsed := time.Since(readStart)
expectedMin := time.Duration(20*5) * time.Millisecond
if elapsed < expectedMin {
t.Logf("backpressure test: elapsed=%.0fms (expected ~%.0fms)", elapsed.Seconds()*1000, expectedMin.Seconds()*1000)
}
}
// TestConcurrentSSEStreams verifies that HTTP/2 multiplexing handles multiple concurrent streams.
// Issue #32: Multiple LLM requests should not block each other.
func TestConcurrentSSEStreams(t *testing.T) {
upstreamServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("X-Accel-Buffering", "no")
w.WriteHeader(http.StatusOK)
rc := http.NewResponseController(w)
// Each request sends unique identifier
reqID := r.URL.Query().Get("id")
for i := 0; i < 5; i++ {
fmt.Fprintf(w, "data: [%s] event %d\n\n", reqID, i)
if err := rc.Flush(); err != nil {
return
}
time.Sleep(10 * time.Millisecond)
}
}))
defer upstreamServer.Close()
upstreamAddr := strings.TrimPrefix(upstreamServer.URL, "http://")
cfg := &config.Config{
Routes: map[string]*config.Route{
"concurrent-route": {
Name: "concurrent-route",
Upstream: config.Upstream{
Address: upstreamAddr,
ConnectTimeout: 5 * time.Second,
ReadTimeout: 10 * time.Second,
WriteTimeout: 5 * time.Second,
MaxBodySize: 1024 * 1024,
AuthRequired: false,
},
},
},
}
handler := New(cfg)
defer handler.Close()
server := httptest.NewServer(handler)
defer server.Close()
// Launch multiple concurrent requests
var wg sync.WaitGroup
results := make(map[string][]string)
resultsMu := sync.Mutex{}
for id := 0; id < 3; id++ {
wg.Add(1)
go func(streamID int) {
defer wg.Done()
url := fmt.Sprintf("%s/concurrent?id=stream%d", server.URL, streamID)
resp, err := http.Get(url)
if err != nil {
t.Errorf("request failed: %v", err)
return
}
defer resp.Body.Close()
reader := bufio.NewReader(resp.Body)
var events []string
for {
line, err := reader.ReadString('\n')
if err != nil {
if err == io.EOF {
break
}
t.Errorf("read failed: %v", err)
return
}
line = strings.TrimSpace(line)
if strings.HasPrefix(line, "data:") {
events = append(events, line)
}
}
resultsMu.Lock()
results[fmt.Sprintf("stream%d", streamID)] = events
resultsMu.Unlock()
}(id)
}
wg.Wait()
// Verify all streams got their events
for i := 0; i < 3; i++ {
key := fmt.Sprintf("stream%d", i)
events, ok := results[key]
if !ok {
t.Errorf("stream%d: no results", i)
continue
}
if len(events) != 5 {
t.Errorf("stream%d: expected 5 events, got %d", i, len(events))
}
// Verify all events belong to this stream
for _, event := range events {
if !strings.Contains(event, key) {
t.Errorf("stream%d: event from wrong stream: %s", i, event)
}
}
}
}
// TestClientDisconnectCancelsUpstream verifies that when a client closes mid-stream, // TestClientDisconnectCancelsUpstream verifies that when a client closes mid-stream,
// the upstream request context is cancelled immediately and no goroutines are leaked. // the upstream request context is cancelled immediately and no goroutines are leaked.
func TestClientDisconnectCancelsUpstream(t *testing.T) { func TestClientDisconnectCancelsUpstream(t *testing.T) {
+484
View File
@@ -0,0 +1,484 @@
// Package proxy provides request routing and forwarding.
package proxy
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"time"
)
// WorkflowRequest represents a workflow execution request
type WorkflowRequest struct {
// Workflow ID or name
Workflow string `json:"workflow"`
// Input parameters for the workflow
Input map[string]interface{} `json:"input"`
// Optional: timeout in seconds
Timeout int `json:"timeout,omitempty"`
// Optional: wait for result (default: true)
Wait *bool `json:"wait,omitempty"`
}
// WorkflowResponse represents the response from workflow execution
type WorkflowResponse struct {
// Workflow execution ID
ID string `json:"id"`
// Workflow name
Workflow string `json:"workflow"`
// Execution status: pending, running, completed, failed
Status string `json:"status"`
// Output of the workflow
Output interface{} `json:"output,omitempty"`
// Error message if workflow failed
Error string `json:"error,omitempty"`
// Timestamp when workflow was created
CreatedAt time.Time `json:"created_at"`
// Timestamp when workflow completed
CompletedAt *time.Time `json:"completed_at,omitempty"`
}
// PredefinedWorkflow defines a workflow template that combines multiple API calls
type PredefinedWorkflow struct {
Name string
Description string
Handler func(*http.Request, *Handler, map[string]interface{}) (interface{}, error)
}
// handleWorkflow handles the /workflows endpoint
// It accepts workflow definitions and orchestrates API calls
func (h *Handler) handleWorkflow(w http.ResponseWriter, r *http.Request) {
// Only POST is supported
if r.Method != "POST" {
w.Header().Set("Content-Type", "application/problem+json")
w.WriteHeader(http.StatusMethodNotAllowed)
fmt.Fprintf(w, `{"type":"https://api.example.com/problems/method-not-allowed","title":"Method Not Allowed","status":405,"detail":"Only POST is supported for /workflows"}`)
return
}
// Parse request body
var workflowReq WorkflowRequest
if err := json.NewDecoder(r.Body).Decode(&workflowReq); err != nil {
writeProblemDetail(w, http.StatusBadRequest, "https://api.example.com/problems/invalid-workflow-request", "Invalid Workflow Request", "Failed to parse workflow request: "+err.Error(), nil)
return
}
// Validate workflow name
if workflowReq.Workflow == "" {
writeProblemDetail(w, http.StatusBadRequest, "https://api.example.com/problems/missing-workflow", "Missing Workflow", "The 'workflow' field is required", nil)
return
}
// Get predefined workflow
workflow, ok := h.getWorkflow(workflowReq.Workflow)
if !ok {
availableWorkflows := h.getAvailableWorkflows()
writeProblemDetail(w, http.StatusBadRequest, "https://api.example.com/problems/unknown-workflow", "Unknown Workflow", fmt.Sprintf("Workflow %q is not available", workflowReq.Workflow), availableWorkflows)
return
}
// Default wait to true
wait := true
if workflowReq.Wait != nil {
wait = *workflowReq.Wait
}
// Set default timeout if not provided
timeout := time.Duration(30) * time.Second
if workflowReq.Timeout > 0 {
timeout = time.Duration(workflowReq.Timeout) * time.Second
}
// Create a context with timeout for workflow execution
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
// Execute workflow
output, err := workflow.Handler(r.WithContext(ctx), h, workflowReq.Input)
// Build response
workflowResp := WorkflowResponse{
ID: generateWorkflowID(),
Workflow: workflowReq.Workflow,
CreatedAt: time.Now(),
}
if err != nil {
workflowResp.Status = "failed"
workflowResp.Error = err.Error()
} else {
if wait {
workflowResp.Status = "completed"
workflowResp.Output = output
now := time.Now()
workflowResp.CompletedAt = &now
} else {
workflowResp.Status = "pending"
}
}
// Write response
w.Header().Set("Content-Type", "application/json")
if err != nil {
w.WriteHeader(http.StatusInternalServerError)
} else {
w.WriteHeader(http.StatusOK)
}
json.NewEncoder(w).Encode(workflowResp)
}
// getWorkflow returns a predefined workflow by name
func (h *Handler) getWorkflow(name string) (*PredefinedWorkflow, bool) {
workflows := h.getPredefinedWorkflows()
for _, wf := range workflows {
if wf.Name == name {
return &wf, true
}
}
return nil, false
}
// getPredefinedWorkflows returns all available workflows
func (h *Handler) getPredefinedWorkflows() []PredefinedWorkflow {
return []PredefinedWorkflow{
{
Name: "chat-and-embed",
Description: "Chat with a model and then embed the response",
Handler: h.chatAndEmbedWorkflow,
},
{
Name: "multi-model-chat",
Description: "Chat with multiple models sequentially",
Handler: h.multiModelChatWorkflow,
},
{
Name: "rag-pipeline",
Description: "RAG pipeline: embed query, rerank, then chat with context",
Handler: h.ragPipelineWorkflow,
},
{
Name: "batch-embeddings",
Description: "Generate embeddings for multiple texts",
Handler: h.batchEmbeddingsWorkflow,
},
}
}
// getAvailableWorkflows returns a list of available workflow names
func (h *Handler) getAvailableWorkflows() []string {
workflows := h.getPredefinedWorkflows()
names := make([]string, len(workflows))
for i, wf := range workflows {
names[i] = wf.Name
}
return names
}
// Workflow implementations
// chatAndEmbedWorkflow: Chat with a model, then embed the response
func (h *Handler) chatAndEmbedWorkflow(r *http.Request, handler *Handler, input map[string]interface{}) (interface{}, error) {
model, ok := input["model"].(string)
if !ok || model == "" {
return nil, fmt.Errorf("missing required parameter: model")
}
embedModel, ok := input["embed_model"].(string)
if !ok {
embedModel = "nomic-ai/nomic-embed-text-v2-moe"
}
messages, ok := input["messages"].([]interface{})
if !ok {
return nil, fmt.Errorf("missing required parameter: messages")
}
// Step 1: Chat
chatReq := map[string]interface{}{
"model": model,
"messages": messages,
}
chatBody, _ := json.Marshal(chatReq)
chatHTTPReq, _ := http.NewRequest("POST", "/v1/chat/completions", io.NopCloser(bytes.NewReader(chatBody)))
chatHTTPReq.Header.Set("Content-Type", "application/json")
// Create a response writer to capture the chat response
chatResp := &responseCapture{}
handler.ServeHTTP(chatResp, chatHTTPReq)
var chatResult map[string]interface{}
if err := json.Unmarshal(chatResp.body.Bytes(), &chatResult); err != nil {
return nil, fmt.Errorf("failed to parse chat response: %v", err)
}
// Extract message content
var messageContent string
if choices, ok := chatResult["choices"].([]interface{}); ok && len(choices) > 0 {
if choice, ok := choices[0].(map[string]interface{}); ok {
if message, ok := choice["message"].(map[string]interface{}); ok {
if content, ok := message["content"].(string); ok {
messageContent = content
}
}
}
}
// Step 2: Embed the response
embedReq := map[string]interface{}{
"model": embedModel,
"input": messageContent,
}
embedBody, _ := json.Marshal(embedReq)
embedHTTPReq, _ := http.NewRequest("POST", "/v1/embeddings", io.NopCloser(bytes.NewReader(embedBody)))
embedHTTPReq.Header.Set("Content-Type", "application/json")
embedResp := &responseCapture{}
handler.ServeHTTP(embedResp, embedHTTPReq)
var embedResult map[string]interface{}
if err := json.Unmarshal(embedResp.body.Bytes(), &embedResult); err != nil {
return nil, fmt.Errorf("failed to parse embedding response: %v", err)
}
return map[string]interface{}{
"chat_response": chatResult,
"embedding_response": embedResult,
}, nil
}
// multiModelChatWorkflow: Chat with multiple models sequentially
func (h *Handler) multiModelChatWorkflow(r *http.Request, handler *Handler, input map[string]interface{}) (interface{}, error) {
models, ok := input["models"].([]interface{})
if !ok || len(models) == 0 {
return nil, fmt.Errorf("missing required parameter: models (array)")
}
messages, ok := input["messages"].([]interface{})
if !ok {
return nil, fmt.Errorf("missing required parameter: messages")
}
results := make([]map[string]interface{}, 0)
for _, modelInterface := range models {
model, ok := modelInterface.(string)
if !ok {
continue
}
chatReq := map[string]interface{}{
"model": model,
"messages": messages,
}
chatBody, _ := json.Marshal(chatReq)
chatHTTPReq, _ := http.NewRequest("POST", "/v1/chat/completions", io.NopCloser(bytes.NewReader(chatBody)))
chatHTTPReq.Header.Set("Content-Type", "application/json")
chatResp := &responseCapture{}
handler.ServeHTTP(chatResp, chatHTTPReq)
var chatResult map[string]interface{}
if err := json.Unmarshal(chatResp.body.Bytes(), &chatResult); err != nil {
results = append(results, map[string]interface{}{
"model": model,
"error": err.Error(),
})
continue
}
results = append(results, map[string]interface{}{
"model": model,
"result": chatResult,
})
}
return results, nil
}
// ragPipelineWorkflow: RAG pipeline - embed query, rerank, chat with context
func (h *Handler) ragPipelineWorkflow(r *http.Request, handler *Handler, input map[string]interface{}) (interface{}, error) {
query, ok := input["query"].(string)
if !ok || query == "" {
return nil, fmt.Errorf("missing required parameter: query")
}
documents, ok := input["documents"].([]interface{})
if !ok {
return nil, fmt.Errorf("missing required parameter: documents")
}
model, ok := input["model"].(string)
if !ok {
model = "reasoning"
}
rerankModel, ok := input["rerank_model"].(string)
if !ok {
rerankModel = "BAAI/bge-reranker-base"
}
topK := 3
if tk, ok := input["top_k"].(float64); ok {
topK = int(tk)
}
// Step 1: Rerank documents based on query
rerankReq := map[string]interface{}{
"model": rerankModel,
"query": query,
"texts": documents,
"top_k": topK,
}
rerankBody, _ := json.Marshal(rerankReq)
rerankHTTPReq, _ := http.NewRequest("POST", "/v1/rerank", io.NopCloser(bytes.NewReader(rerankBody)))
rerankHTTPReq.Header.Set("Content-Type", "application/json")
rerankResp := &responseCapture{}
handler.ServeHTTP(rerankResp, rerankHTTPReq)
var rerankResult map[string]interface{}
if err := json.Unmarshal(rerankResp.body.Bytes(), &rerankResult); err != nil {
return nil, fmt.Errorf("failed to parse rerank response: %v", err)
}
// Extract top documents
var topDocs []string
if results, ok := rerankResult["results"].([]interface{}); ok {
for i, resultInterface := range results {
if i >= topK {
break
}
if result, ok := resultInterface.(map[string]interface{}); ok {
if text, ok := result["text"].(string); ok {
topDocs = append(topDocs, text)
}
}
}
}
// Step 2: Chat with context
context := fmt.Sprintf("Context from documents:\n%v\n\nQuery: %s", topDocs, query)
chatReq := map[string]interface{}{
"model": model,
"messages": []interface{}{
map[string]interface{}{
"role": "user",
"content": context,
},
},
}
chatBody, _ := json.Marshal(chatReq)
chatHTTPReq, _ := http.NewRequest("POST", "/v1/chat/completions", io.NopCloser(bytes.NewReader(chatBody)))
chatHTTPReq.Header.Set("Content-Type", "application/json")
chatResp := &responseCapture{}
handler.ServeHTTP(chatResp, chatHTTPReq)
var chatResult map[string]interface{}
if err := json.Unmarshal(chatResp.body.Bytes(), &chatResult); err != nil {
return nil, fmt.Errorf("failed to parse chat response: %v", err)
}
return map[string]interface{}{
"reranked_documents": topDocs,
"chat_response": chatResult,
}, nil
}
// batchEmbeddingsWorkflow: Generate embeddings for multiple texts
func (h *Handler) batchEmbeddingsWorkflow(r *http.Request, handler *Handler, input map[string]interface{}) (interface{}, error) {
texts, ok := input["texts"].([]interface{})
if !ok || len(texts) == 0 {
return nil, fmt.Errorf("missing required parameter: texts (array)")
}
model, ok := input["model"].(string)
if !ok {
model = "nomic-ai/nomic-embed-text-v2-moe"
}
// Convert interface{} to []string
textStrings := make([]string, 0)
for _, t := range texts {
if str, ok := t.(string); ok {
textStrings = append(textStrings, str)
}
}
if len(textStrings) == 0 {
return nil, fmt.Errorf("no valid text strings in texts array")
}
embedReq := map[string]interface{}{
"model": model,
"input": textStrings,
}
embedBody, _ := json.Marshal(embedReq)
embedHTTPReq, _ := http.NewRequest("POST", "/v1/embeddings", io.NopCloser(bytes.NewReader(embedBody)))
embedHTTPReq.Header.Set("Content-Type", "application/json")
embedResp := &responseCapture{}
handler.ServeHTTP(embedResp, embedHTTPReq)
var embedResult map[string]interface{}
if err := json.Unmarshal(embedResp.body.Bytes(), &embedResult); err != nil {
return nil, fmt.Errorf("failed to parse embedding response: %v", err)
}
return embedResult, nil
}
// Utility functions
// responseCapture captures HTTP response for reuse within workflows
type responseCapture struct {
status int
header http.Header
body bytes.Buffer
}
func (w *responseCapture) Header() http.Header {
if w.header == nil {
w.header = make(http.Header)
}
return w.header
}
func (w *responseCapture) Write(b []byte) (int, error) {
if w.status == 0 {
w.status = http.StatusOK
}
return w.body.Write(b)
}
func (w *responseCapture) WriteHeader(statusCode int) {
if w.status == 0 {
w.status = statusCode
}
}
// generateWorkflowID generates a unique workflow execution ID
func generateWorkflowID() string {
return fmt.Sprintf("wf_%d", time.Now().UnixNano())
}
+258
View File
@@ -0,0 +1,258 @@
package proxy
import (
"bytes"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"forgejo.riotpiao.com/rock/homelab-frontend/internal/config"
)
func TestWorkflowEndpointNotFound(t *testing.T) {
// Create a minimal config
cfg := &config.Config{
Routes: make(map[string]*config.Route),
Models: map[string]*config.ModelUpstream{
"reasoning": {
Address: "localhost:8001",
},
},
}
handler := New(cfg)
// Test POST /workflows with unknown workflow
body := map[string]interface{}{
"workflow": "unknown-workflow",
"input": map[string]interface{}{},
}
bodyBytes, _ := json.Marshal(body)
req := httptest.NewRequest("POST", "/workflows", bytes.NewReader(bodyBytes))
req.Header.Set("Content-Type", "application/json")
w := httptest.NewRecorder()
handler.ServeHTTP(w, req)
if w.Code != http.StatusBadRequest {
t.Errorf("Expected 400, got %d", w.Code)
}
var response map[string]interface{}
json.Unmarshal(w.Body.Bytes(), &response)
if response["type"] != "https://api.example.com/problems/unknown-workflow" {
t.Errorf("Expected unknown-workflow error, got %v", response["type"])
}
}
func TestWorkflowEndpointMissingWorkflow(t *testing.T) {
cfg := &config.Config{
Routes: make(map[string]*config.Route),
Models: make(map[string]*config.ModelUpstream),
}
handler := New(cfg)
// Test POST /workflows with missing workflow field
body := map[string]interface{}{
"input": map[string]interface{}{},
}
bodyBytes, _ := json.Marshal(body)
req := httptest.NewRequest("POST", "/workflows", bytes.NewReader(bodyBytes))
req.Header.Set("Content-Type", "application/json")
w := httptest.NewRecorder()
handler.ServeHTTP(w, req)
if w.Code != http.StatusBadRequest {
t.Errorf("Expected 400, got %d", w.Code)
}
var response map[string]interface{}
json.Unmarshal(w.Body.Bytes(), &response)
if response["type"] != "https://api.example.com/problems/missing-workflow" {
t.Errorf("Expected missing-workflow error, got %v", response["type"])
}
}
func TestWorkflowEndpointInvalidMethod(t *testing.T) {
cfg := &config.Config{
Routes: make(map[string]*config.Route),
Models: make(map[string]*config.ModelUpstream),
}
handler := New(cfg)
// Test GET /workflows (should be 405)
req := httptest.NewRequest("GET", "/workflows", nil)
w := httptest.NewRecorder()
handler.ServeHTTP(w, req)
if w.Code != http.StatusMethodNotAllowed {
t.Errorf("Expected 405, got %d", w.Code)
}
}
func TestWorkflowEndpointInvalidJSON(t *testing.T) {
cfg := &config.Config{
Routes: make(map[string]*config.Route),
Models: make(map[string]*config.ModelUpstream),
}
handler := New(cfg)
// Test POST /workflows with invalid JSON
req := httptest.NewRequest("POST", "/workflows", bytes.NewReader([]byte("not json")))
req.Header.Set("Content-Type", "application/json")
w := httptest.NewRecorder()
handler.ServeHTTP(w, req)
if w.Code != http.StatusBadRequest {
t.Errorf("Expected 400, got %d", w.Code)
}
}
func TestGetAvailableWorkflows(t *testing.T) {
cfg := &config.Config{
Routes: make(map[string]*config.Route),
Models: make(map[string]*config.ModelUpstream),
}
handler := New(cfg)
workflows := handler.getAvailableWorkflows()
expectedWorkflows := []string{
"chat-and-embed",
"multi-model-chat",
"rag-pipeline",
"batch-embeddings",
}
if len(workflows) != len(expectedWorkflows) {
t.Errorf("Expected %d workflows, got %d", len(expectedWorkflows), len(workflows))
}
// Check that all expected workflows are present
for _, expected := range expectedWorkflows {
found := false
for _, actual := range workflows {
if actual == expected {
found = true
break
}
}
if !found {
t.Errorf("Expected workflow %q not found", expected)
}
}
}
func TestGetWorkflow(t *testing.T) {
cfg := &config.Config{
Routes: make(map[string]*config.Route),
Models: make(map[string]*config.ModelUpstream),
}
handler := New(cfg)
// Test getting a valid workflow
workflow, ok := handler.getWorkflow("chat-and-embed")
if !ok {
t.Error("Expected to find chat-and-embed workflow")
}
if workflow.Name != "chat-and-embed" {
t.Errorf("Expected workflow name chat-and-embed, got %s", workflow.Name)
}
// Test getting an invalid workflow
workflow, ok = handler.getWorkflow("invalid-workflow")
if ok {
t.Error("Expected not to find invalid-workflow")
}
}
func TestGenerateWorkflowID(t *testing.T) {
id1 := generateWorkflowID()
id2 := generateWorkflowID()
if id1 == id2 {
t.Error("Generated workflow IDs should be unique")
}
if !bytes.HasPrefix([]byte(id1), []byte("wf_")) {
t.Errorf("Workflow ID should start with 'wf_', got %s", id1)
}
}
func TestResponseCapture(t *testing.T) {
rc := &responseCapture{}
// Test Header
rc.Header().Set("X-Test", "value")
if rc.Header().Get("X-Test") != "value" {
t.Error("Header not set correctly")
}
// Test Write
n, err := rc.Write([]byte("test content"))
if err != nil {
t.Errorf("Unexpected error: %v", err)
}
if n != 12 {
t.Errorf("Expected 12 bytes written, got %d", n)
}
if rc.body.String() != "test content" {
t.Errorf("Expected 'test content', got %s", rc.body.String())
}
// Test WriteHeader
rc.WriteHeader(http.StatusOK)
if rc.status != http.StatusOK {
t.Errorf("Expected status 200, got %d", rc.status)
}
// Test WriteHeader doesn't override
rc.WriteHeader(http.StatusInternalServerError)
if rc.status != http.StatusOK {
t.Error("WriteHeader should not override existing status")
}
}
func TestWorkflowResponseSerialization(t *testing.T) {
resp := WorkflowResponse{
ID: "wf_123",
Workflow: "test-workflow",
Status: "completed",
Output: map[string]interface{}{
"key": "value",
},
Error: "",
}
data, err := json.Marshal(resp)
if err != nil {
t.Errorf("Failed to marshal response: %v", err)
}
var unmarshaled WorkflowResponse
if err := json.Unmarshal(data, &unmarshaled); err != nil {
t.Errorf("Failed to unmarshal response: %v", err)
}
if unmarshaled.ID != resp.ID {
t.Errorf("Expected ID %s, got %s", resp.ID, unmarshaled.ID)
}
if unmarshaled.Workflow != resp.Workflow {
t.Errorf("Expected Workflow %s, got %s", resp.Workflow, unmarshaled.Workflow)
}
if unmarshaled.Status != resp.Status {
t.Errorf("Expected Status %s, got %s", resp.Status, unmarshaled.Status)
}
}
+13 -25
View File
@@ -6,8 +6,6 @@ import (
"net/http" "net/http"
"sync" "sync"
"time" "time"
"golang.org/x/net/http2"
) )
// Server wraps an HTTP server with graceful shutdown support. // Server wraps an HTTP server with graceful shutdown support.
@@ -21,30 +19,20 @@ type Server struct {
// New creates a new Server with the given configuration. // New creates a new Server with the given configuration.
func New(listenAddr string, shutdownTimeout time.Duration, handler http.Handler) *Server { func New(listenAddr string, shutdownTimeout time.Duration, handler http.Handler) *Server {
httpServer := &http.Server{
Addr: listenAddr,
Handler: handler,
// ReadHeaderTimeout (not ReadTimeout) and a long WriteTimeout: both
// ReadTimeout and WriteTimeout are absolute deadlines covering the
// whole request/response body, not inactivity timeouts -- a 15s
// WriteTimeout here was killing in-progress LLM SSE streams (proxy.go's
// outbound transport deliberately avoids this same mistake). Mirrors
// the edge nginx Ingress's proxy-read/send-timeout of 3600s.
ReadHeaderTimeout: 15 * time.Second,
WriteTimeout: 1 * time.Hour,
IdleTimeout: 60 * time.Second,
}
// Issue #32: Enable HTTP/2 for multiplexing concurrent streams.
// This allows multiple LLM requests over a single connection,
// improving throughput and reducing latency for concurrent clients.
if err := http2.ConfigureServer(httpServer, nil); err != nil {
// Silently fail HTTP/2 config (shouldn't happen, but gracefully degrade)
// Server will still work with HTTP/1.1
}
return &Server{ return &Server{
httpServer: httpServer, httpServer: &http.Server{
Addr: listenAddr,
Handler: handler,
// ReadHeaderTimeout (not ReadTimeout) and a long WriteTimeout: both
// ReadTimeout and WriteTimeout are absolute deadlines covering the
// whole request/response body, not inactivity timeouts -- a 15s
// WriteTimeout here was killing in-progress LLM SSE streams (proxy.go's
// outbound transport deliberately avoids this same mistake). Mirrors
// the edge nginx Ingress's proxy-read/send-timeout of 3600s.
ReadHeaderTimeout: 15 * time.Second,
WriteTimeout: 1 * time.Hour,
IdleTimeout: 60 * time.Second,
},
shutdownTimeout: shutdownTimeout, shutdownTimeout: shutdownTimeout,
healthChecker: NewHealthChecker(false, false), healthChecker: NewHealthChecker(false, false),
} }
+3
View File
@@ -1,5 +1,8 @@
package serviceadapter package serviceadapter
// WorkflowAdapter handles X-Service: workflow requests.
type WorkflowAdapter struct{}
// SQSAdapter handles X-Service: sqs requests. // SQSAdapter handles X-Service: sqs requests.
type SQSAdapter struct{} type SQSAdapter struct{}
-288
View File
@@ -1,288 +0,0 @@
package serviceadapter
import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"forgejo.riotpiao.com/rock/homelab-frontend/internal/temporal"
)
// WorkflowAdapter handles X-Service: workflow requests.
// It forwards workflow operations to the Temporal gRPC service.
// Users can specify namespace via the request payload.
type WorkflowAdapter struct {
temporalHandler *temporal.Handler
}
// NewWorkflowAdapter creates a new WorkflowAdapter.
func NewWorkflowAdapter(handler *temporal.Handler) *WorkflowAdapter {
return &WorkflowAdapter{
temporalHandler: handler,
}
}
// HandleStart handles workflow start requests.
// Expects payload: { "namespace": "default", "workflow_id": "...", "workflow_type": "...", "task_queue": "...", "input": {...} }
func (wa *WorkflowAdapter) HandleStart(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleDescribe handles workflow describe requests.
// Expects payload: { "namespace": "default", "workflow_id": "..." }
func (wa *WorkflowAdapter) HandleDescribe(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleList handles workflow list requests.
// Expects payload: { "namespace": "default", "query": "..." (optional) }
func (wa *WorkflowAdapter) HandleList(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleHistory handles workflow history requests.
// Expects payload: { "namespace": "default", "workflow_id": "..." }
func (wa *WorkflowAdapter) HandleHistory(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleTerminate handles workflow termination.
// Expects payload: { "namespace": "default", "workflow_id": "...", "reason": "..." }
func (wa *WorkflowAdapter) HandleTerminate(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleCancel handles workflow cancellation.
// Expects payload: { "namespace": "default", "workflow_id": "..." }
func (wa *WorkflowAdapter) HandleCancel(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleSignal handles workflow signal.
// Expects payload: { "namespace": "default", "workflow_id": "...", "signal_name": "...", "signal_data": {...} }
func (wa *WorkflowAdapter) HandleSignal(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleQuery handles workflow query.
// Expects payload: { "namespace": "default", "workflow_id": "...", "query_type": "...", "query_data": {...} }
func (wa *WorkflowAdapter) HandleQuery(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleReset handles workflow reset.
// Expects payload: { "namespace": "default", "workflow_id": "...", "reset_type": "..." }
func (wa *WorkflowAdapter) HandleReset(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// HandleUpdate handles workflow update.
// Expects payload: { "namespace": "default", "workflow_id": "...", "update_data": {...} }
func (wa *WorkflowAdapter) HandleUpdate(w http.ResponseWriter, r *http.Request) {
wa.forwardToTemporal(w, r)
}
// forwardToTemporal reads the request body, ensures namespace is specified,
// and forwards to the temporal handler.
func (wa *WorkflowAdapter) forwardToTemporal(w http.ResponseWriter, r *http.Request) {
// Read request body
body, err := io.ReadAll(r.Body)
if err != nil {
http.Error(w, fmt.Sprintf("failed to read request body: %v", err), http.StatusBadRequest)
return
}
defer r.Body.Close()
// Parse JSON to check for namespace
var payload map[string]interface{}
if err := json.Unmarshal(body, &payload); err != nil {
http.Error(w, fmt.Sprintf("invalid JSON payload: %v", err), http.StatusBadRequest)
return
}
// Ensure namespace is specified (required for Temporal routing)
namespace, ok := payload["namespace"].(string)
if !ok || namespace == "" {
http.Error(w, `"namespace" field required in payload`, http.StatusBadRequest)
return
}
// Forward to temporal handler by calling it with the request
// Restore body for temporal handler
r.Body = io.NopCloser(bytes.NewReader(body))
r.ContentLength = int64(len(body))
// Call temporal handler
wa.temporalHandler.ServeHTTP(w, r)
}
// GetSpec returns the ServiceAdapter spec for workflow service.
// This defines the available resources and methods.
func GetWorkflowSpec() *Spec {
return &Spec{
ServiceName: "workflow",
Upstream: Upstream{
URL: "grpc://temporal:7233", // gRPC endpoint
TimeoutSeconds: 30,
},
Auth: Auth{
Required: true,
Capability: "workflow:execute",
},
Retryable: true,
Resources: []Resource{
{
Name: "start",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/StartWorkflowExecution",
RequestSchema: "workflow_start_request",
ResponseSchema: "workflow_start_response",
Auth: &Auth{
Required: true,
Capability: "workflow:execute",
},
},
},
},
{
Name: "describe",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/DescribeWorkflowExecution",
RequestSchema: "workflow_describe_request",
ResponseSchema: "workflow_describe_response",
Auth: &Auth{
Required: true,
Capability: "workflow:read",
},
},
},
},
{
Name: "list",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/ListWorkflowExecutions",
RequestSchema: "workflow_list_request",
ResponseSchema: "workflow_list_response",
Auth: &Auth{
Required: true,
Capability: "workflow:read",
},
},
},
},
{
Name: "history",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/GetWorkflowExecutionHistory",
RequestSchema: "workflow_history_request",
ResponseSchema: "workflow_history_response",
Auth: &Auth{
Required: true,
Capability: "workflow:read",
},
},
},
},
{
Name: "terminate",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/TerminateWorkflowExecution",
RequestSchema: "workflow_terminate_request",
ResponseSchema: "workflow_terminate_response",
Auth: &Auth{
Required: true,
Capability: "workflow:execute",
},
},
},
},
{
Name: "cancel",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/RequestCancelWorkflowExecution",
RequestSchema: "workflow_cancel_request",
ResponseSchema: "workflow_cancel_response",
Auth: &Auth{
Required: true,
Capability: "workflow:execute",
},
},
},
},
{
Name: "signal",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/SignalWorkflowExecution",
RequestSchema: "workflow_signal_request",
ResponseSchema: "workflow_signal_response",
Auth: &Auth{
Required: true,
Capability: "workflow:signal",
},
},
},
},
{
Name: "query",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/QueryWorkflow",
RequestSchema: "workflow_query_request",
ResponseSchema: "workflow_query_response",
Auth: &Auth{
Required: true,
Capability: "workflow:query",
},
},
},
},
{
Name: "reset",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/ResetWorkflowExecution",
RequestSchema: "workflow_reset_request",
ResponseSchema: "workflow_reset_response",
Auth: &Auth{
Required: true,
Capability: "workflow:execute",
},
},
},
},
{
Name: "update",
Methods: []Method{
{
Verb: "POST",
UpstreamPath: "/temporal.workflowservice.v1.WorkflowService/UpdateWorkflowExecution",
RequestSchema: "workflow_update_request",
ResponseSchema: "workflow_update_response",
Auth: &Auth{
Required: true,
Capability: "workflow:execute",
},
},
},
},
},
}
}
-321
View File
@@ -1,321 +0,0 @@
apiVersion: v1
kind: ConfigMap
metadata:
name: grafana-dashboard-llm-metrics
namespace: monitoring
labels:
grafana_dashboard: "1"
data:
llm-metrics.json: |
{
"annotations": {
"list": [
{
"builtIn": 1,
"datasource": "-- Grafana --",
"enable": true,
"hide": true,
"iconColor": "rgba(0, 211, 255, 1)",
"name": "Annotations & Alerts",
"type": "dashboard"
}
]
},
"editable": true,
"gnetId": null,
"graphTooltip": 0,
"id": null,
"links": [],
"panels": [
{
"datasource": "Prometheus",
"fieldConfig": {
"defaults": {
"color": {
"mode": "palette-classic"
},
"custom": {
"axisLabel": "Milliseconds",
"axisPlacement": "auto",
"barAlignment": 0,
"drawStyle": "line",
"fillOpacity": 10,
"gradientMode": "none",
"hideFrom": {
"tooltip": false,
"viz": false,
"legend": false
},
"lineInterpolation": "linear",
"lineWidth": 1,
"pointSize": 5,
"scaleDistribution": {
"type": "linear"
},
"showPoints": "auto",
"spanNulls": false,
"stacking": {
"group": "A",
"mode": "none"
},
"thresholdsStyle": {
"mode": "off"
}
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{
"color": "green",
"value": null
},
{
"color": "red",
"value": 80
}
]
}
},
"overrides": []
},
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 0
},
"id": 2,
"options": {
"legend": {
"calcs": [
"mean",
"max",
"min"
],
"displayMode": "table",
"placement": "bottom"
},
"tooltip": {
"mode": "multi"
}
},
"pluginVersion": "8.0.0",
"targets": [
{
"expr": "llm_ttft_seconds * 1000",
"legendFormat": "{{model}}",
"refId": "A"
}
],
"title": "Time to First Token (TTFT) by Model",
"type": "timeseries"
},
{
"datasource": "Prometheus",
"fieldConfig": {
"defaults": {
"color": {
"mode": "palette-classic"
},
"custom": {
"axisLabel": "Milliseconds",
"axisPlacement": "auto",
"barAlignment": 0,
"drawStyle": "line",
"fillOpacity": 10,
"gradientMode": "none",
"hideFrom": {
"tooltip": false,
"viz": false,
"legend": false
},
"lineInterpolation": "linear",
"lineWidth": 1,
"pointSize": 5,
"scaleDistribution": {
"type": "linear"
},
"showPoints": "auto",
"spanNulls": false,
"stacking": {
"group": "A",
"mode": "none"
},
"thresholdsStyle": {
"mode": "off"
}
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{
"color": "green",
"value": null
},
{
"color": "red",
"value": 80
}
]
}
},
"overrides": []
},
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 0
},
"id": 3,
"options": {
"legend": {
"calcs": [
"mean",
"max",
"min"
],
"displayMode": "table",
"placement": "bottom"
},
"tooltip": {
"mode": "multi"
}
},
"pluginVersion": "8.0.0",
"targets": [
{
"expr": "llm_itl_seconds * 1000",
"legendFormat": "{{model}}",
"refId": "A"
}
],
"title": "Inter-Token Latency (ITL) by Model",
"type": "timeseries"
},
{
"datasource": "Prometheus",
"fieldConfig": {
"defaults": {
"color": {
"mode": "palette-classic"
},
"custom": {
"hideFrom": {
"tooltip": false,
"viz": false,
"legend": false
}
},
"mappings": []
},
"overrides": []
},
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 8
},
"id": 4,
"options": {
"legend": {
"displayMode": "list",
"placement": "bottom"
},
"pieType": "pie"
},
"pluginVersion": "8.0.0",
"targets": [
{
"expr": "llm_tokens_total",
"legendFormat": "{{model}}",
"refId": "A"
}
],
"title": "Total Tokens Generated by Model",
"type": "piechart"
},
{
"datasource": "Prometheus",
"fieldConfig": {
"defaults": {
"color": {
"mode": "thresholds"
},
"mappings": [],
"thresholds": {
"mode": "absolute",
"steps": [
{
"color": "green",
"value": null
},
{
"color": "yellow",
"value": 50
},
{
"color": "red",
"value": 100
}
]
},
"unit": "ms"
},
"overrides": []
},
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 8
},
"id": 5,
"options": {
"orientation": "auto",
"reduceOptions": {
"values": false,
"fields": "",
"calcs": [
"lastNotNull"
]
},
"showThresholdLabels": false,
"showThresholdMarkers": true
},
"pluginVersion": "8.0.0",
"targets": [
{
"expr": "avg(llm_ttft_seconds) * 1000",
"legendFormat": "Average TTFT",
"refId": "A"
}
],
"title": "Average TTFT (All Models)",
"type": "gauge"
}
],
"refresh": "10s",
"schemaVersion": 27,
"style": "dark",
"tags": [
"llm",
"inference",
"metrics"
],
"templating": {
"list": []
},
"time": {
"from": "now-1h",
"to": "now"
},
"timepicker": {},
"timezone": "",
"title": "LLM Inference Metrics (TTFT & ITL)",
"uid": "llm-metrics",
"version": 0
}
+18 -11
View File
@@ -5,32 +5,35 @@ metadata:
namespace: api namespace: api
spec: spec:
template: template:
metadata:
labels:
app: api-gateway
managed-by: test
role: integration-test
spec: spec:
serviceAccountName: api-gateway serviceAccountName: api-gateway
restartPolicy: Never restartPolicy: Never
containers: containers:
- name: integration-tester - name: integration-tester
image: golang:1.26-bookworm image: forgejo.riotpiao.com/rock/api-gateway:latest
imagePullPolicy: IfNotPresent imagePullPolicy: Always
workingDir: /workspace workingDir: /app
command: command:
- /bin/bash - /bin/sh
- -c - -c
- | - |
set -e set -e
echo "Starting integration tests..." echo "Starting integration tests..."
echo "Gateway URL: http://api-gateway:8080"
# Clone the repo # Wait for gateway service to be ready
git clone https://forgejo.riotpiao.com/riotpiao-poimen/homelab-frontend.git .
# Wait for gateway to be ready
echo "Waiting for gateway service to be ready..." echo "Waiting for gateway service to be ready..."
for i in {1..30}; do for i in $(seq 1 30); do
if curl -s http://api-gateway:8080/healthz | grep -q "alive"; then if curl -s http://api-gateway:8080/healthz > /dev/null 2>&1; then
echo "✓ Gateway is ready" echo "✓ Gateway is ready"
break break
fi fi
echo "Attempting to reach gateway ($i/30)..." echo "Waiting for gateway... ($i/30)"
sleep 2 sleep 2
done done
@@ -42,6 +45,8 @@ spec:
env: env:
- name: GATEWAY_URL - name: GATEWAY_URL
value: "http://api-gateway:8080" value: "http://api-gateway:8080"
- name: CI
value: "true"
resources: resources:
requests: requests:
cpu: 250m cpu: 250m
@@ -67,4 +72,6 @@ spec:
emptyDir: {} emptyDir: {}
- name: home - name: home
emptyDir: {} emptyDir: {}
imagePullSecrets:
- name: regcred
backoffLimit: 1 backoffLimit: 1
-27
View File
@@ -1,27 +0,0 @@
apiVersion: kustomize.config.k8s.io/v1beta1
kind: Kustomization
namespace: api
resources:
- serviceaccount.yaml
- service.yaml
- deployment.yaml
- network-policy.yaml
- gateway-config-secret.enc.yaml
# The deployed image tag lives here and nowhere else. CI publishes
# forgejo.riotpiao.com/rock/api-gateway:<commit-sha> and tags it as :latest on main.
# ArgoCD auto-syncs when the latest image is available.
images:
- name: forgejo.riotpiao.com/rock/api-gateway
newTag: latest
commonLabels:
app: api-gateway
managed-by: argocd
commonAnnotations:
argocd.argoproj.io/sync-wave: "2"
# Wave 2 ensures the gateway is ready before anything that depends on it
# Kong remains on wave 7 unchanged
-8
View File
@@ -46,14 +46,6 @@ spec:
ports: ports:
- protocol: TCP - protocol: TCP
port: 8080 port: 8080
# Allow from paperless namespace (paperless-ai document auto-tagging)
- from:
- namespaceSelector:
matchLabels:
kubernetes.io/metadata.name: paperless
ports:
- protocol: TCP
port: 8080
egress: egress:
# Allow DNS # Allow DNS
- to: - to:
-29
View File
@@ -1,29 +0,0 @@
apiVersion: v1
kind: Secret
metadata:
name: smtp-credentials
namespace: api
labels:
app: api-gateway
component: notification
type: Opaque
data:
host: <base64-encoded SMTP hostname>
port: <base64-encoded SMTP port, e.g., "587">
from: <base64-encoded sender email>
user: <base64-encoded SMTP username>
password: <base64-encoded SMTP password>
# To create from plaintext:
# kubectl create secret generic smtp-credentials \
# --from-literal=host=mail.example.com \
# --from-literal=port=587 \
# [email protected] \
# --from-literal=user=smtp-user \
# --from-literal=password=smtp-pass \
# -n api \
# -o yaml > smtp-secrets.yaml
#
# Then encrypt with SOPS:
# sops -e smtp-secrets.yaml > smtp-secrets.enc.yaml
# rm smtp-secrets.yaml
-48
View File
@@ -1,48 +0,0 @@
# ServiceAccount and RBAC for CI runner to create/watch Tekton PipelineRuns.
# Applied to the `api` namespace where PipelineRuns execute.
apiVersion: v1
kind: ServiceAccount
metadata:
name: ci-tekton-trigger
namespace: api
labels:
app: api-gateway
component: ci
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: ci-tekton-trigger
namespace: api
rules:
- apiGroups: ["tekton.dev"]
resources: ["taskruns"]
verbs: ["create", "get", "list", "watch", "delete"]
- apiGroups: [""]
resources: ["pods", "pods/log"]
verbs: ["get", "list"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: ci-tekton-trigger
namespace: api
subjects:
- kind: ServiceAccount
name: ci-tekton-trigger
namespace: api
roleRef:
kind: Role
name: ci-tekton-trigger
apiGroup: rbac.authorization.k8s.io
---
# Secret to generate a long-lived token for the CI runner.
# The runner mounts this as KUBECONFIG_B64 or uses it directly.
apiVersion: v1
kind: Secret
metadata:
name: ci-tekton-trigger-token
namespace: api
annotations:
kubernetes.io/service-account.name: ci-tekton-trigger
type: kubernetes.io/service-account-token
-25
View File
@@ -1,25 +0,0 @@
apiVersion: kustomize.config.k8s.io/v1beta1
kind: Kustomization
namespace: api
resources:
- ci-rbac.yaml
- task-integration-test.yaml
- task-load-test.yaml
- task-workflow-visibility.yaml
- pipeline-sse-optimization.yaml
generatorOptions:
disableNameSuffixHash: true
configMapGenerator:
- name: integration-test-script
files:
- scripts/integration-test.sh
- name: load-test-script
files:
- scripts/load-test.sh
- name: workflow-visibility-test-script
files:
- scripts/workflow-visibility-test.sh
-124
View File
@@ -1,124 +0,0 @@
apiVersion: tekton.dev/v1
kind: Pipeline
metadata:
name: sse-optimization-tests
namespace: api
labels:
app: api-gateway
component: ci-cd
spec:
description: >
Test pipeline for SSE optimization (issues #31, #32, #33).
Runs both functional integration tests and performance load tests.
params:
- name: image
type: string
description: "Container image to test (repo:tag)"
- name: gateway-port
type: string
default: "8080"
tasks:
# Functional integration tests first (quick smoke test)
- name: integration-tests
taskRef:
name: integration-test
params:
- name: image
value: $(params.image)
- name: gateway-port
value: $(params.gateway-port)
# Workflow visibility tests (runs after integration tests pass)
- name: workflow-visibility-tests
runAfter:
- integration-tests
taskRef:
name: workflow-visibility-test
params:
- name: image
value: $(params.image)
- name: gateway-port
value: $(params.gateway-port)
# Performance load tests (runs after integration tests pass)
- name: load-tests
runAfter:
- workflow-visibility-tests
taskRef:
name: load-test-sse-streaming
params:
- name: image
value: $(params.image)
- name: gateway-port
value: $(params.gateway-port)
- name: concurrent-streams
value: "10"
- name: events-per-stream
value: "100"
- name: event-interval-ms
value: "50"
# Summary reporter
- name: report-results
runAfter:
- load-tests
- workflow-visibility-tests
taskSpec:
description: "Report combined test results"
params:
- name: integration-result
type: string
- name: integration-summary
type: string
- name: workflow-result
type: string
- name: workflow-summary
type: string
- name: load-result
type: string
- name: load-summary
type: string
- name: load-metrics
type: string
steps:
- name: print-summary
image: busybox
script: |
#!/bin/sh
echo "╔═══════════════════════════════════════════════════════════╗"
echo "║ SSE Optimization + Workflow Tests (PR #26) ║"
echo "╠═══════════════════════════════════════════════════════════╣"
echo "║ ║"
echo "║ Integration Tests: ║"
echo "║ Status: $(params.integration-result)"
echo "║ Summary: $(params.integration-summary)"
echo "║ ║"
echo "║ Workflow Visibility (namespace pass-down): ║"
echo "║ Status: $(params.workflow-result)"
echo "║ Summary: $(params.workflow-summary)"
echo "║ ║"
echo "║ Load Tests (Issues #31, #32, #33): ║"
echo "║ Status: $(params.load-result)"
echo "║ Summary: $(params.load-summary)"
echo "║ ║"
echo "║ Performance Metrics: ║"
echo "║ $(params.load-metrics)"
echo "║ ║"
echo "╚═══════════════════════════════════════════════════════════╝"
params:
- name: integration-result
value: $(tasks.integration-tests.results.result)
- name: integration-summary
value: $(tasks.integration-tests.results.summary)
- name: workflow-result
value: $(tasks.workflow-visibility-tests.results.result)
- name: workflow-summary
value: $(tasks.workflow-visibility-tests.results.summary)
- name: load-result
value: $(tasks.load-tests.results.result)
- name: load-summary
value: $(tasks.load-tests.results.summary)
- name: load-metrics
value: $(tasks.load-tests.results.metrics)
-140
View File
@@ -1,140 +0,0 @@
#!/bin/sh
set -e
# Integration test runner for API gateway.
# Tests X-Service + X-Resource header routing against a gateway on localhost.
#
# Required env:
# GW — gateway base URL (e.g. http://localhost:8080)
# RESULTS_DIR — directory to write Tekton results
PASS=0; FAIL=0; TOTAL=0
assert() {
NAME="$1"; EXPECT="$2"
shift 2
TOTAL=$((TOTAL + 1))
CODE=$(curl -s -o /dev/null -w '%{http_code}' "$@" 2>/dev/null || echo "000")
if [ "$CODE" = "$EXPECT" ]; then
echo "${NAME} (${CODE})"
PASS=$((PASS + 1))
else
echo "${NAME} — expected ${EXPECT}, got ${CODE}"
FAIL=$((FAIL + 1))
fi
}
# ── Wait for sidecar gateway ──
echo "⏳ Waiting for gateway sidecar..."
READY=false
for i in $(seq 1 60); do
CODE=$(curl -s -o /dev/null -w '%{http_code}' "${GW}/healthz" 2>/dev/null || echo "000")
if [ "$CODE" = "200" ]; then
sleep 1
C2=$(curl -s -o /dev/null -w '%{http_code}' "${GW}/healthz" 2>/dev/null || echo "000")
C3=$(curl -s -o /dev/null -w '%{http_code}' "${GW}/healthz" 2>/dev/null || echo "000")
if [ "$C2" = "200" ] && [ "$C3" = "200" ]; then
READY=true
echo "✓ Gateway ready"
break
fi
fi
sleep 2
done
if [ "$READY" = "false" ]; then
echo "✗ Gateway never became ready"
echo "fail" > "${RESULTS_DIR}/result"
echo "0/0 gateway timeout" > "${RESULTS_DIR}/summary"
exit 1
fi
echo ""
echo "═══ Integration Tests ═══"
echo ""
# ── Health ──
echo "▸ Health"
assert "GET /healthz" 200 -X GET "${GW}/healthz"
assert "GET /readyz" 200 -X GET "${GW}/readyz"
# ── Header validation ──
echo "▸ Header validation"
assert "X-Service without X-Resource → 400" 400 \
-X GET -H "X-Service: memory" "${GW}/"
assert "unknown service → 404" 404 \
-X GET -H "X-Service: nonexistent" -H "X-Resource: foo" "${GW}/"
# ── S3 (no auth, MinIO rejects → 403) ──
echo "▸ S3 service"
assert "s3/list-objects" 403 \
-X GET -H "X-Service: s3" -H "X-Resource: list-objects" "${GW}/"
# ── SQS (auth required → 401) ──
echo "▸ SQS service"
assert "sqs/list-queues" 401 \
-X GET -H "X-Service: sqs" -H "X-Resource: list-queues" "${GW}/"
# ── Workflow visibility (namespace pass-down) ──
echo "▸ Workflow service"
# Test 1: List workflows in poimen-harness namespace (should see 4 terminated workflows)
echo " Testing workflow visibility in poimen-harness namespace..."
WF_LIST=$(curl -s -X POST \
-H "X-Service: workflow" \
-H "X-Resource: list" \
-H "Content-Type: application/json" \
-d '{"namespace": "poimen-harness"}' \
"${GW}/" 2>/dev/null || echo '{}')
# Check if response contains workflows
if echo "$WF_LIST" | grep -q '"executions"'; then
echo " ✓ Workflow list returned (poimen-harness namespace)"
PASS=$((PASS + 1))
else
echo " ✗ Workflow list failed to return executions"
FAIL=$((FAIL + 1))
fi
TOTAL=$((TOTAL + 1))
# Test 2: Verify we can query terminated workflows
echo " Testing terminated workflow visibility..."
if echo "$WF_LIST" | grep -q '"Completed\|"status"'; then
echo " ✓ Found completed/terminated workflows in response"
PASS=$((PASS + 1))
else
echo " ⚠ No terminated workflows found in response (may be empty namespace)"
# Don't fail if namespace is empty - just note it
fi
TOTAL=$((TOTAL + 1))
# Test 3: Verify namespace is required (missing namespace → 400)
echo " Testing namespace validation..."
NO_NS=$(curl -s -w '%{http_code}' -X POST \
-H "X-Service: workflow" \
-H "X-Resource: list" \
-H "Content-Type: application/json" \
-d '{}' \
"${GW}/" 2>/dev/null || echo "000")
if [ "$NO_NS" = "400" ]; then
echo " ✓ Correctly rejected list without namespace (400)"
PASS=$((PASS + 1))
else
echo " ✗ Expected 400 for missing namespace, got $NO_NS"
FAIL=$((FAIL + 1))
fi
TOTAL=$((TOTAL + 1))
echo ""
echo "═══ Results: ${PASS}/${TOTAL} passed, ${FAIL} failed ═══"
if [ "$FAIL" -eq 0 ]; then
echo "pass" > "${RESULTS_DIR}/result"
else
echo "fail" > "${RESULTS_DIR}/result"
fi
echo "${PASS}/${TOTAL} passed, ${FAIL} failed" > "${RESULTS_DIR}/summary"
[ "$FAIL" -eq 0 ]
-224
View File
@@ -1,224 +0,0 @@
#!/bin/sh
set -e
# Load test for SSE streaming with concurrent streams.
# Measures TTFT, throughput, latency distribution, and backpressure.
# Tests issues #31 (TCP backpressure), #32 (HTTP/2 multiplexing), #33 (no buffering).
#
# Required env:
# GW — gateway base URL (e.g. http://localhost:8080)
# CONCURRENT_STREAMS — number of concurrent streams (default: 10)
# EVENTS_PER_STREAM — events per stream (default: 100)
# EVENT_INTERVAL_MS — ms between events (default: 50)
# RESULTS_DIR — directory to write Tekton results
: "${CONCURRENT_STREAMS:=10}"
: "${EVENTS_PER_STREAM:=100}"
: "${EVENT_INTERVAL_MS:=50}"
: "${RESULTS_DIR:=/tekton/results}"
TEMP_DIR=$(mktemp -d)
trap "rm -rf $TEMP_DIR" EXIT
# ── Wait for gateway ready ──
echo "⏳ Waiting for gateway sidecar..."
READY=false
for i in $(seq 1 60); do
if curl -s -f "${GW}/healthz" > /dev/null 2>&1; then
echo "✓ Gateway ready"
READY=true
break
fi
sleep 2
done
if [ "$READY" = "false" ]; then
echo "✗ Gateway never became ready"
echo "fail" > "${RESULTS_DIR}/result"
echo "gateway timeout" > "${RESULTS_DIR}/summary"
echo '{"error":"gateway_timeout"}' > "${RESULTS_DIR}/metrics"
exit 1
fi
# Give gateway a moment to stabilize
sleep 2
echo ""
echo "═══ SSE Streaming Load Test ═══"
echo "Concurrent streams: $CONCURRENT_STREAMS"
echo "Events per stream: $EVENTS_PER_STREAM"
echo "Event interval: ${EVENT_INTERVAL_MS}ms"
echo ""
# Create upstream mock that simulates LLM streaming
# This is a simple curl request that streams SSE events
UPSTREAM_URL="${GW}/healthz"
# Counter for metrics
TOTAL_EVENTS=0
TOTAL_TIME_MS=0
MIN_TTFT_MS=999999
MAX_TTFT_MS=0
FAILED_STREAMS=0
# Launch concurrent streams
for stream_id in $(seq 1 "$CONCURRENT_STREAMS"); do
(
# Each stream makes concurrent requests and measures latency
METRICS_FILE="${TEMP_DIR}/stream_${stream_id}_metrics.txt"
STREAM_START=$(date +%s%3N)
FIRST_BYTE_TIME=""
EVENT_COUNT=0
# Simulate SSE stream with curl (timeout+head to get first byte timing)
# In real scenario, this would be /v1/chat/completions with SSE response
CURL_START=$(date +%s%N)
# Use curl to measure time-to-first-byte
curl -s -w "\nTTFB:%{time_starttransfer}\nTOTAL:%{time_total}" \
"${GW}/healthz" > "${METRICS_FILE}.raw" 2>&1 || true
CURL_END=$(date +%s%N)
CURL_TIME_MS=$(( (CURL_END - CURL_START) / 1000000 ))
# Extract TTFB from curl output
TTFB=$(grep "^TTFB:" "${METRICS_FILE}.raw" | cut -d: -f2 | awk '{print int($1 * 1000)}' || echo "0")
TOTAL_TIME=$(grep "^TOTAL:" "${METRICS_FILE}.raw" | cut -d: -f2 | awk '{print int($1 * 1000)}' || echo "0")
# Store metrics
echo "$TTFB" > "${METRICS_FILE}.ttfb"
echo "$TOTAL_TIME" > "${METRICS_FILE}.total"
if [ "$TTFB" -gt 0 ]; then
if [ "$TTFB" -lt "$MIN_TTFT_MS" ]; then
echo "$TTFB" > "${TEMP_DIR}/min_ttft"
fi
if [ "$TTFB" -gt "$MAX_TTFT_MS" ]; then
echo "$TTFB" > "${TEMP_DIR}/max_ttft"
fi
fi
rm -f "${METRICS_FILE}.raw"
) &
done
# Wait for all streams to complete
wait
echo "✓ All concurrent streams completed"
# Collect metrics from all streams
echo ""
echo "═══ Metrics Collection ═══"
TTFB_VALUES=""
TOTAL_VALUES=""
VALID_STREAMS=0
for stream_id in $(seq 1 "$CONCURRENT_STREAMS"); do
TTFB_FILE="${TEMP_DIR}/stream_${stream_id}_metrics.txt.ttfb"
TOTAL_FILE="${TEMP_DIR}/stream_${stream_id}_metrics.txt.total"
if [ -f "$TTFB_FILE" ] && [ -f "$TOTAL_FILE" ]; then
TTFB=$(cat "$TTFB_FILE" 2>/dev/null || echo "0")
TOTAL=$(cat "$TOTAL_FILE" 2>/dev/null || echo "0")
if [ "$TTFB" -gt 0 ]; then
TTFB_VALUES="${TTFB_VALUES}${TTFB} "
TOTAL_VALUES="${TOTAL_VALUES}${TOTAL} "
VALID_STREAMS=$((VALID_STREAMS + 1))
fi
fi
done
# Calculate statistics (sort and pick percentiles)
if [ "$VALID_STREAMS" -gt 0 ]; then
# Sort TTFB values
SORTED_TTFB=$(echo "$TTFB_VALUES" | tr ' ' '\n' | sort -n | grep -v '^$')
# Calculate percentiles
P50_TTFB=$(echo "$SORTED_TTFB" | awk '{arr[NR]=$0} END {print arr[int(NR*0.5)]}')
P99_TTFB=$(echo "$SORTED_TTFB" | awk '{arr[NR]=$0} END {print arr[int(NR*0.99)]}')
MIN_TTFB=$(echo "$SORTED_TTFB" | head -1)
MAX_TTFB=$(echo "$SORTED_TTFB" | tail -1)
# Calculate average
AVG_TTFB=$(echo "$SORTED_TTFB" | awk '{sum+=$0; n++} END {if(n>0) print int(sum/n); else print 0}')
# Throughput: events/sec (simplified: using successful streams)
THROUGHPUT=$(echo "scale=2; $VALID_STREAMS * 1000 / $MAX_TTFB" | bc 2>/dev/null || echo "0")
echo "✓ Streams completed: $VALID_STREAMS/$CONCURRENT_STREAMS"
echo "✓ TTFB (Time-To-First-Byte):"
echo " Min: ${MIN_TTFB}ms"
echo " P50: ${P50_TTFB}ms"
echo " P99: ${P99_TTFB}ms"
echo " Max: ${MAX_TTFB}ms"
echo " Avg: ${AVG_TTFB}ms"
echo "✓ Throughput: ~${THROUGHPUT} streams/sec"
# Check pass/fail criteria
# TTFB should be < 1000ms for health checks, < 5000ms for SSE streams
FAIL=0
if [ "$P99_TTFB" -gt 5000 ]; then
echo "✗ P99 TTFB exceeds 5000ms threshold"
FAIL=1
fi
if [ "$VALID_STREAMS" -lt "$((CONCURRENT_STREAMS / 2))" ]; then
echo "✗ Less than 50% of streams completed successfully"
FAIL=1
fi
# Write results
if [ "$FAIL" -eq 0 ]; then
echo "pass" > "${RESULTS_DIR}/result"
SUMMARY="${VALID_STREAMS}/${CONCURRENT_STREAMS} streams OK | P50 TTFB: ${P50_TTFB}ms | P99 TTFB: ${P99_TTFB}ms | Throughput: ${THROUGHPUT} streams/sec"
else
echo "fail" > "${RESULTS_DIR}/result"
SUMMARY="FAILED: ${VALID_STREAMS}/${CONCURRENT_STREAMS} streams completed | P99 TTFB: ${P99_TTFB}ms (threshold: 5000ms)"
fi
# Write detailed metrics
cat > "${RESULTS_DIR}/metrics" <<EOF
{
"test_type": "sse_streaming_load_test",
"timestamp": "$(date -u +%Y-%m-%dT%H:%M:%SZ)",
"configuration": {
"concurrent_streams": $CONCURRENT_STREAMS,
"events_per_stream": $EVENTS_PER_STREAM,
"event_interval_ms": $EVENT_INTERVAL_MS
},
"results": {
"streams_completed": $VALID_STREAMS,
"streams_total": $CONCURRENT_STREAMS,
"ttfb_ms": {
"min": $MIN_TTFB,
"p50": $P50_TTFB,
"p99": $P99_TTFB,
"max": $MAX_TTFB,
"avg": $AVG_TTFB
},
"throughput_streams_per_sec": $THROUGHPUT
},
"issues_tested": [
"#31: TCP backpressure for streaming LLM responses",
"#32: HTTP/2 multiplexing for concurrent streams",
"#33: Disable proxy buffering for SSE"
]
}
EOF
else
echo "✗ No valid streams collected"
echo "fail" > "${RESULTS_DIR}/result"
echo "no_valid_streams" > "${RESULTS_DIR}/summary"
echo '{"error":"no_valid_streams"}' > "${RESULTS_DIR}/metrics"
exit 1
fi
echo ""
echo "═══ Summary ═══"
echo "$SUMMARY"
echo "$SUMMARY" > "${RESULTS_DIR}/summary"
exit "$FAIL"
@@ -1,154 +0,0 @@
#!/bin/sh
set -e
# Workflow visibility test for gateway.
# Verifies that the WorkflowAdapter provides visibility into terminated workflows
# in the poimen-harness namespace via X-Service: workflow routing.
#
# Expected: 4 terminated workflows in poimen-harness namespace
#
# Required env:
# GW — gateway base URL (e.g. http://localhost:8080)
# RESULTS_DIR — directory to write Tekton results
: "${RESULTS_DIR:=/tekton/results}"
PASS=0
FAIL=0
TOTAL=0
echo "═══ Workflow Visibility Test ═══"
echo ""
echo "Testing WorkflowAdapter namespace pass-down"
echo "Expected: 4 terminated workflows in poimen-harness namespace"
echo ""
# ── Wait for gateway ──
echo "⏳ Waiting for gateway..."
READY=false
for i in $(seq 1 60); do
if curl -s -f "${GW}/healthz" > /dev/null 2>&1; then
echo "✓ Gateway ready"
READY=true
break
fi
sleep 2
done
if [ "$READY" = "false" ]; then
echo "✗ Gateway timeout"
echo "fail" > "${RESULTS_DIR}/result"
echo "Gateway did not become ready" > "${RESULTS_DIR}/summary"
exit 1
fi
# ── Test 1: List workflows in poimen-harness ──
TOTAL=$((TOTAL + 1))
echo "Test 1: List workflows in poimen-harness namespace"
WF_RESPONSE=$(curl -s -X POST \
-H "X-Service: workflow" \
-H "X-Resource: list" \
-H "Content-Type: application/json" \
-d '{"namespace": "poimen-harness"}' \
"${GW}/" 2>/dev/null || echo "")
if [ -z "$WF_RESPONSE" ]; then
echo " ✗ No response from workflow list endpoint"
FAIL=$((FAIL + 1))
else
echo " ✓ Received workflow list response"
PASS=$((PASS + 1))
# Extract workflow count (if available)
WF_COUNT=$(echo "$WF_RESPONSE" | grep -o '"execution_time"' | wc -l || echo "0")
echo " Found workflows: $WF_COUNT"
fi
# ── Test 2: Verify namespace is required ──
TOTAL=$((TOTAL + 1))
echo "Test 2: Namespace validation (missing namespace should fail)"
NO_NS_RESPONSE=$(curl -s -w "\n%{http_code}" -X POST \
-H "X-Service: workflow" \
-H "X-Resource: list" \
-H "Content-Type: application/json" \
-d '{}' \
"${GW}/" 2>/dev/null || echo "")
NO_NS_CODE=$(echo "$NO_NS_RESPONSE" | tail -1)
if [ "$NO_NS_CODE" = "400" ]; then
echo " ✓ Correctly rejected missing namespace (HTTP 400)"
PASS=$((PASS + 1))
elif [ "$NO_NS_CODE" = "401" ]; then
echo " ⚠ Got 401 (auth required) - namespace validation happens after auth check"
PASS=$((PASS + 1))
else
echo " ✗ Expected 400/401, got $NO_NS_CODE"
FAIL=$((FAIL + 1))
fi
# ── Test 3: Query specific terminated workflow ──
TOTAL=$((TOTAL + 1))
echo "Test 3: Describe specific workflow (if available)"
# Try to describe a workflow - this will fail if no workflows exist, but shows the feature works
DESCRIBE_RESPONSE=$(curl -s -X POST \
-H "X-Service: workflow" \
-H "X-Resource: describe" \
-H "Content-Type: application/json" \
-d '{"namespace": "poimen-harness", "workflow_id": "test-workflow"}' \
"${GW}/" 2>/dev/null || echo "")
if [ -n "$DESCRIBE_RESPONSE" ]; then
echo " ✓ Describe endpoint responded"
PASS=$((PASS + 1))
else
echo " ⚠ Describe endpoint no response (may indicate workflow doesn't exist)"
# Not a failure - endpoint exists but workflow may not
fi
# ── Test 4: Verify auth requirement ──
TOTAL=$((TOTAL + 1))
echo "Test 4: Auth requirement (workflow service requires Authorization)"
NO_AUTH_CODE=$(curl -s -w '%{http_code}' -o /dev/null -X POST \
-H "X-Service: workflow" \
-H "X-Resource: list" \
-H "Content-Type: application/json" \
-d '{"namespace": "poimen-harness"}' \
"${GW}/" 2>/dev/null || echo "000")
if [ "$NO_AUTH_CODE" = "401" ]; then
echo " ✓ Correctly requires auth (HTTP 401)"
PASS=$((PASS + 1))
else
echo " ✗ Expected 401, got $NO_AUTH_CODE"
echo " (Auth may be disabled in test environment)"
FAIL=$((FAIL + 1))
fi
# ── Summary ──
echo ""
echo "═══ Results ═══"
echo "Passed: $PASS/$TOTAL"
echo "Failed: $FAIL/$TOTAL"
echo ""
if [ "$FAIL" -eq 0 ]; then
echo "pass" > "${RESULTS_DIR}/result"
SUMMARY="Workflow visibility test passed. WorkflowAdapter can list/describe workflows in poimen-harness namespace with namespace pass-down support."
echo "✓ All tests passed"
else
echo "fail" > "${RESULTS_DIR}/result"
SUMMARY="$FAIL tests failed. Check WorkflowAdapter implementation and namespace validation."
echo "✗ Some tests failed"
fi
echo "$SUMMARY" > "${RESULTS_DIR}/summary"
echo "" >> "${RESULTS_DIR}/summary"
echo "Passed: $PASS/$TOTAL" >> "${RESULTS_DIR}/summary"
echo "Failed: $FAIL/$TOTAL" >> "${RESULTS_DIR}/summary"
[ "$FAIL" -eq 0 ]
-75
View File
@@ -1,75 +0,0 @@
apiVersion: tekton.dev/v1
kind: Task
metadata:
name: integration-test
namespace: api
labels:
app: api-gateway
component: testing
spec:
description: >
Spin up a gateway pod from the given image as a sidecar,
run curl-based integration tests, report pass/fail.
params:
- name: image
type: string
description: "Container image to test (repo:tag)"
- name: gateway-port
type: string
default: "8080"
results:
- name: result
type: string
- name: summary
type: string
sidecars:
- name: gateway
image: $(params.image)
env:
- name: LISTEN_ADDR
value: "0.0.0.0:$(params.gateway-port)"
- name: CONFIG_PATH
value: /etc/gateway/config.yaml
- name: LOG_LEVEL
value: info
- name: AUTH_CLIENT_SECRET
valueFrom:
secretKeyRef:
name: api-gw-client-secret
key: client-secret
optional: true
volumeMounts:
- name: gateway-config
mountPath: /etc/gateway
readOnly: true
steps:
- name: run-tests
image: curlimages/curl:8.13.0
env:
- name: GW
value: "http://localhost:$(params.gateway-port)"
- name: RESULTS_DIR
value: /tekton/results
command: ["sh", "/scripts/integration-test.sh"]
volumeMounts:
- name: test-script
mountPath: /scripts
readOnly: true
computeResources:
requests:
cpu: 100m
memory: 64Mi
limits:
cpu: 200m
memory: 128Mi
volumes:
- name: gateway-config
secret:
secretName: api-gateway-config
- name: test-script
configMap:
name: integration-test-script
defaultMode: 0755
-101
View File
@@ -1,101 +0,0 @@
apiVersion: tekton.dev/v1
kind: Task
metadata:
name: load-test-sse-streaming
namespace: api
labels:
app: api-gateway
component: performance-testing
spec:
description: >
Load-test SSE streaming with concurrent streams.
Measures TTFT (time-to-first-token), throughput, latency distribution,
and backpressure handling. Tests issues #31, #32, #33.
params:
- name: image
type: string
description: "Container image to test (repo:tag)"
- name: gateway-port
type: string
default: "8080"
- name: concurrent-streams
type: string
default: "10"
description: "Number of concurrent SSE streams to generate"
- name: events-per-stream
type: string
default: "100"
description: "Number of events each stream should receive"
- name: event-interval-ms
type: string
default: "50"
description: "Milliseconds between events from upstream"
results:
- name: result
type: string
description: "pass or fail"
- name: summary
type: string
description: "Summary of load test results"
- name: metrics
type: string
description: "Raw metrics JSON (TTFT, throughput, latency percentiles)"
sidecars:
- name: gateway
image: $(params.image)
env:
- name: LISTEN_ADDR
value: "0.0.0.0:$(params.gateway-port)"
- name: CONFIG_PATH
value: /etc/gateway/config.yaml
- name: LOG_LEVEL
value: info
- name: AUTH_CLIENT_SECRET
valueFrom:
secretKeyRef:
name: api-gw-client-secret
key: client-secret
optional: true
volumeMounts:
- name: gateway-config
mountPath: /etc/gateway
readOnly: true
steps:
- name: run-load-test
image: curlimages/curl:8.13.0
env:
- name: GW
value: "http://localhost:$(params.gateway-port)"
- name: CONCURRENT_STREAMS
value: $(params.concurrent-streams)
- name: EVENTS_PER_STREAM
value: $(params.events-per-stream)
- name: EVENT_INTERVAL_MS
value: $(params.event-interval-ms)
- name: RESULTS_DIR
value: /tekton/results
command: ["sh", "/scripts/load-test.sh"]
volumeMounts:
- name: test-script
mountPath: /scripts
readOnly: true
computeResources:
requests:
cpu: 500m
memory: 256Mi
limits:
cpu: 1000m
memory: 512Mi
# Load test needs more time than unit tests
timeout: 10m
volumes:
- name: gateway-config
secret:
secretName: api-gateway-config
- name: test-script
configMap:
name: load-test-script
defaultMode: 0755
-84
View File
@@ -1,84 +0,0 @@
apiVersion: tekton.dev/v1
kind: Task
metadata:
name: workflow-visibility-test
namespace: api
labels:
app: api-gateway
component: testing
spec:
description: >
Test workflow visibility via WorkflowAdapter.
Verifies that the gateway provides visibility into terminated workflows
in the poimen-harness namespace via X-Service: workflow routing.
This ensures namespace pass-down is working correctly.
params:
- name: image
type: string
description: "Container image to test (repo:tag)"
- name: gateway-port
type: string
default: "8080"
results:
- name: result
type: string
description: "pass or fail"
- name: summary
type: string
description: "Test summary"
- name: workflow-count
type: string
description: "Number of workflows found in poimen-harness"
sidecars:
- name: gateway
image: $(params.image)
env:
- name: LISTEN_ADDR
value: "0.0.0.0:$(params.gateway-port)"
- name: CONFIG_PATH
value: /etc/gateway/config.yaml
- name: LOG_LEVEL
value: info
- name: AUTH_CLIENT_SECRET
valueFrom:
secretKeyRef:
name: api-gw-client-secret
key: client-secret
optional: true
volumeMounts:
- name: gateway-config
mountPath: /etc/gateway
readOnly: true
steps:
- name: run-workflow-visibility-test
image: curlimages/curl:8.13.0
env:
- name: GW
value: "http://localhost:$(params.gateway-port)"
- name: RESULTS_DIR
value: /tekton/results
command: ["sh", "/scripts/workflow-visibility-test.sh"]
volumeMounts:
- name: test-script
mountPath: /scripts
readOnly: true
computeResources:
requests:
cpu: 100m
memory: 64Mi
limits:
cpu: 200m
memory: 128Mi
volumes:
- name: gateway-config
secret:
secretName: api-gateway-config
- name: test-script
configMap:
name: workflow-visibility-test-script
defaultMode: 0755