feat(reels): file d'attente persistante des rendus (reel_jobs)

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018Ze4bs7tpF1KGWUk6ZZSZ4
This commit is contained in:
Claude committed 2026-09-24 09:31:31 +00:00
1 parent fcbc6b706f
commit d24ab0d4b7
19 files changed
+1071 -901

No files matched your search

@@ -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)); }
+2
View File
@@ -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)); }
+7 -3
View File
@@ -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:
+7
View File
@@ -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);
});
});
})();
+21
View File
@@ -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`],
];
+35 -465
View File
@@ -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',
});
}
+81 -405
View File
@@ -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<string, RenderJob>();
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<string> {
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<string> {
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<string, Express.Multer.File[]>;
const uploaded = [...(fields["images"] ?? []), ...(fields["music"] ?? [])];
try {
const fields = req.files as Record<string, Express.Multer.File[]>;
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 <Img>)
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<void> {
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) });
}
});
+9
View File
@@ -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) {
+65
View File
@@ -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<string | undefined> {
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<string | undefined> {
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<string | undefined> {
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<string | undefined> {
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;
}
}
+234
View File
@@ -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<string> | null = null;
function getBundle(): Promise<string> {
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<string, string> = {
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<string> {
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<ImagesReelResult> {
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<typeof computeWordTimings> | 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 */ });
}
}
}
+14
View File
@@ -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<void> {
registerReelJobHandler("video", runVideoReelJob);
registerReelJobHandler("images", runImagesReelJob);
return startReelWorker();
}
+133
View File
@@ -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<Media> {
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<PagePublishResult[]> {
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}` : ""}`;
}
+277
View File
@@ -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<void>;
}
export type JobHandler = (ctx: JobContext) => Promise<unknown>;
const handlers = new Map<ReelJobKind, JobHandler>();
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<ReelJob> {
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<ReelJob | undefined> {
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<number> {
const [row] = await db
.select({ count: sql<number>`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<void> {
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<void> {
for (;;) {
const job = await claimNextJob();
if (!job) return;
await runJob(job);
}
}
async function claimNextJob(): Promise<ReelJob | undefined> {
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<void> {
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<void> {
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<void> {
const patch: Partial<typeof posts.$inferInsert> = {
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<void> {
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<void> {
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`);
}
}
+78
View File
@@ -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 };
}
+6 -4
View File
@@ -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<boolean> {
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);
+8 -3
View File
@@ -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<SyncTiming> {
async calculateSyncTiming(
text: string,
voice: string,
ttsEngine?: string,
geminiApiKey?: string
): Promise<SyncTiming> {
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'));
}
-20
View File
@@ -153,8 +153,6 @@ export interface IStorage {
updatePostGenerationStatus(id: string, status: string, progress: number, error?: string): Promise<Post>;
deletePost(id: string): Promise<void>;
updatePostMedia(postId: string, mediaIds: string[]): Promise<void>;
countProcessingReels(): Promise<number>;
getNextPendingReel(): Promise<Post | undefined>;
// Scheduled Posts
getScheduledPosts(userId: string, startDate?: Date, endDate?: Date): Promise<ScheduledPost[]>;
@@ -459,24 +457,6 @@ export class DatabaseStorage implements IStorage {
}
}
async countProcessingReels(): Promise<number> {
const processing = await db
.select()
.from(posts)
.where(eq(posts.generationStatus, "processing"));
return processing.length;
}
async getNextPendingReel(): Promise<Post | undefined> {
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<any[]> {
return this.queryScheduledPosts(eq(posts.userId, userId), startDate, endDate);
+65
View File
@@ -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<typeof videoReelParamsSchema>;
/** 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<typeof imagesReelParamsSchema>;
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;
}
+27 -1
View File
@@ -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<typeof insertTiktokConfigSchema>;
export type UpdateTiktokConfig = z.infer<typeof updateTiktokConfigSchema>;
export type ReelJob = typeof reelJobs.$inferSelect;
export type InsertReelJob = typeof reelJobs.$inferInsert;