package server import ( "context" "crypto/rand" "crypto/sha256" "crypto/subtle" "encoding/base64" "encoding/json" "errors" "fmt" "io" iofs "io/fs" "log" "math" "net/http" "net/url" "os" "path/filepath" gort "runtime" "strconv" "strings" "sync" "sync/atomic" "time" "neuralhunt/internal/artifact" "neuralhunt/internal/auth" "neuralhunt/internal/core" "neuralhunt/internal/data" rtx "neuralhunt/internal/runtime" "neuralhunt/internal/settings" "neuralhunt/internal/webui" wsx "neuralhunt/internal/ws" "github.com/go-chi/chi/v5" "github.com/gorilla/websocket" ) type Server struct { store *data.Store auth *auth.Manager settings *settings.Manager hub *wsx.Hub placeHub *wsx.Hub runtime *rtx.State artifactWorker *artifact.Worker lottery *guessLottery adminUser, adminPass, staticDir, artifactDir string internalServiceSecret string customerServiceInternalURL string customerServiceHTTP *http.Client upgrader websocket.Upgrader wsAllowedOrigins map[string]struct{} maxUserWS, maxLeaderboardWS, maxPlaceWS int64 userWSCount, leaderboardWSCount, placeWSCount, adminWSCount atomic.Int64 adminSessions sync.Map // sid -> exp unix seconds placeTimeMu sync.Mutex placeTimeClocks map[string]placeTimeClock placeTimeSweep time.Time } func New(store *data.Store, a *auth.Manager, sm *settings.Manager, hub *wsx.Hub, runtimeState *rtx.State, artifactDir string, artifactWorker *artifact.Worker) *Server { s := &Server{ store: store, auth: a, settings: sm, hub: hub, placeHub: wsx.New(), runtime: runtimeState, placeTimeClocks: make(map[string]placeTimeClock), artifactWorker: artifactWorker, lottery: newGuessLottery(func(d beaconDrawAudit) { _ = store.RecordBeaconDraw(context.Background(), d.TaskID, d.WindowEnd, d.BeaconID, d.BeaconRound, d.Randomness, d.Signature, d.BoostedPath, d.Tickets, d.Selected) }), adminUser: env("ADMIN_USER", "admin"), adminPass: env("ADMIN_PASSWORD", "change-me"), staticDir: env("STATIC_DIR", ""), artifactDir: artifactDir, internalServiceSecret: strings.TrimSpace(os.Getenv("CUSTOMER_SERVICE_SHARED_SECRET")), customerServiceInternalURL: strings.TrimRight(strings.TrimSpace(os.Getenv("CUSTOMER_SERVICE_INTERNAL_URL")), "/"), customerServiceHTTP: &http.Client{Timeout: 5 * time.Second}, wsAllowedOrigins: parseOriginAllowlist(os.Getenv("WS_ALLOWED_ORIGINS")), maxUserWS: int64(envIntServer("WS_MAX_USER_CONNECTIONS", 5000)), maxLeaderboardWS: int64(envIntServer("WS_MAX_LEADERBOARD_CONNECTIONS", 500)), maxPlaceWS: int64(envIntServer("WS_MAX_PLACE_CONNECTIONS", 1500)), } s.upgrader = websocket.Upgrader{CheckOrigin: s.checkWSOrigin, Subprotocols: []string{"neuralhunt.v1"}} return s } func envIntServer(k string, d int) int { if raw := strings.TrimSpace(os.Getenv(k)); raw != "" { if n, err := strconv.Atoi(raw); err == nil && n >= 0 { return n } } return d } func parseOriginAllowlist(raw string) map[string]struct{} { out := map[string]struct{}{} for _, item := range strings.Split(raw, ",") { item = strings.TrimSpace(strings.TrimRight(item, "/")) if item != "" { out[strings.ToLower(item)] = struct{}{} } } return out } func (s *Server) checkWSOrigin(r *http.Request) bool { origin := strings.TrimSpace(r.Header.Get("Origin")) if origin == "" { // Native clients do not normally send Origin. Authentication and global // connection caps still apply; browser cross-site websockets do send it. return true } u, err := url.Parse(origin) if err != nil || u.Host == "" { return false } if strings.EqualFold(u.Host, r.Host) { return true } _, ok := s.wsAllowedOrigins[strings.ToLower(strings.TrimRight(origin, "/"))] return ok } func env(k, d string) string { if v := os.Getenv(k); v != "" { return v } return d } func jsonOut(w http.ResponseWriter, status int, v any) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(status) _ = json.NewEncoder(w).Encode(v) } func jsonAPIError(w http.ResponseWriter, status int, code, message string, extra map[string]any) { body := map[string]any{"error": message, "code": code} for k, v := range extra { body[k] = v } jsonOut(w, status, body) } func decode(r *http.Request, v any) error { d := json.NewDecoder(io.LimitReader(r.Body, 1<<20)) d.DisallowUnknownFields() if err := d.Decode(v); err != nil { return err } // Reject a second JSON value while still allowing insignificant trailing // whitespace. This keeps the request format unambiguous. var extra any if err := d.Decode(&extra); err != io.EOF { if err == nil { return errors.New("multiple JSON values") } return err } return nil } // Authentication receives a standards-compliant JWK exported by the browser. // Different WebCrypto implementations may add optional JWK members such as // alg/use/kid (and future implementations may add more). Those members are not // security relevant here: ClientID/PublicKey only consume kty, crv, x and y. // Therefore auth payloads intentionally accept unknown nested JWK properties. func decodeAuth(r *http.Request, v any) error { d := json.NewDecoder(io.LimitReader(r.Body, 1<<20)) if err := d.Decode(v); err != nil { return err } var extra any if err := d.Decode(&extra); err != io.EOF { if err == nil { return errors.New("multiple JSON values") } return err } return nil } type ctxKey string const claimsKey ctxKey = "claims" func (s *Server) bearer(r *http.Request) (auth.Claims, error) { h := r.Header.Get("Authorization") if !strings.HasPrefix(h, "Bearer ") { return auth.Claims{}, errors.New("missing bearer") } return s.auth.Parse(strings.TrimPrefix(h, "Bearer ")) } const adminSessionCookie = "neuralhunt_admin_session" func adminCookieSecure(r *http.Request) bool { raw := strings.ToLower(strings.TrimSpace(os.Getenv("ADMIN_COOKIE_SECURE"))) switch raw { case "0", "false", "no", "off": return false case "auto": if r.TLS != nil { return true } return strings.EqualFold(strings.TrimSpace(strings.Split(r.Header.Get("X-Forwarded-Proto"), ",")[0]), "https") default: // Public deployments should terminate HTTPS in front of this listener. // Defaulting to Secure avoids a proxy-header mistake silently weakening // the admin session. Local plain-HTTP development can set false explicitly. return true } } func (s *Server) adminCookieClaims(r *http.Request) (auth.Claims, error) { cookie, err := r.Cookie(adminSessionCookie) if err != nil || strings.TrimSpace(cookie.Value) == "" { return auth.Claims{}, errors.New("missing admin session") } c, err := s.auth.Parse(cookie.Value) if err != nil || c.Role != "admin" || c.SessionID == "" { return auth.Claims{}, errors.New("invalid admin session") } expRaw, ok := s.adminSessions.Load(c.SessionID) exp, typed := expRaw.(int64) if !ok || !typed || exp < time.Now().Unix() { s.adminSessions.Delete(c.SessionID) return auth.Claims{}, errors.New("admin session revoked or expired") } return c, nil } func (s *Server) require(role string, next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { var c auth.Claims var err error if role == "admin" { c, err = s.adminCookieClaims(r) } else { c, err = s.bearer(r) } if err != nil || (role != "" && c.Role != role) { jsonOut(w, 401, map[string]string{"error": "unauthorized"}) return } next.ServeHTTP(w, r.WithContext(context.WithValue(r.Context(), claimsKey, c))) }) } func claims(r *http.Request) auth.Claims { return r.Context().Value(claimsKey).(auth.Claims) } // PublicRoutes is safe to expose to the Internet. The admin UI and every // /api/admin endpoint are deliberately absent from this listener. func (s *Server) PublicRoutes() http.Handler { next := s.Routes() return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { p := strings.ToLower(r.URL.Path) if p == "/admin" || strings.HasPrefix(p, "/admin/") || p == "/api/admin" || strings.HasPrefix(p, "/api/admin/") || p == "/api/internal" || strings.HasPrefix(p, "/api/internal/") { http.NotFound(w, r) return } next.ServeHTTP(w, r) }) } // AdminRoutes is intended for the private/VPN listener. It serves only the // control plane, its static frontend assets and health check. func (s *Server) AdminRoutes() http.Handler { next := s.Routes() return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { p := strings.ToLower(r.URL.Path) if p == "/" { http.Redirect(w, r, "/admin", http.StatusTemporaryRedirect) return } allowed := p == "/admin" || strings.HasPrefix(p, "/admin/") || p == "/api/healthz" || p == "/api/admin" || strings.HasPrefix(p, "/api/admin/") || p == "/api/internal" || strings.HasPrefix(p, "/api/internal/") || p == "/app.js" || p == "/styles.css" || p == "/index.html" if !allowed { http.NotFound(w, r) return } next.ServeHTTP(w, r) }) } func (s *Server) Routes() http.Handler { r := chi.NewRouter() r.Use(func(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { w.Header().Set("X-Content-Type-Options", "nosniff") w.Header().Set("X-Frame-Options", "DENY") w.Header().Set("Referrer-Policy", "no-referrer") w.Header().Set("Permissions-Policy", "camera=(), microphone=(), geolocation=()") w.Header().Set("Content-Security-Policy", "default-src 'self'; base-uri 'none'; object-src 'none'; frame-ancestors 'none'; script-src 'self'; style-src 'self' 'unsafe-inline'; img-src 'self' data: blob:; connect-src 'self' ws: wss:; font-src 'self' data:; form-action 'self'") w.Header().Set("Strict-Transport-Security", "max-age=31536000") w.Header().Set("Cross-Origin-Opener-Policy", "same-origin") path := strings.ToLower(req.URL.Path) if path == "/" || path == "/admin" || path == "/leaderboard" || path == "/place" || strings.HasPrefix(path, "/api/admin") || strings.HasPrefix(path, "/api/auth") || strings.HasSuffix(path, ".html") || strings.HasSuffix(path, ".js") || strings.HasSuffix(path, ".css") { w.Header().Set("Cache-Control", "no-store") } next.ServeHTTP(w, req) }) }) r.Get("/api/healthz", func(w http.ResponseWriter, r *http.Request) { jsonOut(w, 200, map[string]bool{"ok": true}) }) r.Get("/api/public/leaderboard", s.publicLeaderboard) r.Get("/api/public/artifacts", s.publicArtifacts) r.Get("/api/public/beacon/{id}/latest", s.latestBeaconDraw) r.Get("/api/public/tasks", s.publicTaskCatalog) r.Get("/api/public/place", s.publicPlace) r.Get("/api/public/place/changes", s.publicPlaceChanges) r.Get("/api/public/artifacts/{id}/preview", s.publicArtifactPreview) r.Get("/api/public/tasks/{id}/style-reference", s.publicTaskStyleReference) r.Get("/api/leaderboard/ws", s.leaderboardWS) r.Get("/api/place/ws", s.placeWS) r.Post("/api/auth/challenge", s.challenge) r.Post("/api/auth/login", s.login) r.Post("/api/admin/login", s.adminLogin) r.Post("/api/admin/logout", s.adminLogout) r.Post("/api/internal/delegations", s.internalDelegation) r.Post("/api/internal/identity-exists", s.internalIdentityExists) r.Post("/api/internal/customer-link/consume", s.internalCustomerLinkConsume) r.Get("/api/internal/customer-service/health", s.internalCustomerServiceHealth) r.Group(func(r chi.Router) { r.Use(func(n http.Handler) http.Handler { return s.require("user", n) }) r.Get("/api/tasks", s.clientTasks) r.Get("/api/tasks/current", s.currentTask) r.Post("/api/tasks/select", s.selectTask) r.Post("/api/tasks/{id}/guess", s.guess) r.Get("/api/tasks/{id}/points", s.points) r.Get("/api/me", s.me) r.Get("/api/place/me", s.placeMe) r.Post("/api/place/pixel", s.placePixel) r.Get("/api/me/artifacts", s.myArtifacts) r.Get("/api/me/artifacts/{id}/download", s.myArtifactDownload) r.Post("/api/me/customer-link", s.customerLinkCode) r.Get("/api/leaderboard", s.leaderboard) }) r.Group(func(r chi.Router) { r.Use(func(n http.Handler) http.Handler { return s.require("admin", n) }) r.Get("/api/admin/session", s.adminSession) r.Get("/api/admin/ws", s.adminWS) r.Get("/api/admin/overview", s.adminOverview) r.Get("/api/admin/performance", s.adminPerformance) r.Get("/api/admin/profiles/cleanup-preview", s.adminProfileCleanupPreview) r.Post("/api/admin/profiles/cleanup", s.adminProfileCleanup) r.Get("/api/admin/settings", s.adminSettingsGet) r.Put("/api/admin/settings", s.adminSettingsPut) r.Get("/api/admin/tasks", s.adminTasks) r.Get("/api/admin/clients", s.adminClients) r.Get("/api/admin/place/earnings", s.adminPlaceEarnings) r.Post("/api/admin/place/points", s.adminPlaceGrant) r.Get("/api/admin/tasks/{id}/points", s.adminPoints) r.Get("/api/admin/tasks/{id}/actions", s.adminTaskActions) r.Put("/api/admin/tasks/{id}/config", s.adminTaskConfigPut) r.Post("/api/admin/tasks/{id}/artifact/transfer", s.adminArtifactTransfer) r.Post("/api/admin/artifacts/drop", s.adminArtifactDrop) r.Get("/api/admin/artifacts/status", s.adminArtifactStatus) r.Post("/api/admin/tasks/{id}/actions", s.adminScheduleAction) r.Post("/api/admin/actions/{id}/cancel", s.adminCancelAction) r.Get("/api/admin/artifact/providers", s.adminArtifactProviders) r.Get("/api/admin/artifact/usage", s.adminArtifactUsage) r.Post("/api/admin/artifact/character-anchor", s.adminCreateCharacterAnchor) r.Get("/api/admin/artifact/character-anchor", s.adminCharacterAnchorFile) r.Put("/api/admin/tasks/{id}/style-reference", s.adminTaskStyleReferencePut) r.Delete("/api/admin/tasks/{id}/style-reference", s.adminTaskStyleReferenceDelete) r.Get("/api/admin/tasks/{id}/style-reference", s.adminTaskStyleReferenceFile) r.Post("/api/admin/tasks/{id}/pipeline-test", s.adminCreatePipelineTestCard) r.Get("/api/admin/tasks/{id}/pipeline-test/card", s.adminPipelineTestCardFile) r.Get("/api/admin/tasks/{id}/artifact", s.adminArtifactFile) r.Get("/api/admin/tasks/{id}/manifest", s.adminArtifactManifest) r.Post("/api/admin/tasks/{id}/close", s.adminCloseTask) r.Post("/api/admin/tasks/ensure", s.adminEnsure) }) r.Get("/api/ws", s.ws) // Original winner artifacts are intentionally not publicly file-served. // Public viewers only receive /api/public/artifacts/{id}/preview, which is // watermarked. Keep the old namespace as an explicit 404 instead of letting // the SPA fallback accidentally return index.html for an artifact URL. r.Handle("/artifacts/*", http.NotFoundHandler()) // The production UI is embedded in the Go binary. STATIC_DIR remains an // optional development override, but Node.js/npm are never required to run // the application. var staticFS iofs.FS if s.staticDir != "" { if st, err := os.Stat(s.staticDir); err == nil && st.IsDir() { staticFS = os.DirFS(s.staticDir) log.Printf("serving frontend override from %q", s.staticDir) } } if staticFS == nil { embedded, err := iofs.Sub(webui.Dist, "dist") if err != nil { panic(fmt.Errorf("embedded frontend: %w", err)) } staticFS = embedded } fileServer := http.FileServer(http.FS(staticFS)) spa := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { name := strings.TrimPrefix(r.URL.Path, "/") if name == "" { name = "index.html" } if info, err := iofs.Stat(staticFS, name); err == nil && !info.IsDir() { fileServer.ServeHTTP(w, r) return } // Browser-side routes such as /admin receive the SPA entry point. b, err := iofs.ReadFile(staticFS, "index.html") if err != nil { http.Error(w, "embedded frontend unavailable", http.StatusInternalServerError) return } w.Header().Set("Content-Type", "text/html; charset=utf-8") _, _ = w.Write(b) }) r.Handle("/", spa) r.Handle("/*", spa) return r } func (s *Server) challenge(w http.ResponseWriter, r *http.Request) { var in struct { PublicJWK auth.PublicJWK `json:"public_jwk"` } if err := decodeAuth(r, &in); err != nil { log.Printf("auth challenge decode: %v", err) jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } cid, err := auth.ClientID(in.PublicJWK) if err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } proofBits := 0 if !s.store.ClientExists(r.Context(), cid) { proofBits = s.settings.Get().SybilProofOfWorkBits } c, err := s.auth.NewChallengeWithProof(r.Context(), cid, proofBits) if err != nil { log.Printf("auth challenge for %s: %v", cid, err) jsonOut(w, 500, map[string]string{"error": "challenge failed"}) return } jsonOut(w, 200, map[string]any{"client_id": cid, "challenge": c, "proof_of_work_bits": proofBits}) } func (s *Server) login(w http.ResponseWriter, r *http.Request) { var in struct { PublicJWK auth.PublicJWK `json:"public_jwk"` Challenge string `json:"challenge"` Signature string `json:"signature"` ProofOfWorkCounter string `json:"proof_of_work_counter"` } if err := decodeAuth(r, &in); err != nil { log.Printf("auth login decode: %v", err) jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } cid, err := auth.ClientID(in.PublicJWK) if err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } pub, err := auth.PublicKey(in.PublicJWK) if err != nil || !auth.VerifyRaw(pub, "login|"+in.Challenge+"|"+cid, in.Signature) { jsonOut(w, 401, map[string]string{"error": "invalid signature"}) return } if err = s.auth.ConsumeChallengeWithProof(r.Context(), cid, in.Challenge, in.ProofOfWorkCounter); err != nil { log.Printf("auth login challenge for %s: %v", cid, err) jsonOut(w, 401, map[string]string{"error": err.Error()}) return } if err = s.store.UpsertClient(r.Context(), cid, in.PublicJWK); err != nil { log.Printf("auth store client %s: %v", cid, err) jsonOut(w, 500, map[string]string{"error": "client store failed"}) return } tok, _, err := s.auth.Issue(cid, "user", 24*time.Hour) if err != nil { log.Printf("auth issue token for %s: %v", cid, err) jsonOut(w, 500, map[string]string{"error": "token creation failed"}) return } jsonOut(w, 200, map[string]string{"token": tok, "client_id": cid}) } func (s *Server) adminLogin(w http.ResponseWriter, r *http.Request) { var in struct { User string `json:"user"` Password string `json:"password"` } if decode(r, &in) != nil || in.User != s.adminUser || in.Password != s.adminPass { jsonOut(w, 401, map[string]string{"error": "invalid credentials"}) return } tok, sessionClaims, err := s.auth.Issue("admin", "admin", 8*time.Hour) if err != nil { log.Printf("admin auth issue token: %v", err) jsonOut(w, 500, map[string]string{"error": "session creation failed"}) return } nowUnix := time.Now().Unix() s.adminSessions.Range(func(key, value any) bool { if exp, ok := value.(int64); !ok || exp < nowUnix { s.adminSessions.Delete(key) } return true }) s.adminSessions.Store(sessionClaims.SessionID, sessionClaims.Exp) http.SetCookie(w, &http.Cookie{ Name: adminSessionCookie, Value: tok, Path: "/", MaxAge: int((8 * time.Hour).Seconds()), HttpOnly: true, Secure: adminCookieSecure(r), SameSite: http.SameSiteStrictMode, }) jsonOut(w, 200, map[string]bool{"ok": true}) } func (s *Server) adminLogout(w http.ResponseWriter, r *http.Request) { if cookie, err := r.Cookie(adminSessionCookie); err == nil { if c, err := s.auth.Parse(cookie.Value); err == nil && c.SessionID != "" { s.adminSessions.Delete(c.SessionID) } } http.SetCookie(w, &http.Cookie{ Name: adminSessionCookie, Value: "", Path: "/", MaxAge: -1, HttpOnly: true, Secure: adminCookieSecure(r), SameSite: http.SameSiteStrictMode, }) jsonOut(w, 200, map[string]bool{"ok": true}) } func (s *Server) adminSession(w http.ResponseWriter, r *http.Request) { c := claims(r) jsonOut(w, 200, map[string]any{"ok": true, "role": c.Role, "expires_at": time.Unix(c.Exp, 0).UTC()}) } func taskIntervals(t data.Task, sm settings.Runtime) (int, int) { serverMin := sm.GuessMinIntervalSec clientSubmit := sm.ClientSubmitIntervalSec if t.GuessMinIntervalSec != nil { serverMin = *t.GuessMinIntervalSec } if t.ClientSubmitIntervalSec != nil { clientSubmit = *t.ClientSubmitIntervalSec } return serverMin, clientSubmit } func taskDTO(t data.Task, next int64, sm settings.Runtime) map[string]any { serverMin, clientSubmit := taskIntervals(t, sm) return map[string]any{ "id": t.ID, "public_seed": t.PublicSeed, "range_bits": t.RangeBits, "next_seq": next, "server_min_interval_sec": serverMin, "client_submit_interval_sec": clientSubmit, "guess_lottery_window_sec": sm.GuessLotteryWindowSec, "guess_lottery_max_accepted": sm.GuessLotteryMaxAccepted, "beacon_hunt_enabled": sm.BeaconHuntEnabled, "beacon_bonus_weight": sm.BeaconBonusWeight, "beacon_paths": beaconPaths, "default_max_nodes": sm.DefaultMaxNodes, "paused": t.Paused, "revision": t.Revision, "display_name": t.DisplayName, "description": t.Description, "parent_task_id": t.ParentTaskID, "retire_after_completion": t.RetireAfterCompletion, "created_at": t.CreatedAt, } } func (s *Server) clientTasks(w http.ResponseWriter, r *http.Request) { c := claims(r) if !s.store.ClientExists(r.Context(), c.ClientID) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "identity not registered in this database"}) return } items, err := s.store.ActiveTasksForClient(r.Context(), c.ClientID) if err != nil { jsonOut(w, 500, map[string]string{"error": "tasks failed"}) return } if len(items) == 0 { _ = s.store.EnsureActiveTasks(r.Context(), s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits) items, err = s.store.ActiveTasksForClient(r.Context(), c.ClientID) if err != nil { jsonOut(w, 500, map[string]string{"error": "tasks failed"}) return } } // Keep score disclosure consistent with the regular client endpoints. prec := s.settings.Get().PublicScorePrecision for i := range items { items[i].OwnScore = round(items[i].OwnScore, prec) } jsonOut(w, 200, items) } func (s *Server) selectTask(w http.ResponseWriter, r *http.Request) { c := claims(r) var in struct { TaskID string `json:"task_id"` } if err := decode(r, &in); err != nil { jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } in.TaskID = strings.TrimSpace(in.TaskID) if in.TaskID == "" { jsonOut(w, 400, map[string]string{"error": "task_id required"}) return } if err := s.store.SetClientTaskSelection(r.Context(), c.ClientID, in.TaskID); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } s.runtime.SetTaskSelection(c.ClientID, in.TaskID) t, err := s.store.TaskForClient(r.Context(), c.ClientID) if err != nil { jsonOut(w, 500, map[string]string{"error": "task selection failed"}) return } s.runtime.SetTaskSelection(c.ClientID, t.ID) snap, _ := s.store.LoadGuessState(r.Context(), t.ID, c.ClientID) g := s.runtime.InitGuess(t, c.ClientID, rtx.GuessState{NextSeq: snap.NextSeq, LastGuess: snap.LastGuess, BestScore: snap.BestScore, GuessCount: snap.GuessCount}) jsonOut(w, 200, taskDTO(t, g.NextSeq, s.settings.Get())) } func (s *Server) currentTask(w http.ResponseWriter, r *http.Request) { c := claims(r) if !s.store.ClientExists(r.Context(), c.ClientID) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "identity not registered in this database"}) return } t, err := s.store.TaskForClient(r.Context(), c.ClientID) if data.IsNoRows(err) { _ = s.store.EnsureActiveTasks(r.Context(), s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits) t, err = s.store.TaskForClient(r.Context(), c.ClientID) } if err != nil { jsonOut(w, 503, map[string]string{"error": "no active task"}) return } s.runtime.SetTaskSelection(c.ClientID, t.ID) snap, _ := s.store.LoadGuessState(r.Context(), t.ID, c.ClientID) g := s.runtime.InitGuess(t, c.ClientID, rtx.GuessState{NextSeq: snap.NextSeq, LastGuess: snap.LastGuess, BestScore: snap.BestScore, GuessCount: snap.GuessCount}) jsonOut(w, 200, taskDTO(t, g.NextSeq, s.settings.Get())) } func guessMsg(taskID string, seq int64, guess, beaconPath string, beaconEnabled bool) string { if beaconEnabled { return fmt.Sprintf("guess|%s|%d|%s|%s", taskID, seq, guess, normalizeBeaconPath(beaconPath)) } return fmt.Sprintf("guess|%s|%d|%s", taskID, seq, guess) } func round(v float64, p int) float64 { m := math.Pow10(p) return math.Round(v*m) / m } func (s *Server) guess(w http.ResponseWriter, r *http.Request) { c := claims(r) id := chi.URLParam(r, "id") if selected := s.runtime.TaskSelection(c.ClientID); selected != id { durable, err := s.store.SelectedTaskID(r.Context(), c.ClientID) if err != nil || durable != id { jsonAPIError(w, http.StatusConflict, "selection_conflict", "task selection changed", nil) return } s.runtime.SetTaskSelection(c.ClientID, id) } if !s.runtime.HasPresence(c.ClientID, c.SessionID) { jsonAPIError(w, http.StatusConflict, "presence_required", "live connection lost; reconnect websocket", nil) return } var in struct { Seq int64 `json:"seq"` Guess string `json:"guess"` Signature string `json:"signature"` BeaconPath string `json:"beacon_path,omitempty"` } if decode(r, &in) != nil { jsonOut(w, 400, false) return } if warmup := s.settings.Get().SybilWarmupSec; warmup > 0 { created, err := s.store.ClientCreatedAt(r.Context(), c.ClientID) if err != nil { jsonOut(w, 401, false) return } readyAt := created.Add(time.Duration(warmup) * time.Second) if wait := time.Until(readyAt); wait > 0 { retry := int(math.Ceil(wait.Seconds())) jsonAPIError(w, http.StatusTooManyRequests, "identity_warmup", "new identity is still in anti-sybil warmup", map[string]any{"retry_after_sec": retry}) return } } t, err := s.store.SecretTask(r.Context(), id) if err != nil || t.Status != "active" { jsonAPIError(w, http.StatusConflict, "task_inactive", "task is no longer active", nil) return } if t.Paused { jsonOut(w, 423, false) return } if core.ExpectedGuess(id, t.PublicSeed, c.ClientID, in.Seq, t.RangeBits) != in.Guess { // This most commonly means an admin rerolled/changed the task between the // client's config read and its submit. Treat it as recoverable config drift. jsonAPIError(w, http.StatusConflict, "task_config_changed", "task configuration changed; resync required", map[string]any{"revision": t.Revision}) return } jwk, err := s.store.ClientPublicJWK(r.Context(), c.ClientID) if err != nil { jsonOut(w, 401, false) return } pub, _ := auth.PublicKey(jwk) beaconEnabled := s.settings.Get().BeaconHuntEnabled == 1 && s.settings.Get().GuessLotteryMaxAccepted > 0 if beaconEnabled && normalizeBeaconPath(in.BeaconPath) == "" { jsonAPIError(w, http.StatusBadRequest, "beacon_path_required", "choose PULSE, FLUX or ORBIT before entering the draw", map[string]any{"paths": beaconPaths}) return } if pub == nil || !auth.VerifyRaw(pub, guessMsg(id, in.Seq, in.Guess, in.BeaconPath, beaconEnabled), in.Signature) { jsonOut(w, 401, false) return } if _, ok := s.runtime.Current(t.Task, c.ClientID); !ok { snap, _ := s.store.LoadGuessState(r.Context(), t.ID, c.ClientID) s.runtime.InitGuess(t.Task, c.ClientID, rtx.GuessState{NextSeq: snap.NextSeq, LastGuess: snap.LastGuess, BestScore: snap.BestScore, GuessCount: snap.GuessCount}) } cfg := s.settings.Get() serverMin, _ := taskIntervals(t.Task, cfg) minInterval := time.Duration(serverMin) * time.Second if err := s.runtime.CanSubmit(t.Task, c.ClientID, in.Seq, minInterval); err != nil { switch { case errors.Is(err, rtx.ErrRateLimited): jsonOut(w, 429, false) case errors.Is(err, rtx.ErrBadSequence): expected := int64(0) if cur, ok := s.runtime.Current(t.Task, c.ClientID); ok { expected = cur.NextSeq } jsonAPIError(w, http.StatusConflict, "sequence_mismatch", "guess sequence is stale", map[string]any{"next_seq": expected}) default: jsonOut(w, 500, false) } return } // Active-time Place rewards are driven only by valid, signed Hunt activity // with an active presence. Duplicate sequence numbers are ignored by the // in-memory activity clock, and long offline gaps are never backfilled. if cfg.PlaceEnabled == 1 && cfg.PlaceTimePoints > 0 { if err := s.rewardPlaceActiveTime(r.Context(), t.ID, c.ClientID, t.Revision, in.Seq, cfg.PlaceTimePoints, cfg.PlaceTimeIntervalSec, cfg.PlaceTimeMaxGapSec); err != nil { log.Printf("place active-time reward %s/%s seq=%d: %v", t.ID, c.ClientID, in.Seq, err) } } var draw lotteryResult if cfg.GuessLotteryMaxAccepted > 0 { drawResult, drawErr := s.lottery.enter(r.Context(), t.ID, c.ClientID, in.Seq, in.BeaconPath, time.Duration(cfg.GuessLotteryWindowSec)*time.Second, cfg.GuessLotteryMaxAccepted, cfg.BeaconHuntEnabled == 1, cfg.BeaconBonusWeight) draw = drawResult if drawErr != nil { switch { case errors.Is(drawErr, errLotteryDuplicate): jsonAPIError(w, http.StatusConflict, "lottery_duplicate", "guess is already waiting for the current draw", nil) case errors.Is(drawErr, errLotteryFull): jsonAPIError(w, http.StatusTooManyRequests, "lottery_full", "guess lottery window is full", nil) case errors.Is(drawErr, errBeaconUnavailable): jsonAPIError(w, http.StatusServiceUnavailable, "beacon_unavailable", "external randomness beacon is temporarily unavailable; ticket was not evaluated", nil) case errors.Is(drawErr, context.Canceled), errors.Is(drawErr, context.DeadlineExceeded): return default: jsonOut(w, 500, false) } return } // A task may have been rerolled, paused or completed while this request // waited for the draw. Never consume/evaluate a stale lottery ticket. fresh, freshErr := s.store.SecretTask(r.Context(), id) if freshErr != nil || fresh.Status != "active" { jsonAPIError(w, http.StatusConflict, "task_inactive", "task is no longer active", nil) return } if fresh.Paused { jsonOut(w, 423, false) return } if fresh.PublicSeed != t.PublicSeed || fresh.Revision != t.Revision { jsonAPIError(w, http.StatusConflict, "task_config_changed", "task configuration changed during lottery; resync required", map[string]any{"revision": fresh.Revision}) return } t = fresh if draw.Selected { rewardCfg := s.settings.Get() if rewardCfg.PlaceEnabled == 1 && rewardCfg.PlaceDrawPoints > 0 { weight := draw.Weight if weight < 1 { weight = 1 } if _, _, created, err := s.store.PlaceRewardDraw(r.Context(), t.ID, c.ClientID, t.Revision, in.Seq, rewardCfg.PlaceDrawPoints, weight, rewardCfg.PlaceDrawBeaconMultiplier == 1); err != nil { log.Printf("place draw reward %s/%s seq=%d: %v", t.ID, c.ClientID, in.Seq, err) } else if created { s.runtime.MarkSQLiteWrite() } } } if !draw.Selected { next, skipErr := s.runtime.SkipLottery(t.Task, c.ClientID, in.Seq, minInterval) if skipErr != nil { if errors.Is(skipErr, rtx.ErrBadSequence) { jsonAPIError(w, http.StatusConflict, "sequence_mismatch", "guess sequence changed while waiting for lottery", map[string]any{"next_seq": next}) } else if errors.Is(skipErr, rtx.ErrRateLimited) { jsonOut(w, 429, false) } else { jsonOut(w, 500, false) } return } extra := map[string]any{"next_seq": next} if draw.BeaconEnabled { extra["chosen_path"] = draw.ChosenPath extra["boosted_path"] = draw.BoostedPath extra["beacon_round"] = draw.BeaconRound extra["beacon_id"] = draw.BeaconID extra["weight"] = draw.Weight } jsonAPIError(w, http.StatusTooManyRequests, "lottery_not_selected", "guess was not selected in this lottery window", extra) return } } d, err := core.Distance(in.Guess, t.Secret) if err != nil { jsonOut(w, 400, false) return } correct := d.Sign() == 0 score := core.Score(d, t.RangeBits) accepted, err := s.runtime.Accept(t.Task, c.ClientID, in.Seq, score, minInterval) if err != nil { switch { case errors.Is(err, rtx.ErrRateLimited): jsonOut(w, 429, false) case errors.Is(err, rtx.ErrBadSequence): expected := int64(0) if cur, ok := s.runtime.Current(t.Task, c.ClientID); ok { expected = cur.NextSeq } jsonAPIError(w, http.StatusConflict, "sequence_mismatch", "guess sequence is stale", map[string]any{"next_seq": expected}) default: jsonOut(w, 500, false) } return } // Losing tips are intentionally ephemeral: no SQLite write and no websocket event. if accepted.Improved || correct { rewardOwner := c.ClientID if correct { rewardOwner = s.store.RewardOwnerForWorker(r.Context(), c.ClientID) } beaconPath, beaconBoost, beaconRound := "", "", uint64(0) if correct && draw.BeaconEnabled { beaconPath, beaconBoost, beaconRound = draw.ChosenPath, draw.BoostedPath, draw.BeaconRound } placePointsPerScore := 0 if placeCfg := s.settings.Get(); placeCfg.PlaceEnabled == 1 { placePointsPerScore = placeCfg.PlacePointsPerScore } p, err := s.store.PersistImprovement(r.Context(), t, c.ClientID, rewardOwner, accepted.State.NextSeq, accepted.State.GuessCount, accepted.State.LastGuess, accepted.State.BestScore, in.Guess, in.Signature, correct, beaconPath, beaconBoost, beaconRound, placePointsPerScore) if err != nil { s.runtime.Restore(t.Task, c.ClientID, accepted.State.NextSeq, accepted.Previous) switch { case errors.Is(err, data.ErrTaskCompleted): jsonAPIError(w, http.StatusConflict, "task_inactive", "task completed while submitting", nil) case errors.Is(err, data.ErrTaskPaused): jsonOut(w, 423, false) default: jsonOut(w, 500, false) } return } s.runtime.MarkSQLiteWrite() // A positive hosted-worker tip means a genuine personal-best improvement, // not merely a submitted or lottery-selected guess. Record it durably and // let the scheduler deliver the idempotent credit event to Customer Service. if accepted.Improved { rewardOwner := s.store.RewardOwnerForWorker(r.Context(), c.ClientID) if _, created, err := s.store.RecordHostedPositiveCreditEvent(r.Context(), id, c.ClientID, rewardOwner, in.Seq, accepted.State.BestScore); err != nil { log.Printf("hosted positive-tip outbox: %v", err) } else if created { s.runtime.MarkSQLiteWrite() } } p.Score = round(p.Score, s.settings.Get().PublicScorePrecision) s.hub.PublishPoint(id, c.ClientID, p) } // Admin-only live telemetry: expose the score of every actually evaluated // guess, even when it is below the client's personal best. This event is // ephemeral and never reaches public/user websocket clients. _ = s.hub.PublishAdmin(r.Context(), wsx.Event{Type: "guess_signal", TaskID: id, Data: map[string]any{ "client_id": c.ClientID, "seq": in.Seq, "score": round(score, 2), "best_score": round(accepted.State.BestScore, 2), "improved": accepted.Improved, "correct": correct, }}) if correct { rewardOwner := s.store.RewardOwnerForWorker(r.Context(), c.ClientID) dataOut := map[string]string{"winner_client_id": rewardOwner, "winner_worker_client_id": c.ClientID} if successor, succErr := s.store.EnsureSuccessorTask(r.Context(), id, s.settings.Get().TaskRangeBits); succErr == nil { dataOut["successor_task_id"] = successor.ID s.runtime.ReplaceTaskSelection(id, successor.ID) } else if !errors.Is(succErr, data.ErrTaskRetiring) { log.Printf("successor for %s: %v", id, succErr) } else { dataOut["retired"] = "true" } _ = s.hub.Publish(r.Context(), wsx.Event{Type: "task_completed", TaskID: id, Data: dataOut}) _ = s.store.EnsureActiveTasks(r.Context(), s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits) } if draw.BeaconEnabled { w.Header().Set("X-NeuralHunt-Beacon-Path", draw.BoostedPath) w.Header().Set("X-NeuralHunt-Beacon-Round", strconv.FormatUint(draw.BeaconRound, 10)) } jsonOut(w, 200, correct) } func serviceTokenOK(secret, header string) bool { provided := strings.TrimSpace(strings.TrimPrefix(header, "Bearer ")) return secret != "" && len(secret) == len(provided) && subtle.ConstantTimeCompare([]byte(secret), []byte(provided)) == 1 } func customerLinkHash(code string) string { h := sha256.Sum256([]byte("nh-customer-link-v1|" + strings.TrimSpace(code))) return fmt.Sprintf("%x", h[:]) } // customerLinkCode issues a short-lived one-shot proof that the authenticated // browser/CLI controls this exact P-256 identity. Customer Service redeems the // code over the private 8081 control plane; the private key never leaves the // owner device. func (s *Server) customerLinkCode(w http.ResponseWriter, r *http.Request) { c := claims(r) b := make([]byte, 24) if _, err := rand.Read(b); err != nil { jsonOut(w, 500, map[string]string{"error": "could not create pairing code"}) return } code := "nhlink_" + base64.RawURLEncoding.EncodeToString(b) expires := time.Now().UTC().Add(10 * time.Minute) if err := s.store.CreateCustomerLinkToken(r.Context(), customerLinkHash(code), c.ClientID, expires); err != nil { jsonOut(w, 500, map[string]string{"error": "could not store pairing code"}) return } w.Header().Set("Cache-Control", "no-store") jsonOut(w, 201, map[string]any{"code": code, "client_id": c.ClientID, "expires_at": expires}) } func (s *Server) internalCustomerServiceHealth(w http.ResponseWriter, r *http.Request) { secret := strings.TrimSpace(s.internalServiceSecret) if !serviceTokenOK(secret, r.Header.Get("Authorization")) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "unauthorized"}) return } w.Header().Set("Cache-Control", "no-store") jsonOut(w, http.StatusOK, map[string]any{ "ok": true, "service": "neuralhunt-game-control-plane", "customer_link_supported": true, }) } func (s *Server) internalCustomerLinkConsume(w http.ResponseWriter, r *http.Request) { secret := strings.TrimSpace(s.internalServiceSecret) if !serviceTokenOK(secret, r.Header.Get("Authorization")) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "unauthorized"}) return } var in struct { Code string `json:"code"` } if err := decode(r, &in); err != nil || !strings.HasPrefix(strings.TrimSpace(in.Code), "nhlink_") { jsonOut(w, 400, map[string]string{"error": "valid pairing code required"}) return } cid, err := s.store.ConsumeCustomerLinkToken(r.Context(), customerLinkHash(in.Code)) if err != nil { jsonOut(w, 404, map[string]string{"error": "pairing code expired, invalid, or already used"}) return } jsonOut(w, 200, map[string]string{"client_id": cid}) } func (s *Server) internalIdentityExists(w http.ResponseWriter, r *http.Request) { secret := strings.TrimSpace(s.internalServiceSecret) if !serviceTokenOK(secret, r.Header.Get("Authorization")) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "unauthorized"}) return } var in struct { ClientID string `json:"client_id"` } if err := decode(r, &in); err != nil || strings.TrimSpace(in.ClientID) == "" { jsonOut(w, 400, map[string]string{"error": "client_id required"}) return } jsonOut(w, 200, map[string]bool{"exists": s.store.ClientExists(r.Context(), strings.TrimSpace(in.ClientID))}) } func (s *Server) internalDelegation(w http.ResponseWriter, r *http.Request) { secret := strings.TrimSpace(s.internalServiceSecret) if !serviceTokenOK(secret, r.Header.Get("Authorization")) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "unauthorized"}) return } var in struct { WorkerClientID string `json:"worker_client_id"` OwnerClientID string `json:"owner_client_id"` } if err := decode(r, &in); err != nil { jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } if err := s.store.SetIdentityDelegation(r.Context(), in.WorkerClientID, in.OwnerClientID); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, map[string]any{"ok": true, "worker_client_id": in.WorkerClientID, "owner_client_id": in.OwnerClientID}) } func (s *Server) publicTaskCatalog(w http.ResponseWriter, r *http.Request) { items, err := s.store.ActiveTasksForClient(r.Context(), "") if err != nil { jsonOut(w, 500, map[string]string{"error": "tasks failed"}) return } cfg := s.settings.Get() out := make([]map[string]any, 0, len(items)) for _, t := range items { out = append(out, map[string]any{ "id": t.ID, "display_name": t.DisplayName, "description": t.Description, "range_bits": t.RangeBits, "paused": t.Paused, "guess_lottery_window_sec": cfg.GuessLotteryWindowSec, "guess_lottery_max_accepted": cfg.GuessLotteryMaxAccepted, "beacon_hunt_enabled": cfg.BeaconHuntEnabled, "beacon_bonus_weight": cfg.BeaconBonusWeight, "beacon_paths": beaconPaths, "style_reference_uri": "/api/public/tasks/" + url.PathEscape(t.ID) + "/style-reference", }) } jsonOut(w, 200, out) } func (s *Server) latestBeaconDraw(w http.ResponseWriter, r *http.Request) { taskID := chi.URLParam(r, "id") d, err := s.store.LatestBeaconDraw(r.Context(), taskID) if err != nil { if data.IsNoRows(err) { jsonOut(w, 404, map[string]string{"error": "no beacon draw yet"}) return } jsonOut(w, 500, map[string]string{"error": "beacon draw lookup failed"}) return } jsonOut(w, 200, d) } func (s *Server) points(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) if limit <= 0 { limit = s.settings.Get().DefaultMaxNodes * 3 } if limit > 10000 { limit = 10000 } ps, err := s.store.PointsForClient(r.Context(), chi.URLParam(r, "id"), claims(r).ClientID, limit) if err != nil { jsonOut(w, 500, map[string]string{"error": "points failed"}) return } for i := range ps { ps[i].Score = round(ps[i].Score, s.settings.Get().PublicScorePrecision) } jsonOut(w, 200, ps) } func (s *Server) me(w http.ResponseWriter, r *http.Request) { c := claims(r) // A JWT can outlive a replaced/empty SQLite database. Do not treat such a // token as a valid registered browser identity; force challenge/login again // so the public key is upserted into this database. if !s.store.ClientExists(r.Context(), c.ClientID) { jsonOut(w, http.StatusUnauthorized, map[string]string{"error": "identity not registered in this database"}) return } t, err := s.store.TaskForClient(r.Context(), c.ClientID) if err != nil { jsonOut(w, 200, data.Me{ClientID: c.ClientID, Unlocks: []string{}}) return } m, _ := s.store.Me(r.Context(), c.ClientID, t.ID) m.Score = round(m.Score, s.settings.Get().PublicScorePrecision) jsonOut(w, 200, m) } func (s *Server) myArtifacts(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) items, err := s.store.OwnedArtifacts(r.Context(), claims(r).ClientID, limit) if err != nil { jsonOut(w, 500, map[string]string{"error": "owned artifacts failed"}) return } jsonOut(w, 200, items) } func (s *Server) myArtifactDownload(w http.ResponseWriter, r *http.Request) { id := strings.TrimSpace(chi.URLParam(r, "id")) if id == "" { jsonOut(w, 400, map[string]string{"error": "task id required"}) return } uri, ok, err := s.store.OwnedArtifactSource(r.Context(), id, claims(r).ClientID) if err != nil { jsonOut(w, 500, map[string]string{"error": "artifact lookup failed"}) return } if !ok { // Deliberately use 404 rather than revealing that another identity owns it. jsonOut(w, 404, map[string]string{"error": "artifact not found"}) return } path, err := artifactLocalPath(s.artifactDir, uri) if err != nil { jsonOut(w, 404, map[string]string{"error": "artifact file unavailable"}) return } ext := filepath.Ext(path) if ext == "" { ext = ".bin" } name := "neuralhunt-" + id + ext w.Header().Set("Cache-Control", "private, no-store") w.Header().Set("Content-Disposition", fmt.Sprintf(`attachment; filename=%q`, name)) http.ServeFile(w, r, path) } func (s *Server) leaderboard(w http.ResponseWriter, r *http.Request) { l, err := s.store.Leaderboard(r.Context(), 100) if err != nil { jsonOut(w, 500, map[string]string{"error": "leaderboard failed"}) return } for i := range l { l[i].Connected = s.runtime.IsConnected(l[i].ClientID) } jsonOut(w, 200, l) } func (s *Server) publicLeaderboard(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) mode := strings.ToLower(strings.TrimSpace(r.URL.Query().Get("mode"))) var ( l []data.Leader err error ) if mode == "live" { l, err = s.store.LiveLeaderboard(r.Context(), limit) } else { l, err = s.store.Leaderboard(r.Context(), limit) } if err != nil { jsonOut(w, 500, map[string]string{"error": "leaderboard failed"}) return } prec := s.settings.Get().PublicScorePrecision for i := range l { l[i].BestScore = round(l[i].BestScore, prec) l[i].LiveScore = round(l[i].LiveScore, prec) l[i].Connected = s.runtime.IsConnected(l[i].ClientID) } jsonOut(w, 200, l) } func (s *Server) publicArtifacts(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) winner := strings.TrimSpace(r.URL.Query().Get("winner")) items, err := s.store.PublicArtifacts(r.Context(), limit, winner) if err != nil { jsonOut(w, 500, map[string]string{"error": "artifact gallery failed"}) return } jsonOut(w, 200, items) } func (s *Server) publicArtifactPreview(w http.ResponseWriter, r *http.Request) { id := strings.TrimSpace(chi.URLParam(r, "id")) if id == "" { jsonOut(w, 400, map[string]string{"error": "task id required"}) return } artifactURI, _, ok, err := s.store.PublicArtifactSource(r.Context(), id) if err != nil { jsonOut(w, 500, map[string]string{"error": "artifact lookup failed"}) return } if !ok { jsonOut(w, 404, map[string]string{"error": "artifact not found"}) return } path, err := artifactLocalPath(s.artifactDir, artifactURI) if err != nil { log.Printf("artifact preview path %s: %v", id, err) jsonOut(w, 404, map[string]string{"error": "artifact file unavailable"}) return } preview, contentType, err := watermarkPreviewFile(path, "NEURAL HUNT PREVIEW") if err != nil { log.Printf("artifact preview render %s: %v", id, err) jsonOut(w, 500, map[string]string{"error": "preview generation failed"}) return } w.Header().Set("Content-Type", contentType) w.Header().Set("Cache-Control", "public, max-age=300") w.Header().Set("X-Neural-Hunt-Watermark", "leaderboard-preview") w.Header().Set("Content-Disposition", "inline") w.WriteHeader(http.StatusOK) _, _ = w.Write(preview) } func (s *Server) adminWS(w http.ResponseWriter, r *http.Request) { // The route itself is behind require("admin"), so the HttpOnly admin session // cookie is validated before the websocket upgrade. Keep a small independent // cap because this stream can contain one event for every evaluated guess. if !acquireWSCap(&s.adminWSCount, 32) { http.Error(w, "admin websocket capacity reached", http.StatusServiceUnavailable) return } defer s.adminWSCount.Add(-1) conn, err := s.upgrader.Upgrade(w, r, nil) if err != nil { return } cl := wsx.NewAdminClient(conn) s.hub.Add(cl) defer s.hub.Remove(cl) conn.SetReadLimit(1024) _ = conn.SetReadDeadline(time.Now().Add(90 * time.Second)) conn.SetPongHandler(func(string) error { return conn.SetReadDeadline(time.Now().Add(90 * time.Second)) }) ping := time.NewTicker(30 * time.Second) defer ping.Stop() done := make(chan struct{}) go func() { defer close(done) for { if _, _, err := conn.ReadMessage(); err != nil { return } } }() for { select { case <-done: return case <-r.Context().Done(): return case <-ping.C: if err := cl.Ping(); err != nil { return } } } } func (s *Server) leaderboardWS(w http.ResponseWriter, r *http.Request) { if !acquireWSCap(&s.leaderboardWSCount, s.maxLeaderboardWS) { http.Error(w, "leaderboard websocket capacity reached", http.StatusServiceUnavailable) return } defer s.leaderboardWSCount.Add(-1) conn, err := s.upgrader.Upgrade(w, r, nil) if err != nil { return } cl := wsx.NewClient(conn, "", "leaderboard", true) s.hub.Add(cl) defer s.hub.Remove(cl) initial, _ := s.store.LiveLeaderboard(r.Context(), 200) cl.Enqueue(wsx.Event{Type: "leaderboard", Data: initial}) conn.SetReadLimit(1024) _ = conn.SetReadDeadline(time.Now().Add(90 * time.Second)) conn.SetPongHandler(func(string) error { return conn.SetReadDeadline(time.Now().Add(90 * time.Second)) }) ping := time.NewTicker(30 * time.Second) defer ping.Stop() done := make(chan struct{}) go func() { defer close(done) for { if _, _, err := conn.ReadMessage(); err != nil { return } } }() for { select { case <-done: return case <-r.Context().Done(): return case <-ping.C: if err := cl.Ping(); err != nil { return } } } } func (s *Server) adminArtifactFile(w http.ResponseWriter, r *http.Request) { s.serveAdminArtifactPart(w, r, false) } func (s *Server) adminArtifactManifest(w http.ResponseWriter, r *http.Request) { s.serveAdminArtifactPart(w, r, true) } func (s *Server) serveAdminArtifactPart(w http.ResponseWriter, r *http.Request, manifest bool) { id := strings.TrimSpace(chi.URLParam(r, "id")) imageURI, manifestURI, ok, err := s.store.TaskArtifactURIs(r.Context(), id) if err != nil { jsonOut(w, 500, map[string]string{"error": "artifact lookup failed"}) return } if !ok { jsonOut(w, 404, map[string]string{"error": "artifact not found"}) return } uri := imageURI if manifest { uri = manifestURI if uri == "" { jsonOut(w, 404, map[string]string{"error": "manifest not found"}) return } } path, err := artifactLocalPath(s.artifactDir, uri) if err != nil { jsonOut(w, 404, map[string]string{"error": "artifact file unavailable"}) return } if manifest { w.Header().Set("Content-Type", "application/json; charset=utf-8") w.Header().Set("Content-Disposition", `inline; filename="manifest.json"`) } else { w.Header().Set("Content-Disposition", `inline`) } w.Header().Set("Cache-Control", "private, no-store") http.ServeFile(w, r, path) } type profileCleanupPreview struct { CutoffMS int64 `json:"cutoff_ms"` InactiveForSeconds int64 `json:"inactive_for_seconds"` Eligible int `json:"eligible"` ProtectedWinners int64 `json:"protected_winners"` ProtectedPlace int64 `json:"protected_place"` ProtectedConnected int `json:"protected_connected"` OldestEligibleMS int64 `json:"oldest_eligible_ms,omitempty"` NewestEligibleMS int64 `json:"newest_eligible_ms,omitempty"` } func profileCleanupDurationSeconds(raw string) (int64, error) { seconds, err := strconv.ParseInt(strings.TrimSpace(raw), 10, 64) if err != nil { return 0, errors.New("inactive_for_seconds must be an integer") } // A one-hour minimum prevents an accidental near-live purge while still // allowing short-lived development/test deployments to clean up quickly. if seconds < 3600 || seconds > 10*365*24*60*60 { return 0, errors.New("inactive_for_seconds must be between 3600 seconds and 10 years") } return seconds, nil } func (s *Server) profileCleanupPreview(ctx context.Context, inactiveForSeconds int64) (profileCleanupPreview, []string, error) { cutoff := time.Now().UTC().Add(-time.Duration(inactiveForSeconds) * time.Second).UnixMilli() candidates, err := s.store.InactiveNonWinnerClients(ctx, cutoff) if err != nil { return profileCleanupPreview{}, nil, err } protectedWinners, err := s.store.OldWinnerCount(ctx, cutoff) if err != nil { return profileCleanupPreview{}, nil, err } protectedPlace, err := s.store.OldPlaceParticipantCount(ctx, cutoff) if err != nil { return profileCleanupPreview{}, nil, err } ids := make([]string, 0, len(candidates)) out := profileCleanupPreview{CutoffMS: cutoff, InactiveForSeconds: inactiveForSeconds, ProtectedWinners: protectedWinners, ProtectedPlace: protectedPlace} for _, c := range candidates { if s.runtime.IsConnected(c.ClientID) { out.ProtectedConnected++ continue } ids = append(ids, c.ClientID) if out.OldestEligibleMS == 0 || c.LastSeen < out.OldestEligibleMS { out.OldestEligibleMS = c.LastSeen } if c.LastSeen > out.NewestEligibleMS { out.NewestEligibleMS = c.LastSeen } } out.Eligible = len(ids) return out, ids, nil } func (s *Server) adminProfileCleanupPreview(w http.ResponseWriter, r *http.Request) { seconds, err := profileCleanupDurationSeconds(r.URL.Query().Get("inactive_for_seconds")) if err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } preview, _, err := s.profileCleanupPreview(r.Context(), seconds) if err != nil { jsonOut(w, 500, map[string]string{"error": "profile cleanup preview failed: " + err.Error()}) return } jsonOut(w, 200, preview) } func (s *Server) adminProfileCleanup(w http.ResponseWriter, r *http.Request) { var in struct { InactiveForSeconds int64 `json:"inactive_for_seconds"` } if err := decode(r, &in); err != nil { jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } seconds, err := profileCleanupDurationSeconds(strconv.FormatInt(in.InactiveForSeconds, 10)) if err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } preview, ids, err := s.profileCleanupPreview(r.Context(), seconds) if err != nil { jsonOut(w, 500, map[string]string{"error": "profile cleanup check failed: " + err.Error()}) return } deleted, err := s.store.DeleteInactiveNonWinnerClients(r.Context(), preview.CutoffMS, ids) if err != nil { jsonOut(w, 500, map[string]string{"error": "profile cleanup failed: " + err.Error()}) return } for _, id := range deleted { s.runtime.ForgetClient(id) } jsonOut(w, 200, map[string]any{ "deleted": len(deleted), "eligible_before": preview.Eligible, "protected_winners": preview.ProtectedWinners, "protected_place": preview.ProtectedPlace, "protected_connected": preview.ProtectedConnected, "cutoff_ms": preview.CutoffMS, "inactive_for_seconds": seconds, }) } func (s *Server) adminOverview(w http.ResponseWriter, r *http.Request) { var clients, activeTasks, completedTasks, guesses, artifacts int64 _ = s.store.DB.QueryRowContext(r.Context(), `SELECT count(*) FROM clients`).Scan(&clients) _ = s.store.DB.QueryRowContext(r.Context(), `SELECT count(*) FROM tasks WHERE status='active'`).Scan(&activeTasks) _ = s.store.DB.QueryRowContext(r.Context(), `SELECT count(*) FROM tasks WHERE status='completed'`).Scan(&completedTasks) _ = s.store.DB.QueryRowContext(r.Context(), `SELECT COALESCE(sum(guess_count),0) FROM task_points`).Scan(&guesses) _ = s.store.DB.QueryRowContext(r.Context(), `SELECT count(*) FROM tasks WHERE artifact_status='ready'`).Scan(&artifacts) connected := s.runtime.ConnectedCount() jsonOut(w, 200, map[string]int64{"clients": clients, "connected": connected, "active_tasks": activeTasks, "completed_tasks": completedTasks, "guesses": guesses, "artifacts_ready": artifacts}) } func (s *Server) adminPerformance(w http.ResponseWriter, r *http.Request) { rm := s.runtime.Metrics() wm := s.hub.Metrics() pm := wsx.Metrics{} if s.placeHub != nil { pm = s.placeHub.Metrics() } var ms gort.MemStats gort.ReadMemStats(&ms) jsonOut(w, 200, map[string]any{ "runtime": rm, "websocket": wm, "place_websocket": pm, "process": map[string]any{ "goroutines": gort.NumGoroutine(), "heap_bytes": ms.HeapAlloc, "heap_objects": ms.HeapObjects, "gc_cycles": ms.NumGC, }, }) } func (s *Server) adminSettingsGet(w http.ResponseWriter, r *http.Request) { jsonOut(w, 200, s.settings.Get()) } func (s *Server) adminSettingsPut(w http.ResponseWriter, r *http.Request) { var v settings.Runtime if decode(r, &v) != nil { jsonOut(w, 400, map[string]string{"error": "bad json"}) return } if err := settings.Validate(v); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } current := s.settings.Get() if v.PlaceWidth < current.PlaceWidth || v.PlaceHeight < current.PlaceHeight { outside, err := s.store.PlacePixelsOutside(r.Context(), v.PlaceWidth, v.PlaceHeight) if err != nil { jsonOut(w, 500, map[string]string{"error": "place canvas bounds check failed"}) return } if outside > 0 { jsonOut(w, 400, map[string]string{"error": fmt.Sprintf("place canvas cannot shrink: %d existing pixels would fall outside the new dimensions", outside)}) return } } if err := s.settings.Update(r.Context(), v); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, v) } func (s *Server) adminTasks(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) ts, err := s.store.AdminTasks(r.Context(), r.URL.Query().Get("status"), r.URL.Query().Get("q"), limit) if err != nil { jsonOut(w, 500, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, ts) } func (s *Server) adminClients(w http.ResponseWriter, r *http.Request) { q := strings.TrimSpace(r.URL.Query().Get("q")) limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) if limit < 1 || limit > 500 { limit = 200 } rows, err := s.store.DB.QueryContext(r.Context(), `SELECT c.id,c.last_seen, (SELECT count(*) FROM tasks t WHERE t.winner_client_id=c.id AND t.status='completed') AS wins, (SELECT count(*) FROM tasks t WHERE COALESCE(t.artifact_owner_client_id,t.winner_client_id)=c.id AND t.status='completed' AND t.artifact_status='ready') AS nft_count, COALESCE((SELECT d.owner_client_id FROM identity_delegations d WHERE d.worker_client_id=c.id),c.id) AS place_owner_client_id, COALESCE((SELECT pw.balance_milli FROM place_wallets pw WHERE pw.client_id=COALESCE((SELECT d2.owner_client_id FROM identity_delegations d2 WHERE d2.worker_client_id=c.id),c.id)),0) AS place_balance_milli FROM clients c WHERE (?='' OR lower(c.id) LIKE '%'||lower(?)||'%') ORDER BY c.last_seen DESC LIMIT ?`, q, q, limit) if err != nil { jsonOut(w, 500, map[string]string{"error": err.Error()}) return } defer rows.Close() out := make([]map[string]any, 0) for rows.Next() { var id string var lastSeen int64 var wins, nfts int var placeOwner string var placeBalanceMilli int64 if err := rows.Scan(&id, &lastSeen, &wins, &nfts, &placeOwner, &placeBalanceMilli); err != nil { jsonOut(w, 500, map[string]string{"error": err.Error()}) return } out = append(out, map[string]any{"client_id": id, "last_seen": lastSeen, "wins": wins, "nft_count": nfts, "place_owner_client_id": placeOwner, "place_balance_points": float64(placeBalanceMilli) / 1000.0}) } jsonOut(w, 200, out) } func (s *Server) adminArtifactTransfer(w http.ResponseWriter, r *http.Request) { var in struct { TargetClientID string `json:"target_client_id"` Reason string `json:"reason"` } if err := decode(r, &in); err != nil || strings.TrimSpace(in.TargetClientID) == "" { jsonOut(w, 400, map[string]string{"error": "target_client_id required"}) return } taskID := chi.URLParam(r, "id") from, err := s.store.TransferArtifact(r.Context(), taskID, in.TargetClientID, in.Reason) if err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } _ = s.hub.Publish(r.Context(), wsx.Event{Type: "task_changed", TaskID: taskID, Data: map[string]string{"action": "artifact_transfer"}}) jsonOut(w, 200, map[string]any{"ok": true, "from_client_id": from, "to_client_id": strings.TrimSpace(in.TargetClientID)}) } func (s *Server) adminArtifactDrop(w http.ResponseWriter, r *http.Request) { var in struct { TargetClientID string `json:"target_client_id"` TemplateTaskID string `json:"template_task_id"` Rarity string `json:"rarity"` Count int `json:"count"` } if err := decode(r, &in); err != nil || strings.TrimSpace(in.TargetClientID) == "" { jsonOut(w, 400, map[string]string{"error": "target_client_id required"}) return } if in.Count == 0 { in.Count = 1 } if in.Count < 1 || in.Count > 20 { jsonOut(w, 400, map[string]string{"error": "count must be 1..20"}) return } created := make([]string, 0, in.Count) for i := 0; i < in.Count; i++ { t, err := s.store.CreateAdminArtifactDrop(r.Context(), in.TargetClientID, in.TemplateTaskID, in.Rarity, s.settings.Get().TaskRangeBits) if err != nil { jsonOut(w, 400, map[string]any{"error": err.Error(), "created_task_ids": created}) return } created = append(created, t.ID) } if s.artifactWorker != nil { s.artifactWorker.Notify() } log.Printf("admin artifact drop queued: owner=%s count=%d rarity=%q tasks=%s", strings.TrimSpace(in.TargetClientID), len(created), strings.TrimSpace(in.Rarity), strings.Join(created, ",")) jsonOut(w, 201, map[string]any{"ok": true, "created_task_ids": created, "artifact_status": "pending", "rarity": strings.TrimSpace(in.Rarity)}) } func (s *Server) adminArtifactStatus(w http.ResponseWriter, r *http.Request) { ids := append([]string(nil), r.URL.Query()["id"]...) if raw := strings.TrimSpace(r.URL.Query().Get("ids")); raw != "" { ids = append(ids, strings.Split(raw, ",")...) } if len(ids) == 0 { jsonOut(w, 400, map[string]string{"error": "at least one id is required"}) return } items, err := s.store.AdminArtifactQueueItems(r.Context(), ids) if err != nil { jsonOut(w, 500, map[string]string{"error": "artifact status lookup failed: " + err.Error()}) return } jsonOut(w, 200, items) } func (s *Server) adminTaskConfigPut(w http.ResponseWriter, r *http.Request) { id := chi.URLParam(r, "id") var in struct { DisplayName string `json:"display_name"` Description string `json:"description"` NFTPromptInstructions string `json:"nft_prompt_instructions"` NFTNegativePrompt string `json:"nft_negative_prompt"` RetireAfterCompletion bool `json:"retire_after_completion"` } if err := decode(r, &in); err != nil { jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } if err := s.store.UpdateTaskConfig(r.Context(), id, in.DisplayName, in.Description, in.NFTPromptInstructions, in.NFTNegativePrompt); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } if err := s.store.SetTaskRetireAfterCompletion(r.Context(), id, in.RetireAfterCompletion); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } _ = s.hub.Publish(r.Context(), wsx.Event{Type: "task_changed", TaskID: id, Data: map[string]string{"action": "config"}}) jsonOut(w, 200, map[string]bool{"ok": true}) } func (s *Server) adminPoints(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) ps, err := s.store.Points(r.Context(), chi.URLParam(r, "id"), limit) if err != nil { jsonOut(w, 500, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, ps) } func (s *Server) adminTaskActions(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) actions, err := s.store.TaskActions(r.Context(), chi.URLParam(r, "id"), limit) if err != nil { jsonOut(w, 500, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, actions) } type scheduleActionRequest struct { ActionType string `json:"action_type"` Payload json.RawMessage `json:"payload"` ExecuteAt *time.Time `json:"execute_at"` } func validateAction(actionType string, raw json.RawMessage) error { actionType = strings.ToLower(strings.TrimSpace(actionType)) switch actionType { case "set_range_bits": var p struct { Bits int `json:"bits"` Mode string `json:"mode"` } if len(raw) == 0 || json.Unmarshal(raw, &p) != nil { return errors.New("set_range_bits requires payload {bits,mode}") } if p.Bits < 8 || p.Bits > 128 { return errors.New("bits must be 8..128") } if p.Mode != "preserve" && p.Mode != "reroll" { return errors.New("mode must be preserve or reroll") } case "set_intervals": var p struct { Server int `json:"server_min_interval_sec"` Client int `json:"client_submit_interval_sec"` } if len(raw) == 0 || json.Unmarshal(raw, &p) != nil { return errors.New("set_intervals requires interval payload") } if p.Server < 1 || p.Server > 3600 || p.Client <= p.Server || p.Client > 7200 { return errors.New("intervals require server 1..3600 and client > server <=7200") } case "pause", "resume", "reroll", "clear_intervals", "close", "regenerate_artifact": // no payload required default: return fmt.Errorf("unsupported action_type %q", actionType) } return nil } func (s *Server) adminScheduleAction(w http.ResponseWriter, r *http.Request) { id := chi.URLParam(r, "id") var in scheduleActionRequest if err := decode(r, &in); err != nil { jsonOut(w, 400, map[string]string{"error": "bad json: " + err.Error()}) return } in.ActionType = strings.ToLower(strings.TrimSpace(in.ActionType)) if err := validateAction(in.ActionType, in.Payload); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } executeAt := time.Now().UTC() if in.ExecuteAt != nil { executeAt = in.ExecuteAt.UTC() } if executeAt.Before(time.Now().Add(-2 * time.Minute)) { jsonOut(w, 400, map[string]string{"error": "execute_at is in the past"}) return } if executeAt.After(time.Now().Add(366 * 24 * time.Hour)) { jsonOut(w, 400, map[string]string{"error": "execute_at is more than one year away"}) return } payload := any(map[string]any{}) if len(in.Payload) > 0 { var x any if err := json.Unmarshal(in.Payload, &x); err != nil { jsonOut(w, 400, map[string]string{"error": "invalid payload"}) return } payload = x } a, err := s.store.ScheduleTaskAction(r.Context(), id, in.ActionType, payload, executeAt) if err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } // "Run now" feels immediate in the admin UI while still going through the // same persisted action/audit path as scheduled changes. if !executeAt.After(time.Now().Add(1500 * time.Millisecond)) { s.runTaskAction(r.Context(), a) actions, _ := s.store.TaskActions(r.Context(), id, 1) if len(actions) > 0 { a = actions[0] } } jsonOut(w, 200, a) } func (s *Server) adminCancelAction(w http.ResponseWriter, r *http.Request) { if err := s.store.CancelTaskAction(r.Context(), chi.URLParam(r, "id")); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, map[string]bool{"ok": true}) } func (s *Server) adminArtifactProviders(w http.ResponseWriter, r *http.Request) { cfg := s.settings.Get() anchorPath := filepath.Join(s.artifactDir, "_collection", "character_anchor.png") anchorReady := false if st, err := os.Stat(anchorPath); err == nil && st.Size() > 1024 { anchorReady = true } jsonOut(w, 200, map[string]any{ "current": cfg.ArtifactProvider, "preset": cfg.ArtifactPreset, "model": cfg.ArtifactModel, "character_anchor": anchorReady, "providers": map[string]bool{ "local": true, "openai": strings.TrimSpace(os.Getenv("OPENAI_API_KEY")) != "", "comfyui": strings.TrimSpace(os.Getenv("COMFYUI_URL")) != "" && strings.TrimSpace(os.Getenv("COMFYUI_WORKFLOW_PATH")) != "", "a1111": strings.TrimSpace(os.Getenv("A1111_URL")) != "", }, }) } func (s *Server) executeTaskAction(ctx context.Context, a data.TaskAction) error { switch a.ActionType { case "set_range_bits": var p struct { Bits int `json:"bits"` Mode string `json:"mode"` } if err := json.Unmarshal(a.Payload, &p); err != nil { return err } return s.store.SetTaskRangeBits(ctx, a.TaskID, p.Bits, p.Mode) case "set_intervals": var p struct { Server int `json:"server_min_interval_sec"` Client int `json:"client_submit_interval_sec"` } if err := json.Unmarshal(a.Payload, &p); err != nil { return err } return s.store.SetTaskIntervals(ctx, a.TaskID, p.Server, p.Client) case "clear_intervals": return s.store.ClearTaskIntervals(ctx, a.TaskID) case "pause": return s.store.SetTaskPaused(ctx, a.TaskID, true) case "resume": return s.store.SetTaskPaused(ctx, a.TaskID, false) case "reroll": return s.store.RerollTask(ctx, a.TaskID) case "close": return s.store.CloseTask(ctx, a.TaskID) case "regenerate_artifact": return s.store.QueueArtifact(ctx, a.TaskID) default: return fmt.Errorf("unsupported action %q", a.ActionType) } } func (s *Server) runTaskAction(ctx context.Context, a data.TaskAction) { if !s.store.StartTaskAction(ctx, a.ID) { return } err := s.executeTaskAction(ctx, a) s.store.FinishTaskAction(ctx, a.ID, err) if err != nil { log.Printf("task action %s (%s): %v", a.ID, a.ActionType, err) return } if a.ActionType == "close" { dataOut := map[string]string{"reason": "scheduled_close"} if successor, succErr := s.store.EnsureSuccessorTask(ctx, a.TaskID, s.settings.Get().TaskRangeBits); succErr == nil { dataOut["successor_task_id"] = successor.ID s.runtime.ReplaceTaskSelection(a.TaskID, successor.ID) } else if !errors.Is(succErr, data.ErrTaskRetiring) { log.Printf("successor for %s: %v", a.TaskID, succErr) } else { dataOut["retired"] = "true" } _ = s.hub.Publish(ctx, wsx.Event{Type: "task_completed", TaskID: a.TaskID, Data: dataOut}) _ = s.store.EnsureActiveTasks(ctx, s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits) return } _ = s.hub.Publish(ctx, wsx.Event{Type: "task_changed", TaskID: a.TaskID, Data: map[string]string{"action": a.ActionType}}) } func (s *Server) runDueActions(ctx context.Context) { actions, err := s.store.DueTaskActions(ctx, 20) if err != nil { log.Printf("task actions: %v", err) return } for _, a := range actions { s.runTaskAction(ctx, a) } } func (s *Server) adminCloseTask(w http.ResponseWriter, r *http.Request) { id := chi.URLParam(r, "id") if err := s.store.CloseTask(r.Context(), id); err != nil { jsonOut(w, 400, map[string]string{"error": err.Error()}) return } dataOut := map[string]string{"reason": "closed_by_admin"} if successor, succErr := s.store.EnsureSuccessorTask(r.Context(), id, s.settings.Get().TaskRangeBits); succErr == nil { dataOut["successor_task_id"] = successor.ID s.runtime.ReplaceTaskSelection(id, successor.ID) } else if !errors.Is(succErr, data.ErrTaskRetiring) { log.Printf("successor for %s: %v", id, succErr) } else { dataOut["retired"] = "true" } _ = s.hub.Publish(r.Context(), wsx.Event{Type: "task_completed", TaskID: id, Data: dataOut}) _ = s.store.EnsureActiveTasks(r.Context(), s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits) jsonOut(w, 200, map[string]bool{"ok": true}) } func (s *Server) adminEnsure(w http.ResponseWriter, r *http.Request) { err := s.store.EnsureActiveTasks(r.Context(), s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits) if err != nil { jsonOut(w, 500, map[string]string{"error": err.Error()}) return } jsonOut(w, 200, map[string]bool{"ok": true}) } func websocketAuthToken(r *http.Request) string { if h := strings.TrimSpace(r.Header.Get("Authorization")); strings.HasPrefix(h, "Bearer ") { return strings.TrimSpace(strings.TrimPrefix(h, "Bearer ")) } for _, part := range strings.Split(r.Header.Get("Sec-WebSocket-Protocol"), ",") { part = strings.TrimSpace(part) if strings.HasPrefix(part, "nh-auth.") { return strings.TrimPrefix(part, "nh-auth.") } } if strings.EqualFold(strings.TrimSpace(os.Getenv("WS_ALLOW_QUERY_TOKEN")), "1") { return r.URL.Query().Get("token") } return "" } func acquireWSCap(counter *atomic.Int64, max int64) bool { if max <= 0 { counter.Add(1) return true } for { cur := counter.Load() if cur >= max { return false } if counter.CompareAndSwap(cur, cur+1) { return true } } } func (s *Server) ws(w http.ResponseWriter, r *http.Request) { tok := websocketAuthToken(r) c, err := s.auth.Parse(tok) if err != nil || c.Role != "user" { http.Error(w, "unauthorized", 401) return } if !s.store.ClientExists(r.Context(), c.ClientID) { http.Error(w, "identity not registered", http.StatusUnauthorized) return } // WebSocket use counts as recent profile activity. This also makes a race // with the admin cleanup safe: a newly connecting identity no longer // matches an old last_seen cutoff. s.store.TouchClient(r.Context(), c.ClientID) t, err := s.store.TaskForClient(r.Context(), c.ClientID) if err != nil { http.Error(w, "no task", 503) return } s.runtime.SetTaskSelection(c.ClientID, t.ID) if !acquireWSCap(&s.userWSCount, s.maxUserWS) { http.Error(w, "websocket capacity reached", http.StatusServiceUnavailable) return } defer s.userWSCount.Add(-1) leaseID, err := s.runtime.AcquirePresence(c.ClientID, c.SessionID) if err != nil { http.Error(w, "identity already connected", 409) return } conn, err := s.upgrader.Upgrade(w, r, nil) if err != nil { s.runtime.ReleasePresence(c.ClientID, c.SessionID, leaseID) return } cl := wsx.NewClient(conn, t.ID, c.ClientID, false) s.hub.Add(cl) defer func() { s.hub.Remove(cl) s.runtime.ReleasePresence(c.ClientID, c.SessionID, leaseID) // Mark the disconnect time as last activity. A client that stayed online // for days therefore starts its inactivity window only after disconnect. s.store.TouchClient(context.Background(), c.ClientID) }() // A map point is durable only once per client/task. Subsequent losing guesses // stay in memory; improvements are checkpointed by the guess handler. if p, err := s.store.EnsurePoint(r.Context(), t.ID, c.ClientID); err == nil { p.Score = round(p.Score, s.settings.Get().PublicScorePrecision) s.hub.PublishPoint(t.ID, c.ClientID, p) } snap, _ := s.store.LoadGuessState(r.Context(), t.ID, c.ClientID) s.runtime.InitGuess(t, c.ClientID, rtx.GuessState{NextSeq: snap.NextSeq, LastGuess: snap.LastGuess, BestScore: snap.BestScore, GuessCount: snap.GuessCount}) maxNodes := s.settings.Get().DefaultMaxNodes if q, _ := strconv.Atoi(r.URL.Query().Get("max_nodes")); q > 0 { maxNodes = q } // Give the browser a bounded overscan set for local LOD, not the entire task. limit := maxNodes * 3 if limit < 300 { limit = 300 } if limit > 10000 { limit = 10000 } ps, _ := s.store.PointsForClient(r.Context(), t.ID, c.ClientID, limit) for i := range ps { ps[i].Score = round(ps[i].Score, s.settings.Get().PublicScorePrecision) } cl.Enqueue(wsx.Event{Type: "snapshot", TaskID: t.ID, Data: ps}) conn.SetReadLimit(4 << 10) _ = conn.SetReadDeadline(time.Now().Add(90 * time.Second)) conn.SetPongHandler(func(string) error { return conn.SetReadDeadline(time.Now().Add(90 * time.Second)) }) ping := time.NewTicker(30 * time.Second) defer ping.Stop() done := make(chan struct{}) go func() { defer close(done) for { if _, _, err := conn.ReadMessage(); err != nil { return } } }() for { select { case <-done: return case <-r.Context().Done(): return case <-ping.C: // Real WebSocket control ping. The browser answers with pong // automatically and the PongHandler extends the 90s read deadline. if err := cl.Ping(); err != nil { return } } } } func (s *Server) Scheduler(ctx context.Context) { t := time.NewTicker(time.Second) defer t.Stop() maintenance := 0 for { select { case <-ctx.Done(): return case <-t.C: s.runDueActions(ctx) maintenance++ if maintenance%5 == 0 { s.dispatchHostedCreditEvents(ctx) if err := s.store.EnsureActiveTasks(ctx, s.settings.Get().ActiveTaskCount, s.settings.Get().TaskRangeBits); err != nil { log.Printf("scheduler: %v", err) } } } } }