diff --git a/..bfg-report/2026-09-13/11-02-26/cache-stats.txt b/..bfg-report/2026-09-13/11-02-26/cache-stats.txt new file mode 100644 index 0000000..31eb65d --- /dev/null +++ b/..bfg-report/2026-09-13/11-02-26/cache-stats.txt @@ -0,0 +1,4 @@ +(apply,CacheStats{hitCount=337, missCount=199, loadSuccessCount=199, loadExceptionCount=0, totalLoadTime=581291927, evictionCount=0}) +(tree,CacheStats{hitCount=986, missCount=352, loadSuccessCount=299, loadExceptionCount=0, totalLoadTime=821650758, evictionCount=0}) +(commit,CacheStats{hitCount=108, missCount=107, loadSuccessCount=107, loadExceptionCount=0, totalLoadTime=78983052, evictionCount=0}) +(tag,CacheStats{hitCount=0, missCount=2, loadSuccessCount=2, loadExceptionCount=0, totalLoadTime=319542, evictionCount=0}) diff --git a/..bfg-report/2026-09-13/11-02-26/changed-files.txt b/..bfg-report/2026-09-13/11-02-26/changed-files.txt new file mode 100644 index 0000000..dbab69a --- /dev/null +++ b/..bfg-report/2026-09-13/11-02-26/changed-files.txt @@ -0,0 +1,4 @@ +e71e5b78236a67327c678490cb50b46981f19de0 bbcbb68b91e786eb71bbb0a4443d7b8a26140e1b .sops.yaml +4189696f5581ac0ffdc125c3bf9b9f664b3ddfb0 7cd3f1ee4865c563d141464f6fc185436993b84b .sops.yaml +635630e73152a5f22e6cbd42322ec55d79f8d9c0 297e94a89d73d18c4f47013bb0e8303f123715f3 configmap.yaml +29e515e7b46742fab8c3fcc2189af7010a6ccc62 6869fa11f96e03f7ec76a0ea14a4ddaf604004a4 gateway-config-secret.enc.yaml diff --git a/..bfg-report/2026-09-13/11-02-26/object-id-map.old-new.txt b/..bfg-report/2026-09-13/11-02-26/object-id-map.old-new.txt new file mode 100644 index 0000000..121c322 --- /dev/null +++ b/..bfg-report/2026-09-13/11-02-26/object-id-map.old-new.txt @@ -0,0 +1,12 @@ +0a95af80c0051bacbeb8483c1632e47acd3db5be 40207e487cfb63409a976fb2a0b9e1e62c8b1513 +27428d910111299d0699f429190284a9ca6e50b7 3318daf758349402aef43b095482743ab96b37f9 +329a495af4c935529fdae17229314101c0c77876 67f24ea76359c8dba4b56267790aad76bbc58464 +4c8bc6c920b6b75399555827022f69ef0c4f7d15 1fa839b41975fa3f0ac9052355ffb625f5a8f324 +528545f414c83217408edfea234dcd1f3edee0c2 b8f95506ca1545b876b5531cd385172e9ca5b4b0 +81038e1cf7567a9133d7c233a97b1e2f19fa1c82 4a00312906ba725f3968187656fde2663b1763ab +a5b3b5c44a406896bcb414df6c6426c277715706 2ab47a9dbe5ba36dfa0e275991ef7b7656908410 +ce27643667a0399115cd1f2b6d38123fdcf2b4f1 6ff0a50de8efbad105fa588245f22fdb26afddc4 +d49756886a46542b38533b913a1f776b5145f5ec d82cc5a6970a1fb32e21dda9a737b987a8668111 +db3a30fbcf1f139c667fb68a91762582c49b8cee 04619a269fed9eeea53ab4d4d73131e3713f40a0 +eb54715e4dec0fb35402576fcc224a09808b00c1 d53b7632cf9646dda1c978a5f94615dc9eaed5e8 +ef72b5bbccf2df89aa1c86dee29311c63f33bf62 ba55d184fefef1a73a50409ca4fb1f7b27f5b075 diff --git a/examples/sendmsg-email.sh b/examples/sendmsg-email.sh new file mode 100644 index 0000000..e8e1ed4 --- /dev/null +++ b/examples/sendmsg-email.sh @@ -0,0 +1,36 @@ +#!/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": "admin@example.com", + "cc": "ops-team@example.com" + } +} +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 "" diff --git a/internal/notification/handler.go b/internal/notification/handler.go new file mode 100644 index 0000000..4493257 --- /dev/null +++ b/internal/notification/handler.go @@ -0,0 +1,160 @@ +package notification + +import ( + "encoding/json" + "fmt" + "log" + "net/http" + "net/smtp" + "os" +) + +// 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": "user@example.com", "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) + 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", + } +} diff --git a/internal/proxy/proxy.go b/internal/proxy/proxy.go index aaaeeab..9373558 100644 --- a/internal/proxy/proxy.go +++ b/internal/proxy/proxy.go @@ -11,6 +11,7 @@ import ( "net/url" "sort" "strings" + "syscall" "time" "forgejo.riotpiao.com/rock/homelab-frontend/internal/auth" @@ -127,6 +128,15 @@ func (h *Handler) getOrCreateTransport(addr string, up *config.Upstream) *http.T dialer := &net.Dialer{ Timeout: up.ConnectTimeout, 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{ @@ -134,8 +144,15 @@ func (h *Handler) getOrCreateTransport(addr string, up *config.Upstream) *http.T DialContext: dialer.DialContext, MaxIdleConns: 100, 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 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 @@ -429,6 +446,12 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { // A streaming response that's continuously sending should not be cut off. // 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 proxy.ServeHTTP(w, r) } diff --git a/internal/proxy/streaming_test.go b/internal/proxy/streaming_test.go index a8103e5..db3dbcd 100644 --- a/internal/proxy/streaming_test.go +++ b/internal/proxy/streaming_test.go @@ -476,6 +476,214 @@ 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, // the upstream request context is cancelled immediately and no goroutines are leaked. func TestClientDisconnectCancelsUpstream(t *testing.T) { diff --git a/internal/server/server.go b/internal/server/server.go index b616f7e..8a9fdcd 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -6,6 +6,8 @@ import ( "net/http" "sync" "time" + + "golang.org/x/net/http2" ) // Server wraps an HTTP server with graceful shutdown support. @@ -19,20 +21,30 @@ type Server struct { // New creates a new Server with the given configuration. 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{ - 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, - }, + httpServer: httpServer, shutdownTimeout: shutdownTimeout, healthChecker: NewHealthChecker(false, false), } diff --git a/k8s/kustomization.yaml.bak b/k8s/kustomization.yaml.bak new file mode 100644 index 0000000..4bd1abd --- /dev/null +++ b/k8s/kustomization.yaml.bak @@ -0,0 +1,27 @@ +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: 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 diff --git a/k8s/smtp-secrets.example.yaml b/k8s/smtp-secrets.example.yaml new file mode 100644 index 0000000..2032c89 --- /dev/null +++ b/k8s/smtp-secrets.example.yaml @@ -0,0 +1,29 @@ +apiVersion: v1 +kind: Secret +metadata: + name: smtp-credentials + namespace: api + labels: + app: api-gateway + component: notification +type: Opaque +data: + host: + port: + from: + user: + password: + +# To create from plaintext: +# kubectl create secret generic smtp-credentials \ +# --from-literal=host=mail.example.com \ +# --from-literal=port=587 \ +# --from-literal=from=noreply@example.com \ +# --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