Skip to content
Merged

Eda #17

Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
66 commits
Select commit Hold shift + click to select a range
6c4dc7b
Update dynamic_pool.go
RyazanovAlexander Mar 10, 2026
6627825
Update post-start.sh
RyazanovAlexander Mar 10, 2026
7f6f10c
Update justfile
RyazanovAlexander Mar 10, 2026
6c2e4ee
Update justfile
RyazanovAlexander Mar 10, 2026
f9b0634
u
RyazanovAlexander Mar 10, 2026
7a5c2ea
Update devcontainer.json
RyazanovAlexander Mar 11, 2026
e1eea23
u
RyazanovAlexander Mar 11, 2026
3aa3c75
u
RyazanovAlexander Mar 12, 2026
edafc42
u
RyazanovAlexander Mar 21, 2026
f6ac26a
u
RyazanovAlexander Mar 22, 2026
25cea5e
u
RyazanovAlexander Mar 22, 2026
8fd4aef
u
RyazanovAlexander Mar 22, 2026
7b18c06
Update proto.go
RyazanovAlexander Mar 23, 2026
f4023f9
u
RyazanovAlexander Mar 23, 2026
6a90576
Update README.md
RyazanovAlexander Mar 23, 2026
af5c7e3
Update logger_test.go
RyazanovAlexander Mar 23, 2026
3298a8e
u
RyazanovAlexander Mar 23, 2026
658b3cd
Update handler_test.go
RyazanovAlexander Mar 23, 2026
4dfa463
Update feedback_load.js
RyazanovAlexander Mar 23, 2026
803b899
u
RyazanovAlexander Mar 23, 2026
a331d83
u
RyazanovAlexander Mar 23, 2026
8fac4e9
u
RyazanovAlexander Mar 23, 2026
cc69fe0
Update main.go
RyazanovAlexander Mar 25, 2026
3bd58a6
u
RyazanovAlexander Mar 25, 2026
190832b
Delete 0.make.global.skill..md
RyazanovAlexander Mar 25, 2026
29879d3
u
RyazanovAlexander Mar 25, 2026
8d92617
u
RyazanovAlexander Mar 25, 2026
d15e2d7
u
RyazanovAlexander Mar 25, 2026
d80d38b
u
RyazanovAlexander Mar 25, 2026
00d1e2d
u
RyazanovAlexander Mar 25, 2026
297bf31
u
RyazanovAlexander Mar 25, 2026
9a5d2f9
Update repository.go
RyazanovAlexander Mar 25, 2026
072b919
Update service.go
RyazanovAlexander Mar 25, 2026
43ec8ff
u
RyazanovAlexander Mar 25, 2026
7385fd0
Update carrier.go
RyazanovAlexander Mar 25, 2026
2fa50f7
Update model.go
RyazanovAlexander Mar 25, 2026
38f72d5
u
RyazanovAlexander Mar 25, 2026
0e58484
u
RyazanovAlexander Mar 25, 2026
a1be478
Update proto.go
RyazanovAlexander Mar 25, 2026
3164d92
u
RyazanovAlexander Mar 26, 2026
1f7c4e5
u
RyazanovAlexander Mar 26, 2026
47e7cab
u
RyazanovAlexander Apr 2, 2026
b145501
u
RyazanovAlexander Apr 2, 2026
31f73a4
u
RyazanovAlexander Apr 2, 2026
72bfa92
u
RyazanovAlexander Apr 2, 2026
d8f3748
u
RyazanovAlexander Apr 2, 2026
c536eb7
Update logger_test.go
RyazanovAlexander Apr 12, 2026
82b4360
u
RyazanovAlexander Apr 12, 2026
b3e77ad
u
RyazanovAlexander Apr 12, 2026
29a91c8
u
RyazanovAlexander Apr 16, 2026
74bdd22
u
RyazanovAlexander Apr 20, 2026
8a57f14
Update justfile
RyazanovAlexander Apr 20, 2026
e316c6c
Update config.go
RyazanovAlexander Apr 20, 2026
50bcea4
u
RyazanovAlexander Apr 20, 2026
d5781a5
Update proto.go
RyazanovAlexander Apr 20, 2026
11e6fd9
Update consumer.go
RyazanovAlexander Apr 20, 2026
b90a421
Update proto.go
RyazanovAlexander Apr 20, 2026
9f38b24
Update justfile
RyazanovAlexander Apr 20, 2026
4882537
u
RyazanovAlexander Apr 20, 2026
e43fad4
Update skaffold.yaml
RyazanovAlexander Apr 20, 2026
b1918be
Update skaffold.yaml
RyazanovAlexander Apr 20, 2026
86c1fa3
Update Dockerfile
RyazanovAlexander Apr 20, 2026
9880efe
Update justfile
RyazanovAlexander Apr 20, 2026
5ac7eb3
Update justfile
RyazanovAlexander Apr 20, 2026
24d3073
Update justfile
RyazanovAlexander Apr 20, 2026
7305508
Update network-policies.yaml
RyazanovAlexander Apr 20, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .claude/CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ GET /swagger/*any — Swagger UI
- **Secrets rule:** ONLY via Vault Kubernetes auth. No `envFrom.secretRef`, no env var passwords, no ESO ExternalSecrets.
- DB creds: Vault Database Secrets Engine `creds/mentor-api-role` → `{username, password}` (new Vault lease each call)
- Static secrets: Vault KV `local/mathtrail-mentor`
- **CDC:** Debezium monitors the `feedback` table, publishes events to Kafka — app does NOT publish events
- **CDC:** RisingWave monitors the `feedback` table via PostgreSQL logical replication (WAL), publishes events to AutoMQ
- Helm chart uses `mathtrail-service-lib` library chart from `https://MathTrail.github.io/charts/charts`
- The service-lib provides: ServiceAccount, RBAC, ConfigMap, Migration Job, Deployment, Service, HPA

Expand Down
1 change: 1 addition & 0 deletions .devcontainer/devcontainer.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
"build": {
"dockerfile": "Dockerfile"
},
"initializeCommand": "docker pull ghcr.io/mathtrail/platform-env:1",
"features": {
"ghcr.io/devcontainers/features/go:1": {
"version": "1.26.0"
Expand Down
2 changes: 2 additions & 0 deletions .devcontainer/post-start.sh
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ if sudo test -f "$KUBECONFIG_SRC"; then
sudo install -o vscode -g vscode -m 600 "$KUBECONFIG_SRC" /home/vscode/.kube/config
# k3d uses 0.0.0.0 on Linux host, but from inside a container the host is reached via host.docker.internal
sed -i 's|https://0.0.0.0:|https://host.docker.internal:|g' /home/vscode/.kube/config
# TLS cert doesn't include host.docker.internal SAN — replace CA data with insecure flag
sed -i 's/certificate-authority-data:.*/insecure-skip-tls-verify: true/' /home/vscode/.kube/config
echo "Kubeconfig ready"
else
echo "Warning: kubeconfig not found at $KUBECONFIG_SRC"
Expand Down
32 changes: 32 additions & 0 deletions .mockery.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
with-expecter: true
disable-version-string: true
resolve-type-alias: false
issue-845-fix: true
dir: "{{.InterfaceDir}}/mocks"
outpkg: mocks
filename: "mock_{{.InterfaceName | snakecase}}.go"
mockname: "Mock{{.InterfaceName}}"

