// Package server exposes the scraper as an HTTP service. // // Endpoints: // // POST /scrape — enqueue a full catalogue scrape // POST /scrape/book — enqueue a single-book scrape (JSON body: {"url":"..."}) // GET /health — liveness probe package server import ( "bytes" "context" "encoding/json" "fmt" "io" "log/slog" "net/http" "os" "strconv" "sync" "time" "github.com/libnovel/scraper/internal/orchestrator" "github.com/libnovel/scraper/internal/scraper" "github.com/libnovel/scraper/internal/writer" ) // Server wraps an HTTP mux with the scraping endpoints. type Server struct { addr string oCfg orchestrator.Config novel scraper.NovelScraper log *slog.Logger writer *writer.Writer mu sync.Mutex running bool 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, audioInFlight: make(map[string]chan struct{}), } } // ListenAndServe starts the HTTP server and blocks until the provided context // is cancelled. func (s *Server) ListenAndServe(ctx context.Context) error { mux := http.NewServeMux() mux.HandleFunc("GET /health", s.handleHealth) mux.HandleFunc("POST /scrape", s.handleScrapeCatalogue) mux.HandleFunc("POST /scrape/book", s.handleScrapeBook) // UI routes mux.HandleFunc("GET /", s.handleHome) mux.HandleFunc("GET /ranking", s.handleRanking) mux.HandleFunc("POST /ranking/refresh", s.handleRankingRefresh) mux.HandleFunc("GET /ranking/view", s.handleRankingView) mux.HandleFunc("GET /books/{slug}", s.handleBook) mux.HandleFunc("GET /books/{slug}/chapters/{n}", s.handleChapter) mux.HandleFunc("GET /books/{slug}/chapters-page", s.handleBookChaptersPage) mux.HandleFunc("POST /ui/scrape/book", s.handleUIScrapeBook) mux.HandleFunc("GET /ui/scrape/status", s.handleUIScrapeStatus) 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, Handler: mux, ReadTimeout: 15 * time.Second, WriteTimeout: 60 * time.Second, IdleTimeout: 60 * time.Second, } errCh := make(chan error, 1) go func() { errCh <- srv.ListenAndServe() }() s.log.Info("HTTP server listening", "addr", s.addr) select { case <-ctx.Done(): shutCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() return srv.Shutdown(shutCtx) case err := <-errCh: return err } } func (s *Server) handleHealth(w http.ResponseWriter, _ *http.Request) { w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]string{"status": "ok"}) } // handleChapterText returns the plain text of a chapter (markdown stripped) // for browser-side TTS. The browser POSTs this directly to Kokoro-FastAPI. func (s *Server) handleChapterText(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 } raw, err := s.writer.ReadChapter(slug, n) if err != nil { http.NotFound(w, r) return } w.Header().Set("Content-Type", "text/plain; charset=utf-8") w.Header().Set("Cache-Control", "no-store") 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 s.runAsync(w, cfg) } func (s *Server) handleScrapeBook(w http.ResponseWriter, r *http.Request) { var body struct { URL string `json:"url"` } if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.URL == "" { http.Error(w, `{"error":"request body must be JSON with \"url\" field"}`, http.StatusBadRequest) return } cfg := s.oCfg cfg.SingleBookURL = body.URL s.runAsync(w, cfg) } // runAsync launches an orchestrator in the background and returns 202 Accepted. // Only one scrape job runs at a time; concurrent requests receive 409 Conflict. func (s *Server) runAsync(w http.ResponseWriter, cfg orchestrator.Config) { s.mu.Lock() if s.running { s.mu.Unlock() http.Error(w, `{"error":"a scrape job is already running"}`, http.StatusConflict) return } s.running = true s.mu.Unlock() w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusAccepted) _ = json.NewEncoder(w).Encode(map[string]string{"status": "accepted"}) go func() { defer func() { s.mu.Lock() s.running = false s.mu.Unlock() }() ctx, cancel := context.WithTimeout(context.Background(), 24*time.Hour) defer cancel() o := orchestrator.New(cfg, s.novel, s.log) if err := o.Run(ctx); err != nil { s.log.Error("scrape job failed", "err", fmt.Sprintf("%v", err)) } }() }