diff --git a/scraper/internal/server/server.go b/scraper/internal/server/server.go index 8b37107..631c200 100644 --- a/scraper/internal/server/server.go +++ b/scraper/internal/server/server.go @@ -17,6 +17,7 @@ import ( "net/http" "os" "strconv" + "strings" "sync" "time" @@ -89,7 +90,9 @@ func (s *Server) ListenAndServe(ctx context.Context) error { `{"error":"audio generation timed out"}`, ) mux.Handle("POST /ui/audio/{slug}/{n}", audioGenHandler) + mux.HandleFunc("GET /ui/audio/{slug}/{n}/status", s.handleAudioStatus) mux.HandleFunc("GET /ui/audio-file/{slug}/{n}", s.handleAudioFile) + mux.HandleFunc("GET /ui/audio-file/{slug}/{n}/part/{p}", s.handleAudioFilePart) srv := &http.Server{ Addr: s.addr, @@ -138,15 +141,25 @@ 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=…"} +// ─── Chunked audio generation ──────────────────────────────────────────────── // -// 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). +// handleAudioGenerate handles POST /ui/audio/{slug}/{n}. +// +// Flow: +// 1. If the merged MP3 already exists on disk → return it immediately. +// 2. Otherwise split the chapter text into up to audioParts equal-ish chunks +// (by paragraph), generate part 0 synchronously (so the browser can start +// playing right away), then launch a background goroutine that generates +// parts 1…N and merges them into the final file. +// 3. Response: {"url":"","parts":N,"merged":false} +// (or {"url":"","parts":1,"merged":true} on a cache hit). +// +// Deduplication: if another request is already generating the *merged* file +// for the same (slug,n,voice,speed) key, the new request blocks until it +// finishes and then serves the cached result. +const audioParts = 10 + +// handleAudioGenerate handles POST /ui/audio/{slug}/{n}. func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { slug := r.PathValue("slug") n, err := strconv.Atoi(r.PathValue("n")) @@ -174,20 +187,17 @@ func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { audioPath := s.writer.AudioPath(slug, n, voice, speed) - // Idempotent: return immediately if already cached. + // Fast path: merged file already on disk. if _, err := os.Stat(audioPath); err == nil { - s.writeAudioURL(w, slug, n, voice, speed) + s.writeAudioResponse(w, slug, n, voice, speed, 1, true) 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. + // Deduplicate concurrent generation requests for the same merged file. 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: @@ -195,20 +205,17 @@ func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { 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) + s.writeAudioResponse(w, slug, n, voice, speed, 1, true) } 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) @@ -216,7 +223,7 @@ func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { close(ch) }() - // Load chapter text. + // Load and validate chapter text. raw, err := s.writer.ReadChapter(slug, n) if err != nil { http.Error(w, `{"error":"chapter not found"}`, http.StatusNotFound) @@ -227,13 +234,105 @@ func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { 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. + // Ensure audio dir exists. + if err := os.MkdirAll(s.writer.AudioDir(slug), 0o755); err != nil { + http.Error(w, `{"error":"failed to create audio dir"}`, http.StatusInternalServerError) + return + } + + parts := splitTextIntoParts(text, audioParts) + totalParts := len(parts) + + // Generate part 0 synchronously so the browser can start playing immediately. + if err := s.generateAudioPart(r.Context(), slug, n, voice, speed, 0, parts[0]); err != nil { + s.log.Error("part 0 generation failed", "err", err) + http.Error(w, `{"error":"part 0 generation failed"}`, http.StatusBadGateway) + return + } + + if totalParts == 1 { + // Only one part — rename it to the final path directly. + partPath := s.writer.AudioPartPath(slug, n, voice, speed, 0) + if err := os.Rename(partPath, audioPath); err != nil { + http.Error(w, `{"error":"failed to save audio"}`, http.StatusInternalServerError) + return + } + s.log.Info("audio generated (single part)", "slug", slug, "chapter", n) + s.writeAudioResponse(w, slug, n, voice, speed, 1, true) + return + } + + // Return part-0 URL immediately; generate the rest in the background. + s.writeAudioResponse(w, slug, n, voice, speed, totalParts, false) + + // Background: generate parts 1…N then merge. + go func() { + bgCtx := context.Background() + for p := 1; p < totalParts; p++ { + if err := s.generateAudioPart(bgCtx, slug, n, voice, speed, p, parts[p]); err != nil { + s.log.Error("background part generation failed", "part", p, "err", err) + return + } + } + if err := s.mergeAudioParts(slug, n, voice, speed, totalParts); err != nil { + s.log.Error("audio merge failed", "slug", slug, "chapter", n, "err", err) + return + } + s.log.Info("audio merged", "slug", slug, "chapter", n, "parts", totalParts) + }() +} + +// splitTextIntoParts divides text (paragraphs separated by blank lines) into +// at most n equal-ish chunks. Returns at least 1 element. +func splitTextIntoParts(text string, n int) []string { + // Split into paragraphs on blank lines. + raw := strings.Split(text, "\n\n") + var paras []string + for _, p := range raw { + p = strings.TrimSpace(p) + if p != "" { + paras = append(paras, p) + } + } + if len(paras) == 0 { + return []string{text} + } + if n > len(paras) { + n = len(paras) + } + if n < 1 { + n = 1 + } + + chunks := make([]string, n) + chunkSize := (len(paras) + n - 1) / n // ceiling division + for i := 0; i < n; i++ { + start := i * chunkSize + end := start + chunkSize + if start >= len(paras) { + // Fewer paragraphs than requested parts: return what we have. + chunks = chunks[:i] + break + } + if end > len(paras) { + end = len(paras) + } + chunks[i] = strings.Join(paras[start:end], "\n\n") + } + if len(chunks) == 0 { + return []string{text} + } + return chunks +} + +// generateAudioPart calls Kokoro for a single text chunk and writes the result +// atomically to AudioPartPath(…, part). +func (s *Server) generateAudioPart(ctx context.Context, slug string, n int, voice string, speed float64, part int, text string) error { reqBody, _ := json.Marshal(map[string]interface{}{ "model": "kokoro", "input": text, @@ -242,67 +341,140 @@ func (s *Server) handleAudioGenerate(w http.ResponseWriter, r *http.Request) { "speed": speed, "stream": false, }) - kokoroReq, err := http.NewRequestWithContext(r.Context(), http.MethodPost, + req, err := http.NewRequestWithContext(ctx, 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 + return fmt.Errorf("build request: %w", err) } - kokoroReq.Header.Set("Content-Type", "application/json") + req.Header.Set("Content-Type", "application/json") - resp, err := http.DefaultClient.Do(kokoroReq) + resp, err := http.DefaultClient.Do(req) if err != nil { - s.log.Error("kokoro request failed", "err", err) - http.Error(w, `{"error":"kokoro unavailable"}`, http.StatusBadGateway) - return + return fmt.Errorf("kokoro request: %w", err) } 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 + body, _ := io.ReadAll(resp.Body) + return fmt.Errorf("kokoro status %d: %s", resp.StatusCode, string(body)) } - // 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" + partPath := s.writer.AudioPartPath(slug, n, voice, speed, part) + tmpPath := partPath + ".tmp" f, err := os.Create(tmpPath) if err != nil { - http.Error(w, `{"error":"failed to create temp file"}`, http.StatusInternalServerError) - return + return fmt.Errorf("create temp file: %w", err) } 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 + return fmt.Errorf("write audio: %w", err) } f.Close() - if err := os.Rename(tmpPath, audioPath); err != nil { + if err := os.Rename(tmpPath, partPath); err != nil { os.Remove(tmpPath) - http.Error(w, `{"error":"failed to save audio"}`, http.StatusInternalServerError) - return + return fmt.Errorf("rename temp file: %w", err) } - - s.log.Info("audio generated", "slug", slug, "chapter", n, "voice", voice, "speed", speed) - s.writeAudioURL(w, slug, n, voice, speed) + return nil } -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) +// mergeAudioParts concatenates totalParts part files in order into AudioPath, +// then removes the individual part files. +func (s *Server) mergeAudioParts(slug string, n int, voice string, speed float64, totalParts int) error { + audioPath := s.writer.AudioPath(slug, n, voice, speed) + tmpPath := audioPath + ".tmp" + + out, err := os.Create(tmpPath) + if err != nil { + return fmt.Errorf("create merged temp: %w", err) + } + + for p := 0; p < totalParts; p++ { + partPath := s.writer.AudioPartPath(slug, n, voice, speed, p) + data, err := os.ReadFile(partPath) + if err != nil { + out.Close() + os.Remove(tmpPath) + return fmt.Errorf("read part %d: %w", p, err) + } + if _, err := out.Write(data); err != nil { + out.Close() + os.Remove(tmpPath) + return fmt.Errorf("write merged part %d: %w", p, err) + } + } + out.Close() + + if err := os.Rename(tmpPath, audioPath); err != nil { + os.Remove(tmpPath) + return fmt.Errorf("rename merged: %w", err) + } + + // Clean up part files (best-effort). + for p := 0; p < totalParts; p++ { + os.Remove(s.writer.AudioPartPath(slug, n, voice, speed, p)) + } + return nil +} + +func (s *Server) writeAudioResponse(w http.ResponseWriter, slug string, n int, voice string, speed float64, parts int, merged bool) { + var url string + if merged { + url = fmt.Sprintf("/ui/audio-file/%s/%d?voice=%s&speed=%.1f", slug, n, voice, speed) + } else { + url = fmt.Sprintf("/ui/audio-file/%s/%d/part/0?voice=%s&speed=%.1f", slug, n, voice, speed) + } w.Header().Set("Content-Type", "application/json") - _ = json.NewEncoder(w).Encode(map[string]string{"url": url}) + _ = json.NewEncoder(w).Encode(map[string]interface{}{ + "url": url, + "parts": parts, + "merged": merged, + }) +} + +// writeAudioURL is kept for backward compatibility (used by dedup waiters). +func (s *Server) writeAudioURL(w http.ResponseWriter, slug string, n int, voice string, speed float64) { + s.writeAudioResponse(w, slug, n, voice, speed, 1, true) +} + +// handleAudioStatus handles GET /ui/audio/{slug}/{n}/status. +// Returns {"merged":true/false,"url":"..."} so the browser can poll for the +// merged file after receiving a parts response from handleAudioGenerate. +func (s *Server) handleAudioStatus(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 + } + 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) + merged := false + if _, err := os.Stat(audioPath); err == nil { + merged = true + } + + mergedURL := 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]interface{}{ + "merged": merged, + "url": mergedURL, + }) } // handleAudioFile handles GET /ui/audio-file/{slug}/{n}. -// Serves the cached MP3 file identified by voice and speed query params. +// Serves the cached merged 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")) @@ -333,6 +505,43 @@ func (s *Server) handleAudioFile(w http.ResponseWriter, r *http.Request) { http.ServeFile(w, r, audioPath) } +// handleAudioFilePart handles GET /ui/audio-file/{slug}/{n}/part/{p}. +// Serves a specific MP3 part file. +func (s *Server) handleAudioFilePart(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 + } + p, err := strconv.Atoi(r.PathValue("p")) + if err != nil || p < 0 { + 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 + } + } + + partPath := s.writer.AudioPartPath(slug, n, voice, speed, p) + if _, err := os.Stat(partPath); err != nil { + http.NotFound(w, r) + return + } + + w.Header().Set("Content-Type", "audio/mpeg") + w.Header().Set("Cache-Control", "no-store") + http.ServeFile(w, r, partPath) +} + 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 dd3e58a..e165a51 100644 --- a/scraper/internal/server/ui.go +++ b/scraper/internal/server/ui.go @@ -3,6 +3,7 @@ package server import ( "bytes" "context" + "encoding/json" "fmt" "html/template" "net/http" @@ -137,15 +138,25 @@ const homeTmpl = `

Scrape a new book

-
- +
+ + + + + +