Files
silo-server/internal/api/middleware/request_logger.go
T
CoffeeKnyte 0217cce2df feat(playback): stream kill switch + async over-cap enforcer
Add the enforcement layer on top of server-observed monitoring: a revocation
kill switch that stops any stream within ~120s and keeps it dead, plus an async
over-cap enforcer that drives kills off the live monitoring picture — entirely
off the per-segment hot path and with no client-protocol change.

- internal/streamrevoke: the central kill list. IsRevoked is a pure in-memory
  lookup safe on the request hot path; a Redis pub/sub + poll mirror keeps edge
  caches current, and a Postgres durable mirror lets kills survive a server
  restart AND a Redis flush so a restart-resilient stream cannot be reconstructed
  and re-served after being killed. A user revocation is a cutoff (kills tokens
  minted before it, spares post-reauth tokens), not a 24h ban.
- internal/streamenforcer: async over-cap brain — reads the monitoring snapshot
  and per-user limits, selects victims, and collapses every reason (exceeded
  limit, admin terminate, abuse) to the same action: write a revocation.
- Edge + native + jellycompat enforcement: proxy refuses revoked sessions on
  every request and cuts long direct-play/remux pours mid-stream; the transcode
  node guards both serve and the reconstruct path so a killed session is never
  re-spawned after a node restart; jellycompat serve surfaces close their
  kill-switch coverage holes.
- streamtoken.IssuedTime exposes the token iat the user-kill cutoff compares
  against; token IssuedTime + revocation guards wire through router, downloads,
  and admin terminate-by-id (with admin-list dedupe).
- Restore sendfile zero-copy on direct-play/remux byte counting so the monitor's
  served-byte accounting does not cost the sendfile fast path.
- migrations/sql: stream_revocations durable table.

Part of the stream monitoring & kill-switch epic.
2026-07-29 12:26:03 +00:00

123 lines
3.4 KiB
Go

package middleware
import (
"bufio"
"fmt"
"log/slog"
"net"
"net/http"
"strings"
"time"
"github.com/go-chi/chi/v5"
chimw "github.com/go-chi/chi/v5/middleware"
"github.com/Silo-Server/silo-server/internal/activitylog"
"github.com/Silo-Server/silo-server/internal/clientip"
)
func RequestLogger(nodeID string) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
for _, prefix := range []string{"/api/v1/health", "/api/v1/ready", "/api/v1/admin/logs"} {
if strings.HasPrefix(r.URL.Path, prefix) {
next.ServeHTTP(w, r)
return
}
}
start := time.Now()
wrapped := &requestStatusWriter{ResponseWriter: w, status: http.StatusOK}
lc := activitylog.GetLogContext(r.Context())
if lc == nil {
lc = &activitylog.LogContext{}
r = r.WithContext(activitylog.SetLogContext(r.Context(), lc))
}
playbackLC := activitylog.GetPlaybackLogContext(r.Context())
if playbackLC == nil {
playbackLC = &activitylog.PlaybackLogContext{}
r = r.WithContext(activitylog.SetPlaybackLogContext(r.Context(), playbackLC))
}
next.ServeHTTP(wrapped, r)
pathPattern := r.URL.Path
if routeCtx := chi.RouteContext(r.Context()); routeCtx != nil {
if route := routeCtx.RoutePattern(); route != "" {
pathPattern = route
}
}
attrs := []any{
"component", "api",
"request_id", chimw.GetReqID(r.Context()),
"method", r.Method,
"path", activitylog.RedactSecretPathParams(r, r.URL.Path),
"path_pattern", pathPattern,
"status", wrapped.status,
"duration_ms", time.Since(start).Milliseconds(),
"client_ip", clientip.FromContext(r.Context()),
"user_agent", r.UserAgent(),
"node_id", nodeID,
}
if lc.UserID != nil {
attrs = append(attrs, "user_id", *lc.UserID)
}
if lc.SessionID != "" {
attrs = append(attrs, "session_id", lc.SessionID)
}
if playbackLC.PlaybackSessionID != "" {
attrs = append(attrs, "playback_session_id", playbackLC.PlaybackSessionID)
}
slog.InfoContext(r.Context(), "api request", append([]any{"component", "api"}, attrs...)...)
})
}
}
type requestStatusWriter struct {
http.ResponseWriter
status int
wroteHeader bool
}
func (w *requestStatusWriter) WriteHeader(status int) {
if !w.wroteHeader {
w.status = status
w.wroteHeader = true
}
w.ResponseWriter.WriteHeader(status)
}
func (w *requestStatusWriter) Write(b []byte) (int, error) {
if !w.wroteHeader {
w.wroteHeader = true
}
return w.ResponseWriter.Write(b)
}
func (w *requestStatusWriter) Hijack() (net.Conn, *bufio.ReadWriter, error) {
if hj, ok := w.ResponseWriter.(http.Hijacker); ok {
return hj.Hijack()
}
return nil, nil, fmt.Errorf("underlying ResponseWriter does not implement http.Hijacker")
}
// Flush implements http.Flusher so progressive responses (streamed subtitle
// extracts, remux output) keep flushing through the logging wrapper instead of
// silently buffering until the handler returns.
func (w *requestStatusWriter) Flush() {
if !w.wroteHeader {
w.WriteHeader(http.StatusOK)
}
if f, ok := w.ResponseWriter.(http.Flusher); ok {
f.Flush()
}
}
// Unwrap exposes the wrapped ResponseWriter so http.ResponseController (used by
// the stream kill switch's in-flight cut via SetWriteDeadline) can reach the
// underlying socket instead of stopping at this wrapper and no-oping.
func (w *requestStatusWriter) Unwrap() http.ResponseWriter {
return w.ResponseWriter
}