// Package api exposes the gateway's HTTP surface. // // The API is shaped for one TV screen at a time rather than mirroring Emby: /v1/home // returns everything the launcher renders in a single round trip, which is the whole // point of putting a gateway in front of Emby. package api import ( "context" "crypto/rand" "crypto/sha256" "crypto/subtle" "encoding/base64" "encoding/hex" "encoding/json" "errors" "log/slog" "net/http" "slices" "strconv" "strings" "sync" "time" "github.com/ponzischeme89/memby/server/internal/cache" "github.com/ponzischeme89/memby/server/internal/config" "github.com/ponzischeme89/memby/server/internal/emby" "github.com/ponzischeme89/memby/server/internal/foryou" serverlogging "github.com/ponzischeme89/memby/server/internal/logging" "github.com/ponzischeme89/memby/server/internal/mdblist" "github.com/ponzischeme89/memby/server/internal/radarr" "github.com/ponzischeme89/memby/server/internal/recommend" "github.com/ponzischeme89/memby/server/internal/sonarr" "github.com/ponzischeme89/memby/server/internal/store" ) type Server struct { cfg config.Config emby *emby.Client store *store.Store cache *cache.Cache recommender *recommend.Engine forYou *foryou.Service sonarr *sonarr.Client radarr *radarr.Client mdblist *mdblist.Client syncer syncerHandle log *slog.Logger events *serverlogging.Buffer sonarrMu sync.Mutex radarrMu sync.Mutex mdblistMu sync.Mutex // alertMu serialises the read-modify-write of the shared alert list. Its producers // are events — a webhook, a finished sync, a health probe — none of them paced by // this server, so two can land at once. alertMu sync.Mutex recommendationBuilds recommendationBuilds maintenance maintenanceState updatePolicy updatePolicyCache } // Deps are the collaborators the API needs. A struct rather than positional arguments: // this list has grown three times already. type Deps struct { Emby *emby.Client Store *store.Store Cache *cache.Cache Recommender *recommend.Engine ForYou *foryou.Service Sonarr *sonarr.Client Radarr *radarr.Client MDBList *mdblist.Client Syncer syncerHandle Log *slog.Logger Events *serverlogging.Buffer } func New(cfg config.Config, deps Deps) *Server { return &Server{ cfg: cfg, emby: deps.Emby, store: deps.Store, cache: deps.Cache, recommender: deps.Recommender, forYou: deps.ForYou, sonarr: deps.Sonarr, radarr: deps.Radarr, mdblist: deps.MDBList, syncer: deps.Syncer, log: deps.Log, events: deps.Events, } } func (s *Server) Routes() http.Handler { // The client API lives on its own mux so maintenance mode can gate all of it at // once, without the gate ever touching health checks or the admin page. v1 := http.NewServeMux() v1.HandleFunc("POST /v1/auth/login", s.handleLogin) v1.Handle("POST /v1/auth/logout", s.authed(s.handleLogout)) v1.Handle("GET /v1/auth/session", s.authed(s.handleSession)) v1.Handle("GET /v1/auth/devices", s.authed(s.handleDevices)) v1.Handle("PUT /v1/auth/devices/{deviceID}", s.authed(s.handleRenameDevice)) v1.Handle("DELETE /v1/auth/devices/{deviceID}", s.authed(s.handleDeleteDevice)) v1.Handle("GET /v1/home", s.authed(s.handleHome)) v1.Handle("GET /v1/screensaver", s.authed(s.handleScreensaver)) v1.Handle("GET /v1/search", s.authed(s.handleSearch)) v1.Handle("GET /v1/search/history", s.authed(s.handleRecentSearches)) v1.Handle("POST /v1/search/history", s.authed(s.handleSearchHistory)) v1.Handle("GET /v1/requests/lookup", s.authed(s.handleRequestLookup)) v1.Handle("POST /v1/requests", s.authed(s.handleRequest)) v1.Handle("GET /v1/recommendations", s.authed(s.handleRecommendations)) v1.Handle("PUT /v1/recommendations/{id}/action", s.authed(s.handleRecommendationAction)) v1.Handle("DELETE /v1/recommendations/{id}/action", s.authed(s.handleRecommendationAction)) v1.Handle("PUT /v1/recommendations/preferences", s.authed(s.handleRecommendationPreferences)) v1.Handle("GET /v1/recommendations/preferences", s.authed(s.handleRecommendationPreferences)) v1.Handle("GET /v1/for-you", s.authed(s.handleForYou)) v1.Handle("GET /v1/preroll", s.authed(s.handlePreroll)) v1.Handle("GET /v1/my-shows", s.authed(s.handleMyShows)) v1.Handle("POST /v1/my-shows", s.authed(s.handleMyShows)) v1.Handle("DELETE /v1/my-shows/{id}", s.authed(s.handleMyShow)) v1.Handle("GET /v1/notifications", s.authed(s.handleNotifications)) v1.Handle("PUT /v1/notifications", s.authed(s.handleNotifications)) v1.Handle("POST /v1/notifications/{id}/{action}", s.authed(s.handleNotificationAction)) v1.Handle("GET /v1/features", s.authed(s.handleFeatures)) v1.Handle("GET /v1/items/{id}", s.authed(s.handleItem)) v1.Handle("GET /v1/items/{id}/ratings", s.authed(s.handleMovieRatings)) v1.Handle("GET /v1/items/{id}/season-finale", s.authed(s.handleSeasonFinale)) v1.Handle("GET /v1/items/{id}/episodes", s.authed(s.handleSeriesEpisodes)) v1.Handle("GET /v1/items/{id}/related", s.authed(s.handleRelated)) v1.Handle("POST /v1/items/{id}/favorite", s.authed(s.handleFavorite)) v1.Handle("POST /v1/items/{id}/played", s.authed(s.handlePlayed)) v1.Handle("GET /v1/items/{id}/playback", s.authed(s.handlePlayback)) v1.Handle("GET /v1/items/{id}/next", s.authed(s.handleNextEpisode)) v1.Handle("GET /v1/items/{id}/trailer", s.authed(s.handleTrailer)) v1.Handle("POST /v1/playback/{phase}", s.authed(s.handlePlaybackReport)) v1.Handle("POST /v1/analytics/rows", s.authed(s.handleRowAnalytics)) v1.Handle("GET /v1/images/{itemId}/{imageType}", s.authed(s.handleImage)) mux := http.NewServeMux() mux.HandleFunc("GET /healthz", s.handleHealth) mux.HandleFunc("GET /readyz", s.handleReady) // Update policy is app-scoped, not user-scoped. Keep it outside authentication and // maintenance so a fresh install, a signed-out TV, and a retired build can all learn // whether the server requires an update without touching a viewer session. mux.HandleFunc("GET /v1/update", s.handleUpdate) // Exact route outside the maintenance gate: signed-in clients poll this lightweight // status even while every normal /v1 operation is deliberately unavailable. mux.Handle("GET /v1/status", s.authed(s.handleServiceStatus)) mux.Handle("/v1/", s.maintenanceGate(v1)) // Radarr pushes here when an import finishes. Outside the gate on purpose: an event // arriving during maintenance would otherwise be lost rather than delayed. mux.HandleFunc("POST /hooks/radarr", s.handleRadarrWebhook) mux.Handle("/admin/", s.adminRoutes()) mux.HandleFunc("GET /{$}", s.handleInstallPage) mux.HandleFunc("GET /install", s.handleInstallPage) mux.HandleFunc("GET /install/{$}", s.handleInstallPage) mux.HandleFunc("POST /install/login", s.handleInstallLogin) mux.HandleFunc("POST /install/logout", s.handleInstallLogout) mux.HandleFunc("GET /robots.txt", handleRobots) mux.HandleFunc("GET /updates/latest.apk", s.handleLatestReleaseDownload) mux.HandleFunc("GET /updates/{filename}", s.handleReleaseDownload) return s.withLogging(mux) } // --- middleware ------------------------------------------------------------- type authedFunc func(http.ResponseWriter, *http.Request, store.Session) // authed resolves the bearer token to a session before running h. // // Images are also accepted with a `t=` query parameter: Coil builds plain URLs from the // repository's helpers and cannot attach headers to them. func (s *Server) authed(h authedFunc) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { token := bearerToken(r) if token == "" { writeError(w, http.StatusUnauthorized, "missing token") return } sess, err := s.sessionFor(r.Context(), token) if err != nil { if errors.Is(err, store.ErrNotFound) { writeError(w, http.StatusUnauthorized, "invalid token") return } s.log.Error("session lookup failed", "error", err) writeError(w, http.StatusInternalServerError, "session lookup failed") return } sess = s.captureClientIdentity(r, sess) h(w, r, sess) }) } // captureClientIdentity makes the session the durable source of attribution. Normal API // calls refresh it from headers; authenticated artwork requests, which can only carry a // query token, inherit the last identity reported by that same TV. func (s *Server) captureClientIdentity(r *http.Request, sess store.Session) store.Session { changed := mergeClientIdentity(r, &sess) if changed { if err := s.store.UpdateSessionClientIdentity( r.Context(), sess.TokenHash, sess.ClientVersion, sess.ClientProtocol, sess.ClientCapabilities, ); err != nil { s.log.Warn("client identity update failed", "error", err) } else { s.cacheSession(r.Context(), sess) } } return sess } func mergeClientIdentity(r *http.Request, sess *store.Session) bool { version := clientVersion(r) protocol := clientProtocol(r) capabilities := clientCapabilities(r) changed := false if version != "" && version != sess.ClientVersion { sess.ClientVersion = version changed = true } if protocol != "" && protocol != sess.ClientProtocol { sess.ClientProtocol = protocol changed = true } if len(capabilities) > 0 && !slices.Equal(capabilities, sess.ClientCapabilities) { sess.ClientCapabilities = capabilities changed = true } if version == "" && sess.ClientVersion != "" { r.Header.Set("X-Memby-Version", sess.ClientVersion) } if protocol == "" && sess.ClientProtocol != "" { r.Header.Set("X-Memby-Protocol", sess.ClientProtocol) } return changed } func (s *Server) withLogging(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { start := time.Now() rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK} next.ServeHTTP(rec, r) // Polling the live-log endpoint must not create another live-log record and // become a self-sustaining stream. if r.URL.Path == "/admin/api/events" { return } // Path only: query strings can carry image tokens. level := requestLogLevel(r.URL.Path, rec.status) s.log.Log(r.Context(), level, "HTTP request", "method", r.Method, "path", r.URL.Path, "status", rec.status, "duration", time.Since(start).Round(time.Millisecond), "client_version", clientLogValue(clientVersion(r)), "client_protocol", clientLogValue(clientProtocol(r)), ) }) } func clientLogValue(value string) string { if value == "" { return "unknown" } return value } // Successful high-frequency probes and artwork fetches stay available at DEBUG without // overwhelming the normal Docker log. Failures are always promoted so they remain // visible regardless of path. func requestLogLevel(path string, status int) slog.Level { switch { case status >= http.StatusInternalServerError: return slog.LevelError case status >= http.StatusBadRequest: return slog.LevelWarn case path == "/healthz", path == "/readyz", path == "/v1/status", strings.HasPrefix(path, "/v1/images/"): return slog.LevelDebug default: return slog.LevelInfo } } type statusRecorder struct { http.ResponseWriter status int } func (r *statusRecorder) WriteHeader(code int) { r.status = code r.ResponseWriter.WriteHeader(code) } // --- sessions --------------------------------------------------------------- func bearerToken(r *http.Request) string { if h := r.Header.Get("Authorization"); strings.HasPrefix(h, "Bearer ") { return strings.TrimSpace(strings.TrimPrefix(h, "Bearer ")) } if h := r.Header.Get("X-Memby-Token"); h != "" { return strings.TrimSpace(h) } return strings.TrimSpace(r.URL.Query().Get("t")) } func hashToken(token string) []byte { sum := sha256.Sum256([]byte(token)) return sum[:] } func newToken() (string, error) { buf := make([]byte, 32) if _, err := rand.Read(buf); err != nil { return "", err } return base64.RawURLEncoding.EncodeToString(buf), nil } type cachedSession struct { EmbyUserID string `json:"u"` EmbyToken string `json:"t"` Username string `json:"n"` ServerID string `json:"s"` DeviceID string `json:"d"` DeviceName string `json:"dn,omitempty"` ClientVersion string `json:"v,omitempty"` ClientProtocol string `json:"p,omitempty"` } // sessionFor resolves a token, using Redis to keep the hot path off Postgres. func (s *Server) sessionFor(ctx context.Context, token string) (store.Session, error) { hash := hashToken(token) key := cache.SessionKey(hex.EncodeToString(hash)) if raw, err := s.cache.Get(ctx, key); err == nil { var cs cachedSession if json.Unmarshal(raw, &cs) == nil { return store.Session{ TokenHash: hash, EmbyUserID: cs.EmbyUserID, EmbyToken: cs.EmbyToken, Username: cs.Username, ServerID: cs.ServerID, DeviceID: cs.DeviceID, DeviceName: cs.DeviceName, ClientVersion: cs.ClientVersion, ClientProtocol: cs.ClientProtocol, }, nil } } sess, err := s.store.SessionByTokenHash(ctx, hash) if err != nil { return store.Session{}, err } // Constant-time confirmation that the stored hash matches the presented token. if subtle.ConstantTimeCompare(sess.TokenHash, hash) != 1 { return store.Session{}, store.ErrNotFound } s.cacheSession(ctx, sess) // Best-effort activity stamp; a failure here must not fail the request. if err := s.store.Touch(ctx, hash); err != nil { s.log.Warn("touch session failed", "error", err) } return sess, nil } func (s *Server) cacheSession(ctx context.Context, sess store.Session) { if raw, err := json.Marshal(cachedSession{ EmbyUserID: sess.EmbyUserID, EmbyToken: sess.EmbyToken, Username: sess.Username, ServerID: sess.ServerID, DeviceID: sess.DeviceID, DeviceName: sess.DeviceName, ClientVersion: sess.ClientVersion, ClientProtocol: sess.ClientProtocol, }); err == nil { _ = s.cache.Set( ctx, cache.SessionKey(hex.EncodeToString(sess.TokenHash)), raw, s.cfg.SessionTTL, ) } } func credentials(sess store.Session) emby.Credentials { return emby.Credentials{ UserID: sess.EmbyUserID, Token: sess.EmbyToken, DeviceID: sess.DeviceID, DeviceName: sess.DeviceName, } } // --- responses -------------------------------------------------------------- func writeJSON(w http.ResponseWriter, status int, body any) { w.Header().Set("Content-Type", "application/json; charset=utf-8") w.WriteHeader(status) if err := json.NewEncoder(w).Encode(body); err != nil { // Headers are already out; nothing useful left to do but stop. return } } func writeRaw(w http.ResponseWriter, status int, body []byte) { w.Header().Set("Content-Type", "application/json; charset=utf-8") w.WriteHeader(status) _, _ = w.Write(body) } func writeError(w http.ResponseWriter, status int, message string) { writeJSON(w, status, map[string]string{"error": message}) } // writeUpstreamError mirrors Emby's status so the TV can tell "signed out" (401) from // "server is unwell" (5xx) without parsing strings. func (s *Server) writeUpstreamError(w http.ResponseWriter, err error, message string) { var apiErr *emby.APIError if errors.As(err, &apiErr) { switch { case apiErr.StatusCode == http.StatusUnauthorized, apiErr.StatusCode == http.StatusForbidden: writeError(w, http.StatusUnauthorized, "emby rejected the session") return case apiErr.StatusCode == http.StatusNotFound: writeError(w, http.StatusNotFound, "not found on the emby server") return } } s.log.Error(message, "error", err) writeError(w, http.StatusBadGateway, message) } func queryInt(r *http.Request, key string, fallback, max int) int { raw := r.URL.Query().Get(key) if raw == "" { return fallback } v, err := strconv.Atoi(raw) if err != nil || v <= 0 { return fallback } if v > max { return max } return v }