Garmin's rate-limiting means every unattended sync attempt risks a ban; syncing should only ever happen when explicitly triggered via "Sync now" (POST /api/sync/run), never on an unattended timer. Removes the periodic background loop (main.go's runIncrementalSyncLoop, api.Server.RunIncrementalSyncForAllUsers), its GENIUSRUN_INCREMENTAL_SYNC_EVERY config, and store.DB.ListUsers (which existed solely to feed it). The manual "Sync now" flow (Backfill/IncrementalSync/FillPendingDetails via FullSync) is untouched.
433 lines
14 KiB
Go
433 lines
14 KiB
Go
// Package api is geniusrun's HTTP layer: REST handlers over internal/store,
|
|
// internal/garmin, and internal/sync.
|
|
package api
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"path/filepath"
|
|
"strconv"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/go-chi/chi/v5"
|
|
"github.com/go-chi/chi/v5/middleware"
|
|
|
|
"geniusrun/backend/internal/applog"
|
|
"geniusrun/backend/internal/auth"
|
|
"geniusrun/backend/internal/garmin"
|
|
"geniusrun/backend/internal/store"
|
|
appsync "geniusrun/backend/internal/sync"
|
|
)
|
|
|
|
// Server wires the HTTP handlers to the app's dependencies. garmin.Client
|
|
// and sync.Service are per-user (each user might have their own Garmin
|
|
// account), built lazily on first use via GarminFactory and cached.
|
|
type Server struct {
|
|
DB *store.DB
|
|
Auth auth.Verifier
|
|
Session SessionConfig
|
|
|
|
// GarminFactory builds a real (or fake, in tests) garmin.Client from a
|
|
// fully-resolved per-user Config. Production wiring passes
|
|
// garmin.NewClient; tests inject a factory returning a shared
|
|
// *mock.Client (see newTestServer in api_test.go).
|
|
GarminFactory func(garmin.Config) garmin.Client
|
|
// GarminBase holds the plumbing shared by every user's garmin.Config
|
|
// (the python interpreter path + the token-store root directory); only
|
|
// GarminEmail/GarminPassword/TokenStorePath vary per user, filled in by
|
|
// garminFor.
|
|
GarminBase garmin.Config
|
|
SyncConfig appsync.Config
|
|
|
|
mu sync.Mutex
|
|
userGarmin map[int64]garmin.Client
|
|
userSync map[int64]*appsync.Service
|
|
userAuthStatus map[int64]garmin.AuthStatus
|
|
userAuthMessage map[int64]string
|
|
userSyncRunning map[int64]bool
|
|
setupGarmin map[string]*setupSession
|
|
}
|
|
|
|
// NewServer builds a Server.
|
|
func NewServer(db *store.DB, garminFactory func(garmin.Config) garmin.Client, garminBase garmin.Config, syncConfig appsync.Config, authVerifier auth.Verifier, session SessionConfig) *Server {
|
|
return &Server{
|
|
DB: db, GarminFactory: garminFactory, GarminBase: garminBase, SyncConfig: syncConfig,
|
|
Auth: authVerifier, Session: session,
|
|
userGarmin: map[int64]garmin.Client{},
|
|
userSync: map[int64]*appsync.Service{},
|
|
userAuthStatus: map[int64]garmin.AuthStatus{},
|
|
userAuthMessage: map[int64]string{},
|
|
userSyncRunning: map[int64]bool{},
|
|
setupGarmin: map[string]*setupSession{},
|
|
}
|
|
}
|
|
|
|
// setupSession is a temporary, not-yet-persisted Garmin authentication
|
|
// attempt made during onboarding, before any users/profile row exists --
|
|
// keyed by OIDC subject (the only stable identifier available pre-account)
|
|
// rather than a user id. Promoted into Server.userGarmin once
|
|
// /api/setup/complete actually creates the account; evicted lazily (the
|
|
// next setup-endpoint touch for that subject checks staleness first) once
|
|
// idle past setupSessionIdleTimeout.
|
|
type setupSession struct {
|
|
Client garmin.Client
|
|
Email, Password string
|
|
Status garmin.AuthStatus
|
|
Message string
|
|
LastUsed time.Time
|
|
}
|
|
|
|
// setupSessionIdleTimeout bounds how long an onboarding Garmin session
|
|
// survives without being touched (login, MFA, or complete) before it's
|
|
// evicted -- long enough to check email for an MFA code, short enough that
|
|
// an abandoned attempt doesn't leave a subprocess running indefinitely.
|
|
const setupSessionIdleTimeout = 15 * time.Minute
|
|
|
|
// setupTokenStoreDir returns the token-store directory an ephemeral setup
|
|
// session for sub should use -- hashed rather than sub itself, since sub is
|
|
// an opaque string from the identity provider and using it verbatim in a
|
|
// filesystem path would be a directory-traversal risk if it ever contained
|
|
// path separators.
|
|
func setupTokenStoreDir(root, sub string) string {
|
|
h := sha256.Sum256([]byte(sub))
|
|
return filepath.Join(root, "setup", hex.EncodeToString(h[:]))
|
|
}
|
|
|
|
// setupSessionFor returns sub's in-progress ephemeral Garmin session, if
|
|
// any and not stale. A stale session is closed and evicted first, so the
|
|
// caller always either gets a fresh, live session or none.
|
|
func (s *Server) setupSessionFor(sub string) (*setupSession, bool) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
sess, ok := s.setupGarmin[sub]
|
|
if !ok {
|
|
return nil, false
|
|
}
|
|
if time.Since(sess.LastUsed) > setupSessionIdleTimeout {
|
|
sess.Client.Close()
|
|
delete(s.setupGarmin, sub)
|
|
if s.GarminBase.TokenStorePath != "" {
|
|
os.RemoveAll(setupTokenStoreDir(s.GarminBase.TokenStorePath, sub))
|
|
}
|
|
return nil, false
|
|
}
|
|
return sess, true
|
|
}
|
|
|
|
// replaceSetupSession closes and replaces sub's ephemeral Garmin session
|
|
// (if any) with a freshly built one for the given credentials -- same
|
|
// "close old, spawn new" semantics as garmin.Client.UpdateCredentials.
|
|
func (s *Server) replaceSetupSession(sub, email, password string) *setupSession {
|
|
s.mu.Lock()
|
|
old, ok := s.setupGarmin[sub]
|
|
s.mu.Unlock()
|
|
if ok {
|
|
old.Client.Close()
|
|
}
|
|
|
|
cfg := s.GarminBase
|
|
cfg.GarminEmail = email
|
|
cfg.GarminPassword = password
|
|
if cfg.TokenStorePath != "" {
|
|
cfg.TokenStorePath = setupTokenStoreDir(s.GarminBase.TokenStorePath, sub)
|
|
}
|
|
sess := &setupSession{Client: s.GarminFactory(cfg), Email: email, Password: password, LastUsed: time.Now()}
|
|
|
|
s.mu.Lock()
|
|
s.setupGarmin[sub] = sess
|
|
s.mu.Unlock()
|
|
return sess
|
|
}
|
|
|
|
// recordSetupAuthResult updates sub's ephemeral session after a login or
|
|
// MFA attempt. A no-op if the session is gone (e.g. evicted concurrently).
|
|
func (s *Server) recordSetupAuthResult(sub string, res garmin.AuthResult) {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if sess, ok := s.setupGarmin[sub]; ok {
|
|
sess.Status = res.Status
|
|
sess.Message = res.Message
|
|
sess.LastUsed = time.Now()
|
|
}
|
|
}
|
|
|
|
// removeSetupSession closes and drops sub's ephemeral Garmin session, if
|
|
// any, and best-effort removes its token-store directory. Used when
|
|
// abandoning onboarding (logout) -- NOT used during promotion in
|
|
// handleSetupComplete, which transfers ownership of the client (and
|
|
// renames the directory) instead of discarding them.
|
|
func (s *Server) removeSetupSession(sub string) {
|
|
s.mu.Lock()
|
|
sess, ok := s.setupGarmin[sub]
|
|
delete(s.setupGarmin, sub)
|
|
s.mu.Unlock()
|
|
if !ok {
|
|
return
|
|
}
|
|
sess.Client.Close()
|
|
if s.GarminBase.TokenStorePath != "" {
|
|
os.RemoveAll(setupTokenStoreDir(s.GarminBase.TokenStorePath, sub))
|
|
}
|
|
}
|
|
|
|
// garminFor returns userID's garmin.Client, building and caching it (from
|
|
// userID's own profile row) on first use.
|
|
func (s *Server) garminFor(ctx context.Context, userID int64) (garmin.Client, error) {
|
|
s.mu.Lock()
|
|
if c, ok := s.userGarmin[userID]; ok {
|
|
s.mu.Unlock()
|
|
return c, nil
|
|
}
|
|
s.mu.Unlock()
|
|
|
|
profile, err := s.DB.GetProfile(ctx, userID)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("load profile for garmin client (user %d): %w", userID, err)
|
|
}
|
|
cfg := s.GarminBase
|
|
cfg.GarminEmail = profile.GarminEmail
|
|
cfg.GarminPassword = profile.GarminPassword
|
|
if cfg.TokenStorePath != "" {
|
|
cfg.TokenStorePath = filepath.Join(cfg.TokenStorePath, strconv.FormatInt(userID, 10))
|
|
}
|
|
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if c, ok := s.userGarmin[userID]; ok {
|
|
return c, nil // built concurrently by another request between our unlock and re-lock
|
|
}
|
|
client := s.GarminFactory(cfg)
|
|
s.userGarmin[userID] = client
|
|
return client, nil
|
|
}
|
|
|
|
// syncFor returns userID's sync.Service, building and caching it on first use.
|
|
func (s *Server) syncFor(ctx context.Context, userID int64) (*appsync.Service, error) {
|
|
s.mu.Lock()
|
|
if svc, ok := s.userSync[userID]; ok {
|
|
s.mu.Unlock()
|
|
return svc, nil
|
|
}
|
|
s.mu.Unlock()
|
|
|
|
client, err := s.garminFor(ctx, userID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if svc, ok := s.userSync[userID]; ok {
|
|
return svc, nil
|
|
}
|
|
svc := appsync.NewService(client, s.DB, userID, s.SyncConfig, nil)
|
|
s.userSync[userID] = svc
|
|
return svc, nil
|
|
}
|
|
|
|
// removeGarminClient drops userID's cached garmin.Client/sync.Service (if
|
|
// any) and every other per-user in-memory entry for userID, terminating the
|
|
// client's subprocess and best-effort removing its on-disk token-store
|
|
// directory. Called when a user's account has just been deleted from the
|
|
// DB, so nothing in memory keeps referencing a userID that no longer
|
|
// exists.
|
|
func (s *Server) removeGarminClient(userID int64) {
|
|
s.mu.Lock()
|
|
client, ok := s.userGarmin[userID]
|
|
delete(s.userGarmin, userID)
|
|
delete(s.userSync, userID)
|
|
delete(s.userAuthStatus, userID)
|
|
delete(s.userAuthMessage, userID)
|
|
delete(s.userSyncRunning, userID)
|
|
s.mu.Unlock()
|
|
|
|
if ok {
|
|
if err := client.Close(); err != nil {
|
|
log.Printf("api: close garmin client for deleted user %d: %v", userID, err)
|
|
}
|
|
}
|
|
if s.GarminBase.TokenStorePath == "" {
|
|
return
|
|
}
|
|
tokenStoreDir := filepath.Join(s.GarminBase.TokenStorePath, strconv.FormatInt(userID, 10))
|
|
if err := os.RemoveAll(tokenStoreDir); err != nil {
|
|
log.Printf("api: remove token store dir for deleted user %d: %v", userID, err)
|
|
}
|
|
}
|
|
|
|
var requestIDCounter atomic.Int64
|
|
|
|
// requestLoggingMiddleware logs one JSON line per HTTP request (method,
|
|
// path, status, duration) and attaches a per-request logger (tagged with a
|
|
// request_id) to the request context, so any downstream call this request
|
|
// triggers -- e.g. a Garmin wrapper round-trip -- logs with the same
|
|
// correlating id (see internal/applog, internal/garmin's roundTrip).
|
|
func requestLoggingMiddleware(next http.Handler) http.Handler {
|
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
id := fmt.Sprintf("req-%d", requestIDCounter.Add(1))
|
|
logger := applog.FromContext(r.Context()).With("request_id", id)
|
|
r = r.WithContext(applog.WithLogger(r.Context(), logger))
|
|
|
|
ww := middleware.NewWrapResponseWriter(w, r.ProtoMajor)
|
|
start := time.Now()
|
|
next.ServeHTTP(ww, r)
|
|
|
|
level := slog.LevelInfo
|
|
if ww.Status() >= 500 {
|
|
level = slog.LevelWarn
|
|
}
|
|
logger.LogAttrs(r.Context(), level, "http request",
|
|
slog.String("method", r.Method),
|
|
slog.String("path", r.URL.Path),
|
|
slog.Int("status", ww.Status()),
|
|
slog.Int64("duration_ms", time.Since(start).Milliseconds()),
|
|
)
|
|
})
|
|
}
|
|
|
|
// Router builds the HTTP routes.
|
|
func (s *Server) Router() http.Handler {
|
|
r := chi.NewRouter()
|
|
r.Use(requestLoggingMiddleware)
|
|
r.Use(corsMiddleware)
|
|
r.Route("/api", func(r chi.Router) {
|
|
r.Get("/health", s.handleHealth)
|
|
|
|
// Unprotected: these two ARE the login flow, so they can't require
|
|
// a session yet.
|
|
r.Get("/session/login", s.handleSessionLogin)
|
|
r.Get("/session/callback", s.handleSessionCallback)
|
|
|
|
r.Group(func(r chi.Router) {
|
|
r.Use(auth.RequireSession(s.Session.Secret))
|
|
r.Use(s.resolveUser)
|
|
|
|
r.Get("/session/me", s.handleSessionMe)
|
|
r.Post("/session/logout", s.handleSessionLogout)
|
|
r.Route("/setup", func(r chi.Router) {
|
|
r.Post("/complete", s.handleSetupComplete)
|
|
r.Route("/garmin", func(r chi.Router) {
|
|
r.Post("/login", s.handleSetupGarminLogin)
|
|
r.Post("/mfa", s.handleSetupGarminMFA)
|
|
})
|
|
})
|
|
|
|
r.Group(func(r chi.Router) {
|
|
r.Use(requireProvisionedUser)
|
|
|
|
r.Route("/profile", func(r chi.Router) {
|
|
r.Get("/", s.handleGetProfile)
|
|
r.Put("/", s.handleUpdateProfile)
|
|
r.Delete("/", s.handleDeleteProfile)
|
|
})
|
|
|
|
r.Route("/auth", func(r chi.Router) {
|
|
r.Post("/login", s.handleAuthLogin)
|
|
r.Post("/mfa", s.handleAuthMFA)
|
|
r.Get("/status", s.handleAuthStatus)
|
|
})
|
|
|
|
r.Route("/sync", func(r chi.Router) {
|
|
r.Post("/run", s.handleSyncRun)
|
|
r.Post("/reset", s.handleSyncReset)
|
|
r.Get("/runs", s.handleSyncRuns)
|
|
r.Get("/status", s.handleSyncStatus)
|
|
})
|
|
|
|
r.Route("/activities", func(r chi.Router) {
|
|
r.Get("/", s.handleListActivities)
|
|
r.Get("/{id}", s.handleGetActivity)
|
|
})
|
|
|
|
r.Route("/workout-kinds", func(r chi.Router) {
|
|
r.Get("/", s.handleListWorkoutKinds)
|
|
r.Get("/{id}", s.handleGetWorkoutKind)
|
|
r.Put("/{id}", s.handleUpdateWorkoutKind)
|
|
})
|
|
|
|
r.Post("/reclassify", s.handleReclassifyAll)
|
|
|
|
r.Route("/review-queue", func(r chi.Router) {
|
|
r.Get("/", s.handleReviewQueue)
|
|
r.Post("/{activityID}/resolve", s.handleResolveReview)
|
|
r.Post("/{activityID}/unlock", s.handleUnlockReview)
|
|
r.Post("/{activityID}/unassign", s.handleUnassignReview)
|
|
})
|
|
|
|
r.Get("/progression/{kindID}", s.handleProgression)
|
|
})
|
|
})
|
|
})
|
|
return r
|
|
}
|
|
|
|
// corsMiddleware allows the frontend dev server (a different port) to call
|
|
// this API. Reflecting any origin back is safe even with credentials
|
|
// enabled: this remains a single-operator app whose real access control is
|
|
// the OIDC login gate (internal/auth), not origin-based CSRF defense.
|
|
func corsMiddleware(next http.Handler) http.Handler {
|
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
if origin := r.Header.Get("Origin"); origin != "" {
|
|
w.Header().Set("Access-Control-Allow-Origin", origin)
|
|
w.Header().Set("Access-Control-Allow-Methods", "GET, POST, PUT, DELETE, OPTIONS")
|
|
w.Header().Set("Access-Control-Allow-Headers", "Content-Type")
|
|
w.Header().Set("Access-Control-Allow-Credentials", "true")
|
|
}
|
|
if r.Method == http.MethodOptions {
|
|
w.WriteHeader(http.StatusNoContent)
|
|
return
|
|
}
|
|
next.ServeHTTP(w, r)
|
|
})
|
|
}
|
|
|
|
func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) {
|
|
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
|
|
}
|
|
|
|
func writeJSON(w http.ResponseWriter, status int, v any) {
|
|
w.Header().Set("Content-Type", "application/json")
|
|
w.WriteHeader(status)
|
|
if err := json.NewEncoder(w).Encode(v); err != nil {
|
|
log.Printf("api: encode response: %v", err)
|
|
}
|
|
}
|
|
|
|
func writeError(w http.ResponseWriter, status int, msg string) {
|
|
writeJSON(w, status, map[string]string{"error": msg})
|
|
}
|
|
|
|
// backgroundSync runs fn in a goroutine with a fresh context, guarded so
|
|
// only one sync operation per userID runs at a time. Returns false if one
|
|
// is already in progress for that user.
|
|
func (s *Server) backgroundSync(userID int64, fn func(ctx context.Context) error) bool {
|
|
s.mu.Lock()
|
|
if s.userSyncRunning[userID] {
|
|
s.mu.Unlock()
|
|
return false
|
|
}
|
|
s.userSyncRunning[userID] = true
|
|
s.mu.Unlock()
|
|
|
|
go func() {
|
|
defer func() {
|
|
s.mu.Lock()
|
|
s.userSyncRunning[userID] = false
|
|
s.mu.Unlock()
|
|
}()
|
|
if err := fn(context.Background()); err != nil {
|
|
log.Printf("api: background sync error (user %d): %v", userID, err)
|
|
}
|
|
}()
|
|
return true
|
|
}
|