diff --git a/scraper/internal/server/server.go b/scraper/internal/server/server.go index 48f48d5..8b37107 100644 --- a/scraper/internal/server/server.go +++ b/scraper/internal/server/server.go @@ -8,11 +8,14 @@ package server import ( + "bytes" "context" "encoding/json" "fmt" + "io" "log/slog" "net/http" + "os" "strconv" "sync" "time" @@ -34,18 +37,26 @@ type Server struct { rankingRunning bool kokoroURL string // Kokoro-FastAPI base URL, e.g. http://kokoro:8880 kokoroVoice string // default voice, e.g. af_bella + + // audioMu guards audioInFlight. + // audioInFlight maps an audio cache key to a channel that is closed when + // the in-flight Kokoro request for that key finishes (successfully or not). + // This prevents duplicate concurrent TTS generation for the same file. + audioMu sync.Mutex + audioInFlight map[string]chan struct{} } // New creates a new Server. func New(addr string, oCfg orchestrator.Config, novel scraper.NovelScraper, log *slog.Logger, kokoroURL, kokoroVoice string) *Server { return &Server{ - addr: addr, - oCfg: oCfg, - novel: novel, - log: log, - writer: writer.New(oCfg.StaticRoot), - kokoroURL: kokoroURL, - kokoroVoice: kokoroVoice, + addr: addr, + oCfg: oCfg, + novel: novel, + log: log, + writer: writer.New(oCfg.StaticRoot), + kokoroURL: kokoroURL, + kokoroVoice: kokoroVoice, + audioInFlight: make(map[string]chan struct{}), } } @@ -69,6 +80,16 @@ func (s *Server) ListenAndServe(ctx context.Context) error { mux.HandleFunc("GET /ui/ranking/status", s.handleRankingStatus) // Plain-text chapter content for browser-side TTS mux.HandleFunc("GET /ui/chapter-text/{slug}/{n}", s.handleChapterText) + // Server-side audio generation and serving. + // Audio generation can take several minutes for long chapters, so wrap it + // in its own timeout handler instead of relying on the server WriteTimeout. + audioGenHandler := http.TimeoutHandler( + http.HandlerFunc(s.handleAudioGenerate), + 10*time.Minute, + `{"error":"audio generation timed out"}`, + ) + mux.Handle("POST /ui/audio/{slug}/{n}", audioGenHandler) + mux.HandleFunc("GET /ui/audio-file/{slug}/{n}", s.handleAudioFile) srv := &http.Server{ Addr: s.addr, @@ -117,6 +138,201 @@ func (s *Server) handleChapterText(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, stripMarkdown(raw)) } +// handleAudioGenerate handles POST /ui/audio/{slug}/{n}. +// It accepts an optional JSON body {voice, speed} (falling back to server +// defaults). If the MP3 is already cached on disk it returns immediately; +// otherwise it calls Kokoro-FastAPI to generate and save the file. +// Response: JSON {"url": "/ui/audio-file/{slug}/{n}?voice=…&speed=…"} +// +// Concurrent requests for the same (slug, n, voice, speed) are deduplicated: +// the first caller does the work; subsequent callers block until it finishes +// and then serve the cached file (or receive the same error). +func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { + slug := r.PathValue("slug") + n, err := strconv.Atoi(r.PathValue("n")) + if err != nil || n < 1 { + http.Error(w, `{"error":"invalid chapter"}`, http.StatusBadRequest) + return + } + + // Parse optional voice/speed from JSON body. + voice := s.kokoroVoice + speed := 1.0 + var body struct { + Voice string `json:"voice"` + Speed float64 `json:"speed"` + } + if r.Body != nil { + _ = json.NewDecoder(r.Body).Decode(&body) + } + if body.Voice != "" { + voice = body.Voice + } + if body.Speed > 0 { + speed = body.Speed + } + + audioPath := s.writer.AudioPath(slug, n, voice, speed) + + // Idempotent: return immediately if already cached. + if _, err := os.Stat(audioPath); err == nil { + s.writeAudioURL(w, slug, n, voice, speed) + return + } + + // Deduplicate concurrent generation requests for the same file. + // If another goroutine is already generating this file, wait for it and + // then serve the (now-cached) result. + cacheKey := fmt.Sprintf("%s/%d/%s/%.2f", slug, n, voice, speed) + + s.audioMu.Lock() + if ch, ok := s.audioInFlight[cacheKey]; ok { + // Someone else is already generating — wait for them. + s.audioMu.Unlock() + select { + case <-ch: + case <-r.Context().Done(): + http.Error(w, `{"error":"request cancelled"}`, http.StatusServiceUnavailable) + return + } + // Serve the cached file (or 404 if generation failed). + if _, err := os.Stat(audioPath); err == nil { + s.writeAudioURL(w, slug, n, voice, speed) + } else { + http.Error(w, `{"error":"audio generation failed"}`, http.StatusInternalServerError) + } + return + } + // Register ourselves as the in-flight generator. + ch := make(chan struct{}) + s.audioInFlight[cacheKey] = ch + s.audioMu.Unlock() + + // Always close the channel (unblocking waiters) and remove our entry. + defer func() { + s.audioMu.Lock() + delete(s.audioInFlight, cacheKey) + s.audioMu.Unlock() + close(ch) + }() + + // Load chapter text. + raw, err := s.writer.ReadChapter(slug, n) + if err != nil { + http.Error(w, `{"error":"chapter not found"}`, http.StatusNotFound) + return + } + text := stripMarkdown(raw) + if text == "" { + http.Error(w, `{"error":"chapter text is empty"}`, http.StatusUnprocessableEntity) + return + } + + if s.kokoroURL == "" { + http.Error(w, `{"error":"kokoro not configured"}`, http.StatusServiceUnavailable) + return + } + + // Call Kokoro-FastAPI. + reqBody, _ := json.Marshal(map[string]interface{}{ + "model": "kokoro", + "input": text, + "voice": voice, + "response_format": "mp3", + "speed": speed, + "stream": false, + }) + kokoroReq, err := http.NewRequestWithContext(r.Context(), http.MethodPost, + s.kokoroURL+"/v1/audio/speech", bytes.NewReader(reqBody)) + if err != nil { + http.Error(w, `{"error":"failed to build kokoro request"}`, http.StatusInternalServerError) + return + } + kokoroReq.Header.Set("Content-Type", "application/json") + + resp, err := http.DefaultClient.Do(kokoroReq) + if err != nil { + s.log.Error("kokoro request failed", "err", err) + http.Error(w, `{"error":"kokoro unavailable"}`, http.StatusBadGateway) + return + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + body2, _ := io.ReadAll(resp.Body) + s.log.Error("kokoro returned error", "status", resp.StatusCode, "body", string(body2)) + http.Error(w, fmt.Sprintf(`{"error":"kokoro error %d"}`, resp.StatusCode), http.StatusBadGateway) + return + } + + // Ensure the audio directory exists. + if err := os.MkdirAll(s.writer.AudioDir(slug), 0o755); err != nil { + http.Error(w, `{"error":"failed to create audio dir"}`, http.StatusInternalServerError) + return + } + + // Write to a temp file then rename atomically. + tmpPath := audioPath + ".tmp" + f, err := os.Create(tmpPath) + if err != nil { + http.Error(w, `{"error":"failed to create temp file"}`, http.StatusInternalServerError) + return + } + if _, err := io.Copy(f, resp.Body); err != nil { + f.Close() + os.Remove(tmpPath) + http.Error(w, `{"error":"failed to write audio"}`, http.StatusInternalServerError) + return + } + f.Close() + if err := os.Rename(tmpPath, audioPath); err != nil { + os.Remove(tmpPath) + http.Error(w, `{"error":"failed to save audio"}`, http.StatusInternalServerError) + return + } + + s.log.Info("audio generated", "slug", slug, "chapter", n, "voice", voice, "speed", speed) + s.writeAudioURL(w, slug, n, voice, speed) +} + +func (s *Server) writeAudioURL(w http.ResponseWriter, slug string, n int, voice string, speed float64) { + url := fmt.Sprintf("/ui/audio-file/%s/%d?voice=%s&speed=%.1f", slug, n, voice, speed) + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]string{"url": url}) +} + +// handleAudioFile handles GET /ui/audio-file/{slug}/{n}. +// Serves the cached MP3 file identified by voice and speed query params. +func (s *Server) handleAudioFile(w http.ResponseWriter, r *http.Request) { + slug := r.PathValue("slug") + n, err := strconv.Atoi(r.PathValue("n")) + if err != nil || n < 1 { + http.NotFound(w, r) + return + } + voice := r.URL.Query().Get("voice") + if voice == "" { + voice = s.kokoroVoice + } + speedStr := r.URL.Query().Get("speed") + speed := 1.0 + if speedStr != "" { + if v, err := strconv.ParseFloat(speedStr, 64); err == nil && v > 0 { + speed = v + } + } + + audioPath := s.writer.AudioPath(slug, n, voice, speed) + if _, err := os.Stat(audioPath); err != nil { + http.NotFound(w, r) + return + } + + w.Header().Set("Content-Type", "audio/mpeg") + w.Header().Set("Cache-Control", "public, max-age=86400") + http.ServeFile(w, r, audioPath) +} + func (s *Server) handleScrapeCatalogue(w http.ResponseWriter, r *http.Request) { cfg := s.oCfg cfg.SingleBookURL = "" // full catalogue diff --git a/scraper/internal/server/ui.go b/scraper/internal/server/ui.go index 31f82bc..dd3e58a 100644 --- a/scraper/internal/server/ui.go +++ b/scraper/internal/server/ui.go @@ -162,7 +162,7 @@ const homeTmpl = `