packages:
github.com/MathTrail/mentor-api/internal/domain/onboarding:
interfaces:
DLQPublisher:
github.com/MathTrail/mentor-api/internal/domain/feedback:
interfaces:
Service:
Repository:
github.com/MathTrail/mentor-api/internal/domain/roadmap:
interfaces:
Service:
github.com/MathTrail/mentor-api/internal/clients:
interfaces:
FeedbackClient:
ProfileClient:
github.com/MathTrail/mentor-api/internal/runner:
interfaces:
Worker:
github.com/MathTrail/mentor-api/internal/infra/postgres:
interfaces:
DB:
config:
filename: "mock_db.go"
6 changes: 4 additions & 2 deletions Dockerfile
Original file line number Diff line number Diff line change
@@ -1,17 +1,19 @@
FROM golang:1.26.0-alpine AS builder
FROM docker.io/library/golang:1.26.0-alpine AS builder
RUN apk add --no-cache git
WORKDIR /build
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux go build -ldflags="-s -w" -o /app ./cmd/server
RUN CGO_ENABLED=0 GOOS=linux go build -ldflags="-s -w" -o /migrate ./cmd/migrate
RUN CGO_ENABLED=0 GOOS=linux go build -ldflags="-s -w" -o /kafka-setup ./cmd/kafka-setup

