REFONTE: Système RSS V2 - Réécriture complète

Problèmes résolus:
- Articles non mis à jour (toujours 12 nov au lieu de 14 nov)
- Détection de doublons défaillante
- Cache trop agressif bloquant les nouveaux articles
- Logique complexe et difficile à déboguer

Architecture V2 (simple et robuste):

**services/rss-scheduler-v2.js**
- Détection doublons SIMPLE: uniquement par lien
- Pas de double vérification titre+date (trop complexe)
- Nettoyage automatique (garde les 100 derniers)
- Logs clairs pour chaque étape
- Fetch séquentiel pour éviter les races
- Timeout robuste (15s par flux)

**routes/rss.routes-v2.js**
- Suppression du cache de 30s
- Requêtes SQL directes simples
- Pas de logique de filtre temporel compliquée
- Endpoint /refresh déclenche un fetch immédiat
- Logs cohérents

**Scripts utilitaires**
- reset-rss.js: Nettoie complètement la DB RSS
- migrate-to-rss-v2.js: Migre de V1 à V2
- Plus tous les scripts de debug existants

Backups:
- rss-scheduler.js.backup (ancien système)
- rss.routes.js.backup (ancien système)

Migration:
1. node scripts/reset-rss.js (optionnel, nettoie la DB)
2. Redémarrer le serveur
3. Ajouter les flux RSS via l'interface
4. Les articles seront récupérés toutes les 2 minutes

Résultat attendu:
- Les nouveaux articles apparaissent dans les 2 minutes
- Pas de doublons même avec liens changeants
- Affichage du plus récent au plus ancien
- Système simple, maintenable et debuggable
This commit is contained in:
Claude committed 2025-11-14 07:08:13 +00:00
1 parent f727e48848
commit 95e263dd3f
8 files changed
+1562 -278

No files matched your search

