From d24ab0d4b71c7b5bb01afae020fd168dbc9e1e3d Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 24 Sep 2026 09:31:31 +0000 Subject: [PATCH] feat(reels): file d'attente persistante des rendus (reel_jobs) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit L'ancienne file comptait les posts « processing » : un redémarrage pendant un rendu la bloquait définitivement. Les rendus passent désormais par une table reel_jobs, réservée avec FOR UPDATE SKIP LOCKED, avec heartbeat et reprise des jobs interrompus au démarrage (2 tentatives maximum). - Pipeline vidéo et pipeline images (Remotion) extraits des routes vers server/services/reels/, les routes ne font que valider (zod) et mettre en file - Rendus d'images suivis en base au lieu d'une Map en mémoire - storeName conservé pour les jobs mis en attente - Publication Facebook en binaire : l'URL relative /uploads/... transmise auparavant n'était pas téléchargeable par Facebook - Job en échec si toutes les pages échouent ou si la voix demandée manque (le service Python renseigne enfin tts_error) - La voix choisie sur la page images est réellement utilisée - sync-info mesure la voix Gemini avec sa clé - Timeouts sur les appels au service FFmpeg, vignettes via execFile, /debug-ffmpeg protégé par la clé API, plus de voix anglaise en secours Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_018Ze4bs7tpF1KGWUk6ZZSZ4 --- client/src/pages/mobile/remotion-video.tsx | 2 + client/src/pages/remotion-video.tsx | 2 + ffmpeg-service/main.py | 10 +- server/index.ts | 7 + server/migrate.ts | 21 + server/routes/reels.ts | 500 ++------------------- server/routes/remotion.ts | 486 ++++---------------- server/services/ffmpeg.ts | 9 + server/services/reels/assets.ts | 65 +++ server/services/reels/imagesPipeline.ts | 234 ++++++++++ server/services/reels/index.ts | 14 + server/services/reels/publish.ts | 133 ++++++ server/services/reels/queue.ts | 277 ++++++++++++ server/services/reels/videoPipeline.ts | 78 ++++ server/services/thumbnail.ts | 10 +- server/services/ttsSync.ts | 11 +- server/storage.ts | 20 - shared/reel.ts | 65 +++ shared/schema.ts | 28 +- 19 files changed, 1071 insertions(+), 901 deletions(-) create mode 100644 server/services/reels/assets.ts create mode 100644 server/services/reels/imagesPipeline.ts create mode 100644 server/services/reels/index.ts create mode 100644 server/services/reels/publish.ts create mode 100644 server/services/reels/queue.ts create mode 100644 server/services/reels/videoPipeline.ts create mode 100644 shared/reel.ts diff --git a/client/src/pages/mobile/remotion-video.tsx b/client/src/pages/mobile/remotion-video.tsx index cadc0e1..b6b302d 100644 --- a/client/src/pages/mobile/remotion-video.tsx +++ b/client/src/pages/mobile/remotion-video.tsx @@ -168,6 +168,8 @@ export default function MobileRemotionVideoPage() { images.forEach(img => formData.append("images", img)); selectedLibraryImages.forEach(m => formData.append("existingImageUrls", m.originalUrl)); if (overlayText) formData.append("overlayText", overlayText); + formData.append("ttsEngine", ttsEngine); + formData.append("ttsVoice", ttsVoice); if (selectedPageIds[0]) formData.append("selectedPageId", selectedPageIds[0]); if (musicFile) { formData.append("music", musicFile); formData.append("musicVolume", String(musicVolume)); } else if (selectedTrack) { formData.append("musicTrackUrl", selectedTrack.url); formData.append("musicVolume", String(musicVolume)); } diff --git a/client/src/pages/remotion-video.tsx b/client/src/pages/remotion-video.tsx index d520f41..7c3e4c5 100644 --- a/client/src/pages/remotion-video.tsx +++ b/client/src/pages/remotion-video.tsx @@ -134,6 +134,8 @@ export default function RemotionVideoPage() { images.forEach(img => formData.append("images", img)); selectedLibraryImages.forEach(m => formData.append("existingImageUrls", m.originalUrl)); if (overlayText) formData.append("overlayText", overlayText); + formData.append("ttsEngine", ttsEngine); + formData.append("ttsVoice", ttsVoice); if (selectedPageIds[0]) formData.append("selectedPageId", selectedPageIds[0]); if (musicFile) { formData.append("music", musicFile); formData.append("musicVolume", String(musicVolume)); } else if (selectedTrack) { formData.append("musicTrackUrl", selectedTrack.url); formData.append("musicVolume", String(musicVolume)); } diff --git a/ffmpeg-service/main.py b/ffmpeg-service/main.py index 7417cb3..abe0963 100644 --- a/ffmpeg-service/main.py +++ b/ffmpeg-service/main.py @@ -46,7 +46,9 @@ run_diagnostics() @app.get("/debug-ffmpeg") -async def debug_ffmpeg(): +async def debug_ffmpeg(x_api_key: str = Header(None)): + if x_api_key != API_KEY: + raise HTTPException(status_code=401, detail="Invalid API Key") try: filters = subprocess.run( ["ffmpeg", "-filters"], capture_output=True, text=True @@ -57,7 +59,6 @@ async def debug_ffmpeg(): "subtitles": "subtitles" in filters, "drawtext": "drawtext" in filters, }, - "env": {k: v for k, v in os.environ.items() if "API" not in k}, "fonts": fonts.splitlines()[:50], # First 50 "raw_filters_hint": filters[:500], } @@ -337,7 +338,6 @@ async def generate_tts_with_subs( "fr-FR-DeniseNeural", ] - fallback_voices.append("en-US-JennyNeural") fallback_voices = list(dict.fromkeys(fallback_voices)) last_error = None @@ -902,6 +902,7 @@ async def process_reel(request: ReelRequest, x_api_key: str = Header(None)): # 3. Generate TTS (if enabled) has_tts = False + tts_error = None tts_clean_text = "" if request.tts_enabled and request.text: @@ -966,12 +967,14 @@ async def process_reel(request: ReelRequest, x_api_key: str = Header(None)): has_tts = True else: print("❌ TTS audio file missing or empty!") + tts_error = "Fichier audio de la voix vide" else: print("⚠️ TTS text is empty after cleaning, skipping.") except Exception as e: import traceback print(f"❌ Failed to generate TTS: {e}") traceback.print_exc() + tts_error = str(e) or type(e).__name__ stats["tts_duration"] = time.time() - start_step start_step = time.time() @@ -1298,6 +1301,7 @@ async def process_reel(request: ReelRequest, x_api_key: str = Header(None)): "output_base64": out_b64, "duration": duration, "processing_stats": stats, + "tts_error": tts_error, } except Exception as e: diff --git a/server/index.ts b/server/index.ts index 7ba3bec..8f52097 100644 --- a/server/index.ts +++ b/server/index.ts @@ -14,6 +14,7 @@ import { schedulerService } from "./services/scheduler"; import { ensureAdminUserExists } from "./init-admin"; import { startTokenCron } from "./cron"; import { migrate } from "./migrate"; +import { startReelProcessing } from "./services/reels"; import { jamendoService } from "./services/jamendo"; import { ffmpegService } from "./services/ffmpeg"; @@ -230,5 +231,11 @@ app.use((req, res, next) => { // Start Token Refresh Cron startTokenCron(); + + // File des rendus de Reels : démarrée une fois le serveur à l'écoute, car + // le rendu d'images lit la voix et la musique via http://localhost. + startReelProcessing().catch((error) => { + console.error("❌ Démarrage de la file des Reels impossible :", error); + }); }); })(); diff --git a/server/migrate.ts b/server/migrate.ts index 3a9aa85..9f8259a 100644 --- a/server/migrate.ts +++ b/server/migrate.ts @@ -213,6 +213,26 @@ export async function migrate() { ADD COLUMN IF NOT EXISTS "publish_status" text; `); + // reel_jobs : file d'attente persistante des rendus de Reels + await client.query(` + CREATE TABLE IF NOT EXISTS "reel_jobs" ( + "id" varchar PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL, + "user_id" varchar NOT NULL REFERENCES "users"("id") ON DELETE cascade, + "post_id" varchar REFERENCES "posts"("id") ON DELETE cascade, + "kind" text NOT NULL, + "params" jsonb NOT NULL, + "status" text NOT NULL DEFAULT 'pending', + "step" text, + "progress" integer NOT NULL DEFAULT 0, + "attempts" integer NOT NULL DEFAULT 0, + "result" jsonb, + "error" text, + "heartbeat_at" timestamp, + "created_at" timestamp NOT NULL DEFAULT now(), + "updated_at" timestamp NOT NULL DEFAULT now() + ); + `); + // 4. Index de performance // // Postgres n'indexe pas les clés étrangères tout seul : sans ces index, le @@ -249,6 +269,7 @@ export async function migrate() { ["idx_scheduled_posts_publish_id", `CREATE INDEX IF NOT EXISTS "idx_scheduled_posts_publish_id" ON "scheduled_posts" ("publish_id") WHERE "publish_id" IS NOT NULL`], // Reels en cours de génération + ["idx_reel_jobs_queue", `CREATE INDEX IF NOT EXISTS "idx_reel_jobs_queue" ON "reel_jobs" ("created_at") WHERE "status" IN ('pending', 'processing')`], ["idx_posts_generation_status", `CREATE INDEX IF NOT EXISTS "idx_posts_generation_status" ON "posts" ("generation_status") WHERE "generation_status" IS NOT NULL`], ]; diff --git a/server/routes/reels.ts b/server/routes/reels.ts index c3b35a6..4d7b1c4 100644 --- a/server/routes/reels.ts +++ b/server/routes/reels.ts @@ -3,17 +3,14 @@ */ import { Router, Request, Response } from 'express'; -import type { User, Media, SocialPage } from '@shared/schema'; +import type { User } from '@shared/schema'; import { storage } from '../storage'; import { ffmpegService } from '../services/ffmpeg'; -import { facebookService } from '../services/facebook'; -import { tiktokService } from '../services/tiktok'; -import { createVideoThumbnail } from '../services/thumbnail'; -import { minioService as cloudinaryService, buildMinioUrl, resolveInternalUrl } from '../services/minio'; +import { resolveInternalUrl } from '../services/minio'; +import { videoReelParamsSchema, type VideoReelParams } from '@shared/reel'; +import { enqueueReelJob, countActiveReelJobs } from '../services/reels/queue'; +import { resolveGeminiApiKey, resolveLogoPath, resolveMusicUrl, resolveStoreName } from '../services/reels/assets'; import { openRouterService, describeGenerationError } from '../services/openrouter'; -import { db } from '../db'; -import { cloudinaryConfig } from '@shared/schema'; -import { eq } from 'drizzle-orm'; import { ttsSyncService } from '../services/ttsSync'; /** Piste musicale telle qu'attendue par le client. */ @@ -306,42 +303,12 @@ reelsRouter.post('/reels/preview', async (req: Request, res: Response) => { return res.status(400).json({ error: 'Le média doit être une vidéo' }); } - // Récupérer l'URL de la musique si trackId fourni - let finalMusicUrl = musicUrl; - if (musicTrackId && !musicUrl) { - if (musicTrackId.startsWith('internal_')) { - const internalId = musicTrackId.replace('internal_', ''); - const track = await storage.getAudioTrack(internalId); - if (track) { - // Utiliser INTERNAL_APP_URL pour que le container ffmpeg (réseau Docker interne) - // puisse télécharger le fichier audio. APP_URL (HTTPS public) n'est pas accessible - // depuis le réseau internal Docker. - const internalBaseUrl = process.env.INTERNAL_APP_URL || process.env.APP_URL || 'http://localhost:5555'; - finalMusicUrl = track.url.startsWith('http') - ? track.url - : `${internalBaseUrl}${track.url}`; - console.log(`🎵 Audio URL for ffmpeg: ${finalMusicUrl}`); - } - } else { - } - } - - let watermarkUrl: string | undefined = undefined; - try { - const config = await storage.getCloudinaryConfig(); - if (config && config.logoPublicId) { - watermarkUrl = resolveInternalUrl(buildMinioUrl(config.cloudName, config.logoPublicId, config.publicUrl)); - } - } catch (e) { - console.error('Error fetching watermark configuration', e); - } - - // Get Gemini API key if using Gemini TTS - let geminiApiKey: string | undefined = undefined; - if (ttsEngine === "gemini") { - const appCfg = await storage.getAppConfig(); - geminiApiKey = appCfg?.geminiApiKey ?? process.env.GEMINI_API_KEY ?? undefined; - } + const [finalMusicUrl, logoPath, geminiApiKey] = await Promise.all([ + resolveMusicUrl(musicTrackId, musicUrl), + resolveLogoPath(), + resolveGeminiApiKey(ttsEngine), + ]); + const watermarkUrl = logoPath ? resolveInternalUrl(logoPath) : undefined; const finalWordDuration = wordDuration; @@ -391,12 +358,7 @@ reelsRouter.post('/reels/tts-preview', async (req: Request, res: Response) => { return res.status(400).json({ error: 'Texte requis' }); } - let geminiApiKey: string | undefined = undefined; - if (ttsEngine === "gemini") { - const appCfg = await storage.getAppConfig(); - geminiApiKey = appCfg?.geminiApiKey ?? process.env.GEMINI_API_KEY ?? undefined; - } - + const geminiApiKey = await resolveGeminiApiKey(ttsEngine); const result = await ffmpegService.previewTTS(text, ttsVoice, ttsEngine, geminiApiKey); if (!result.success) { @@ -420,7 +382,8 @@ reelsRouter.post('/reels/sync-info', async (req: Request, res: Response) => { if (!text || !ttsVoice) { return res.status(400).json({ error: 'Texte et voix requis' }); } - const sync = await ttsSyncService.calculateSyncTiming(text, ttsVoice, ttsEngine); + const geminiApiKey = await resolveGeminiApiKey(ttsEngine); + const sync = await ttsSyncService.calculateSyncTiming(text, ttsVoice, ttsEngine, geminiApiKey); res.json(sync); } catch (error) { console.error('❌ Error calculating sync:', error); @@ -429,444 +392,51 @@ reelsRouter.post('/reels/sync-info', async (req: Request, res: Response) => { }); /** - * Traitement d'arrière-plan pour les Reels - * Gère le pipeline FFmpeg -> Cloudinary -> Facebook de manière asynchrone - */ - -/** - * Vérifie la file d'attente et lance le prochain job si disponible - */ -async function checkQueueAndProcessNext() { - try { - const nextPost = await storage.getNextPendingReel(); - if (nextPost) { - console.log(`📥 [Queue] Found pending job: ${nextPost.id}. Starting processing...`); - - // Marquer comme processing immédiatement - await storage.updatePostGenerationStatus(nextPost.id, 'processing', 0); - - // Reconstruire le payload depuis les données stockées (si stockées) ou les defaults - // Note: En mode "pending", on a perdu le body de la requête initiale car on ne stocke pas tout dans Post. - // Pour une solution robuste, il faudrait stocker les paramètres de génération dans une table 'reel_jobs'. - // ICI: HACK PROVISOIRE -> On suppose que les données sont stockées dans 'content' ou 'productInfo' mais ce n'est pas le cas. - // SOLUTION: On ne peut pas relancer processReelBackground sans les arguments (ttsVoice, music, etc). - // - // FIX: Pour le MVP, comme on n'a pas de table 'jobs', on va devoir stocker les paramètres requis dans 'productInfo' (jsonb) du Post - // lors de la création en mode 'pending'. - - // Récupérer les paramètres stockés - const jobData = nextPost.productInfo as any; - - if (!jobData || !jobData.videoMediaId) { - console.error(`❌ [Queue] Job ${nextPost.id} has no stored job data in productInfo.`); - await storage.updatePostGenerationStatus(nextPost.id, 'failed', 0, "Données de job manquantes"); - return; - } - - // Lancer le traitement - processReelBackground(nextPost.userId, nextPost.id, jobData) - .catch(err => console.error('🔥 [Queue] Unhandled error starting queued job:', err)); - } else { - console.log('🏁 [Queue] No more pending jobs.'); - } - } catch (error) { - console.error('❌ [Queue] Error checking queue:', error); - } -} - -/** - * Traitement d'arrière-plan pour les Reels - * Gère le pipeline FFmpeg -> Cloudinary -> Facebook de manière asynchrone - */ -async function processReelBackground( - userId: string, - postId: string, - data: { - videoMediaId: string; - musicTrackId?: string; - musicUrl?: string; - overlayText?: string; - description?: string; - ttsEnabled?: boolean; - ttsVoice?: string; - ttsEngine?: string; - pageIds: string[]; - scheduledFor?: string; - wordDuration?: number; - fontSize?: number; - musicVolume?: number; - drawText?: boolean; - stabilize?: boolean; - enableEndingEffect?: boolean; - }, - storeName?: string -) { - const { - videoMediaId, - musicTrackId, - musicUrl, - overlayText, - description, - ttsEnabled, - ttsVoice, - ttsEngine, - pageIds, - scheduledFor, - wordDuration, - fontSize, - musicVolume, - drawText, - stabilize, - enableEndingEffect, - } = data; - - console.log(`🔄 [Background] Starting processing for Post ${postId}`); - - // Helper to update generation progress - const updateProgress = async (progress: number, status: string = 'processing') => { - try { - await storage.updatePostGenerationStatus(postId, status, progress); - } catch (e) { - console.error(`⚠️ [Background] Failed to update progress for ${postId}:`, e); - } - }; - - try { - await updateProgress(5); - // 1. Récupérer le média vidéo source - const media = await storage.getMediaById(videoMediaId); - if (!media || media.type !== 'video') { - throw new Error('Vidéo source introuvable ou invalide'); - } - - // 2. Récupérer l'URL de la musique - let finalMusicUrl = musicUrl; - if (musicTrackId && !musicUrl) { - if (musicTrackId.startsWith('internal_')) { - try { - const internalId = musicTrackId.replace('internal_', ''); - const track = await storage.getAudioTrack(internalId); - if (track) { - const internalBaseUrl = process.env.INTERNAL_APP_URL || process.env.APP_URL || 'http://localhost:5555'; - finalMusicUrl = track.url.startsWith('http') - ? track.url - : `${internalBaseUrl}${track.url}`; - console.log(`🎵 [Background] Audio URL for ffmpeg: ${finalMusicUrl}`); - } - } catch (e) { console.error('Error fetching internal track', e); } - } - } - - console.log('🎬 [Background] FFmpeg Processing:', { - postId, - videoUrl: media.originalUrl, - hasMusic: !!finalMusicUrl, // Log boolean to avoid long URL - hasText: !!overlayText - }); - - // 3. Traiter la vidéo via FFmpeg - await updateProgress(15); - const startTime = Date.now(); - - let watermarkUrl: string | undefined = undefined; - try { - const config = await storage.getCloudinaryConfig(); - if (config && config.logoPublicId) { - watermarkUrl = resolveInternalUrl(buildMinioUrl(config.cloudName, config.logoPublicId, config.publicUrl)); - } - } catch (e) { - console.error('Error fetching watermark configuration', e); - } - - // Get Gemini API key if using Gemini TTS - let geminiApiKey: string | undefined = undefined; - if (ttsEngine === "gemini") { - const appCfg = await storage.getAppConfig(); - geminiApiKey = appCfg?.geminiApiKey ?? process.env.GEMINI_API_KEY ?? undefined; - } - - const finalWordDuration = wordDuration ?? 0.6; - - console.log('🔊 [Background] TTS config:', { ttsEnabled, ttsVoice, ttsEngine }); - - const ffmpegResult = await ffmpegService.processReelFromUrl(resolveInternalUrl(media.originalUrl), { - text: overlayText, - musicUrl: finalMusicUrl, - ttsEnabled, - ttsVoice, - ttsEngine, - geminiApiKey, - wordDuration: finalWordDuration, - fontSize, - musicVolume, - drawText, - stabilize, - watermarkUrl, - storeName, - enableEndingEffect, - }); - - console.log(`⏱️ [Background] FFmpeg took ${(Date.now() - startTime) / 1000}s`); - - if (ffmpegResult.ttsError) { - console.error(`❌ [Background] TTS FAILED: ${ffmpegResult.ttsError}`); - await storage.updatePost(postId, { generationError: `TTS échoué: ${ffmpegResult.ttsError}` }); - } - - if (!ffmpegResult.success || !ffmpegResult.videoBase64) { - throw new Error(ffmpegResult.error || 'Erreur de traitement vidéo FFmpeg'); - } - - // 4. Upload sur Cloudinary - await updateProgress(65); - console.log('☁️ [Background] Uploading to Cloudinary...'); - const videoBuffer = Buffer.from(ffmpegResult.videoBase64, 'base64'); - const cloudinaryResult = await cloudinaryService.uploadMedia( - videoBuffer, - `reel-${Date.now()}.mp4`, - userId, - 'video/mp4' - ); - console.log('✅ [Background] Uploaded:', cloudinaryResult.originalUrl); - - // 5. Créer l'enregistrement Media pour la vidéo traitée - // La vignette est extraite maintenant : la vidéo sera supprimée du disque - // après publication, mais l'historique doit rester illustré. - const thumbnailUrl = await createVideoThumbnail(videoBuffer); - - const processedMedia = await storage.createMedia({ - userId: userId, - type: 'video', - cloudinaryPublicId: cloudinaryResult.publicId, - originalUrl: cloudinaryResult.originalUrl, - facebookFeedUrl: cloudinaryResult.facebookFeedUrl || null, - instagramFeedUrl: cloudinaryResult.instagramFeedUrl || null, - instagramStoryUrl: cloudinaryResult.instagramStoryUrl || null, - thumbnailUrl, - fileName: `reel-processed-${Date.now()}.mp4`, - fileSize: videoBuffer.length, - }); - - // 6. Lier le média traité au Post existant - await storage.updatePostMedia(postId, [processedMedia.id]); - await updateProgress(85); - - // 7. Publier sur les pages - await updateProgress(90); - const results: { pageId: string; success: boolean; reelId?: string; error?: string }[] = []; - - for (const pageId of pageIds) { - try { - const page = await storage.getSocialPage(pageId); - if (!page) { - results.push({ pageId, success: false, error: 'Page non trouvée' }); - continue; - } - - if (page.platform !== 'facebook' && page.platform !== 'tiktok') { - results.push({ - pageId, - success: false, - error: `Plateforme non supportée pour les reels : ${page.platform}`, - }); - continue; - } - - // Créer l'entrée scheduled_post (log de publication) - const scheduledPost = await storage.createScheduledPost({ - postId: postId, - pageId: page.id, - postType: 'reel', - scheduledAt: scheduledFor ? new Date(scheduledFor) : new Date(), - }); - - if (!scheduledFor) { - // Publication immédiate — la même vidéo part sur chaque destination - console.log(`🚀 [Background] Publishing to ${page.platform} ${page.pageName}...`); - const finalDescription = description || overlayText || ''; - - if (page.platform === 'facebook') { - const reelId = await facebookService.publishReel( - page, - cloudinaryResult.originalUrl, - finalDescription - ); - - // Mise à jour succès - await storage.updateScheduledPost(scheduledPost.id, { - publishedAt: new Date(), - externalPostId: reelId, - }); - - results.push({ pageId, success: true, reelId }); - } else { - const publishId = await tiktokService.publishVideoFromBuffer( - page, - videoBuffer, - finalDescription - ); - - // TikTok finalise la publication de son côté : on marque l'envoi - // comme effectué pour ne pas republier, et le poller de statut - // renseignera l'identifiant définitif du post. - await storage.updateScheduledPost(scheduledPost.id, { - publishedAt: new Date(), - publishId, - publishStatus: 'PROCESSING_UPLOAD', - }); - - results.push({ pageId, success: true, reelId: publishId }); - } - } else { - // Planifié - results.push({ pageId, success: true, reelId: 'scheduled' }); - } - - } catch (pageError: any) { - console.error(`❌ [Background] Error publishing to page ${pageId}:`, pageError); - results.push({ - pageId, - success: false, - error: pageError.message || 'Erreur inconnue', - }); - } - } - - // 8. Mettre à jour le statut global du Post - const allSuccess = results.every(r => r.success); - const anySuccess = results.some(r => r.success); - - // Si au moins une réussite, on considère "published" (ou partial), sinon "failed" - // Si planifié, reste "scheduled". - let finalStatus: 'failed' | 'scheduled' | 'published' = 'failed'; - if (scheduledFor) { - finalStatus = 'scheduled'; - } else if (allSuccess) { - finalStatus = 'published'; - } else if (anySuccess) { - finalStatus = 'published'; // Partiellement publié - } - - await storage.updatePost(postId, { - status: finalStatus, - }); - await updateProgress(100, 'completed'); - - console.log(`✅ [Background] Processing complete for Post ${postId}. Status: ${finalStatus}`); - - } catch (error: any) { - console.error(`❌ [Background] Critical error for Post ${postId}:`, error); - await storage.updatePost(postId, { - status: 'failed', - }); - try { - await storage.updatePostGenerationStatus( - postId, 'failed', 0, - error?.message || 'Erreur inconnue lors du traitement' - ); - } catch (e) { - console.error(`⚠️ [Background] Failed to update error status for ${postId}:`, e); - } - } finally { - // IMPORTANT: Toujours vérifier la file d'attente à la fin (succès ou échec) - await checkQueueAndProcessNext(); - } -} - -/** - * Créer et publier un Reel (Asynchrone avec File d'Attente) + * Créer et publier un Reel : le rendu part dans la file persistante. * POST /api/reels */ reelsRouter.post('/reels', async (req: Request, res: Response) => { try { const user = req.user as User; - const { - videoMediaId, - musicTrackId, - musicUrl, - overlayText, - description, - ttsEnabled, - ttsVoice, - ttsEngine, - pageIds, - scheduledFor, - wordDuration = 0.6, - fontSize = 64, - musicVolume = 0.25, - drawText = true, - stabilize = false, - enableEndingEffect = true, - } = req.body; - - // Validation immédiate - if (!videoMediaId) { - return res.status(400).json({ error: 'Vidéo requise' }); + const parsed = videoReelParamsSchema.safeParse(req.body); + if (!parsed.success) { + return res.status(400).json({ error: parsed.error.issues[0]?.message ?? 'Paramètres invalides' }); } - if (!pageIds || pageIds.length === 0) { - return res.status(400).json({ error: 'Au moins une page requise' }); - } + const params: VideoReelParams = { + ...parsed.data, + storeName: await resolveStoreName(user.id, parsed.data.pageIds[0]), + }; - // Vérifier le nombre de jobs en cours - const processingCount = await storage.countProcessingReels(); - // MAX_CONCURRENT = 1 - const isQueueBusy = processingCount >= 1; + const waiting = await countActiveReelJobs(); - const initialStatus = isQueueBusy ? 'pending' : 'processing'; - const initialMessage = isQueueBusy - ? "File d'attente pleine. Votre vidéo sera traitée dès que possible." - : "Traitement démarré en arrière-plan."; - - // Récupérer le nom de la première page pour l'utiliser comme storeName - let storeName: string | undefined = undefined; - try { - const page = await storage.getSocialPage(pageIds[0]); - if (page) { - storeName = page.pageName; - } - } catch (e) { - console.error('Erreur récupération nom de page', e); - } - - // Créer immédiatement le Post en base - // ON STOCKE LES PARAMS DU JOB DANS productInfo POUR POUVOIR LE REPRENDRE PLUS TARD - // C'est un hack car on n'a pas de table params_job, mais ça marche car productInfo est jsonb const post = await storage.createPost({ userId: user.id, - content: description || overlayText || '', + content: params.description || params.overlayText || '', aiGenerated: 'false', - status: scheduledFor ? 'scheduled' : 'draft', - scheduledFor: scheduledFor ? new Date(scheduledFor) : undefined, - generationStatus: initialStatus, + status: params.scheduledFor ? 'scheduled' : 'draft', + scheduledFor: params.scheduledFor ? new Date(params.scheduledFor) : undefined, + generationStatus: 'pending', generationProgress: 0, - productInfo: req.body, // Stockage complet des paramètres }); - console.log(`✨ Reel Request accepted. Post ID: ${post.id}. Status: ${initialStatus}`); + await enqueueReelJob({ kind: 'video', userId: user.id, postId: post.id, params }); + console.log(`✨ Reel en file. Post ${post.id}, ${waiting} job(s) avant lui.`); - if (!isQueueBusy) { - // Démarrer le traitement en arrière-plan (Fire & Forget) - processReelBackground(user.id, post.id, req.body, storeName).catch(err => { - console.error('🔥 Unhandled background error:', err); - }); - } else { - console.log(`⏳ [Queue] Worker busy (count=${processingCount}). Job ${post.id} is queued.`); - } - - // Réponse immédiate au client res.json({ success: true, postId: post.id, - message: initialMessage, - queued: isQueueBusy, + queued: waiting > 0, + message: waiting > 0 + ? "File d'attente occupée. Votre vidéo sera traitée dès que possible." + : 'Traitement démarré en arrière-plan.', results: [], - videoUrl: "" + videoUrl: '', }); - } catch (error) { console.error('❌ Error initiating Reel:', error); res.status(500).json({ + success: false, error: error instanceof Error ? error.message : 'Erreur lors de l\'initialisation du Reel', }); } diff --git a/server/routes/remotion.ts b/server/routes/remotion.ts index 076d4cc..1fc3205 100644 --- a/server/routes/remotion.ts +++ b/server/routes/remotion.ts @@ -2,358 +2,104 @@ import { Router } from "express"; import multer from "multer"; import path from "path"; import fs from "fs"; -import { bundle } from "@remotion/bundler"; -import { renderMedia, selectComposition } from "@remotion/renderer"; -import { ffmpegService } from "../services/ffmpeg"; +import type { User } from "@shared/schema"; +import { imagesReelParamsSchema, type ImagesReelResult } from "@shared/reel"; import { storage as dbStorage } from "../storage"; -import { minioService as cloudinaryService, buildMinioUrl } from "../services/minio"; -import { facebookService } from "../services/facebook"; -import { tiktokService } from "../services/tiktok"; -import { generateVideoThumbnail, createVideoThumbnail } from "../services/thumbnail"; -import * as musicMetadata from "music-metadata"; - +import { enqueueReelJob, getReelJob } from "../services/reels/queue"; +import { REMOTION_TEMP_DIR } from "../services/reels/imagesPipeline"; +import { resolveStoreName } from "../services/reels/assets"; +import { publishReelToPages, storeRenderedVideo } from "../services/reels/publish"; export const remotionRouter = Router(); -interface RenderJob { - status: 'processing' | 'done' | 'error'; - url?: string; - thumbnailUrl?: string | null; - error?: string; -} -const renderJobs = new Map(); - -remotionRouter.get("/render/status/:jobId", (req, res) => { - const job = renderJobs.get(req.params.jobId); - if (!job) return res.status(404).json({ error: "Job introuvable" }); - res.json(job); -}); - -const uploadDir = path.join(process.cwd(), "uploads", "temp"); +const uploadDir = REMOTION_TEMP_DIR; if (!fs.existsSync(uploadDir)) { fs.mkdirSync(uploadDir, { recursive: true }); } -const storage = multer.diskStorage({ - destination: uploadDir, - filename: (_req, file, cb) => { - const uniqueSuffix = Date.now() + "-" + Math.round(Math.random() * 1e9); - cb(null, `remotion-${uniqueSuffix}${path.extname(file.originalname)}`); - }, +const upload = multer({ + storage: multer.diskStorage({ + destination: uploadDir, + filename: (_req, file, cb) => { + const uniqueSuffix = Date.now() + "-" + Math.round(Math.random() * 1e9); + cb(null, `remotion-${uniqueSuffix}${path.extname(file.originalname)}`); + }, + }), }); -const upload = multer({ storage }); - -/** Cached Remotion bundle URL — only bundle once per process lifecycle */ -let remotionBundleCache: string | null = null; - -async function getBundle(): Promise { - if (remotionBundleCache) return remotionBundleCache; - // Use process.cwd() (= /app in Docker, project root in dev) — works in both CJS and ESM - const entryPoint = path.resolve(process.cwd(), "client/src/remotion/index.ts"); - console.log("📦 Bundling Remotion from:", entryPoint); - remotionBundleCache = await bundle({ entryPoint }); - console.log("📦 Bundle cached at:", remotionBundleCache); - return remotionBundleCache; -} - /** - * Strips hashtags and emojis from text for TTS synthesis. - * Uses Unicode-aware regex to handle French accented chars in hashtags - * (e.g. #AménagementExtérieur must be fully removed, not just #Am). + * GET /api/remotion/render/status/:jobId + * État d'un rendu, lu dans la file persistante (survit à un redémarrage). */ -function stripForTTS(text: string): string { - return text - // Hashtags with Unicode letters (covers all French accented chars) - .replace(/#[\w\u00C0-\u024F\u1E00-\u1EFF]*/g, ' ') - // Emoji: surrogate pairs (covers virtually all emoji in UTF-16 strings) - .replace(/[\uD800-\uDFFF][\uDC00-\uDFFF]/g, ' ') // surrogate pairs (most emojis) - .replace(/[\u2600-\u27BF]/g, ' ') // misc symbols & arrows - .replace(/[\u2B00-\u2BFF]/g, ' ') // misc symbols extended - .replace(/\s+/g, ' ') - .trim(); -} - -/** Estimates syllable count for a French word (vowel-group method). */ -function countSyllablesFr(word: string): number { - const clean = word.replace(/[^a-zàâéèêëîïôùûüç]/gi, '').toLowerCase(); - if (!clean) return 1; - const groups = clean.match(/[aeiouyàâéèêëîïôùûü]+/gi); - return Math.max(1, groups?.length ?? 1); -} - -/** Purely-punctuation token (should not appear in overlay). */ -const PUNCT_ONLY = /^[.,!?;:…\-—«»"''()\[\]]+$/; - -/** - * Calculates word timings (in frames) for SPOKEN words only. - * Duration per word is weighted by syllable count → much tighter sync with TTS voice. - */ -function computeWordTimings( - displayText: string, - audioDurationSeconds: number, - fps: number, - startFrame: number -): Array<{ word: string; startFrame: number; endFrame: number }> { - const ttsText = stripForTTS(displayText); - // Filter out standalone punctuation tokens - const spokenWords = ttsText.split(/\s+/).filter(w => w && !PUNCT_ONLY.test(w)); - if (spokenWords.length === 0) return []; - - // Strip leading/trailing punctuation from each word for clean display - const cleanWords = spokenWords.map(w => w.replace(/^[.,!?;:…«»"''()\[\]]+|[.,!?;:…«»"''()\[\]]+$/g, '') || w); - - // Syllable-weighted distribution: longer words get proportionally more screen time - const syllables = cleanWords.map(countSyllablesFr); - const totalSyllables = syllables.reduce((a, b) => a + b, 0); - const totalFrames = audioDurationSeconds * fps; - - let currentFrame = startFrame; - return cleanWords.map((word, i) => { - const wordStart = currentFrame; - const wordFrames = Math.round((syllables[i] / totalSyllables) * totalFrames); - currentFrame += wordFrames; - return { word, startFrame: wordStart, endFrame: currentFrame }; - }); -} - -// Duration constants -const FPS = 30; -const ENDING_SECONDS = 3; // ending slide (logo + store name) -const MIN_CONTENT_SECONDS = 22; // so total >= 25s -const MAX_CONTENT_SECONDS = 27; // so total <= 30s - -/** - * Converts a local /uploads/... relative path to a base64 data URL. - * Chromium inside Docker cannot reliably fetch http://localhost via HTTP, - * so we embed images directly — zero network dependency. - */ -async function toDataUrl(relativeUrl: string): Promise { - if (!relativeUrl || !relativeUrl.startsWith('/uploads/')) return relativeUrl; - const filePath = path.join(process.cwd(), relativeUrl); - try { - const buffer = await fs.promises.readFile(filePath); - const ext = path.extname(filePath).slice(1).toLowerCase(); - const mime = ext === 'png' ? 'image/png' - : ext === 'gif' ? 'image/gif' - : ext === 'webp' ? 'image/webp' - : 'image/jpeg'; - return `data:${mime};base64,${buffer.toString('base64')}`; - } catch { - console.warn(`⚠️ toDataUrl: file not found: ${filePath}`); - return relativeUrl; +remotionRouter.get("/render/status/:jobId", async (req, res) => { + const user = req.user as User; + const job = await getReelJob(req.params.jobId); + if (!job || job.kind !== "images" || (job.userId !== user.id && user.role !== "admin")) { + return res.status(404).json({ error: "Job introuvable" }); } -} + if (job.status === "completed") { + const result = job.result as ImagesReelResult; + return res.json({ status: "done", url: result.url, thumbnailUrl: result.thumbnailUrl }); + } + if (job.status === "failed") { + return res.json({ status: "error", error: "Erreur lors du rendu: " + (job.error ?? "inconnue") }); + } + res.json({ status: "processing", progress: job.progress, step: job.step, queued: job.status === "pending" }); +}); + +/** + * POST /api/remotion/render + * Met en file le rendu d'un Reel à partir d'images (1 à 4). + */ remotionRouter.post("/render", upload.fields([{ name: "images", maxCount: 4 }, { name: "music", maxCount: 1 }]), async (req, res) => { + const fields = (req.files ?? {}) as Record; + const uploaded = [...(fields["images"] ?? []), ...(fields["music"] ?? [])]; + try { - const fields = req.files as Record; - const files = fields["images"] ?? []; - let imageUrls: string[] = []; - - // For audio/music we still need HTTP URLs (data URLs for audio are too large) - const port = process.env.PORT || "5555"; - const host = `http://localhost:${port}`; - - // Library images: resolve relative paths → data URLs (avoids Chromium HTTP issues) + const user = req.user as User; const existing = req.body.existingImageUrls; - if (existing) { - const rawUrls = Array.isArray(existing) ? existing : [existing]; - const resolved = await Promise.all(rawUrls.map(toDataUrl)); - imageUrls.push(...resolved); - } - // Newly uploaded temp files: read from disk → data URL - if (files && files.length > 0) { - const resolved = await Promise.all(files.map(async f => { - const buffer = await fs.promises.readFile(f.path); - const ext = path.extname(f.originalname).slice(1).toLowerCase(); - const mime = ext === 'png' ? 'image/png' : ext === 'gif' ? 'image/gif' : ext === 'webp' ? 'image/webp' : 'image/jpeg'; - return `data:${mime};base64,${buffer.toString('base64')}`; - })); - imageUrls.push(...resolved); - } + const libraryUrls: string[] = existing ? (Array.isArray(existing) ? existing : [existing]) : []; + const uploadedImageUrls = (fields["images"] ?? []).map((f) => `/uploads/temp/${path.basename(f.path)}`); - if (imageUrls.length === 0) { - return res.status(400).json({ error: "Aucune image fournie" }); - } - - const overlayText: string | undefined = typeof req.body.overlayText === 'string' && req.body.overlayText.trim() - ? req.body.overlayText.trim() - : undefined; - - const musicVolume: number = parseFloat(req.body.musicVolume ?? "0.3"); - - // Music: HTTP URLs are fine for audio (Remotion uses Web Audio API, not ) const musicFile = fields["music"]?.[0]; - const rawMusicTrackUrl = req.body.musicTrackUrl as string | undefined; - const musicUrl: string | undefined = musicFile - ? `${host}/uploads/temp/${path.basename(musicFile.path)}` - : rawMusicTrackUrl?.startsWith('/') ? `${host}${rawMusicTrackUrl}` : rawMusicTrackUrl || undefined; + const musicVolume = parseFloat(req.body.musicVolume ?? "0.3"); - // --- Fetch logo and store name from config --- - let logoUrl: string | undefined; - let storeName: string | undefined; - try { - const cloudinaryConfig = await dbStorage.getCloudinaryConfig(); - if (cloudinaryConfig?.cloudName && cloudinaryConfig?.logoPublicId) { - const relLogo = buildMinioUrl(cloudinaryConfig.cloudName, cloudinaryConfig.logoPublicId, cloudinaryConfig.publicUrl); - logoUrl = await toDataUrl(relLogo); - console.log("🏢 Logo embedded as data URL:", logoUrl.slice(0, 40) + "..."); - } - } catch (e) { - console.warn("⚠️ Could not fetch logo config:", e); - } - try { - const user = req.user as any; - if (user?.id) { - const pages = await dbStorage.getSocialPages(user.id); - if (pages.length > 0) { - // Use the page selected by the user, or fall back to first page - const selectedPageId = req.body.selectedPageId as string | undefined; - const page = selectedPageId ? pages.find(p => p.id === selectedPageId) : pages[0]; - storeName = (page ?? pages[0]).pageName; - console.log("🏪 Store name:", storeName); - } - } - } catch (e) { - console.warn("⚠️ Could not fetch store name:", e); - } - - // --- TTS Audio Generation --- - let audioUrl: string | undefined; - let wordTimings: Array<{ word: string; startFrame: number; endFrame: number }> | undefined; - let estimatedAudioDuration = 0; - - if (overlayText) { - const ttsText = stripForTTS(overlayText); - if (ttsText) { - console.log("🎙️ Generating TTS for:", ttsText); - try { - const ttsResult = await ffmpegService.previewTTS(ttsText); - if (ttsResult.success && ttsResult.audioBase64) { - const audioFilename = `tts-${Date.now()}.mp3`; - const audioPath = path.join(uploadDir, audioFilename); - const audioBuffer = Buffer.from(ttsResult.audioBase64, "base64"); - fs.writeFileSync(audioPath, audioBuffer); - audioUrl = `${host}/uploads/temp/${audioFilename}`; - - // Get exact audio duration from MP3 metadata - try { - const meta = await musicMetadata.parseFile(audioPath); - estimatedAudioDuration = meta.format.duration ?? 0; - } catch { - // Fallback: rough estimate from byte size (128kbps = 16KB/s) - estimatedAudioDuration = audioBuffer.length / 16000; - } - estimatedAudioDuration = Math.max(estimatedAudioDuration, ttsText.split(/\s+/).length * 0.35); - - wordTimings = computeWordTimings(overlayText, estimatedAudioDuration, FPS, 0); - console.log(`✅ TTS generated, ~${estimatedAudioDuration.toFixed(1)}s, ${wordTimings.length} words`); - } else { - console.warn("⚠️ TTS failed:", ttsResult.error); - } - } catch (ttsErr) { - console.warn("⚠️ TTS error (rendering without audio):", ttsErr); - } - } - } - - // --- Duration: 25-30s total --- - // content = max(tts duration, images*3s), clamped to 22-27s; ending = 3s - const naturalContent = Math.max(estimatedAudioDuration, imageUrls.length * 3); - const contentSeconds = Math.min(Math.max(naturalContent, MIN_CONTENT_SECONDS), MAX_CONTENT_SECONDS); - const endingFrames = (logoUrl || storeName) ? ENDING_SECONDS * FPS : 0; - const totalFrames = Math.round(contentSeconds * FPS) + endingFrames; - - console.log(`🎬 Duration: ${contentSeconds.toFixed(1)}s content + ${ENDING_SECONDS}s ending = ${(totalFrames / FPS).toFixed(1)}s total`); - - console.log("🎬 Getting Remotion bundle (cached)..."); - const bundleLocation = await getBundle(); - - const inputProps = { - images: imageUrls, - overlayText, - audioUrl, - wordTimings, - musicUrl, - musicVolume, - logoUrl, - storeName, - endingFrames, - }; - - console.log("🎬 Selecting composition..."); - const composition = await selectComposition({ - serveUrl: bundleLocation, - id: "ImageVideo", - inputProps, + const parsed = imagesReelParamsSchema.safeParse({ + imageUrls: [...libraryUrls, ...uploadedImageUrls], + overlayText: req.body.overlayText, + musicUrl: musicFile ? `/uploads/temp/${path.basename(musicFile.path)}` : req.body.musicTrackUrl, + musicVolume: Number.isFinite(musicVolume) ? musicVolume : undefined, + ttsEngine: req.body.ttsEngine || undefined, + ttsVoice: req.body.ttsVoice, + storeName: await resolveStoreName(user.id, req.body.selectedPageId), + tempFiles: uploaded.map((f) => f.path), }); + if (!parsed.success) { + await removeFiles(uploaded); + return res.status(400).json({ error: parsed.error.issues[0]?.message ?? "Paramètres invalides" }); + } - const jobId = Date.now().toString(); - renderJobs.set(jobId, { status: 'processing' }); - res.json({ jobId, message: "Rendu démarré en arrière-plan" }); - - // Run async to prevent 504 Gateway Timeout - (async () => { - try { - const outputFilename = `out-${Date.now()}.mp4`; - const outputLocation = path.join(uploadDir, outputFilename); - - console.log("🎬 Rendering media...", totalFrames, "frames (Job:", jobId, ")"); - - await renderMedia({ - composition: { ...composition, durationInFrames: totalFrames }, - serveUrl: bundleLocation, - codec: "h264", - outputLocation, - inputProps, - // Ne pas y remettre "args" : Remotion ne lit jamais cette clé (elle est - // absente de ChromiumOptions et de son code), les drapeaux passés là - // n'atteignaient donc jamais le navigateur. - chromiumOptions: { - disableWebSecurity: true, - ignoreCertificateErrors: true, - }, - concurrency: 1, // Limit concurrency to prevent Docker memory exhaustion - }); - - console.log("✅ Render completed:", outputFilename); - - // Generate thumbnail for the video - const thumbnailFilename = outputFilename.replace('.mp4', '-thumb.jpg'); - const thumbnailPath = path.join(uploadDir, thumbnailFilename); - const thumbnailGenerated = await generateVideoThumbnail(outputLocation, thumbnailPath, 2); - - if (thumbnailGenerated) { - console.log("🖼️ Thumbnail generated:", thumbnailFilename); - } - - renderJobs.set(jobId, { - status: 'done', - url: `/uploads/temp/${outputFilename}`, - thumbnailUrl: thumbnailGenerated ? `/uploads/temp/${thumbnailFilename}` : null, - }); - - } catch (err: any) { - console.error("❌ Remotion render error for Job", jobId, ":", err); - renderJobs.set(jobId, { status: 'error', error: "Erreur lors du rendu: " + err.message }); - } - })(); - - } catch (err: any) { + const job = await enqueueReelJob({ kind: "images", userId: user.id, params: parsed.data }); + res.json({ jobId: job.id, message: "Rendu démarré en arrière-plan" }); + } catch (err) { console.error("❌ Remotion pre-render error:", err); - res.status(500).json({ error: "Erreur lors du lancement du rendu: " + err.message }); + await removeFiles(uploaded); + res.status(500).json({ error: "Erreur lors du lancement du rendu: " + (err instanceof Error ? err.message : err) }); } }); +async function removeFiles(files: Express.Multer.File[]): Promise { + await Promise.all(files.map((f) => fs.promises.unlink(f.path).catch(() => { /* déjà absent */ }))); +} + /** * POST /api/remotion/publish - * Uploads the rendered MP4 to Cloudinary and publishes (or schedules) it as a Reel. + * Range la vidéo rendue dans la médiathèque puis la publie (ou la planifie). */ remotionRouter.post("/publish", async (req, res) => { try { - const user = req.user as any; - if (!user?.id) return res.status(401).json({ error: "Non authentifié" }); + const user = req.user as User; const { videoUrl, pageIds, scheduledFor, description } = req.body as { videoUrl: string; @@ -365,106 +111,36 @@ remotionRouter.post("/publish", async (req, res) => { if (!videoUrl) return res.status(400).json({ error: "videoUrl requis" }); if (!pageIds?.length) return res.status(400).json({ error: "Au moins une page requise" }); - // Resolve the local file path from the temp URL + // Chemin local du rendu à partir de son URL temporaire const filename = path.basename(videoUrl.split("?")[0]); const localPath = path.join(uploadDir, filename); if (!fs.existsSync(localPath)) { return res.status(400).json({ error: "Fichier vidéo introuvable (expiré ?)" }); } - console.log("☁️ Uploading Remotion video to Cloudinary..."); - const cloudinaryResult = await cloudinaryService.uploadMedia( - localPath, - filename, - user.id, - "video/mp4" - ); - console.log("✅ Cloudinary upload:", cloudinaryResult.originalUrl); + const videoBuffer = await fs.promises.readFile(localPath); + const mediaRecord = await storeRenderedVideo(user.id, videoBuffer, filename); - // Create a Post record to track this publication const post = await dbStorage.createPost({ userId: user.id, content: description || "", aiGenerated: "false", - status: scheduledFor ? "scheduled" : "published", + status: scheduledFor ? "scheduled" : "draft", scheduledFor: scheduledFor ? new Date(scheduledFor) : undefined, }); - - // Create a Media record for the processed video - // Vignette extraite avant que la vidéo ne soit purgée du disque - const thumbnailUrl = await createVideoThumbnail(localPath); - - const mediaRecord = await dbStorage.createMedia({ - userId: user.id, - type: "video", - cloudinaryPublicId: cloudinaryResult.publicId, - originalUrl: cloudinaryResult.originalUrl, - facebookFeedUrl: null, - instagramFeedUrl: null, - instagramStoryUrl: null, - thumbnailUrl, - fileName: filename, - fileSize: fs.statSync(localPath).size, - }); - await dbStorage.updatePostMedia(post.id, [mediaRecord.id]); - const results: { pageId: string; success: boolean; reelId?: string; error?: string }[] = []; + const results = await publishReelToPages({ + postId: post.id, + pageIds, + videoBuffer, + description: description || "", + scheduledFor, + }); - for (const pageId of pageIds) { - try { - const page = await dbStorage.getSocialPage(pageId); - if (!page) { results.push({ pageId, success: false, error: "Page introuvable" }); continue; } - if (page.platform !== "facebook" && page.platform !== "tiktok") { - results.push({ pageId, success: false, error: `Plateforme non supportée pour les reels : ${page.platform}` }); - continue; - } - - const scheduledPost = await dbStorage.createScheduledPost({ - postId: post.id, - pageId: page.id, - postType: "reel", - scheduledAt: scheduledFor ? new Date(scheduledFor) : new Date(), - }); - - if (!scheduledFor) { - console.log(`🚀 Publishing to ${page.platform} ${page.pageName} (direct binary upload)...`); - // Upload video bytes directly — les APIs ne peuvent pas lire une URL locale - const videoBuffer = await fs.promises.readFile(localPath); - - if (page.platform === "facebook") { - const reelId = await facebookService.publishVideoFromBuffer(page, videoBuffer, description || ""); - await dbStorage.updateScheduledPost(scheduledPost.id, { publishedAt: new Date(), externalPostId: reelId }); - results.push({ pageId, success: true, reelId }); - } else { - // TikTok finalise la publication de son côté : le poller de statut - // renseignera l'identifiant définitif du post. - const publishId = await tiktokService.publishVideoFromBuffer(page, videoBuffer, description || ""); - await dbStorage.updateScheduledPost(scheduledPost.id, { - publishedAt: new Date(), - publishId, - publishStatus: "PROCESSING_UPLOAD", - }); - results.push({ pageId, success: true, reelId: publishId }); - } - } else { - results.push({ pageId, success: true, reelId: "scheduled" }); - } - } catch (pageErr: any) { - console.error(`❌ Error publishing to page ${pageId}:`, pageErr); - results.push({ pageId, success: false, error: pageErr.message }); - } - } - - const anySuccess = results.some(r => r.success); - if (!anySuccess) { - await dbStorage.updatePost(post.id, { status: "failed" }); - } - - res.json({ success: anySuccess, results, postId: post.id }); - - } catch (err: any) { + res.json({ success: results.some((r) => r.success), results, postId: post.id }); + } catch (err) { console.error("❌ Remotion publish error:", err); - res.status(500).json({ error: "Erreur lors de la publication: " + err.message }); + res.status(500).json({ error: "Erreur lors de la publication: " + (err instanceof Error ? err.message : err) }); } }); diff --git a/server/services/ffmpeg.ts b/server/services/ffmpeg.ts index 823dddd..94a87a3 100644 --- a/server/services/ffmpeg.ts +++ b/server/services/ffmpeg.ts @@ -33,6 +33,11 @@ interface FFmpegReelResponse { tts_error?: string; } +/** Un rendu long (stabilisation + encodage) peut dépasser plusieurs minutes. */ +const PROCESS_TIMEOUT_MS = 15 * 60_000; +const TTS_TIMEOUT_MS = 2 * 60_000; +const HEALTH_TIMEOUT_MS = 5_000; + interface FFmpegConfig { apiUrl: string; apiKey: string; @@ -131,6 +136,7 @@ export class FFmpegService { 'X-API-Key': config.apiKey, }, body: JSON.stringify(requestBody), + signal: AbortSignal.timeout(PROCESS_TIMEOUT_MS), }); if (!response.ok) { @@ -240,6 +246,7 @@ export class FFmpegService { 'X-API-Key': config.apiKey, }, body: JSON.stringify(requestBody), + signal: AbortSignal.timeout(PROCESS_TIMEOUT_MS), }); if (!response.ok) { @@ -292,6 +299,7 @@ export class FFmpegService { headers: { 'X-API-Key': config.apiKey, }, + signal: AbortSignal.timeout(HEALTH_TIMEOUT_MS), }); return response.ok; } catch { @@ -320,6 +328,7 @@ export class FFmpegService { tts_engine: ttsEngine, gemini_api_key: geminiApiKey, }), + signal: AbortSignal.timeout(TTS_TIMEOUT_MS), }); if (!response.ok) { diff --git a/server/services/reels/assets.ts b/server/services/reels/assets.ts new file mode 100644 index 0000000..a8c2794 --- /dev/null +++ b/server/services/reels/assets.ts @@ -0,0 +1,65 @@ +/** + * Résolution des ressources d'un Reel (musique, logo, clé Gemini, nom du + * magasin), jusqu'ici recopiée dans chaque route. + */ + +import { storage } from "../../storage"; +import { buildMinioUrl, resolveInternalUrl } from "../minio"; + +/** + * URL de la musique téléchargeable par le service FFmpeg (réseau Docker + * interne : l'URL publique HTTPS n'y est pas joignable). + */ +export async function resolveMusicUrl( + musicTrackId?: string, + musicUrl?: string, +): Promise { + if (musicUrl) return resolveInternalUrl(musicUrl); + if (!musicTrackId?.startsWith("internal_")) return undefined; + + try { + const track = await storage.getAudioTrack(musicTrackId.slice("internal_".length)); + return track ? resolveInternalUrl(track.url) : undefined; + } catch (error) { + console.error("⚠️ [Reels] Piste audio introuvable :", error); + return undefined; + } +} + +/** Chemin relatif (/uploads/...) du logo configuré, s'il existe. */ +export async function resolveLogoPath(): Promise { + try { + const config = await storage.getCloudinaryConfig(); + if (config?.logoPublicId) { + return buildMinioUrl(config.cloudName, config.logoPublicId, config.publicUrl); + } + } catch (error) { + console.error("⚠️ [Reels] Configuration du logo illisible :", error); + } + return undefined; +} + +/** Clé Gemini : configuration de l'application, sinon variable d'environnement. */ +export async function resolveGeminiApiKey(ttsEngine?: string): Promise { + if (ttsEngine !== "gemini") return undefined; + const appConfig = await storage.getAppConfig(); + return appConfig?.geminiApiKey ?? process.env.GEMINI_API_KEY ?? undefined; +} + +/** Nom affiché dans l'outro : celui de la page choisie, sinon de la première page de l'utilisateur. */ +export async function resolveStoreName( + userId: string, + preferredPageId?: string, +): Promise { + try { + if (preferredPageId) { + const page = await storage.getSocialPage(preferredPageId); + if (page) return page.pageName; + } + const pages = await storage.getSocialPages(userId); + return pages[0]?.pageName; + } catch (error) { + console.error("⚠️ [Reels] Nom du magasin introuvable :", error); + return undefined; + } +} diff --git a/server/services/reels/imagesPipeline.ts b/server/services/reels/imagesPipeline.ts new file mode 100644 index 0000000..da7ed58 --- /dev/null +++ b/server/services/reels/imagesPipeline.ts @@ -0,0 +1,234 @@ +/** + * Rendu d'un Reel à partir d'images avec Remotion. + * Exécuté par la file `reel_jobs` ; la publication reste une étape distincte, + * déclenchée par l'utilisateur après l'aperçu. + */ + +import fs from "fs"; +import path from "path"; +import * as musicMetadata from "music-metadata"; +import { bundle } from "@remotion/bundler"; +import { renderMedia, selectComposition } from "@remotion/renderer"; +import { imagesReelParamsSchema, type ImagesReelResult } from "@shared/reel"; +import { ffmpegService } from "../ffmpeg"; +import { generateVideoThumbnail } from "../thumbnail"; +import type { JobContext } from "./queue"; +import { resolveGeminiApiKey, resolveLogoPath } from "./assets"; + +export const REMOTION_TEMP_DIR = path.join(process.cwd(), "uploads", "temp"); + +// Durées +const FPS = 30; +const ENDING_SECONDS = 3; // diapositive de fin (logo + nom du magasin) +const MIN_CONTENT_SECONDS = 22; // total >= 25 s +const MAX_CONTENT_SECONDS = 27; // total <= 30 s + +/** Bundle Remotion mis en cache : un seul bundle par processus. */ +let bundleCache: Promise | null = null; + +function getBundle(): Promise { + if (!bundleCache) { + // process.cwd() = /app dans Docker, racine du projet en dev + const entryPoint = path.resolve(process.cwd(), "client/src/remotion/index.ts"); + console.log("📦 Bundling Remotion from:", entryPoint); + bundleCache = bundle({ entryPoint }).catch((error) => { + bundleCache = null; // réessayer au prochain rendu + throw error; + }); + } + return bundleCache; +} + +/** URL HTTP locale : l'audio est lu par Chromium, trop lourd pour une data URL. */ +function localHttpUrl(url: string): string { + if (!url.startsWith("/")) return url; + const port = process.env.PORT || "5555"; + return `http://localhost:${port}${url}`; +} + +const MIME_BY_EXT: Record = { + png: "image/png", + gif: "image/gif", + webp: "image/webp", +}; + +/** + * Convertit un chemin local /uploads/... en data URL : Chromium dans Docker ne + * lit pas http://localhost de façon fiable, les images sont donc embarquées. + */ +async function toDataUrl(relativeUrl: string): Promise { + if (!relativeUrl.startsWith("/uploads/")) return relativeUrl; + const uploadsRoot = path.join(process.cwd(), "uploads") + path.sep; + const filePath = path.resolve(process.cwd(), "." + relativeUrl); + if (!filePath.startsWith(uploadsRoot)) { + throw new Error(`Chemin d'image refusé : ${relativeUrl}`); + } + try { + const buffer = await fs.promises.readFile(filePath); + const ext = path.extname(filePath).slice(1).toLowerCase(); + const mime = MIME_BY_EXT[ext] ?? "image/jpeg"; + return `data:${mime};base64,${buffer.toString("base64")}`; + } catch { + console.warn(`⚠️ toDataUrl: fichier introuvable : ${filePath}`); + return relativeUrl; + } +} + +/** + * Retire hashtags et emojis avant la synthèse vocale. Les lettres accentuées + * font partie du hashtag (#AménagementExtérieur est retiré en entier). + */ +export function stripForTTS(text: string): string { + return text + .replace(/#[\wÀ-ɏḀ-ỿ]*/g, " ") + .replace(/[\uD800-\uDFFF][\uDC00-\uDFFF]/g, " ") // paires de substitution (la plupart des emojis) + .replace(/[☀-➿]/g, " ") // symboles divers + .replace(/[⬀-⯿]/g, " ") + .replace(/\s+/g, " ") + .trim(); +} + +/** Nombre de syllabes d'un mot français (groupes de voyelles). */ +function countSyllablesFr(word: string): number { + const clean = word.replace(/[^a-zàâéèêëîïôùûüç]/gi, "").toLowerCase(); + if (!clean) return 1; + return Math.max(1, clean.match(/[aeiouyàâéèêëîïôùûü]+/gi)?.length ?? 1); +} + +const PUNCT_ONLY = /^[.,!?;:…\-—«»"''()\[\]]+$/; +const EDGE_PUNCT = /^[.,!?;:…«»"''()\[\]]+|[.,!?;:…«»"''()\[\]]+$/g; + +/** + * Timings des mots prononcés (en images), répartis au prorata des syllabes. + * Estimation provisoire : remplacée par les vrais timings de la voix au lot + * « voix et sous-titres ». + */ +export function computeWordTimings( + displayText: string, + audioDurationSeconds: number, + fps: number, + startFrame: number, +): Array<{ word: string; startFrame: number; endFrame: number }> { + const spokenWords = stripForTTS(displayText) + .split(/\s+/) + .filter((w) => w && !PUNCT_ONLY.test(w)); + if (spokenWords.length === 0) return []; + + const cleanWords = spokenWords.map((w) => w.replace(EDGE_PUNCT, "") || w); + const syllables = cleanWords.map(countSyllablesFr); + const totalSyllables = syllables.reduce((a, b) => a + b, 0); + const totalFrames = audioDurationSeconds * fps; + + let currentFrame = startFrame; + return cleanWords.map((word, i) => { + const wordStart = currentFrame; + currentFrame += Math.round((syllables[i] / totalSyllables) * totalFrames); + return { word, startFrame: wordStart, endFrame: currentFrame }; + }); +} + +export async function runImagesReelJob({ job, progress }: JobContext): Promise { + const params = imagesReelParamsSchema.parse(job.params); + const jobTempFiles: string[] = []; + + try { + await progress(5, "prepare"); + const images = await Promise.all(params.imageUrls.map(toDataUrl)); + const logoPath = await resolveLogoPath(); + const logoUrl = logoPath ? await toDataUrl(logoPath) : undefined; + const musicUrl = params.musicUrl ? localHttpUrl(params.musicUrl) : undefined; + const { overlayText, storeName } = params; + + // --- Voix --- + let audioUrl: string | undefined; + let wordTimings: ReturnType | undefined; + let audioDuration = 0; + + const ttsText = overlayText ? stripForTTS(overlayText) : ""; + if (overlayText && ttsText) { + await progress(15, "voice"); + const geminiApiKey = await resolveGeminiApiKey(params.ttsEngine); + const tts = await ffmpegService.previewTTS(ttsText, params.ttsVoice, params.ttsEngine, geminiApiKey); + if (!tts.success || !tts.audioBase64) { + throw new Error(`La voix n'a pas pu être générée : ${tts.error ?? "réponse vide"}`); + } + + const audioFilename = `tts-${job.id}.mp3`; + const audioPath = path.join(REMOTION_TEMP_DIR, audioFilename); + const audioBuffer = Buffer.from(tts.audioBase64, "base64"); + await fs.promises.writeFile(audioPath, audioBuffer); + jobTempFiles.push(audioPath); + audioUrl = localHttpUrl(`/uploads/temp/${audioFilename}`); + + try { + audioDuration = (await musicMetadata.parseFile(audioPath)).format.duration ?? 0; + } catch { + audioDuration = audioBuffer.length / 16000; // estimation à 128 kb/s + } + audioDuration = Math.max(audioDuration, ttsText.split(/\s+/).length * 0.35); + wordTimings = computeWordTimings(overlayText, audioDuration, FPS, 0); + } + + // --- Durée : 25 à 30 s --- + const naturalContent = Math.max(audioDuration, images.length * 3); + const contentSeconds = Math.min(Math.max(naturalContent, MIN_CONTENT_SECONDS), MAX_CONTENT_SECONDS); + const endingFrames = logoUrl || storeName ? ENDING_SECONDS * FPS : 0; + const totalFrames = Math.round(contentSeconds * FPS) + endingFrames; + + await progress(25, "render"); + const serveUrl = await getBundle(); + const inputProps = { + images, + overlayText, + audioUrl, + wordTimings, + musicUrl, + musicVolume: params.musicVolume, + logoUrl, + storeName, + endingFrames, + }; + const composition = await selectComposition({ serveUrl, id: "ImageVideo", inputProps }); + + const outputFilename = `out-${job.id}.mp4`; + const outputLocation = path.join(REMOTION_TEMP_DIR, outputFilename); + let lastReported = 25; + + await renderMedia({ + composition: { ...composition, durationInFrames: totalFrames }, + serveUrl, + codec: "h264", + outputLocation, + inputProps, + chromiumOptions: { + disableWebSecurity: true, + ignoreCertificateErrors: true, + }, + concurrency: 1, // évite d'épuiser la mémoire du conteneur + onProgress: ({ progress: ratio }) => { + const pct = 25 + Math.floor(ratio * 70); + if (pct >= lastReported + 5) { + lastReported = pct; + progress(pct, "render").catch(() => { /* non bloquant */ }); + } + }, + }); + + const thumbnailFilename = outputFilename.replace(".mp4", "-thumb.jpg"); + const thumbnailOk = await generateVideoThumbnail( + outputLocation, + path.join(REMOTION_TEMP_DIR, thumbnailFilename), + 2, + ); + + return { + url: `/uploads/temp/${outputFilename}`, + thumbnailUrl: thumbnailOk ? `/uploads/temp/${thumbnailFilename}` : null, + }; + } finally { + // Images, musique importée et voix ne servent plus une fois le rendu fini + for (const file of [...params.tempFiles, ...jobTempFiles]) { + await fs.promises.unlink(file).catch(() => { /* déjà absent */ }); + } + } +} diff --git a/server/services/reels/index.ts b/server/services/reels/index.ts new file mode 100644 index 0000000..a244559 --- /dev/null +++ b/server/services/reels/index.ts @@ -0,0 +1,14 @@ +/** + * Point d'entrée des Reels : enregistre les traitements de la file et démarre + * le worker. + */ + +import { registerReelJobHandler, startReelWorker } from "./queue"; +import { runVideoReelJob } from "./videoPipeline"; +import { runImagesReelJob } from "./imagesPipeline"; + +export function startReelProcessing(): Promise { + registerReelJobHandler("video", runVideoReelJob); + registerReelJobHandler("images", runImagesReelJob); + return startReelWorker(); +} diff --git a/server/services/reels/publish.ts b/server/services/reels/publish.ts new file mode 100644 index 0000000..0a7a084 --- /dev/null +++ b/server/services/reels/publish.ts @@ -0,0 +1,133 @@ +/** + * Enregistrement et publication d'un Reel rendu, commun aux Reels vidéo et + * images. + * + * La vidéo est toujours envoyée en binaire : les fichiers vivent sous + * /uploads, une URL relative que Facebook ne peut pas télécharger (c'était le + * cas de l'ancien `publishReel(originalUrl)`). + */ + +import type { Media } from "@shared/schema"; +import { storage } from "../../storage"; +import { minioService } from "../minio"; +import { createVideoThumbnail } from "../thumbnail"; +import { facebookService } from "../facebook"; +import { tiktokService } from "../tiktok"; + +export interface PagePublishResult { + pageId: string; + success: boolean; + reelId?: string; + error?: string; +} + +/** Range la vidéo rendue dans la médiathèque, avec sa vignette. */ +export async function storeRenderedVideo( + userId: string, + videoBuffer: Buffer, + fileName: string, +): Promise { + const uploaded = await minioService.uploadMedia(videoBuffer, fileName, userId, "video/mp4"); + // Vignette extraite maintenant : la vidéo sera supprimée du disque après + // publication, mais l'historique doit rester illustré. + const thumbnailUrl = await createVideoThumbnail(videoBuffer); + + return storage.createMedia({ + userId, + type: "video", + cloudinaryPublicId: uploaded.publicId, + originalUrl: uploaded.originalUrl, + facebookFeedUrl: uploaded.facebookFeedUrl, + instagramFeedUrl: uploaded.instagramFeedUrl, + instagramStoryUrl: uploaded.instagramStoryUrl, + thumbnailUrl, + fileName, + fileSize: videoBuffer.length, + }); +} + +/** + * Publie (ou planifie) la vidéo sur chaque page, crée les entrées + * `scheduled_posts` et met à jour le statut du post. + */ +export async function publishReelToPages(options: { + postId: string; + pageIds: string[]; + videoBuffer: Buffer; + description: string; + scheduledFor?: string; +}): Promise { + const { postId, pageIds, videoBuffer, description, scheduledFor } = options; + const results: PagePublishResult[] = []; + + for (const pageId of pageIds) { + try { + const page = await storage.getSocialPage(pageId); + if (!page) { + results.push({ pageId, success: false, error: "Page non trouvée" }); + continue; + } + if (page.platform !== "facebook" && page.platform !== "tiktok") { + results.push({ + pageId, + success: false, + error: `Plateforme non supportée pour les reels : ${page.platform}`, + }); + continue; + } + + const scheduledPost = await storage.createScheduledPost({ + postId, + pageId: page.id, + postType: "reel", + scheduledAt: scheduledFor ? new Date(scheduledFor) : new Date(), + }); + + if (scheduledFor) { + // Le planificateur publiera à l'heure dite + results.push({ pageId, success: true, reelId: "scheduled" }); + continue; + } + + console.log(`🚀 [Reels] Publication sur ${page.platform} ${page.pageName}...`); + if (page.platform === "facebook") { + const reelId = await facebookService.publishVideoFromBuffer(page, videoBuffer, description); + await storage.updateScheduledPost(scheduledPost.id, { + publishedAt: new Date(), + externalPostId: reelId, + }); + results.push({ pageId, success: true, reelId }); + } else { + // TikTok finalise la publication de son côté : l'envoi est marqué comme + // effectué pour ne pas republier, et le poller de statut renseignera + // l'identifiant définitif du post. + const publishId = await tiktokService.publishVideoFromBuffer(page, videoBuffer, description); + await storage.updateScheduledPost(scheduledPost.id, { + publishedAt: new Date(), + publishId, + publishStatus: "PROCESSING_UPLOAD", + }); + results.push({ pageId, success: true, reelId: publishId }); + } + } catch (error) { + console.error(`❌ [Reels] Publication sur la page ${pageId} impossible :`, error); + results.push({ + pageId, + success: false, + error: error instanceof Error ? error.message : "Erreur inconnue", + }); + } + } + + const succeeded = results.filter((r) => r.success).length; + const status = scheduledFor && succeeded > 0 ? "scheduled" : succeeded > 0 ? "published" : "failed"; + await storage.updatePost(postId, { status }); + + return results; +} + +/** Message d'erreur lisible lorsque toutes les publications ont échoué. */ +export function describePublishFailure(results: PagePublishResult[]): string { + const details = results.map((r) => r.error).filter(Boolean).join(" ; "); + return `Publication impossible sur toutes les pages${details ? ` : ${details}` : ""}`; +} diff --git a/server/services/reels/queue.ts b/server/services/reels/queue.ts new file mode 100644 index 0000000..f96a6cc --- /dev/null +++ b/server/services/reels/queue.ts @@ -0,0 +1,277 @@ +/** + * File d'attente persistante des rendus de Reels. + * + * L'ancienne file comptait les posts « processing » pour décider de lancer un + * rendu : un redémarrage pendant un job laissait ce compteur à 1 pour toujours + * et bloquait tous les Reels suivants. Ici, chaque job est une ligne de + * `reel_jobs` : + * - le worker la réserve atomiquement (`FOR UPDATE SKIP LOCKED`), ce qui exclut + * tout double traitement, même avec plusieurs instances ; + * - un battement de cœur est écrit pendant le rendu ; au démarrage, les jobs + * dont le cœur ne bat plus sont remis en file (ou échouent après + * MAX_ATTEMPTS tentatives). + */ + +import { and, eq, sql } from "drizzle-orm"; +import { db } from "../../db"; +import { posts, reelJobs, type ReelJob } from "@shared/schema"; +import type { ReelJobKind, ReelJobStatus } from "@shared/reel"; + +const HEARTBEAT_INTERVAL_MS = 15_000; +/** Au-delà, un job « processing » est considéré comme orphelin. */ +const STALE_AFTER_MS = 2 * 60_000; +const POLL_INTERVAL_MS = 10_000; +const MAX_ATTEMPTS = 2; + +/** Outils fournis à un handler pour rendre compte de son avancement. */ +export interface JobContext { + job: ReelJob; + progress(progress: number, step?: string): Promise; +} + +export type JobHandler = (ctx: JobContext) => Promise; + +const handlers = new Map(); +let running = false; +let started = false; +let pollTimer: NodeJS.Timeout | null = null; + +export function registerReelJobHandler(kind: ReelJobKind, handler: JobHandler): void { + handlers.set(kind, handler); +} + +/** Ajoute un job à la file et réveille le worker. */ +export async function enqueueReelJob(input: { + kind: ReelJobKind; + userId: string; + postId?: string; + params: unknown; +}): Promise { + const [job] = await db + .insert(reelJobs) + .values({ + kind: input.kind, + userId: input.userId, + postId: input.postId, + params: input.params, + }) + .returning(); + + if (job.postId) { + await syncPost(job.postId, "pending", 0); + } + + kick(); + return job; +} + +export async function getReelJob(id: string): Promise { + const [job] = await db.select().from(reelJobs).where(eq(reelJobs.id, id)); + return job; +} + +/** Nombre de jobs qui passeront avant un nouveau job (utile pour informer l'utilisateur). */ +export async function countActiveReelJobs(): Promise { + const [row] = await db + .select({ count: sql`count(*)::int` }) + .from(reelJobs) + .where(sql`${reelJobs.status} IN ('pending', 'processing')`); + return row?.count ?? 0; +} + +/** + * Démarre le worker : récupère les jobs orphelins puis traite la file. + * Idempotent. + */ +export async function startReelWorker(): Promise { + if (started) return; + started = true; + + try { + await recoverStaleJobs(); + await failLegacyStuckPosts(); + } catch (error) { + console.error("❌ [ReelQueue] Récupération des jobs orphelins impossible :", error); + } + + pollTimer = setInterval(kick, POLL_INTERVAL_MS); + pollTimer.unref(); + kick(); +} + +/** Relance la boucle si elle est à l'arrêt. */ +function kick(): void { + if (!started || running) return; + running = true; + drain() + .catch((error) => console.error("❌ [ReelQueue] Erreur de la boucle :", error)) + .finally(() => { + running = false; + }); +} + +async function drain(): Promise { + for (;;) { + const job = await claimNextJob(); + if (!job) return; + await runJob(job); + } +} + +async function claimNextJob(): Promise { + const result = await db.execute<{ id: string }>(sql` + UPDATE reel_jobs + SET status = 'processing', + attempts = attempts + 1, + heartbeat_at = now(), + updated_at = now(), + error = NULL + WHERE id = ( + SELECT id FROM reel_jobs + WHERE status = 'pending' + ORDER BY created_at + FOR UPDATE SKIP LOCKED + LIMIT 1 + ) + RETURNING id + `); + const id = result.rows[0]?.id; + return id ? getReelJob(id) : undefined; +} + +async function runJob(job: ReelJob): Promise { + const handler = handlers.get(job.kind as ReelJobKind); + console.log(`🎬 [ReelQueue] Job ${job.id} (${job.kind}), tentative ${job.attempts}`); + + const heartbeat = setInterval(() => { + db.update(reelJobs) + .set({ heartbeatAt: new Date() }) + .where(eq(reelJobs.id, job.id)) + .catch((error) => console.warn(`⚠️ [ReelQueue] Heartbeat ${job.id} :`, error)); + }, HEARTBEAT_INTERVAL_MS); + + const ctx: JobContext = { + job, + async progress(progress, step) { + await updateJob(job, { progress, step }); + }, + }; + + try { + if (!handler) { + throw new Error(`Aucun traitement enregistré pour les jobs « ${job.kind} »`); + } + if (job.postId) { + await syncPost(job.postId, "processing", 0); + } + const result = await handler(ctx); + await updateJob(job, { status: "completed", progress: 100, step: "done", result: result ?? null }); + console.log(`✅ [ReelQueue] Job ${job.id} terminé`); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + console.error(`❌ [ReelQueue] Job ${job.id} en échec :`, error); + await updateJob(job, { status: "failed", error: message }).catch((e) => + console.error(`⚠️ [ReelQueue] Statut d'échec non enregistré pour ${job.id} :`, e), + ); + } finally { + clearInterval(heartbeat); + } +} + +async function updateJob( + job: ReelJob, + patch: { + status?: ReelJobStatus; + progress?: number; + step?: string; + result?: unknown; + error?: string; + }, +): Promise { + await db + .update(reelJobs) + .set({ ...patch, heartbeatAt: new Date(), updatedAt: new Date() }) + .where(eq(reelJobs.id, job.id)); + + if (job.postId) { + const status = patch.status ?? "processing"; + const progress = patch.status === "failed" ? 0 : patch.progress; + await syncPost(job.postId, status, progress, patch.error); + } +} + +/** + * Reflète l'état du job sur le post : l'interface (`/reels/ongoing`) lit les + * colonnes `generation_*` du post. + */ +async function syncPost( + postId: string, + status: ReelJobStatus, + progress?: number, + error?: string, +): Promise { + const patch: Partial = { + generationStatus: status, + updatedAt: new Date(), + }; + if (progress !== undefined) patch.generationProgress = progress; + if (error !== undefined) patch.generationError = error; + if (status === "failed") patch.status = "failed"; + await db.update(posts).set(patch).where(eq(posts.id, postId)); +} + +/** + * Remet en file les jobs interrompus (serveur arrêté pendant un rendu), ou les + * fait échouer s'ils ont déjà épuisé leurs tentatives : un job qui fait tomber + * le serveur ne doit pas le faire tomber en boucle. + */ +async function recoverStaleJobs(): Promise { + const staleBefore = new Date(Date.now() - STALE_AFTER_MS); + const stale = await db + .select() + .from(reelJobs) + .where( + and( + eq(reelJobs.status, "processing"), + sql`(${reelJobs.heartbeatAt} IS NULL OR ${reelJobs.heartbeatAt} < ${staleBefore})`, + ), + ); + + for (const job of stale) { + if (job.attempts >= MAX_ATTEMPTS) { + console.warn(`⚠️ [ReelQueue] Job ${job.id} abandonné après ${job.attempts} tentatives`); + await updateJob(job, { + status: "failed", + error: "Le rendu a été interrompu à plusieurs reprises (redémarrage du serveur).", + }); + } else { + console.warn(`♻️ [ReelQueue] Job ${job.id} interrompu, remis en file`); + await db + .update(reelJobs) + .set({ status: "pending", step: null, progress: 0, updatedAt: new Date() }) + .where(eq(reelJobs.id, job.id)); + if (job.postId) await syncPost(job.postId, "pending", 0); + } + } +} + +/** + * Posts laissés « processing »/« pending » par l'ancienne file (avant + * reel_jobs) : ils n'ont aucun job pour les reprendre et bloquaient l'affichage + * des Reels en cours. On les marque en échec pour que l'utilisateur relance. + */ +async function failLegacyStuckPosts(): Promise { + const result = await db.execute(sql` + UPDATE posts + SET generation_status = 'failed', + generation_progress = 0, + generation_error = 'Génération interrompue, veuillez relancer le Reel.', + status = 'failed', + updated_at = now() + WHERE generation_status IN ('processing', 'pending') + AND NOT EXISTS (SELECT 1 FROM reel_jobs j WHERE j.post_id = posts.id) + `); + if (result.rowCount) { + console.warn(`⚠️ [ReelQueue] ${result.rowCount} Reel(s) de l'ancienne file marqués en échec`); + } +} diff --git a/server/services/reels/videoPipeline.ts b/server/services/reels/videoPipeline.ts new file mode 100644 index 0000000..8562d1f --- /dev/null +++ b/server/services/reels/videoPipeline.ts @@ -0,0 +1,78 @@ +/** + * Rendu d'un Reel à partir d'une vidéo : FFmpeg → médiathèque → publication. + * Exécuté par la file `reel_jobs` (voir queue.ts). + */ + +import { videoReelParamsSchema } from "@shared/reel"; +import { storage } from "../../storage"; +import { ffmpegService } from "../ffmpeg"; +import { resolveInternalUrl } from "../minio"; +import type { JobContext } from "./queue"; +import { resolveGeminiApiKey, resolveLogoPath, resolveMusicUrl } from "./assets"; +import { describePublishFailure, publishReelToPages, storeRenderedVideo } from "./publish"; + +export async function runVideoReelJob({ job, progress }: JobContext) { + if (!job.postId) throw new Error("Job vidéo sans post associé"); + const postId = job.postId; + const params = videoReelParamsSchema.parse(job.params); + + await progress(5, "prepare"); + const media = await storage.getMediaById(params.videoMediaId); + if (!media || media.type !== "video") { + throw new Error("Vidéo source introuvable ou invalide"); + } + + const [musicUrl, logoPath, geminiApiKey] = await Promise.all([ + resolveMusicUrl(params.musicTrackId, params.musicUrl), + resolveLogoPath(), + resolveGeminiApiKey(params.ttsEngine), + ]); + + await progress(15, "render"); + const startedAt = Date.now(); + const rendered = await ffmpegService.processReelFromUrl(resolveInternalUrl(media.originalUrl), { + text: params.overlayText, + musicUrl, + ttsEnabled: params.ttsEnabled, + ttsVoice: params.ttsVoice, + ttsEngine: params.ttsEngine, + geminiApiKey, + wordDuration: params.wordDuration, + fontSize: params.fontSize, + musicVolume: params.musicVolume, + drawText: params.drawText, + stabilize: params.stabilize, + watermarkUrl: logoPath ? resolveInternalUrl(logoPath) : undefined, + storeName: params.storeName, + enableEndingEffect: params.enableEndingEffect, + }); + console.log(`⏱️ [Reels] FFmpeg : ${((Date.now() - startedAt) / 1000).toFixed(1)} s`); + + if (!rendered.success || !rendered.videoBase64) { + throw new Error(rendered.error || "Erreur de traitement vidéo FFmpeg"); + } + if (rendered.ttsError) { + // Voix demandée mais absente : ne jamais publier un Reel muet sans le dire + throw new Error(`La voix n'a pas pu être générée : ${rendered.ttsError}`); + } + + await progress(65, "store"); + const videoBuffer = Buffer.from(rendered.videoBase64, "base64"); + const processedMedia = await storeRenderedVideo(job.userId, videoBuffer, `reel-${Date.now()}.mp4`); + await storage.updatePostMedia(postId, [processedMedia.id]); + + await progress(85, "publish"); + const results = await publishReelToPages({ + postId, + pageIds: params.pageIds, + videoBuffer, + description: params.description || params.overlayText || "", + scheduledFor: params.scheduledFor, + }); + + if (!results.some((r) => r.success)) { + throw new Error(describePublishFailure(results)); + } + + return { mediaId: processedMedia.id, results }; +} diff --git a/server/services/thumbnail.ts b/server/services/thumbnail.ts index 1123e96..445fe71 100644 --- a/server/services/thumbnail.ts +++ b/server/services/thumbnail.ts @@ -13,11 +13,11 @@ import fs from 'fs'; import path from 'path'; -import { exec } from 'child_process'; +import { execFile } from 'child_process'; import { promisify } from 'util'; import { minioService } from './minio'; -const execAsync = promisify(exec); +const execFileAsync = promisify(execFile); const TEMP_DIR = path.join(process.cwd(), 'uploads', 'temp'); @@ -35,8 +35,10 @@ export async function generateVideoThumbnail( seekTime: number = 1 ): Promise { try { - const cmd = `ffmpeg -y -ss ${seekTime} -i "${videoPath}" -vframes 1 -q:v 2 "${outputPath}"`; - await execAsync(cmd); + // Arguments séparés : aucun passage par le shell, les chemins ne sont jamais interprétés + await execFileAsync('ffmpeg', [ + '-y', '-ss', String(seekTime), '-i', videoPath, '-vframes', '1', '-q:v', '2', outputPath, + ]); return fs.existsSync(outputPath); } catch (error) { console.warn('⚠️ Failed to generate video thumbnail:', error); diff --git a/server/services/ttsSync.ts b/server/services/ttsSync.ts index 9bfd641..8113d28 100644 --- a/server/services/ttsSync.ts +++ b/server/services/ttsSync.ts @@ -15,11 +15,16 @@ export class TtsSyncService { * Calcule le word_duration optimal pour synchroniser l'affichage du texte * avec la durée réelle de la voix TTS générée. */ - async calculateSyncTiming(text: string, voice: string, ttsEngine?: string): Promise { + async calculateSyncTiming( + text: string, + voice: string, + ttsEngine?: string, + geminiApiKey?: string + ): Promise { const cleanText = this.cleanText(text); - // 1. Générer le TTS preview et mesurer sa durée exacte - const ttsResult = await ffmpegService.previewTTS(cleanText, voice, ttsEngine); + // 1. Générer le TTS preview (avec la même voix que le rendu) et mesurer sa durée exacte + const ttsResult = await ffmpegService.previewTTS(cleanText, voice, ttsEngine, geminiApiKey); if (!ttsResult.success || !ttsResult.audioBase64) { throw new Error('TTS preview failed: ' + (ttsResult.error || 'unknown')); } diff --git a/server/storage.ts b/server/storage.ts index bf874a2..d6cdcf2 100644 --- a/server/storage.ts +++ b/server/storage.ts @@ -153,8 +153,6 @@ export interface IStorage { updatePostGenerationStatus(id: string, status: string, progress: number, error?: string): Promise; deletePost(id: string): Promise; updatePostMedia(postId: string, mediaIds: string[]): Promise; - countProcessingReels(): Promise; - getNextPendingReel(): Promise; // Scheduled Posts getScheduledPosts(userId: string, startDate?: Date, endDate?: Date): Promise; @@ -459,24 +457,6 @@ export class DatabaseStorage implements IStorage { } } - async countProcessingReels(): Promise { - const processing = await db - .select() - .from(posts) - .where(eq(posts.generationStatus, "processing")); - return processing.length; - } - - async getNextPendingReel(): Promise { - const [post] = await db - .select() - .from(posts) - .where(eq(posts.generationStatus, "pending")) - .orderBy(asc(posts.createdAt)) - .limit(1); - return post || undefined; - } - // Scheduled Posts async getScheduledPosts(userId: string, startDate?: Date, endDate?: Date): Promise { return this.queryScheduledPosts(eq(posts.userId, userId), startDate, endDate); diff --git a/shared/reel.ts b/shared/reel.ts new file mode 100644 index 0000000..6796313 --- /dev/null +++ b/shared/reel.ts @@ -0,0 +1,65 @@ +/** + * Paramètres des rendus de Reels, validés à l'entrée des routes et stockés tels + * quels dans `reel_jobs.params` : un job repris après redémarrage retrouve ainsi + * exactement ce que l'utilisateur avait demandé. + */ + +import { z } from "zod"; + +const optionalText = z + .string() + .trim() + .optional() + .transform((value) => (value ? value : undefined)); + +/** Reel construit à partir d'une vidéo de la médiathèque. */ +export const videoReelParamsSchema = z.object({ + videoMediaId: z.string({ required_error: "Vidéo requise" }).min(1, "Vidéo requise"), + pageIds: z + .array(z.string().min(1), { required_error: "Au moins une page requise" }) + .min(1, "Au moins une page requise"), + musicTrackId: optionalText, + musicUrl: optionalText, + overlayText: optionalText, + description: optionalText, + ttsEnabled: z.boolean().default(false), + ttsVoice: optionalText, + ttsEngine: z.enum(["gemini", "edge"]).optional(), + scheduledFor: optionalText, + wordDuration: z.number().positive().max(5).default(0.6), + fontSize: z.number().int().min(16).max(200).default(64), + musicVolume: z.number().min(0).max(2).default(0.25), + drawText: z.boolean().default(true), + stabilize: z.boolean().default(false), + enableEndingEffect: z.boolean().default(true), + // Déterminé par le serveur à partir de la première page, jamais par le client + storeName: z.string().optional(), +}); + +export type VideoReelParams = z.infer; + +/** Reel construit à partir d'images (rendu Remotion). */ +export const imagesReelParamsSchema = z.object({ + // Chemins relatifs /uploads/... ou URL absolues + imageUrls: z.array(z.string().min(1)).min(1, "Aucune image fournie").max(4), + overlayText: optionalText, + musicUrl: optionalText, + musicVolume: z.number().min(0).max(2).default(0.3), + ttsEngine: z.enum(["gemini", "edge"]).optional(), + ttsVoice: optionalText, + storeName: z.string().optional(), + // Fichiers temporaires à supprimer une fois le rendu terminé + tempFiles: z.array(z.string()).default([]), +}); + +export type ImagesReelParams = z.infer; + +export type ReelJobKind = "video" | "images"; + +export type ReelJobStatus = "pending" | "processing" | "completed" | "failed"; + +/** Résultat d'un rendu d'images, consulté par le client pour l'aperçu. */ +export interface ImagesReelResult { + url: string; + thumbnailUrl: string | null; +} diff --git a/shared/schema.ts b/shared/schema.ts index ce17633..fc24006 100644 --- a/shared/schema.ts +++ b/shared/schema.ts @@ -281,6 +281,31 @@ export const userPagePermissions = pgTable("user_page_permissions", { createdAt: timestamp("created_at").defaultNow(), }); +/** + * File d'attente persistante des rendus de Reels. + * + * Un job survit à un redémarrage du serveur : le worker le réserve avec + * `FOR UPDATE SKIP LOCKED`, entretient `heartbeat_at` pendant le rendu, et les + * jobs dont le cœur ne bat plus sont remis en file au démarrage suivant. + */ +export const reelJobs = pgTable("reel_jobs", { + id: varchar("id").primaryKey().default(sql`gen_random_uuid()`), + userId: varchar("user_id").notNull().references(() => users.id, { onDelete: "cascade" }), + // Post suivi par l'interface (reel vidéo). Absent pour un rendu d'images non publié. + postId: varchar("post_id").references(() => posts.id, { onDelete: "cascade" }), + kind: text("kind").notNull(), // 'video' | 'images' + params: jsonb("params").notNull(), + status: text("status").notNull().default("pending"), // 'pending' | 'processing' | 'completed' | 'failed' + step: text("step"), + progress: integer("progress").notNull().default(0), + attempts: integer("attempts").notNull().default(0), + result: jsonb("result"), + error: text("error"), + heartbeatAt: timestamp("heartbeat_at"), + createdAt: timestamp("created_at").notNull().defaultNow(), + updatedAt: timestamp("updated_at").notNull().defaultNow(), +}); + export const socialPagesRelations = relations(socialPages, ({ one, many }) => ({ user: one(users, { fields: [socialPages.userId], @@ -558,4 +583,5 @@ export type TiktokConfig = typeof tiktokConfig.$inferSelect; export type InsertTiktokConfig = z.infer; export type UpdateTiktokConfig = z.infer; - +export type ReelJob = typeof reelJobs.$inferSelect; +export type InsertReelJob = typeof reelJobs.$inferInsert;