Files
silo-server/internal/audiobooks/abs/session_local.go
203a18ae83 feat(observability): OpenTelemetry logs+traces with secret redaction and slog standardization (#290)
* feat(observability): OpenTelemetry logs+traces with secret redaction

Part of #265. Adds opt-in OpenTelemetry (logs + traces) alongside the existing
stderr + opslog pipeline, plus secret redaction on all sinks. Default-off: with
no OTEL_* / SILO_OTEL_ENABLED config, behavior is unchanged.

Bootstrap (internal/telemetry):
- Setup() builds one shared resource, a TracerProvider (parent-based trace-id
  ratio sampler), a LoggerProvider, and the W3C TraceContext+Baggage propagator
  from env. It installs NO MeterProvider — metrics stay on Prometheus, and the
  built-in no-op global MeterProvider keeps the trace instrumentation libs from
  double-emitting. Shutdown is deferred with a flush timeout.
- Logs are bridged via otelslog fan-out (slog.MultiHandler), level-gated by the
  shared LevelVar and best-effort so a failing collector can't break the console
  or DB branches. stderr + opslog stay untouched.

Secret redaction (internal/logredact):
- A slog.Handler masks secret-keyed attributes (password, token, api_key,
  authorization, cookie, ...) — including .With-bound attrs, nested groups,
  secret-keyed group subtrees, and values behind a LogValuer — on the console
  and OTLP sinks, with a no-op fast path when a record has no secret keys.
  opslog.shouldRedact delegates to logredact.SecretKey so all sinks share one
  marker list.

Rotation is infra-managed (no custom file sink): container runtime for stderr,
collector/backend for OTLP, opslog partition-pruning for the DB. Documented in
docs/architecture/observability.md.

Verification: go build ./..., go vet, gofmt -l — clean; go test
./internal/telemetry/ ./internal/logredact/ -race pass.

AI-use disclosure: implemented with AI assistance (Claude Code), including
adversarial reviews that hardened the bootstrap and fixed two redaction leak
paths; reviewed by the author.

* refactor(observability): slog context+component sweep, sloglint gate (phase 3)

Part of #265. Builds on the OTel bootstrap + redaction commit.

Standardizes every log call site onto the context-carrying slog variants so
records correlate with the active OpenTelemetry trace, and locks the standard
in with a machine gate so future code (human- or AI-authored) can't drift back.

- Call-site sweep: converted the remaining slog.<Level>(...) calls to the
  slog.<Level>Context(ctx, ...) form wherever a context.Context is in scope
  (background/init calls with no ctx are left as-is), across 183 files. Applied
  via a type-aware AST codemod. Log levels and message strings are preserved
  verbatim; a component attr (canonical per-package name) is added to direct
  package-level slog calls. Bound-logger calls keep their existing .With
  bindings. The main.go and telemetry package conversions rode with their file
  in the previous commit to keep each file within a single commit.
- Enforcement (.golangci.yml): enable sloglint with context=scope, static-msg,
  key-naming-case=snake, no-mixed-args. After the sweep all four report zero
  violations repo-wide (tests included), so make lint / CI now blocks any
  regression to the non-context form. The gate ships with the sweep because it
  cannot be green until the legacy sites are converted.

Metrics remain on Prometheus; no behavior change to /metrics or Grafana.

Verification: go build ./..., go vet ./..., gofmt -l — clean; sloglint (all 4
rules) 0 violations repo-wide; log levels verified unchanged.

AI-use disclosure: implemented with AI assistance (Claude Code), including the
codemod; reviewed by the author.

* fix(observability): honor per-signal OTLP protocol and secret WithGroup names

Two Codex review findings on PR #290:

- telemetry: OTEL_EXPORTER_OTLP_{TRACES,LOGS}_PROTOCOL now override the
  generic OTEL_EXPORTER_OTLP_PROTOCOL per signal, so mixed collector
  setups (e.g. HTTP logs + gRPC traces) build the right exporter.
- logredact: entering a group whose name is secret-bearing (e.g.
  WithGroup("authorization")) now masks every leaf in that subtree,
  matching how slog.Group("authorization", ...) is masked as a whole.

* fix(observability): address review feedback on telemetry bootstrap

- Telemetry setup failure no longer kills boot: Setup returns usable
  no-op providers alongside the error and main logs and continues with
  telemetry disabled, honoring the best-effort contract.
- Honor OTEL_TRACES_SAMPLER (always_on/off, traceidratio, parentbased_*
  variants); unsupported values fall back to parentbased_traceidratio.
- Attach node identity as semconv service.instance.id instead of the
  non-semconv node.name.
- Rename opslog retention-scope log attrs to target_component/target_level
  so they no longer collide with the canonical component routing key, and
  tag those lines with component=opslog.
- Fix stale levelGated comment casing; use WarnContext in the telemetry
  shutdown defer; document the LogValuer double-resolve on the redaction
  slow path.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Quick <31828688+Quick104@users.noreply.github.com>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2026-07-09 08:53:52 -04:00

214 lines
8.6 KiB
Go

package abs
import (
"context"
"encoding/json"
"io"
"log/slog"
"net/http"
"time"
"github.com/Silo-Server/silo-server/internal/catalog"
)
// ---------------------------------------------------------------------------
// Offline session sync (ABS SessionController.syncLocal / syncLocalSessions)
// ---------------------------------------------------------------------------
//
// The official ABS mobile app records playback while offline into local
// PlaybackSession objects, then POSTs them back on reconnect so the server's
// media progress catches up. Two endpoints implement this:
//
// POST /session/local — one session (SessionController.syncLocal)
// POST /session/local-all — many sessions (SessionController.syncLocalSessions)
//
// Real ABS (server/managers/PlaybackSessionManager.js):
// - syncLocalSessionRequest: syncLocalSession(one) → 200 on success,
// 500 + error text on failure.
// - syncLocalSessionsRequest: reads req.body.sessions, loops each through
// syncLocalSession, replies { results: [ {id, success, error?,
// progressSynced} ] } with HTTP 200 regardless of per-session outcome.
// - syncLocalSession returns {id, success:false, error} when the library
// item can't be found, else {id, success:true, progressSynced}.
//
// silo does not persist arbitrary client-supplied sessions, so we do not
// create a server-side playback-session row here. We mirror the *effect* that
// matters — the caller's resume position — into user_watch_progress and emit
// the same user_item_progress_updated realtime event handleSessionSync does.
// An existing progress row is advanced monotonically (UpdateProgressPosition);
// a book listened to entirely offline has no row yet, so one is created
// (UpsertProgress) rather than silently dropping the position. Accumulated
// offline listening time is not persisted: there is no store method to add
// standalone listening time without an existing session row. Position sync is
// the client-visible behaviour offline sync exists to restore.
// localPlaybackSession is the subset of the ABS PlaybackSession payload the
// client POSTs for offline sync that silo acts on. EpisodeID is a pointer so a
// present-but-null value (audiobook) is distinguishable from a podcast episode.
type localPlaybackSession struct {
ID string `json:"id"`
LibraryItemID string `json:"libraryItemId"`
EpisodeID *string `json:"episodeId"`
CurrentTime float64 `json:"currentTime"`
TimeListening float64 `json:"timeListening"`
DisplayTitle string `json:"displayTitle"`
}
// localSyncResult mirrors the per-session object ABS returns from
// syncLocalSession: {id, success, error?, progressSynced}.
type localSyncResult struct {
ID string `json:"id"`
Success bool `json:"success"`
Error string `json:"error,omitempty"`
ProgressSynced bool `json:"progressSynced"`
}
// handleSyncLocalSession — POST /session/local
// Syncs a single offline-recorded session. Matches ABS syncLocalSessionRequest:
// 200 on success, 500 + error text when the item can't be resolved.
func (h *Handler) handleSyncLocalSession(w http.ResponseWriter, r *http.Request) {
a, ok := absAuthFrom(r)
if !ok || a.UserID == "" {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
var sess localPlaybackSession
if err := json.NewDecoder(io.LimitReader(r.Body, 1<<20)).Decode(&sess); err != nil {
http.Error(w, "invalid body", http.StatusBadRequest)
return
}
access, err := h.accessFilterForAuth(r.Context(), a)
if err != nil {
http.Error(w, "resolve access: "+err.Error(), http.StatusForbidden)
return
}
res := h.syncOneLocalSession(r.Context(), a, access, sess)
if !res.Success {
// Real ABS: res.status(500).send(result.error).
http.Error(w, res.Error, http.StatusInternalServerError)
return
}
// Real ABS: res.sendStatus(200) — plain 200, no JSON body.
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte("OK"))
}
// handleSyncLocalSessions — POST /session/local-all
// Batch-syncs offline-recorded sessions. Matches ABS syncLocalSessionsRequest:
// reads {sessions: [...]}, loops each, always replies 200 with
// {results: [...]}. A bad session never fails the whole batch.
func (h *Handler) handleSyncLocalSessions(w http.ResponseWriter, r *http.Request) {
a, ok := absAuthFrom(r)
if !ok || a.UserID == "" {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
// Decode into raw messages so one malformed session doesn't sink the batch.
var body struct {
Sessions []json.RawMessage `json:"sessions"`
}
if err := json.NewDecoder(io.LimitReader(r.Body, 1<<20)).Decode(&body); err != nil {
http.Error(w, "invalid body", http.StatusBadRequest)
return
}
access, err := h.accessFilterForAuth(r.Context(), a)
if err != nil {
http.Error(w, "resolve access: "+err.Error(), http.StatusForbidden)
return
}
results := make([]localSyncResult, 0, len(body.Sessions))
for _, raw := range body.Sessions {
var sess localPlaybackSession
if err := json.Unmarshal(raw, &sess); err != nil {
slog.WarnContext(r.Context(), "abs local session sync: skipping malformed session", "component", "audiobooks", "error", err)
results = append(results, localSyncResult{Success: false, Error: "invalid session"})
continue
}
results = append(results, h.syncOneLocalSession(r.Context(), a, access, sess))
}
writeJSON(w, http.StatusOK, map[string]any{"results": results})
}
// syncOneLocalSession applies one offline session's position to the caller's
// media progress and returns the ABS-shaped per-session result. It never
// panics and never propagates an error to fail a batch.
func (h *Handler) syncOneLocalSession(ctx context.Context, a ctxAuth, access catalog.AccessFilter, sess localPlaybackSession) localSyncResult {
res := localSyncResult{ID: sess.ID}
// Podcast episode sessions are out of scope for silo's audiobook-only
// catalog. Accept them as a no-op success so the client clears its queue.
if sess.EpisodeID != nil && *sess.EpisodeID != "" {
res.Success = true
return res
}
if sess.LibraryItemID == "" {
res.Error = "Media item not found"
return res
}
// Ownership / existence / access gate — the item must be visible to the
// caller. Mirrors handleSessionSync's access check.
if h.deps.MediaStore == nil {
res.Error = "Media item not found"
return res
}
item, err := h.deps.MediaStore.GetAudiobookByID(ctx, sess.LibraryItemID, access)
if err != nil || item == nil {
if err != nil {
slog.WarnContext(ctx, "abs local session sync: media lookup failed", "component", "audiobooks",
"library_item_id", sess.LibraryItemID, "error", err)
}
res.Error = "Media item not found"
return res
}
res.Success = true
// Persist the offline resume position into user_watch_progress.
//
// For an existing row we use UpdateProgressPosition (not a full upsert) so a
// stale offline tick can't overwrite is_finished / progress_pct the user set
// explicitly — it advances position monotonically and skips completed rows.
//
// But UpdateProgressPosition is UPDATE-only: when a book was listened to
// entirely offline no row exists yet, so it would affect zero rows, drop the
// position, and still report ProgressSynced=true — the client then clears its
// local session with nothing saved. Create the row in that case so the resume
// point survives.
if h.deps.ProgressStore != nil {
var syncErr error
existing, getErr := h.deps.ProgressStore.GetProgress(ctx, a.UserID, a.ProfileID, sess.LibraryItemID)
if getErr == nil && existing == nil {
syncErr = h.deps.ProgressStore.UpsertProgress(ctx, ProgressRow{
UserID: a.UserID,
ProfileID: a.ProfileID,
ContentID: sess.LibraryItemID,
CurrentSeconds: sess.CurrentTime,
// For audiobooks MediaItem.Runtime holds total seconds (set by the
// scanner), matching what the ABS libraries handler reads.
DurationSeconds: float64(item.Runtime),
UpdatedAt: time.Now(),
})
} else {
syncErr = h.deps.ProgressStore.UpdateProgressPosition(
ctx, a.UserID, a.ProfileID, sess.LibraryItemID, sess.CurrentTime,
)
}
if syncErr != nil {
slog.WarnContext(ctx, "abs local session sync: persist progress position failed", "component", "audiobooks",
"library_item_id", sess.LibraryItemID, "error", syncErr)
} else {
res.ProgressSynced = true
// Realtime push so other connected clients see the caught-up
// position — same event handleSessionSync emits.
h.publish(a.UserID, "user_item_progress_updated", map[string]any{
"data": map[string]any{
"libraryItemId": sess.LibraryItemID,
"currentTime": sess.CurrentTime,
},
})
}
}
return res
}