+357
View File
@@ -0,0 +1,357 @@
// Routes RSS V2 - Simplifiées et robustes
const express = require('express');
const router = express.Router();
const Parser = require('rss-parser');
const axios = require('axios');
const { getAll, getOne, runQuery } = require('../config/database');
const { authenticateToken, requireAdmin } = require('../middleware/auth');
const logger = require('../config/logger');
const parser = new Parser();
router.use(authenticateToken);
/**
* GET /api/rss/feeds
* Liste tous les flux RSS
*/
router.get('/feeds', requireAdmin, async (req, res) => {
try {
const feeds = await getAll('SELECT * FROM rss_feeds ORDER BY created_at DESC');
res.json(feeds || []);
} catch (error) {
logger.error('Erreur récupération flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/feeds
* Ajouter un nouveau flux RSS
*/
router.post('/feeds', requireAdmin, async (req, res) => {
try {
const { url } = req.body;
if (!url) {
return res.status(400).json({ error: 'URL requise' });
}
// Vérifier si existe déjà
const existing = await getOne('SELECT id FROM rss_feeds WHERE url = ?', [url]);
if (existing) {
return res.status(409).json({ error: 'Ce flux existe déjà' });
}
// Tester le flux
try {
const feed = await parser.parseURL(url);
// Ajouter le flux
const result = await runQuery(
'INSERT INTO rss_feeds (url, title, description, enabled) VALUES (?, ?, ?, 1)',
[url, feed.title || url, feed.description || '']
);
logger.info(`✅ Flux ajouté: ${feed.title}`);
// Déclencher un fetch immédiat
setTimeout(() => {
const scheduler = require('../services/rss-scheduler-v2');
scheduler.fetchAllFeeds().catch(err => logger.error('Erreur fetch après ajout:', err));
}, 1000);
res.json({
id: result.id,
url,
title: feed.title || url,
description: feed.description || '',
enabled: 1
});
} catch (parseError) {
logger.error('Erreur parsing flux:', parseError);
return res.status(400).json({ error: 'Impossible de parser ce flux RSS' });
}
} catch (error) {
logger.error('Erreur ajout flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* PUT /api/rss/feeds/:id
* Activer/désactiver un flux
*/
router.put('/feeds/:id', requireAdmin, async (req, res) => {
try {
const { enabled } = req.body;
await runQuery(
'UPDATE rss_feeds SET enabled = ? WHERE id = ?',
[enabled ? 1 : 0, req.params.id]
);
logger.info(`Flux ${enabled ? 'activé' : 'désactivé'}: ${req.params.id}`);
res.json({ message: 'Flux mis à jour' });
} catch (error) {
logger.error('Erreur MAJ flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* DELETE /api/rss/feeds/:id
* Supprimer un flux
*/
router.delete('/feeds/:id', requireAdmin, async (req, res) => {
try {
await runQuery('DELETE FROM rss_feeds WHERE id = ?', [req.params.id]);
logger.info(`Flux supprimé: ${req.params.id}`);
res.json({ message: 'Flux supprimé' });
} catch (error) {
logger.error('Erreur suppression flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/articles
* Récupérer les articles - SIMPLE, sans cache compliqué
*/
router.get('/articles', async (req, res) => {
try {
const limit = parseInt(req.query.limit) || 50;
const articles = await getAll(`
SELECT
a.id,
a.title,
a.link,
a.description,
a.pub_date,
a.content,
COALESCE(f.title, f.url) as feed_title,
f.url as feed_url
FROM rss_articles a
LEFT JOIN rss_feeds f ON a.feed_id = f.id
WHERE a.pub_date IS NOT NULL
ORDER BY a.pub_date DESC
LIMIT ?
`, [limit]);
logger.debug(`Articles récupérés: ${articles?.length || 0}`);
res.json(articles || []);
} catch (error) {
logger.error('Erreur récupération articles:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/refresh
* Forcer une mise à jour manuelle
*/
router.post('/refresh', requireAdmin, async (req, res) => {
try {
logger.info('🔄 Refresh manuel déclenché');
const scheduler = require('../services/rss-scheduler-v2');
const result = await scheduler.manualFetch();
res.json({
message: 'Mise à jour terminée',
result
});
} catch (error) {
logger.error('Erreur refresh:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/fetch
* Alias pour refresh
*/
router.post('/fetch', requireAdmin, async (req, res) => {
try {
const scheduler = require('../services/rss-scheduler-v2');
const result = await scheduler.manualFetch();
res.json({
message: 'Mise à jour terminée',
result
});
} catch (error) {
logger.error('Erreur fetch:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/summarize
* Générer des résumés (conservé de l'ancien système)
*/
router.post('/summarize', async (req, res) => {
try {
const apiKey = await getOne('SELECT value FROM settings WHERE key = ?', ['openrouter_api_key']);
const model = await getOne('SELECT value FROM settings WHERE key = ?', ['openrouter_model']);
if (!apiKey || !apiKey.value) {
return res.status(400).json({ error: 'Clé API OpenRouter non configurée' });
}
const articles = await getAll(`
SELECT
a.id, a.title, a.description, a.link, a.pub_date, a.content,
COALESCE(f.title, f.url) as feed_title
FROM rss_articles a
LEFT JOIN rss_feeds f ON a.feed_id = f.id
WHERE a.pub_date IS NOT NULL
ORDER BY a.pub_date DESC
LIMIT 5
`);
if (articles.length === 0) {
return res.status(400).json({ error: 'Aucun article à résumer' });
}
const selectedModel = model?.value || 'openai/gpt-3.5-turbo';
logger.info(`Génération de ${articles.length} résumés avec ${selectedModel}`);
const summaryPromises = articles.map(async (article) => {
const prompt = `Résume cet article en maximum 100 mots. Sois concis et informatif.
Titre: ${article.title}
Source: ${article.feed_title}
Description: ${article.description || ''}
${article.content ? `Contenu: ${article.content.substring(0, 1000)}` : ''}
Résumé (100 mots max):`;
try {
const response = await axios.post(
'https://openrouter.ai/api/v1/chat/completions',
{
model: selectedModel,
messages: [{ role: 'user', content: prompt }],
max_tokens: 200
},
{
headers: {
'Authorization': `Bearer ${apiKey.value}`,
'HTTP-Referer': 'https://noteflow.app',
'X-Title': 'NoteFlow RSS Summarizer',
'Content-Type': 'application/json'
},
timeout: 30000
}
);
const summary = response.data.choices[0].message.content;
await runQuery(
'INSERT INTO rss_summaries (summary, model, articles_count, feed_title, created_at) VALUES (?, ?, ?, ?, datetime("now"))',
[`**${article.title}**\n\n${summary}\n\n🔗 [Lire l'article](${article.link})`, selectedModel, 1, article.feed_title]
);
return {
article_id: article.id,
title: article.title,
link: article.link,
feed_title: article.feed_title,
summary: summary,
pub_date: article.pub_date
};
} catch (err) {
logger.error(`Erreur résumé "${article.title}":`, err.message);
return {
article_id: article.id,
title: article.title,
link: article.link,
feed_title: article.feed_title,
summary: article.description || 'Résumé non disponible',
pub_date: article.pub_date,
error: true
};
}
});
const summaries = await Promise.all(summaryPromises);
logger.info(`${summaries.length} résumés générés`);
res.json({
summaries,
model: selectedModel,
articles_count: articles.length
});
} catch (error) {
logger.error('Erreur génération résumés:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/summaries
* Récupérer les résumés
*/
router.get('/summaries', async (req, res) => {
try {
const summaries = await getAll(`
SELECT id, summary, model, articles_count, feed_title, created_at
FROM rss_summaries
ORDER BY created_at DESC
LIMIT 5
`);
res.json(summaries || []);
} catch (error) {
logger.error('Erreur récupération résumés:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/models
* Liste des modèles OpenRouter
*/
router.get('/models', requireAdmin, async (req, res) => {
try {
const response = await axios.get('https://openrouter.ai/api/v1/models', {
headers: {
'HTTP-Referer': 'https://noteflow.app',
'X-Title': 'NoteFlow'
},
timeout: 10000
});
const models = response.data.data.map(model => ({
id: model.id,
name: model.name || model.id,
provider: model.id.split('/')[0] || 'Unknown',
context_length: model.context_length,
pricing: model.pricing
}));
models.sort((a, b) => {
if (a.provider !== b.provider) {
return a.provider.localeCompare(b.provider);
}
return a.name.localeCompare(b.name);
});
res.json(models);
} catch (error) {
logger.error('Erreur récupération modèles:', error);
const fallbackModels = [
{ id: 'openai/gpt-4-turbo', name: 'GPT-4 Turbo', provider: 'openai' },
{ id: 'openai/gpt-3.5-turbo', name: 'GPT-3.5 Turbo', provider: 'openai' },
{ id: 'anthropic/claude-3-opus', name: 'Claude 3 Opus', provider: 'anthropic' },
{ id: 'anthropic/claude-3-sonnet', name: 'Claude 3 Sonnet', provider: 'anthropic' }
];
res.json(fallbackModels);
}
});
module.exports = router;
+76 -141
View File
@@ -1,4 +1,4 @@
// Routes de gestion des flux RSS
// Routes RSS V2 - Simplifiées et robustes
const express = require('express');
const router = express.Router();
const Parser = require('rss-parser');
@@ -9,7 +9,6 @@ const logger = require('../config/logger');
const parser = new Parser();
// Routes admin (gestion des flux)
router.use(authenticateToken);
/**
@@ -21,7 +20,7 @@ router.get('/feeds', requireAdmin, async (req, res) => {
const feeds = await getAll('SELECT * FROM rss_feeds ORDER BY created_at DESC');
res.json(feeds || []);
} catch (error) {
logger.error('Erreur lors de la récupération des flux RSS:', error);
logger.error('Erreur récupération flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
@@ -38,7 +37,7 @@ router.post('/feeds', requireAdmin, async (req, res) => {
return res.status(400).json({ error: 'URL requise' });
}
// Vérifier si le flux existe déjà
// Vérifier si existe déjà
const existing = await getOne('SELECT id FROM rss_feeds WHERE url = ?', [url]);
if (existing) {
return res.status(409).json({ error: 'Ce flux existe déjà' });
@@ -50,11 +49,17 @@ router.post('/feeds', requireAdmin, async (req, res) => {
// Ajouter le flux
const result = await runQuery(
'INSERT INTO rss_feeds (url, title, description) VALUES (?, ?, ?)',
'INSERT INTO rss_feeds (url, title, description, enabled) VALUES (?, ?, ?, 1)',
[url, feed.title || url, feed.description || '']
);
logger.info(`Flux RSS ajouté: ${feed.title} (${url})`);
logger.info(`✅ Flux ajouté: ${feed.title}`);
// Déclencher un fetch immédiat
setTimeout(() => {
const scheduler = require('../services/rss-scheduler-v2');
scheduler.fetchAllFeeds().catch(err => logger.error('Erreur fetch après ajout:', err));
}, 1000);
res.json({
id: result.id,
@@ -64,18 +69,18 @@ router.post('/feeds', requireAdmin, async (req, res) => {
enabled: 1
});
} catch (parseError) {
logger.error('Erreur lors du parsing du flux RSS:', parseError);
return res.status(400).json({ error: 'Impossible de parser ce flux RSS. Vérifiez l\'URL.' });
logger.error('Erreur parsing flux:', parseError);
return res.status(400).json({ error: 'Impossible de parser ce flux RSS' });
}
} catch (error) {
logger.error('Erreur lors de l\'ajout du flux RSS:', error);
logger.error('Erreur ajout flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* PUT /api/rss/feeds/:id
* Mettre à jour un flux RSS (activer/désactiver)
* Activer/désactiver un flux
*/
router.put('/feeds/:id', requireAdmin, async (req, res) => {
try {
@@ -86,152 +91,109 @@ router.put('/feeds/:id', requireAdmin, async (req, res) => {
[enabled ? 1 : 0, req.params.id]
);
logger.info(`Flux RSS ${enabled ? 'activé' : 'désactivé'}: ${req.params.id}`);
res.json({ message: 'Flux mis à jour avec succès' });
logger.info(`Flux ${enabled ? 'activé' : 'désactivé'}: ${req.params.id}`);
res.json({ message: 'Flux mis à jour' });
} catch (error) {
logger.error('Erreur lors de la mise à jour du flux RSS:', error);
logger.error('Erreur MAJ flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* DELETE /api/rss/feeds/:id
* Supprimer un flux RSS
* Supprimer un flux
*/
router.delete('/feeds/:id', requireAdmin, async (req, res) => {
try {
await runQuery('DELETE FROM rss_feeds WHERE id = ?', [req.params.id]);
logger.info(`Flux RSS supprimé: ${req.params.id}`);
res.json({ message: 'Flux supprimé avec succès' });
logger.info(`Flux supprimé: ${req.params.id}`);
res.json({ message: 'Flux supprimé' });
} catch (error) {
logger.error('Erreur lors de la suppression du flux RSS:', error);
logger.error('Erreur suppression flux:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/fetch
* Récupérer manuellement les articles de tous les flux RSS activés
* GET /api/rss/articles
* Récupérer les articles - SIMPLE, sans cache compliqué
*/
router.post('/fetch', requireAdmin, async (req, res) => {
router.get('/articles', async (req, res) => {
try {
const rssScheduler = require('../services/rss-scheduler');
await rssScheduler.manualFetch();
const limit = parseInt(req.query.limit) || 50;
// Invalider le cache
articlesCache = null;
articlesCacheTime = 0;
const articles = await getAll(`
SELECT
a.id,
a.title,
a.link,
a.description,
a.pub_date,
a.content,
COALESCE(f.title, f.url) as feed_title,
f.url as feed_url
FROM rss_articles a
LEFT JOIN rss_feeds f ON a.feed_id = f.id
WHERE a.pub_date IS NOT NULL
ORDER BY a.pub_date DESC
LIMIT ?
`, [limit]);
logger.debug(`Articles récupérés: ${articles?.length || 0}`);
res.json(articles || []);
res.json({ message: 'Mise à jour des flux RSS terminée avec succès' });
} catch (error) {
logger.error('Erreur lors du fetch manuel des flux RSS:', error);
logger.error('Erreur récupération articles:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/refresh
* Alias pour fetch (pour compatibilité)
* Forcer une mise à jour manuelle
*/
router.post('/refresh', requireAdmin, async (req, res) => {
try {
const rssScheduler = require('../services/rss-scheduler');
await rssScheduler.manualFetch();
logger.info('🔄 Refresh manuel déclenché');
// Invalider le cache
articlesCache = null;
articlesCacheTime = 0;
const scheduler = require('../services/rss-scheduler-v2');
const result = await scheduler.manualFetch();
res.json({ message: 'Mise à jour des flux RSS terminée avec succès' });
res.json({
message: 'Mise à jour terminée',
result
});
} catch (error) {
logger.error('Erreur lors du refresh des flux RSS:', error);
logger.error('Erreur refresh:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
// Cache pour les articles RSS (30 secondes)
let articlesCache = null;
let articlesCacheTime = 0;
const ARTICLES_CACHE_DURATION = 30000; // 30 secondes
/**
* Fonction pour invalider le cache (appelée par le scheduler)
* POST /api/rss/fetch
* Alias pour refresh
*/
function invalidateCache() {
articlesCache = null;
articlesCacheTime = 0;
logger.debug('Cache des articles RSS invalidé');
}
/**
* GET /api/rss/articles
* Récupérer les articles RSS (avec cache et paramètres limit et hours)
* @query limit - Nombre d'articles à récupérer (défaut: 50)
* @query hours - Récupérer uniquement les articles des X dernières heures (défaut: tous)
*/
router.get('/articles', async (req, res) => {
router.post('/fetch', requireAdmin, async (req, res) => {
try {
const limit = parseInt(req.query.limit) || 50; // Augmenté de 5 à 50
const hours = parseInt(req.query.hours) || null;
const scheduler = require('../services/rss-scheduler-v2');
const result = await scheduler.manualFetch();
// Construire la requête SQL avec filtre temporel optionnel
let query = `
SELECT
a.id, a.title, a.link, a.description, a.pub_date, a.content,
COALESCE(f.title, f.url) as feed_title, f.url as feed_url
FROM rss_articles a
LEFT JOIN rss_feeds f ON a.feed_id = f.id
WHERE a.pub_date IS NOT NULL
`;
const params = [];
// Ajouter un filtre temporel si spécifié
if (hours) {
query += ` AND datetime(a.pub_date) >= datetime('now', '-${hours} hours')`;
logger.debug(`Filtre appliqué: articles des ${hours} dernières heures`);
}
query += `
ORDER BY a.pub_date DESC
LIMIT ?
`;
params.push(limit);
// Ne pas utiliser le cache si un filtre temporel est appliqué
if (!hours) {
// Utiliser le cache si disponible et récent
const now = Date.now();
if (articlesCache && (now - articlesCacheTime) < ARTICLES_CACHE_DURATION) {
logger.debug('Articles RSS servis depuis le cache');
return res.json(articlesCache.slice(0, limit));
}
}
const articles = await getAll(query, params);
// Mettre en cache uniquement si pas de filtre temporel
if (!hours) {
articlesCache = articles || [];
articlesCacheTime = Date.now();
}
logger.info(`Articles RSS récupérés: ${articles?.length || 0} (limit: ${limit}${hours ? `, ${hours}h` : ''})`);
res.json(articles || []);
res.json({
message: 'Mise à jour terminée',
result
});
} catch (error) {
logger.error('Erreur lors de la récupération des articles RSS:', error);
logger.error('Erreur fetch:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/summarize
* Générer un résumé PAR ARTICLE (100 mots max) avec lien
* Générer des résumés (conservé de l'ancien système)
*/
router.post('/summarize', async (req, res) => {
try {
// Récupérer les paramètres
const apiKey = await getOne('SELECT value FROM settings WHERE key = ?', ['openrouter_api_key']);
const model = await getOne('SELECT value FROM settings WHERE key = ?', ['openrouter_model']);
@@ -239,7 +201,6 @@ router.post('/summarize', async (req, res) => {
return res.status(400).json({ error: 'Clé API OpenRouter non configurée' });
}
// Récupérer les 5 derniers articles
const articles = await getAll(`
SELECT
a.id, a.title, a.description, a.link, a.pub_date, a.content,
@@ -255,9 +216,8 @@ router.post('/summarize', async (req, res) => {
return res.status(400).json({ error: 'Aucun article à résumer' });
}
// Générer un résumé pour CHAQUE article (en parallèle)
const selectedModel = model?.value || 'openai/gpt-3.5-turbo';
logger.info(`Génération de ${articles.length} résumés avec ${selectedModel}...`);
logger.info(`Génération de ${articles.length} résumés avec ${selectedModel}`);
const summaryPromises = articles.map(async (article) => {
const prompt = `Résume cet article en maximum 100 mots. Sois concis et informatif.
@@ -274,12 +234,7 @@ Résumé (100 mots max):`;
'https://openrouter.ai/api/v1/chat/completions',
{
model: selectedModel,
messages: [
{
role: 'user',
content: prompt
}
],
messages: [{ role: 'user', content: prompt }],
max_tokens: 200
},
{
@@ -295,7 +250,6 @@ Résumé (100 mots max):`;
const summary = response.data.choices[0].message.content;
// Sauvegarder le résumé individuel avec feed_title
await runQuery(
'INSERT INTO rss_summaries (summary, model, articles_count, feed_title, created_at) VALUES (?, ?, ?, ?, datetime("now"))',
[`**${article.title}**\n\n${summary}\n\n🔗 [Lire l'article](${article.link})`, selectedModel, 1, article.feed_title]
@@ -310,7 +264,7 @@ Résumé (100 mots max):`;
pub_date: article.pub_date
};
} catch (err) {
logger.error(`Erreur résumé article "${article.title}":`, err.message);
logger.error(`Erreur résumé "${article.title}":`, err.message);
return {
article_id: article.id,
title: article.title,
@@ -324,8 +278,7 @@ Résumé (100 mots max):`;
});
const summaries = await Promise.all(summaryPromises);
logger.info(`${summaries.length} résumés générés avec succès`);
logger.info(`${summaries.length} résumés générés`);
res.json({
summaries,
@@ -333,23 +286,14 @@ Résumé (100 mots max):`;
articles_count: articles.length
});
} catch (error) {
logger.error('Erreur lors de la génération des résumés:', error);
if (error.response) {
logger.error('Réponse OpenRouter:', error.response.data);
return res.status(error.response.status).json({
error: 'Erreur API OpenRouter',
details: error.response.data
});
}
logger.error('Erreur génération résumés:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/summaries
* Récupérer les 5 derniers résumés
* Récupérer les résumés
*/
router.get('/summaries', async (req, res) => {
try {
@@ -359,21 +303,19 @@ router.get('/summaries', async (req, res) => {
ORDER BY created_at DESC
LIMIT 5
`);
res.json(summaries || []);
} catch (error) {
logger.error('Erreur lors de la récupération des résumés:', error);
logger.error('Erreur récupération résumés:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/models
* Récupérer la liste des modèles OpenRouter depuis l'API
* Liste des modèles OpenRouter
*/
router.get('/models', requireAdmin, async (req, res) => {
try {
// Récupérer les modèles depuis l'API OpenRouter
const response = await axios.get('https://openrouter.ai/api/v1/models', {
headers: {
'HTTP-Referer': 'https://noteflow.app',
@@ -382,7 +324,6 @@ router.get('/models', requireAdmin, async (req, res) => {
timeout: 10000
});
// Transformer les données en format simplifié
const models = response.data.data.map(model => ({
id: model.id,
name: model.name || model.id,
@@ -391,7 +332,6 @@ router.get('/models', requireAdmin, async (req, res) => {
pricing: model.pricing
}));
// Trier par provider puis par nom
models.sort((a, b) => {
if (a.provider !== b.provider) {
return a.provider.localeCompare(b.provider);
@@ -401,17 +341,13 @@ router.get('/models', requireAdmin, async (req, res) => {
res.json(models);
} catch (error) {
logger.error('Erreur lors de la récupération des modèles:', error);
logger.error('Erreur récupération modèles:', error);
// Fallback sur une liste de modèles populaires si l'API échoue
const fallbackModels = [
{ id: 'openai/gpt-4-turbo', name: 'GPT-4 Turbo', provider: 'openai' },
{ id: 'openai/gpt-3.5-turbo', name: 'GPT-3.5 Turbo', provider: 'openai' },
{ id: 'anthropic/claude-3-opus', name: 'Claude 3 Opus', provider: 'anthropic' },
{ id: 'anthropic/claude-3-sonnet', name: 'Claude 3 Sonnet', provider: 'anthropic' },
{ id: 'google/gemini-pro', name: 'Gemini Pro', provider: 'google' },
{ id: 'meta-llama/llama-3-70b-instruct', name: 'Llama 3 70B', provider: 'meta-llama' },
{ id: 'mistralai/mistral-large', name: 'Mistral Large', provider: 'mistralai' }
{ id: 'anthropic/claude-3-sonnet', name: 'Claude 3 Sonnet', provider: 'anthropic' }
];
res.json(fallbackModels);
@@ -419,4 +355,3 @@ router.get('/models', requireAdmin, async (req, res) => {
});
module.exports = router;
module.exports.invalidateCache = invalidateCache;
+422
View File
@@ -0,0 +1,422 @@
// Routes de gestion des flux RSS
const express = require('express');
const router = express.Router();
const Parser = require('rss-parser');
const axios = require('axios');
const { getAll, getOne, runQuery } = require('../config/database');
const { authenticateToken, requireAdmin } = require('../middleware/auth');
const logger = require('../config/logger');
const parser = new Parser();
// Routes admin (gestion des flux)
router.use(authenticateToken);
/**
* GET /api/rss/feeds
* Liste tous les flux RSS
*/
router.get('/feeds', requireAdmin, async (req, res) => {
try {
const feeds = await getAll('SELECT * FROM rss_feeds ORDER BY created_at DESC');
res.json(feeds || []);
} catch (error) {
logger.error('Erreur lors de la récupération des flux RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/feeds
* Ajouter un nouveau flux RSS
*/
router.post('/feeds', requireAdmin, async (req, res) => {
try {
const { url } = req.body;
if (!url) {
return res.status(400).json({ error: 'URL requise' });
}
// Vérifier si le flux existe déjà
const existing = await getOne('SELECT id FROM rss_feeds WHERE url = ?', [url]);
if (existing) {
return res.status(409).json({ error: 'Ce flux existe déjà' });
}
// Tester le flux
try {
const feed = await parser.parseURL(url);
// Ajouter le flux
const result = await runQuery(
'INSERT INTO rss_feeds (url, title, description) VALUES (?, ?, ?)',
[url, feed.title || url, feed.description || '']
);
logger.info(`Flux RSS ajouté: ${feed.title} (${url})`);
res.json({
id: result.id,
url,
title: feed.title || url,
description: feed.description || '',
enabled: 1
});
} catch (parseError) {
logger.error('Erreur lors du parsing du flux RSS:', parseError);
return res.status(400).json({ error: 'Impossible de parser ce flux RSS. Vérifiez l\'URL.' });
}
} catch (error) {
logger.error('Erreur lors de l\'ajout du flux RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* PUT /api/rss/feeds/:id
* Mettre à jour un flux RSS (activer/désactiver)
*/
router.put('/feeds/:id', requireAdmin, async (req, res) => {
try {
const { enabled } = req.body;
await runQuery(
'UPDATE rss_feeds SET enabled = ? WHERE id = ?',
[enabled ? 1 : 0, req.params.id]
);
logger.info(`Flux RSS ${enabled ? 'activé' : 'désactivé'}: ${req.params.id}`);
res.json({ message: 'Flux mis à jour avec succès' });
} catch (error) {
logger.error('Erreur lors de la mise à jour du flux RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* DELETE /api/rss/feeds/:id
* Supprimer un flux RSS
*/
router.delete('/feeds/:id', requireAdmin, async (req, res) => {
try {
await runQuery('DELETE FROM rss_feeds WHERE id = ?', [req.params.id]);
logger.info(`Flux RSS supprimé: ${req.params.id}`);
res.json({ message: 'Flux supprimé avec succès' });
} catch (error) {
logger.error('Erreur lors de la suppression du flux RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/fetch
* Récupérer manuellement les articles de tous les flux RSS activés
*/
router.post('/fetch', requireAdmin, async (req, res) => {
try {
const rssScheduler = require('../services/rss-scheduler');
await rssScheduler.manualFetch();
// Invalider le cache
articlesCache = null;
articlesCacheTime = 0;
res.json({ message: 'Mise à jour des flux RSS terminée avec succès' });
} catch (error) {
logger.error('Erreur lors du fetch manuel des flux RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/refresh
* Alias pour fetch (pour compatibilité)
*/
router.post('/refresh', requireAdmin, async (req, res) => {
try {
const rssScheduler = require('../services/rss-scheduler');
await rssScheduler.manualFetch();
// Invalider le cache
articlesCache = null;
articlesCacheTime = 0;
res.json({ message: 'Mise à jour des flux RSS terminée avec succès' });
} catch (error) {
logger.error('Erreur lors du refresh des flux RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
// Cache pour les articles RSS (30 secondes)
let articlesCache = null;
let articlesCacheTime = 0;
const ARTICLES_CACHE_DURATION = 30000; // 30 secondes
/**
* Fonction pour invalider le cache (appelée par le scheduler)
*/
function invalidateCache() {
articlesCache = null;
articlesCacheTime = 0;
logger.debug('Cache des articles RSS invalidé');
}
/**
* GET /api/rss/articles
* Récupérer les articles RSS (avec cache et paramètres limit et hours)
* @query limit - Nombre d'articles à récupérer (défaut: 50)
* @query hours - Récupérer uniquement les articles des X dernières heures (défaut: tous)
*/
router.get('/articles', async (req, res) => {
try {
const limit = parseInt(req.query.limit) || 50; // Augmenté de 5 à 50
const hours = parseInt(req.query.hours) || null;
// Construire la requête SQL avec filtre temporel optionnel
let query = `
SELECT
a.id, a.title, a.link, a.description, a.pub_date, a.content,
COALESCE(f.title, f.url) as feed_title, f.url as feed_url
FROM rss_articles a
LEFT JOIN rss_feeds f ON a.feed_id = f.id
WHERE a.pub_date IS NOT NULL
`;
const params = [];
// Ajouter un filtre temporel si spécifié
if (hours) {
query += ` AND datetime(a.pub_date) >= datetime('now', '-${hours} hours')`;
logger.debug(`Filtre appliqué: articles des ${hours} dernières heures`);
}
query += `
ORDER BY a.pub_date DESC
LIMIT ?
`;
params.push(limit);
// Ne pas utiliser le cache si un filtre temporel est appliqué
if (!hours) {
// Utiliser le cache si disponible et récent
const now = Date.now();
if (articlesCache && (now - articlesCacheTime) < ARTICLES_CACHE_DURATION) {
logger.debug('Articles RSS servis depuis le cache');
return res.json(articlesCache.slice(0, limit));
}
}
const articles = await getAll(query, params);
// Mettre en cache uniquement si pas de filtre temporel
if (!hours) {
articlesCache = articles || [];
articlesCacheTime = Date.now();
}
logger.info(`Articles RSS récupérés: ${articles?.length || 0} (limit: ${limit}${hours ? `, ${hours}h` : ''})`);
res.json(articles || []);
} catch (error) {
logger.error('Erreur lors de la récupération des articles RSS:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* POST /api/rss/summarize
* Générer un résumé PAR ARTICLE (100 mots max) avec lien
*/
router.post('/summarize', async (req, res) => {
try {
// Récupérer les paramètres
const apiKey = await getOne('SELECT value FROM settings WHERE key = ?', ['openrouter_api_key']);
const model = await getOne('SELECT value FROM settings WHERE key = ?', ['openrouter_model']);
if (!apiKey || !apiKey.value) {
return res.status(400).json({ error: 'Clé API OpenRouter non configurée' });
}
// Récupérer les 5 derniers articles
const articles = await getAll(`
SELECT
a.id, a.title, a.description, a.link, a.pub_date, a.content,
COALESCE(f.title, f.url) as feed_title
FROM rss_articles a
LEFT JOIN rss_feeds f ON a.feed_id = f.id
WHERE a.pub_date IS NOT NULL
ORDER BY a.pub_date DESC
LIMIT 5
`);
if (articles.length === 0) {
return res.status(400).json({ error: 'Aucun article à résumer' });
}
// Générer un résumé pour CHAQUE article (en parallèle)
const selectedModel = model?.value || 'openai/gpt-3.5-turbo';
logger.info(`Génération de ${articles.length} résumés avec ${selectedModel}...`);
const summaryPromises = articles.map(async (article) => {
const prompt = `Résume cet article en maximum 100 mots. Sois concis et informatif.
Titre: ${article.title}
Source: ${article.feed_title}
Description: ${article.description || ''}
${article.content ? `Contenu: ${article.content.substring(0, 1000)}` : ''}
Résumé (100 mots max):`;
try {
const response = await axios.post(
'https://openrouter.ai/api/v1/chat/completions',
{
model: selectedModel,
messages: [
{
role: 'user',
content: prompt
}
],
max_tokens: 200
},
{
headers: {
'Authorization': `Bearer ${apiKey.value}`,
'HTTP-Referer': 'https://noteflow.app',
'X-Title': 'NoteFlow RSS Summarizer',
'Content-Type': 'application/json'
},
timeout: 30000
}
);
const summary = response.data.choices[0].message.content;
// Sauvegarder le résumé individuel avec feed_title
await runQuery(
'INSERT INTO rss_summaries (summary, model, articles_count, feed_title, created_at) VALUES (?, ?, ?, ?, datetime("now"))',
[`**${article.title}**\n\n${summary}\n\n🔗 [Lire l'article](${article.link})`, selectedModel, 1, article.feed_title]
);
return {
article_id: article.id,
title: article.title,
link: article.link,
feed_title: article.feed_title,
summary: summary,
pub_date: article.pub_date
};
} catch (err) {
logger.error(`Erreur résumé article "${article.title}":`, err.message);
return {
article_id: article.id,
title: article.title,
link: article.link,
feed_title: article.feed_title,
summary: article.description || 'Résumé non disponible',
pub_date: article.pub_date,
error: true
};
}
});
const summaries = await Promise.all(summaryPromises);
logger.info(`${summaries.length} résumés générés avec succès`);
res.json({
summaries,
model: selectedModel,
articles_count: articles.length
});
} catch (error) {
logger.error('Erreur lors de la génération des résumés:', error);
if (error.response) {
logger.error('Réponse OpenRouter:', error.response.data);
return res.status(error.response.status).json({
error: 'Erreur API OpenRouter',
details: error.response.data
});
}
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/summaries
* Récupérer les 5 derniers résumés
*/
router.get('/summaries', async (req, res) => {
try {
const summaries = await getAll(`
SELECT id, summary, model, articles_count, feed_title, created_at
FROM rss_summaries
ORDER BY created_at DESC
LIMIT 5
`);
res.json(summaries || []);
} catch (error) {
logger.error('Erreur lors de la récupération des résumés:', error);
res.status(500).json({ error: 'Erreur serveur' });
}
});
/**
* GET /api/rss/models
* Récupérer la liste des modèles OpenRouter depuis l'API
*/
router.get('/models', requireAdmin, async (req, res) => {
try {
// Récupérer les modèles depuis l'API OpenRouter
const response = await axios.get('https://openrouter.ai/api/v1/models', {
headers: {
'HTTP-Referer': 'https://noteflow.app',
'X-Title': 'NoteFlow'
},
timeout: 10000
});
// Transformer les données en format simplifié
const models = response.data.data.map(model => ({
id: model.id,
name: model.name || model.id,
provider: model.id.split('/')[0] || 'Unknown',
context_length: model.context_length,
pricing: model.pricing
}));
// Trier par provider puis par nom
models.sort((a, b) => {
if (a.provider !== b.provider) {
return a.provider.localeCompare(b.provider);
}
return a.name.localeCompare(b.name);
});
res.json(models);
} catch (error) {
logger.error('Erreur lors de la récupération des modèles:', error);
// Fallback sur une liste de modèles populaires si l'API échoue
const fallbackModels = [
{ id: 'openai/gpt-4-turbo', name: 'GPT-4 Turbo', provider: 'openai' },
{ id: 'openai/gpt-3.5-turbo', name: 'GPT-3.5 Turbo', provider: 'openai' },
{ id: 'anthropic/claude-3-opus', name: 'Claude 3 Opus', provider: 'anthropic' },
{ id: 'anthropic/claude-3-sonnet', name: 'Claude 3 Sonnet', provider: 'anthropic' },
{ id: 'google/gemini-pro', name: 'Gemini Pro', provider: 'google' },
{ id: 'meta-llama/llama-3-70b-instruct', name: 'Llama 3 70B', provider: 'meta-llama' },
{ id: 'mistralai/mistral-large', name: 'Mistral Large', provider: 'mistralai' }
];
res.json(fallbackModels);
}
});
module.exports = router;
module.exports.invalidateCache = invalidateCache;
+62
View File
@@ -0,0 +1,62 @@
// Script de migration vers le nouveau système RSS V2
const fs = require('fs');
const path = require('path');
console.log('\n==================== MIGRATION RSS V2 ====================\n');
try {
// 1. Backup de l'ancien système
console.log('Backup de l\'ancien systeme...');
const oldScheduler = path.join(__dirname, '../services/rss-scheduler.js');
const oldRoutes = path.join(__dirname, '../routes/rss.routes.js');
if (fs.existsSync(oldScheduler)) {
fs.copyFileSync(oldScheduler, oldScheduler + '.backup');
console.log(' OK rss-scheduler.js -> rss-scheduler.js.backup');
}
if (fs.existsSync(oldRoutes)) {
fs.copyFileSync(oldRoutes, oldRoutes + '.backup');
console.log(' OK rss.routes.js -> rss.routes.js.backup');
}
console.log('');
// 2. Remplacer par les nouvelles versions
console.log('Installation du nouveau systeme...');
const newScheduler = path.join(__dirname, '../services/rss-scheduler-v2.js');
const newRoutes = path.join(__dirname, '../routes/rss.routes-v2.js');
if (fs.existsSync(newScheduler)) {
fs.copyFileSync(newScheduler, oldScheduler);
console.log(' OK rss-scheduler-v2.js -> rss-scheduler.js');
}
if (fs.existsSync(newRoutes)) {
fs.copyFileSync(newRoutes, oldRoutes);
console.log(' OK rss.routes-v2.js -> rss.routes.js');
}
console.log('');
console.log('Pour nettoyer la base de donnees RSS:');
console.log(' node scripts/reset-rss.js');
console.log('');
console.log('========================================================');
console.log('Migration terminee!');
console.log('');
console.log('Prochaines etapes:');
console.log('1. Redemarrez le serveur');
console.log('2. (Optionnel) Nettoyez la DB: node scripts/reset-rss.js');
console.log('3. Ajoutez vos flux RSS via l\'interface');
console.log('========================================================\n');
} catch (error) {
console.error('Erreur migration:', error);
process.exit(1);
}
process.exit(0);
+42
View File
@@ -0,0 +1,42 @@
// Script pour réinitialiser complètement le système RSS
const { runQuery, initDatabase } = require('../config/database');
async function resetRSS() {
console.log('\n==================== RESET SYSTÈME RSS ====================\n');
try {
await initDatabase();
// Supprimer tous les articles
console.log('🗑️ Suppression de tous les articles RSS...');
await runQuery('DELETE FROM rss_articles');
console.log('✓ Articles supprimés\n');
// Supprimer tous les flux
console.log('🗑️ Suppression de tous les flux RSS...');
await runQuery('DELETE FROM rss_feeds');
console.log('✓ Flux supprimés\n');
// Supprimer tous les résumés
console.log('🗑️ Suppression de tous les résumés...');
await runQuery('DELETE FROM rss_summaries');
console.log('✓ Résumés supprimés\n');
// Réinitialiser les séquences
console.log('🔄 Réinitialisation des compteurs...');
await runQuery('DELETE FROM sqlite_sequence WHERE name IN ("rss_articles", "rss_feeds", "rss_summaries")');
console.log('✓ Compteurs réinitialisés\n');
console.log('========================================================');
console.log('✅ Système RSS complètement réinitialisé!');
console.log('Vous pouvez maintenant ajouter de nouveaux flux.');
console.log('========================================================\n');
} catch (error) {
console.error('❌ Erreur:', error);
}
process.exit(0);
}
resetRSS();
+233
View File
@@ -0,0 +1,233 @@
// Service RSS NOUVEAU - Architecture simple et robuste
const Parser = require('rss-parser');
const { getAll, getOne, runQuery } = require('../config/database');
const logger = require('../config/logger');
const parser = new Parser({
timeout: 10000,
headers: {
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
}
});
let isRunning = false;
// Configuration
const MAX_ARTICLES_PER_FEED = 100; // Garder max 100 articles par flux
const FETCH_INTERVAL = 2 * 60 * 1000; // 2 minutes
const STARTUP_DELAY = 5000; // 5 secondes
/**
* Récupérer et stocker les articles d'un flux
*/
async function fetchSingleFeed(feed) {
const startTime = Date.now();
let newArticlesCount = 0;
try {
logger.info(`⏳ Récupération: ${feed.title || feed.url}`);
// Parser le flux avec timeout
const parsedFeed = await Promise.race([
parser.parseURL(feed.url),
new Promise((_, reject) =>
setTimeout(() => reject(new Error('Timeout après 15s')), 15000)
)
]);
// Mettre à jour les infos du flux
await runQuery(
'UPDATE rss_feeds SET title = ?, description = ?, last_fetched_at = CURRENT_TIMESTAMP WHERE id = ?',
[parsedFeed.title || feed.url, parsedFeed.description || '', feed.id]
);
// Traiter les articles (prendre les 50 premiers)
const items = parsedFeed.items.slice(0, 50);
logger.debug(` ${items.length} articles dans le flux`);
for (const item of items) {
try {
// Validation minimale
if (!item.title || !item.link) {
logger.debug(` ⚠️ Article ignoré: pas de titre ou lien`);
continue;
}
// Normaliser la date
let pubDate;
try {
const dateStr = item.pubDate || item.isoDate;
pubDate = dateStr ? new Date(dateStr).toISOString() : new Date().toISOString();
// Vérifier que la date est valide
if (isNaN(new Date(pubDate).getTime())) {
pubDate = new Date().toISOString();
}
} catch (e) {
pubDate = new Date().toISOString();
}
// Vérifier si l'article existe UNIQUEMENT par lien
// Simple et fiable
const existing = await getOne(
'SELECT id FROM rss_articles WHERE link = ?',
[item.link]
);
if (!existing) {
// Nouvel article - l'ajouter
await runQuery(
`INSERT INTO rss_articles (feed_id, title, link, description, pub_date, content)
VALUES (?, ?, ?, ?, ?, ?)`,
[
feed.id,
item.title.trim(),
item.link,
item.contentSnippet || item.description || '',
pubDate,
item.content || item['content:encoded'] || ''
]
);
newArticlesCount++;
logger.debug(` ✓ Ajouté: ${item.title.substring(0, 60)}...`);
}
} catch (articleError) {
// Log mais continue
if (articleError.message.includes('UNIQUE')) {
logger.debug(` - Doublon: ${item.title?.substring(0, 60)}...`);
} else {
logger.warn(` ⚠️ Erreur article: ${articleError.message}`);
}
}
}
// Nettoyer les vieux articles (garder les N derniers)
const articlesCount = await getOne(
'SELECT COUNT(*) as count FROM rss_articles WHERE feed_id = ?',
[feed.id]
);
if (articlesCount.count > MAX_ARTICLES_PER_FEED) {
const toDelete = articlesCount.count - MAX_ARTICLES_PER_FEED;
await runQuery(`
DELETE FROM rss_articles
WHERE id IN (
SELECT id FROM rss_articles
WHERE feed_id = ?
ORDER BY pub_date ASC
LIMIT ?
)
`, [feed.id, toDelete]);
logger.debug(` 🗑️ ${toDelete} anciens articles supprimés`);
}
const duration = ((Date.now() - startTime) / 1000).toFixed(2);
logger.info(`✅ ${feed.title}: ${newArticlesCount} nouveaux (${duration}s)`);
return { success: true, newArticles: newArticlesCount };
} catch (error) {
const duration = ((Date.now() - startTime) / 1000).toFixed(2);
logger.error(`❌ ${feed.title || feed.url}: ${error.message} (${duration}s)`);
return { success: false, error: error.message };
}
}
/**
* Récupérer tous les flux activés
*/
async function fetchAllFeeds() {
if (isRunning) {
logger.debug('⏭️ Fetch déjà en cours, skip');
return { skipped: true };
}
isRunning = true;
const globalStart = Date.now();
try {
logger.info('🔄 === Début mise à jour RSS ===');
// Récupérer tous les flux activés
const feeds = await getAll('SELECT * FROM rss_feeds WHERE enabled = 1');
if (!feeds || feeds.length === 0) {
logger.info('⚠️ Aucun flux RSS activé');
return { feeds: 0, newArticles: 0 };
}
logger.info(`📰 ${feeds.length} flux à traiter`);
let totalNew = 0;
let successCount = 0;
let errorCount = 0;
// Traiter chaque flux séquentiellement
for (const feed of feeds) {
const result = await fetchSingleFeed(feed);
if (result.success) {
successCount++;
totalNew += result.newArticles;
} else {
errorCount++;
}
}
const totalDuration = ((Date.now() - globalStart) / 1000).toFixed(2);
logger.info(`✅ === Fin: ${totalNew} nouveaux articles (${successCount} OK, ${errorCount} erreurs, ${totalDuration}s) ===`);
return {
feeds: feeds.length,
newArticles: totalNew,
success: successCount,
errors: errorCount
};
} catch (error) {
logger.error('❌ Erreur globale fetch RSS:', error);
return { error: error.message };
} finally {
isRunning = false;
}
}
/**
* Démarrer le scheduler
*/
function startScheduler() {
logger.info('📰 === RSS Scheduler V2 démarré ===');
logger.info(`Configuration: fetch toutes les ${FETCH_INTERVAL / 60000} minutes`);
logger.info(`Max articles par flux: ${MAX_ARTICLES_PER_FEED}`);
// Premier fetch après le démarrage
setTimeout(() => {
logger.info('🚀 Premier fetch RSS...');
fetchAllFeeds().catch(err => {
logger.error('Erreur premier fetch:', err);
});
}, STARTUP_DELAY);
// Fetch régulier
setInterval(() => {
fetchAllFeeds().catch(err => {
logger.error('Erreur fetch périodique:', err);
});
}, FETCH_INTERVAL);
}
/**
* Fetch manuel (API)
*/
async function manualFetch() {
logger.info('🔄 Fetch manuel déclenché');
return await fetchAllFeeds();
}
module.exports = {
startScheduler,
manualFetch,
fetchAllFeeds
};
+151 -137
View File
@@ -1,4 +1,4 @@
// Service de mise à jour automatique des flux RSS
// Service RSS NOUVEAU - Architecture simple et robuste
const Parser = require('rss-parser');
const { getAll, getOne, runQuery } = require('../config/database');
const logger = require('../config/logger');
@@ -6,214 +6,228 @@ const logger = require('../config/logger');
const parser = new Parser({
timeout: 10000,
headers: {
'User-Agent': 'NoteFlow RSS Reader'
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
}
});
let isRunning = false;
// Configuration
const MAX_ARTICLES_PER_FEED = 100; // Garder max 100 articles par flux
const FETCH_INTERVAL = 2 * 60 * 1000; // 2 minutes
const STARTUP_DELAY = 5000; // 5 secondes
/**
* Initialiser des flux RSS par défaut
* Récupérer et stocker les articles d'un flux
*/
async function initializeDefaultFeeds() {
async function fetchSingleFeed(feed) {
const startTime = Date.now();
let newArticlesCount = 0;
try {
const existingFeeds = await getAll('SELECT COUNT(*) as count FROM rss_feeds');
logger.info(`⏳ Récupération: ${feed.title || feed.url}`);
if (existingFeeds[0].count === 0) {
logger.info('🔧 Aucun flux RSS trouvé, ajout de flux par défaut...');
// Parser le flux avec timeout
const parsedFeed = await Promise.race([
parser.parseURL(feed.url),
new Promise((_, reject) =>
setTimeout(() => reject(new Error('Timeout après 15s')), 15000)
)
]);
const defaultFeeds = [
'https://www.lemonde.fr/rss/une.xml',
'https://feeds.bbci.co.uk/news/world/rss.xml',
'https://www.lefigaro.fr/rss/figaro_actualites.xml'
];
// Mettre à jour les infos du flux
await runQuery(
'UPDATE rss_feeds SET title = ?, description = ?, last_fetched_at = CURRENT_TIMESTAMP WHERE id = ?',
[parsedFeed.title || feed.url, parsedFeed.description || '', feed.id]
);
for (const url of defaultFeeds) {
// Traiter les articles (prendre les 50 premiers)
const items = parsedFeed.items.slice(0, 50);
logger.debug(` ${items.length} articles dans le flux`);
for (const item of items) {
try {
// Validation minimale
if (!item.title || !item.link) {
logger.debug(` ⚠️ Article ignoré: pas de titre ou lien`);
continue;
}
// Normaliser la date
let pubDate;
try {
const feed = await parser.parseURL(url);
const dateStr = item.pubDate || item.isoDate;
pubDate = dateStr ? new Date(dateStr).toISOString() : new Date().toISOString();
// Vérifier que la date est valide
if (isNaN(new Date(pubDate).getTime())) {
pubDate = new Date().toISOString();
}
} catch (e) {
pubDate = new Date().toISOString();
}
// Vérifier si l'article existe UNIQUEMENT par lien
// Simple et fiable
const existing = await getOne(
'SELECT id FROM rss_articles WHERE link = ?',
[item.link]
);
if (!existing) {
// Nouvel article - l'ajouter
await runQuery(
'INSERT INTO rss_feeds (url, title, description, enabled) VALUES (?, ?, ?, 1)',
[url, feed.title || url, feed.description || '']
`INSERT INTO rss_articles (feed_id, title, link, description, pub_date, content)
VALUES (?, ?, ?, ?, ?, ?)`,
[
feed.id,
item.title.trim(),
item.link,
item.contentSnippet || item.description || '',
pubDate,
item.content || item['content:encoded'] || ''
]
);
logger.info(`✓ Flux ajouté: ${feed.title || url}`);
} catch (error) {
logger.warn(`⚠️ Impossible d'ajouter ${url}: ${error.message}`);
newArticlesCount++;
logger.debug(` ✓ Ajouté: ${item.title.substring(0, 60)}...`);
}
} catch (articleError) {
// Log mais continue
if (articleError.message.includes('UNIQUE')) {
logger.debug(` - Doublon: ${item.title?.substring(0, 60)}...`);
} else {
logger.warn(` ⚠️ Erreur article: ${articleError.message}`);
}
}
logger.info('✓ Flux RSS par défaut initialisés');
return true;
}
return false;
// Nettoyer les vieux articles (garder les N derniers)
const articlesCount = await getOne(
'SELECT COUNT(*) as count FROM rss_articles WHERE feed_id = ?',
[feed.id]
);
if (articlesCount.count > MAX_ARTICLES_PER_FEED) {
const toDelete = articlesCount.count - MAX_ARTICLES_PER_FEED;
await runQuery(`
DELETE FROM rss_articles
WHERE id IN (
SELECT id FROM rss_articles
WHERE feed_id = ?
ORDER BY pub_date ASC
LIMIT ?
)
`, [feed.id, toDelete]);
logger.debug(` 🗑️ ${toDelete} anciens articles supprimés`);
}
const duration = ((Date.now() - startTime) / 1000).toFixed(2);
logger.info(`✅ ${feed.title}: ${newArticlesCount} nouveaux (${duration}s)`);
return { success: true, newArticles: newArticlesCount };
} catch (error) {
logger.error('Erreur lors de l\'initialisation des flux par défaut:', error);
return false;
const duration = ((Date.now() - startTime) / 1000).toFixed(2);
logger.error(`❌ ${feed.title || feed.url}: ${error.message} (${duration}s)`);
return { success: false, error: error.message };
}
}
/**
* Récupérer et mettre à jour tous les flux RSS activés
* Récupérer tous les flux activés
*/
async function fetchAllFeeds() {
if (isRunning) {
logger.info('Fetch RSS déjà en cours, skip...');
return;
logger.debug('⏭️ Fetch déjà en cours, skip');
return { skipped: true };
}
isRunning = true;
const startTime = Date.now();
const globalStart = Date.now();
try {
logger.info('🔄 Début de la mise à jour des flux RSS...');
logger.info('🔄 === Début mise à jour RSS ===');
// Récupérer tous les flux activés
const feeds = await getAll('SELECT * FROM rss_feeds WHERE enabled = 1');
if (!feeds || feeds.length === 0) {
logger.info('⚠️ Aucun flux RSS activé, initialisation...');
const initialized = await initializeDefaultFeeds();
if (initialized) {
// Réessayer avec les nouveaux flux
isRunning = false;
return await fetchAllFeeds();
}
isRunning = false;
return;
logger.info('⚠️ Aucun flux RSS activé');
return { feeds: 0, newArticles: 0 };
}
logger.info(`📰 Mise à jour de ${feeds.length} flux RSS...`);
logger.info(`📰 ${feeds.length} flux à traiter`);
let totalArticles = 0;
let totalErrors = 0;
let totalNew = 0;
let successCount = 0;
let errorCount = 0;
// Traiter chaque flux
// Traiter chaque flux séquentiellement
for (const feed of feeds) {
try {
logger.info(`⏳ Fetch: ${feed.url}`);
const result = await fetchSingleFeed(feed);
// Parser le flux avec timeout
const parsedFeed = await Promise.race([
parser.parseURL(feed.url),
new Promise((_, reject) =>
setTimeout(() => reject(new Error('Timeout')), 15000)
)
]);
// Mettre à jour le titre et description du flux
await runQuery(
'UPDATE rss_feeds SET title = ?, description = ?, last_fetched_at = CURRENT_TIMESTAMP WHERE id = ?',
[parsedFeed.title || feed.url, parsedFeed.description || '', feed.id]
);
// Ajouter les articles (limiter à 100 par flux pour capturer plus d'articles)
const items = parsedFeed.items.slice(0, 100);
let feedArticles = 0;
for (const item of items) {
try {
if (!item.link) continue; // Skip articles sans lien
// Normaliser la date pour la comparaison
const pubDate = item.pubDate || item.isoDate || new Date().toISOString();
// Vérifier si l'article existe déjà par lien OU par titre+date
// Permet de gérer les liens qui changent (tracking) et les vrais doublons
const existingByLink = await getOne(
'SELECT id FROM rss_articles WHERE link = ?',
[item.link]
);
const existingByTitleDate = await getOne(
'SELECT id FROM rss_articles WHERE feed_id = ? AND title = ? AND DATE(pub_date) = DATE(?)',
[feed.id, item.title, pubDate]
);
// Ajouter seulement si n'existe ni par lien ni par titre+date
if (!existingByLink && !existingByTitleDate) {
await runQuery(
'INSERT INTO rss_articles (feed_id, title, link, description, pub_date, content) VALUES (?, ?, ?, ?, ?, ?)',
[
feed.id,
item.title || 'Sans titre',
item.link,
item.contentSnippet || item.description || '',
pubDate,
item.content || item['content:encoded'] || ''
]
);
feedArticles++;
totalArticles++;
}
} catch (articleError) {
// Ignorer les articles en double (contrainte UNIQUE sur link)
if (!articleError.message.includes('UNIQUE')) {
logger.debug(`Article ignoré: ${articleError.message}`);
}
}
}
logger.info(`✓ ${feed.title || feed.url}: ${feedArticles} nouveaux articles`);
} catch (feedError) {
totalErrors++;
logger.error(`✗ Erreur fetch ${feed.url}: ${feedError.message}`);
if (result.success) {
successCount++;
totalNew += result.newArticles;
} else {
errorCount++;
}
}
const duration = ((Date.now() - startTime) / 1000).toFixed(2);
logger.info(`✅ Mise à jour terminée: ${totalArticles} nouveaux articles, ${totalErrors} erreurs (${duration}s)`);
const totalDuration = ((Date.now() - globalStart) / 1000).toFixed(2);
logger.info(`✅ === Fin: ${totalNew} nouveaux articles (${successCount} OK, ${errorCount} erreurs, ${totalDuration}s) ===`);
// Invalider le cache des articles dans les routes
try {
const rssRoutes = require('../routes/rss.routes');
if (rssRoutes && rssRoutes.invalidateCache) {
rssRoutes.invalidateCache();
logger.debug('Cache des articles RSS invalidé');
}
} catch (err) {
// Ignore si la fonction n'existe pas encore
}
return {
feeds: feeds.length,
newArticles: totalNew,
success: successCount,
errors: errorCount
};
} catch (error) {
logger.error('Erreur lors de la mise à jour automatique des flux RSS:', error);
logger.error('❌ Erreur globale fetch RSS:', error);
return { error: error.message };
} finally {
isRunning = false;
}
}
/**
* Démarrer le scheduler (toutes les 2 minutes pour mises à jour fréquentes)
* Démarrer le scheduler
*/
function startScheduler() {
logger.info('📰 Scheduler RSS démarré (mise à jour toutes les 2 minutes)');
logger.info('📰 === RSS Scheduler V2 démarré ===');
logger.info(`Configuration: fetch toutes les ${FETCH_INTERVAL / 60000} minutes`);
logger.info(`Max articles par flux: ${MAX_ARTICLES_PER_FEED}`);
// Initialiser les flux par défaut si nécessaire, puis première exécution
setTimeout(async () => {
await initializeDefaultFeeds();
// Premier fetch après le démarrage
setTimeout(() => {
logger.info('🚀 Premier fetch RSS...');
fetchAllFeeds().catch(err => {
logger.error('Erreur lors de la première mise à jour RSS:', err);
logger.error('Erreur premier fetch:', err);
});
}, 5000); // Attendre 5 secondes après le démarrage du serveur
}, STARTUP_DELAY);
// Ensuite toutes les 2 minutes (réduit de 5 minutes pour mises à jour plus fréquentes)
// Fetch régulier
setInterval(() => {
fetchAllFeeds().catch(err => {
logger.error('Erreur lors de la mise à jour RSS:', err);
logger.error('Erreur fetch périodique:', err);
});
}, 2 * 60 * 1000); // 2 minutes
}, FETCH_INTERVAL);
}
/**
* Fetch manuel (utilisé par la route API)
* Fetch manuel (API)
*/
async function manualFetch() {
logger.info('🔄 Fetch manuel déclenché');
return await fetchAllFeeds();
}
module.exports = {
startScheduler,
manualFetch,
fetchAllFeeds,
initializeDefaultFeeds
fetchAllFeeds
};
+219
View File
@@ -0,0 +1,219 @@
// Service de mise à jour automatique des flux RSS
const Parser = require('rss-parser');
const { getAll, getOne, runQuery } = require('../config/database');
const logger = require('../config/logger');
const parser = new Parser({
timeout: 10000,
headers: {
'User-Agent': 'NoteFlow RSS Reader'
}
});
let isRunning = false;
/**
* Initialiser des flux RSS par défaut
*/
async function initializeDefaultFeeds() {
try {
const existingFeeds = await getAll('SELECT COUNT(*) as count FROM rss_feeds');
if (existingFeeds[0].count === 0) {
logger.info('🔧 Aucun flux RSS trouvé, ajout de flux par défaut...');
const defaultFeeds = [
'https://www.lemonde.fr/rss/une.xml',
'https://feeds.bbci.co.uk/news/world/rss.xml',
'https://www.lefigaro.fr/rss/figaro_actualites.xml'
];
for (const url of defaultFeeds) {
try {
const feed = await parser.parseURL(url);
await runQuery(
'INSERT INTO rss_feeds (url, title, description, enabled) VALUES (?, ?, ?, 1)',
[url, feed.title || url, feed.description || '']
);
logger.info(`✓ Flux ajouté: ${feed.title || url}`);
} catch (error) {
logger.warn(`⚠️ Impossible d'ajouter ${url}: ${error.message}`);
}
}
logger.info('✓ Flux RSS par défaut initialisés');
return true;
}
return false;
} catch (error) {
logger.error('Erreur lors de l\'initialisation des flux par défaut:', error);
return false;
}
}
/**
* Récupérer et mettre à jour tous les flux RSS activés
*/
async function fetchAllFeeds() {
if (isRunning) {
logger.info('Fetch RSS déjà en cours, skip...');
return;
}
isRunning = true;
const startTime = Date.now();
try {
logger.info('🔄 Début de la mise à jour des flux RSS...');
// Récupérer tous les flux activés
const feeds = await getAll('SELECT * FROM rss_feeds WHERE enabled = 1');
if (!feeds || feeds.length === 0) {
logger.info('⚠️ Aucun flux RSS activé, initialisation...');
const initialized = await initializeDefaultFeeds();
if (initialized) {
// Réessayer avec les nouveaux flux
isRunning = false;
return await fetchAllFeeds();
}
isRunning = false;
return;
}
logger.info(`📰 Mise à jour de ${feeds.length} flux RSS...`);
let totalArticles = 0;
let totalErrors = 0;
// Traiter chaque flux
for (const feed of feeds) {
try {
logger.info(`⏳ Fetch: ${feed.url}`);
// Parser le flux avec timeout
const parsedFeed = await Promise.race([
parser.parseURL(feed.url),
new Promise((_, reject) =>
setTimeout(() => reject(new Error('Timeout')), 15000)
)
]);
// Mettre à jour le titre et description du flux
await runQuery(
'UPDATE rss_feeds SET title = ?, description = ?, last_fetched_at = CURRENT_TIMESTAMP WHERE id = ?',
[parsedFeed.title || feed.url, parsedFeed.description || '', feed.id]
);
// Ajouter les articles (limiter à 100 par flux pour capturer plus d'articles)
const items = parsedFeed.items.slice(0, 100);
let feedArticles = 0;
for (const item of items) {
try {
if (!item.link) continue; // Skip articles sans lien
// Normaliser la date pour la comparaison
const pubDate = item.pubDate || item.isoDate || new Date().toISOString();
// Vérifier si l'article existe déjà par lien OU par titre+date
// Permet de gérer les liens qui changent (tracking) et les vrais doublons
const existingByLink = await getOne(
'SELECT id FROM rss_articles WHERE link = ?',
[item.link]
);
const existingByTitleDate = await getOne(
'SELECT id FROM rss_articles WHERE feed_id = ? AND title = ? AND DATE(pub_date) = DATE(?)',
[feed.id, item.title, pubDate]
);
// Ajouter seulement si n'existe ni par lien ni par titre+date
if (!existingByLink && !existingByTitleDate) {
await runQuery(
'INSERT INTO rss_articles (feed_id, title, link, description, pub_date, content) VALUES (?, ?, ?, ?, ?, ?)',
[
feed.id,
item.title || 'Sans titre',
item.link,
item.contentSnippet || item.description || '',
pubDate,
item.content || item['content:encoded'] || ''
]
);
feedArticles++;
totalArticles++;
}
} catch (articleError) {
// Ignorer les articles en double (contrainte UNIQUE sur link)
if (!articleError.message.includes('UNIQUE')) {
logger.debug(`Article ignoré: ${articleError.message}`);
}
}
}
logger.info(`✓ ${feed.title || feed.url}: ${feedArticles} nouveaux articles`);
} catch (feedError) {
totalErrors++;
logger.error(`✗ Erreur fetch ${feed.url}: ${feedError.message}`);
}
}
const duration = ((Date.now() - startTime) / 1000).toFixed(2);
logger.info(`✅ Mise à jour terminée: ${totalArticles} nouveaux articles, ${totalErrors} erreurs (${duration}s)`);
// Invalider le cache des articles dans les routes
try {
const rssRoutes = require('../routes/rss.routes');
if (rssRoutes && rssRoutes.invalidateCache) {
rssRoutes.invalidateCache();
logger.debug('Cache des articles RSS invalidé');
}
} catch (err) {
// Ignore si la fonction n'existe pas encore
}
} catch (error) {
logger.error('Erreur lors de la mise à jour automatique des flux RSS:', error);
} finally {
isRunning = false;
}
}
/**
* Démarrer le scheduler (toutes les 2 minutes pour mises à jour fréquentes)
*/
function startScheduler() {
logger.info('📰 Scheduler RSS démarré (mise à jour toutes les 2 minutes)');
// Initialiser les flux par défaut si nécessaire, puis première exécution
setTimeout(async () => {
await initializeDefaultFeeds();
fetchAllFeeds().catch(err => {
logger.error('Erreur lors de la première mise à jour RSS:', err);
});
}, 5000); // Attendre 5 secondes après le démarrage du serveur
// Ensuite toutes les 2 minutes (réduit de 5 minutes pour mises à jour plus fréquentes)
setInterval(() => {
fetchAllFeeds().catch(err => {
logger.error('Erreur lors de la mise à jour RSS:', err);
});
}, 2 * 60 * 1000); // 2 minutes
}
/**
* Fetch manuel (utilisé par la route API)
*/
async function manualFetch() {
return await fetchAllFeeds();
}
module.exports = {
startScheduler,
manualFetch,
fetchAllFeeds,
initializeDefaultFeeds
};