FROM alpine:3.21
FROM docker.io/library/alpine:3.21
RUN apk add --no-cache ca-certificates tzdata \
&& adduser -D -u 10001 appuser
COPY --from=builder /app /app
COPY --from=builder /migrate /migrate
COPY --from=builder /kafka-setup /kafka-setup
USER 10001
EXPOSE 8080
ENTRYPOINT ["/app"]
25 changes: 14 additions & 11 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,9 @@ Mentor API is the intelligence hub of the MathTrail platform, responsible for ad
## System Architecture

[![PostgreSQL](https://img.shields.io/badge/PostgreSQL-336791?style=for-the-badge&logo=postgresql&logoColor=white)](https://www.postgresql.org/)
[![Debezium](https://img.shields.io/badge/Debezium-FF6A00?style=for-the-badge&logo=redhat&logoColor=white)](https://debezium.io/)
[![Apache Kafka](https://img.shields.io/badge/Kafka-000000?style=for-the-badge&logo=apachekafka&logoColor=white)](https://kafka.apache.org/)
[![Apache Flink](https://img.shields.io/badge/Flink-E6526F?style=for-the-badge&logo=apacheflink&logoColor=white)](https://flink.apache.org/)
[![Apicurio](https://img.shields.io/badge/Apicurio-0066CC?style=for-the-badge&logo=redhat&logoColor=white)](https://www.apicur.io/)
[![AutoMQ](https://img.shields.io/badge/AutoMQ-FF6A00?style=for-the-badge&logo=apachekafka&logoColor=white)](https://www.automq.com/)
[![RisingWave CDC](https://img.shields.io/badge/RisingWave-CDC-1E90FF?style=for-the-badge&logo=postgresql&logoColor=white)](https://risingwave.com/)

[![Architecture: EDA](https://img.shields.io/badge/Architecture-Event--Driven-8A2BE2?style=for-the-badge&logo=eventstore)](https://aws.amazon.com/event-driven-architecture/)
[![Kubernetes](https://img.shields.io/badge/Kubernetes-326CE5?style=for-the-badge&logo=kubernetes&logoColor=white)](./infra/helm/mentor-api)
Expand All @@ -40,7 +40,7 @@ Mentor API is the intelligence hub of the MathTrail platform, responsible for ad
graph LR
User([Student UI]) -- "Auth" --> OK[Oathkeeper]

subgraph MentorService [Mentor API Platform]
subgraph MentorService [Mentor API]
direction LR
App["Mentor API"]
end
Expand All @@ -54,15 +54,18 @@ graph LR
end

App -- "SQL" --> PGB
PG -- "CDC" --> Deb[Debezium]
PG -- "CDC" --> RW[RisingWave CDC]

subgraph Bus [Event Bus]
Kfk{Kafka}
direction TB
AMQ{AutoMQ}
Apr["Apicurio"]
end

Deb -- "feedback.created" --> Kfk
Kfk -- "progress / profile" --> App
App -- "strategy / roadmap" --> Kfk
RW -- "feedback.created" --> AMQ
AMQ -- "progress / profile" --> App
App -- "strategy / roadmap" --> AMQ
App -. "schema" .-> Apr

subgraph Support [Infra Support]
direction TB
Expand All @@ -84,7 +87,7 @@ graph LR
classDef actorCls fill:#1e1b4b,stroke:#818cf8,color:#fff

class App svc; class OK authCls;
class PGB,PG,Mig dataCls; class Deb cdcCls; class Kfk eventCls;
class PGB,PG,Mig dataCls; class RW cdcCls; class AMQ,Apr eventCls;
class Vault,ESO secretCls; class Obs obsCls; class User actorCls;
```

Expand Down
234 changes: 234 additions & 0 deletions cmd/kafka-setup/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,234 @@
package main

import (
"context"
"errors"
"fmt"
"io"
"net"
"os"
"strconv"
"strings"
"time"

"github.com/twmb/franz-go/pkg/kadm"
"github.com/twmb/franz-go/pkg/kerr"
"github.com/twmb/franz-go/pkg/kgo"
"github.com/twmb/franz-go/pkg/sasl/scram"
"go.uber.org/zap"

"github.com/MathTrail/mentor-api/internal/logger"
)

// topicSpec holds the desired configuration for a single Kafka topic.
type topicSpec struct {
Name string
Partitions int32
ReplicationFactor int16
}

// config holds all startup configuration resolved from environment variables.
type config struct {
BootstrapServers string
Topics []topicSpec
SASLUser string
SASLPass string
}

// loadConfig reads and validates all required environment variables.
// Calls log.Fatal on the first missing or malformed value.
func loadConfig(log *zap.Logger) config {
bootstrapServers, err := requiredEnv("KAFKA_BOOTSTRAP_SERVERS")
if err != nil {
log.Fatal("missing required config", zap.Error(err))
}
topicsRaw, err := requiredEnv("KAFKA_TOPICS")
if err != nil {
log.Fatal("missing required config", zap.Error(err))
}
topics, err := parseTopics(topicsRaw)
if err != nil {
log.Fatal("invalid KAFKA_TOPICS", zap.Error(err))
}
return config{
BootstrapServers: bootstrapServers,
Topics: topics,
SASLUser: envOrDefault("KAFKA_SASL_USERNAME", ""),
SASLPass: envOrDefault("KAFKA_SASL_PASSWORD", ""),
}
}

func main() {
log := logger.NewLogger("kafka-setup", "info", "json")

cfg := loadConfig(log)

brokers := splitBrokers(cfg.BootstrapServers)

opts := []kgo.Opt{kgo.SeedBrokers(brokers...)}

// SASL is opt-in: only enabled when both credentials are provided.
// Brokers configured with PLAINTEXT listeners (e.g. local dev AutoMQ on port 9092)
// reject any SASL handshake with ILLEGAL_SASL_STATE. Keeping SASL unconditional
// would break plaintext-only environments even when credentials are empty strings.
if cfg.SASLUser != "" && cfg.SASLPass != "" {
auth := scram.Auth{User: cfg.SASLUser, Pass: cfg.SASLPass}
opts = append(opts, kgo.SASL(auth.AsSha512Mechanism()))
}

client, err := kgo.NewClient(opts...)
if err != nil {
log.Fatal("failed to create kafka client", zap.Error(err))
}
defer client.Close()

adm := kadm.NewClient(client)

// Retry loop: AutoMQ controller may not be fully initialised at the time
// this Job starts (especially during first cluster boot or CI). Transient
// network errors are retried up to 3 times with a 2-second pause.
const (
maxAttempts = 3
retryPause = 2 * time.Second
)

ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()

var lastErr error
for attempt := 1; attempt <= maxAttempts; attempt++ {
lastErr = createTopics(ctx, adm, cfg.Topics, log)
if lastErr == nil {
log.Info("kafka topics ready")
return
}

// Do not retry on configuration or auth errors — only transient ones.
if !isTransient(lastErr) {
log.Fatal("permanent error creating topics", zap.Error(lastErr))
}

if attempt < maxAttempts {
log.Warn("transient error creating topics, retrying",
zap.Error(lastErr),
zap.Int("attempt", attempt),
zap.Duration("retry_in", retryPause),
)
select {
case <-time.After(retryPause):
// continue retry loop
case <-ctx.Done():
log.Fatal("context expired while waiting to retry", zap.Error(ctx.Err()))
}
}
}

log.Fatal("failed to create topics after retries", zap.Error(lastErr))
}

// createTopics creates all specified topics. TOPIC_ALREADY_EXISTS is treated as
// success — this makes the Job safe to run on every deploy (idempotent).
func createTopics(ctx context.Context, adm *kadm.Client, topics []topicSpec, log *zap.Logger) error {
for _, t := range topics {
resp, err := adm.CreateTopics(ctx, t.Partitions, t.ReplicationFactor, nil, t.Name)
if err != nil {
return fmt.Errorf("create topic %q: %w", t.Name, err)
}
for _, r := range resp {
if r.Err != nil && !errors.Is(r.Err, kerr.TopicAlreadyExists) {
return fmt.Errorf("create topic %q: %w", r.Topic, r.Err)
}
if errors.Is(r.Err, kerr.TopicAlreadyExists) {
log.Debug("topic already exists, skipping", zap.String("topic", r.Topic))
} else {
log.Info("topic created", zap.String("topic", r.Topic),
zap.Int32("partitions", t.Partitions),
zap.Int16("replication_factor", t.ReplicationFactor),
)
}
}
}
return nil
}

// parseTopics parses the KAFKA_TOPICS env var value.
// Format: "name:partitions:replication[,name:partitions:replication,...]"
// Example: "students.onboarding.ready.dlq:1:1"
func parseTopics(raw string) ([]topicSpec, error) {
var topics []topicSpec
for _, entry := range strings.Split(raw, ",") {
entry = strings.TrimSpace(entry)
if entry == "" {
continue
}
parts := strings.Split(entry, ":")
if len(parts) != 3 {
return nil, fmt.Errorf("entry %q must have format name:partitions:replication", entry)
}
name := strings.TrimSpace(parts[0])
if name == "" {
return nil, fmt.Errorf("entry %q has empty topic name", entry)
}
partitions, err := strconv.ParseInt(strings.TrimSpace(parts[1]), 10, 32)
if err != nil || partitions <= 0 {
return nil, fmt.Errorf("entry %q has invalid partitions %q: must be a positive integer", entry, parts[1])
}
replication, err := strconv.ParseInt(strings.TrimSpace(parts[2]), 10, 16)
if err != nil || replication <= 0 {
return nil, fmt.Errorf("entry %q has invalid replication factor %q: must be a positive integer", entry, parts[2])
}
topics = append(topics, topicSpec{
Name: name,
Partitions: int32(partitions),
ReplicationFactor: int16(replication),
})
}
if len(topics) == 0 {
return nil, fmt.Errorf("KAFKA_TOPICS is empty or contains no valid entries")
}
return topics, nil
}

// isTransient returns true for errors that are worth retrying (network-level
// failures). Auth errors, invalid configs, and Kafka protocol errors that
// indicate a permanent condition are not retried.
func isTransient(err error) bool {
if err == nil {
return false
}
// Prefer typed checks over string matching: net.Error covers dial failures,
// connection refused, i/o timeouts, and DNS errors.
var netErr net.Error
if errors.As(err, &netErr) {
return true
}
// io.EOF can occur when the broker closes the connection before responding.
return errors.Is(err, io.EOF)
}

func splitBrokers(s string) []string {
var brokers []string
for _, b := range strings.Split(s, ",") {
if addr := strings.TrimSpace(b); addr != "" {
brokers = append(brokers, addr)
}
}
return brokers
}

// requiredEnv reads an environment variable or returns an error if it is empty.
// A missing variable means the K8s Job is misconfigured — fail fast.
func requiredEnv(key string) (string, error) {
v := os.Getenv(key)
if v == "" {
return "", fmt.Errorf("required environment variable %s is not set", key)
}
return v, nil
}

func envOrDefault(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}
Loading