From 8093cc1e40436e859d87d73b9fd44e4b5fde36df Mon Sep 17 00:00:00 2001 From: lordpietre Date: Wed, 1 Apr 2026 03:29:09 +0000 Subject: [PATCH] feat: add remote worker system for distributed translation - Add remote translator worker (Docker) that connects via WebSocket - Add backend Go handlers for remote worker CRUD (REST + WebSocket) - Add frontend Admin panel for managing remote workers - Add ProtectedRoute component to secure admin routes - Add PostgreSQL schema for remote_workers table - Fix duplicate Load function in config.go - Fix frontend duplicate fetchStatus function - Fix Dockerfiles missing go.sum (use go mod download + tidy) Features: - Workers connect to server via WebSocket (no ports needed) - API key authentication per worker - Job assigner automatically distributes translation jobs - Worker status tracking (online/offline/disabled) - Timeout handling for stuck jobs --- Dockerfile.discovery | 6 +- Dockerfile.qdrant | 6 +- Dockerfile.related | 6 +- Dockerfile.remote-worker | 50 ++ Dockerfile.scraper | 6 +- Dockerfile.topics | 6 +- Dockerfile.wiki | 6 +- backend/Dockerfile | 3 +- backend/cmd/server/main.go | 49 +- backend/internal/config/config.go | 39 +- backend/internal/handlers/remote_worker.go | 187 ++++++ backend/internal/handlers/worker_ws.go | 297 ++++++++++ backend/internal/middleware/auth.go | 52 +- backend/internal/models/remote_worker.go | 41 ++ docker-compose.yml | 32 ++ frontend/src/App.tsx | 23 +- frontend/src/components/ProtectedRoute.tsx | 47 ++ frontend/src/pages/AdminWorkers.tsx | 212 ++++++- init-db/35-remote-workers.sql | 31 + remote-worker.md | 635 +++++++++++++++++++++ workers/remote_translator_worker.py | 389 +++++++++++++ 21 files changed, 2065 insertions(+), 58 deletions(-) create mode 100644 Dockerfile.remote-worker create mode 100644 backend/internal/handlers/remote_worker.go create mode 100644 backend/internal/handlers/worker_ws.go create mode 100644 backend/internal/models/remote_worker.go create mode 100644 frontend/src/components/ProtectedRoute.tsx create mode 100644 init-db/35-remote-workers.sql create mode 100644 remote-worker.md create mode 100644 workers/remote_translator_worker.py diff --git a/Dockerfile.discovery b/Dockerfile.discovery index 90e405d..d31dbbc 100644 --- a/Dockerfile.discovery +++ b/Dockerfile.discovery @@ -6,11 +6,15 @@ RUN apk add --no-cache git WORKDIR /app -COPY backend/go.mod backend/go.sum ./ +COPY backend/go.mod ./ RUN go mod download COPY backend/ ./ +RUN go mod tidy + +COPY backend/ ./ + RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/discovery ./cmd/discovery FROM alpine:3.19 diff --git a/Dockerfile.qdrant b/Dockerfile.qdrant index e80bfae..f67c02b 100644 --- a/Dockerfile.qdrant +++ b/Dockerfile.qdrant @@ -6,11 +6,15 @@ RUN apk add --no-cache git WORKDIR /app -COPY backend/go.mod backend/go.sum ./ +COPY backend/go.mod ./ RUN go mod download COPY backend/ ./ +RUN go mod tidy + +COPY backend/ ./ + RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/qdrant-worker ./cmd/qdrant FROM alpine:3.19 diff --git a/Dockerfile.related b/Dockerfile.related index 12e011d..a5edc77 100644 --- a/Dockerfile.related +++ b/Dockerfile.related @@ -6,11 +6,15 @@ RUN apk add --no-cache git WORKDIR /app -COPY backend/go.mod backend/go.sum ./ +COPY backend/go.mod ./ RUN go mod download COPY backend/ ./ +RUN go mod tidy + +COPY backend/ ./ + RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/related ./cmd/related FROM alpine:3.19 diff --git a/Dockerfile.remote-worker b/Dockerfile.remote-worker new file mode 100644 index 0000000..44fcb11 --- /dev/null +++ b/Dockerfile.remote-worker @@ -0,0 +1,50 @@ +FROM python:3.11-slim-bookworm + +SHELL ["/bin/bash", "-c"] + +RUN apt-get update && apt-get install -y --no-install-recommends \ + patchelf libpq-dev gcc git curl wget bash \ + && rm -rf /var/lib/apt/lists/* + +ENV PYTHONUNBUFFERED=1 \ + PIP_DISABLE_PIP_VERSION_CHECK=1 \ + TOKENIZERS_PARALLELISM=false \ + HF_HOME=/root/.cache/huggingface \ + DEBIAN_FRONTEND=noninteractive + +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir --upgrade pip + +RUN pip install --no-cache-dir torch==2.1.0 torchvision==0.16.0 --index-url https://download.pytorch.org/whl/cpu + +RUN pip install --no-cache-dir \ + ctranslate2==3.24.0 \ + sentencepiece \ + transformers==4.36.0 \ + protobuf==3.20.3 \ + "numpy<2" \ + psycopg2-binary \ + langdetect \ + websocket-client + +RUN find /usr/local/lib/python3.11/site-packages/ctranslate2* \ + -name "libctranslate2-*.so.*" -o -name "libctranslate2.so*" | \ + xargs -I {} patchelf --clear-execstack {} || true + +COPY workers/ ./workers/ +COPY init-db/ ./init-db/ +COPY migrations/ ./migrations/ + +ENV DB_HOST=db +ENV DB_PORT=5432 +ENV DB_NAME=rss +ENV DB_USER=rss +ENV DB_PASS=x + +ENV WORKER_SERVER=ws://localhost:8080/ws/worker +ENV CT2_DEVICE=cpu +ENV CT2_COMPUTE_TYPE=int8 + +CMD ["python", "-m", "workers.remote_translator_worker"] \ No newline at end of file diff --git a/Dockerfile.scraper b/Dockerfile.scraper index 02d380e..587600d 100644 --- a/Dockerfile.scraper +++ b/Dockerfile.scraper @@ -6,13 +6,17 @@ RUN apk add --no-cache git WORKDIR /app -COPY backend/go.mod backend/go.sum ./ +COPY backend/go.mod ./ RUN go mod download COPY backend/ ./ RUN go mod tidy +COPY backend/ ./ + +RUN go mod tidy + RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/scraper ./cmd/scraper FROM alpine:3.19 diff --git a/Dockerfile.topics b/Dockerfile.topics index fc82ea7..1178211 100644 --- a/Dockerfile.topics +++ b/Dockerfile.topics @@ -6,11 +6,15 @@ RUN apk add --no-cache git WORKDIR /app -COPY backend/go.mod backend/go.sum ./ +COPY backend/go.mod ./ RUN go mod download COPY backend/ ./ +RUN go mod tidy + +COPY backend/ ./ + RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/topics ./cmd/topics FROM alpine:3.19 diff --git a/Dockerfile.wiki b/Dockerfile.wiki index fbd84e0..f59df1c 100644 --- a/Dockerfile.wiki +++ b/Dockerfile.wiki @@ -6,13 +6,17 @@ RUN apk add --no-cache git WORKDIR /app -COPY backend/go.mod backend/go.sum ./ +COPY backend/go.mod ./ RUN go mod download COPY backend/ ./ RUN go mod tidy +COPY backend/ ./ + +RUN go mod tidy + RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/wiki_worker ./cmd/wiki_worker FROM alpine:3.19 diff --git a/backend/Dockerfile b/backend/Dockerfile index 6d232b9..5f04721 100644 --- a/backend/Dockerfile +++ b/backend/Dockerfile @@ -4,10 +4,11 @@ WORKDIR /app RUN apt-get update && apt-get install -y gcc musl-dev git -COPY go.mod go.sum ./ +COPY go.mod ./ RUN go mod download COPY . . +RUN go mod tidy RUN CGO_ENABLED=0 GOOS=linux go build -buildvcs=false -o /server ./cmd/server diff --git a/backend/cmd/server/main.go b/backend/cmd/server/main.go index cf13d80..2f3abf7 100644 --- a/backend/cmd/server/main.go +++ b/backend/cmd/server/main.go @@ -9,6 +9,7 @@ import ( "syscall" "github.com/gin-gonic/gin" + "github.com/rss2/backend/internal/auth" "github.com/rss2/backend/internal/cache" "github.com/rss2/backend/internal/config" "github.com/rss2/backend/internal/db" @@ -74,6 +75,39 @@ func initDB() { INSERT INTO config (key, value) VALUES ('translator_status', 'stopped') ON CONFLICT (key) DO NOTHING `) + + // Crear tabla de remote_workers si no existe + _, err = db.GetPool().Exec(ctx, ` + CREATE TABLE IF NOT EXISTS remote_workers ( + id SERIAL PRIMARY KEY, + name VARCHAR(255) NOT NULL, + api_key VARCHAR(64) UNIQUE NOT NULL, + capabilities VARCHAR(50) DEFAULT 'cpu', + status VARCHAR(20) DEFAULT 'offline', + last_seen TIMESTAMP, + created_at TIMESTAMP DEFAULT NOW() + ) + `) + if err != nil { + log.Printf("Warning: Could not create remote_workers table: %v", err) + } else { + log.Println("Table remote_workers ready") + } + + // Añadir columnas a traducciones para workers remotos + _, err = db.GetPool().Exec(ctx, ` + ALTER TABLE traducciones ADD COLUMN IF NOT EXISTS worker_id INTEGER REFERENCES remote_workers(id) + `) + if err != nil { + log.Printf("Warning: Could not add worker_id column: %v", err) + } + + _, err = db.GetPool().Exec(ctx, ` + ALTER TABLE traducciones ADD COLUMN IF NOT EXISTS assigned_at TIMESTAMP + `) + if err != nil { + log.Printf("Warning: Could not add assigned_at column: %v", err) + } } func main() { @@ -109,7 +143,7 @@ func main() { api := r.Group("/api") { // Serve static images downloaded by wiki_worker - api.StaticFS("/wiki-images", gin.Dir("/app/data/wiki_images", false)) + api.StaticFS("/wiki-images", gin.Dir(cfg.WikiImagesPath, false)) api.POST("/auth/login", handlers.Login) api.POST("/auth/register", handlers.Register) @@ -161,8 +195,17 @@ func main() { admin.POST("/workers/config", handlers.SetWorkerConfig) admin.POST("/workers/start", handlers.StartWorkers) admin.POST("/workers/stop", handlers.StopWorkers) + + admin.GET("/workers/remote", handlers.ListRemoteWorkers) + admin.POST("/workers/remote", handlers.CreateRemoteWorker) + admin.GET("/workers/remote/:id", handlers.GetRemoteWorker) + admin.DELETE("/workers/remote/:id", handlers.DeleteRemoteWorker) + admin.POST("/workers/remote/:id/toggle", handlers.ToggleRemoteWorker) + admin.POST("/workers/remote/:id/regenerate-key", handlers.RegenerateAPIKey) } + r.GET("/ws/worker", handlers.HandleWorkerWS) + auth := api.Group("/auth") auth.Use(middleware.AuthRequired()) { @@ -170,7 +213,7 @@ func main() { } } - middleware.SetJWTSecret(cfg.SecretKey) + auth.SetJWTSecret(cfg.SecretKey) port := cfg.ServerPort addr := fmt.Sprintf(":%s", port) @@ -182,6 +225,8 @@ func main() { } }() + handlers.StartJobAssigner() + quit := make(chan os.Signal, 1) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) <-quit diff --git a/backend/internal/config/config.go b/backend/internal/config/config.go index 269ae7b..f218100 100644 --- a/backend/internal/config/config.go +++ b/backend/internal/config/config.go @@ -20,23 +20,34 @@ type Config struct { DefaultLang string NewsPerPage int RateLimitPerMinute int + DockerComposeDir string + WikiImagesPath string + AllowedOrigins string } func Load() *Config { + secretKey := os.Getenv("SECRET_KEY") + if secretKey == "" { + secretKey = "change-this-secret-key" + } + return &Config{ - ServerPort: getEnv("SERVER_PORT", "8080"), - DatabaseURL: getEnv("DATABASE_URL", "postgres://rss:rss@localhost:5432/rss"), - RedisURL: getEnv("REDIS_URL", "redis://localhost:6379"), - QdrantHost: getEnv("QDRANT_HOST", "localhost"), - QdrantPort: getEnvInt("QDRANT_PORT", 6333), - SecretKey: getEnv("SECRET_KEY", "change-this-secret-key"), - JWTExpiration: getEnvDuration("JWT_EXPIRATION", 24*time.Hour), - TranslationURL: getEnv("TRANSLATION_URL", "http://libretranslate:7790"), - OllamaURL: getEnv("OLLAMA_URL", "http://ollama:11434"), - SpacyURL: getEnv("SPACY_URL", "http://spacy:8000"), - DefaultLang: getEnv("DEFAULT_LANG", "es"), - NewsPerPage: getEnvInt("NEWS_PER_PAGE", 30), - RateLimitPerMinute: getEnvInt("RATE_LIMIT_PER_MINUTE", 60), + ServerPort: getEnv("SERVER_PORT", "8080"), + DatabaseURL: getEnv("DATABASE_URL", "postgres://rss:rss@localhost:5432/rss"), + RedisURL: getEnv("REDIS_URL", "redis://localhost:6379"), + QdrantHost: getEnv("QDRANT_HOST", "localhost"), + QdrantPort: getEnvInt("QDRANT_PORT", 6333), + SecretKey: secretKey, + JWTExpiration: getEnvDuration("JWT_EXPIRATION", 24*time.Hour), + TranslationURL: getEnv("TRANSLATION_URL", "http://libretranslate:7790"), + OllamaURL: getEnv("OLLAMA_URL", "http://ollama:11434"), + SpacyURL: getEnv("SPACY_URL", "http://spacy:8000"), + DefaultLang: getEnv("DEFAULT_LANG", "es"), + NewsPerPage: getEnvInt("NEWS_PER_PAGE", 30), + RateLimitPerMinute: getEnvInt("RATE_LIMIT_PER_MINUTE", 60), + DockerComposeDir: getEnv("DOCKER_COMPOSE_DIR", "/datos/rss2"), + WikiImagesPath: getEnv("WIKI_IMAGES_PATH", "/app/data/wiki_images"), + AllowedOrigins: getEnv("ALLOWED_ORIGINS", "*"), } } @@ -63,4 +74,4 @@ func getEnvDuration(key string, defaultValue time.Duration) time.Duration { } } return defaultValue -} +} \ No newline at end of file diff --git a/backend/internal/handlers/remote_worker.go b/backend/internal/handlers/remote_worker.go new file mode 100644 index 0000000..6e37d17 --- /dev/null +++ b/backend/internal/handlers/remote_worker.go @@ -0,0 +1,187 @@ +package handlers + +import ( + "crypto/rand" + "encoding/hex" + "log" + "net/http" + + "github.com/gin-gonic/gin" + "github.com/rss2/backend/internal/db" + "github.com/rss2/backend/internal/models" +) + +func CreateRemoteWorker(c *gin.Context) { + var req struct { + Name string `json:"name" binding:"required"` + Capabilities string `json:"capabilities"` + } + + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": "Invalid request", "message": err.Error()}) + return + } + + if req.Capabilities == "" { + req.Capabilities = "cpu" + } + + apiKey := generateAPIKey() + + ctx := c.Request.Context() + var workerID int + err := db.GetPool().QueryRow(ctx, ` + INSERT INTO remote_workers (name, api_key, capabilities, status, created_at) + VALUES ($1, $2, $3, 'offline', NOW()) + RETURNING id + `, req.Name, apiKey, req.Capabilities).Scan(&workerID) + + if err != nil { + log.Printf("Error creating remote worker: %v", err) + c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to create worker"}) + return + } + + c.JSON(http.StatusCreated, gin.H{ + "id": workerID, + "name": req.Name, + "api_key": apiKey, + "capabilities": req.Capabilities, + "status": "offline", + }) +} + +func ListRemoteWorkers(c *gin.Context) { + ctx := c.Request.Context() + + rows, err := db.GetPool().Query(ctx, ` + SELECT id, name, api_key, capabilities, status, last_seen, created_at + FROM remote_workers + ORDER BY created_at DESC + `) + if err != nil { + log.Printf("Error listing remote workers: %v", err) + c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to list workers"}) + return + } + defer rows.Close() + + var workers []models.RemoteWorker + for rows.Next() { + var w models.RemoteWorker + if err := rows.Scan(&w.ID, &w.Name, &w.APIKey, &w.Capabilities, &w.Status, &w.LastSeen, &w.CreatedAt); err != nil { + log.Printf("Error scanning worker: %v", err) + continue + } + workers = append(workers, w) + } + + if workers == nil { + workers = []models.RemoteWorker{} + } + + c.JSON(http.StatusOK, workers) +} + +func GetRemoteWorker(c *gin.Context) { + id := c.Param("id") + ctx := c.Request.Context() + + var w models.RemoteWorker + err := db.GetPool().QueryRow(ctx, ` + SELECT id, name, api_key, capabilities, status, last_seen, created_at + FROM remote_workers + WHERE id = $1 + `, id).Scan(&w.ID, &w.Name, &w.APIKey, &w.Capabilities, &w.Status, &w.LastSeen, &w.CreatedAt) + + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Worker not found"}) + return + } + + var stats struct { + JobsCompleted int `json:"jobs_completed"` + JobsPending int `json:"jobs_pending"` + } + + db.GetPool().QueryRow(ctx, ` + SELECT COUNT(*) FROM traducciones WHERE worker_id = $1 AND status = 'done' + `, id).Scan(&stats.JobsCompleted) + + db.GetPool().QueryRow(ctx, ` + SELECT COUNT(*) FROM traducciones WHERE worker_id = $1 AND status = 'pending' + `, id).Scan(&stats.JobsPending) + + c.JSON(http.StatusOK, gin.H{ + "worker": w, + "stats": stats, + }) +} + +func DeleteRemoteWorker(c *gin.Context) { + id := c.Param("id") + ctx := c.Request.Context() + + result, err := db.GetPool().Exec(ctx, ` + DELETE FROM remote_workers WHERE id = $1 + `, id) + + if err != nil { + log.Printf("Error deleting remote worker: %v", err) + c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to delete worker"}) + return + } + + if result.RowsAffected() == 0 { + c.JSON(http.StatusNotFound, gin.H{"error": "Worker not found"}) + return + } + + c.JSON(http.StatusOK, gin.H{"message": "Worker deleted"}) +} + +func ToggleRemoteWorker(c *gin.Context) { + id := c.Param("id") + ctx := c.Request.Context() + + var currentStatus string + err := db.GetPool().QueryRow(ctx, `SELECT status FROM remote_workers WHERE id = $1`, id).Scan(¤tStatus) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "Worker not found"}) + return + } + + newStatus := "disabled" + if currentStatus == "disabled" { + newStatus = "offline" + } + + _, err = db.GetPool().Exec(ctx, `UPDATE remote_workers SET status = $1 WHERE id = $2`, newStatus, id) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to update worker status"}) + return + } + + c.JSON(http.StatusOK, gin.H{"status": newStatus}) +} + +func RegenerateAPIKey(c *gin.Context) { + id := c.Param("id") + ctx := c.Request.Context() + + newAPIKey := generateAPIKey() + + _, err := db.GetPool().Exec(ctx, `UPDATE remote_workers SET api_key = $1 WHERE id = $2`, newAPIKey, id) + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": "Failed to regenerate API key"}) + return + } + + c.JSON(http.StatusOK, gin.H{"api_key": newAPIKey}) +} + +func generateAPIKey() string { + bytes := make([]byte, 32) + rand.Read(bytes) + return hex.EncodeToString(bytes) +} \ No newline at end of file diff --git a/backend/internal/handlers/worker_ws.go b/backend/internal/handlers/worker_ws.go new file mode 100644 index 0000000..15eb698 --- /dev/null +++ b/backend/internal/handlers/worker_ws.go @@ -0,0 +1,297 @@ +package handlers + +import ( + "context" + "encoding/json" + "log" + "net/http" + "sync" + "time" + + "github.com/gin-gonic/gin" + "github.com/gorilla/websocket" + "github.com/rss2/backend/internal/db" + "github.com/rss2/backend/internal/models" +) + +var upgrader = websocket.Upgrader{ + CheckOrigin: func(r *http.Request) bool { return true }, +} + +type WSWorker struct { + conn *websocket.Conn + workerID int + capabilities string + workerName string + lastHeartbeat time.Time + mu sync.Mutex +} + +var ( + workers = make(map[int]*WSWorker) + workersByAPI = make(map[string]*WSWorker) + workersMu sync.RWMutex +) + +func HandleWorkerWS(c *gin.Context) { + apiKey := c.Query("api_key") + if apiKey == "" { + apiKey = c.GetHeader("X-API-Key") + } + + if apiKey == "" { + c.JSON(http.StatusUnauthorized, gin.H{"error": "API key required"}) + return + } + + ctx := c.Request.Context() + + var workerID int + var capabilities string + err := db.GetPool().QueryRow(ctx, ` + SELECT id, capabilities FROM remote_workers + WHERE api_key = $1 AND status != 'disabled' + `, apiKey).Scan(&workerID, &capabilities) + + if err != nil { + c.JSON(http.StatusUnauthorized, gin.H{"error": "Invalid API key"}) + return + } + + conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) + if err != nil { + log.Printf("WebSocket upgrade error: %v", err) + return + } + + wsWorker := &WSWorker{ + conn: conn, + workerID: workerID, + capabilities: capabilities, + lastHeartbeat: time.Now(), + } + + workersMu.Lock() + workers[workerID] = wsWorker + workersByAPI[apiKey] = wsWorker + workersMu.Unlock() + + db.GetPool().Exec(ctx, ` + UPDATE remote_workers SET status = 'online', last_seen = NOW() WHERE id = $1 + `, workerID) + + go heartbeatLoop(wsWorker) + go writeLoop(wsWorker) + readLoop(wsWorker, apiKey) +} + +func readLoop(wsWorker *WSWorker, apiKey string) { + defer func() { + cleanupWorker(wsWorker, apiKey) + }() + + for { + var msg models.WSClientMessage + err := wsWorker.conn.ReadJSON(&msg) + if err != nil { + log.Printf("Worker %d read error: %v", wsWorker.workerID, err) + break + } + + switch msg.Type { + case "register": + wsWorker.mu.Lock() + if msg.WorkerName != "" { + wsWorker.workerName = msg.WorkerName + } + if msg.Capabilities != "" { + wsWorker.capabilities = msg.Capabilities + } + wsWorker.mu.Unlock() + + sendWS(wsWorker, models.WSServerMessage{Type: "ack", Job: nil}) + + case "heartbeat": + wsWorker.mu.Lock() + wsWorker.lastHeartbeat = time.Now() + wsWorker.mu.Unlock() + sendWS(wsWorker, models.WSServerMessage{Type: "ack", Job: nil}) + + case "result": + if msg.Result != nil { + handleTranslationResult(wsWorker.workerID, msg.Result) + } + } + } +} + +func writeLoop(wsWorker *WSWorker) { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + if err := wsWorker.conn.WriteMessage(websocket.PingMessage, nil); err != nil { + return + } + } + } +} + +func heartbeatLoop(wsWorker *WSWorker) { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + wsWorker.mu.Lock() + elapsed := time.Since(wsWorker.lastHeartbeat) + wsWorker.mu.Unlock() + + if elapsed > 60*time.Second { + log.Printf("Worker %d heartbeat timeout", wsWorker.workerID) + wsWorker.conn.Close() + return + } + + sendWS(wsWorker, models.WSServerMessage{Type: "ping", Job: nil}) + } + } +} + +func sendWS(wsWorker *WSWorker, msg models.WSServerMessage) { + data, _ := json.Marshal(msg) + wsWorker.conn.WriteMessage(websocket.TextMessage, data) +} + +func cleanupWorker(wsWorker *WSWorker, apiKey string) { + workersMu.Lock() + delete(workers, wsWorker.workerID) + delete(workersByAPI, apiKey) + workersMu.Unlock() + + ctx := context.Background() + db.GetPool().Exec(ctx, ` + UPDATE remote_workers SET status = 'offline', last_seen = NOW() WHERE id = $1 + `, wsWorker.workerID) + + wsWorker.conn.Close() + + db.GetPool().Exec(ctx, ` + UPDATE traducciones + SET status = 'pending', worker_id = NULL, assigned_at = NULL + WHERE worker_id = $1 AND status = 'assigned' + `, wsWorker.workerID) + + log.Printf("Worker %d disconnected", wsWorker.workerID) +} + +func handleTranslationResult(workerID int, result *models.TranslationResult) { + ctx := context.Background() + + if result.Error != "" { + db.GetPool().Exec(ctx, ` + UPDATE traducciones + SET status = 'error', worker_id = NULL, assigned_at = NULL + WHERE id = $1 + `, result.JobID) + log.Printf("Job %d failed: %s", result.JobID, result.Error) + return + } + + _, err := db.GetPool().Exec(ctx, ` + UPDATE traducciones + SET titulo_trad = $1, resumen_trad = $2, status = 'done', + worker_id = $3, assigned_at = NULL + WHERE id = $4 + `, result.TitleTr, result.SummaryTr, workerID, result.JobID) + + if err != nil { + log.Printf("Error updating translation result: %v", err) + } + + db.GetPool().Exec(ctx, ` + UPDATE remote_workers SET last_seen = NOW() WHERE id = $1 + `, workerID) + + log.Printf("Job %d completed by worker %d", result.JobID, workerID) +} + +func AssignJobToWorker(workerID int) *models.TranslationJob { + ctx := context.Background() + + workersMu.RLock() + wsWorker, exists := workers[workerID] + workersMu.RUnlock() + + if !exists { + return nil + } + + wsWorker.mu.Lock() + workerCapabilities := wsWorker.capabilities + wsWorker.mu.Unlock() + + _ = workerCapabilities + + var job models.TranslationJob + err := db.GetPool().QueryRow(ctx, ` + SELECT t.id, t.noticia_id, t.lang_from, t.lang_to, n.titulo, n.resumen + FROM traducciones t + JOIN noticias n ON n.id = t.noticia_id + WHERE t.status = 'pending' + AND (t.worker_id IS NULL OR t.status = 'assigned' AND t.assigned_at < NOW() - INTERVAL '5 minutes') + AND t.lang_to = 'es' + AND (t.titulo_trad IS NULL OR t.resumen_trad IS NULL) + ORDER BY n.fecha DESC + LIMIT 1 + FOR UPDATE SKIP LOCKED + `).Scan(&job.ID, &job.NewsID, &job.LangFrom, &job.LangTo, &job.Title, &job.Summary) + + if err != nil { + return nil + } + + db.GetPool().Exec(ctx, ` + UPDATE traducciones + SET status = 'assigned', worker_id = $1, assigned_at = NOW() + WHERE id = $2 + `, workerID, job.ID) + + sendWS(wsWorker, models.WSServerMessage{ + Type: "job", + Job: &job, + }) + + return &job +} + +func StartJobAssigner() { + go func() { + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + workersMu.RLock() + for workerID := range workers { + workersMu.RUnlock() + AssignJobToWorker(workerID) + workersMu.RLock() + } + workersMu.RUnlock() + + ctx := context.Background() + db.GetPool().Exec(ctx, ` + UPDATE traducciones + SET status = 'pending', worker_id = NULL, assigned_at = NULL + WHERE status = 'assigned' + AND assigned_at < NOW() - INTERVAL '10 minutes' + `) + } + } + }() +} \ No newline at end of file diff --git a/backend/internal/middleware/auth.go b/backend/internal/middleware/auth.go index 57df7c7..4c35c35 100644 --- a/backend/internal/middleware/auth.go +++ b/backend/internal/middleware/auth.go @@ -5,23 +5,10 @@ import ( "strings" "github.com/gin-gonic/gin" - "github.com/golang-jwt/jwt/v5" + "github.com/rss2/backend/internal/auth" + "github.com/rss2/backend/internal/config" ) -var jwtSecret []byte - -func SetJWTSecret(secret string) { - jwtSecret = []byte(secret) -} - -type Claims struct { - UserID int64 `json:"user_id"` - Email string `json:"email"` - Username string `json:"username"` - IsAdmin bool `json:"is_admin"` - jwt.RegisteredClaims -} - func AuthRequired() gin.HandlerFunc { return func(c *gin.Context) { authHeader := c.GetHeader("Authorization") @@ -38,12 +25,8 @@ func AuthRequired() gin.HandlerFunc { return } - claims := &Claims{} - token, err := jwt.ParseWithClaims(tokenString, claims, func(token *jwt.Token) (interface{}, error) { - return jwtSecret, nil - }) - - if err != nil || !token.Valid { + claims, err := auth.ValidateToken(tokenString) + if err != nil { c.JSON(http.StatusUnauthorized, gin.H{"error": "Invalid token"}) c.Abort() return @@ -63,7 +46,7 @@ func AdminRequired() gin.HandlerFunc { return } - claims := userVal.(*Claims) + claims := userVal.(*auth.Claims) if !claims.IsAdmin { c.JSON(http.StatusForbidden, gin.H{"error": "Admin access required"}) c.Abort() @@ -75,9 +58,30 @@ func AdminRequired() gin.HandlerFunc { } func CORSMiddleware() gin.HandlerFunc { + cfg := config.Load() + allowedOrigins := strings.Split(cfg.AllowedOrigins, ",") + return func(c *gin.Context) { - c.Writer.Header().Set("Access-Control-Allow-Origin", "*") - c.Writer.Header().Set("Access-Control-Allow-Credentials", "true") + origin := c.Request.Header.Get("Origin") + + allowed := false + for _, o := range allowedOrigins { + o = strings.TrimSpace(o) + if o == "*" || o == origin { + allowed = true + break + } + } + + if allowed { + if cfg.AllowedOrigins == "*" { + c.Writer.Header().Set("Access-Control-Allow-Origin", "*") + } else { + c.Writer.Header().Set("Access-Control-Allow-Origin", origin) + c.Writer.Header().Set("Access-Control-Allow-Credentials", "true") + } + } + c.Writer.Header().Set("Access-Control-Allow-Headers", "Content-Type, Content-Length, Accept-Encoding, X-CSRF-Token, Authorization, accept, origin, Cache-Control, X-Requested-With") c.Writer.Header().Set("Access-Control-Allow-Methods", "POST, OPTIONS, GET, PUT, DELETE, PATCH") diff --git a/backend/internal/models/remote_worker.go b/backend/internal/models/remote_worker.go new file mode 100644 index 0000000..ce77213 --- /dev/null +++ b/backend/internal/models/remote_worker.go @@ -0,0 +1,41 @@ +package models + +import "time" + +type RemoteWorker struct { + ID int `json:"id"` + Name string `json:"name"` + APIKey string `json:"api_key,omitempty"` + Capabilities string `json:"capabilities"` + Status string `json:"status"` + LastSeen *time.Time `json:"last_seen,omitempty"` + CreatedAt time.Time `json:"created_at"` +} + +type TranslationJob struct { + ID int64 `json:"id"` + NewsID int64 `json:"noticia_id"` + LangFrom string `json:"lang_from"` + LangTo string `json:"lang_to"` + Title string `json:"title"` + Summary string `json:"summary"` +} + +type TranslationResult struct { + JobID int64 `json:"job_id"` + TitleTr string `json:"title_trad"` + SummaryTr string `json:"resumen_trad"` + Error string `json:"error,omitempty"` +} + +type WSClientMessage struct { + Type string `json:"type"` + Capabilities string `json:"capabilities,omitempty"` + WorkerName string `json:"worker_name,omitempty"` + Result *TranslationResult `json:"result,omitempty"` +} + +type WSServerMessage struct { + Type string `json:"type"` + Job *TranslationJob `json:"job,omitempty"` +} \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml index f022eed..e90b51b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -312,6 +312,38 @@ services: condition: service_healthy restart: unless-stopped + # ================================================================================== + # REMOTE TRANSLATOR WORKER - Worker remoto que se conecta por WebSocket + # ================================================================================== + remote-translator: + build: + context: . + dockerfile: Dockerfile.remote-worker + container_name: rss2_remote_translator + environment: + WORKER_NAME: ${WORKER_NAME:-remote-worker-1} + WORKER_API_KEY: ${WORKER_API_KEY:-} + WORKER_SERVER: ${WORKER_SERVER:-ws://backend-go:8080/ws/worker} + CT2_DEVICE: ${CT2_DEVICE:-cpu} + CT2_MODEL_PATH: /app/models/nllb-ct2 + CT2_COMPUTE_TYPE: ${CT2_COMPUTE_TYPE:-int8} + UNIVERSAL_MODEL: facebook/nllb-200-distilled-600M + HF_HOME: /app/hf_cache + TZ: Europe/Madrid + volumes: + - ./hf_cache:/app/hf_cache + - ./models:/app/models + networks: + - backend + profiles: + - remote + deploy: + resources: + limits: + cpus: '1' + memory: 2G + restart: unless-stopped + # ================================================================================== # TRANSLATION SCHEDULER - Creates translation jobs # ================================================================================== diff --git a/frontend/src/App.tsx b/frontend/src/App.tsx index 38af056..3f3c3d3 100644 --- a/frontend/src/App.tsx +++ b/frontend/src/App.tsx @@ -1,6 +1,7 @@ import { useState, useEffect } from 'react' import { Routes, Route } from 'react-router-dom' import { Layout } from './components/layout/Layout' +import { ProtectedRoute } from './components/ProtectedRoute' import { Home } from './pages/Home' import { News } from './pages/News' import { Feeds } from './pages/Feeds' @@ -64,13 +65,25 @@ function App() { } /> } /> } /> - } /> - } /> - } /> - } /> + {}} + /> + {}} + /> + {}} + /> + {}} + /> ) } -export default App +export default App \ No newline at end of file diff --git a/frontend/src/components/ProtectedRoute.tsx b/frontend/src/components/ProtectedRoute.tsx new file mode 100644 index 0000000..5bca3f1 --- /dev/null +++ b/frontend/src/components/ProtectedRoute.tsx @@ -0,0 +1,47 @@ +import { Navigate } from 'react-router-dom' +import { useState, useEffect } from 'react' + +interface ProtectedRouteProps { + children: React.ReactNode + requireAdmin?: boolean +} + +export function ProtectedRoute({ children, requireAdmin = false }: ProtectedRouteProps) { + const [isAuthenticated, setIsAuthenticated] = useState(null) + const [isAdmin, setIsAdmin] = useState(false) + + useEffect(() => { + const token = localStorage.getItem('token') + const userData = localStorage.getItem('user') + + if (token && userData) { + try { + const user = JSON.parse(userData) + setIsAuthenticated(true) + setIsAdmin(user.is_admin === true) + } catch { + setIsAuthenticated(false) + } + } else { + setIsAuthenticated(false) + } + }, []) + + if (isAuthenticated === null) { + return ( +
+
+
+ ) + } + + if (!isAuthenticated) { + return + } + + if (requireAdmin && !isAdmin) { + return + } + + return <>{children} +} \ No newline at end of file diff --git a/frontend/src/pages/AdminWorkers.tsx b/frontend/src/pages/AdminWorkers.tsx index 16dc347..cd5e0f2 100644 --- a/frontend/src/pages/AdminWorkers.tsx +++ b/frontend/src/pages/AdminWorkers.tsx @@ -1,6 +1,6 @@ import { useState, useEffect } from 'react' import { api } from '../services/api' -import { Cpu, Play, Square, Settings, RefreshCw, Loader2, Server } from 'lucide-react' +import { Cpu, Play, Square, Settings, RefreshCw, Loader2, Server, Wifi, Plus, Trash2, Key } from 'lucide-react' interface WorkerStatus { type: string @@ -9,6 +9,16 @@ interface WorkerStatus { running: number } +interface RemoteWorker { + id: number + name: string + api_key?: string + capabilities: string + status: string + last_seen?: string + created_at: string +} + export function AdminWorkers() { const [status, setStatus] = useState(null) const [loading, setLoading] = useState(true) @@ -16,24 +26,50 @@ export function AdminWorkers() { const [configType, setConfigType] = useState<'cpu' | 'gpu'>('cpu') const [configWorkers, setConfigWorkers] = useState(2) - useEffect(() => { - fetchStatus() - }, []) + const [remoteWorkers, setRemoteWorkers] = useState([]) + const [showAddForm, setShowAddForm] = useState(false) + const [newWorkerName, setNewWorkerName] = useState('') + const [newWorkerCapabilities, setNewWorkerCapabilities] = useState('cpu') + const [newWorkerKey, setNewWorkerKey] = useState(null) const fetchStatus = async () => { setLoading(true) try { const res = await api.get('/admin/workers/status') + console.log('Status response:', res.data) setStatus(res.data) setConfigType(res.data.type) setConfigWorkers(res.data.workers) - } catch (err) { - console.error(err) + } catch (err: any) { + console.error('Error fetching status:', err) } finally { setLoading(false) } } + const fetchRemoteWorkers = async () => { + try { + const res = await api.get('/admin/workers/remote') + console.log('Remote workers response:', res.data) + setRemoteWorkers(res.data) + } catch (err: any) { + console.error('Error fetching remote workers:', err) + // If 401, user is not authenticated - try refreshing + if (err.response?.status === 401) { + const userData = localStorage.getItem('user') + if (userData) { + const user = JSON.parse(userData) + console.log('Current user from localStorage:', user) + } + } + } + } + + useEffect(() => { + fetchStatus() + fetchRemoteWorkers() + }, []) + const handleStart = async () => { setActionLoading('start') try { @@ -76,6 +112,53 @@ export function AdminWorkers() { } } + const handleAddRemoteWorker = async () => { + if (!newWorkerName.trim()) return + setActionLoading('add') + try { + const res = await api.post('/admin/workers/remote', { + name: newWorkerName, + capabilities: newWorkerCapabilities + }) + setNewWorkerKey(res.data.api_key) + setNewWorkerName('') + fetchRemoteWorkers() + } catch (err: any) { + alert(err.response?.data?.error || 'Error al crear worker') + } finally { + setActionLoading(null) + } + } + + const handleDeleteRemoteWorker = async (id: number) => { + if (!confirm('¿Eliminar este worker remoto?')) return + try { + await api.delete(`/admin/workers/remote/${id}`) + fetchRemoteWorkers() + } catch (err: any) { + alert(err.response?.data?.error || 'Error al eliminar worker') + } + } + + const handleToggleRemoteWorker = async (id: number) => { + try { + await api.post(`/admin/workers/remote/${id}/toggle`) + fetchRemoteWorkers() + } catch (err: any) { + alert(err.response?.data?.error || 'Error al cambiar estado') + } + } + + const handleRegenerateKey = async (id: number) => { + if (!confirm('¿Regenerar API key? La anterior deixará de funcionar.')) return + try { + const res = await api.post(`/admin/workers/remote/${id}/regenerate-key`) + setNewWorkerKey(res.data.api_key) + } catch (err: any) { + alert(err.response?.data?.error || 'Error al regenerar key') + } + } + if (loading) { return (
@@ -256,6 +339,123 @@ export function AdminWorkers() {
  • • Recomendado: 2-4 workers para CPU, 1-2 para GPU
  • + +
    +
    +

    + + Workers Remotos +

    + +
    + + {showAddForm && ( +
    +

    Nuevo Worker Remoto

    +
    + setNewWorkerName(e.target.value)} + className="flex-1 px-3 py-2 border rounded-lg dark:bg-gray-600 dark:border-gray-600" + /> + +
    +
    + + +
    +
    + )} + + {newWorkerKey && ( +
    +
    + + API Key creada +
    + + {newWorkerKey} + +

    + Copia esta clave. No se mostrará de nuevo. +

    + +
    + )} + + {remoteWorkers.length === 0 ? ( +

    No hay workers remotos configurados

    + ) : ( +
    + {remoteWorkers.map((worker) => ( +
    +
    +
    +
    + {worker.name} + {worker.capabilities} +
    +
    +
    + + + +
    +
    + ))} +
    + )} +
    ) } diff --git a/init-db/35-remote-workers.sql b/init-db/35-remote-workers.sql new file mode 100644 index 0000000..625be13 --- /dev/null +++ b/init-db/35-remote-workers.sql @@ -0,0 +1,31 @@ +-- Tabla de workers remotos +CREATE TABLE IF NOT EXISTS remote_workers ( + id SERIAL PRIMARY KEY, + name VARCHAR(255) NOT NULL, + api_key VARCHAR(64) UNIQUE NOT NULL, + capabilities VARCHAR(50) DEFAULT 'cpu', + status VARCHAR(20) DEFAULT 'offline', + last_seen TIMESTAMP, + created_at TIMESTAMP DEFAULT NOW() +); + +-- Añadir columnas a traducciones para workers remotos +ALTER TABLE traducciones ADD COLUMN IF NOT EXISTS worker_id INTEGER REFERENCES remote_workers(id); +ALTER TABLE traducciones ADD COLUMN IF NOT EXISTS assigned_at TIMESTAMP; + +-- Índices para workers remotos +CREATE INDEX IF NOT EXISTS idx_remote_workers_status ON remote_workers(status); +CREATE INDEX IF NOT EXISTS idx_remote_workers_api_key ON remote_workers(api_key); + +-- Índice para buscar jobs disponibles por capabilities +CREATE INDEX IF NOT EXISTS idx_traducciones_pending_worker + ON traducciones(lang_to, status, worker_id) + WHERE status = 'pending' AND worker_id IS NULL; + +-- Índice para jobs asignados (para cleanup de timeout) +CREATE INDEX IF NOT EXISTS idx_traducciones_assigned_timeout + ON traducciones(assigned_at) + WHERE status = 'assigned' AND assigned_at IS NOT NULL; + +-- Tabla de workers remotos conectados (en memoria = no persistir, pero referenciar) +-- El status se actualiza desde el WebSocket handler \ No newline at end of file diff --git a/remote-worker.md b/remote-worker.md new file mode 100644 index 0000000..b39cfc0 --- /dev/null +++ b/remote-worker.md @@ -0,0 +1,635 @@ +# Plan: Sistema de Workers Remotos para Traducción + +## Visión General + +Sistema que permite conectar workers de traducción desde equipos externos (con GPU) al servidor central mediante WebSockets, eliminando la necesidad de exponer el servidor a internet y permitiendo procesamiento distribuido. + +--- + +## Arquitectura del Sistema + +``` +┌─────────────────────────────────────────────────────────────────────────────┐ +│ SERVidor VPS │ +│ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ ┌─────────────────┐ │ +│ │ API REST │ │ WebSocket │ │ Scheduler │ │ PostgreSQL │ │ +│ │ (Go) │◄─┤ Server │ │ (Go) │ │ (Cola Jobs) │ │ +│ │ │ │ (Go) │ │ │ │ │ │ +│ └─────────────┘ └─────────────┘ └─────────────┘ └─────────────────┘ │ +│ │ │ │ │ +│ └────────────────┼────────────────┘ │ +│ │ │ +│ PUERTO 80/443 (ya expuesto) │ +└─────────────────────────────────────────────────────────────────────────────┘ + ▲ + │ WebSocket (conexión saliente) + │ API Key autenticación + │ +┌─────────────────────────────────────────────────────────────────────────────┐ +│ EQUIPO LOCAL CON GPU │ +│ ┌─────────────────────────────────────────────────────────────────┐ │ +│ │ Worker Client (Go) │ │ +│ │ ┌─────────────┐ ┌─────────────┐ ┌─────────────────────────┐ │ │ +│ │ │ WS Client │ │ Lógica │ │ CTranslate2 (subproc) │ │ │ +│ │ │ (Go) │ │ (Go) │ │ + Modelo NLLB-1.3B │ │ │ +│ │ └─────────────┘ └─────────────┘ └─────────────────────────┘ │ │ +│ └─────────────────────────────────────────────────────────────────┘ │ +│ │ +│ Conexión saliente (no requiere puertos abiertos) │ +└─────────────────────────────────────────────────────────────────────────────┘ +``` + +### Flujo de Datos + +1. **Scheduler (servidor)** crea jobs en `traducciones` con status `pending` +2. **Worker remoto** se conecta por WS, envía capabilities (`gpu`/`cpu`) +3. **Servidor** asigna job disponible que coincida con capabilities → status `assigned` +4. **Worker** recibe job → traduce con NLLB-1.3B → envía resultado por WS +5. **Servidor** actualiza `traducciones` con resultado → status `completed` + +--- + +## Análisis de la Aplicación Actual + +### Componentes Relevantes Identificados + +| Componente | Ubicación | Función | +|------------|-----------|---------| +| Handlers admin | `backend/internal/handlers/admin.go` | API workers locales (Docker) | +| Modelos | `backend/internal/models/models.go` | Estructuras de datos | +| Tabla traducciones | `init-db/00-complete-schema.sql` | Cola de traducciones | +| Worker Python | `workers/ctranslator_worker.py` | Traducción actual | +| Scheduler Python | `workers/translation_scheduler.py` | Creador de jobs | +| Frontend Admin | `frontend/src/pages/AdminWorkers.tsx` | Panel de control | + +### Estados Actuales de traducción + +La tabla `traducciones` tiene: +- `pending` - disponible +- `done` - completado +- `error` - fallido + +**Problema identificado:** No hay estado `assigned` → varios workers pueden tomar el mismo job. + +### Query Actual del Worker + +```python +# workers/ctranslator_worker.py:361-378 +SELECT ... FROM traducciones t +WHERE t.lang_to = %s + AND (t.titulo_trad IS NULL OR t.resumen_trad IS NULL) + AND (t.locked_at IS NULL OR t.locked_at < NOW() - INTERVAL '10 minutes') +ORDER BY n.fecha DESC +LIMIT %s +FOR UPDATE SKIP LOCKED +``` + +Usa `locked_at` como mecanismo de locking, pero no actualiza un `status`. + +--- + +## Decisiones de Diseño + +### 1. Conexión: WebSocket + +- Bidireccional, tiempo real +- El worker solo necesita conexión saliente (como navegación web) +- No requiere exponer puertos adicionales en el VPS + +### 2. Autenticación: API Key + +- El servidor genera una API key única por worker +- El worker la envía en cada conexión WS o como header +- Más simple que usuario/password + +### 3. Cola de Trabajos: PostgreSQL + +- Persistencia estable (no Redis, para simplificar) +- El servidor escribe resultados (más control que acceso directo del worker) +- Si el worker falla, el servidor puede reasignar el job + +### 4. ML en Worker: CTranslate2 + +- El worker llama a CTranslate2 como subproceso +- Mismo mecanismo que el worker actual (Docker) pero standalone +- No requiere bindings Go nativos + +--- + +## Problemas Potenciales y Soluciones + +| Problema | Severity | Solución | +|----------|----------|----------| +| Worker sin internet | N/A | No es problema: el worker inicia conexión saliente | +| Worker se desconecta durante procesamiento | Alta | Timeout en servidor → job vuelve a `pending` | +| Job asignado a dos workers | Alta | Nuevo estado `assigned` + columna `worker_id` | +| Worker escribe directamente a BD | Media | Enviar resultados por WS → servidor controla escritura | +| Modelo NLLB no disponible localmente | Baja | Descargar modelo convertido previamente | + +--- + +## Plan de Implementación + +### Fase 1: Base de Datos + +**Archivo:** `init-db/35-remote-workers.sql` + +```sql +-- Tabla de workers remotos +CREATE TABLE IF NOT EXISTS remote_workers ( + id SERIAL PRIMARY KEY, + name VARCHAR(255) NOT NULL, + api_key VARCHAR(64) UNIQUE NOT NULL, + capabilities VARCHAR(50) DEFAULT 'cpu', -- 'cpu' o 'gpu' + status VARCHAR(20) DEFAULT 'offline', -- 'online', 'offline', 'disabled' + last_seen TIMESTAMP, + created_at TIMESTAMP DEFAULT NOW() +); + +-- Añadir columnas a traducciones para workers remotos +ALTER TABLE traducciones ADD COLUMN IF NOT EXISTS worker_id INTEGER REFERENCES remote_workers(id); +ALTER TABLE traducciones ADD COLUMN IF NOT EXISTS assigned_at TIMESTAMP; + +-- Nuevo índice para buscar jobs disponibles por capabilities +CREATE INDEX IF NOT EXISTS idx_traducciones_pending_worker + ON traducciones(lang_to, status, worker_id) + WHERE status = 'pending' AND worker_id IS NULL; + +-- Índice para jobs asignados (para cleanup de timeout) +CREATE INDEX IF NOT EXISTS idx_traducciones_assigned_timeout + ON traducciones(assigned_at) + WHERE status = 'assigned' AND assigned_at IS NOT NULL; +``` + +### Fase 2: Modelos (Backend Go) + +**Nuevo archivo:** `backend/internal/models/worker.go` + +```go +package models + +import "time" + +// Worker remoto registrado +type RemoteWorker struct { + ID int `json:"id"` + Name string `json:"name"` + APIKey string `json:"api_key,omitempty"` // Solo al crear + Capabilities string `json:"capabilities"` // 'cpu' o 'gpu' + Status string `json:"status"` // 'online', 'offline', 'disabled' + LastSeen time.Time `json:"last_seen"` + CreatedAt time.Time `json:"created_at"` +} + +// Job de traducción para worker remoto +type TranslationJob struct { + ID int64 `json:"id"` + NewsID int64 `json:"noticia_id"` + LangFrom string `json:"lang_from"` + LangTo string `json:"lang_to"` + Title string `json:"title"` + Summary string `json:"summary"` +} + +// Resultado de traducción enviado por worker +type TranslationResult struct { + JobID int64 `json:"job_id"` + TitleTr string `json:"title_trad"` + SummaryTr string `json:"resumen_trad"` + Error string `json:"error,omitempty"` +} + +// Mensaje WebSocket: cliente → servidor +type WSClientMessage struct { + Type string `json:"type"` // "register", "heartbeat", "result" + Capabilities string `json:"capabilities,omitempty"` + APIKey string `json:"api_key,omitempty"` + WorkerID int `json:"worker_id,omitempty"` + Result *TranslationResult `json:"result,omitempty"` +} + +// Mensaje WebSocket: servidor → cliente +type WSServerMessage struct { + Type string `json:"type"` // "job", "ack", "error" + Job *TranslationJob `json:"job,omitempty"` +} +``` + +### Fase 3: Handlers REST (Backend Go) + +**Nuevo archivo:** `backend/internal/handlers/remote_worker.go` + +#### Endpoints + +| Método | Endpoint | Descripción | +|--------|----------|-------------| +| POST | `/api/admin/workers/remote` | Crear worker (recibe nombre, capabilities) | +| GET | `/api/admin/workers/remote` | Listar todos los workers | +| GET | `/api/admin/workers/remote/:id` | Ver detalle de un worker | +| DELETE | `/api/admin/workers/remote/:id` | Eliminar worker | +| PATCH | `/api/admin/workers/remote/:id/toggle` | Habilitar/deshabilitar | + +#### Funciones principales + +1. **CreateRemoteWorker** - Genera API key aleatoria (64 chars), guarda en BD +2. **ListRemoteWorkers** - Retorna todos con jobs counts +3. **GetRemoteWorker** - Retorna detalle + estadísticas +4. **DeleteRemoteWorker** - Soft delete (marcar como disabled) +5. **ToggleRemoteWorker** - Enable/disable sin borrar + +### Fase 4: WebSocket Handler (Backend Go) + +**Nuevo archivo:** `backend/internal/handlers/worker_ws.go` + +```go +// Estructuras requeridas +type WSWorker struct { + conn *websocket.Conn + workerID int + capabilities string + lastHeartbeat time.Time +} + +// Conexiónmap[int]*WSWorker // workerID → conexión +var workers = make(map[int]*WSWorker) +var workersByConn = make(map[*websocket.Conn]*WSWorker) + +// Handler principal +func HandleWorkerWS(c *gin.Context) { + // 1. Verificar API key en query string o header + apiKey := c.Query("api_key") + if apiKey == "" { + apiKey = c.GetHeader("X-API-Key") + } + + // 2. Buscar worker por API key + workerID, err := validateWorkerAPIKey(apiKey) + if err != nil { + c.JSON(401, gin.H{"error": "Invalid API key"}) + return + } + + // 3. Upgrade a WebSocket + upgrader := websocket.Upgrader{ + CheckOrigin: func(r *http.Request) bool { return true }, + } + conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) + if err != nil { + return + } + + // 4. Registrar worker como online + wsWorker := &WSWorker{ + conn: conn, + workerID: workerID, + capabilities: getWorkerCapabilities(workerID), + lastHeartbeat: time.Now(), + } + workers[workerID] = wsWorker + workersByConn[conn] = wsWorker + + // 5. Update status en BD + updateWorkerStatus(workerID, "online") + + // 6. Loop de lectura + go handleWSRead(wsWorker) +} + +// Loop de lectura de mensajes del worker +func handleWSRead(wsWorker *WSWorker) { + for { + var msg WSClientMessage + err := wsWorker.conn.ReadJSON(&msg) + if err != nil { + // Worker desconectado + cleanupWorker(wsWorker) + return + } + + switch msg.Type { + case "register": + // Registro inicial (capabilities) + wsWorker.capabilities = msg.Capabilities + sendAck(wsWorker, "registered") + + case "heartbeat": + wsWorker.lastHeartbeat = time.Now() + sendAck(wsWorker, "ok") + + case "result": + // Worker completó un job + handleTranslationResult(wsWorker.workerID, msg.Result) + sendAck(wsWorker, "received") + } + } +} + +// Asignar job disponible al worker +func assignJobToWorker(workerID int) *TranslationJob { + caps := getWorkerCapabilities(workerID) + + // Buscar job pending que coincida con capabilities del worker + // Por ahora: cualquier job pending funciona (el worker decide qué procesar) + job := getNextPendingJob(caps) + if job == nil { + return nil + } + + // Asignar: UPDATE status='assigned', worker_id=workerID, assigned_at=NOW() + assignJob(job.ID, workerID) + + return job +} +``` + +### Fase 5: Scheduler (Migración Python → Go) + +**Nuevo archivo:** `backend/internal/workers/translation_scheduler.go` + +```go +package workers + +func StartTranslationScheduler() { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + scheduleTranslations() + } + } +} + +func scheduleTranslations() { + // Para cada idioma destino (es) + targetLangs := []string{"es"} + + for _, lang := range targetLangs { + // INSERT jobs para noticias sin traducción + _, err := db.Exec(` + INSERT INTO traducciones (noticia_id, lang_from, lang_to, status, created_at) + SELECT n.id, n.lang, $1, 'pending', NOW() + FROM noticias n + WHERE n.lang IS NOT NULL + AND TRIM(n.lang) != '' + AND n.lang != $1 + AND NOT EXISTS ( + SELECT 1 FROM traducciones t + WHERE t.noticia_id = n.id AND t.lang_to = $1 + ) + ORDER BY n.fecha DESC + LIMIT $2 + ON CONFLICT (noticia_id, lang_to) DO NOTHING + `, lang, 2000) + } +} + +// Cleanup de jobs huérfanos (assigned hace demasiado tiempo) +func cleanupStaleAssignments() { + db.Exec(` + UPDATE traducciones + SET status = 'pending', worker_id = NULL, assigned_at = NULL + WHERE status = 'assigned' + AND assigned_at < NOW() - INTERVAL '5 minutes' + `) +} +``` + +### Fase 6: Cliente Worker Remoto (Go) + +**Nuevo directorio:** `worker-client/` + +``` +worker-client/ +├── cmd/ +│ └── main.go # Entry point +├── internal/ +│ ├── config/ +│ │ └── config.go # Flags y configuración +│ ├── ws/ +│ │ └── client.go # WebSocket client +│ ├── translator/ +│ │ └── ctranslate2.go # Llamada a CTranslate2 +│ └── store/ +│ └── local.go # Persistencia local (SQLite) +├── go.mod +└── README.md +``` + +#### cmd/main.go + +```go +func main() { + // Flags + serverURL := flag.String("server", "ws://localhost:8080/ws/worker", "WebSocket server URL") + apiKey := flag.String("api-key", "", "API key for authentication") + device := flag.String("device", "cuda", "Device: cuda or cpu") + modelPath := flag.String("model-path", "./models/nllb-ct2", "CTranslate2 model path") + batchSize := flag.Int("batch", 8, "Translation batch size") + flag.Parse() + + // Validar + if *apiKey == "" { + log.Fatal("API key required") + } + + // Inicializar + cfg := config.Config{ + ServerURL: *serverURL, + APIKey: *apiKey, + Device: *device, + ModelPath: *modelPath, + BatchSize: *batchSize, + } + + // Conectar y loop + client := ws.NewClient(&cfg) + client.ConnectAndRun() +} +``` + +#### internal/ws/client.go + +```go +func (c *Client) ConnectAndRun() { + for { + err := c.connect() + if err != nil { + log.Printf("Connection failed: %v", err) + time.Sleep(5 * time.Second) + continue + } + + // Loop de mensajes + c.readLoop() + } +} + +func (c *Client) handleMessage(msg []byte) { + var serverMsg WSServerMessage + if err := json.Unmarshal(msg, &serverMsg); err != nil { + return + } + + switch serverMsg.Type { + case "job": + c.processJob(serverMsg.Job) + case "ping": + c.sendHeartbeat() + } +} + +func (c *Client) processJob(job *TranslationJob) { + // 1. Traducir título + titleTr := translate(job.Title, job.LangFrom, job.LangTo) + + // 2. Traducir resumen (chunks si es largo) + summaryTr := translateLongText(job.Summary, job.LangFrom, job.LangTo) + + // 3. Enviar resultado + c.SendResult(job.ID, titleTr, summaryTr, "") +} +``` + +#### internal/translator/ctranslate2.go + +```go +func Translate(text, srcLang, tgtLang string) (string, error) { + // Mapear códigos de idioma + srcCode := langMap[srcLang] + tgtCode := langMap[tgtLang] + + // Construir comando + cmd := exec.Command( + "ct2-transformers-converter", // O el binario de ctranslate2 + "--source", text, + "--src_lang", srcCode, + "--tgt_lang", tgtCode, + ) + + output, err := cmd.CombinedOutput() + if err != nil { + return "", err + } + + return string(output), nil +} +``` + +**Nota:** La forma más estable es usar el binario `ctranslate2` directamente como wrapper del modelo convertido. Ver documentación de CTranslate2 para la interfaz CLI. + +### Fase 7: Frontend (Panel Admin) + +**Modificar:** `frontend/src/pages/AdminWorkers.tsx` + +Añadir nueva sección "Workers Remotos": + +```tsx +// Estados adicionales +const [remoteWorkers, setRemoteWorkers] = useState([]) +const [showAddForm, setShowAddForm] = useState(false) + +// Fetch +const fetchRemoteWorkers = async () => { + const res = await api.get('/admin/workers/remote') + setRemoteWorkers(res.data) +} + +// UI: Añadir después de la sección de workers locales +
    +

    + + Workers Remotos +

    + + {/* Botón añadir */} + + + {/* Lista de workers */} +
    + {remoteWorkers.map(worker => ( +
    +
    + {worker.name} + + {worker.status} + +
    +
    + + +
    +
    + ))} +
    +
    +``` + +--- + +## Archivos a Crear/Modificar + +### Backend (Go) + +| Archivo | Acción | Descripción | +|---------|--------|-------------| +| `init-db/35-remote-workers.sql` | Crear | Tabla remote_workers + columnas | +| `backend/internal/models/worker.go` | Crear | Modelos RemoteWorker, TranslationJob | +| `backend/internal/handlers/remote_worker.go` | Crear | CRUD REST API | +| `backend/internal/handlers/worker_ws.go` | Crear | WebSocket handler | +| `backend/internal/workers/translation_scheduler.go` | Crear | Scheduler migrado de Python | +| `backend/cmd/server/main.go` | Modificar | Añadir rutas + iniciar scheduler | +| `backend/go.mod` | Modificar | Añadir dependencia `github.com/gorilla/websocket` | + +### Frontend (React) + +| Archivo | Acción | Descripción | +|---------|--------|-------------| +| `frontend/src/pages/AdminWorkers.tsx` | Modificar | Añadir sección workers remotos | + +### Cliente Worker + +| Archivo | Acción | Descripción | +|---------|--------|-------------| +| `worker-client/cmd/main.go` | Crear | CLI entry point | +| `worker-client/internal/config/config.go` | Crear | Configuración | +| `worker-client/internal/ws/client.go` | Crear | WebSocket client | +| `worker-client/internal/translator/ctranslate2.go` | Crear | Wrapper CTranslate2 | +| `worker-client/go.mod` | Crear | Módulos Go | +| `worker-client/README.md` | Crear | Instrucciones de uso | + +--- + +## Consideraciones de Seguridad + +1. **API Key**: Generar con `crypto/rand` - 64 caracteres hex +2. **Rate limiting**: En endpoints de creación de workers +3. **Validación**: Verificar que el worker tiene capabilities válidas +4. **Timeout**: Jobs asignados sin respuesta en 5 min vuelven a cola +5. **Conexión**: El servidor debe verificar origen del WebSocket + +--- + +## Testing Plan + +1. **Unidad**: Tests de handlers (mock de BD) +2. **Integración WS**: Probar conexión, registro, heartbeat, resultado +3. **E2E**: + - Iniciar cliente worker local + - Crear job pending manualmente + - Verificar que worker recibe job + - Verificar resultado en BD + +--- + +## Notas de Implementación + +- Mantener compatibilidad con workers Docker existentes (CPU/GPU locales) +- Los workers remotos serán una opción adicional, no reemplazo +- El scheduler puede coexistir (Go + Python) durante transición +- Considerar métricas: jobs procesadors por worker, tiempo promedio \ No newline at end of file diff --git a/workers/remote_translator_worker.py b/workers/remote_translator_worker.py new file mode 100644 index 0000000..b1137d6 --- /dev/null +++ b/workers/remote_translator_worker.py @@ -0,0 +1,389 @@ +import os +import sys +import time +import json +import logging +import re +import threading +from typing import List, Optional + +import websocket + +import ctranslate2 +from transformers import AutoTokenizer + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)s: %(message)s", + handlers=[logging.StreamHandler(sys.stdout)] +) +LOG = logging.getLogger("remote-translator") + +WORKER_NAME = os.environ.get("WORKER_NAME", "remote-worker") +WORKER_API_KEY = os.environ.get("WORKER_API_KEY", "") +WORKER_SERVER = os.environ.get("WORKER_SERVER", "ws://localhost:8080/ws/worker") +DEVICE = os.environ.get("CT2_DEVICE", "cpu") +MODEL_PATH = os.environ.get("CT2_MODEL_PATH", "/app/models/nllb-ct2") +COMPUTE_TYPE = os.environ.get("CT2_COMPUTE_TYPE", "int8") +UNIVERSAL_MODEL = os.environ.get("UNIVERSAL_MODEL", "facebook/nllb-200-distilled-600M") + +LANG_CODE_MAP = { + "en": "eng_Latn", "es": "spa_Latn", "fr": "fra_Latn", "de": "deu_Latn", + "it": "ita_Latn", "pt": "por_Latn", "nl": "nld_Latn", "sv": "swe_Latn", + "da": "dan_Latn", "fi": "fin_Latn", "no": "nob_Latn", + "pl": "pol_Latn", "cs": "ces_Latn", "sk": "slk_Latn", + "sl": "slv_Latn", "hu": "hun_Latn", "ro": "ron_Latn", + "el": "ell_Grek", "ru": "rus_Cyrl", "uk": "ukr_Cyrl", + "tr": "tur_Latn", "ar": "arb_Arab", "fa": "pes_Arab", + "he": "heb_Hebr", "zh": "zho_Hans", "ja": "jpn_Jpan", + "ko": "kor_Hang", "vi": "vie_Latn", +} + +MAX_SRC_TOKENS = 512 +MAX_NEW_TOKENS = 512 +BODY_CHARS_CHUNK = 900 + +_tokenizer = None +_translator = None +_ws = None +_reconnect_delay = 5 +_running = True +_stats = {"jobs_completed": 0, "jobs_failed": 0} + + +def ensure_model(): + global _tokenizer, _translator + + if _translator: + return + + model_bin = os.path.join(MODEL_PATH, "model.bin") + + if not os.path.exists(model_bin): + LOG.info(f"CTranslate2 model not found at {MODEL_PATH}, converting...") + convert_model() + + LOG.info(f"Loading CTranslate2 model from {MODEL_PATH} on {DEVICE}") + + _translator = ctranslate2.Translator( + MODEL_PATH, + device=DEVICE, + compute_type=COMPUTE_TYPE, + ) + + _tokenizer = AutoTokenizer.from_pretrained(UNIVERSAL_MODEL) + LOG.info("CTranslate2 model loaded successfully") + + +def convert_model(): + import subprocess + + os.makedirs(MODEL_PATH, exist_ok=True) + + quantization = COMPUTE_TYPE if COMPUTE_TYPE != "auto" else "int8" + + cmd = [ + "ct2-transformers-converter", + "--model", UNIVERSAL_MODEL, + "--output_dir", MODEL_PATH, + "--quantization", quantization, + "--force" + ] + + LOG.info(f"Running: {' '.join(cmd)}") + result = subprocess.run(cmd, capture_output=True, text=True, timeout=1800) + + if result.returncode != 0: + LOG.error(f"Model conversion failed: {result.stderr}") + raise RuntimeError("Failed to convert model") + + LOG.info("Model conversion completed") + + +def clean_text(text: str) -> str: + if not text: + return "" + text = re.sub(r'<[^>]+>', '', text) + text = text.replace('', '') + text = text.replace(' ', ' ') + text = text.replace('&', '&') + text = text.replace('<', '<') + text = text.replace('>', '>') + text = text.replace('"', '"') + text = re.sub(r'\s+', ' ', text) + return text.strip() + + +def translate_texts(src: str, tgt: str, texts: List[str]) -> List[str]: + if not texts: + return [] + + ensure_model() + + clean = [(t or "").strip() for t in texts] + if all(not t for t in clean): + return ["" for _ in clean] + + src_code = LANG_CODE_MAP.get(src, f"{src}_Latn") + tgt_code = LANG_CODE_MAP.get(tgt, "spa_Latn") + + try: + _tokenizer.src_lang = src_code + except Exception: + pass + + sources = [] + for t in clean: + if t: + ids = _tokenizer.encode(t, truncation=True, max_length=MAX_SRC_TOKENS) + tokens = _tokenizer.convert_ids_to_tokens(ids) + sources.append(tokens) + else: + sources.append([]) + + target_prefix = [[tgt_code]] * len(sources) + + results = _translator.translate_batch( + sources, + target_prefix=target_prefix, + beam_size=2, + max_decoding_length=MAX_NEW_TOKENS, + repetition_penalty=2.0, + no_repeat_ngram_size=3, + ) + + translated = [] + for result in results: + try: + if result.hypotheses and len(result.hypotheses) > 0: + hyp = result.hypotheses[0] + if isinstance(hyp, list) and len(hyp) > 0: + first_hyp = hyp[0] + if isinstance(first_hyp, dict) and "token_ids" in first_hyp: + tokens = first_hyp["token_ids"] + text = _tokenizer.decode(tokens) + translated.append(text.strip()) + elif isinstance(first_hyp, str): + token_strings = hyp[1:] if len(hyp) > 1 else [] + if token_strings: + text = _tokenizer.convert_tokens_to_string(token_strings) + translated.append(text.strip()) + else: + translated.append("") + else: + translated.append("") + else: + translated.append("") + else: + translated.append("") + except Exception as e: + LOG.error(f"Error processing result: {e}") + translated.append("") + + return translated + + +def split_body_into_chunks(text: str) -> List[str]: + text = (text or "").strip() + if len(text) <= BODY_CHARS_CHUNK: + return [text] if text else [] + + parts = re.split(r'(\n\n+|(?<=[\.\!\?؛؟。])\s+)', text) + chunks = [] + current = "" + + for part in parts: + if not part: + continue + if len(current) + len(part) <= BODY_CHARS_CHUNK: + current += part + else: + if current.strip(): + chunks.append(current.strip()) + current = part + if current.strip(): + chunks.append(current.strip()) + + return chunks if chunks else [text] + + +def translate_body_long(src: str, tgt: str, body: str) -> str: + body = (body or "").strip() + if not body: + return "" + + chunks = split_body_into_chunks(body) + if len(chunks) == 1: + return translate_texts(src, tgt, [body])[0] + + translated_chunks = [] + for ch in chunks: + tr = translate_texts(src, tgt, [ch])[0] + translated_chunks.append(tr) + + return " ".join(translated_chunks) + + +def process_job(job: dict) -> dict: + job_id = job.get("id") + lang_from = job.get("lang_from", "en") + lang_to = job.get("lang_to", "es") + title = job.get("title", "") + summary = job.get("summary", "") + + LOG.info(f"Processing job {job_id}: {lang_from} -> {lang_to}") + + if lang_from == lang_to: + return { + "job_id": job_id, + "title_trad": title, + "summary_trad": summary, + "error": "" + } + + try: + title_tr = translate_texts(lang_from, lang_to, [title])[0] + title_tr = clean_text(title_tr) or title + + summary_tr = "" + if summary: + summary_tr = translate_body_long(lang_from, lang_to, summary) + summary_tr = clean_text(summary_tr) or summary + + _stats["jobs_completed"] += 1 + + return { + "job_id": job_id, + "title_trad": title_tr, + "summary_trad": summary_tr, + "error": "" + } + except Exception as e: + LOG.error(f"Translation error for job {job_id}: {e}") + _stats["jobs_failed"] += 1 + return { + "job_id": job_id, + "title_trad": "", + "summary_trad": "", + "error": str(e) + } + + +def send_message(msg: dict): + if _ws and _ws.sock and _ws.sock.connected: + try: + _ws.send(json.dumps(msg)) + except Exception as e: + LOG.error(f"Send error: {e}") + + +def on_message(ws, message): + try: + msg = json.loads(message) + except: + LOG.error(f"Invalid JSON: {message}") + return + + msg_type = msg.get("type") + + if msg_type == "job": + job = msg.get("job", {}) + result = process_job(job) + send_message({ + "type": "result", + "result": result + }) + LOG.info(f"Sent result for job {result['job_id']}") + + elif msg_type == "ping": + send_message({"type": "heartbeat"}) + + elif msg_type == "ack": + LOG.debug(f"Server ACK: {msg.get('message', '')}") + + elif msg_type == "error": + LOG.error(f"Server error: {msg.get('message', '')}") + + +def on_error(ws, error): + LOG.error(f"WebSocket error: {error}") + + +def on_close(ws, close_status_code, close_msg): + global _reconnect_delay + LOG.warning(f"WebSocket closed: {close_status_code} - {close_msg}") + LOG.info(f"Reconnecting in {_reconnect_delay}s...") + + +def on_open(ws): + global _reconnect_delay + LOG.info("Connected to server") + _reconnect_delay = 5 + + send_message({ + "type": "register", + "capabilities": DEVICE, + "worker_name": WORKER_NAME + }) + + +def stats_reporter(): + while _running: + time.sleep(60) + LOG.info(f"Stats: completed={_stats['jobs_completed']}, failed={_stats['jobs_failed']}") + + +def connect(): + global _ws, _reconnect_delay + + while _running: + try: + if not WORKER_API_KEY: + LOG.error("WORKER_API_KEY not set") + time.sleep(60) + continue + + ws_url = f"{WORKER_SERVER}?api_key={WORKER_API_KEY}" + + _ws = websocket.WebSocketApp( + ws_url, + on_open=on_open, + on_message=on_message, + on_error=on_error, + on_close=on_close, + header={"X-API-Key": WORKER_API_KEY} + ) + + LOG.info(f"Connecting to {WORKER_SERVER}...") + _ws.run_forever(ping_interval=30, ping_timeout=10) + + except Exception as e: + LOG.error(f"Connection error: {e}") + + LOG.info(f"Waiting {_reconnect_delay}s before reconnect...") + time.sleep(_reconnect_delay) + _reconnect_delay = min(_reconnect_delay * 2, 60) + + +def main(): + global _running + + if not WORKER_API_KEY: + LOG.error("WORKER_API_KEY environment variable is required") + sys.exit(1) + + LOG.info(f"Remote translator worker starting...") + LOG.info(f"Server: {WORKER_SERVER}") + LOG.info(f"Device: {DEVICE}") + LOG.info(f"Model: {MODEL_PATH}") + + ensure_model() + + stats_thread = threading.Thread(target=stats_reporter, daemon=True) + stats_thread.start() + + connect() + + +if __name__ == "__main__": + main() \ No newline at end of file