Topic Creation Methods:
1. Manual CLI (Fastest - 2 min):
└─ kubectl port-forward + curl POST /v1/queues
└─ k8s/config/create-kmsvc-topics.sh (interactive)
2. Kubernetes Job (Automated - 1 min):
└─ kubectl apply kmsvc-topics-job.yaml
└─ Runs once, creates topics if not exist
└─ Can re-run safely
3. Terraform (IaC - 2 min):
└─ terraform apply -target=null_resource.create_kmsvc_topics
└─ Tracks topic creation in .tfstate
└─ Idempotent
4. Shell Script (Interactive - 1 min):
└─ ./create-kmsvc-topics.sh
└─ Auto port-forward or manual mode
└─ Color output + progress logging
Topics Created:
1. poimen-memory-dlq (DLQ for extraction + webhook + agent)
├─ Retention: 14 days (1,209,600 seconds)
├─ Visibility: 5 minutes (300 seconds)
├─ Messages: {id, type, workflow_id, error, timestamp, ...}
└─ Consumer: queue_worker_dlq.rs::DlqHandler
2. poimen-memory-metric-dlq (DLQ for metrics failures)
├─ Retention: 14 days
├─ Visibility: 5 minutes
├─ Messages: {id, type, agent_id, error, timestamp, ...}
└─ Consumer: (future) metrics replay handler
Files Added:
1. k8s/config/create-kmsvc-topics.sh (executable)
├─ 90 lines
├─ Auto port-forward + retry logic
├─ Color output + error handling
└─ Usage: ./create-kmsvc-topics.sh [manual]
2. k8s/config/kmsvc-topics-job.yaml (Kubernetes)
├─ Job resource (one-time execution)
├─ Uses curl container
├─ Waits for management-service readiness
├─ 30-second retry loop
└─ Non-fatal on existing topics
3. k8s/config/terraform-kmsvc-topics.tf (Terraform)
├─ null_resource with local-exec
├─ Variables for endpoint + namespace
├─ Idempotent + traceable
└─ Outputs: created_topics + test_commands
4. docs/PHASE_6_6_KMSVC_TOPICS.md (Complete Guide)
├─ Table of topics + config
├─ 4 creation methods with examples
├─ Verification commands
├─ Message format specs
├─ Monitoring + alerts
├─ Troubleshooting guide
└─ Next steps checklist
Verification Commands:
✅ List all topics:
curl http://localhost:8080/v1/queues
✅ Check specific topic:
curl http://localhost:8080/v1/queues/poimen-memory-dlq
✅ Send test message:
curl -X POST http://localhost:8080/v1/queues/poimen-memory-dlq/messages -H "Content-Type: application/json" -d '{"body": "{\"type\": \"test\"}"}'
✅ Receive messages:
curl -X POST http://localhost:8080/v1/queues/poimen-memory-dlq/messages/receive -H "Content-Type: application/json" -d '{"maxNumberOfMessages": 10}'
Integration Points:
Phase 6.6 Code → kmsvc Topics:
1. webhook_executor.rs::send_dlq_message()
└─ On max retries: Send to poimen-memory-dlq
└─ Payload: {workflow_id, webhook_url, status, error, ...}
2. metrics_persistence.rs::send_persistence_dlq()
└─ On DB failure: Send to poimen-memory-metric-dlq
└─ Payload: {agent_id, error, timestamp}
3. queue_worker_dlq.rs::DlqHandler
└─ Processes poimen-memory-dlq messages
└─ Retries extraction failures
Production Checklist:
✅ Topics defined (2 topics)
✅ Configuration documented (retention, visibility)
✅ Creation methods (4 options)
✅ Verification commands
✅ Message formats specified
✅ Monitoring guide
✅ Troubleshooting guide
✅ Ready for deployment
Next Steps:
1. Choose creation method (recommend Method 1 for fast testing)
2. Create topics: ./create-kmsvc-topics.sh or kubectl apply job
3. Verify: curl http://localhost:8080/v1/queues
4. Deploy memory service (Phase 6.5)
5. Test webhook + metrics failures send to DLQ
6. Monitor DLQ lag + message rate (Phase 7)
Phase 6.6 Complete: ✅
- Webhook execution: ✅
- Metrics persistence: ✅
- kmsvc DLQ integration: ✅
- Topic creation (4 methods): ✅
- Documentation: ✅
6.4 KiB
6.4 KiB
Phase 6.6: kmsvc Topic Creation Guide
Topics
| Topic Name | Purpose | Retention | Visibility Timeout |
|---|---|---|---|
poimen-memory-dlq |
Extraction, webhook, agent failures | 14 days (1,209,600s) | 5 min (300s) |
poimen-memory-metric-dlq |
Metrics persistence failures | 14 days | 5 min |
Methods to Create Topics
Method 1: Manual via CLI (Fastest)
# Port-forward to management-service
kubectl -n sqs port-forward svc/management-service 8080:8080 &
# Create poimen-memory-dlq
curl -X POST http://localhost:8080/v1/queues \
-H "Content-Type: application/json" \
-d '{
"name": "poimen-memory-dlq",
"fifoQueue": false,
"visibilityTimeoutSeconds": 300,
"messageRetentionPeriodSeconds": 1209600
}'
# Create poimen-memory-metric-dlq
curl -X POST http://localhost:8080/v1/queues \
-H "Content-Type: application/json" \
-d '{
"name": "poimen-memory-metric-dlq",
"fifoQueue": false,
"visibilityTimeoutSeconds": 300,
"messageRetentionPeriodSeconds": 1209600
}'
# Verify topics were created
curl http://localhost:8080/v1/queues | jq
Method 2: Kubernetes Job (Automated)
# Deploy the job (creates topics automatically)
kubectl apply -f k8s/config/kmsvc-topics-job.yaml
# Watch job progress
kubectl -n poimen logs -f job/create-poimen-memory-kmsvc-topics
# Verify job completed
kubectl -n poimen get job create-poimen-memory-kmsvc-topics
Method 3: Terraform (IaC)
# Navigate to config directory
cd k8s/config
# Initialize Terraform
terraform init
# Create topics
terraform apply -target=null_resource.create_kmsvc_topics \
-var="management_service_endpoint=http://localhost:8080"
# Verify creation
curl http://localhost:8080/v1/queues | jq
Method 4: Shell Script (Interactive)
# Make script executable
chmod +x k8s/config/create-kmsvc-topics.sh
# Run with auto port-forward
./k8s/config/create-kmsvc-topics.sh
# Or manual port-forward
./k8s/config/create-kmsvc-topics.sh manual
Verification
List all queues
curl http://localhost:8080/v1/queues | jq
Check specific queue exists
curl http://localhost:8080/v1/queues/poimen-memory-dlq | jq
Send test message to DLQ
curl -X POST http://localhost:8080/v1/queues/poimen-memory-dlq/messages \
-H "Content-Type: application/json" \
-d '{
"body": "{\"type\": \"test\", \"workflow_id\": \"test-123\"}"
}'
Receive messages from DLQ
curl -X POST http://localhost:8080/v1/queues/poimen-memory-dlq/messages/receive \
-H "Content-Type: application/json" \
-d '{"maxNumberOfMessages": 10}'
Queue Configuration Details
poimen-memory-dlq
Used by:
- Extraction failures (contradiction detection, validation)
- Webhook execution failures (after all retries)
- Agent initialization failures
Message Format:
{
"id": "uuid",
"type": "extraction_failure|webhook_failure|agent_failure",
"workflow_id": "wf-123",
"webhook_url": "http://...",
"error": "error message",
"timestamp": "2025-01-30T10:00:00Z",
"retry_count": 0,
"max_retries": 3,
"topic": "poimen-memory-dlq"
}
Processing:
- Consumer:
queue_worker_dlq.rs::DlqHandler - Retry policy: Up to 3 retries via exponential backoff
- TTL: 14 days (can be re-processed manually)
poimen-memory-metric-dlq
Used by:
- Metrics persistence failures (DB down, connection errors)
Message Format:
{
"id": "uuid",
"type": "metrics_persistence_failure",
"agent_id": "agent-123",
"error": "DB connection failed",
"timestamp": "2025-01-30T10:00:00Z",
"topic": "poimen-memory-metric-dlq"
}
Processing:
- Consumer:
metrics_persistence.rs::send_persistence_dlq() - Retry policy: Manual intervention or scheduled replay
- TTL: 14 days (allows recovery window)
Connection Pooling
If using kmsvc producer client in Rust:
use rdkafka::producer::FutureProducer;
use rdkafka::ClientConfig;
let producer: FutureProducer = ClientConfig::new()
.set("bootstrap.servers", "kafka.sqs.svc.cluster.local:9092")
.set("client.id", "poimen-memory-producer")
.create()
.expect("Producer creation failed");
// Send message to DLQ
producer
.send(
FutureRecord::to("poimen-memory-dlq")
.key(&workflow_id)
.payload(&dlq_message),
Duration::from_secs(30),
)
.await?;
Monitoring
Prometheus Metrics
# Messages sent to DLQ
dlq_messages_sent_total{topic="poimen-memory-dlq"} 5
dlq_messages_sent_total{topic="poimen-memory-metric-dlq"} 2
# Message lag
dlq_consumer_lag{topic="poimen-memory-dlq"} 0
Alerts
Set up alerts if:
- DLQ message lag > 100 (backlog accumulating)
- DLQ message rate > 10/min (high failure rate)
- DLQ messages older than 7 days (not being processed)
Troubleshooting
Topic already exists error
# This is OK - it means the topic was already created
# You can safely ignore this error
# To force recreation, delete first:
curl -X DELETE http://localhost:8080/v1/queues/poimen-memory-dlq
Connection refused to management-service
# Verify port-forward is running
ps aux | grep "port-forward.*management-service"
# Check service is running
kubectl -n sqs get svc management-service
# Restart port-forward
pkill -f "port-forward.*management-service"
kubectl -n sqs port-forward svc/management-service 8080:8080 &
Topics not appearing in list
# Verify topic creation response was successful
curl -v -X POST http://localhost:8080/v1/queues \
-H "Content-Type: application/json" \
-d '{"name": "test-topic", "fifoQueue": false}'
# Check for HTTP 200/201 response
# If 400/409, check error message
Next Steps
- ✅ Create topics via one of the methods above
- ✅ Verify topics exist with
curl http://localhost:8080/v1/queues - ✅ Deploy memory service (will start sending DLQ messages on failures)
- ⏳ Monitor DLQ message rate + latency
- ⏳ Set up consumer group for DLQ replay
- ⏳ Add alerting for DLQ backlog (Phase 7)
References
- Management Service API:
https://kmsvc.riotpiao.com/docs - kmsvc Kafka Broker:
kafka.sqs.svc.cluster.local:9092 - Phase 6.6 Code:
crates/mem-cli/src/queue/kmsvc_topics.rs - Webhook Executor:
crates/mem-cli/src/handlers/webhook_executor.rs - Metrics Persistence:
crates/mem-cli/src/handlers/metrics_persistence.rs
Status: 🟢 Ready to create topics. Choose Method 1 (fastest) or Method 2 (automated).