diff --git a/backend/backend b/backend/backend index 9dad2a8..d74d232 100755 Binary files a/backend/backend and b/backend/backend differ diff --git a/backend/cmd/backend/main.go b/backend/cmd/backend/main.go index 553489c..746050e 100644 --- a/backend/cmd/backend/main.go +++ b/backend/cmd/backend/main.go @@ -150,18 +150,19 @@ func run() error { Commit: commit, }, backend.Dependencies{ - BookReader: store, - RankingStore: store, - AudioStore: store, - PresignStore: store, - ProgressStore: store, - CoverStore: store, - Producer: producer, - TaskReader: store, - SearchIndex: searchIndex, - Kokoro: kokoroClient, - PocketTTS: pocketTTSClient, - Log: log, + BookReader: store, + RankingStore: store, + AudioStore: store, + TranslationStore: store, + PresignStore: store, + ProgressStore: store, + CoverStore: store, + Producer: producer, + TaskReader: store, + SearchIndex: searchIndex, + Kokoro: kokoroClient, + PocketTTS: pocketTTSClient, + Log: log, }, ) diff --git a/backend/cmd/runner/main.go b/backend/cmd/runner/main.go index 4c68bad..c17bd5e 100644 --- a/backend/cmd/runner/main.go +++ b/backend/cmd/runner/main.go @@ -24,6 +24,7 @@ import ( "github.com/libnovel/backend/internal/browser" "github.com/libnovel/backend/internal/config" "github.com/libnovel/backend/internal/kokoro" + "github.com/libnovel/backend/internal/libretranslate" "github.com/libnovel/backend/internal/meili" "github.com/libnovel/backend/internal/novelfire" "github.com/libnovel/backend/internal/otelsetup" @@ -128,6 +129,14 @@ func run() error { log.Warn("POCKET_TTS_URL not set — pocket-tts voice tasks will fail") } + // ── LibreTranslate ────────────────────────────────────────────────────── + ltClient := libretranslate.New(cfg.LibreTranslate.URL, cfg.LibreTranslate.APIKey) + if ltClient != nil { + log.Info("libretranslate enabled", "url", cfg.LibreTranslate.URL) + } else { + log.Info("LIBRETRANSLATE_URL not set — machine translation disabled") + } + // ── Meilisearch ───────────────────────────────────────────────────────── var searchIndex meili.Client if cfg.Meilisearch.URL != "" { @@ -149,6 +158,7 @@ func run() error { PollInterval: cfg.Runner.PollInterval, MaxConcurrentScrape: cfg.Runner.MaxConcurrentScrape, MaxConcurrentAudio: cfg.Runner.MaxConcurrentAudio, + MaxConcurrentTranslation: cfg.Runner.MaxConcurrentTranslation, OrchestratorWorkers: workers, MetricsAddr: cfg.Runner.MetricsAddr, CatalogueRefreshInterval: cfg.Runner.CatalogueRefreshInterval, @@ -170,16 +180,18 @@ func run() error { } deps := runner.Dependencies{ - Consumer: consumer, - BookWriter: store, - BookReader: store, - AudioStore: store, - CoverStore: store, - SearchIndex: searchIndex, - Novel: novel, - Kokoro: kokoroClient, - PocketTTS: pocketTTSClient, - Log: log, + Consumer: consumer, + BookWriter: store, + BookReader: store, + AudioStore: store, + CoverStore: store, + TranslationStore: store, + SearchIndex: searchIndex, + Novel: novel, + Kokoro: kokoroClient, + PocketTTS: pocketTTSClient, + LibreTranslate: ltClient, + Log: log, } r := runner.New(rCfg, deps) diff --git a/backend/go.mod b/backend/go.mod index 0be5bc0..e43a1a7 100644 --- a/backend/go.mod +++ b/backend/go.mod @@ -43,6 +43,7 @@ require ( github.com/rs/xid v1.6.0 // indirect github.com/spf13/cast v1.10.0 // indirect github.com/tinylib/msgp v1.6.1 // indirect + github.com/yuin/goldmark v1.8.2 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/bridges/otelslog v0.17.0 // indirect go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 // indirect diff --git a/backend/go.sum b/backend/go.sum index 39bddd0..241143c 100644 --- a/backend/go.sum +++ b/backend/go.sum @@ -84,6 +84,8 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu github.com/tinylib/msgp v1.6.1 h1:ESRv8eL3u+DNHUoSAAQRE50Hm162zqAnBoGv9PzScPY= github.com/tinylib/msgp v1.6.1/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA= github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= +github.com/yuin/goldmark v1.8.2 h1:kEGpgqJXdgbkhcOgBxkC0X0PmoPG1ZyoZ117rDVp4zE= +github.com/yuin/goldmark v1.8.2/go.mod h1:ip/1k0VRfGynBgxOz0yCqHrbZXhcjxyuS66Brc7iBKg= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= go.opentelemetry.io/contrib/bridges/otelslog v0.17.0 h1:NFIS6x7wyObQ7cR84x7bt1sr8nYBx89s3x3GwRjw40k= diff --git a/backend/internal/asynqqueue/consumer.go b/backend/internal/asynqqueue/consumer.go index 188cf2c..09979fd 100644 --- a/backend/internal/asynqqueue/consumer.go +++ b/backend/internal/asynqqueue/consumer.go @@ -37,6 +37,10 @@ func (c *Consumer) FinishAudioTask(ctx context.Context, id string, result domain return c.pb.FinishAudioTask(ctx, id, result) } +func (c *Consumer) FinishTranslationTask(ctx context.Context, id string, result domain.TranslationResult) error { + return c.pb.FinishTranslationTask(ctx, id, result) +} + func (c *Consumer) FailTask(ctx context.Context, id, errMsg string) error { return c.pb.FailTask(ctx, id, errMsg) } @@ -51,6 +55,10 @@ func (c *Consumer) ClaimNextAudioTask(_ context.Context, _ string) (domain.Audio return domain.AudioTask{}, false, nil } +func (c *Consumer) ClaimNextTranslationTask(_ context.Context, _ string) (domain.TranslationTask, bool, error) { + return domain.TranslationTask{}, false, nil +} + func (c *Consumer) HeartbeatTask(_ context.Context, _ string) error { return nil } func (c *Consumer) ReapStaleTasks(_ context.Context, _ time.Duration) (int, error) { return 0, nil } diff --git a/backend/internal/asynqqueue/producer.go b/backend/internal/asynqqueue/producer.go index dec69ef..b2c5568 100644 --- a/backend/internal/asynqqueue/producer.go +++ b/backend/internal/asynqqueue/producer.go @@ -73,6 +73,12 @@ func (p *Producer) CreateAudioTask(ctx context.Context, slug string, chapter int return id, nil } +// CreateTranslationTask creates a PocketBase record. Translation tasks are +// not currently dispatched via Asynq — the runner picks them up via polling. +func (p *Producer) CreateTranslationTask(ctx context.Context, slug string, chapter int, lang string) (string, error) { + return p.pb.CreateTranslationTask(ctx, slug, chapter, lang) +} + // CancelTask delegates to PocketBase; Asynq jobs may already be running and // cannot be reliably cancelled, so we only update the audit record. func (p *Producer) CancelTask(ctx context.Context, id string) error { diff --git a/backend/internal/backend/handlers.go b/backend/internal/backend/handlers.go index 871fdb9..d2e7f63 100644 --- a/backend/internal/backend/handlers.go +++ b/backend/internal/backend/handlers.go @@ -32,6 +32,7 @@ package backend // directly (no runner task, no store writes). Used for unscraped books. import ( + "bytes" "context" "encoding/json" "fmt" @@ -49,6 +50,7 @@ import ( "github.com/libnovel/backend/internal/novelfire/htmlutil" "github.com/libnovel/backend/internal/pockettts" "github.com/libnovel/backend/internal/scraper" + "github.com/yuin/goldmark" ) const ( @@ -701,9 +703,252 @@ func (s *Server) handleAudioProxy(w http.ResponseWriter, r *http.Request) { http.Redirect(w, r, presignURL, http.StatusFound) } -// ── Voices ───────────────────────────────────────────────────────────────────── +// ── Translation ──────────────────────────────────────────────────────────────── -// handleVoices handles GET /api/voices. +// supportedTranslationLangs is the set of target locales the backend accepts. +// Source is always "en". +var supportedTranslationLangs = map[string]bool{ + "ru": true, "id": true, "pt": true, "fr": true, +} + +// handleTranslationGenerate handles POST /api/translation/{slug}/{n}. +// Query params: lang (required, one of ru|id|pt|fr) +// +// Returns 200 immediately if translation already exists in MinIO. +// Returns 202 with task_id if a new task was created. +// Returns 503 if TranslationStore is nil (feature disabled). +func (s *Server) handleTranslationGenerate(w http.ResponseWriter, r *http.Request) { + if s.deps.TranslationStore == nil { + jsonError(w, http.StatusServiceUnavailable, "machine translation not configured") + return + } + + slug := r.PathValue("slug") + n, err := strconv.Atoi(r.PathValue("n")) + if err != nil || n < 1 { + jsonError(w, http.StatusBadRequest, "invalid chapter") + return + } + + lang := r.URL.Query().Get("lang") + if !supportedTranslationLangs[lang] { + jsonError(w, http.StatusBadRequest, "unsupported lang; use ru, id, pt, or fr") + return + } + + cacheKey := fmt.Sprintf("%s/%d/%s", slug, n, lang) + + // Fast path: translation already in MinIO + key := s.deps.TranslationStore.TranslationObjectKey(lang, slug, n) + if s.deps.TranslationStore.TranslationExists(r.Context(), key) { + writeJSON(w, 0, map[string]string{"status": "done", "lang": lang}) + return + } + + // Check if a task is already pending/running + task, found, _ := s.deps.TaskReader.GetTranslationTask(r.Context(), cacheKey) + if found && (task.Status == domain.TaskStatusPending || task.Status == domain.TaskStatusRunning) { + writeJSON(w, http.StatusAccepted, map[string]string{ + "task_id": task.ID, + "status": string(task.Status), + "lang": lang, + }) + return + } + + // Create a new translation task + taskID, err := s.deps.Producer.CreateTranslationTask(r.Context(), slug, n, lang) + if err != nil { + s.deps.Log.Error("handleTranslationGenerate: CreateTranslationTask failed", "err", err) + jsonError(w, http.StatusInternalServerError, "failed to create translation task") + return + } + + writeJSON(w, http.StatusAccepted, map[string]string{ + "task_id": taskID, + "status": "pending", + "lang": lang, + }) +} + +// handleTranslationStatus handles GET /api/translation/status/{slug}/{n}. +// Query params: lang (required) +func (s *Server) handleTranslationStatus(w http.ResponseWriter, r *http.Request) { + if s.deps.TranslationStore == nil { + writeJSON(w, 0, map[string]string{"status": "unavailable"}) + return + } + + slug := r.PathValue("slug") + n, err := strconv.Atoi(r.PathValue("n")) + if err != nil || n < 1 || slug == "" { + jsonError(w, http.StatusBadRequest, "invalid params") + return + } + + lang := r.URL.Query().Get("lang") + if !supportedTranslationLangs[lang] { + jsonError(w, http.StatusBadRequest, "unsupported lang") + return + } + + // Fast path: translation exists in MinIO + key := s.deps.TranslationStore.TranslationObjectKey(lang, slug, n) + if s.deps.TranslationStore.TranslationExists(r.Context(), key) { + writeJSON(w, 0, map[string]string{"status": "done", "lang": lang}) + return + } + + cacheKey := fmt.Sprintf("%s/%d/%s", slug, n, lang) + task, found, _ := s.deps.TaskReader.GetTranslationTask(r.Context(), cacheKey) + if !found { + writeJSON(w, 0, map[string]string{"status": "idle", "lang": lang}) + return + } + + resp := map[string]string{ + "status": string(task.Status), + "task_id": task.ID, + "lang": lang, + } + if task.Status == domain.TaskStatusFailed && task.ErrorMessage != "" { + resp["error"] = task.ErrorMessage + } + writeJSON(w, 0, resp) +} + +// handleTranslationRead handles GET /api/translation/{slug}/{n}. +// Query params: lang (required) +// +// Returns {"html": "

...

", "lang": "ru"} from the MinIO-cached translation. +// Returns 404 when the translation has not been generated yet. +func (s *Server) handleTranslationRead(w http.ResponseWriter, r *http.Request) { + if s.deps.TranslationStore == nil { + http.Error(w, `{"error":"machine translation not configured"}`, http.StatusServiceUnavailable) + return + } + + slug := r.PathValue("slug") + n, err := strconv.Atoi(r.PathValue("n")) + if err != nil || n < 1 || slug == "" { + jsonError(w, http.StatusBadRequest, "invalid params") + return + } + + lang := r.URL.Query().Get("lang") + if !supportedTranslationLangs[lang] { + jsonError(w, http.StatusBadRequest, "unsupported lang") + return + } + + key := s.deps.TranslationStore.TranslationObjectKey(lang, slug, n) + md, err := s.deps.TranslationStore.GetTranslation(r.Context(), key) + if err != nil { + s.deps.Log.Warn("handleTranslationRead: translation not found", "slug", slug, "n", n, "lang", lang, "err", err) + jsonError(w, http.StatusNotFound, "translation not available") + return + } + + var buf bytes.Buffer + if err := goldmark.Convert([]byte(md), &buf); err != nil { + s.deps.Log.Error("handleTranslationRead: markdown conversion failed", "err", err) + jsonError(w, http.StatusInternalServerError, "failed to render translation") + return + } + + writeJSON(w, 0, map[string]string{"html": buf.String(), "lang": lang}) +} + +// handleAdminTranslationJobs handles GET /api/admin/translation/jobs. +// Returns the full list of translation jobs sorted by started descending. +func (s *Server) handleAdminTranslationJobs(w http.ResponseWriter, r *http.Request) { + tasks, err := s.deps.TaskReader.ListTranslationTasks(r.Context()) + if err != nil { + s.deps.Log.Error("handleAdminTranslationJobs: ListTranslationTasks failed", "err", err) + jsonError(w, http.StatusInternalServerError, "failed to list translation jobs") + return + } + type jobRow struct { + ID string `json:"id"` + CacheKey string `json:"cache_key"` + Slug string `json:"slug"` + Chapter int `json:"chapter"` + Lang string `json:"lang"` + Status string `json:"status"` + WorkerID string `json:"worker_id"` + ErrorMessage string `json:"error_message"` + Started string `json:"started"` + Finished string `json:"finished"` + } + rows := make([]jobRow, 0, len(tasks)) + for _, t := range tasks { + rows = append(rows, jobRow{ + ID: t.ID, + CacheKey: t.CacheKey, + Slug: t.Slug, + Chapter: t.Chapter, + Lang: t.Lang, + Status: string(t.Status), + WorkerID: t.WorkerID, + ErrorMessage: t.ErrorMessage, + Started: t.Started.Format(time.RFC3339), + Finished: t.Finished.Format(time.RFC3339), + }) + } + writeJSON(w, 0, map[string]any{"jobs": rows}) +} + +// handleAdminTranslationBulk handles POST /api/admin/translation/bulk. +// Body: {"slug": "...", "lang": "ru", "from": 1, "to": 50} +// Enqueues one translation task per chapter in the range [from, to] inclusive. +func (s *Server) handleAdminTranslationBulk(w http.ResponseWriter, r *http.Request) { + var body struct { + Slug string `json:"slug"` + Lang string `json:"lang"` + From int `json:"from"` + To int `json:"to"` + } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + jsonError(w, http.StatusBadRequest, "invalid JSON body") + return + } + if body.Slug == "" { + jsonError(w, http.StatusBadRequest, "slug is required") + return + } + if !supportedTranslationLangs[body.Lang] { + jsonError(w, http.StatusBadRequest, "unsupported lang; use ru, id, pt, or fr") + return + } + if body.From < 1 || body.To < body.From { + jsonError(w, http.StatusBadRequest, "from must be >= 1 and to must be >= from") + return + } + if body.To-body.From > 999 { + jsonError(w, http.StatusBadRequest, "range too large; max 1000 chapters per request") + return + } + + var taskIDs []string + for n := body.From; n <= body.To; n++ { + id, err := s.deps.Producer.CreateTranslationTask(r.Context(), body.Slug, n, body.Lang) + if err != nil { + s.deps.Log.Error("handleAdminTranslationBulk: CreateTranslationTask failed", + "slug", body.Slug, "chapter", n, "lang", body.Lang, "err", err) + jsonError(w, http.StatusInternalServerError, + fmt.Sprintf("failed to create task for chapter %d: %s", n, err)) + return + } + taskIDs = append(taskIDs, id) + } + + writeJSON(w, http.StatusAccepted, map[string]any{ + "enqueued": len(taskIDs), + "task_ids": taskIDs, + }) +} + +// ── Voices ───────────────────────────────────────────────────────────────────── // Returns {"voices": [...]} — merged list from Kokoro and pocket-tts. func (s *Server) handleVoices(w http.ResponseWriter, r *http.Request) { writeJSON(w, 0, map[string]any{"voices": s.voices(r.Context())}) diff --git a/backend/internal/backend/server.go b/backend/internal/backend/server.go index c4d1d6e..5854df3 100644 --- a/backend/internal/backend/server.go +++ b/backend/internal/backend/server.go @@ -47,6 +47,8 @@ type Dependencies struct { RankingStore bookstore.RankingStore // AudioStore checks audio object existence and computes MinIO keys. AudioStore bookstore.AudioStore + // TranslationStore checks translation existence and reads/writes translated markdown. + TranslationStore bookstore.TranslationStore // PresignStore generates short-lived MinIO URLs. PresignStore bookstore.PresignStore // ProgressStore reads/writes per-session reading progress. @@ -160,6 +162,15 @@ func (s *Server) ListenAndServe(ctx context.Context) error { mux.HandleFunc("GET /api/audio/status/{slug}/{n}", s.handleAudioStatus) mux.HandleFunc("GET /api/audio-proxy/{slug}/{n}", s.handleAudioProxy) + // Translation task creation (backend creates task; runner executes via LibreTranslate) + mux.HandleFunc("POST /api/translation/{slug}/{n}", s.handleTranslationGenerate) + mux.HandleFunc("GET /api/translation/status/{slug}/{n}", s.handleTranslationStatus) + mux.HandleFunc("GET /api/translation/{slug}/{n}", s.handleTranslationRead) + + // Admin translation endpoints + mux.HandleFunc("GET /api/admin/translation/jobs", s.handleAdminTranslationJobs) + mux.HandleFunc("POST /api/admin/translation/bulk", s.handleAdminTranslationBulk) + // Voices list mux.HandleFunc("GET /api/voices", s.handleVoices) diff --git a/backend/internal/bookstore/bookstore.go b/backend/internal/bookstore/bookstore.go index d638d36..99b3a6b 100644 --- a/backend/internal/bookstore/bookstore.go +++ b/backend/internal/bookstore/bookstore.go @@ -141,3 +141,19 @@ type CoverStore interface { // CoverExists returns true when a cover image is stored for slug. CoverExists(ctx context.Context, slug string) bool } + +// TranslationStore covers machine-translated chapter storage in MinIO. +// The runner writes translations; the backend reads them. +type TranslationStore interface { + // TranslationObjectKey returns the MinIO object key for a cached translation. + TranslationObjectKey(lang, slug string, n int) string + + // TranslationExists returns true when the translation object is present in MinIO. + TranslationExists(ctx context.Context, key string) bool + + // PutTranslation stores raw translated markdown under the given MinIO object key. + PutTranslation(ctx context.Context, key string, data []byte) error + + // GetTranslation retrieves translated markdown from MinIO. + GetTranslation(ctx context.Context, key string) (string, error) +} diff --git a/backend/internal/config/config.go b/backend/internal/config/config.go index 998da9f..2232901 100644 --- a/backend/internal/config/config.go +++ b/backend/internal/config/config.go @@ -46,6 +46,8 @@ type MinIO struct { BucketAvatars string // BucketBrowse is the bucket that holds cached browse page snapshots (JSON). BucketBrowse string + // BucketTranslations is the bucket that holds machine-translated chapter markdown. + BucketTranslations string } // Kokoro holds connection settings for the Kokoro-FastAPI TTS service. @@ -64,6 +66,16 @@ type PocketTTS struct { URL string } +// LibreTranslate holds connection settings for a self-hosted LibreTranslate instance. +type LibreTranslate struct { + // URL is the base URL of the LibreTranslate instance, e.g. https://translate.libnovel.cc + // An empty string disables machine translation entirely. + URL string + // APIKey is the optional API key for the LibreTranslate instance. + // Leave empty if the instance runs without authentication. + APIKey string +} + // HTTP holds settings for the HTTP server (backend only). type HTTP struct { // Addr is the listen address, e.g. ":8080" @@ -107,6 +119,8 @@ type Runner struct { MaxConcurrentScrape int // MaxConcurrentAudio limits simultaneous audio-generation goroutines. MaxConcurrentAudio int + // MaxConcurrentTranslation limits simultaneous translation goroutines. + MaxConcurrentTranslation int // WorkerID is a unique identifier for this runner instance. // Defaults to the system hostname. WorkerID string @@ -135,15 +149,16 @@ type Runner struct { // Config is the top-level configuration struct consumed by both binaries. type Config struct { - PocketBase PocketBase - MinIO MinIO - Kokoro Kokoro - PocketTTS PocketTTS - HTTP HTTP - Runner Runner - Meilisearch Meilisearch - Valkey Valkey - Redis Redis + PocketBase PocketBase + MinIO MinIO + Kokoro Kokoro + PocketTTS PocketTTS + LibreTranslate LibreTranslate + HTTP HTTP + Runner Runner + Meilisearch Meilisearch + Valkey Valkey + Redis Redis // LogLevel is one of "debug", "info", "warn", "error". LogLevel string } @@ -166,16 +181,17 @@ func Load() Config { }, MinIO: MinIO{ - Endpoint: envOr("MINIO_ENDPOINT", "localhost:9000"), - PublicEndpoint: envOr("MINIO_PUBLIC_ENDPOINT", ""), - AccessKey: envOr("MINIO_ACCESS_KEY", "admin"), - SecretKey: envOr("MINIO_SECRET_KEY", "changeme123"), - UseSSL: envBool("MINIO_USE_SSL", false), - PublicUseSSL: envBool("MINIO_PUBLIC_USE_SSL", true), - BucketChapters: envOr("MINIO_BUCKET_CHAPTERS", "chapters"), - BucketAudio: envOr("MINIO_BUCKET_AUDIO", "audio"), - BucketAvatars: envOr("MINIO_BUCKET_AVATARS", "avatars"), - BucketBrowse: envOr("MINIO_BUCKET_BROWSE", "catalogue"), + Endpoint: envOr("MINIO_ENDPOINT", "localhost:9000"), + PublicEndpoint: envOr("MINIO_PUBLIC_ENDPOINT", ""), + AccessKey: envOr("MINIO_ACCESS_KEY", "admin"), + SecretKey: envOr("MINIO_SECRET_KEY", "changeme123"), + UseSSL: envBool("MINIO_USE_SSL", false), + PublicUseSSL: envBool("MINIO_PUBLIC_USE_SSL", true), + BucketChapters: envOr("MINIO_BUCKET_CHAPTERS", "chapters"), + BucketAudio: envOr("MINIO_BUCKET_AUDIO", "audio"), + BucketAvatars: envOr("MINIO_BUCKET_AVATARS", "avatars"), + BucketBrowse: envOr("MINIO_BUCKET_BROWSE", "catalogue"), + BucketTranslations: envOr("MINIO_BUCKET_TRANSLATIONS", "translations"), }, Kokoro: Kokoro{ @@ -195,6 +211,7 @@ func Load() Config { PollInterval: envDuration("RUNNER_POLL_INTERVAL", 30*time.Second), MaxConcurrentScrape: envInt("RUNNER_MAX_CONCURRENT_SCRAPE", 1), MaxConcurrentAudio: envInt("RUNNER_MAX_CONCURRENT_AUDIO", 1), + MaxConcurrentTranslation: envInt("RUNNER_MAX_CONCURRENT_TRANSLATION", 1), WorkerID: envOr("RUNNER_WORKER_ID", workerID), Workers: envInt("RUNNER_WORKERS", 0), // 0 → runtime.NumCPU() Timeout: envDuration("RUNNER_TIMEOUT", 90*time.Second), diff --git a/backend/internal/domain/domain.go b/backend/internal/domain/domain.go index 226ef44..f6ede8e 100644 --- a/backend/internal/domain/domain.go +++ b/backend/internal/domain/domain.go @@ -149,3 +149,23 @@ type AudioResult struct { ObjectKey string `json:"object_key,omitempty"` ErrorMessage string `json:"error_message,omitempty"` } + +// TranslationTask represents a machine-translation job stored in PocketBase. +type TranslationTask struct { + ID string `json:"id"` + CacheKey string `json:"cache_key"` // "{slug}/{chapter}/{lang}" + Slug string `json:"slug"` + Chapter int `json:"chapter"` + Lang string `json:"lang"` + WorkerID string `json:"worker_id,omitempty"` + Status TaskStatus `json:"status"` + ErrorMessage string `json:"error_message,omitempty"` + Started time.Time `json:"started"` + Finished time.Time `json:"finished,omitempty"` +} + +// TranslationResult is the outcome reported by the runner after finishing a TranslationTask. +type TranslationResult struct { + ObjectKey string `json:"object_key,omitempty"` + ErrorMessage string `json:"error_message,omitempty"` +} diff --git a/backend/internal/libretranslate/client.go b/backend/internal/libretranslate/client.go new file mode 100644 index 0000000..3417864 --- /dev/null +++ b/backend/internal/libretranslate/client.go @@ -0,0 +1,181 @@ +// Package libretranslate provides an HTTP client for a self-hosted +// LibreTranslate instance. It handles text chunking, concurrent translation, +// and reassembly so callers can pass arbitrarily long markdown strings. +package libretranslate + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "strings" + "sync" + "time" +) + +const ( + // maxChunkBytes is the target maximum size of each chunk sent to + // LibreTranslate. LibreTranslate's default limit is 5000 characters; + // we stay comfortably below that. + maxChunkBytes = 4500 + // concurrency is the number of simultaneous translation requests per chapter. + concurrency = 3 +) + +// Client translates text via LibreTranslate. +// A nil Client is valid — all calls return the original text unchanged. +type Client interface { + // Translate translates text from sourceLang to targetLang. + // text is a raw markdown string. The returned string is the translated + // markdown, reassembled in original paragraph order. + Translate(ctx context.Context, text, sourceLang, targetLang string) (string, error) +} + +// New returns a Client for the given LibreTranslate URL. +// Returns nil when url is empty, which disables translation. +func New(url, apiKey string) Client { + if url == "" { + return nil + } + return &httpClient{ + url: strings.TrimRight(url, "/"), + apiKey: apiKey, + http: &http.Client{Timeout: 60 * time.Second}, + } +} + +type httpClient struct { + url string + apiKey string + http *http.Client +} + +// Translate splits text into paragraph chunks, translates them concurrently +// (up to concurrency goroutines), and reassembles in order. +func (c *httpClient) Translate(ctx context.Context, text, sourceLang, targetLang string) (string, error) { + paragraphs := splitParagraphs(text) + if len(paragraphs) == 0 { + return text, nil + } + chunks := binChunks(paragraphs, maxChunkBytes) + + translated := make([]string, len(chunks)) + errs := make([]error, len(chunks)) + + sem := make(chan struct{}, concurrency) + var wg sync.WaitGroup + + for i, chunk := range chunks { + wg.Add(1) + sem <- struct{}{} + go func(idx int, chunkText string) { + defer wg.Done() + defer func() { <-sem }() + result, err := c.translateChunk(ctx, chunkText, sourceLang, targetLang) + translated[idx] = result + errs[idx] = err + }(i, chunk) + } + wg.Wait() + + for _, err := range errs { + if err != nil { + return "", err + } + } + + return strings.Join(translated, "\n\n"), nil +} + +// translateChunk sends a single POST /translate request. +func (c *httpClient) translateChunk(ctx context.Context, text, sourceLang, targetLang string) (string, error) { + reqBody := map[string]string{ + "q": text, + "source": sourceLang, + "target": targetLang, + "format": "html", + } + if c.apiKey != "" { + reqBody["api_key"] = c.apiKey + } + + b, err := json.Marshal(reqBody) + if err != nil { + return "", fmt.Errorf("libretranslate: marshal request: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.url+"/translate", bytes.NewReader(b)) + if err != nil { + return "", fmt.Errorf("libretranslate: build request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + + resp, err := c.http.Do(req) + if err != nil { + return "", fmt.Errorf("libretranslate: request: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + var errBody struct { + Error string `json:"error"` + } + _ = json.NewDecoder(resp.Body).Decode(&errBody) + return "", fmt.Errorf("libretranslate: status %d: %s", resp.StatusCode, errBody.Error) + } + + var result struct { + TranslatedText string `json:"translatedText"` + } + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + return "", fmt.Errorf("libretranslate: decode response: %w", err) + } + return result.TranslatedText, nil +} + +// splitParagraphs splits markdown text on blank lines, preserving non-empty paragraphs. +func splitParagraphs(text string) []string { + // Normalise line endings. + text = strings.ReplaceAll(text, "\r\n", "\n") + // Split on double newlines (blank lines between paragraphs). + parts := strings.Split(text, "\n\n") + var paragraphs []string + for _, p := range parts { + p = strings.TrimSpace(p) + if p != "" { + paragraphs = append(paragraphs, p) + } + } + return paragraphs +} + +// binChunks groups paragraphs into chunks each at most maxBytes in length. +// Each chunk is a single string with paragraphs joined by "\n\n". +func binChunks(paragraphs []string, maxBytes int) []string { + var chunks []string + var current strings.Builder + + for _, p := range paragraphs { + needed := len(p) + if current.Len() > 0 { + needed += 2 // for the "\n\n" separator + } + + if current.Len()+needed > maxBytes && current.Len() > 0 { + // Flush current chunk. + chunks = append(chunks, current.String()) + current.Reset() + } + + if current.Len() > 0 { + current.WriteString("\n\n") + } + current.WriteString(p) + } + + if current.Len() > 0 { + chunks = append(chunks, current.String()) + } + return chunks +} diff --git a/backend/internal/runner/runner.go b/backend/internal/runner/runner.go index 9acfda3..bdb19ce 100644 --- a/backend/internal/runner/runner.go +++ b/backend/internal/runner/runner.go @@ -29,6 +29,7 @@ import ( "github.com/libnovel/backend/internal/bookstore" "github.com/libnovel/backend/internal/domain" "github.com/libnovel/backend/internal/kokoro" + "github.com/libnovel/backend/internal/libretranslate" "github.com/libnovel/backend/internal/meili" "github.com/libnovel/backend/internal/orchestrator" "github.com/libnovel/backend/internal/pockettts" @@ -48,6 +49,8 @@ type Config struct { MaxConcurrentScrape int // MaxConcurrentAudio limits simultaneous audio-generation goroutines. MaxConcurrentAudio int + // MaxConcurrentTranslation limits simultaneous translation goroutines. + MaxConcurrentTranslation int // OrchestratorWorkers is the chapter-scraping parallelism inside each book run. OrchestratorWorkers int // HeartbeatInterval is how often active tasks PATCH their heartbeat_at @@ -95,6 +98,8 @@ type Dependencies struct { BookReader bookstore.BookReader // AudioStore persists generated audio and checks key existence. AudioStore bookstore.AudioStore + // TranslationStore persists translated markdown and checks key existence. + TranslationStore bookstore.TranslationStore // CoverStore stores book cover images in MinIO. CoverStore bookstore.CoverStore // SearchIndex indexes books in Meilisearch after scraping. @@ -107,6 +112,9 @@ type Dependencies struct { // PocketTTS is the pocket-tts client (CPU, kyutai voices: alba, marius, etc.). // If nil, pocket-tts voice tasks will fail with a clear error. PocketTTS pockettts.Client + // LibreTranslate is the machine translation client. + // If nil, translation tasks will fail with a clear error. + LibreTranslate libretranslate.Client // Log is the structured logger. Log *slog.Logger } @@ -137,6 +145,9 @@ func New(cfg Config, deps Dependencies) *Runner { if cfg.MaxConcurrentAudio <= 0 { cfg.MaxConcurrentAudio = 1 } + if cfg.MaxConcurrentTranslation <= 0 { + cfg.MaxConcurrentTranslation = 1 + } if cfg.WorkerID == "" { cfg.WorkerID = "runner" } @@ -175,6 +186,7 @@ func (r *Runner) Run(ctx context.Context) error { "mode", r.mode(), "max_scrape", r.cfg.MaxConcurrentScrape, "max_audio", r.cfg.MaxConcurrentAudio, + "max_translation", r.cfg.MaxConcurrentTranslation, "catalogue_refresh_interval", r.cfg.CatalogueRefreshInterval, "metrics_addr", r.cfg.MetricsAddr, ) @@ -208,6 +220,7 @@ func (r *Runner) mode() string { func (r *Runner) runPoll(ctx context.Context) error { scrapeSem := make(chan struct{}, r.cfg.MaxConcurrentScrape) audioSem := make(chan struct{}, r.cfg.MaxConcurrentAudio) + translationSem := make(chan struct{}, r.cfg.MaxConcurrentTranslation) var wg sync.WaitGroup tick := time.NewTicker(r.cfg.PollInterval) @@ -227,7 +240,7 @@ func (r *Runner) runPoll(ctx context.Context) error { // Run one poll immediately on startup, then on each tick. for { - r.poll(ctx, scrapeSem, audioSem, &wg) + r.poll(ctx, scrapeSem, audioSem, translationSem, &wg) select { case <-ctx.Done(): @@ -252,7 +265,7 @@ func (r *Runner) runPoll(ctx context.Context) error { } // poll claims all available pending tasks and dispatches them to goroutines. -func (r *Runner) poll(ctx context.Context, scrapeSem, audioSem chan struct{}, wg *sync.WaitGroup) { +func (r *Runner) poll(ctx context.Context, scrapeSem, audioSem, translationSem chan struct{}, wg *sync.WaitGroup) { // ── Heartbeat file ──────────────────────────────────────────────────── // Touch /tmp/runner.alive so the Docker health check can confirm the // runner is actively polling. Failure is non-fatal — just log it. @@ -335,6 +348,39 @@ audioLoop: r.runAudioTask(ctx, t) }(task) } + + // ── Translation tasks ───────────────────────────────────────────────── +translationLoop: + for { + if ctx.Err() != nil { + return + } + select { + case translationSem <- struct{}{}: + // Slot acquired — proceed to claim a task. + default: + // All slots busy; leave remaining pending tasks for next tick. + break translationLoop + } + task, ok, err := r.deps.Consumer.ClaimNextTranslationTask(ctx, r.cfg.WorkerID) + if err != nil { + <-translationSem + r.deps.Log.Error("runner: ClaimNextTranslationTask failed", "err", err) + break + } + if !ok { + <-translationSem + break + } + r.tasksRunning.Add(1) + wg.Add(1) + go func(t domain.TranslationTask) { + defer wg.Done() + defer func() { <-translationSem }() + defer r.tasksRunning.Add(-1) + r.runTranslationTask(ctx, t) + }(task) + } } // newOrchestrator builds an orchestrator with the Meilisearch post-hook wired in. diff --git a/backend/internal/runner/runner_test.go b/backend/internal/runner/runner_test.go index 9770089..c6a8a0a 100644 --- a/backend/internal/runner/runner_test.go +++ b/backend/internal/runner/runner_test.go @@ -48,6 +48,10 @@ func (s *stubConsumer) ClaimNextAudioTask(_ context.Context, _ string) (domain.A return t, true, nil } +func (s *stubConsumer) ClaimNextTranslationTask(_ context.Context, _ string) (domain.TranslationTask, bool, error) { + return domain.TranslationTask{}, false, nil +} + func (s *stubConsumer) FinishScrapeTask(_ context.Context, id string, _ domain.ScrapeResult) error { s.finished = append(s.finished, id) return nil @@ -58,6 +62,11 @@ func (s *stubConsumer) FinishAudioTask(_ context.Context, id string, _ domain.Au return nil } +func (s *stubConsumer) FinishTranslationTask(_ context.Context, id string, _ domain.TranslationResult) error { + s.finished = append(s.finished, id) + return nil +} + func (s *stubConsumer) FailTask(_ context.Context, id, _ string) error { s.failCalled = append(s.failCalled, id) return nil diff --git a/backend/internal/runner/translation.go b/backend/internal/runner/translation.go new file mode 100644 index 0000000..9b1926f --- /dev/null +++ b/backend/internal/runner/translation.go @@ -0,0 +1,97 @@ +package runner + +import ( + "context" + "fmt" + "time" + + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + + "github.com/libnovel/backend/internal/domain" +) + +// runTranslationTask executes one machine-translation task end-to-end and +// reports the result back to PocketBase. +func (r *Runner) runTranslationTask(ctx context.Context, task domain.TranslationTask) { + ctx, span := otel.Tracer("runner").Start(ctx, "runner.translation_task") + defer span.End() + span.SetAttributes( + attribute.String("task.id", task.ID), + attribute.String("book.slug", task.Slug), + attribute.Int("chapter.number", task.Chapter), + attribute.String("translation.lang", task.Lang), + ) + + log := r.deps.Log.With("task_id", task.ID, "slug", task.Slug, "chapter", task.Chapter, "lang", task.Lang) + log.Info("runner: translation task starting") + + // Heartbeat goroutine — keeps the task alive while translation runs. + hbCtx, hbCancel := context.WithCancel(ctx) + defer hbCancel() + go func() { + tick := time.NewTicker(r.cfg.HeartbeatInterval) + defer tick.Stop() + for { + select { + case <-hbCtx.Done(): + return + case <-tick.C: + if err := r.deps.Consumer.HeartbeatTask(ctx, task.ID); err != nil { + log.Warn("runner: heartbeat failed", "err", err) + } + } + } + }() + + fail := func(msg string) { + log.Error("runner: translation task failed", "reason", msg) + r.tasksFailed.Add(1) + span.SetStatus(codes.Error, msg) + result := domain.TranslationResult{ErrorMessage: msg} + if err := r.deps.Consumer.FinishTranslationTask(ctx, task.ID, result); err != nil { + log.Error("runner: FinishTranslationTask failed", "err", err) + } + } + + // Guard: LibreTranslate must be configured. + if r.deps.LibreTranslate == nil { + fail("libretranslate client not configured (LIBRETRANSLATE_URL is empty)") + return + } + + // 1. Read raw markdown chapter. + raw, err := r.deps.BookReader.ReadChapter(ctx, task.Slug, task.Chapter) + if err != nil { + fail(fmt.Sprintf("read chapter: %v", err)) + return + } + if raw == "" { + fail("chapter text is empty") + return + } + + // 2. Translate (chunked, concurrent). + translated, err := r.deps.LibreTranslate.Translate(ctx, raw, "en", task.Lang) + if err != nil { + fail(fmt.Sprintf("translate: %v", err)) + return + } + + // 3. Store translated markdown in MinIO. + key := r.deps.TranslationStore.TranslationObjectKey(task.Lang, task.Slug, task.Chapter) + if err := r.deps.TranslationStore.PutTranslation(ctx, key, []byte(translated)); err != nil { + fail(fmt.Sprintf("put translation: %v", err)) + return + } + + // 4. Report success. + r.tasksCompleted.Add(1) + span.SetStatus(codes.Ok, "") + result := domain.TranslationResult{ObjectKey: key} + if err := r.deps.Consumer.FinishTranslationTask(ctx, task.ID, result); err != nil { + log.Error("runner: FinishTranslationTask failed", "err", err) + } + log.Info("runner: translation task finished", "key", key) +} diff --git a/backend/internal/storage/minio.go b/backend/internal/storage/minio.go index 3f5217a..af84577 100644 --- a/backend/internal/storage/minio.go +++ b/backend/internal/storage/minio.go @@ -17,12 +17,13 @@ import ( // minioClient wraps the official minio-go client with bucket names. type minioClient struct { - client *minio.Client // internal — all read/write operations - pubClient *minio.Client // presign-only — initialised against the public endpoint - bucketChapters string - bucketAudio string - bucketAvatars string - bucketBrowse string + client *minio.Client // internal — all read/write operations + pubClient *minio.Client // presign-only — initialised against the public endpoint + bucketChapters string + bucketAudio string + bucketAvatars string + bucketBrowse string + bucketTranslations string } func newMinioClient(cfg config.MinIO) (*minioClient, error) { @@ -74,18 +75,19 @@ func newMinioClient(cfg config.MinIO) (*minioClient, error) { } return &minioClient{ - client: internal, - pubClient: pub, - bucketChapters: cfg.BucketChapters, - bucketAudio: cfg.BucketAudio, - bucketAvatars: cfg.BucketAvatars, - bucketBrowse: cfg.BucketBrowse, + client: internal, + pubClient: pub, + bucketChapters: cfg.BucketChapters, + bucketAudio: cfg.BucketAudio, + bucketAvatars: cfg.BucketAvatars, + bucketBrowse: cfg.BucketBrowse, + bucketTranslations: cfg.BucketTranslations, }, nil } // ensureBuckets creates all required buckets if they don't already exist. func (m *minioClient) ensureBuckets(ctx context.Context) error { - for _, bucket := range []string{m.bucketChapters, m.bucketAudio, m.bucketAvatars, m.bucketBrowse} { + for _, bucket := range []string{m.bucketChapters, m.bucketAudio, m.bucketAvatars, m.bucketBrowse, m.bucketTranslations} { exists, err := m.client.BucketExists(ctx, bucket) if err != nil { return fmt.Errorf("minio: check bucket %q: %w", bucket, err) @@ -125,6 +127,12 @@ func CoverObjectKey(slug string) string { return fmt.Sprintf("covers/%s.jpg", slug) } +// TranslationObjectKey returns the MinIO object key for a translated chapter. +// Format: {lang}/{slug}/{n:06d}.md +func TranslationObjectKey(lang, slug string, n int) string { + return fmt.Sprintf("%s/%s/%06d.md", lang, slug, n) +} + // chapterNumberFromKey extracts the chapter number from a MinIO object key. // e.g. "my-book/chapter-000042.md" → 42 func chapterNumberFromKey(key string) int { diff --git a/backend/internal/storage/store.go b/backend/internal/storage/store.go index a6eb5e8..d4f84d0 100644 --- a/backend/internal/storage/store.go +++ b/backend/internal/storage/store.go @@ -51,6 +51,7 @@ var _ bookstore.AudioStore = (*Store)(nil) var _ bookstore.PresignStore = (*Store)(nil) var _ bookstore.ProgressStore = (*Store)(nil) var _ bookstore.CoverStore = (*Store)(nil) +var _ bookstore.TranslationStore = (*Store)(nil) var _ taskqueue.Producer = (*Store)(nil) var _ taskqueue.Consumer = (*Store)(nil) var _ taskqueue.Reader = (*Store)(nil) @@ -535,13 +536,36 @@ func (s *Store) CreateAudioTask(ctx context.Context, slug string, chapter int, v return rec.ID, nil } +func (s *Store) CreateTranslationTask(ctx context.Context, slug string, chapter int, lang string) (string, error) { + cacheKey := fmt.Sprintf("%s/%d/%s", slug, chapter, lang) + payload := map[string]any{ + "cache_key": cacheKey, + "slug": slug, + "chapter": chapter, + "lang": lang, + "status": string(domain.TaskStatusPending), + "started": time.Now().UTC().Format(time.RFC3339), + } + var rec struct { + ID string `json:"id"` + } + if err := s.pb.post(ctx, "/api/collections/translation_jobs/records", payload, &rec); err != nil { + return "", err + } + return rec.ID, nil +} + func (s *Store) CancelTask(ctx context.Context, id string) error { - // Try scraping_tasks first, then audio_jobs. + // Try scraping_tasks first, then audio_jobs, then translation_jobs. if err := s.pb.patch(ctx, fmt.Sprintf("/api/collections/scraping_tasks/records/%s", id), map[string]string{"status": string(domain.TaskStatusCancelled)}); err == nil { return nil } - return s.pb.patch(ctx, fmt.Sprintf("/api/collections/audio_jobs/records/%s", id), + if err := s.pb.patch(ctx, fmt.Sprintf("/api/collections/audio_jobs/records/%s", id), + map[string]string{"status": string(domain.TaskStatusCancelled)}); err == nil { + return nil + } + return s.pb.patch(ctx, fmt.Sprintf("/api/collections/translation_jobs/records/%s", id), map[string]string{"status": string(domain.TaskStatusCancelled)}) } @@ -571,6 +595,18 @@ func (s *Store) ClaimNextAudioTask(ctx context.Context, workerID string) (domain return task, err == nil, err } +func (s *Store) ClaimNextTranslationTask(ctx context.Context, workerID string) (domain.TranslationTask, bool, error) { + raw, err := s.pb.claimRecord(ctx, "translation_jobs", workerID, nil) + if err != nil { + return domain.TranslationTask{}, false, err + } + if raw == nil { + return domain.TranslationTask{}, false, nil + } + task, err := parseTranslationTask(raw) + return task, err == nil, err +} + func (s *Store) FinishScrapeTask(ctx context.Context, id string, result domain.ScrapeResult) error { status := string(domain.TaskStatusDone) if result.ErrorMessage != "" { @@ -599,6 +635,18 @@ func (s *Store) FinishAudioTask(ctx context.Context, id string, result domain.Au }) } +func (s *Store) FinishTranslationTask(ctx context.Context, id string, result domain.TranslationResult) error { + status := string(domain.TaskStatusDone) + if result.ErrorMessage != "" { + status = string(domain.TaskStatusFailed) + } + return s.pb.patch(ctx, fmt.Sprintf("/api/collections/translation_jobs/records/%s", id), map[string]any{ + "status": status, + "error_message": result.ErrorMessage, + "finished": time.Now().UTC().Format(time.RFC3339), + }) +} + func (s *Store) FailTask(ctx context.Context, id, errMsg string) error { payload := map[string]any{ "status": string(domain.TaskStatusFailed), @@ -608,11 +656,14 @@ func (s *Store) FailTask(ctx context.Context, id, errMsg string) error { if err := s.pb.patch(ctx, fmt.Sprintf("/api/collections/scraping_tasks/records/%s", id), payload); err == nil { return nil } - return s.pb.patch(ctx, fmt.Sprintf("/api/collections/audio_jobs/records/%s", id), payload) + if err := s.pb.patch(ctx, fmt.Sprintf("/api/collections/audio_jobs/records/%s", id), payload); err == nil { + return nil + } + return s.pb.patch(ctx, fmt.Sprintf("/api/collections/translation_jobs/records/%s", id), payload) } // HeartbeatTask updates the heartbeat_at field on a running task. -// Tries scraping_tasks first, then audio_jobs (same pattern as FailTask). +// Tries scraping_tasks first, then audio_jobs, then translation_jobs. func (s *Store) HeartbeatTask(ctx context.Context, id string) error { payload := map[string]any{ "heartbeat_at": time.Now().UTC().Format(time.RFC3339), @@ -620,7 +671,10 @@ func (s *Store) HeartbeatTask(ctx context.Context, id string) error { if err := s.pb.patch(ctx, fmt.Sprintf("/api/collections/scraping_tasks/records/%s", id), payload); err == nil { return nil } - return s.pb.patch(ctx, fmt.Sprintf("/api/collections/audio_jobs/records/%s", id), payload) + if err := s.pb.patch(ctx, fmt.Sprintf("/api/collections/audio_jobs/records/%s", id), payload); err == nil { + return nil + } + return s.pb.patch(ctx, fmt.Sprintf("/api/collections/translation_jobs/records/%s", id), payload) } // ReapStaleTasks finds all running tasks whose heartbeat_at is either missing @@ -638,7 +692,7 @@ func (s *Store) ReapStaleTasks(ctx context.Context, staleAfter time.Duration) (i } total := 0 - for _, collection := range []string{"scraping_tasks", "audio_jobs"} { + for _, collection := range []string{"scraping_tasks", "audio_jobs", "translation_jobs"} { items, err := s.pb.listAll(ctx, collection, filter, "") if err != nil { return total, fmt.Errorf("ReapStaleTasks list %s: %w", collection, err) @@ -715,6 +769,31 @@ func (s *Store) GetAudioTask(ctx context.Context, cacheKey string) (domain.Audio return t, err == nil, err } +func (s *Store) ListTranslationTasks(ctx context.Context) ([]domain.TranslationTask, error) { + items, err := s.pb.listAll(ctx, "translation_jobs", "", "-started") + if err != nil { + return nil, err + } + tasks := make([]domain.TranslationTask, 0, len(items)) + for _, raw := range items { + t, err := parseTranslationTask(raw) + if err == nil { + tasks = append(tasks, t) + } + } + return tasks, nil +} + +func (s *Store) GetTranslationTask(ctx context.Context, cacheKey string) (domain.TranslationTask, bool, error) { + filter := fmt.Sprintf(`cache_key='%s'`, cacheKey) + items, err := s.pb.listAll(ctx, "translation_jobs", filter, "-started") + if err != nil || len(items) == 0 { + return domain.TranslationTask{}, false, err + } + t, err := parseTranslationTask(items[0]) + return t, err == nil, err +} + // ── Parsers ─────────────────────────────────────────────────────────────────── func parseScrapeTask(raw json.RawMessage) (domain.ScrapeTask, error) { @@ -789,6 +868,38 @@ func parseAudioTask(raw json.RawMessage) (domain.AudioTask, error) { }, nil } +func parseTranslationTask(raw json.RawMessage) (domain.TranslationTask, error) { + var rec struct { + ID string `json:"id"` + CacheKey string `json:"cache_key"` + Slug string `json:"slug"` + Chapter int `json:"chapter"` + Lang string `json:"lang"` + WorkerID string `json:"worker_id"` + Status string `json:"status"` + ErrorMessage string `json:"error_message"` + Started string `json:"started"` + Finished string `json:"finished"` + } + if err := json.Unmarshal(raw, &rec); err != nil { + return domain.TranslationTask{}, err + } + started, _ := time.Parse(time.RFC3339, rec.Started) + finished, _ := time.Parse(time.RFC3339, rec.Finished) + return domain.TranslationTask{ + ID: rec.ID, + CacheKey: rec.CacheKey, + Slug: rec.Slug, + Chapter: rec.Chapter, + Lang: rec.Lang, + WorkerID: rec.WorkerID, + Status: domain.TaskStatus(rec.Status), + ErrorMessage: rec.ErrorMessage, + Started: started, + Finished: finished, + }, nil +} + // ── CoverStore ───────────────────────────────────────────────────────────────── func (s *Store) PutCover(ctx context.Context, slug string, data []byte, contentType string) error { @@ -818,3 +929,25 @@ func (s *Store) GetCover(ctx context.Context, slug string) ([]byte, string, bool func (s *Store) CoverExists(ctx context.Context, slug string) bool { return s.mc.coverExists(ctx, CoverObjectKey(slug)) } + +// ── TranslationStore ─────────────────────────────────────────────────────────── + +func (s *Store) TranslationObjectKey(lang, slug string, n int) string { + return TranslationObjectKey(lang, slug, n) +} + +func (s *Store) TranslationExists(ctx context.Context, key string) bool { + return s.mc.objectExists(ctx, s.mc.bucketTranslations, key) +} + +func (s *Store) PutTranslation(ctx context.Context, key string, data []byte) error { + return s.mc.putObject(ctx, s.mc.bucketTranslations, key, "text/markdown; charset=utf-8", data) +} + +func (s *Store) GetTranslation(ctx context.Context, key string) (string, error) { + data, err := s.mc.getObject(ctx, s.mc.bucketTranslations, key) + if err != nil { + return "", fmt.Errorf("GetTranslation: %w", err) + } + return string(data), nil +} diff --git a/backend/internal/taskqueue/taskqueue.go b/backend/internal/taskqueue/taskqueue.go index 1ea1a32..05dca6e 100644 --- a/backend/internal/taskqueue/taskqueue.go +++ b/backend/internal/taskqueue/taskqueue.go @@ -29,6 +29,10 @@ type Producer interface { // returns the assigned PocketBase record ID. CreateAudioTask(ctx context.Context, slug string, chapter int, voice string) (string, error) + // CreateTranslationTask inserts a new translation task with status=pending and + // returns the assigned PocketBase record ID. + CreateTranslationTask(ctx context.Context, slug string, chapter int, lang string) (string, error) + // CancelTask transitions a pending task to status=cancelled. // Returns ErrNotFound if the task does not exist. CancelTask(ctx context.Context, id string) error @@ -46,13 +50,21 @@ type Consumer interface { // Returns (zero, false, nil) when the queue is empty. ClaimNextAudioTask(ctx context.Context, workerID string) (domain.AudioTask, bool, error) + // ClaimNextTranslationTask atomically finds the oldest pending translation task, + // sets its status=running and worker_id=workerID, and returns it. + // Returns (zero, false, nil) when the queue is empty. + ClaimNextTranslationTask(ctx context.Context, workerID string) (domain.TranslationTask, bool, error) + // FinishScrapeTask marks a running scrape task as done and records the result. FinishScrapeTask(ctx context.Context, id string, result domain.ScrapeResult) error // FinishAudioTask marks a running audio task as done and records the result. FinishAudioTask(ctx context.Context, id string, result domain.AudioResult) error - // FailTask marks a task (scrape or audio) as failed with an error message. + // FinishTranslationTask marks a running translation task as done and records the result. + FinishTranslationTask(ctx context.Context, id string, result domain.TranslationResult) error + + // FailTask marks a task (scrape, audio, or translation) as failed with an error message. FailTask(ctx context.Context, id, errMsg string) error // HeartbeatTask updates the heartbeat_at timestamp on a running task. @@ -81,4 +93,11 @@ type Reader interface { // GetAudioTask returns the most recent audio task for cacheKey. // Returns (zero, false, nil) if not found. GetAudioTask(ctx context.Context, cacheKey string) (domain.AudioTask, bool, error) + + // ListTranslationTasks returns all translation tasks sorted by started descending. + ListTranslationTasks(ctx context.Context) ([]domain.TranslationTask, error) + + // GetTranslationTask returns the most recent translation task for cacheKey. + // Returns (zero, false, nil) if not found. + GetTranslationTask(ctx context.Context, cacheKey string) (domain.TranslationTask, bool, error) } diff --git a/backend/internal/taskqueue/taskqueue_test.go b/backend/internal/taskqueue/taskqueue_test.go index 4b3eb17..e2372d3 100644 --- a/backend/internal/taskqueue/taskqueue_test.go +++ b/backend/internal/taskqueue/taskqueue_test.go @@ -23,6 +23,9 @@ func (s *stubStore) CreateScrapeTask(_ context.Context, _, _ string, _, _ int) ( func (s *stubStore) CreateAudioTask(_ context.Context, _ string, _ int, _ string) (string, error) { return "audio-1", nil } +func (s *stubStore) CreateTranslationTask(_ context.Context, _ string, _ int, _ string) (string, error) { + return "translation-1", nil +} func (s *stubStore) CancelTask(_ context.Context, _ string) error { return nil } func (s *stubStore) ClaimNextScrapeTask(_ context.Context, _ string) (domain.ScrapeTask, bool, error) { @@ -31,12 +34,18 @@ func (s *stubStore) ClaimNextScrapeTask(_ context.Context, _ string) (domain.Scr func (s *stubStore) ClaimNextAudioTask(_ context.Context, _ string) (domain.AudioTask, bool, error) { return domain.AudioTask{ID: "audio-1", Status: domain.TaskStatusRunning}, true, nil } +func (s *stubStore) ClaimNextTranslationTask(_ context.Context, _ string) (domain.TranslationTask, bool, error) { + return domain.TranslationTask{ID: "translation-1", Status: domain.TaskStatusRunning}, true, nil +} func (s *stubStore) FinishScrapeTask(_ context.Context, _ string, _ domain.ScrapeResult) error { return nil } func (s *stubStore) FinishAudioTask(_ context.Context, _ string, _ domain.AudioResult) error { return nil } +func (s *stubStore) FinishTranslationTask(_ context.Context, _ string, _ domain.TranslationResult) error { + return nil +} func (s *stubStore) FailTask(_ context.Context, _, _ string) error { return nil } func (s *stubStore) HeartbeatTask(_ context.Context, _ string) error { return nil } @@ -53,6 +62,12 @@ func (s *stubStore) ListAudioTasks(_ context.Context) ([]domain.AudioTask, error func (s *stubStore) GetAudioTask(_ context.Context, _ string) (domain.AudioTask, bool, error) { return domain.AudioTask{}, false, nil } +func (s *stubStore) ListTranslationTasks(_ context.Context) ([]domain.TranslationTask, error) { + return nil, nil +} +func (s *stubStore) GetTranslationTask(_ context.Context, _ string) (domain.TranslationTask, bool, error) { + return domain.TranslationTask{}, false, nil +} // Verify the stub satisfies all three interfaces at compile time. var _ taskqueue.Producer = (*stubStore)(nil) diff --git a/backend/runner b/backend/runner new file mode 100755 index 0000000..106fd96 Binary files /dev/null and b/backend/runner differ diff --git a/ui/src/lib/server/pocketbase.ts b/ui/src/lib/server/pocketbase.ts index bb93653..7f3f4a1 100644 --- a/ui/src/lib/server/pocketbase.ts +++ b/ui/src/lib/server/pocketbase.ts @@ -917,6 +917,24 @@ export async function listAudioJobs(): Promise { return listAll('audio_jobs', '', '-started'); } +// ─── Translation jobs ───────────────────────────────────────────────────────── + +export interface TranslationJob { + id: string; + cache_key: string; // "slug/chapter/lang" + slug: string; + chapter: number; + lang: string; + status: string; // "pending" | "running" | "done" | "failed" + error_message: string; + started: string; + finished: string; +} + +export async function listTranslationJobs(): Promise { + return listAll('translation_jobs', '', '-started'); +} + export async function getAudioTime( sessionId: string, slug: string, diff --git a/ui/src/routes/admin/+layout.svelte b/ui/src/routes/admin/+layout.svelte index f001e5f..17637e1 100644 --- a/ui/src/routes/admin/+layout.svelte +++ b/ui/src/routes/admin/+layout.svelte @@ -5,6 +5,7 @@ const internalLinks = [ { href: '/admin/scrape', label: 'Scrape' }, { href: '/admin/audio', label: 'Audio' }, + { href: '/admin/translation', label: 'Translation' }, { href: '/admin/changelog', label: 'Changelog' } ]; diff --git a/ui/src/routes/admin/translation/+page.server.ts b/ui/src/routes/admin/translation/+page.server.ts new file mode 100644 index 0000000..ab2ce21 --- /dev/null +++ b/ui/src/routes/admin/translation/+page.server.ts @@ -0,0 +1,63 @@ +import { redirect } from '@sveltejs/kit'; +import type { Actions, PageServerLoad } from './$types'; +import { listBooks, listTranslationJobs, type TranslationJob } from '$lib/server/pocketbase'; +import { backendFetch } from '$lib/server/scraper'; +import { log } from '$lib/server/logger'; + +export const load: PageServerLoad = async ({ locals }) => { + if (locals.user?.role !== 'admin') { + redirect(302, '/'); + } + + const [books, jobs] = await Promise.all([ + listBooks().catch((e): Awaited> => { + log.warn('admin/translation', 'failed to load books', { err: String(e) }); + return []; + }), + listTranslationJobs().catch((e): TranslationJob[] => { + log.warn('admin/translation', 'failed to load translation jobs', { err: String(e) }); + return []; + }) + ]); + + return { books, jobs }; +}; + +export const actions: Actions = { + bulk: async ({ request, locals }) => { + if (locals.user?.role !== 'admin') { + return { success: false, error: 'Unauthorized' }; + } + + const form = await request.formData(); + const slug = form.get('slug')?.toString().trim() ?? ''; + const lang = form.get('lang')?.toString().trim() ?? ''; + const from = parseInt(form.get('from')?.toString() ?? '1', 10); + const to = parseInt(form.get('to')?.toString() ?? '1', 10); + + if (!slug || !lang) { + return { success: false, error: 'slug and lang are required' }; + } + if (isNaN(from) || isNaN(to) || from < 1 || to < from) { + return { success: false, error: 'Invalid chapter range' }; + } + + try { + const res = await backendFetch('/api/admin/translation/bulk', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ slug, lang, from, to }) + }); + if (!res.ok) { + const body = await res.text().catch(() => ''); + log.error('admin/translation', 'bulk enqueue failed', { status: res.status, body }); + return { success: false, error: `Backend error ${res.status}: ${body}` }; + } + const data = await res.json(); + return { success: true, enqueued: data.enqueued as number }; + } catch (e) { + log.error('admin/translation', 'bulk enqueue fetch error', { err: String(e) }); + return { success: false, error: String(e) }; + } + } +}; diff --git a/ui/src/routes/admin/translation/+page.svelte b/ui/src/routes/admin/translation/+page.svelte new file mode 100644 index 0000000..c4153e6 --- /dev/null +++ b/ui/src/routes/admin/translation/+page.svelte @@ -0,0 +1,340 @@ + + + + Translation — Admin + + +
+ +
+

