feat(aep): M2 — worker adapté (Bifrost + scrape léger + ntfy)
Reprend worker/enrich.js (fusion, pas réécriture) pour le nouveau format de
soumission libre et l'infra en prod (décisions MOE du 27/09) :
- LLM Mistral direct → Bifrost (${BIFROST_URL}/v1/chat/completions,
x-bf-vk), modèle WORKER_MODEL (défaut groq/llama-3.1-8b-instant).
- Scrape crawl4ai/Python → fetch natif Node 22 : timeout 8s, corps
plafonné 500 Ko, titre + meta description + og:* + texte tronqué 4000c.
- Email Resend → notification ntfy, jamais l'email ni le texte du
contributeur (id NocoDB, nom suggéré, type, confiance uniquement).
- Sortie LLM {nom, description, type_suggere, ville, tags, confiance} ;
le worker ne réécrit `nom`/`submission_type` que si le formulaire assoupli
a laissé un placeholder ("[à qualifier]" / "Type : non précisé").
- Seuil « 5 fiches pending » retiré (décision volume faible, notée si le
volume remonte).
- --dry-run : fixture locale, mock Bifrost/scrape par défaut, aucune
écriture NocoDB, DRY_RUN_LIVE=1 pour forcer de vrais appels réseau.
- worker/daily-digest.js supprimé (non référencé, cron purgé le 15/07).
deploy/aep-worker/ : service + timer (15 min) + README avec les étapes
exactes de déploiement pour la session qui déploiera (backup NocoDB,
copie, systemctl, test bout-en-bout M3) — rien exécuté ni copié ici.
PIPE-IA-DOC.md §11 : documente tous les écarts vs la version NAV V2 d'origine.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+286
-184
@@ -1,14 +1,28 @@
|
||||
#!/usr/bin/env node
|
||||
/**
|
||||
* NAV V2 — Worker enrichissement IA
|
||||
* Lancé via systemd timer toutes les 5 minutes
|
||||
* Pipeline : fetch pending → scrape crawl4ai → Mistral Nemo → update NocoDB → log stats
|
||||
* AEP — Worker enrichissement IA (B5-M2)
|
||||
* Lancé via systemd timer toutes les 15 minutes (deploy/aep-worker/aep-worker.timer)
|
||||
* Pipeline : fetch pending → scrape léger (fetch natif) → Bifrost (LLM) → update NocoDB → ntfy
|
||||
*
|
||||
* Écarts à la version NAV V2 d'origine (PIPE-IA-DOC.md, fusionné pas réécrit) :
|
||||
* - LLM : Bifrost (passerelle OpenAI-compatible en prod) au lieu de Mistral direct.
|
||||
* - Scrape : fetch natif Node 22 (timeout 8s, 500 Ko max, texte tronqué 4000 car.)
|
||||
* au lieu de crawl4ai/Python (disque VPS tendu, décision MOE 27/09).
|
||||
* - Notification : ntfy (https://ntfy.sh/$NTFY_TOPIC) au lieu de Resend (abandonné 15/07).
|
||||
* Le message ne contient JAMAIS l'email ni le texte libre du contributeur.
|
||||
* - Le formulaire assoupli (B5-M1) écrit un `nom` placeholder ("[à qualifier] host") et
|
||||
* éventuellement "Type : non précisé" en tête de description_user. Le worker propose un
|
||||
* vrai nom et un type — voir §Sortie JSON ci-dessous et PIPE-IA-DOC.md.
|
||||
* - Seuil « email à 5 fiches pending » PAS réactivé (décision MOE : volume faible, une
|
||||
* notif par fiche traitée suffit — cf. B5-proposer-pipe.md §Points d'attention).
|
||||
*/
|
||||
|
||||
import { spawnSync, execSync } from 'child_process';
|
||||
import { existsSync, writeFileSync, unlinkSync } from 'fs';
|
||||
import { tmpdir } from 'os';
|
||||
import { join } from 'path';
|
||||
import { existsSync, writeFileSync, unlinkSync, readFileSync, mkdirSync } from 'fs';
|
||||
import { dirname, join } from 'path';
|
||||
import { fileURLToPath } from 'url';
|
||||
|
||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||
|
||||
// ─── CONFIG DEPUIS .env ───────────────────────────────────────────────────────
|
||||
const NOCODB_URL = process.env.NOCODB_URL || 'http://localhost:8070';
|
||||
@@ -16,25 +30,44 @@ const NOCODB_TOKEN = process.env.NOCODB_TOKEN;
|
||||
const NOCODB_BASE = process.env.NOCODB_BASE;
|
||||
const NOCODB_TABLE_ORGAS = process.env.NOCODB_TABLE_ORGAS;
|
||||
const NOCODB_TABLE_STATS = process.env.NOCODB_TABLE_STATS;
|
||||
const MISTRAL_API_KEY = process.env.MISTRAL_API_KEY;
|
||||
const RESEND_API_KEY = process.env.RESEND_API_KEY;
|
||||
const RESEND_FROM = process.env.RESEND_FROM || 'contact@trans-former.fr';
|
||||
const EMAIL_JULES = process.env.EMAIL_JULES || 'jules@trans-former.fr';
|
||||
|
||||
const BIFROST_URL = process.env.BIFROST_URL || 'http://127.0.0.1:8080';
|
||||
const BIFROST_VK = process.env.BIFROST_VK;
|
||||
const WORKER_MODEL = process.env.WORKER_MODEL || 'groq/llama-3.1-8b-instant';
|
||||
|
||||
const NTFY_TOPIC = process.env.NTFY_TOPIC;
|
||||
|
||||
const BUDGET_MAX_EUR = parseFloat(process.env.BUDGET_MAX_EUR || '20');
|
||||
const WORKER_LIMIT = parseInt(process.env.WORKER_LIMIT || '5');
|
||||
const LOCK_FILE = '/tmp/nav-worker.lock';
|
||||
const LOCK_FILE = process.env.WORKER_LOCK_FILE || '/tmp/aep-worker.lock';
|
||||
|
||||
// ─── PRIX MISTRAL NEMO (USD → EUR) ───────────────────────────────────────────
|
||||
const NEMO_PRICE_IN = 0.02 / 1_000_000; // $0.02 / 1M tokens input
|
||||
const NEMO_PRICE_OUT = 0.04 / 1_000_000; // $0.04 / 1M tokens output
|
||||
// Prix par 1M tokens (USD) — par défaut 0 (Groq gratuit dans le tier RAPIDE Bifrost).
|
||||
// Réglable si le worker bascule un jour sur un modèle payant.
|
||||
const WORKER_PRICE_IN = parseFloat(process.env.WORKER_PRICE_IN_USD_PER_M || '0') / 1_000_000;
|
||||
const WORKER_PRICE_OUT = parseFloat(process.env.WORKER_PRICE_OUT_USD_PER_M || '0') / 1_000_000;
|
||||
const USD_TO_EUR = 0.93;
|
||||
|
||||
// Scrape léger
|
||||
const SCRAPE_TIMEOUT_MS = 8_000;
|
||||
const SCRAPE_MAX_BYTES = 500_000;
|
||||
const SCRAPE_TEXT_MAX_CHARS = 4_000;
|
||||
const SCRAPE_USER_AGENT = 'AEP/2.0 contact@trans-former.fr';
|
||||
|
||||
// ─── CLI ──────────────────────────────────────────────────────────────────────
|
||||
const DRY_RUN = process.argv.includes('--dry-run');
|
||||
// DRY_RUN_LIVE=1 : en dry-run, tente quand même de vrais appels réseau (scrape + Bifrost)
|
||||
// pour un test manuel avec accès réseau — ne touche JAMAIS NocoDB en dry-run, dans tous les cas.
|
||||
const DRY_RUN_LIVE = process.env.DRY_RUN_LIVE === '1';
|
||||
const DRY_RUN_FIXTURE = process.env.DRY_RUN_FIXTURE || join(__dirname, 'fixtures', 'dry-run-rows.json');
|
||||
|
||||
// ─── TAXONOMIE VALIDE (apostrophe typographique U+2019 comme NocoDB) ─────────
|
||||
const VALID_FONCTIONS = [
|
||||
'Juridique', 'Technique', 'Économique', 'Administratif', 'Chantier',
|
||||
'Comptabilité', 'Développement', 'Formation', 'Gestion d\u2019agence', 'Santé mentale'
|
||||
'Comptabilité', 'Développement', 'Formation', 'Gestion d’agence', 'Santé mentale'
|
||||
];
|
||||
|
||||
const VALID_SUBMISSION_TYPES = ['ecosysteme', 'reseau', 'job', 'outil'];
|
||||
|
||||
// ─── MAPPING NORMALISATION TAGS ───────────────────────────────────────────────
|
||||
const TAG_MAP = [
|
||||
[['juridique', 'droit', 'litige', 'contrat', 'déontologie', 'décennale', 'médiation', 'pi ', 'propriété intellectuelle', 'ccag', 'marchés publics droit'], 'Juridique'],
|
||||
@@ -45,16 +78,14 @@ const TAG_MAP = [
|
||||
[['comptabilité', 'fiscal', 'tva', 'bnc', 'bic', 'expert-comptable', 'bilan', 'trésorerie', 'transmission agence', 'création agence', 'micro'], 'Comptabilité'],
|
||||
[['développement', 'prospection', 'commercial', 'client', 'réseau', 'candidature', 'consultation', 'acquisition', 'marketing', 'notoriété', 'ao '], 'Développement'],
|
||||
[['formation', 'école', 'mooc', 'organisme', 'formation continue', 'cpf', 'dpc', 'cfaa'], 'Formation'],
|
||||
[['gestion d\u2019agence', 'gestion d\'agence', 'rh', 'recrutement', 'emploi', 'salaire', 'ccn', 'convention collective', 'idcc', 'temps de travail', 'management'], 'Gestion d\u2019agence'],
|
||||
[['gestion d’agence', 'gestion d\'agence', 'rh', 'recrutement', 'emploi', 'salaire', 'ccn', 'convention collective', 'idcc', 'temps de travail', 'management'], 'Gestion d’agence'],
|
||||
[['santé mentale', 'burn-out', 'épuisement', 'souffrance', 'bien-être', 'harcèlement', 'stress', 'psychologique', 'équilibre'], 'Santé mentale'],
|
||||
];
|
||||
|
||||
function normalizeTag(raw) {
|
||||
const t = raw.toLowerCase().trim();
|
||||
// D'abord chercher correspondance exacte dans les valeurs valides
|
||||
const t = String(raw).toLowerCase().trim();
|
||||
const exact = VALID_FONCTIONS.find(v => v.toLowerCase() === t);
|
||||
if (exact) return exact;
|
||||
// Sinon chercher par mots-clés
|
||||
for (const [patterns, normalized] of TAG_MAP) {
|
||||
if (patterns.some(p => t.includes(p))) return normalized;
|
||||
}
|
||||
@@ -117,8 +148,18 @@ async function nocodbPost(tableId, data) {
|
||||
return res.json();
|
||||
}
|
||||
|
||||
/** Écrit une mise à jour de fiche — no-op loggé en dry-run (jamais de PATCH réel). */
|
||||
async function patchRow(rowId, data) {
|
||||
if (DRY_RUN) {
|
||||
log(`[dry-run] PATCH fiche ${rowId} :`, JSON.stringify(data));
|
||||
return;
|
||||
}
|
||||
return nocodbPatch(NOCODB_TABLE_ORGAS, rowId, data);
|
||||
}
|
||||
|
||||
// ─── BUDGET CIRCUIT BREAKER ──────────────────────────────────────────────────
|
||||
async function getBudgetMoisCourant() {
|
||||
if (DRY_RUN) return 0;
|
||||
const now = new Date();
|
||||
const year = now.getFullYear();
|
||||
const month = now.getMonth(); // 0-indexed
|
||||
@@ -142,7 +183,12 @@ async function getBudgetMoisCourant() {
|
||||
async function logUsage(usage, model, endpoint, orgaId) {
|
||||
const tokensIn = usage?.prompt_tokens || 0;
|
||||
const tokensOut = usage?.completion_tokens || 0;
|
||||
const coutEur = ((tokensIn * NEMO_PRICE_IN) + (tokensOut * NEMO_PRICE_OUT)) * USD_TO_EUR;
|
||||
const coutEur = ((tokensIn * WORKER_PRICE_IN) + (tokensOut * WORKER_PRICE_OUT)) * USD_TO_EUR;
|
||||
|
||||
if (DRY_RUN) {
|
||||
log(`[dry-run] Usage (non loggé) : ${tokensIn}in + ${tokensOut}out = €${coutEur.toFixed(6)} (${model})`);
|
||||
return coutEur;
|
||||
}
|
||||
|
||||
await nocodbPost(NOCODB_TABLE_STATS, {
|
||||
model,
|
||||
@@ -166,111 +212,168 @@ async function fetchPendingRows() {
|
||||
return data.list || [];
|
||||
}
|
||||
|
||||
// ─── SCRAPING CRAWL4AI (mode HTTP statique, sans Playwright) ─────────────────
|
||||
async function scrapeWithCrawl4ai(url) {
|
||||
log(`Scraping: ${url}`);
|
||||
function loadFixtureRows() {
|
||||
if (!existsSync(DRY_RUN_FIXTURE)) {
|
||||
throw new Error(`Fixture dry-run introuvable : ${DRY_RUN_FIXTURE}`);
|
||||
}
|
||||
const raw = JSON.parse(readFileSync(DRY_RUN_FIXTURE, 'utf-8'));
|
||||
return Array.isArray(raw) ? raw : [raw];
|
||||
}
|
||||
|
||||
// Script Python temporaire pour crawl4ai
|
||||
const urlSafe = url.replace(/\\/g, '\\\\').replace(/'/g, "\\'");
|
||||
const script = `
|
||||
import asyncio, sys
|
||||
from crawl4ai import AsyncWebCrawler, CrawlerRunConfig
|
||||
from crawl4ai.async_crawler_strategy import AsyncHTTPCrawlerStrategy
|
||||
|
||||
async def scrape():
|
||||
strategy = AsyncHTTPCrawlerStrategy()
|
||||
run_cfg = CrawlerRunConfig(
|
||||
word_count_threshold=20,
|
||||
excluded_tags=['nav', 'footer', 'script', 'style', 'head'],
|
||||
remove_overlay_elements=True
|
||||
)
|
||||
async with AsyncWebCrawler(crawler_strategy=strategy, verbose=False) as crawler:
|
||||
result = await crawler.arun(url='${urlSafe}', config=run_cfg)
|
||||
if result.success and result.markdown:
|
||||
content = result.markdown[:16000]
|
||||
sys.stdout.buffer.write(content.encode('utf-8'))
|
||||
else:
|
||||
sys.stderr.write(f"Scrape failed: success={result.success}\\n")
|
||||
sys.exit(1)
|
||||
|
||||
asyncio.run(scrape())
|
||||
`;
|
||||
|
||||
const scriptPath = join(tmpdir(), `nav-scrape-${Date.now()}.py`);
|
||||
writeFileSync(scriptPath, script, 'utf-8');
|
||||
// ─── LIENS — extraits de url + description_user (format B5-M1 "Liens :\n...") ─
|
||||
function extractLinks(row) {
|
||||
const text = `${row.url || ''}\n${row.description_user || row.description || ''}`;
|
||||
const matches = text.match(/https?:\/\/[^\s)"'<>]+/g) || [];
|
||||
const cleaned = matches.map(u => u.replace(/[.,;:!?]+$/, ''));
|
||||
return [...new Set(cleaned)];
|
||||
}
|
||||
|
||||
// ─── SCRAPE LÉGER (fetch natif, pas de crawl4ai) ─────────────────────────────
|
||||
async function fetchCapped(url) {
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), SCRAPE_TIMEOUT_MS);
|
||||
try {
|
||||
const result = spawnSync('python3', [scriptPath], {
|
||||
timeout: 180_000,
|
||||
maxBuffer: 20 * 1024 * 1024,
|
||||
env: { ...process.env, PYTHONDONTWRITEBYTECODE: '1', PYTHONIOENCODING: 'utf-8' }
|
||||
const res = await fetch(url, {
|
||||
signal: controller.signal,
|
||||
redirect: 'follow',
|
||||
headers: { 'User-Agent': SCRAPE_USER_AGENT },
|
||||
});
|
||||
if (!res.ok) throw new Error(`HTTP ${res.status}`);
|
||||
|
||||
if (result.error) throw result.error;
|
||||
const reader = res.body?.getReader();
|
||||
if (!reader) return await res.text();
|
||||
|
||||
const stdout = result.stdout?.toString('utf-8') || '';
|
||||
const stderr = result.stderr?.toString('utf-8') || '';
|
||||
|
||||
if (result.status === 0 && stdout.length > 50) {
|
||||
log(`Scrape OK: ${stdout.length} chars`);
|
||||
return stdout.trim();
|
||||
} else {
|
||||
throw new Error(`Scrape failed (code ${result.status}): ${stderr.slice(0, 300)}`);
|
||||
const chunks = [];
|
||||
let total = 0;
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
chunks.push(value);
|
||||
total += value.length;
|
||||
if (total >= SCRAPE_MAX_BYTES) {
|
||||
await reader.cancel().catch(() => {});
|
||||
break;
|
||||
}
|
||||
}
|
||||
const buf = Buffer.concat(chunks.map(c => Buffer.from(c)));
|
||||
return buf.subarray(0, SCRAPE_MAX_BYTES).toString('utf-8');
|
||||
} finally {
|
||||
try { unlinkSync(scriptPath); } catch {}
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
// ─── APPEL MISTRAL NEMO ───────────────────────────────────────────────────────
|
||||
const SYSTEM_PROMPT = `Tu es un assistant spécialisé dans l'écosystème professionnel de l'architecture en France. Tu reçois des informations sur une organisation ou ressource liée au secteur de l'architecture, et tu dois les enrichir pour alimenter une cartographie collaborative.
|
||||
function extractMeta(html, key) {
|
||||
const re1 = new RegExp(`<meta[^>]*(?:name|property)=["']${key}["'][^>]*content=["']([^"']*)["']`, 'i');
|
||||
const re2 = new RegExp(`<meta[^>]*content=["']([^"']*)["'][^>]*(?:name|property)=["']${key}["']`, 'i');
|
||||
return (html.match(re1) || html.match(re2))?.[1]?.trim() || null;
|
||||
}
|
||||
|
||||
function extractOgTags(html) {
|
||||
const og = {};
|
||||
const re1 = /<meta[^>]*property=["']og:([a-zA-Z:_-]+)["'][^>]*content=["']([^"']*)["']/gi;
|
||||
const re2 = /<meta[^>]*content=["']([^"']*)["'][^>]*property=["']og:([a-zA-Z:_-]+)["']/gi;
|
||||
let m;
|
||||
while ((m = re1.exec(html))) og[m[1]] = m[2];
|
||||
while ((m = re2.exec(html))) if (!(m[2] in og)) og[m[2]] = m[1];
|
||||
return og;
|
||||
}
|
||||
|
||||
function decodeEntities(s) {
|
||||
return s
|
||||
.replace(/ /g, ' ')
|
||||
.replace(/&/g, '&')
|
||||
.replace(/</g, '<')
|
||||
.replace(/>/g, '>')
|
||||
.replace(/"/g, '"')
|
||||
.replace(/�?39;/g, "'");
|
||||
}
|
||||
|
||||
function extractVisibleText(html) {
|
||||
let s = html
|
||||
.replace(/<!--[\s\S]*?-->/g, ' ')
|
||||
.replace(/<script[\s\S]*?<\/script>/gi, ' ')
|
||||
.replace(/<style[\s\S]*?<\/style>/gi, ' ')
|
||||
.replace(/<nav[\s\S]*?<\/nav>/gi, ' ')
|
||||
.replace(/<footer[\s\S]*?<\/footer>/gi, ' ')
|
||||
.replace(/<head[\s\S]*?<\/head>/gi, ' ');
|
||||
s = s.replace(/<[^>]+>/g, ' ');
|
||||
s = decodeEntities(s);
|
||||
s = s.replace(/\s+/g, ' ').trim();
|
||||
return s.slice(0, SCRAPE_TEXT_MAX_CHARS);
|
||||
}
|
||||
|
||||
function extractTitle(html) {
|
||||
return html.match(/<title[^>]*>([^<]*)<\/title>/i)?.[1]?.trim() || null;
|
||||
}
|
||||
|
||||
async function scrapeLight(url) {
|
||||
log(`Scraping léger: ${url}`);
|
||||
const html = await fetchCapped(url);
|
||||
return {
|
||||
title: extractTitle(html),
|
||||
metaDescription: extractMeta(html, 'description'),
|
||||
og: extractOgTags(html),
|
||||
text: extractVisibleText(html),
|
||||
};
|
||||
}
|
||||
|
||||
const MOCK_SCRAPE_RESULT = {
|
||||
title: '[mock dry-run] Titre de la page',
|
||||
metaDescription: '[mock dry-run] Meta description factice, aucun réseau contacté.',
|
||||
og: { site_name: '[mock dry-run]' },
|
||||
text: '[mock dry-run] Extrait de texte factice utilisé quand DRY_RUN_LIVE n’est pas activé.',
|
||||
};
|
||||
|
||||
// ─── APPEL LLM VIA BIFROST ────────────────────────────────────────────────────
|
||||
const SYSTEM_PROMPT = `Tu es un assistant qui aide à qualifier des ressources soumises pour une cartographie collaborative de l'écosystème professionnel de l'architecture en France (AEP).
|
||||
|
||||
RÈGLES ABSOLUES :
|
||||
1. Tu ne dois JAMAIS inventer d'informations non présentes dans les sources fournies.
|
||||
2. Si une information est absente ou incertaine, retourne \`null\` pour ce champ.
|
||||
3. Tu dois retourner UNIQUEMENT un objet JSON valide, sans texte avant ou après.
|
||||
4. La description_enrichie doit être neutre, factuelle, en français, sans jugement de valeur.
|
||||
5. Les points_cles sont des phrases courtes (max 12 mots chacune), actionnables pour un architecte.
|
||||
6. Pour les tags_fonction, ne propose que des valeurs parmi la liste autorisée.
|
||||
4. "description" : neutre, factuelle, en français, max 300 caractères, sans jugement de valeur.
|
||||
5. "nom" : un nom de fiche court et identifiable, sans le préfixe "[à qualifier]".
|
||||
6. "tags" : 1 à 5 valeurs, uniquement parmi la liste autorisée.
|
||||
|
||||
TAXONOMIE AUTORISÉE :
|
||||
- Échelle (une seule valeur) : "National" | "Régional" | "Départemental" | "Local"
|
||||
- Territoire (une seule valeur) : "Métropole" | "Guadeloupe" | "Martinique" | "Guyane" | "Réunion" | "Mayotte" | null
|
||||
- Tags fonction (1 à 5 valeurs) : "Juridique" | "Technique" | "Économique" | "Administratif" | "Chantier" | "Comptabilité" | "Développement" | "Formation" | "Gestion d\u2019agence" | "Santé mentale"
|
||||
TAGS AUTORISÉS : "Juridique" | "Technique" | "Économique" | "Administratif" | "Chantier" | "Comptabilité" | "Développement" | "Formation" | "Gestion d’agence" | "Santé mentale"
|
||||
TYPE_SUGGERE AUTORISÉ (une seule valeur ou null) : "ecosysteme" | "reseau" | "job" | "outil"
|
||||
|
||||
FORMAT DE SORTIE JSON :
|
||||
{
|
||||
"description_enrichie": "string (max 300 chars, français, neutre, factuel)",
|
||||
"points_cles": ["string", "string", "string"],
|
||||
"tags_fonction": ["Valeur1", "Valeur2"],
|
||||
"echelle": "National" | "Régional" | "Départemental" | "Local" | null,
|
||||
"territoire": "Métropole" | ... | null,
|
||||
"localisation_ville": "string" | null,
|
||||
"nom": "string | null",
|
||||
"description": "string (max 300 chars, français, neutre, factuel) | null",
|
||||
"type_suggere": "ecosysteme" | "reseau" | "job" | "outil" | null,
|
||||
"ville": "string | null",
|
||||
"tags": ["string", "..."],
|
||||
"confiance": "haute" | "moyenne" | "faible"
|
||||
}
|
||||
|
||||
Le champ "confiance" reflète ta certitude globale sur l'enrichissement :
|
||||
- "haute" : URL scrapée avec contenu riche, informations claires
|
||||
- "moyenne" : URL scrapée mais contenu partiel, ou description_user seule suffisante
|
||||
- "faible" : URL non disponible et description_user vague, inférences importantes`;
|
||||
Le champ "confiance" reflète ta certitude globale :
|
||||
- "haute" : contenu scrapé riche, informations claires.
|
||||
- "moyenne" : contenu partiel, ou texte du contributeur seul mais clair.
|
||||
- "faible" : aucun contenu scrapé et texte du contributeur vague, inférences importantes.`;
|
||||
|
||||
function buildUserPrompt(row, scrapeContent) {
|
||||
return `ORGANISATION À ENRICHIR :
|
||||
function buildUserPrompt(row, scrapeData, autresLiens) {
|
||||
return `RESSOURCE À QUALIFIER :
|
||||
|
||||
Nom : ${row.nom}
|
||||
URL : ${row.url || 'non fournie'}
|
||||
Description soumise par l'utilisateur : ${row.description_user || row.description || 'non fournie'}
|
||||
Nom actuel (placeholder à remplacer) : ${row.nom}
|
||||
Lien principal : ${row.url || 'non fourni'}
|
||||
${autresLiens.length ? `Autres liens mentionnés par le contributeur : ${autresLiens.join(', ')}` : ''}
|
||||
Texte du contributeur (pourquoi c'est pertinent pour lui) : ${row.description_user || row.description || 'non fourni'}
|
||||
|
||||
CONTENU DU SITE WEB (extrait par scraping) :
|
||||
${scrapeContent || 'Site non accessible ou URL non fournie.'}
|
||||
CONTENU EXTRAIT DU SITE (scraping léger) :
|
||||
Titre : ${scrapeData?.title || 'non disponible'}
|
||||
Meta description : ${scrapeData?.metaDescription || 'non disponible'}
|
||||
Open Graph : ${scrapeData && Object.keys(scrapeData.og || {}).length ? JSON.stringify(scrapeData.og) : 'non disponible'}
|
||||
Texte visible (extrait) : ${scrapeData?.text || 'Site non accessible ou lien non fourni.'}
|
||||
|
||||
---
|
||||
|
||||
Enrichis cette fiche selon les règles du system prompt. Retourne uniquement le JSON.`;
|
||||
Qualifie cette ressource selon les règles du system prompt. Retourne uniquement le JSON.`;
|
||||
}
|
||||
|
||||
async function callMistralWithRetry(row, scrapeContent, maxRetries = 2) {
|
||||
const userPrompt = buildUserPrompt(row, scrapeContent);
|
||||
async function callBifrostWithRetry(row, scrapeData, autresLiens, maxRetries = 2) {
|
||||
const userPrompt = buildUserPrompt(row, scrapeData, autresLiens);
|
||||
|
||||
for (let attempt = 0; attempt <= maxRetries; attempt++) {
|
||||
if (attempt > 0) {
|
||||
@@ -279,14 +382,14 @@ async function callMistralWithRetry(row, scrapeContent, maxRetries = 2) {
|
||||
}
|
||||
|
||||
try {
|
||||
const res = await fetch('https://api.mistral.ai/v1/chat/completions', {
|
||||
const res = await fetch(`${BIFROST_URL}/v1/chat/completions`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'Authorization': `Bearer ${MISTRAL_API_KEY}`,
|
||||
'x-bf-vk': BIFROST_VK,
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
body: JSON.stringify({
|
||||
model: 'open-mistral-nemo',
|
||||
model: WORKER_MODEL,
|
||||
temperature: 0.2,
|
||||
max_tokens: 800,
|
||||
response_format: { type: 'json_object' },
|
||||
@@ -300,101 +403,88 @@ async function callMistralWithRetry(row, scrapeContent, maxRetries = 2) {
|
||||
|
||||
if (!res.ok) {
|
||||
const err = await res.text();
|
||||
throw new Error(`Mistral API ${res.status}: ${err}`);
|
||||
throw new Error(`Bifrost API ${res.status}: ${err}`);
|
||||
}
|
||||
|
||||
const data = await res.json();
|
||||
const content = data.choices?.[0]?.message?.content;
|
||||
if (!content) throw new Error('Réponse Mistral vide');
|
||||
if (!content) throw new Error('Réponse Bifrost vide');
|
||||
|
||||
const parsed = JSON.parse(content);
|
||||
// Attacher usage pour logging
|
||||
parsed._usage = data.usage;
|
||||
parsed._raw = content;
|
||||
return parsed;
|
||||
|
||||
} catch (e) {
|
||||
log(`Erreur Mistral (tentative ${attempt + 1}): ${e.message}`);
|
||||
log(`Erreur Bifrost (tentative ${attempt + 1}): ${e.message}`);
|
||||
if (attempt === maxRetries) return null;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
// ─── EMAIL JULES VIA RESEND ───────────────────────────────────────────────────
|
||||
async function sendEmailJules(subject, body) {
|
||||
if (!RESEND_API_KEY) {
|
||||
log('RESEND_API_KEY absent, email skippé');
|
||||
const MOCK_BIFROST_RESULT = {
|
||||
nom: '[mock dry-run] Nom suggéré',
|
||||
description: '[mock dry-run] Description factice générée sans appel réseau.',
|
||||
type_suggere: 'ecosysteme',
|
||||
ville: null,
|
||||
tags: ['Développement'],
|
||||
confiance: 'faible',
|
||||
_usage: { prompt_tokens: 0, completion_tokens: 0 },
|
||||
_raw: '{"mock":"dry-run"}',
|
||||
};
|
||||
|
||||
// ─── NOTIFICATION NTFY (jamais l'email ni le texte du contributeur) ──────────
|
||||
async function notifyNtfy(title, message) {
|
||||
if (DRY_RUN) {
|
||||
log(`[dry-run] ntfy (non envoyé) — ${title} :`, message.replace(/\n/g, ' | '));
|
||||
return;
|
||||
}
|
||||
if (!NTFY_TOPIC) {
|
||||
log('NTFY_TOPIC absent, notification skippée');
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const res = await fetch('https://api.resend.com/emails', {
|
||||
const res = await fetch(`https://ntfy.sh/${NTFY_TOPIC}`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'Authorization': `Bearer ${RESEND_API_KEY}`,
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
body: JSON.stringify({
|
||||
from: RESEND_FROM,
|
||||
to: [EMAIL_JULES],
|
||||
subject,
|
||||
text: body
|
||||
})
|
||||
headers: { Title: title, Priority: '3', Tags: 'inbox_tray' },
|
||||
body: message,
|
||||
});
|
||||
if (!res.ok) log('Email error:', await res.text());
|
||||
else log('Email envoyé à Jules:', subject);
|
||||
if (!res.ok) log('ntfy erreur:', await res.text());
|
||||
else log('ntfy envoyé:', title);
|
||||
} catch (e) {
|
||||
log('Email exception:', e.message);
|
||||
}
|
||||
}
|
||||
|
||||
// ─── VÉRIFICATION SEUIL 5 FICHES PENDING MODÉRATION ─────────────────────────
|
||||
async function checkModerationQueue() {
|
||||
const data = await nocodbGet(
|
||||
`${NOCODB_TABLE_ORGAS}?where=(moderation_status,eq,ai_processed)&limit=100`
|
||||
);
|
||||
const count = data.pageInfo?.totalRows || 0;
|
||||
if (count >= 5) {
|
||||
log(`Seuil modération atteint: ${count} fiches en attente`);
|
||||
await sendEmailJules(
|
||||
`NAV — ${count} fiches à modérer`,
|
||||
`Bonjour Jules,\n\n${count} fiches ont été enrichies par l'IA et attendent ta validation dans NocoDB.\n\nLien NocoDB : http://localhost:8070\nFiltre : moderation_status = ai_processed\n\nBonne modération !`
|
||||
);
|
||||
log('ntfy exception:', e.message);
|
||||
}
|
||||
}
|
||||
|
||||
// ─── MAIN ─────────────────────────────────────────────────────────────────────
|
||||
async function run() {
|
||||
if (!acquireLock()) {
|
||||
if (!DRY_RUN && !acquireLock()) {
|
||||
log('Worker déjà en cours, skip');
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
const startTime = Date.now();
|
||||
log('=== Worker NAV enrichissement démarré ===');
|
||||
log(`=== Worker AEP enrichissement démarré${DRY_RUN ? ' (--dry-run)' : ''} ===`);
|
||||
|
||||
try {
|
||||
// Vérification clés obligatoires
|
||||
if (!NOCODB_TOKEN || !NOCODB_BASE || !NOCODB_TABLE_ORGAS || !MISTRAL_API_KEY) {
|
||||
throw new Error('Variables .env manquantes (NOCODB_TOKEN, NOCODB_BASE, NOCODB_TABLE_ORGAS, MISTRAL_API_KEY)');
|
||||
if (!DRY_RUN) {
|
||||
if (!NOCODB_TOKEN || !NOCODB_BASE || !NOCODB_TABLE_ORGAS || !BIFROST_VK) {
|
||||
throw new Error('Variables .env manquantes (NOCODB_TOKEN, NOCODB_BASE, NOCODB_TABLE_ORGAS, BIFROST_VK)');
|
||||
}
|
||||
}
|
||||
|
||||
// Check budget global
|
||||
const budgetMois = await getBudgetMoisCourant();
|
||||
log(`Budget mois courant: €${budgetMois.toFixed(4)} / €${BUDGET_MAX_EUR}`);
|
||||
|
||||
if (budgetMois >= BUDGET_MAX_EUR) {
|
||||
log('Budget épuisé pour ce mois. Worker en pause.');
|
||||
await sendEmailJules(
|
||||
'NAV — Budget IA épuisé ce mois',
|
||||
`Le budget IA de ${BUDGET_MAX_EUR}€ a été atteint. Le worker est en pause jusqu'au 1er du mois prochain.\n\nConsommation actuelle : €${budgetMois.toFixed(4)}`
|
||||
);
|
||||
await notifyNtfy('[AEP] Budget IA épuisé', `Le budget de ${BUDGET_MAX_EUR}€ a été atteint. Worker en pause jusqu'au 1er du mois prochain.`);
|
||||
return;
|
||||
}
|
||||
|
||||
// Fetch fiches pending
|
||||
const rows = await fetchPendingRows();
|
||||
log(`${rows.length} fiche(s) à traiter`);
|
||||
const rows = DRY_RUN ? loadFixtureRows() : await fetchPendingRows();
|
||||
log(`${rows.length} fiche(s) à traiter${DRY_RUN ? ' (fixture)' : ''}`);
|
||||
|
||||
if (rows.length === 0) {
|
||||
log('Rien à traiter.');
|
||||
@@ -407,71 +497,86 @@ async function run() {
|
||||
const rowStart = Date.now();
|
||||
log(`--- Traitement fiche ${row.Id}: ${row.nom} ---`);
|
||||
|
||||
// Re-check budget avant chaque fiche
|
||||
const budgetCheck = await getBudgetMoisCourant();
|
||||
if (budgetCheck >= BUDGET_MAX_EUR) {
|
||||
log('Budget atteint mid-pipeline, arrêt.');
|
||||
break;
|
||||
}
|
||||
|
||||
// Scraping
|
||||
let scrapeContent = null;
|
||||
const hasUrl = row.url && row.url.trim().length > 0;
|
||||
const shouldScrape = hasUrl && (row.scrape_status === 'pending' || !row.scrape_status);
|
||||
const liens = extractLinks(row);
|
||||
const primaryUrl = row.url && row.url.trim() ? row.url.trim() : null;
|
||||
const autresLiens = liens.filter(l => l !== primaryUrl);
|
||||
|
||||
// Scraping (uniquement le lien principal — un seul champ scrape_content en base)
|
||||
let scrapeData = null;
|
||||
const shouldScrape = primaryUrl && (row.scrape_status === 'pending' || !row.scrape_status);
|
||||
const useNetwork = !DRY_RUN || DRY_RUN_LIVE;
|
||||
|
||||
if (shouldScrape) {
|
||||
try {
|
||||
scrapeContent = await scrapeWithCrawl4ai(row.url);
|
||||
await nocodbPatch(NOCODB_TABLE_ORGAS, row.Id, {
|
||||
scrapeData = useNetwork ? await scrapeLight(primaryUrl) : MOCK_SCRAPE_RESULT;
|
||||
await patchRow(row.Id, {
|
||||
scrape_status: 'scraped',
|
||||
scrape_content: scrapeContent
|
||||
scrape_content: JSON.stringify(scrapeData),
|
||||
});
|
||||
} catch (e) {
|
||||
log(`Scrape échoué: ${e.message}`);
|
||||
await nocodbPatch(NOCODB_TABLE_ORGAS, row.Id, { scrape_status: 'failed' });
|
||||
// Continue avec l'IA sans contenu scrape
|
||||
await patchRow(row.Id, { scrape_status: 'failed' });
|
||||
}
|
||||
} else if (!hasUrl) {
|
||||
await nocodbPatch(NOCODB_TABLE_ORGAS, row.Id, { scrape_status: 'no_link' });
|
||||
} else if (!primaryUrl) {
|
||||
await patchRow(row.Id, { scrape_status: 'no_link' });
|
||||
}
|
||||
|
||||
// Appel Mistral Nemo
|
||||
const enriched = await callMistralWithRetry(row, scrapeContent);
|
||||
// Appel LLM (Bifrost)
|
||||
const enriched = useNetwork
|
||||
? await callBifrostWithRetry(row, scrapeData, autresLiens)
|
||||
: MOCK_BIFROST_RESULT;
|
||||
|
||||
if (!enriched) {
|
||||
log(`Échec Mistral sur fiche ${row.Id}, flag ai_error`);
|
||||
await nocodbPatch(NOCODB_TABLE_ORGAS, row.Id, {
|
||||
moderation_status: 'ai_error',
|
||||
ai_processed: true
|
||||
});
|
||||
log(`Échec Bifrost sur fiche ${row.Id}, flag ai_error`);
|
||||
await patchRow(row.Id, { moderation_status: 'ai_error', ai_processed: true });
|
||||
continue;
|
||||
}
|
||||
|
||||
// Normalisation tags
|
||||
const rawTags = enriched.tags_fonction || [];
|
||||
const rawTags = enriched.tags || enriched.tags_fonction || [];
|
||||
const normalizedTags = [...new Set(rawTags.map(normalizeTag).filter(Boolean))];
|
||||
|
||||
// Update NocoDB
|
||||
const updateData = {
|
||||
description_enrichie: enriched.description_enrichie || null,
|
||||
points_cles: enriched.points_cles ? JSON.stringify(enriched.points_cles) : null,
|
||||
description_enrichie: enriched.description || null,
|
||||
tags_fonction: normalizedTags.join(','),
|
||||
moderation_status: 'ai_processed',
|
||||
ai_processed: true,
|
||||
ai_raw_output: JSON.stringify({ output: enriched, confiance: enriched.confiance })
|
||||
ai_raw_output: JSON.stringify({ output: enriched, confiance: enriched.confiance }),
|
||||
};
|
||||
|
||||
// Conserver echelle/territoire/localisation si l'IA les a enrichis
|
||||
if (enriched.echelle && !row.echelle) updateData.echelle = enriched.echelle;
|
||||
if (enriched.territoire && !row.territoire) updateData.territoire = enriched.territoire;
|
||||
if (enriched.localisation_ville && !row.localisation_ville) {
|
||||
updateData.localisation_ville = enriched.localisation_ville;
|
||||
// Nom : ne remplacer que le placeholder posé par le formulaire assoupli (B5-M1).
|
||||
if (enriched.nom && typeof row.nom === 'string' && row.nom.startsWith('[à qualifier]')) {
|
||||
updateData.nom = enriched.nom;
|
||||
}
|
||||
if (enriched.ville && !row.localisation_ville) {
|
||||
updateData.localisation_ville = enriched.ville;
|
||||
}
|
||||
// submission_type : ne réajuster que si l'utilisateur n'avait pas choisi de chip
|
||||
// (marqueur posé par utils/submitLibre.ts::toOrgaPayload — "Type : non précisé").
|
||||
const typeNonPrecise = typeof row.description_user === 'string'
|
||||
&& row.description_user.startsWith('Type : non précisé');
|
||||
if (typeNonPrecise && VALID_SUBMISSION_TYPES.includes(enriched.type_suggere)) {
|
||||
updateData.submission_type = enriched.type_suggere;
|
||||
}
|
||||
|
||||
await nocodbPatch(NOCODB_TABLE_ORGAS, row.Id, updateData);
|
||||
await patchRow(row.Id, updateData);
|
||||
await logUsage(enriched._usage, WORKER_MODEL, 'enrichissement', row.Id);
|
||||
|
||||
// Log usage tokens
|
||||
await logUsage(enriched._usage, 'open-mistral-nemo', 'enrichissement', row.Id);
|
||||
await notifyNtfy(
|
||||
'[AEP] Fiche enrichie',
|
||||
[
|
||||
`Id NocoDB : ${row.Id}`,
|
||||
`Nom suggéré : ${updateData.nom || row.nom}`,
|
||||
`Type : ${updateData.submission_type || row.submission_type || 'nc'}`,
|
||||
`Confiance : ${enriched.confiance || 'nc'}`,
|
||||
].join('\n'),
|
||||
);
|
||||
|
||||
const elapsed = ((Date.now() - rowStart) / 1000).toFixed(1);
|
||||
log(`Fiche ${row.Id} traitée en ${elapsed}s — confiance: ${enriched.confiance || 'nc'}`);
|
||||
@@ -480,15 +585,12 @@ async function run() {
|
||||
|
||||
log(`=== Run terminé: ${processedCount}/${rows.length} fiches traitées en ${((Date.now() - startTime) / 1000).toFixed(1)}s ===`);
|
||||
|
||||
// Check seuil modération
|
||||
await checkModerationQueue();
|
||||
|
||||
} catch (e) {
|
||||
log('ERREUR WORKER:', e.message);
|
||||
console.error(e.stack);
|
||||
process.exit(1);
|
||||
process.exitCode = 1;
|
||||
} finally {
|
||||
releaseLock();
|
||||
if (!DRY_RUN) releaseLock();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user