Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
51 commits
Select commit Hold shift + click to select a range
826c71f
feat(playback): carry an immutable session creation time in stream to…
CoffeeKnyte Aug 16, 2026
d5d52c4
fix(jellycompat): attribute proxied stream tokens to their owner
CoffeeKnyte Aug 16, 2026
8282533
fix(clientip): resolve viewer addresses on the proxy and ABS listeners
CoffeeKnyte Aug 16, 2026
760287d
fix(httpstream): keep sendfile and the write deadline alive through w…
CoffeeKnyte Aug 16, 2026
a9d54b6
feat(streamtelemetry): add local shadow telemetry for native media ro…
CoffeeKnyte Aug 16, 2026
7a49cca
feat(streamtelemetry): publish snapshots to redis and merge a global …
CoffeeKnyte Aug 17, 2026
d332db9
feat(streamtelemetry): enrol the proxy and transcode-node families
CoffeeKnyte Aug 17, 2026
2908314
feat(streamtelemetry): enrol the jellycompat and ABS families
CoffeeKnyte Aug 17, 2026
b6c9a1c
feat(streamtelemetry): add the P0d admin parity projection
CoffeeKnyte Aug 17, 2026
f6c8b04
docs(streamtelemetry): one working document for what shipped, one app…
CoffeeKnyte Aug 17, 2026
07383f7
feat(playback): add header-authenticated media transport
Quick104 Aug 22, 2026
d762264
feat(playback): negotiate bounded software decode
Quick104 Aug 22, 2026
97b9833
fix(abs): key the login rate limiter on the transport peer
Quick104 Aug 22, 2026
ef565d8
fix(jellycompat): key stream telemetry on the upstream playback session
Quick104 Aug 22, 2026
9bd8c31
fix(streamtelemetry): make Truncated recoverable and hold early realt…
Quick104 Aug 22, 2026
6ec97c1
fix(httpstream): give every ReadFrom slice a full stall window
Quick104 Aug 22, 2026
5a9e6e7
fix(proxy): credit the egress meter often enough to measure slow viewers
Quick104 Aug 22, 2026
ff19ab6
fix(downloads): roll the direct-download deadline and carry the profile
Quick104 Aug 22, 2026
fc46ad8
fix(streamtelemetry): fold ranged transfers, guard delta publishes, s…
Quick104 Aug 22, 2026
813509d
refactor(httpstream): one ForwardReadFrom helper for all nine wrappers
Quick104 Aug 22, 2026
f25ef9b
refactor(streamtelemetry): share viewer-IP, env and client-info helpers
Quick104 Aug 22, 2026
88451ee
docs: document the admin stream-telemetry parity endpoint
Quick104 Aug 22, 2026
a1ba4d6
docs: distill the streaming write-deadline design into architecture
Quick104 Aug 23, 2026
d116d49
Merge branch 'main' into feat/stream-telemetry-enforcer
Quick104 Aug 23, 2026
6b3bca2
chore: ignore skill secrets and state paths
Quick104 Aug 23, 2026
5b7cb63
feat(playback): tokenless V3 playback, DV7 client transforms, admin t…
Quick104 Aug 23, 2026
bbca34f
fix(playback): regenerate conformance matrix for software_video_decod…
Quick104 Aug 23, 2026
8e3f03f
fix(downloads): apply the coarse resolution ceiling when the detailed…
Quick104 Aug 23, 2026
d6ba68b
feat(playback): restore proxy and transcode-node egress for header-au…
Quick104 Aug 23, 2026
5875949
fix(playback): address automated review findings on tokenless proxy e…
Quick104 Aug 23, 2026
9b09ae3
fix(playback): survive transcode-node restarts on tokenless attempts …
Quick104 Aug 23, 2026
413ec6b
test(transcodenode): check CloseProcess error in tokenless reconstruc…
Quick104 Aug 23, 2026
f7ec1e1
Merge pull request #723 from Silo-Server/codex/aether-tokenless-playb…
Quick104 Aug 23, 2026
f12fa1c
Merge branch 'main' into feat/stream-telemetry-enforcer
Quick104 Aug 23, 2026
080c577
fix(streamtelemetry): enrol tokenless /stream/v3 proxy routes
Quick104 Aug 23, 2026
003e9ba
feat(streamtelemetry): enable by default and derive distributed mode …
Quick104 Aug 23, 2026
6301b7e
feat(streamtelemetry): observe every route family by default
Quick104 Aug 23, 2026
cfaab97
Merge pull request #667 from Silo-Server/feat/stream-telemetry-enforcer
Quick104 Aug 23, 2026
5af2a09
feat(scanner): persist H.264 copy-safety verdicts and move analysis o…
Quick104 Aug 23, 2026
f027119
feat(playback): optimistic remux race with server-initiated plan inva…
Quick104 Aug 23, 2026
38ce1a3
fix(playback): sweep sessions that register after a copy-unsafe verdi…
Quick104 Aug 23, 2026
d917f59
fix(playback): harden copy-safety invalidation against review findings
Quick104 Aug 23, 2026
fc09195
fix(playback): validate realtime command ownership before consuming it
Quick104 Aug 23, 2026
cd0c31f
fix(playback): close copy-safety races in replan commits, transport s…
Quick104 Aug 23, 2026
aa0bf35
fix(playback): gate copy-unsafe revivals before reconstruction and pe…
Quick104 Aug 23, 2026
120547b
fix(playback): re-engage the copy-safety race on revival and close ge…
Quick104 Aug 23, 2026
6f8db02
Merge pull request #734 from Silo-Server/feat/optimistic-remux-copy-s…
Quick104 Aug 23, 2026
b6a504d
feat(playback): let original players manage HDR
Quick104 Aug 23, 2026
820eef7
Merge pull request #737 from Silo-Server/codex/aether-managed-original
Quick104 Aug 24, 2026
ed991ad
Merge upstream main into production fork
blurbery Aug 24, 2026
0b98026
chore: satisfy envutil lint
blurbery Aug 24, 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
4 changes: 2 additions & 2 deletions .claude/skills/silo-discord-triage/.gitignore
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
.secrets/
.state/
.secrets
.state
__pycache__/
*.pyc
*.env
4 changes: 2 additions & 2 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -63,8 +63,8 @@ docs/superpowers/
!/.claude/skills/
# Skills may carry per-developer state next to their instructions (the global
# *.local.md rule above covers notes; these cover tokens and run state).
**/.secrets/
**/.state/
**/.secrets
**/.state
.cursor/
.superpowers/
.codex/
Expand Down
27 changes: 25 additions & 2 deletions cmd/playbackfixtures/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -555,6 +555,25 @@ func goldenConformanceMatrix() playback.ConformanceMatrixV3 {
hdr10Request.PlaybackAttemptID = "attempt-hdr10-direct"
planner = append(planner, makePlannerScenario("hdr10_exact_direct", "hdr_dv_matrix", hdr10Request, conformanceHDRFile(), nil, settings, registry))

clientManagedFile := conformanceHDRFile()
clientManagedFile.AudioTracks = append(clientManagedFile.AudioTracks, models.AudioTrack{Codec: codecAAC, Channels: 2, Layout: audioLayoutStereo})
clientManagedRequest := conformanceHDRRequest()
clientManagedRequest.PlaybackAttemptID = "attempt-client-managed-original"
clientManagedRequest.Capabilities.HDR = false
clientManagedRequest.Capabilities.HDRDetails = &playback.HDRCapabilitiesV3{DolbyVisionProfiles: []int{}}
clientManagedRequest.ClientPlaybackContext.Output.HDRDetails = &playback.HDRCapabilitiesV3{DolbyVisionProfiles: []int{}}
clientManagedAudioIndex := 1
clientManagedRequest.AudioTrackIndex = &clientManagedAudioIndex
clientManagedRequest.AudioTrackID = playback.TrackIDV3(clientManagedFile.ID, "audio", clientManagedAudioIndex)
clientManagedDelivery := clientManagedRequest.ClientPlaybackContext.Deliveries[playback.DeliveryClassOriginalHTTPV3]
clientManagedDelivery.HDRDetails = &playback.HDRCapabilitiesV3{DolbyVisionProfiles: []int{}}
clientManagedDelivery.ValidatedClaims = []string{playback.ClaimClientManagedDynamicRangeV3, playback.ClaimClientSelectedAudioTrackV3}
clientManagedRequest.ClientPlaybackContext.Deliveries[playback.DeliveryClassOriginalHTTPV3] = clientManagedDelivery
planner = append(planner, makePlannerScenarioWithAudioIndex(
"client_managed_hdr_selected_audio", "hdr_dv_matrix", clientManagedRequest, clientManagedFile,
clientManagedAudioIndex, nil, settings, registry,
))

dv8File := conformanceHDRFile()
dv8File.VideoTracks[0].DVProfile = 8
dv8File.VideoTracks[0].DVBLCompatID = 1
Expand Down Expand Up @@ -745,8 +764,12 @@ func goldenConformanceMatrix() playback.ConformanceMatrixV3 {
}

func makePlannerScenario(name, category string, request playback.StartRequestV3, file *models.MediaFile, attempted []string, settings playback.PlannerSettingsV3, registry *playback.TransformationRegistryV3) playback.PlannerScenarioV3 {
return makePlannerScenarioWithAudioIndex(name, category, request, file, 0, attempted, settings, registry)
}

func makePlannerScenarioWithAudioIndex(name, category string, request playback.StartRequestV3, file *models.MediaFile, audioTrackIndex int, attempted []string, settings playback.PlannerSettingsV3, registry *playback.TransformationRegistryV3) playback.PlannerScenarioV3 {
result := playback.PlanPlaybackV3(playback.PlannerInputV3{
Request: request, RequestedFile: file, EffectiveFile: file, AudioTrackIndex: 0,
Request: request, RequestedFile: file, EffectiveFile: file, AudioTrackIndex: audioTrackIndex,
Settings: settings, Registry: registry, AttemptedKeys: attempted,
})
expected := playback.PlannerExpectationV3{Outcome: playback.OutcomeAdaptationUnavailableV3}
Expand All @@ -771,7 +794,7 @@ func makePlannerScenario(name, category string, request playback.StartRequestV3,
}
return playback.PlannerScenarioV3{
Name: name, Category: category, Request: request,
Source: playback.SourceDescriptorFromFileV3(file, 0),
Source: playback.SourceDescriptorFromFileV3(file, audioTrackIndex),
AttemptedKeys: append([]string(nil), attempted...), Expected: expected,
}
}
Expand Down
145 changes: 129 additions & 16 deletions cmd/silo/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import (
"github.com/hashicorp/go-hclog"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/redis/go-redis/v9"

pluginv1 "github.com/Silo-Server/silo-plugin-sdk/pkg/pluginproto/silo/plugin/v1"
sdkcapability "github.com/Silo-Server/silo-plugin-sdk/pkg/pluginsdk/capability"
Expand All @@ -40,6 +41,7 @@ import (
"github.com/Silo-Server/silo-server/internal/api"
"github.com/Silo-Server/silo-server/internal/api/handlers"
"github.com/Silo-Server/silo-server/internal/audiobooks"
"github.com/Silo-Server/silo-server/internal/audiobooks/abs"
"github.com/Silo-Server/silo-server/internal/audiobooks/podcastfeed"
"github.com/Silo-Server/silo-server/internal/auth"
"github.com/Silo-Server/silo-server/internal/autoscan"
Expand All @@ -56,6 +58,7 @@ import (
"github.com/Silo-Server/silo-server/internal/ebooks"
evt "github.com/Silo-Server/silo-server/internal/events"
"github.com/Silo-Server/silo-server/internal/historyimport"
"github.com/Silo-Server/silo-server/internal/httpstream"
"github.com/Silo-Server/silo-server/internal/imagecache"
"github.com/Silo-Server/silo-server/internal/intromarkers"
"github.com/Silo-Server/silo-server/internal/jellycompat"
Expand Down Expand Up @@ -97,6 +100,7 @@ import (
"github.com/Silo-Server/silo-server/internal/sections"
"github.com/Silo-Server/silo-server/internal/server"
"github.com/Silo-Server/silo-server/internal/settingscontract"
"github.com/Silo-Server/silo-server/internal/streamtelemetry"
"github.com/Silo-Server/silo-server/internal/subtitles"
"github.com/Silo-Server/silo-server/internal/taskmanager"
taskrepository "github.com/Silo-Server/silo-server/internal/taskmanager/repository"
Expand Down Expand Up @@ -133,6 +137,88 @@ func resolveNodeIdentity() string {
return h
}

func clientIPResolverFromConfig(cfg *config.Config) (*clientip.Resolver, error) {
if cfg == nil {
return nil, fmt.Errorf("config is not loaded")
}
raw := cfg.ClientIP.TrustedProxies
if raw == "" {
raw = clientip.DefaultTrustedProxies
}
cidrs, err := clientip.ParseCIDRs(raw)
if err != nil {
return nil, err
}
return clientip.NewResolver(cidrs), nil
}

func registerClientIPConfigReload(watcher *nodeconfig.Watcher, resolver *clientip.Resolver) {
watcher.OnChange(func(old, updated *config.Config) {
if old != nil && old.ClientIP.TrustedProxies == updated.ClientIP.TrustedProxies {
return
}
raw := updated.ClientIP.TrustedProxies
if raw == "" {
raw = clientip.DefaultTrustedProxies
}
cidrs, err := clientip.ParseCIDRs(raw)
if err != nil {
slog.WarnContext(context.Background(), "clientip config reload failed", "component", "app", "error", err)
return
}
resolver.UpdateTrustedCIDRs(cidrs)
})
}

// newStreamTelemetryRegistry builds the telemetry registry for this process,
// preferring the Redis-backed store in distributed mode. Distributed mode is
// derived from whether Redis is configured unless the operator pinned
// SILO_STREAM_TELEMETRY_DISTRIBUTED: a single-process deployment then stays on
// the local store and a clustered one merges, without either being asked to set
// a variable that only restates its own topology. Every process builds its
// registry here, so the derivation belongs in this function rather than at the
// call sites. It never falls back to LocalStore on a failed ping:
// cache.NewRedisClient builds a lazy client that never dials, and a Redis
// restart mid-deploy must not strand a publisher local-only for the life of the
// process.
func newStreamTelemetryRegistry(ctx context.Context, nodeID string, redisClient *redis.Client) *streamtelemetry.Registry {
streamTelemetryConfig := streamtelemetry.ConfigFromEnv(nodeID)
if !streamTelemetryConfig.DistributedExplicit {
streamTelemetryConfig.Distributed = redisClient != nil
}
store := streamtelemetry.GlobalSnapshotStore(streamtelemetry.NewLocalStore())
if streamTelemetryConfig.Enabled && streamTelemetryConfig.Distributed {
if redisClient != nil {
store = streamtelemetry.NewRedisStore(redisClient, streamTelemetryConfig, slog.Default())
pingCtx, pingCancel := context.WithTimeout(ctx, 2*time.Second)
if pingErr := redisClient.Ping(pingCtx).Err(); pingErr != nil {
slog.ErrorContext(ctx, "stream telemetry distributed mode cannot reach redis; publisher will retry each sweep", "address", redisClient.Options().Addr, "error", pingErr)
}
pingCancel()
} else {
slog.ErrorContext(ctx, "stream telemetry distributed mode requested but redis is not configured; using local store")
}
}
if streamTelemetryConfig.Enabled {
slog.InfoContext(ctx, "stream telemetry observing families",
"families", strings.Join(streamTelemetryConfig.ObservedFamilies(), ","))
}
return streamtelemetry.NewRegistry(streamTelemetryConfig, store, slog.Default())
}

// newStreamTelemetryViewCache builds the bounded-staleness cache the admin
// parity endpoint reads. It shares one cached view across every reader so the
// merged rebuild is paid at most once per TTL, not once per request.
func newStreamTelemetryViewCache(registry *streamtelemetry.Registry) *streamtelemetry.ViewCache {
if registry == nil {
return nil
}
// Reads the TTL off the registry rather than calling ConfigFromEnv again:
// a second parse logs every invalid variable twice and the two calls could
// disagree if the environment changed between them.
return streamtelemetry.NewViewCache(registry, registry.ViewTTL(), slog.Default())
}

func resolvePluginCacheDir() string {
if v := strings.TrimSpace(os.Getenv("SILO_PLUGIN_CACHE_DIR")); v != "" {
return v
Expand Down Expand Up @@ -655,6 +741,8 @@ func main() {

appCtx, appCancel := context.WithCancel(ctx)
defer appCancel()
var streamTelemetryRegistry *streamtelemetry.Registry
var streamTelemetryViewCache *streamtelemetry.ViewCache
restartReqCh := make(chan struct{}, 1)
var restartRequested atomic.Bool

Expand Down Expand Up @@ -683,6 +771,8 @@ func main() {
slog.Error("redis is required for this mode", "mode", mode, "error", err)
os.Exit(1)
}
streamTelemetryRegistry = newStreamTelemetryRegistry(appCtx, nodeID, redisClient)
streamTelemetryRegistry.Start(appCtx)

bootstrap := nodeconfig.BootstrapOverrides{
Listen: cfg.Server.Listen,
Expand Down Expand Up @@ -714,10 +804,29 @@ func main() {
defer cleanupCancel()
tracker.Cleanup(cleanupCtx)
}()
defer func() {
telemetryCtx, telemetryCancel := context.WithTimeout(context.Background(), 5*time.Second)
defer telemetryCancel()
if stopErr := streamTelemetryRegistry.Stop(telemetryCtx); stopErr != nil {
slog.Error("stream telemetry shutdown error", "error", stopErr)
}
}()

var handler http.Handler
if mode == "proxy" {
srv := proxy.NewServer(watcher, tracker)
proxyIPResolver, resolverErr := clientIPResolverFromConfig(watcher.Config())
if resolverErr != nil {
log.Fatalf("load trusted CIDRs: %v", resolverErr)
}
registerClientIPConfigReload(watcher, proxyIPResolver)
srv.SetClientIPResolver(proxyIPResolver)
srv.SetStreamTelemetry(streamTelemetryRegistry)
// Serve header-authenticated sessions: the recipe comes from the
// shared grant store central wrote at plan time, and the caller's
// own access token is re-checked against the live login session in
// Postgres, so a revoked login stops streaming here immediately.
srv.SetMediaGrantAuthority(noderecipe.NewProxyGrantStore(redisClient, 0), auth.NewSessionRepository(pool))
srv.SetRemoteArtifactMissReporter(downloads.NewArtifactManager(
downloads.NewArtifactRepository(pool),
downloads.NewRepository(pool),
Expand All @@ -736,6 +845,7 @@ func main() {
// start, so this node can rebuild a Jellyfin transcode after its own
// restart (the node hop token is recipe-less). Shares the offload Redis.
srv.SetRecipeStore(noderecipe.NewStore(redisClient, 0))
srv.SetStreamTelemetry(streamTelemetryRegistry)
// Reclaim orphaned transcode dirs at boot and hourly thereafter, bound
// to appCtx so it stops on shutdown.
srv.StartOrphanSweeper(appCtx)
Expand Down Expand Up @@ -802,6 +912,12 @@ func main() {
defer func() { _ = apiRedisClient.Close() }()
}

if mode == "" || mode == "integrated" || mode == "api" {
streamTelemetryRegistry = newStreamTelemetryRegistry(appCtx, nodeID, apiRedisClient)
streamTelemetryRegistry.Start(appCtx)
streamTelemetryViewCache = newStreamTelemetryViewCache(streamTelemetryRegistry)
}

// Assigned below once the trusted-proxy config is seeded; captured by the
// OnServerSettingUpdated closure, which only runs on admin requests after
// startup completes.
Expand All @@ -818,6 +934,8 @@ func main() {
BootstrapSensitiveValues: bootstrapSensitiveValues,
RedisBootstrapAvailable: redisBootstrapAvailable,
AppContext: appCtx,
StreamTelemetry: streamTelemetryRegistry,
StreamTelemetryViewCache: streamTelemetryViewCache,
DB: pool,
SecretCipher: dataCipher,
EventBus: eventBus,
Expand Down Expand Up @@ -1849,21 +1967,7 @@ func main() {
})
// The config watcher covers the Redis-less poll/RequestReload path, so
// admin UI edits apply without a restart on single-node deployments too.
configWatcher.OnChange(func(old, updated *config.Config) {
if old != nil && old.ClientIP.TrustedProxies == updated.ClientIP.TrustedProxies {
return
}
raw := updated.ClientIP.TrustedProxies
if raw == "" {
raw = clientip.DefaultTrustedProxies
}
cidrs, parseErr := clientip.ParseCIDRs(raw)
if parseErr != nil {
slog.WarnContext(context.Background(), "clientip config reload failed", "component", "app", "error", parseErr)
return
}
ipResolver.UpdateTrustedCIDRs(cidrs)
})
registerClientIPConfigReload(configWatcher, ipResolver)

// Step 6b: Create rate limiter.
if cfg.RateLimit.Enabled && deps.DB != nil {
Expand Down Expand Up @@ -2319,6 +2423,8 @@ func main() {
SessionSyncer: deps.SessionSyncer,
}
absH := audiobooksService.BuildABSHandler(absHDeps)
// Must precede Mount: Mount is what registers the observed handlers.
absH.SetStreamTelemetry(streamTelemetryRegistry)
deps.ABSHandler = absH
}
_ = audiobooksService
Expand Down Expand Up @@ -2560,6 +2666,7 @@ func main() {
DB: deps.DB,
SecretCipher: dataCipher,
ClientIPResolver: ipResolver,
StreamTelemetry: streamTelemetryRegistry,
NodePlanner: deps.NodePlanner,
JWTSecret: cfg.Auth.JWTSecret,
RecWorker: recWorker,
Expand Down Expand Up @@ -2705,8 +2812,11 @@ func main() {
var absSrv *http.Server
if (mode == "integrated" || mode == "api") && deps.ABSHandler != nil && cfg.AudiobookshelfCompat.Listen != "" {
absRouter := chi.NewRouter()
if ipResolver != nil {
absRouter.Use(clientip.Middleware(ipResolver))
}
absRouter.Use(chimiddleware.Recoverer)
absRouter.Use(chimiddleware.Compress(5))
absRouter.Use(httpstream.CompressExcept(5, abs.SkipMediaCompression))
deps.ABSHandler.Mount(absRouter)
absSrv = &http.Server{
Addr: cfg.AudiobookshelfCompat.Listen,
Expand Down Expand Up @@ -2802,6 +2912,9 @@ func main() {
slog.Error("abs compat shutdown error", "error", shutdownErr)
}
}
if stopErr := streamTelemetryRegistry.Stop(shutdownCtx); stopErr != nil {
slog.Error("stream telemetry shutdown error", "error", stopErr)
}

// 2. Clean up stale sessions.
if sessionCleaner != nil {
Expand Down
1 change: 1 addition & 0 deletions cmd/silo/session_sync.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ func buildLiveSessionSync(s *playback.Session, reportingNode string) worker.Sess
TargetResolution: s.TargetResolution,
TargetVideoCodec: s.TargetVideoCodec,
TargetAudioCodec: s.TargetAudioCodec,
TargetAudioChannels: s.TargetAudioChannels,
TargetBitrateKbps: s.TargetBitrateKbps,
TranscodeHWAccel: s.TranscodeHWAccel,
StartedAt: s.StartedAt,
Expand Down
Loading
Loading