Machine Translation

+

+ {stats.total} job{stats.total !== 1 ? 's' : ''} · + {stats.done} done + {#if stats.failed > 0} + · {stats.failed} failed + {/if} + {#if stats.inFlight > 0} + · {stats.inFlight} in-flight + {/if} +

+
+ + +
+ + +
+ + + {#if activeTab === 'enqueue'} +
+ + {#if form?.success} +
+ Enqueued {form.enqueued} translation job{form.enqueued !== 1 ? 's' : ''} successfully. +
+ {:else if form?.error} +
+ {form.error} +
+ {/if} + +
{ + submitting = true; + return async ({ update }) => { + await update(); + submitting = false; + activeTab = 'jobs'; + }; + }} + class="space-y-4" + > + +
+ + + + {#each data.books as book} + + {/each} + +
+ + +
+ + +
+ + +
+
+ + +
+
+ + +
+
+ +

+ Enqueues {Math.max(0, toInput - fromInput + 1)} task{toInput - fromInput + 1 !== 1 ? 's' : ''} — one per chapter. Max 1000 at a time. +

+ + +
+
+ {/if} + + + {#if activeTab === 'jobs'} + + + {#if filteredJobs.length === 0} +

+ {jobsQ.trim() ? 'No matching jobs.' : 'No translation jobs yet.'} +

+ {:else} + + + + +
+ {#each filteredJobs as job} +
+
+ + {job.slug} + + {job.status} +
+
+ Chapter{job.chapter} + Lang{job.lang} + Started{fmtDate(job.started)} + Duration{duration(job.started, job.finished)} +
+ {#if job.error_message} +

{job.error_message}

+ {/if} +
+ {/each} +
+ {/if} + {/if} +
diff --git a/ui/src/routes/books/[slug]/chapters/[n]/+page.server.ts b/ui/src/routes/books/[slug]/chapters/[n]/+page.server.ts index ce8f2a5..21aa89d 100644 --- a/ui/src/routes/books/[slug]/chapters/[n]/+page.server.ts +++ b/ui/src/routes/books/[slug]/chapters/[n]/+page.server.ts @@ -6,6 +6,8 @@ import { log } from '$lib/server/logger'; import { backendFetch } from '$lib/server/scraper'; import type { Voice } from '$lib/types'; +const SUPPORTED_LANGS = new Set(['ru', 'id', 'pt', 'fr']); + export const load: PageServerLoad = async ({ params, url, locals }) => { const { slug } = params; const n = parseInt(params.n, 10); @@ -15,6 +17,8 @@ export const load: PageServerLoad = async ({ params, url, locals }) => { const isPreview = url.searchParams.get('preview') === '1'; const chapterUrl = url.searchParams.get('chapter_url') ?? ''; const chapterTitle = url.searchParams.get('title') ?? ''; + const lang = url.searchParams.get('lang') ?? ''; + const useTranslation = SUPPORTED_LANGS.has(lang); if (isPreview) { // ── Preview path: scrape chapter live, nothing from PocketBase/MinIO ── @@ -77,7 +81,9 @@ export const load: PageServerLoad = async ({ params, url, locals }) => { next: null as number | null, chapters: [] as { number: number; title: string }[], sessionId: locals.sessionId, - isPreview: true + isPreview: true, + lang: '', + translationStatus: 'unavailable' as string }; } @@ -105,7 +111,37 @@ export const load: PageServerLoad = async ({ params, url, locals }) => { // Non-critical — UI will use store default } - // Fetch chapter markdown directly from the backend (server-side MinIO read) + // ── Translation path: try to serve translated HTML ───────────────────── + if (useTranslation) { + try { + const tRes = await backendFetch( + `/api/translation/${encodeURIComponent(slug)}/${n}?lang=${lang}` + ); + if (tRes.ok) { + const tData = (await tRes.json()) as { html: string; lang: string }; + const prevChapter = chapters.find((c) => c.number === n - 1) ?? null; + const nextChapter = chapters.find((c) => c.number === n + 1) ?? null; + return { + book: { slug: book.slug, title: book.title, cover: book.cover ?? '' }, + chapter: chapterIdx, + html: tData.html, + voices, + prev: prevChapter ? prevChapter.number : null, + next: nextChapter ? nextChapter.number : null, + chapters: chapters.map((c) => ({ number: c.number, title: c.title })), + sessionId: locals.sessionId, + isPreview: false, + lang, + translationStatus: 'done' + }; + } + // 404 = not generated yet — fall through to original, UI can trigger generation + } catch { + // Non-critical — fall through to original content + } + } + + // ── Original content path ────────────────────────────────────────────── let html = ''; try { const res = await backendFetch(`/api/chapter-markdown/${encodeURIComponent(slug)}/${n}`); @@ -122,6 +158,22 @@ export const load: PageServerLoad = async ({ params, url, locals }) => { error(502, 'Could not fetch chapter content'); } + // Check translation status for the UI switcher (non-blocking) + let translationStatus = 'idle'; + if (useTranslation) { + try { + const stRes = await backendFetch( + `/api/translation/status/${encodeURIComponent(slug)}/${n}?lang=${lang}` + ); + if (stRes.ok) { + const stData = (await stRes.json()) as { status: string }; + translationStatus = stData.status ?? 'idle'; + } + } catch { + // Non-critical + } + } + const prevChapter = chapters.find((c) => c.number === n - 1) ?? null; const nextChapter = chapters.find((c) => c.number === n + 1) ?? null; @@ -134,6 +186,8 @@ export const load: PageServerLoad = async ({ params, url, locals }) => { next: nextChapter ? nextChapter.number : null, chapters: chapters.map((c) => ({ number: c.number, title: c.title })), sessionId: locals.sessionId, - isPreview: false + isPreview: false, + lang: useTranslation ? lang : '', + translationStatus }; }; diff --git a/ui/src/routes/books/[slug]/chapters/[n]/+page.svelte b/ui/src/routes/books/[slug]/chapters/[n]/+page.svelte index 2150ffb..96d43a3 100644 --- a/ui/src/routes/books/[slug]/chapters/[n]/+page.svelte +++ b/ui/src/routes/books/[slug]/chapters/[n]/+page.svelte @@ -1,5 +1,7 @@