k8s/messaging: add kafka kmsvc and temporal workflows

- Kafka 3-broker cluster (RF=3, min-ISR=2)
- kmsvc SQS-like API on Kafka
- Redis dedup (standalone, can extend to HA)
- Temporal workflow orchestration (Cassandra backend)
This commit is contained in:
Story Crater Bot
2026-07-11 19:17:42 -07:00
parent 4ab596196e
commit 1c02e2b831
38 changed files with 2072 additions and 0 deletions
@@ -0,0 +1,25 @@
apiVersion: argoproj.io/v1alpha1
kind: Application
metadata:
name: strimzi-operator
namespace: cicd
annotations:
argocd.argoproj.io/sync-wave: "0"
spec:
project: kmsvc
source:
repoURL: https://strimzi.io/charts/
chart: strimzi-kafka-operator
targetRevision: 0.46.0
helm:
values: |
watchNamespaces: ["sqs"]
destination:
server: https://kubernetes.default.svc
namespace: sqs
syncPolicy:
automated:
prune: true
selfHeal: true
syncOptions:
- CreateNamespace=true
+31
View File
@@ -0,0 +1,31 @@
apiVersion: argoproj.io/v1alpha1
kind: Application
metadata:
name: kafka-cluster
namespace: cicd
annotations:
argocd.argoproj.io/sync-wave: "1"
spec:
project: kmsvc
source:
repoURL: https://forgejo.riotpiao.homelab.com/rock/kafaka-management-service.git
targetRevision: main
path: k8s/charts/kafka-cluster
helm:
values: |
namespace: sqs
nodePool:
replicas: 3
storage:
class: longhorn
sizeGi: 50
resources:
memory: 5Gi
cpu: "2"
destination:
server: https://kubernetes.default.svc
namespace: sqs
syncPolicy:
automated:
prune: true
selfHeal: true
+43
View File
@@ -0,0 +1,43 @@
apiVersion: argoproj.io/v1alpha1
kind: Application
metadata:
name: kmsvc-redis
namespace: cicd
annotations:
argocd.argoproj.io/sync-wave: "1"
spec:
project: kmsvc
source:
repoURL: https://charts.bitnami.com/bitnami
chart: redis
targetRevision: 20.6.0
helm:
values: |
architecture: standalone
# docker.io/bitnami stopped publishing version-pinned tags; bitnamilegacy
# mirrors them for free. allowInsecureImages silences the chart's
# container-image allowlist check, which doesn't know about that mirror.
global:
security:
allowInsecureImages: true
image:
repository: bitnamilegacy/redis
auth:
enabled: false
master:
persistence:
enabled: true
storageClass: longhorn
size: 2Gi
resources:
limits:
memory: 1Gi
requests:
memory: 1Gi
destination:
server: https://kubernetes.default.svc
namespace: sqs
syncPolicy:
automated:
prune: true
selfHeal: true
+30
View File
@@ -0,0 +1,30 @@
apiVersion: argoproj.io/v1alpha1
kind: Application
metadata:
name: queue-crd
namespace: cicd
annotations:
argocd.argoproj.io/sync-wave: "2"
spec:
project: kmsvc
source:
repoURL: https://forgejo.riotpiao.homelab.com/rock/kafaka-management-service.git
targetRevision: main
path: k8s/charts/queue-crd
helm:
values: |
namespace: sqs
kafkaBrokers: "kmsvc-kafka-bootstrap.sqs.svc.cluster.local:9092"
redisAddr: "kmsvc-redis-master.sqs.svc.cluster.local:6379"
image:
repository: forgejo.riotpiao.homelab.com/rock/kafka-management-service-queue-operator
# CI (.forgejo/workflows/release.yaml) writes the released git tag
# here and pushes the commit -- ArgoCD picks it up on its next sync.
tag: latest
destination:
server: https://kubernetes.default.svc
namespace: sqs
syncPolicy:
automated:
prune: true
selfHeal: true
@@ -0,0 +1,37 @@
apiVersion: argoproj.io/v1alpha1
kind: Application
metadata:
name: management-service
namespace: cicd
annotations:
argocd.argoproj.io/sync-wave: "2"
spec:
project: kmsvc
source:
repoURL: https://forgejo.riotpiao.homelab.com/rock/kafaka-management-service.git
targetRevision: main
path: k8s/charts/management-service
helm:
values: |
namespace: sqs
image:
repository: forgejo.riotpiao.homelab.com/rock/kafka-management-service
# CI (.forgejo/workflows/release.yaml) writes the released git tag
# here and pushes the commit -- ArgoCD picks it up on its next sync.
tag: latest
env:
kafkaBrokers: "kmsvc-kafka-bootstrap.sqs.svc.cluster.local:9092"
redisAddr: "kmsvc-redis-master.sqs.svc.cluster.local:6379"
authentikIssuerURL: "https://authentik.riotpiao.homelab.com/application/o/kafaka/"
authentikAudience: "QI0gPtR99ar8VvhK8Tqox4SDkTKzbNU7lbgwBNSc"
ingress:
enabled: true
host: kmsvc.riotpiao.homelab.com
clusterIssuer: homelab-ca
destination:
server: https://kubernetes.default.svc
namespace: sqs
syncPolicy:
automated:
prune: true
selfHeal: true
+23
View File
@@ -0,0 +1,23 @@
apiVersion: argoproj.io/v1alpha1
kind: AppProject
metadata:
name: kmsvc
namespace: cicd
spec:
description: Kafka Management Service (design.md) -- Strimzi/Kafka, Redis, queue-operator, message-plane server
sourceRepos:
- https://forgejo.riotpiao.homelab.com/rock/kafaka-management-service.git
- https://strimzi.io/charts/
- https://charts.bitnami.com/bitnami
destinations:
- namespace: sqs
server: https://kubernetes.default.svc
- namespace: cicd
server: https://kubernetes.default.svc
clusterResourceWhitelist:
- group: "apiextensions.k8s.io"
kind: CustomResourceDefinition
- group: "rbac.authorization.k8s.io"
kind: ClusterRole
- group: "rbac.authorization.k8s.io"
kind: ClusterRoleBinding
+22
View File
@@ -0,0 +1,22 @@
apiVersion: argoproj.io/v1alpha1
kind: Application
metadata:
name: kmsvc-root
namespace: cicd
spec:
project: kmsvc
source:
repoURL: https://forgejo.riotpiao.homelab.com/rock/kafaka-management-service.git
targetRevision: main
path: k8s/argocd/apps
directory:
recurse: false
destination:
server: https://kubernetes.default.svc
namespace: cicd
syncPolicy:
automated:
prune: true
selfHeal: true
syncOptions:
- CreateNamespace=true
+5
View File
@@ -0,0 +1,5 @@
apiVersion: v2
name: kafka-cluster
description: Strimzi Kafka/KafkaNodePool CRs for the kmsvc Kafka cluster (design.md §7)
type: application
version: 0.1.0
@@ -0,0 +1,30 @@
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
name: {{ .Values.clusterName }}
namespace: {{ .Values.namespace }}
annotations:
strimzi.io/node-pools: enabled
strimzi.io/kraft: enabled
spec:
kafka:
version: 4.0.0
metadataVersion: 4.0-IV3
listeners:
- name: plain
port: 9092
type: internal
tls: false
- name: tls
port: 9093
type: internal
tls: true
config:
default.replication.factor: {{ .Values.kafka.replicationFactor }}
min.insync.replicas: {{ .Values.kafka.minInsyncReplicas }}
offsets.topic.replication.factor: {{ .Values.kafka.replicationFactor }}
transaction.state.log.replication.factor: {{ .Values.kafka.replicationFactor }}
transaction.state.log.min.isr: {{ .Values.kafka.minInsyncReplicas }}
entityOperator:
topicOperator: {}
userOperator: {}
@@ -0,0 +1,39 @@
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaNodePool
metadata:
name: {{ .Values.clusterName }}-pool
namespace: {{ .Values.namespace }}
labels:
strimzi.io/cluster: {{ .Values.clusterName }}
spec:
replicas: {{ .Values.nodePool.replicas }}
roles:
- controller
- broker
storage:
type: persistent-claim
size: {{ .Values.nodePool.storage.sizeGi }}Gi
class: {{ .Values.nodePool.storage.class }}
deleteClaim: false
resources:
limits:
memory: {{ .Values.nodePool.resources.memory }}
cpu: {{ .Values.nodePool.resources.cpu | quote }}
requests:
memory: {{ .Values.nodePool.resources.memory }}
cpu: {{ .Values.nodePool.resources.cpu | quote }}
template:
pod:
affinity:
podAntiAffinity:
preferredDuringSchedulingIgnoredDuringExecution:
- weight: 100
podAffinityTerm:
topologyKey: {{ .Values.nodePool.antiAffinityTopologyKey }}
labelSelector:
matchLabels:
strimzi.io/cluster: {{ .Values.clusterName }}
kafkaContainer:
env:
- name: KAFKA_HEAP_OPTS
value: {{ .Values.nodePool.heapOpts | quote }}
@@ -0,0 +1,19 @@
{{- if eq .Values.nodePool.storage.class "longhorn-kafka" }}
# The default "longhorn" StorageClass requests 3 replicas across 3 zone-labeled
# nodes (az-a/az-b/az-c). With Longhorn's zone-aware anti-affinity, replicas
# spread 1-per-zone for durability.
apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
name: longhorn-kafka
provisioner: driver.longhorn.io
allowVolumeExpansion: true
reclaimPolicy: Delete
volumeBindingMode: Immediate
parameters:
numberOfReplicas: "3"
staleReplicaTimeout: "30"
fromBackup: ""
fsType: "ext4"
dataLocality: "disabled"
{{- end }}
+26
View File
@@ -0,0 +1,26 @@
clusterName: kmsvc
namespace: sqs
nodePool:
replicas: 3
storage:
class: longhorn-kafka
# Longhorn's per-node scheduling budget on the current 2-node cluster has
# only ~36Gi of headroom left (other PVCs already reserve the rest), and
# each node hosts one replica of all 3 broker volumes -- so 3 * sizeGi
# must fit in that headroom. Revisit once the 3rd node joins.
sizeGi: 10
resources:
memory: 5Gi
cpu: "2"
heapOpts: "-Xms2g -Xmx2g"
# design.md §7: 3 real zones now exist (talos-cp-1=az-a, talos-worker-1=az-b,
# talos-worker-2=az-c), so anti-affinity keys off zone instead of hostname —
# spreads the 3 broker pods one-per-zone/one-per-node (equivalent today,
# but zone is the correct long-term key if a node ever gets replaced within
# the same zone).
antiAffinityTopologyKey: topology.kubernetes.io/zone
kafka:
replicationFactor: 3
minInsyncReplicas: 2
@@ -0,0 +1,5 @@
apiVersion: v2
name: management-service
description: kmsvc message-plane gRPC+REST server (design.md §1, §7a, §9)
type: application
version: 0.1.0
@@ -0,0 +1,12 @@
apiVersion: v1
kind: ConfigMap
metadata:
name: management-service-config
namespace: {{ .Values.namespace }}
data:
KMSVC_KAFKA_BROKERS: {{ .Values.env.kafkaBrokers | quote }}
KMSVC_REDIS_ADDR: {{ .Values.env.redisAddr | quote }}
KMSVC_AUTHENTIK_ISSUER_URL: {{ .Values.env.authentikIssuerURL | quote }}
KMSVC_AUTHENTIK_AUDIENCE: {{ .Values.env.authentikAudience | quote }}
KMSVC_GRPC_LISTEN_ADDR: ":{{ .Values.grpcPort }}"
KMSVC_HTTP_LISTEN_ADDR: ":{{ .Values.httpPort }}"
@@ -0,0 +1,42 @@
apiVersion: apps/v1
kind: Deployment
metadata:
name: management-service
namespace: {{ .Values.namespace }}
spec:
replicas: {{ .Values.replicaCount }}
selector:
matchLabels:
app: management-service
template:
metadata:
labels:
app: management-service
spec:
containers:
- name: management-service
image: "{{ .Values.image.repository }}:{{ .Values.image.tag }}"
imagePullPolicy: {{ .Values.image.pullPolicy }}
ports:
- name: grpc
containerPort: {{ .Values.grpcPort }}
- name: http
containerPort: {{ .Values.httpPort }}
env:
- name: GOMEMLIMIT
value: {{ .Values.goMemLimit | quote }}
envFrom:
- configMapRef:
name: management-service-config
resources:
{{- toYaml .Values.resources | nindent 12 }}
readinessProbe:
tcpSocket:
port: {{ .Values.httpPort }}
initialDelaySeconds: 5
periodSeconds: 10
livenessProbe:
tcpSocket:
port: {{ .Values.httpPort }}
initialDelaySeconds: 10
periodSeconds: 20
@@ -0,0 +1,27 @@
{{- if .Values.hpa.enabled }}
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: management-service
namespace: {{ .Values.namespace }}
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: management-service
minReplicas: {{ .Values.hpa.minReplicas }}
maxReplicas: {{ .Values.hpa.maxReplicas }}
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: {{ .Values.hpa.targetCPUUtilizationPercentage }}
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: {{ .Values.hpa.targetMemoryUtilizationPercentage }}
{{- end }}
@@ -0,0 +1,27 @@
{{- if and .Values.ingress.enabled .Values.ingress.grpcEnabled }}
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: management-service-grpc
namespace: {{ .Values.namespace }}
annotations:
cert-manager.io/cluster-issuer: {{ .Values.ingress.clusterIssuer }}
nginx.ingress.kubernetes.io/backend-protocol: "GRPC"
spec:
ingressClassName: {{ .Values.ingress.className }}
tls:
- hosts:
- {{ .Values.ingress.host }}
secretName: {{ .Values.ingress.tlsSecretName }}
rules:
- host: {{ .Values.ingress.host }}
http:
paths:
- path: {{ .Values.ingress.grpcPathPrefix }}
pathType: Prefix
backend:
service:
name: management-service
port:
number: {{ .Values.grpcPort }}
{{- end }}
@@ -0,0 +1,26 @@
{{- if .Values.ingress.enabled }}
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: management-service
namespace: {{ .Values.namespace }}
annotations:
cert-manager.io/cluster-issuer: {{ .Values.ingress.clusterIssuer }}
spec:
ingressClassName: {{ .Values.ingress.className }}
tls:
- hosts:
- {{ .Values.ingress.host }}
secretName: {{ .Values.ingress.tlsSecretName }}
rules:
- host: {{ .Values.ingress.host }}
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: management-service
port:
number: {{ .Values.httpPort }}
{{- end }}
@@ -0,0 +1,16 @@
apiVersion: v1
kind: Service
metadata:
name: management-service
namespace: {{ .Values.namespace }}
spec:
selector:
app: management-service
ports:
- name: grpc
port: {{ .Values.grpcPort }}
targetPort: {{ .Values.grpcPort }}
- name: http
port: {{ .Values.httpPort }}
targetPort: {{ .Values.httpPort }}
type: ClusterIP
@@ -0,0 +1,50 @@
namespace: sqs
replicaCount: 3
image:
repository: forgejo.riotpiao.homelab.com/rock/kafka-management-service
tag: latest
pullPolicy: Always
grpcPort: 9090
httpPort: 8080
env:
kafkaBrokers: "kmsvc-kafka-bootstrap.sqs.svc.cluster.local:9092"
redisAddr: "kmsvc-redis-master.sqs.svc.cluster.local:6379"
authentikIssuerURL: ""
authentikAudience: ""
resources:
requests:
cpu: 200m
memory: 256Mi
limits:
cpu: "1"
memory: 512Mi
# Go's GC only reacts to GOGC by default and has no idea about the cgroup
# memory limit above -- it'll happily grow heap until the kernel OOMKills it.
# Setting GOMEMLIMIT to ~90% of the container limit makes the GC self-throttle
# before that happens. Keep this in sync with resources.limits.memory.
goMemLimit: "460MiB"
hpa:
enabled: true
minReplicas: 3
maxReplicas: 9
targetCPUUtilizationPercentage: 70
targetMemoryUtilizationPercentage: 80
ingress:
enabled: true
className: nginx
clusterIssuer: homelab-ca
host: kmsvc.riotpiao.homelab.com
tlsSecretName: kmsvc-tls
# kmsvc-cli connects via gRPC directly to --server/KMSVC_SERVER (default
# kmsvc.riotpiao.homelab.com:443, see kmsvc-cli's README), so raw gRPC needs an
# external path too — scoped to the gRPC service's own path prefix on the
# same host/port, rather than opening the whole host to gRPC passthrough.
grpcEnabled: true
grpcPathPrefix: /kafkamgmt.v1.QueueService/
+5
View File
@@ -0,0 +1,5 @@
apiVersion: v2
name: queue-crd
description: Queue CRD definition + queue-operator Deployment/RBAC (design.md §2a)
type: application
version: 0.1.0
+274
View File
@@ -0,0 +1,274 @@
---
apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
annotations:
controller-gen.kubebuilder.io/version: v0.21.0
name: queues.kmsvc.io
spec:
group: kmsvc.io
names:
kind: Queue
listKind: QueueList
plural: queues
shortNames:
- queue
- queues
singular: queue
scope: Namespaced
versions:
- additionalPrinterColumns:
- jsonPath: .spec.fifoQueue
name: FIFO
type: boolean
- jsonPath: .status.phase
name: Phase
type: string
name: v1
schema:
openAPIV3Schema:
description: Queue is the Schema for the queues API — see design.md §2a.
properties:
apiVersion:
description: |-
APIVersion defines the versioned schema of this representation of an object.
Servers should convert recognized schemas to the latest internal value, and
may reject unrecognized values.
More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#resources
type: string
kind:
description: |-
Kind is a string value representing the REST resource this object represents.
Servers may infer this from the endpoint the client submits requests to.
Cannot be updated.
In CamelCase.
More info: https://git.k8s.io/community/contributors/devel/sig-architecture/api-conventions.md#types-kinds
type: string
metadata:
type: object
spec:
description: QueueSpec defines the desired state of a Queue (design.md
§2a).
properties:
deadLetterTargetQueue:
description: |-
DeadLetterTargetQueue is the name of another Queue to route exhausted
messages to. Must not point at itself or at another DLQ (design.md §5).
type: string
delaySeconds:
description: DelaySeconds is the default delivery delay applied to
sent messages.
format: int32
maximum: 900
minimum: 0
type: integer
fifoQueue:
default: false
description: FIFOQueue enables per-MessageGroupId ordering and deduplication
semantics.
type: boolean
isDLQ:
description: |-
IsDLQ marks this queue as itself a dead-letter queue, used to enforce
the no-DLQ-chaining validation rule in design.md §5.
type: boolean
maxReceiveCount:
default: 5
description: |-
MaxReceiveCount is how many times a message may be redelivered before
being routed to DeadLetterTargetQueue.
format: int32
minimum: 1
type: integer
maxShards:
default: 8
description: MaxShards is the ceiling on shard count the operator
may split up to (design.md §2c).
format: int32
minimum: 1
type: integer
messageRetentionPeriodSeconds:
default: 345600
description: MessageRetentionPeriodSeconds maps to the underlying
Kafka topic's retention.ms.
format: int32
maximum: 1209600
minimum: 60
type: integer
minShards:
default: 1
description: MinShards is the floor on shard count; the operator never
merges below this.
format: int32
minimum: 1
type: integer
partitionsPerShard:
default: 6
description: PartitionsPerShard is the Kafka partition count on each
shard's topic.
format: int32
minimum: 1
type: integer
shardSplitCooldownSeconds:
default: 300
description: |-
ShardSplitCooldownSeconds is the minimum age a shard must reach before it
is eligible to be split again, preventing rapid re-splitting of a child
that hasn't yet absorbed its share of traffic.
format: int32
minimum: 0
type: integer
shardSplitThresholdBytesPerSec:
default: 5242880
description: |-
ShardSplitThresholdBytesPerSec is the sustained per-shard throughput that
triggers a split into two child shards (design.md §2c).
format: int64
minimum: 1
type: integer
visibilityTimeoutSeconds:
default: 30
description: |-
VisibilityTimeoutSeconds is how long a received-but-unacked message stays
invisible to other consumers before being redelivered.
format: int32
maximum: 43200
minimum: 0
type: integer
type: object
status:
description: QueueStatus defines the observed state of a Queue.
properties:
conditions:
description: Conditions hold detailed status information.
items:
description: Condition contains details for one aspect of the current
state of this API Resource.
properties:
lastTransitionTime:
description: |-
lastTransitionTime is the last time the condition transitioned from one status to another.
This should be when the underlying condition changed. If that is not known, then using the time when the API field changed is acceptable.
format: date-time
type: string
message:
description: |-
message is a human readable message indicating details about the transition.
This may be an empty string.
maxLength: 32768
type: string
observedGeneration:
description: |-
observedGeneration represents the .metadata.generation that the condition was set based upon.
For instance, if .metadata.generation is currently 12, but the .status.conditions[x].observedGeneration is 9, the condition is out of date
with respect to the current state of the instance.
format: int64
minimum: 0
type: integer
reason:
description: |-
reason contains a programmatic identifier indicating the reason for the condition's last transition.
Producers of specific condition types may define expected values and meanings for this field,
and whether the values are considered a guaranteed API.
The value should be a CamelCase string.
This field may not be empty.
maxLength: 1024
minLength: 1
pattern: ^[A-Za-z]([A-Za-z0-9_,:]*[A-Za-z0-9_])?$
type: string
status:
description: status of the condition, one of True, False, Unknown.
enum:
- "True"
- "False"
- Unknown
type: string
type:
description: type of condition in CamelCase or in foo.example.com/CamelCase.
maxLength: 316
pattern: ^([a-z0-9]([-a-z0-9]*[a-z0-9])?(\.[a-z0-9]([-a-z0-9]*[a-z0-9])?)*/)?(([A-Za-z0-9][-A-Za-z0-9_.]*)?[A-Za-z0-9])$
type: string
required:
- lastTransitionTime
- message
- reason
- status
- type
type: object
type: array
phase:
description: Phase is the current reconciliation phase.
enum:
- Pending
- Ready
- Failed
type: string
shards:
description: |-
Shards lists every shard backing this queue, active or draining
(design.md §2a/§2c).
items:
description: ShardStatus describes one shard backing a Queue (design.md
§2a/§2c).
properties:
availabilityZones:
description: |-
AvailabilityZones lists the topology.kubernetes.io/zone values of every
node currently hosting a Kafka replica of this shard's topic, resolved
from the broker pods' node placement each reconcile. Empty until the
first successful resolution (e.g. node lookup failed transiently).
items:
type: string
type: array
createdAt:
description: |-
CreatedAt timestamps when this shard was created, used to enforce
ShardSplitCooldownSeconds.
format: date-time
type: string
hashRangeEnd:
format: int64
type: integer
hashRangeStart:
description: |-
HashRangeStart/HashRangeEnd define the [start, end) murmur2 hash range
this shard owns over the 32-bit key space. Stored as int64 (not uint32)
because controller-gen maps Go uint32 to OpenAPI format:int32, whose max
(2147483647) is smaller than FullHashRangeEnd (0xFFFFFFFF) and the
apiserver rejects the status update.
format: int64
type: integer
id:
description: ID is the shard's identifier, used in its topic
name (kmsvc.{queue}.shard-{id}).
type: string
parentId:
description: |-
ParentID is the shard ID this shard was split from, empty for the
original shard-0.
type: string
phase:
description: Phase is this shard's lifecycle state.
enum:
- Active
- Closing
- Closed
type: string
topic:
description: Topic is the underlying Kafka topic name for this
shard.
type: string
required:
- hashRangeEnd
- hashRangeStart
- id
- phase
- topic
type: object
type: array
type: object
type: object
served: true
storage: true
subresources:
status: {}
@@ -0,0 +1,37 @@
apiVersion: apps/v1
kind: Deployment
metadata:
name: queue-operator
namespace: {{ .Values.namespace }}
spec:
replicas: 1
selector:
matchLabels:
app: queue-operator
template:
metadata:
labels:
app: queue-operator
spec:
serviceAccountName: queue-operator
containers:
- name: queue-operator
image: "{{ .Values.image.repository }}:{{ .Values.image.tag }}"
imagePullPolicy: {{ .Values.image.pullPolicy }}
env:
- name: KMSVC_KAFKA_BROKERS
value: {{ .Values.kafkaBrokers | quote }}
- name: KMSVC_REDIS_ADDR
value: {{ .Values.redisAddr | quote }}
- name: GOMEMLIMIT
value: {{ .Values.goMemLimit | quote }}
- name: KMSVC_NAMESPACE
valueFrom:
fieldRef:
fieldPath: metadata.namespace
- name: KMSVC_KAFKA_CLUSTER_NAME
value: {{ .Values.kafkaClusterName | quote }}
- name: KMSVC_KAFKA_POOL_NAME
value: {{ .Values.kafkaPoolName | quote }}
resources:
{{- toYaml .Values.resources | nindent 12 }}
@@ -0,0 +1,42 @@
apiVersion: v1
kind: ServiceAccount
metadata:
name: queue-operator
namespace: {{ .Values.namespace }}
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: queue-operator
rules:
- apiGroups: ["kmsvc.io"]
resources: ["queues"]
verbs: ["get", "list", "watch", "update", "patch"]
- apiGroups: ["kmsvc.io"]
resources: ["queues/status"]
verbs: ["get", "update", "patch"]
- apiGroups: ["kmsvc.io"]
resources: ["queues/finalizers"]
verbs: ["update"]
- apiGroups: ["coordination.k8s.io"]
resources: ["leases"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
- apiGroups: [""]
resources: ["events"]
verbs: ["create", "patch"]
- apiGroups: [""]
resources: ["pods", "nodes"]
verbs: ["get"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
name: queue-operator
roleRef:
apiGroup: rbac.authorization.k8s.io
kind: ClusterRole
name: queue-operator
subjects:
- kind: ServiceAccount
name: queue-operator
namespace: {{ .Values.namespace }}
+26
View File
@@ -0,0 +1,26 @@
namespace: sqs
image:
repository: forgejo.riotpiao.homelab.com/rock/kafka-management-service-queue-operator
tag: latest
pullPolicy: Always
kafkaBrokers: "kmsvc-kafka-bootstrap.sqs.svc.cluster.local:9092"
redisAddr: "kmsvc-redis-master.sqs.svc.cluster.local:6379"
# Must match kafka-cluster chart's clusterName/derived pool name -- used to
# resolve "<kafkaClusterName>-<kafkaPoolName>-<brokerID>" broker pod names
# for AZ-aware Queue status (design.md §2a).
kafkaClusterName: kmsvc
kafkaPoolName: kmsvc-pool
resources:
requests:
cpu: 100m
memory: 128Mi
limits:
cpu: 500m
memory: 256Mi
# See management-service/values.yaml's goMemLimit comment -- same reasoning.
goMemLimit: "230MiB"
+29
View File
@@ -0,0 +1,29 @@
# design.md §7b: cluster-specific values for the homelab environment.
# No secrets here — Authentik client secret etc. flow through the existing
# Vault/talos-cli pattern, referenced at deploy time, not inlined.
namespace: sqs
kafkaCluster:
nodePool:
replicas: 3
storage:
# longhorn-kafka now uses numberOfReplicas: 3 across 3 zone-labeled nodes.
class: longhorn-kafka
# Per-node headroom: with 3 nodes and existing PVCs, estimate ~100+ Gi total
# available. Each node hosts one replica of all 3 broker volumes, so 3 *
# sizeGi must fit. Monitor usage during Kafka deployment.
sizeGi: 10
resources:
memory: 5Gi
cpu: "2"
redis:
storageClass: longhorn
memoryLimit: 1Gi
managementService:
ingress:
host: kmsvc.riotpiao.homelab.com
clusterIssuer: homelab-ca
authentikIssuerURL: "https://authentik.riotpiao.homelab.com/application/o/kafaka/"
authentikAudience: "QI0gPtR99ar8VvhK8Tqox4SDkTKzbNU7lbgwBNSc"
+99
View File
@@ -0,0 +1,99 @@
environments:
default:
values:
- environments/homelab.yaml
homelab:
values:
- environments/homelab.yaml
---
helmDefaults:
wait: true
timeout: 600
repositories:
- name: strimzi
url: https://strimzi.io/charts/
- name: bitnami
url: https://charts.bitnami.com/bitnami
releases:
- name: strimzi-operator
namespace: {{ .Values.namespace }}
chart: strimzi/strimzi-kafka-operator
version: 0.46.0
values:
- watchNamespaces: ["{{ .Values.namespace }}"]
- name: kafka-cluster
namespace: {{ .Values.namespace }}
chart: charts/kafka-cluster
needs:
- {{ .Values.namespace }}/strimzi-operator
values:
- namespace: {{ .Values.namespace }}
nodePool:
replicas: {{ .Values.kafkaCluster.nodePool.replicas }}
storage:
class: {{ .Values.kafkaCluster.nodePool.storage.class }}
sizeGi: {{ .Values.kafkaCluster.nodePool.storage.sizeGi }}
resources:
memory: {{ .Values.kafkaCluster.nodePool.resources.memory }}
cpu: {{ .Values.kafkaCluster.nodePool.resources.cpu | quote }}
- name: kmsvc-redis
namespace: {{ .Values.namespace }}
chart: bitnami/redis
version: 20.6.0
values:
- architecture: standalone
# Bitnami stopped publishing version-pinned tags under docker.io/bitnami
# (only `latest` remains there); bitnamilegacy/* mirrors the old
# versioned tags for free, so pin there instead of floating on `latest`.
# The chart's container-image allowlist check doesn't know about the
# legacy mirror, hence allowInsecureImages.
global:
security:
allowInsecureImages: true
image:
repository: bitnamilegacy/redis
auth:
enabled: false
master:
persistence:
enabled: true
storageClass: {{ .Values.redis.storageClass }}
size: 2Gi
resources:
limits:
memory: {{ .Values.redis.memoryLimit }}
requests:
memory: {{ .Values.redis.memoryLimit }}
- name: queue-crd
namespace: {{ .Values.namespace }}
chart: charts/queue-crd
needs:
- {{ .Values.namespace }}/kafka-cluster
- {{ .Values.namespace }}/kmsvc-redis
values:
- namespace: {{ .Values.namespace }}
kafkaBrokers: "kmsvc-kafka-bootstrap.{{ .Values.namespace }}.svc.cluster.local:9092"
redisAddr: "kmsvc-redis-master.{{ .Values.namespace }}.svc.cluster.local:6379"
- name: management-service
namespace: {{ .Values.namespace }}
chart: charts/management-service
needs:
- {{ .Values.namespace }}/kafka-cluster
- {{ .Values.namespace }}/kmsvc-redis
values:
- namespace: {{ .Values.namespace }}
env:
kafkaBrokers: "kmsvc-kafka-bootstrap.{{ .Values.namespace }}.svc.cluster.local:9092"
redisAddr: "kmsvc-redis-master.{{ .Values.namespace }}.svc.cluster.local:6379"
authentikIssuerURL: {{ .Values.managementService.authentikIssuerURL | quote }}
authentikAudience: {{ .Values.managementService.authentikAudience | quote }}
ingress:
enabled: true
host: {{ .Values.managementService.ingress.host | quote }}
clusterIssuer: {{ .Values.managementService.ingress.clusterIssuer | quote }}
+25
View File
@@ -0,0 +1,25 @@
# Ingress for kmsvc REST API — routes to OAuth2-Proxy
# TLS terminated here; oauth2-proxy handles OIDC auth
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: kmsvc
namespace: sqs
spec:
ingressClassName: nginx
tls:
- secretName: kmsvc-tls
hosts:
- kmsvc.riotpiao.homelab.com
rules:
- host: kmsvc.riotpiao.homelab.com
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: oauth2-proxy
port:
number: 4180
+106
View File
@@ -0,0 +1,106 @@
# OAuth2-Proxy for kmsvc REST API
# Protects gRPC-gateway (REST) endpoint with Authentik OIDC
apiVersion: v1
kind: ServiceAccount
metadata:
name: oauth2-proxy
namespace: sqs
---
apiVersion: apps/v1
kind: Deployment
metadata:
name: oauth2-proxy
namespace: sqs
spec:
replicas: 1
selector:
matchLabels:
app: oauth2-proxy
template:
metadata:
labels:
app: oauth2-proxy
annotations:
secret.reloader.stakater.com/reload: "kmsvc-oidc"
spec:
serviceAccountName: oauth2-proxy
containers:
- name: oauth2-proxy
image: quay.io/oauth2-proxy/oauth2-proxy:v7.5.1
imagePullPolicy: IfNotPresent
ports:
- name: http
containerPort: 4180
protocol: TCP
env:
- name: OAUTH2_PROXY_PROVIDER
value: "oidc"
- name: OAUTH2_PROXY_OIDC_ISSUER_URL
value: "https://authentik.riotpiao.homelab.com/application/o/kmsvc/"
- name: OAUTH2_PROXY_CLIENT_ID
value: "kmsvc"
- name: OAUTH2_PROXY_CLIENT_SECRET
valueFrom:
secretKeyRef:
name: kmsvc-oidc
key: clientSecret
- name: OAUTH2_PROXY_COOKIE_SECRET
valueFrom:
secretKeyRef:
name: kmsvc-oidc
key: cookieSecret
- name: OAUTH2_PROXY_REDIRECT_URL
value: "https://kmsvc.riotpiao.homelab.com/oauth2/callback"
- name: OAUTH2_PROXY_UPSTREAM
value: "http://kmsvc-management-service:8080"
- name: OAUTH2_PROXY_COOKIE_SECURE
value: "true"
- name: OAUTH2_PROXY_COOKIE_HTTPONLY
value: "true"
- name: OAUTH2_PROXY_COOKIE_SAMESITE
value: "Lax"
- name: OAUTH2_PROXY_EMAIL_DOMAIN
value: "*"
- name: OAUTH2_PROXY_SKIP_AUTH_REGEX
value: "^/health"
- name: OAUTH2_PROXY_PASS_AUTHORIZATION_HEADER
value: "true"
- name: OAUTH2_PROXY_REVERSE_PROXY
value: "true"
resources:
requests:
cpu: 100m
memory: 128Mi
limits:
cpu: 200m
memory: 256Mi
livenessProbe:
httpGet:
path: /ping
port: http
initialDelaySeconds: 10
periodSeconds: 10
readinessProbe:
httpGet:
path: /ping
port: http
initialDelaySeconds: 5
periodSeconds: 5
---
apiVersion: v1
kind: Service
metadata:
name: oauth2-proxy
namespace: sqs
spec:
type: ClusterIP
ports:
- port: 4180
targetPort: http
protocol: TCP
name: http
selector:
app: oauth2-proxy
+32
View File
@@ -0,0 +1,32 @@
apiVersion: kmsvc.io/v1
kind: Queue
metadata:
name: orders-fifo
namespace: sqs
spec:
fifoQueue: true
visibilityTimeoutSeconds: 30
messageRetentionPeriodSeconds: 345600
maxReceiveCount: 5
deadLetterTargetQueue: orders-fifo-dlq
delaySeconds: 0
partitionsPerShard: 6
minShards: 1
maxShards: 8
shardSplitThresholdBytesPerSec: 5242880
shardSplitCooldownSeconds: 300
---
apiVersion: kmsvc.io/v1
kind: Queue
metadata:
name: orders-fifo-dlq
namespace: sqs
spec:
fifoQueue: true
isDLQ: true
visibilityTimeoutSeconds: 30
messageRetentionPeriodSeconds: 1209600
maxReceiveCount: 5
partitionsPerShard: 6
minShards: 1
maxShards: 1