// Package service orchestrates the odidere voice assistant server. // It coordinates the HTTP server, LLM clients, and handles // graceful shutdown. package service import ( "context" "embed" "encoding/json" "fmt" "html/template" "io/fs" "log/slog" "net/http" "os" "os/signal" "runtime/debug" "strings" "sync" "syscall" "time" "code.chimeric.al/chimerical/odidere/internal/config" "code.chimeric.al/chimerical/odidere/internal/job" "code.chimeric.al/chimerical/odidere/internal/llm" "code.chimeric.al/chimerical/odidere/internal/service/templates" "github.com/google/uuid" "github.com/robfig/cron/v3" ) // Context keys for request-scoped values. type key string const ( logKey key = "log" idKey key = "id" ipKey key = "ip" modelsTimeout = 10 * time.Second ) //go:embed all:static/* var static embed.FS // Service is the main application coordinator. // It owns the HTTP server and all processing clients. type Service struct { cfg *config.Config cron *cron.Cron jobs []*job.Job llms map[string]*llm.Client log *slog.Logger allowedModels map[string]map[string]struct{} mux *http.ServeMux server *http.Server tmpl *template.Template tools *llm.Registry } // New creates a Service from the provided configuration. // It initializes all clients and the HTTP server. func New(cfg *config.Config, log *slog.Logger) (*Service, error) { var svc = &Service{ cfg: cfg, log: log, mux: http.NewServeMux(), } // Setup tool registry. registry, err := llm.NewRegistry(cfg.Tools) if err != nil { return nil, fmt.Errorf("load tools: %v", err) } svc.tools = registry // Create LLM clients for each provider. svc.llms = make(map[string]*llm.Client, len(cfg.Providers)) for _, pc := range cfg.Providers { client, err := llm.NewClient( llm.Config{ Key: pc.Key, SystemMessage: cfg.SystemMessage, Timeout: pc.Timeout, URL: pc.URL, }, registry, log, ) if err != nil { return nil, fmt.Errorf( "create LLM client for provider %q: %w", pc.Name, err, ) } svc.llms[pc.Name] = client log.Info( "LLM provider registered", slog.String("provider", pc.Name), slog.String("url", pc.URL), ) } // Build model allowed list map: provider name -> set of base model IDs. svc.allowedModels = make( map[string]map[string]struct{}, len(cfg.Providers), ) for _, pc := range cfg.Providers { if len(pc.Models) == 0 { continue } allowed := make(map[string]struct{}, len(pc.Models)) for _, m := range pc.Models { base, _, _ := strings.Cut(m, ":") allowed[base] = struct{}{} } svc.allowedModels[pc.Name] = allowed } // Convert job configs to jobs. jobs := make([]*job.Job, 0, len(cfg.Jobs)) for _, jc := range cfg.Jobs { // Resolve the provider and model for this job. var provider, modelID string if jc.Model.Provider != "" { if jc.Model.Model == "" { return nil, fmt.Errorf( "job %q: model.provider set but model.model is empty", jc.Name, ) } provider = jc.Model.Provider modelID = jc.Model.Model } else { provider = svc.cfg.DefaultModel.Provider modelID = svc.cfg.DefaultModel.Model } llmc := svc.llms[provider] if llmc == nil { return nil, fmt.Errorf( "create job %q: no LLM client for provider %q", jc.Name, provider, ) } j, err := job.New(jc, log, llmc, modelID) if err != nil { return nil, fmt.Errorf( "create job %q: %w", jc.Name, err, ) } jobs = append(jobs, j) } svc.jobs = jobs // Parse templates. tmpl, err := templates.Parse() if err != nil { return nil, fmt.Errorf("parse templates: %v", err) } svc.tmpl = tmpl // Setup static file server. staticFS, err := fs.Sub(static, "static") if err != nil { return nil, fmt.Errorf("setup static fs: %v", err) } // Register routes. svc.mux.HandleFunc("GET /", svc.home) svc.mux.HandleFunc("GET /status", svc.status) svc.mux.Handle( "GET /static/", http.StripPrefix( "/static/", http.FileServer(http.FS(staticFS)), ), ) svc.mux.HandleFunc("POST /v1/chat/voice", svc.chat) svc.mux.HandleFunc("POST /v1/chat/voice/stream", svc.chatStream) svc.mux.HandleFunc("GET /v1/models", svc.models) svc.server = &http.Server{ Addr: cfg.Address, Handler: svc, } return svc, nil } // ServeHTTP implements http.Handler. It logs requests, assigns a UUID, // sets context values, handles panics, and delegates to the mux. func (svc *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) { var ( start = time.Now() id = uuid.NewString() ip = func() string { if ip := r.Header.Get("X-Forwarded-For"); ip != "" { before, _, found := strings.Cut(ip, ",") if found { return strings.TrimSpace(before) } return ip } return r.RemoteAddr }() log = svc.log.With(slog.Group( "request", slog.String("id", id), slog.String("ip", ip), slog.String("method", r.Method), slog.String("path", r.URL.Path), )) ) // Log completion time. defer func() { log.InfoContext( r.Context(), "completed", slog.Duration("duration", time.Since(start)), ) }() // Panic recovery. defer func() { if err := recover(); err != nil { log.ErrorContext( r.Context(), "panic recovered", slog.Any("error", err), slog.String("stack", string(debug.Stack())), ) http.Error( w, http.StatusText( http.StatusInternalServerError, ), http.StatusInternalServerError, ) } }() // Enrich context with request-scoped values. ctx := r.Context() ctx = context.WithValue(ctx, logKey, log) ctx = context.WithValue(ctx, idKey, id) ctx = context.WithValue(ctx, ipKey, ip) r = r.WithContext(ctx) log.InfoContext(ctx, "handling") // Pass the request on to the multiplexer. svc.mux.ServeHTTP(w, r) } // Run starts the service and blocks until shutdown. // Shutdown is triggered by SIGINT or SIGTERM. func (svc *Service) Run(ctx context.Context) error { svc.log.Info( "starting odidere", slog.Group( "llm", slog.String( "default_provider", svc.cfg.DefaultModel.Provider, ), slog.String( "default_model", svc.cfg.DefaultModel.Model, ), slog.Int("providers", len(svc.cfg.Providers)), ), slog.Group( "server", slog.String("address", svc.cfg.Address), ), slog.Group( "tools", slog.Int("count", len(svc.tools.List())), slog.Any( "names", strings.Join(svc.tools.List(), ","), ), ), slog.Group( "tts_url", slog.String("url", svc.cfg.TTS.URL), slog.String("default_voice", svc.cfg.TTS.Voice), ), slog.Group( "stt_url", slog.String("url", svc.cfg.STT.URL), ), slog.Group( "jobs", slog.Int("count", svc.JobCount()), ), slog.Any("shutdown_timeout", svc.cfg.GetShutdownTimeout()), ) // Setup signal handling for graceful shutdown. ctx, cancel := signal.NotifyContext( ctx, os.Interrupt, syscall.SIGTERM, ) defer cancel() // Register jobs with cron scheduler. svc.cron = cron.New() for _, j := range svc.jobs { entry, err := svc.cron.AddFunc(j.Schedule(), func() { result := j.Run(ctx) if result.Error != nil { svc.log.Error("job failed", slog.String("name", j.Name()), slog.Any("error", result.Error), slog.Duration( "duration", result.Duration, ), ) } else { svc.log.Info("job completed", slog.String("name", j.Name()), slog.Duration( "duration", result.Duration, ), ) } }) if err != nil { svc.log.Error("failed to schedule job", slog.String("name", j.Name()), slog.Any("error", err), ) continue } svc.log.Info("job scheduled", slog.String("name", j.Name()), slog.String("schedule", j.Schedule()), slog.Int("entry_id", int(entry)), ) } svc.cron.Start() // Start HTTP server in background. var errs = make(chan error, 1) go func() { svc.log.Info( "HTTP server listening", slog.String("address", svc.cfg.Address), ) if err := svc.server.ListenAndServe(); err != nil && err != http.ErrServerClosed { errs <- err } close(errs) }() // Wait for shutdown signal or server error. select { case <-ctx.Done(): svc.log.Info("shutdown signal received") case err := <-errs: if err != nil { return fmt.Errorf("server error: %w", err) } } // Graceful shutdown with timeout. shutdownCtx, shutdownCancel := context.WithTimeout( context.Background(), svc.cfg.GetShutdownTimeout(), ) defer shutdownCancel() // Stop cron scheduler — in-flight jobs continue to completion. svc.log.Info("stopping job scheduler") svc.cron.Stop() svc.log.Info("shutting down HTTP server") if err := svc.server.Shutdown(shutdownCtx); err != nil { svc.log.Warn( "shutdown timeout reached", slog.Any("error", err), ) } svc.log.Info("terminating") return nil } // ListJobs returns all registered jobs. func (svc *Service) ListJobs() []*job.Job { return svc.jobs } // JobCount returns the number of registered jobs. func (svc *Service) JobCount() int { return len(svc.jobs) } func (svc *Service) home(w http.ResponseWriter, r *http.Request) { var ( ctx = r.Context() log = ctx.Value(logKey).(*slog.Logger) ) if r.URL.Path != "/" { http.NotFound(w, r) return } w.Header().Set("Content-Type", "text/html; charset=utf-8") if err := svc.tmpl.ExecuteTemplate( w, "index.gohtml", struct { TTSURL string TTSDefaultVoice string STTURL string }{ TTSURL: svc.cfg.TTS.URL, TTSDefaultVoice: svc.cfg.TTS.Voice, STTURL: svc.cfg.STT.URL, }, ); err != nil { log.ErrorContext( ctx, "template error", slog.Any("error", err), ) http.Error( w, http.StatusText(http.StatusInternalServerError), http.StatusInternalServerError, ) } } // status returns server status. func (svc *Service) status(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) } // Request is the incoming request format for the chat endpoints. type Request struct { // Messages is the conversation history. Messages []llm.Message `json:"messages"` // Provider is the LLM provider name. If empty, the default provider is used. Provider string `json:"provider,omitempty"` // Model is the LLM model ID. If empty, the default model is used. Model string `json:"model,omitempty"` // SystemMessage overrides the configured system message for this request. SystemMessage string `json:"system_message,omitempty"` // MaxIterations caps the number of agent loop turns (0 = unlimited). MaxIterations int `json:"max_iterations,omitempty"` } // Response is the response format for chat and voice endpoints. type Response struct { // Messages is the full list of messages generated during the completions request, // including tool calls and tool results. Messages []llm.Message `json:"messages,omitempty"` // Provider is the LLM provider used for the response. Provider string `json:"used_provider,omitempty"` // Model is the LLM model used for the response. Model string `json:"used_model,omitempty"` // Usage is token usage from the API response. Usage *llm.Usage `json:"usage,omitempty"` // Timings is llama.cpp timing info from the API response. Timings *llm.Timings `json:"timings,omitempty"` } // chat processes text chat requests. func (svc *Service) chat(w http.ResponseWriter, r *http.Request) { var ( ctx = r.Context() log = ctx.Value(logKey).(*slog.Logger) ) // Parse request. r.Body = http.MaxBytesReader(w, r.Body, 32<<20) var req Request if err := json.NewDecoder(r.Body).Decode(&req); err != nil { log.ErrorContext( ctx, "failed to decode request", slog.Any("error", err), ) http.Error(w, "invalid request", http.StatusBadRequest) return } // Validate messages. if len(req.Messages) == 0 { http.Error(w, "messages required", http.StatusBadRequest) return } log.DebugContext( ctx, "messages", slog.Any("data", req.Messages), ) // Resolve provider and model. provider := req.Provider if provider == "" { provider = svc.cfg.DefaultModel.Provider } llmc := svc.llms[provider] if llmc == nil { log.ErrorContext( ctx, "unknown provider", slog.String("provider", provider), ) http.Error( w, fmt.Sprintf("unknown provider: %q", provider), http.StatusBadRequest, ) return } model := req.Model if model == "" { model = svc.cfg.DefaultModel.Model } // Validate model against provider allowed list, if configured. if allowed, ok := svc.allowedModels[provider]; ok { m, _, _ := strings.Cut(model, ":") if _, ok := allowed[m]; !ok { log.ErrorContext( ctx, "model not in allowed list", slog.String("provider", provider), slog.String("model", model), ) http.Error( w, fmt.Sprintf( "model %q not allowed for provider %q", model, provider, ), http.StatusBadRequest, ) return } } res, err := llmc.Completions( ctx, model, req.SystemMessage, req.MaxIterations, req.Messages, ) if err != nil { log.ErrorContext( ctx, "LLM request failed", slog.Any("error", err), ) http.Error(w, "LLM error", http.StatusInternalServerError) return } if len(res.Messages) == 0 { http.Error( w, "no response from LLM", http.StatusInternalServerError, ) return } final := res.Messages[len(res.Messages)-1] log.DebugContext( ctx, "LLM response", slog.String("text", final.Text()), slog.String("provider", provider), slog.String("model", model), ) w.Header().Set("Content-Type", "application/json") if err := json.NewEncoder(w).Encode(Response{ Messages: res.Messages, Provider: provider, Model: model, Usage: res.Usage, Timings: res.Timings, }); err != nil { log.ErrorContext( ctx, "failed to json encode response", slog.Any("error", err), ) } } // StreamMessage is the SSE event payload for the streaming chat endpoint. type StreamMessage struct { // Error is an error message, if any. Error string `json:"error,omitempty"` // Delta is a text fragment from the assistant response. Delta string `json:"delta,omitempty"` // ReasoningDelta is a reasoning text fragment from the assistant response. ReasoningDelta string `json:"reasoning_delta,omitempty"` // Message is the chat completion message. Message *llm.Message `json:"message,omitempty"` // Provider is the LLM provider used for the response. Provider string `json:"provider,omitempty"` // Model is the LLM model used for the response. Model string `json:"model,omitempty"` // Usage is token usage from the API response. Usage *llm.Usage `json:"usage,omitempty"` // Timings is llama.cpp timing info from the API response. Timings *llm.Timings `json:"timings,omitempty"` } // chatStream processes chat requests with streaming SSE output. func (svc *Service) chatStream(w http.ResponseWriter, r *http.Request) { var ( ctx = r.Context() log = ctx.Value(logKey).(*slog.Logger) ) // Check that the response writer supports flushing. flusher, ok := w.(http.Flusher) if !ok { http.Error( w, "streaming not supported", http.StatusInternalServerError, ) return } // Parse request. r.Body = http.MaxBytesReader(w, r.Body, 32<<20) var req Request if err := json.NewDecoder(r.Body).Decode(&req); err != nil { log.ErrorContext( ctx, "failed to decode request", slog.Any("error", err), ) http.Error(w, "invalid request", http.StatusBadRequest) return } // Validate messages. if len(req.Messages) == 0 { http.Error(w, "messages required", http.StatusBadRequest) return } // Resolve provider and model. provider := req.Provider if provider == "" { provider = svc.cfg.DefaultModel.Provider } model := req.Model if model == "" { model = svc.cfg.DefaultModel.Model } llmc := svc.llms[provider] if llmc == nil { log.ErrorContext( ctx, "unknown provider", slog.String("provider", provider), ) http.Error( w, fmt.Sprintf("unknown provider: %q", provider), http.StatusBadRequest, ) return } // Validate model against provider allowed list, if configured. if allowed, ok := svc.allowedModels[provider]; ok { m, _, _ := strings.Cut(model, ":") if _, ok := allowed[m]; !ok { log.ErrorContext( ctx, "model not in allowed list", slog.String("provider", provider), slog.String("model", model), ) http.Error( w, fmt.Sprintf( "model %q not allowed for provider %q", model, provider, ), http.StatusBadRequest, ) return } } // Set SSE headers. w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") w.Header().Set("Content-Type", "text/event-stream") // Helper to send an SSE event. send := func(event string, msg StreamMessage) { data, err := json.Marshal(msg) if err != nil { log.ErrorContext(ctx, "failed to marshal SSE event", slog.Any("error", err), ) return } fmt.Fprintf(w, "event: %s\ndata: %s\n\n", event, data) flusher.Flush() } // Start streaming LLM completions request. var ( events = make(chan llm.StreamEvent) errs = make(chan error, 1) ) go func() { errs <- llmc.CompletionsStream( ctx, model, req.SystemMessage, req.MaxIterations, req.Messages, events, ) close(errs) }() // Consume events and send as SSE. for evt := range events { msg := StreamMessage{ Delta: evt.Delta, ReasoningDelta: evt.ReasoningDelta, Provider: provider, Model: model, } if evt.Message.Role != "" { msg.Message = &evt.Message } if evt.Usage != nil { msg.Usage = evt.Usage } if evt.Timings != nil { msg.Timings = evt.Timings } send(evt.Type, msg) } if err := <-errs; err != nil { log.ErrorContext( ctx, "LLM stream failed", slog.Any("error", err), ) msg := llm.Message{Role: llm.RoleAssistant} send("message", StreamMessage{ Message: &msg, Error: fmt.Sprintf("LLM error: %v", err), }) return } } // ModelStatus represents the availability status of a model. type ModelStatus string const ( // ModelStatusAvailable indicates the model is available for use. ModelStatusAvailable ModelStatus = "available" // ModelStatusNotFound indicates the model is configured but not found // in the provider's model list. ModelStatusNotFound ModelStatus = "not_found" // ModelStatusProviderError indicates the provider could not be reached // to verify the model. ModelStatusProviderError ModelStatus = "provider_error" ) // Model represents a model in the /v1/models response. type Model struct { Model string `json:"model"` Status ModelStatus `json:"status"` ContextSize *int `json:"context_size,omitempty"` } // Models is the response format for the /v1/models endpoint. type Models struct { Providers map[string][]Model `json:"providers"` DefaultProvider string `json:"default_provider"` DefaultModel string `json:"default_model"` } // models returns available LLM models from all configured providers. func (svc *Service) models(w http.ResponseWriter, r *http.Request) { var ( ctx = r.Context() log = ctx.Value(logKey).(*slog.Logger) mu sync.Mutex wg sync.WaitGroup providers = make(map[string][]Model, len(svc.llms)) ) cfg := make(map[string]config.ProviderConfig, len(svc.cfg.Providers)) for _, pc := range svc.cfg.Providers { cfg[pc.Name] = pc } for name, client := range svc.llms { wg.Add(1) go func(name string, c *llm.Client, pc config.ProviderConfig) { defer wg.Done() var result = []Model{} fetchCtx, cancel := context.WithTimeout( ctx, modelsTimeout, ) defer cancel() models, err := c.ListModels(fetchCtx) if err != nil { log.WarnContext( ctx, "failed to list models for provider", slog.String("provider", name), slog.Any("error", err), ) for _, m := range pc.Models { result = append(result, Model{ Model: m, Status: ModelStatusProviderError, }) } mu.Lock() providers[name] = result mu.Unlock() return } available := make(map[string]struct{}, len(models)) contextSizes := make(map[string]int, len(models)) for _, m := range models { baseID, _, _ := strings.Cut(m.ID, ":") available[baseID] = struct{}{} cs := m.ContextLength if cs == 0 && m.Meta != nil { cs = m.Meta.NCtx } if cs > 0 { contextSizes[baseID] = cs } } var toCheck []string if len(pc.Models) > 0 { toCheck = pc.Models } else { toCheck = make([]string, 0, len(models)) for _, m := range models { toCheck = append(toCheck, m.ID) } } for _, m := range toCheck { base, _, _ := strings.Cut(m, ":") mi := Model{Model: m} if _, ok := available[base]; ok { mi.Status = ModelStatusAvailable if cs, ok := contextSizes[base]; ok { mi.ContextSize = &cs } } else { mi.Status = ModelStatusNotFound } result = append(result, mi) } mu.Lock() providers[name] = result mu.Unlock() }(name, client, cfg[name]) } wg.Wait() w.Header().Set("Content-Type", "application/json") if err := json.NewEncoder(w).Encode(Models{ Providers: providers, DefaultProvider: svc.cfg.DefaultModel.Provider, DefaultModel: svc.cfg.DefaultModel.Model, }); err != nil { log.ErrorContext(ctx, "failed to encode models response", slog.Any("error", err)) } }