Files
jbergner 8c67c7a7fa
All checks were successful
release-tag / release-image (push) Successful in 10m52s
Update 1.5.0
2026-08-27 07:51:39 +02:00

912 lines
36 KiB
Go

package httpapi
import (
"context"
"crypto/subtle"
"embed"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net/http"
"os"
"strconv"
"strings"
"sync/atomic"
"time"
"neuroforge/internal/brain"
"neuroforge/internal/core"
"neuroforge/internal/cost"
"neuroforge/internal/provider"
"neuroforge/internal/store"
)
//go:embed index.html
var webFS embed.FS
type Server struct {
store *store.Store
brain *brain.Engine
router *provider.Router
cost *cost.Manager
mux *http.ServeMux
metrics *metricsRegistry
inflight atomic.Int64
readinessOllamaLive bool
}
func New(s *store.Store, b *brain.Engine, r *provider.Router, c *cost.Manager) *Server {
x := &Server{store: s, brain: b, router: r, cost: c, mux: http.NewServeMux(), metrics: newMetricsRegistry()}
x.routes()
return x
}
func (s *Server) SetReadinessOllamaLive(enabled bool) { s.readinessOllamaLive = enabled }
func (s *Server) Handler() http.Handler {
var h http.Handler = s.mux
h = s.requestLimits(h)
h = s.securityHeaders(h)
h = s.logging(h)
return h
}
func (s *Server) routes() {
s.mux.HandleFunc("GET /", s.index)
s.mux.HandleFunc("GET /admin", s.index)
s.mux.HandleFunc("GET /metrics", s.metricsEndpoint)
s.mux.HandleFunc("GET /healthz", s.livez)
s.mux.HandleFunc("GET /livez", s.livez)
s.mux.HandleFunc("GET /readyz", s.readyz)
s.mux.HandleFunc("GET /version", func(w http.ResponseWriter, r *http.Request) { s.json(w, 200, map[string]any{"version": "0.8.2"}) })
s.mux.Handle("POST /api/v1/chat", s.appAuth(http.HandlerFunc(s.chat)))
s.mux.Handle("POST /api/v1/learn", s.appAuth(http.HandlerFunc(s.learn)))
s.mux.Handle("POST /api/v1/search", s.appAuth(http.HandlerFunc(s.search)))
s.mux.Handle("POST /api/v1/search/vector", s.appAuth(http.HandlerFunc(s.searchVector)))
s.mux.Handle("POST /api/v1/memory/import", s.appAuth(http.HandlerFunc(s.importMemory)))
s.mux.Handle("POST /api/v1/feedback", s.appAuth(http.HandlerFunc(s.feedback)))
s.mux.Handle("GET /api/v1/stats", s.controlReadAuth(http.HandlerFunc(s.stats)))
s.mux.Handle("GET /api/v1/goals", s.appAuth(http.HandlerFunc(s.goalsList)))
s.mux.Handle("POST /api/v1/goals", s.appAuth(http.HandlerFunc(s.goalsCreate)))
s.mux.Handle("GET /api/v1/goals/{id}", s.appAuth(http.HandlerFunc(s.goalsGet)))
s.mux.Handle("PUT /api/v1/goals/{id}", s.appAuth(http.HandlerFunc(s.goalsPut)))
s.mux.Handle("DELETE /api/v1/goals/{id}", s.appAuth(http.HandlerFunc(s.goalsDelete)))
s.mux.Handle("POST /api/v1/goals/{id}/pause", s.appAuth(http.HandlerFunc(s.goalPause)))
s.mux.Handle("POST /api/v1/goals/{id}/resume", s.appAuth(http.HandlerFunc(s.goalResume)))
s.mux.Handle("POST /api/v1/goals/{id}/cycle", s.appAuth(http.HandlerFunc(s.goalCycle)))
s.mux.Handle("GET /api/v1/goals/{id}/research/live", s.appAuth(http.HandlerFunc(s.goalResearchLive)))
s.mux.Handle("GET /api/v1/goals/{id}/research/history", s.appAuth(http.HandlerFunc(s.goalResearchHistory)))
s.mux.Handle("GET /api/v1/learning-cycles", s.appAuth(http.HandlerFunc(s.learningCycles)))
s.mux.Handle("GET /api/v1/conflicts", s.appAuth(http.HandlerFunc(s.conflicts)))
s.mux.Handle("POST /api/v1/ingest/text", s.appAuth(http.HandlerFunc(s.ingestText)))
s.mux.Handle("POST /api/v1/ingest/document", s.appAuth(http.HandlerFunc(s.ingestDocument)))
s.mux.Handle("GET /api/v1/sources", s.appAuth(http.HandlerFunc(s.sourcesList)))
s.mux.Handle("GET /api/v1/sources/{id}", s.appAuth(http.HandlerFunc(s.sourceGet)))
s.mux.Handle("POST /api/v1/research", s.appAuth(http.HandlerFunc(s.researchSearch)))
s.mux.Handle("POST /api/v1/integrations/knowledge/upsert", s.integrationAuth(http.HandlerFunc(s.integrationKnowledgeUpsert)))
s.mux.Handle("DELETE /api/v1/integrations/knowledge/{namespace}/{document_id}", s.integrationAuth(http.HandlerFunc(s.integrationKnowledgeDelete)))
s.mux.Handle("POST /api/v1/integrations/knowledge/search", s.integrationAuth(http.HandlerFunc(s.integrationKnowledgeSearch)))
s.mux.Handle("POST /api/v1/integrations/events", s.integrationAuth(http.HandlerFunc(s.integrationEvent)))
s.mux.Handle("POST /api/v1/integrations/outcomes", s.integrationAuth(http.HandlerFunc(s.integrationValidatedOutcome)))
s.mux.Handle("POST /api/v1/integrations/outcomes/search", s.integrationAuth(http.HandlerFunc(s.integrationValidatedOutcomeSearch)))
s.mux.Handle("GET /api/v1/integrations/graph/research", s.controlReadAuth(http.HandlerFunc(s.integrationResearchGraph)))
s.mux.Handle("GET /api/v1/integrations/graph/brain", s.controlReadAuth(http.HandlerFunc(s.integrationBrainGraph)))
s.mux.Handle("POST /internal/v1/cluster/request-vote", s.clusterAuth(http.HandlerFunc(s.clusterRequestVote)))
s.mux.Handle("POST /internal/v1/cluster/heartbeat", s.clusterAuth(http.HandlerFunc(s.clusterHeartbeat)))
s.mux.Handle("POST /internal/v1/cluster/prepare", s.clusterAuth(http.HandlerFunc(s.clusterPrepare)))
s.mux.Handle("POST /internal/v1/cluster/commit", s.clusterAuth(http.HandlerFunc(s.clusterCommit)))
s.mux.Handle("POST /internal/v1/cluster/abort", s.clusterAuth(http.HandlerFunc(s.clusterAbort)))
s.mux.Handle("POST /internal/v1/cluster/propose/memory", s.clusterAuth(http.HandlerFunc(s.clusterProposeMemory)))
s.mux.Handle("GET /internal/v1/cluster/decision/{id}", s.clusterAuth(http.HandlerFunc(s.clusterDecision)))
s.mux.Handle("GET /internal/v1/cluster/status", s.clusterAuth(http.HandlerFunc(s.clusterStatus)))
s.mux.Handle("POST /api/v1/worker/claim", s.workerAuth(http.HandlerFunc(s.workerClaim)))
s.mux.Handle("POST /api/v1/worker/complete", s.workerAuth(http.HandlerFunc(s.workerComplete)))
s.mux.Handle("GET /admin/api/status", s.adminAuth(http.HandlerFunc(s.adminStatus)))
s.mux.Handle("GET /admin/api/config", s.adminAuth(http.HandlerFunc(s.adminGetConfig)))
s.mux.Handle("PUT /admin/api/config", s.adminAuth(http.HandlerFunc(s.adminPutConfig)))
s.mux.Handle("GET /admin/api/model-routing", s.adminAuth(http.HandlerFunc(s.adminGetModelRouting)))
s.mux.Handle("PUT /admin/api/model-routing", s.adminAuth(http.HandlerFunc(s.adminPutModelRouting)))
s.mux.Handle("GET /admin/api/secrets/status", s.adminAuth(http.HandlerFunc(s.adminSecretsStatus)))
s.mux.Handle("GET /admin/api/secrets", s.adminAuth(http.HandlerFunc(s.adminGetSecrets)))
s.mux.Handle("PUT /admin/api/secrets", s.adminAuth(http.HandlerFunc(s.adminPutSecrets)))
s.mux.Handle("POST /admin/api/provider-health", s.adminAuth(http.HandlerFunc(s.adminProviderHealth)))
s.mux.Handle("GET /admin/api/memories", s.adminAuth(http.HandlerFunc(s.adminMemories)))
s.mux.Handle("DELETE /admin/api/memories/{id}", s.adminAuth(http.HandlerFunc(s.adminDeleteMemory)))
s.mux.Handle("GET /admin/api/synapses", s.adminAuth(http.HandlerFunc(s.adminSynapses)))
s.mux.Handle("GET /admin/api/usage", s.adminAuth(http.HandlerFunc(s.adminUsage)))
s.mux.Handle("GET /admin/api/export", s.adminAuth(http.HandlerFunc(s.adminExport)))
s.mux.Handle("POST /admin/api/consolidate", s.adminAuth(http.HandlerFunc(s.adminConsolidate)))
s.mux.Handle("POST /admin/api/retention", s.adminAuth(http.HandlerFunc(s.adminRetention)))
s.mux.Handle("POST /admin/api/autonomy", s.adminAuth(http.HandlerFunc(s.adminAutonomy)))
s.mux.Handle("POST /admin/api/rebalance", s.adminAuth(http.HandlerFunc(s.adminRebalance)))
s.mux.Handle("POST /admin/api/checkpoint", s.adminAuth(http.HandlerFunc(s.adminCheckpoint)))
s.mux.Handle("GET /admin/api/wal", s.adminAuth(http.HandlerFunc(s.adminWAL)))
s.mux.Handle("GET /admin/api/storage", s.adminAuth(http.HandlerFunc(s.adminStorageStatus)))
s.mux.Handle("POST /admin/api/storage/compact", s.adminAuth(http.HandlerFunc(s.adminCompactSegments)))
s.mux.Handle("POST /admin/api/storage/tier", s.adminAuth(http.HandlerFunc(s.adminTierStorage)))
s.mux.Handle("POST /admin/api/index/merge", s.adminAuth(http.HandlerFunc(s.adminMergeIndex)))
s.mux.Handle("GET /admin/api/index/disk", s.adminAuth(http.HandlerFunc(s.adminDiskANNStatus)))
s.mux.Handle("POST /admin/api/index/disk/rebuild", s.adminAuth(http.HandlerFunc(s.adminDiskANNBuild)))
s.mux.Handle("GET /admin/api/cluster", s.adminAuth(http.HandlerFunc(s.clusterStatus)))
s.mux.Handle("POST /admin/api/cluster/repair", s.adminAuth(http.HandlerFunc(s.adminClusterRepair)))
s.mux.Handle("POST /admin/api/conflicts/resolve", s.adminAuth(http.HandlerFunc(s.adminResolveConflict)))
s.mux.Handle("GET /admin/api/knowledge/summary", s.adminAuth(http.HandlerFunc(s.adminKnowledgeSummary)))
s.mux.Handle("GET /admin/api/knowledge/memories", s.adminAuth(http.HandlerFunc(s.adminKnowledgeMemories)))
s.mux.Handle("GET /admin/api/knowledge/memory/{id}", s.adminAuth(http.HandlerFunc(s.adminKnowledgeMemory)))
s.mux.Handle("GET /admin/api/knowledge/graph", s.adminAuth(http.HandlerFunc(s.adminKnowledgeGraph)))
s.mux.Handle("GET /admin/api/knowledge/events", s.adminAuth(http.HandlerFunc(s.adminKnowledgeEvents)))
s.mux.Handle("POST /admin/api/knowledge/search", s.adminAuth(http.HandlerFunc(s.adminKnowledgeSearch)))
s.mux.Handle("GET /admin/api/learning-policy", s.adminAuth(http.HandlerFunc(s.adminGetLearningPolicy)))
s.mux.Handle("PUT /admin/api/learning-policy", s.adminAuth(http.HandlerFunc(s.adminPutLearningPolicy)))
s.mux.Handle("GET /admin/api/research", s.adminAuth(http.HandlerFunc(s.adminResearchGet)))
s.mux.Handle("PUT /admin/api/research", s.adminAuth(http.HandlerFunc(s.adminResearchPut)))
s.mux.Handle("POST /admin/api/research/test", s.adminAuth(http.HandlerFunc(s.adminResearchTest)))
}
func (s *Server) index(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/" && r.URL.Path != "/admin" {
http.NotFound(w, r)
return
}
b, err := webFS.ReadFile("index.html")
if err != nil {
http.Error(w, err.Error(), 500)
return
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
w.Header().Set("Cache-Control", "no-store")
w.Header().Set("Pragma", "no-cache")
w.Write(b)
}
type statusWriter struct {
http.ResponseWriter
status int
bytes int64
}
func (w *statusWriter) Unwrap() http.ResponseWriter { return w.ResponseWriter }
func (w *statusWriter) WriteHeader(code int) {
if w.status != 0 {
return
}
w.status = code
w.ResponseWriter.WriteHeader(code)
}
func (w *statusWriter) Write(p []byte) (int, error) {
if w.status == 0 {
w.status = http.StatusOK
}
n, err := w.ResponseWriter.Write(p)
w.bytes += int64(n)
return n, err
}
func (s *Server) logging(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
start := time.Now()
sw := &statusWriter{ResponseWriter: w}
next.ServeHTTP(sw, r)
d := time.Since(start)
s.metrics.observeHTTP(r.Method, normalizeMetricRoute(r), sw.status, sw.bytes, d)
log.Printf("%s %s %d %s", r.Method, r.URL.Path, func() int {
if sw.status == 0 {
return http.StatusOK
}
return sw.status
}(), d.Round(time.Millisecond))
})
}
func secureEqual(a, b string) bool {
if len(a) == 0 || len(a) != len(b) {
return false
}
return subtle.ConstantTimeCompare([]byte(a), []byte(b)) == 1
}
func bearer(r *http.Request) string {
h := r.Header.Get("Authorization")
if strings.HasPrefix(strings.ToLower(h), "bearer ") {
return strings.TrimSpace(h[7:])
}
return ""
}
func (s *Server) appAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
cfg := s.store.Config()
if cfg.API.RequireKey {
sec := s.store.Secrets()
// The admin dashboard is already authenticated with the stronger admin
// credential. Allow it to call application endpoints directly so the
// browser never needs the App API key (which is masked by default in
// production). External applications still authenticate with Bearer.
adminOK := secureEqual(r.Header.Get("X-Admin-Token"), sec.AdminToken)
appOK := secureEqual(bearer(r), sec.AppAPIKey)
if !adminOK && !appOK {
s.err(w, 401, errors.New("invalid app API key or admin token"))
return
}
}
next.ServeHTTP(w, r)
})
}
func (s *Server) integrationAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
sec := s.store.Secrets()
adminOK := secureEqual(r.Header.Get("X-Admin-Token"), sec.AdminToken)
integrationOK := secureEqual(bearer(r), sec.IntegrationToken)
if !adminOK && !integrationOK {
s.err(w, http.StatusUnauthorized, errors.New("invalid integration token or admin token"))
return
}
next.ServeHTTP(w, r)
})
}
func (s *Server) controlReadAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
sec := s.store.Secrets()
adminOK := secureEqual(r.Header.Get("X-Admin-Token"), sec.AdminToken)
controlOK := secureEqual(bearer(r), sec.ControlReadToken)
appOK := secureEqual(bearer(r), sec.AppAPIKey)
if !adminOK && !controlOK && !appOK {
s.err(w, http.StatusUnauthorized, errors.New("invalid control/app read token or admin token"))
return
}
next.ServeHTTP(w, r)
})
}
func (s *Server) workerAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !secureEqual(bearer(r), s.store.Secrets().WorkerToken) {
s.err(w, 401, errors.New("invalid worker token"))
return
}
next.ServeHTTP(w, r)
})
}
func (s *Server) adminAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !secureEqual(r.Header.Get("X-Admin-Token"), s.store.Secrets().AdminToken) {
s.err(w, 401, errors.New("invalid admin token"))
return
}
next.ServeHTTP(w, r)
})
}
func (s *Server) clusterAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
sec := s.store.Secrets()
if sec.ClusterToken == "" || !secureEqual(r.Header.Get("X-Cluster-Token"), sec.ClusterToken) {
s.err(w, 401, errors.New("invalid cluster token"))
return
}
next.ServeHTTP(w, r)
})
}
func decode(r *http.Request, v any) error {
defer r.Body.Close()
d := json.NewDecoder(io.LimitReader(r.Body, 128<<20))
d.DisallowUnknownFields()
return d.Decode(v)
}
func (s *Server) json(w http.ResponseWriter, status int, v any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(v)
}
func (s *Server) err(w http.ResponseWriter, status int, err error) {
s.json(w, status, map[string]any{"error": err.Error()})
}
func (s *Server) chat(w http.ResponseWriter, r *http.Request) {
var q brain.ChatRequest
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
out, err := s.brain.Chat(r.Context(), q)
if err != nil {
s.err(w, 502, err)
return
}
s.json(w, 200, out)
}
func (s *Server) learn(w http.ResponseWriter, r *http.Request) {
var q brain.LearnRequest
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
m, err := s.brain.Learn(r.Context(), q)
if err != nil {
s.err(w, 502, err)
return
}
s.json(w, 201, m)
}
func (s *Server) search(w http.ResponseWriter, r *http.Request) {
var q struct {
Text string `json:"text"`
K int `json:"k"`
}
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
hits, err := s.brain.Search(r.Context(), q.Text, q.K)
if err != nil {
s.err(w, 502, err)
return
}
s.json(w, 200, hits)
}
func (s *Server) searchVector(w http.ResponseWriter, r *http.Request) {
var q struct {
Vector []float32 `json:"vector"`
K int `json:"k"`
MinSimilarity *float64 `json:"min_similarity,omitempty"`
GraphBonus *float64 `json:"graph_bonus,omitempty"`
}
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
if len(q.Vector) == 0 {
s.err(w, 400, errors.New("vector required"))
return
}
cfg := s.store.Config()
if q.K <= 0 {
q.K = cfg.Brain.RecallK
}
min := cfg.Brain.MinSimilarity
if q.MinSimilarity != nil {
min = *q.MinSimilarity
}
bonus := cfg.Brain.GraphBonus
if q.GraphBonus != nil {
bonus = *q.GraphBonus
}
// Intentionally local-only: shard federation is one hop and must not recurse.
hits := s.store.SearchVector(q.Vector, q.K, min, bonus)
for i := range hits {
if hits[i].Memory.ShardID == "" {
hits[i].Memory.ShardID = cfg.Sharding.LocalShardID
}
}
s.json(w, 200, hits)
}
func (s *Server) importMemory(w http.ResponseWriter, r *http.Request) {
var m core.Memory
if err := decode(r, &m); err != nil {
s.err(w, 400, err)
return
}
out, err := s.brain.ImportMemory(r.Context(), m)
if err != nil {
s.err(w, 400, err)
return
}
s.json(w, 201, out)
}
func (s *Server) feedback(w http.ResponseWriter, r *http.Request) {
var q brain.FeedbackRequest
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
if err := s.brain.Feedback(q); err != nil {
s.err(w, 400, err)
return
}
s.json(w, 200, map[string]bool{"ok": true})
}
func (s *Server) stats(w http.ResponseWriter, r *http.Request) {
s.json(w, 200, map[string]any{"stats": s.store.Stats(), "cost": s.cost.Totals()})
}
func (s *Server) workerClaim(w http.ResponseWriter, r *http.Request) {
var q struct {
WorkerID string `json:"worker_id"`
}
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
if q.WorkerID == "" {
s.err(w, 400, errors.New("worker_id required"))
return
}
lease := s.store.Config().Worker.LeaseSeconds
if lease < 10 {
lease = 120
}
j, err := s.store.ClaimJob(q.WorkerID, time.Duration(lease)*time.Second)
if err != nil {
s.err(w, 500, err)
return
}
if j == nil {
w.WriteHeader(http.StatusNoContent)
return
}
s.json(w, 200, j)
}
func (s *Server) workerComplete(w http.ResponseWriter, r *http.Request) {
var q struct {
WorkerID string `json:"worker_id"`
JobID string `json:"job_id"`
Result json.RawMessage `json:"result"`
Error string `json:"error"`
}
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
j, err := s.store.CompleteJob(q.JobID, q.WorkerID, q.Result, q.Error)
if err != nil {
s.err(w, 400, err)
return
}
if err := s.brain.ApplyJobResult(j); err != nil {
s.err(w, 500, err)
return
}
s.json(w, 200, map[string]bool{"ok": true})
}
func (s *Server) adminStatus(w http.ResponseWriter, r *http.Request) {
obs := s.store.ObservabilitySnapshot()
maint := s.store.MaintenanceStatus()
stats := map[string]any{
"revision": obs.Revision, "memories": obs.Memories, "synapses": obs.Synapses, "goals": obs.Goals, "learning_cycles": obs.LearningCycles,
"pending_jobs": obs.JobsQueued + obs.JobsClaimed, "hnsw_nodes": obs.HNSWNodes, "hnsw_dimensions": obs.HNSWDimensions,
"disk_pq_items": obs.DiskPQItems, "disk_pq_bytes": obs.DiskPQBytes, "index_mode": obs.IndexMode, "remote_shards": obs.RemoteShards, "maintenance": maint,
}
cfg := s.store.Config()
sec := s.store.Secrets()
providers := make([]map[string]any, 0, len(cfg.Ollama)+1)
for _, o := range cfg.Ollama {
providers = append(providers, map[string]any{"provider": "ollama", "id": o.ID, "name": o.Name, "model": o.ChatModel, "enabled": o.Enabled})
}
providers = append(providers, map[string]any{"provider": "openai", "model": cfg.OpenAI.ChatModel, "enabled": cfg.OpenAI.Enabled, "configured": sec.OpenAIAPIKey != ""})
tiering := map[string]any{
"hot_memories": obs.HotMemories, "cold_memories": obs.ColdMemories, "hot_bytes": obs.HotBytes, "tier_evictions_total": obs.TierEvictions,
"page_cache": map[string]any{"enabled": obs.PageCacheEnabled, "max_bytes": obs.PageCacheMaxBytes, "bytes": obs.PageCacheBytes, "entries": obs.PageCacheEntries, "hits": obs.PageCacheHits, "misses": obs.PageCacheMisses, "evictions": obs.PageCacheEvicts},
}
cluster := map[string]any{
"enabled": obs.ClusterEnabled, "node_id": obs.ClusterNodeID, "leader_id": obs.ClusterLeaderID, "role": obs.ClusterRole, "term": obs.ClusterTerm,
"last_index": obs.ClusterLastIndex, "commit_index": obs.ClusterCommitIndex, "peers": obs.ClusterPeers, "voters": obs.ClusterVoters, "quorum": obs.ClusterQuorum,
"replicated_log": obs.ClusterLog,
}
payload := map[string]any{
"stats": stats, "wal": s.store.WALStatus(),
"storage": map[string]any{"memory_segments": obs.Segments, "index_snapshot": map[string]any{"revision": obs.IndexSnapshotRevision, "deltas": obs.IndexDeltaCount, "segmented": cfg.Storage.IndexSegments.Enabled}, "tiering": tiering, "disk_ann": s.store.DiskANNStatus()},
"cluster": cluster, "cost": s.cost.Totals(), "providers": providers, "usage": s.store.RecentUsage(25),
"observability": obs, "runtime": currentRuntimeSnapshot(), "http": s.metrics.dashboardSnapshot(),
}
// Keep the original flat status fields for dashboard and external-client
// compatibility while retaining the richer nested stats object.
for k, v := range stats {
payload[k] = v
}
s.json(w, 200, payload)
}
func (s *Server) adminGetConfig(w http.ResponseWriter, r *http.Request) {
s.json(w, 200, s.store.Config())
}
func (s *Server) adminPutConfig(w http.ResponseWriter, r *http.Request) {
var c core.Config
if err := decode(r, &c); err != nil {
s.err(w, 400, err)
return
}
if err := s.store.ValidateConfig(c); err != nil {
s.err(w, 400, err)
return
}
if err := s.store.UpdateConfig(c); err != nil {
s.err(w, 500, err)
return
}
_ = s.store.AddKnowledgeEvent(core.KnowledgeEvent{Type: "admin.config_changed", Summary: "Runtime configuration updated", Reason: "PUT /admin/api/config", Actor: "admin"})
s.json(w, 200, c)
}
type modelRoutingLearning struct {
AutoRewardEnabled bool `json:"auto_reward_enabled"`
AutoRewardMode string `json:"auto_reward_mode"`
ConsolidationEnabled bool `json:"consolidation_enabled"`
ConsolidationUseLLM bool `json:"consolidation_use_llm"`
AutonomyEnabled bool `json:"autonomy_enabled"`
AutonomyUseLLM bool `json:"autonomy_use_llm"`
}
type modelRoutingSettings struct {
Routing core.RoutingConfig `json:"routing"`
Ollama []core.OllamaServer `json:"ollama"`
Learning modelRoutingLearning `json:"learning"`
}
type modelRoutingUpdate struct {
Routing *core.RoutingConfig `json:"routing,omitempty"`
Ollama *[]core.OllamaServer `json:"ollama,omitempty"`
Learning *modelRoutingLearning `json:"learning,omitempty"`
}
func modelRoutingFromConfig(c core.Config) modelRoutingSettings {
var out modelRoutingSettings
out.Routing = c.Routing
out.Ollama = append([]core.OllamaServer(nil), c.Ollama...)
out.Learning.AutoRewardEnabled = c.Brain.AutoReward.Enabled
out.Learning.AutoRewardMode = c.Brain.AutoReward.Mode
out.Learning.ConsolidationEnabled = c.Brain.Consolidation.Enabled
out.Learning.ConsolidationUseLLM = c.Brain.Consolidation.UseLLM
out.Learning.AutonomyEnabled = c.Autonomy.Enabled
out.Learning.AutonomyUseLLM = c.Autonomy.UseLLM
return out
}
func (s *Server) adminGetModelRouting(w http.ResponseWriter, r *http.Request) {
s.json(w, 200, modelRoutingFromConfig(s.store.Config()))
}
func (s *Server) adminPutModelRouting(w http.ResponseWriter, r *http.Request) {
var q modelRoutingUpdate
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
if q.Routing == nil && q.Ollama == nil && q.Learning == nil {
s.err(w, 400, errors.New("routing, ollama or learning is required"))
return
}
c := s.store.Config()
if q.Routing != nil {
c.Routing = *q.Routing
}
if q.Ollama != nil {
c.Ollama = append([]core.OllamaServer(nil), (*q.Ollama)...)
}
if q.Learning != nil {
c.Brain.AutoReward.Enabled = q.Learning.AutoRewardEnabled
c.Brain.AutoReward.Mode = q.Learning.AutoRewardMode
c.Brain.Consolidation.Enabled = q.Learning.ConsolidationEnabled
c.Brain.Consolidation.UseLLM = q.Learning.ConsolidationUseLLM
c.Autonomy.Enabled = q.Learning.AutonomyEnabled
c.Autonomy.UseLLM = q.Learning.AutonomyUseLLM
}
if err := s.store.ValidateConfig(c); err != nil {
s.err(w, 400, err)
return
}
if err := s.store.UpdateConfig(c); err != nil {
s.err(w, 500, err)
return
}
_ = s.store.AddKnowledgeEvent(core.KnowledgeEvent{Type: "admin.model_routing_changed", Summary: "Model routing / Ollama configuration updated", Reason: "PUT /admin/api/model-routing", Actor: "admin", Metadata: map[string]string{"chat_provider": c.Routing.ChatProvider, "embedding_provider": c.Routing.EmbeddingProvider}})
s.json(w, 200, modelRoutingFromConfig(c))
}
func (s *Server) adminSecretsStatus(w http.ResponseWriter, r *http.Request) {
sec := s.store.Secrets()
s.json(w, 200, map[string]any{"openai_configured": sec.OpenAIAPIKey != "", "app_key_configured": sec.AppAPIKey != "", "integration_token_configured": sec.IntegrationToken != "", "control_read_token_configured": sec.ControlReadToken != "", "worker_token_configured": sec.WorkerToken != "", "metrics_token_configured": sec.MetricsToken != "", "shard_tokens": len(sec.ShardAPIToken), "cluster_token_configured": sec.ClusterToken != ""})
}
func maskedSecret(v string) string {
if v == "" {
return ""
}
if len(v) <= 4 {
return "••••"
}
return "••••••••" + v[len(v)-4:]
}
func (s *Server) adminGetSecrets(w http.ResponseWriter, r *http.Request) {
sec := s.store.Secrets()
reveal := r.URL.Query().Get("reveal") == "1" && s.store.Config().Security.AllowSecretReveal
if reveal {
s.json(w, 200, map[string]any{"revealed": true, "app_api_key": sec.AppAPIKey, "integration_token": sec.IntegrationToken, "control_read_token": sec.ControlReadToken, "worker_token": sec.WorkerToken, "metrics_token": sec.MetricsToken, "shard_api_tokens": sec.ShardAPIToken, "cluster_token": sec.ClusterToken})
return
}
maskedShards := map[string]string{}
for k, v := range sec.ShardAPIToken {
maskedShards[k] = maskedSecret(v)
}
s.json(w, 200, map[string]any{"revealed": false, "reveal_allowed": s.store.Config().Security.AllowSecretReveal, "app_api_key": maskedSecret(sec.AppAPIKey), "integration_token": maskedSecret(sec.IntegrationToken), "control_read_token": maskedSecret(sec.ControlReadToken), "worker_token": maskedSecret(sec.WorkerToken), "metrics_token": maskedSecret(sec.MetricsToken), "shard_api_tokens": maskedShards, "cluster_token": maskedSecret(sec.ClusterToken)})
}
func (s *Server) adminPutSecrets(w http.ResponseWriter, r *http.Request) {
var q struct {
OpenAIAPIKey string `json:"openai_api_key,omitempty"`
AppAPIKey string `json:"app_api_key,omitempty"`
IntegrationToken string `json:"integration_token,omitempty"`
ControlReadToken string `json:"control_read_token,omitempty"`
WorkerToken string `json:"worker_token,omitempty"`
MetricsToken string `json:"metrics_token,omitempty"`
ShardAPIToken map[string]string `json:"shard_api_tokens,omitempty"`
ClusterToken string `json:"cluster_token,omitempty"`
}
if err := decode(r, &q); err != nil {
s.err(w, 400, err)
return
}
sec := s.store.Secrets()
envLocked := func(name string) bool {
_, ok := os.LookupEnv(name)
return ok && strings.TrimSpace(os.Getenv(name)) != ""
}
for name, value := range map[string]string{
"OPENAI_API_KEY": q.OpenAIAPIKey,
"NEUROFORGE_APP_API_KEY": q.AppAPIKey,
"NEUROFORGE_INTEGRATION_TOKEN": q.IntegrationToken,
"NEUROFORGE_CONTROL_READ_TOKEN": q.ControlReadToken,
"NEUROFORGE_WORKER_TOKEN": q.WorkerToken,
"NEUROFORGE_METRICS_TOKEN": q.MetricsToken,
"NEUROFORGE_CLUSTER_TOKEN": q.ClusterToken,
} {
if value != "" && envLocked(name) {
s.err(w, http.StatusConflict, fmt.Errorf("%s is environment-managed and cannot be changed through the admin API", name))
return
}
}
if q.OpenAIAPIKey != "" {
sec.OpenAIAPIKey = q.OpenAIAPIKey
}
if q.AppAPIKey != "" {
sec.AppAPIKey = q.AppAPIKey
}
if q.IntegrationToken != "" {
sec.IntegrationToken = q.IntegrationToken
}
if q.ControlReadToken != "" {
sec.ControlReadToken = q.ControlReadToken
}
if q.WorkerToken != "" {
sec.WorkerToken = q.WorkerToken
}
if q.MetricsToken != "" {
sec.MetricsToken = q.MetricsToken
}
if q.ShardAPIToken != nil {
sec.ShardAPIToken = q.ShardAPIToken
}
if q.ClusterToken != "" {
sec.ClusterToken = q.ClusterToken
}
if err := s.store.UpdateSecrets(sec); err != nil {
s.err(w, 500, err)
return
}
_ = s.store.AddKnowledgeEvent(core.KnowledgeEvent{Type: "admin.secrets_changed", Summary: "One or more service credentials were updated", Reason: "PUT /admin/api/secrets", Actor: "admin"})
s.json(w, 200, map[string]bool{"ok": true})
}
func (s *Server) adminProviderHealth(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), 8*time.Second)
defer cancel()
s.json(w, 200, s.router.Health(ctx))
}
func (s *Server) adminMemories(w http.ResponseWriter, r *http.Request) {
limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
if limit <= 0 {
limit = 100
}
all := s.store.MemoriesSnapshot()
if len(all) > limit {
all = all[len(all)-limit:]
}
s.json(w, 200, all)
}
func (s *Server) adminDeleteMemory(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
m, _ := s.store.GetMemory(id)
if err := s.store.DeleteMemory(id); err != nil {
s.err(w, 500, err)
return
}
if m != nil {
_ = s.store.AddKnowledgeEvent(core.KnowledgeEvent{Type: "memory.deleted", MemoryID: id, Summary: "Memory deleted by administrator", Reason: "DELETE /admin/api/memories/{id}", Actor: "admin", Metadata: map[string]string{"kind": m.Kind, "memory_type": m.MemoryType, "truth_key": m.TruthKey}})
}
s.json(w, 200, map[string]bool{"ok": true})
}
func (s *Server) adminSynapses(w http.ResponseWriter, r *http.Request) {
s.json(w, 200, s.store.SynapsesSnapshot())
}
func (s *Server) adminUsage(w http.ResponseWriter, r *http.Request) {
limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
if limit <= 0 {
limit = 100
}
s.json(w, 200, s.store.RecentUsage(limit))
}
func (s *Server) adminExport(w http.ResponseWriter, r *http.Request) {
s.json(w, 200, s.store.ExportSafe())
}
func (s *Server) adminConsolidate(w http.ResponseWriter, r *http.Request) {
out, err := s.brain.Consolidate(r.Context())
if err != nil {
// A cycle can partially succeed and still report a synthesis/provider error.
s.json(w, 207, map[string]any{"result": out, "warning": err.Error()})
return
}
s.json(w, 200, out)
}
func (s *Server) securityHeaders(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if s.store.Config().Security.SecureHeaders {
w.Header().Set("X-Content-Type-Options", "nosniff")
w.Header().Set("Cross-Origin-Opener-Policy", "same-origin")
w.Header().Set("X-Permitted-Cross-Domain-Policies", "none")
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'; img-src 'self' data:; connect-src 'self'; script-src 'self' 'unsafe-inline'; style-src 'self' 'unsafe-inline'; frame-ancestors 'none'; base-uri 'none'; form-action 'self'")
}
if r.URL.Path == "/admin" || strings.HasPrefix(r.URL.Path, "/admin/api/") {
w.Header().Set("Cache-Control", "no-store")
}
if r.Header.Get("X-Request-ID") == "" {
r.Header.Set("X-Request-ID", store.NewID("req"))
}
w.Header().Set("X-Request-ID", r.Header.Get("X-Request-ID"))
next.ServeHTTP(w, r)
})
}
func (s *Server) requestLimits(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
cfg := s.store.Config().HTTP
if cfg.MaxBodyBytes > 0 && r.Body != nil {
r.Body = http.MaxBytesReader(w, r.Body, cfg.MaxBodyBytes)
}
if r.URL.Path == "/healthz" || r.URL.Path == "/livez" || r.URL.Path == "/readyz" || r.URL.Path == "/metrics" {
next.ServeHTTP(w, r)
return
}
n := s.inflight.Add(1)
defer s.inflight.Add(-1)
max := int64(cfg.MaxConcurrentRequests)
if max <= 0 {
max = 128
}
if n > max {
w.Header().Set("Retry-After", "1")
s.err(w, http.StatusServiceUnavailable, errors.New("server is at the configured concurrent request limit"))
return
}
next.ServeHTTP(w, r)
})
}
func (s *Server) livez(w http.ResponseWriter, r *http.Request) {
s.json(w, http.StatusOK, map[string]any{"ok": true, "status": "alive", "time": time.Now().UTC(), "version": "0.8.2"})
}
func configuredModelAvailable(models map[string]bool, configured string) bool {
configured = strings.TrimSpace(configured)
if configured == "" {
return false
}
if models[configured] {
return true
}
if !strings.Contains(configured, ":") && models[configured+":latest"] {
return true
}
return false
}
func checkConfiguredOllamaModels(ctx context.Context, cfg core.Config) (bool, any) {
type tagsResponse struct {
Models []struct {
Name string `json:"name"`
} `json:"models"`
}
details := map[string]any{}
anyEnabled := false
for _, node := range cfg.Ollama {
if !node.Enabled || strings.TrimSpace(node.BaseURL) == "" {
continue
}
anyEnabled = true
req, err := http.NewRequestWithContext(ctx, http.MethodGet, strings.TrimRight(node.BaseURL, "/")+"/api/tags", nil)
if err != nil {
details[node.ID] = err.Error()
continue
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
details[node.ID] = err.Error()
continue
}
var tags tagsResponse
decodeErr := json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&tags)
resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 || decodeErr != nil {
details[node.ID] = fmt.Sprintf("HTTP %d / invalid tags response", resp.StatusCode)
continue
}
models := map[string]bool{}
for _, model := range tags.Models {
models[strings.TrimSpace(model.Name)] = true
}
chatOK := configuredModelAvailable(models, node.ChatModel)
embedOK := configuredModelAvailable(models, node.EmbeddingModel)
details[node.ID] = map[string]any{"reachable": true, "chat_model": node.ChatModel, "chat_present": chatOK, "embedding_model": node.EmbeddingModel, "embedding_present": embedOK}
if chatOK && embedOK {
return true, details
}
}
if !anyEnabled {
return false, map[string]any{"error": "no enabled Ollama node configured"}
}
return false, details
}
func (s *Server) readyz(w http.ResponseWriter, r *http.Request) {
cfg := s.store.Config()
sec := s.store.Secrets()
obs := s.store.ObservabilitySnapshot()
components := map[string]any{}
configOK := s.store.ValidateConfig(cfg) == nil
components["config"] = configOK
components["app_auth"] = !cfg.API.RequireKey || sec.AppAPIKey != ""
hasOllamaChat, hasOllamaEmbed := false, false
for _, o := range cfg.Ollama {
if !o.Enabled {
continue
}
if strings.TrimSpace(o.ChatModel) != "" {
hasOllamaChat = true
}
if strings.TrimSpace(o.EmbeddingModel) != "" {
hasOllamaEmbed = true
}
}
openAIReady := cfg.OpenAI.Enabled && sec.OpenAIAPIKey != ""
chatReady := (cfg.Routing.ChatProvider == "openai" && openAIReady) || (cfg.Routing.ChatProvider == "ollama" && hasOllamaChat) || ((cfg.Routing.ChatProvider == "auto" || cfg.Routing.ChatProvider == "") && (hasOllamaChat || openAIReady))
embedReady := (cfg.Routing.EmbeddingProvider == "openai" && openAIReady) || (cfg.Routing.EmbeddingProvider == "ollama" && hasOllamaEmbed) || ((cfg.Routing.EmbeddingProvider == "auto" || cfg.Routing.EmbeddingProvider == "") && (hasOllamaEmbed || openAIReady))
components["chat_route_configured"] = chatReady
components["embedding_route_configured"] = embedReady
clusterReady := !cfg.Cluster.Enabled || obs.ClusterLeaderID != ""
components["cluster"] = clusterReady
ollamaLiveReady := true
if s.readinessOllamaLive {
ctx, cancel := context.WithTimeout(r.Context(), 4*time.Second)
var detail any
ollamaLiveReady, detail = checkConfiguredOllamaModels(ctx, cfg)
cancel()
components["ollama_live_models"] = detail
}
stagingReady := true
if s.brain != nil {
ctx, cancel := context.WithTimeout(r.Context(), 4*time.Second)
if err := s.brain.CheckStagingPublisher(ctx); err != nil {
stagingReady = false
components["kb_staging"] = err.Error()
} else {
components["kb_staging"] = true
}
cancel()
}
ready := configOK && chatReady && embedReady && clusterReady && ollamaLiveReady && stagingReady && (!cfg.API.RequireKey || sec.AppAPIKey != "")
status := http.StatusOK
if !ready {
status = http.StatusServiceUnavailable
}
s.json(w, status, map[string]any{"ok": ready, "status": map[bool]string{true: "ready", false: "not_ready"}[ready], "components": components, "revision": obs.Revision, "time": time.Now().UTC()})
}