From 4831c74acc8a641f02fe056ad0d0ddbb99c1ca4d Mon Sep 17 00:00:00 2001 From: Admin Date: Fri, 27 Mar 2026 16:05:22 +0500 Subject: [PATCH] feat(observability+tts): OTel logs for backend/runner/ui; add kokoro-fastapi (GPU) and pocket-tts (CPU) to homelab - otelsetup.Init now returns a *slog.Logger wired to the OTLP log exporter so all slog output is shipped to Loki with embedded trace IDs - backend and runner both adopt the new OTel-bridged logger - runner.runScrapeTask and runAudioTask now emit structured OTel spans - ui/hooks.server.ts adds BatchLogRecordProcessor alongside existing trace exporter - homelab: add kokoro-fastapi GPU service (ghcr.io/remsky/kokoro-fastapi-gpu) using deploy.resources.reservations for NVIDIA GPU, exposed internally on :8880 - homelab: add pocket-tts CPU service (ghcr.io/kyutai-labs/pocket-tts) on :8000 - runner KOKORO_URL hardcoded to http://kokoro-fastapi:8880 (fixes DNS failure for the stale kokoro.kalekber.cc hostname) --- backend/cmd/backend/main.go | 9 ++-- backend/cmd/runner/main.go | 14 +++++ backend/go.mod | 4 ++ backend/go.sum | 8 +++ backend/internal/otelsetup/otelsetup.go | 72 ++++++++++++++++++------- backend/internal/runner/runner.go | 25 +++++++++ homelab/docker-compose.yml | 44 ++++++++++++++- ui/package-lock.json | 1 + ui/package.json | 1 + ui/src/hooks.server.ts | 11 +++- 10 files changed, 164 insertions(+), 25 deletions(-) diff --git a/backend/cmd/backend/main.go b/backend/cmd/backend/main.go index 11ae1ca..eda95a4 100644 --- a/backend/cmd/backend/main.go +++ b/backend/cmd/backend/main.go @@ -71,14 +71,17 @@ func run() error { ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - // ── OpenTelemetry tracing ───────────────────────────────────────────────── - otelShutdown, err := otelsetup.Init(ctx, version) + // ── OpenTelemetry tracing + logs ────────────────────────────────────────── + otelShutdown, otelLog, err := otelsetup.Init(ctx, version) if err != nil { return fmt.Errorf("init otel: %w", err) } if otelShutdown != nil { defer otelShutdown() - log.Info("otel tracing enabled", "endpoint", os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT")) + // Replace the plain slog logger with the OTel-bridged one so all + // structured log lines are forwarded to Loki with trace IDs attached. + log = otelLog + log.Info("otel tracing + logs enabled", "endpoint", os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT")) } // ── Storage ────────────────────────────────────────────────────────────── diff --git a/backend/cmd/runner/main.go b/backend/cmd/runner/main.go index 45ebef4..5b08732 100644 --- a/backend/cmd/runner/main.go +++ b/backend/cmd/runner/main.go @@ -25,6 +25,7 @@ import ( "github.com/libnovel/backend/internal/kokoro" "github.com/libnovel/backend/internal/meili" "github.com/libnovel/backend/internal/novelfire" + "github.com/libnovel/backend/internal/otelsetup" "github.com/libnovel/backend/internal/runner" "github.com/libnovel/backend/internal/storage" ) @@ -70,6 +71,19 @@ func run() error { ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() + // ── OpenTelemetry tracing + logs ───────────────────────────────────────── + otelShutdown, otelLog, err := otelsetup.Init(ctx, version) + if err != nil { + return fmt.Errorf("init otel: %w", err) + } + if otelShutdown != nil { + defer otelShutdown() + // Switch to the OTel-bridged logger so all structured log lines are + // forwarded to Loki with trace IDs attached. + log = otelLog + log.Info("otel tracing + logs enabled", "endpoint", os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT")) + } + // ── Storage ───────────────────────────────────────────────────────────── store, err := storage.NewStore(ctx, cfg, log) if err != nil { diff --git a/backend/go.mod b/backend/go.mod index 3429906..bbf3a10 100644 --- a/backend/go.mod +++ b/backend/go.mod @@ -34,12 +34,16 @@ require ( github.com/rs/xid v1.6.0 // indirect github.com/tinylib/msgp v1.6.1 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/contrib/bridges/otelslog v0.17.0 // indirect go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 // indirect go.opentelemetry.io/otel v1.42.0 // indirect + go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.18.0 // indirect go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.42.0 // indirect go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.42.0 // indirect + go.opentelemetry.io/otel/log v0.18.0 // indirect go.opentelemetry.io/otel/metric v1.42.0 // indirect go.opentelemetry.io/otel/sdk v1.42.0 // indirect + go.opentelemetry.io/otel/sdk/log v0.18.0 // indirect go.opentelemetry.io/otel/trace v1.42.0 // indirect go.opentelemetry.io/proto/otlp v1.9.0 // indirect go.uber.org/atomic v1.11.0 // indirect diff --git a/backend/go.sum b/backend/go.sum index 962286b..bc91505 100644 --- a/backend/go.sum +++ b/backend/go.sum @@ -58,18 +58,26 @@ github.com/tinylib/msgp v1.6.1/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77ro github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/contrib/bridges/otelslog v0.17.0 h1:NFIS6x7wyObQ7cR84x7bt1sr8nYBx89s3x3GwRjw40k= +go.opentelemetry.io/contrib/bridges/otelslog v0.17.0/go.mod h1:39SaByOyDMRMe872AE7uelMuQZidIw7LLFAnQi0FWTE= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 h1:OyrsyzuttWTSur2qN/Lm0m2a8yqyIjUVBZcxFPuXq2o= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0/go.mod h1:C2NGBr+kAB4bk3xtMXfZ94gqFDtg/GkI7e9zqGh5Beg= go.opentelemetry.io/otel v1.42.0 h1:lSQGzTgVR3+sgJDAU/7/ZMjN9Z+vUip7leaqBKy4sho= go.opentelemetry.io/otel v1.42.0/go.mod h1:lJNsdRMxCUIWuMlVJWzecSMuNjE7dOYyWlqOXWkdqCc= +go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.18.0 h1:icqq3Z34UrEFk2u+HMhTtRsvo7Ues+eiJVjaJt62njs= +go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.18.0/go.mod h1:W2m8P+d5Wn5kipj4/xmbt9uMqezEKfBjzVJadfABSBE= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.42.0 h1:THuZiwpQZuHPul65w4WcwEnkX2QIuMT+UFoOrygtoJw= go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.42.0/go.mod h1:J2pvYM5NGHofZ2/Ru6zw/TNWnEQp5crgyDeSrYpXkAw= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.42.0 h1:uLXP+3mghfMf7XmV4PkGfFhFKuNWoCvvx5wP/wOXo0o= go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.42.0/go.mod h1:v0Tj04armyT59mnURNUJf7RCKcKzq+lgJs6QSjHjaTc= +go.opentelemetry.io/otel/log v0.18.0 h1:XgeQIIBjZZrliksMEbcwMZefoOSMI1hdjiLEiiB0bAg= +go.opentelemetry.io/otel/log v0.18.0/go.mod h1:KEV1kad0NofR3ycsiDH4Yjcoj0+8206I6Ox2QYFSNgI= go.opentelemetry.io/otel/metric v1.42.0 h1:2jXG+3oZLNXEPfNmnpxKDeZsFI5o4J+nz6xUlaFdF/4= go.opentelemetry.io/otel/metric v1.42.0/go.mod h1:RlUN/7vTU7Ao/diDkEpQpnz3/92J9ko05BIwxYa2SSI= go.opentelemetry.io/otel/sdk v1.42.0 h1:LyC8+jqk6UJwdrI/8VydAq/hvkFKNHZVIWuslJXYsDo= go.opentelemetry.io/otel/sdk v1.42.0/go.mod h1:rGHCAxd9DAph0joO4W6OPwxjNTYWghRWmkHuGbayMts= +go.opentelemetry.io/otel/sdk/log v0.18.0 h1:n8OyZr7t7otkeTnPTbDNom6rW16TBYGtvyy2Gk6buQw= +go.opentelemetry.io/otel/sdk/log v0.18.0/go.mod h1:C0+wxkTwKpOCZLrlJ3pewPiiQwpzycPI/u6W0Z9fuYk= go.opentelemetry.io/otel/trace v1.42.0 h1:OUCgIPt+mzOnaUTpOQcBiM/PLQ/Op7oq6g4LenLmOYY= go.opentelemetry.io/otel/trace v1.42.0/go.mod h1:f3K9S+IFqnumBkKhRJMeaZeNk9epyhnCmQh/EysQCdc= go.opentelemetry.io/proto/otlp v1.9.0 h1:l706jCMITVouPOqEnii2fIAuO3IVGBRPV5ICjceRb/A= diff --git a/backend/internal/otelsetup/otelsetup.go b/backend/internal/otelsetup/otelsetup.go index 6e053ce..f457501 100644 --- a/backend/internal/otelsetup/otelsetup.go +++ b/backend/internal/otelsetup/otelsetup.go @@ -5,12 +5,13 @@ // OTEL_EXPORTER_OTLP_ENDPOINT — OTLP/HTTP endpoint, e.g. http://otel-collector:4318 // OTEL_SERVICE_NAME — service name reported in traces (default: "backend") // -// When OTEL_EXPORTER_OTLP_ENDPOINT is empty the function is a no-op and -// returns a nil shutdown func, so the caller never needs to branch on it. +// When OTEL_EXPORTER_OTLP_ENDPOINT is empty the function is a no-op: it +// returns a nil shutdown func and the default slog.Logger, so callers never +// need to branch on it. // // Usage in main.go: // -// shutdown, err := otelsetup.Init(ctx, version) +// shutdown, log, err := otelsetup.Init(ctx, version) // if err != nil { return err } // if shutdown != nil { defer shutdown() } package otelsetup @@ -18,22 +19,31 @@ package otelsetup import ( "context" "fmt" + "log/slog" "os" "time" + "go.opentelemetry.io/contrib/bridges/otelslog" "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp" "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp" + otellog "go.opentelemetry.io/otel/log/global" + "go.opentelemetry.io/otel/sdk/log" "go.opentelemetry.io/otel/sdk/resource" sdktrace "go.opentelemetry.io/otel/sdk/trace" semconv "go.opentelemetry.io/otel/semconv/v1.26.0" ) -// Init sets up a TracerProvider that exports spans via OTLP/HTTP. -// Returns a shutdown function (or nil if OTel is disabled) and any error. -func Init(ctx context.Context, version string) (shutdown func(), err error) { +// Init sets up TracerProvider and LoggerProvider that export via OTLP/HTTP. +// +// Returns: +// - shutdown: flushes and stops both providers (nil when OTel is disabled). +// - logger: an slog.Logger bridged to OTel logs (falls back to default when disabled). +// - err: non-nil only on SDK initialisation failure. +func Init(ctx context.Context, version string) (shutdown func(), logger *slog.Logger, err error) { endpoint := os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT") if endpoint == "" { - return nil, nil // OTel disabled — not an error + return nil, slog.Default(), nil // OTel disabled — not an error } serviceName := os.Getenv("OTEL_SERVICE_NAME") @@ -41,14 +51,7 @@ func Init(ctx context.Context, version string) (shutdown func(), err error) { serviceName = "backend" } - exp, err := otlptracehttp.New(ctx, - otlptracehttp.WithEndpoint(endpoint), - otlptracehttp.WithInsecure(), // collector is on the internal Docker network - ) - if err != nil { - return nil, fmt.Errorf("otelsetup: create OTLP exporter: %w", err) - } - + // ── Shared resource ─────────────────────────────────────────────────────── res, err := resource.New(ctx, resource.WithAttributes( semconv.ServiceName(serviceName), @@ -56,19 +59,50 @@ func Init(ctx context.Context, version string) (shutdown func(), err error) { ), ) if err != nil { - return nil, fmt.Errorf("otelsetup: create resource: %w", err) + return nil, slog.Default(), fmt.Errorf("otelsetup: create resource: %w", err) + } + + // ── Trace provider ──────────────────────────────────────────────────────── + traceExp, err := otlptracehttp.New(ctx, + otlptracehttp.WithEndpoint(endpoint), + otlptracehttp.WithInsecure(), // collector is on the internal Docker network + ) + if err != nil { + return nil, slog.Default(), fmt.Errorf("otelsetup: create OTLP trace exporter: %w", err) } tp := sdktrace.NewTracerProvider( - sdktrace.WithBatcher(exp), + sdktrace.WithBatcher(traceExp), sdktrace.WithResource(res), sdktrace.WithSampler(sdktrace.ParentBased(sdktrace.TraceIDRatioBased(0.2))), ) otel.SetTracerProvider(tp) - return func() { + // ── Log provider ────────────────────────────────────────────────────────── + logExp, err := otlploghttp.New(ctx, + otlploghttp.WithEndpoint(endpoint), + otlploghttp.WithInsecure(), + ) + if err != nil { + return nil, slog.Default(), fmt.Errorf("otelsetup: create OTLP log exporter: %w", err) + } + + lp := log.NewLoggerProvider( + log.WithProcessor(log.NewBatchProcessor(logExp)), + log.WithResource(res), + ) + otellog.SetLoggerProvider(lp) + + // Bridge slog → OTel logs. Structured fields and trace IDs are forwarded + // automatically; Grafana can correlate log lines with Tempo traces. + otelLogger := otelslog.NewLogger(serviceName) + + shutdown = func() { shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() _ = tp.Shutdown(shutCtx) - }, nil + _ = lp.Shutdown(shutCtx) + } + + return shutdown, otelLogger, nil } diff --git a/backend/internal/runner/runner.go b/backend/internal/runner/runner.go index f9c85a4..48be6a0 100644 --- a/backend/internal/runner/runner.go +++ b/backend/internal/runner/runner.go @@ -22,6 +22,10 @@ import ( "sync/atomic" "time" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + "github.com/libnovel/backend/internal/bookstore" "github.com/libnovel/backend/internal/domain" "github.com/libnovel/backend/internal/kokoro" @@ -301,6 +305,14 @@ func (r *Runner) newOrchestrator() *orchestrator.Orchestrator { // runScrapeTask executes one scrape task end-to-end and reports the result. func (r *Runner) runScrapeTask(ctx context.Context, task domain.ScrapeTask) { + ctx, span := otel.Tracer("runner").Start(ctx, "runner.scrape_task") + defer span.End() + span.SetAttributes( + attribute.String("task.id", task.ID), + attribute.String("task.kind", task.Kind), + attribute.String("task.url", task.TargetURL), + ) + log := r.deps.Log.With("task_id", task.ID, "kind", task.Kind, "url", task.TargetURL) log.Info("runner: scrape task starting") @@ -340,8 +352,10 @@ func (r *Runner) runScrapeTask(ctx context.Context, task domain.ScrapeTask) { if result.ErrorMessage != "" { r.tasksFailed.Add(1) + span.SetStatus(codes.Error, result.ErrorMessage) } else { r.tasksCompleted.Add(1) + span.SetStatus(codes.Ok, "") } log.Info("runner: scrape task finished", @@ -384,6 +398,15 @@ func (r *Runner) runCatalogueTask(ctx context.Context, task domain.ScrapeTask, o // runAudioTask executes one audio-generation task. func (r *Runner) runAudioTask(ctx context.Context, task domain.AudioTask) { + ctx, span := otel.Tracer("runner").Start(ctx, "runner.audio_task") + defer span.End() + span.SetAttributes( + attribute.String("task.id", task.ID), + attribute.String("book.slug", task.Slug), + attribute.Int("chapter.number", task.Chapter), + attribute.String("audio.voice", task.Voice), + ) + log := r.deps.Log.With("task_id", task.ID, "slug", task.Slug, "chapter", task.Chapter, "voice", task.Voice) log.Info("runner: audio task starting") @@ -407,6 +430,7 @@ func (r *Runner) runAudioTask(ctx context.Context, task domain.AudioTask) { fail := func(msg string) { log.Error("runner: audio task failed", "reason", msg) r.tasksFailed.Add(1) + span.SetStatus(codes.Error, msg) result := domain.AudioResult{ErrorMessage: msg} if err := r.deps.Consumer.FinishAudioTask(ctx, task.ID, result); err != nil { log.Error("runner: FinishAudioTask failed", "err", err) @@ -441,6 +465,7 @@ func (r *Runner) runAudioTask(ctx context.Context, task domain.AudioTask) { } r.tasksCompleted.Add(1) + span.SetStatus(codes.Ok, "") result := domain.AudioResult{ObjectKey: key} if err := r.deps.Consumer.FinishAudioTask(ctx, task.ID, result); err != nil { log.Error("runner: FinishAudioTask failed", "err", err) diff --git a/homelab/docker-compose.yml b/homelab/docker-compose.yml index 2f419e9..8d52dac 100644 --- a/homelab/docker-compose.yml +++ b/homelab/docker-compose.yml @@ -58,7 +58,7 @@ services: VALKEY_ADDR: "" GODEBUG: "preferIPv4=1" - KOKORO_URL: "${KOKORO_URL}" + KOKORO_URL: "http://kokoro-fastapi:8880" KOKORO_VOICE: "${KOKORO_VOICE}" RUNNER_WORKER_ID: "${RUNNER_WORKER_ID}" @@ -386,6 +386,48 @@ services: timeout: 5s retries: 5 + # ── Kokoro-FastAPI (GPU TTS) ──────────────────────────────────────────────── + # OpenAI-compatible TTS service backed by the Kokoro model, running on the + # homelab RTX 3050 (8 GB VRAM). Replaces the broken kokoro.kalekber.cc DNS. + # Voices match existing IDs: af_bella, af_sky, af_heart, etc. + # The runner reaches it at http://kokoro-fastapi:8880 via the Docker network. + kokoro-fastapi: + image: ghcr.io/remsky/kokoro-fastapi-gpu:latest + restart: unless-stopped + deploy: + resources: + reservations: + devices: + - driver: nvidia + count: 1 + capabilities: [gpu] + expose: + - "8880" + healthcheck: + test: ["CMD", "curl", "-sf", "http://localhost:8880/health"] + interval: 30s + timeout: 10s + retries: 5 + start_period: 60s + + # ── pocket-tts (CPU TTS) ──────────────────────────────────────────────────── + # Lightweight CPU-only TTS using kyutai-labs/pocket-tts. + # OpenAI-compatible: POST /v1/audio/speech on port 8000. + # Voices: alba, marius, javert, jean, fantine, cosette, eponine, azelma. + # Not currently used by the runner (runner uses kokoro-fastapi), but available + # for experimentation / fallback. + pocket-tts: + image: ghcr.io/kyutai-labs/pocket-tts:latest + restart: unless-stopped + expose: + - "8000" + healthcheck: + test: ["CMD", "curl", "-sf", "http://localhost:8000/health"] + interval: 30s + timeout: 10s + retries: 5 + start_period: 60s + # ── Watchtower ────────────────────────────────────────────────────────────── # Auto-updates runner image when CI pushes a new tag. # Only watches services with the watchtower label. diff --git a/ui/package-lock.json b/ui/package-lock.json index 1138892..738c7d3 100644 --- a/ui/package-lock.json +++ b/ui/package-lock.json @@ -10,6 +10,7 @@ "dependencies": { "@aws-sdk/client-s3": "^3.1005.0", "@aws-sdk/s3-request-presigner": "^3.1005.0", + "@opentelemetry/exporter-logs-otlp-http": "^0.214.0", "@opentelemetry/exporter-trace-otlp-http": "^0.214.0", "@opentelemetry/resources": "^2.6.1", "@opentelemetry/sdk-node": "^0.214.0", diff --git a/ui/package.json b/ui/package.json index 27e1b4f..c7c8a43 100644 --- a/ui/package.json +++ b/ui/package.json @@ -30,6 +30,7 @@ "dependencies": { "@aws-sdk/client-s3": "^3.1005.0", "@aws-sdk/s3-request-presigner": "^3.1005.0", + "@opentelemetry/exporter-logs-otlp-http": "^0.214.0", "@opentelemetry/exporter-trace-otlp-http": "^0.214.0", "@opentelemetry/resources": "^2.6.1", "@opentelemetry/sdk-node": "^0.214.0", diff --git a/ui/src/hooks.server.ts b/ui/src/hooks.server.ts index bb4d6d7..cc3cafb 100644 --- a/ui/src/hooks.server.ts +++ b/ui/src/hooks.server.ts @@ -9,10 +9,12 @@ import { createUserSession, touchUserSession, isSessionRevoked } from '$lib/serv import { drain as drainPresignCache } from '$lib/server/presignCache'; import { NodeSDK } from '@opentelemetry/sdk-node'; import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-http'; +import { OTLPLogExporter } from '@opentelemetry/exporter-logs-otlp-http'; +import { BatchLogRecordProcessor } from '@opentelemetry/sdk-logs'; import { resourceFromAttributes } from '@opentelemetry/resources'; import { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } from '@opentelemetry/semantic-conventions'; -// ─── OpenTelemetry server-side tracing ──────────────────────────────────────── +// ─── OpenTelemetry server-side tracing + logs ───────────────────────────────── // No-op when OTEL_EXPORTER_OTLP_ENDPOINT is unset (e.g. local dev). const otlpEndpoint = process.env.OTEL_EXPORTER_OTLP_ENDPOINT; if (otlpEndpoint) { @@ -21,7 +23,12 @@ if (otlpEndpoint) { [ATTR_SERVICE_NAME]: process.env.OTEL_SERVICE_NAME ?? 'ui', [ATTR_SERVICE_VERSION]: pubEnv.PUBLIC_BUILD_VERSION ?? 'dev' }), - traceExporter: new OTLPTraceExporter({ url: `${otlpEndpoint}/v1/traces` }) + traceExporter: new OTLPTraceExporter({ url: `${otlpEndpoint}/v1/traces` }), + logRecordProcessors: [ + new BatchLogRecordProcessor( + new OTLPLogExporter({ url: `${otlpEndpoint}/v1/logs` }) + ) + ] }); sdk.start(); process.once('SIGTERM', () => sdk.shutdown().catch(() => {}));