#!/usr/bin/env node
/**
* 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, 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';
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 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/openai/gpt-oss-20b'; // llama-3.1-8b-instant retiré par Groq (28/09)
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 = process.env.WORKER_LOCK_FILE || '/tmp/aep-worker.lock';
// 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’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'],
[['technique', 're2020', 'thermique', 'structure', 'bim', 'dtu', 'acoustique', 'matériaux', 'simulation', 'opr', 'réserves', 'acv', 'pcd', 'eurocodes', 'pmc'], 'Technique'],
[['économique', 'prix', 'tarif', 'honoraire', 'devis', 'roi', 'financement', 'subvention', 'cee', 'maprimerenov', 'anah', 'business plan', 'pricing'], 'Économique'],
[['administratif', 'permis', 'plu', 'plui', 'erp', 'autorisation travaux', 'abf', 'patrimoine', 'marchés publics procédure', 'cctp', 'dpgf', 'concours', 'urbanisme'], 'Administratif'],
[['chantier', 'coordination chantier', 'det', 'suivi travaux', 'sps', 'sécurité chantier', 'planning chantier', 'entreprise', 'sous-traitance', 'réception travaux'], 'Chantier'],
[['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’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 = String(raw).toLowerCase().trim();
const exact = VALID_FONCTIONS.find(v => v.toLowerCase() === t);
if (exact) return exact;
for (const [patterns, normalized] of TAG_MAP) {
if (patterns.some(p => t.includes(p))) return normalized;
}
return null;
}
// ─── UTILITAIRES LOG ─────────────────────────────────────────────────────────
function log(...args) {
const ts = new Date().toISOString();
console.log(`[${ts}]`, ...args);
}
// ─── LOCK ANTI-OVERLAP ────────────────────────────────────────────────────────
function acquireLock() {
if (existsSync(LOCK_FILE)) {
const content = execSync(`cat ${LOCK_FILE}`).toString().trim();
const pid = parseInt(content);
try {
execSync(`kill -0 ${pid} 2>/dev/null`);
return false; // Process encore vivant
} catch {
log('Lock orphelin détecté, suppression');
unlinkSync(LOCK_FILE);
}
}
writeFileSync(LOCK_FILE, process.pid.toString());
return true;
}
function releaseLock() {
try { unlinkSync(LOCK_FILE); } catch {}
}
// ─── NOCODB API ──────────────────────────────────────────────────────────────
async function nocodbGet(path) {
const res = await fetch(`${NOCODB_URL}/api/v1/db/data/noco/${NOCODB_BASE}/${path}`, {
headers: { 'xc-token': NOCODB_TOKEN }
});
if (!res.ok) throw new Error(`NocoDB GET ${path} → ${res.status}: ${await res.text()}`);
return res.json();
}
async function nocodbPatch(tableId, rowId, data) {
const res = await fetch(`${NOCODB_URL}/api/v1/db/data/noco/${NOCODB_BASE}/${tableId}/${rowId}`, {
method: 'PATCH',
headers: { 'xc-token': NOCODB_TOKEN, 'Content-Type': 'application/json' },
body: JSON.stringify(data)
});
if (!res.ok) throw new Error(`NocoDB PATCH ${tableId}/${rowId} → ${res.status}: ${await res.text()}`);
return res.json();
}
async function nocodbPost(tableId, data) {
const res = await fetch(`${NOCODB_URL}/api/v1/db/data/noco/${NOCODB_BASE}/${tableId}`, {
method: 'POST',
headers: { 'xc-token': NOCODB_TOKEN, 'Content-Type': 'application/json' },
body: JSON.stringify(data)
});
if (!res.ok) throw new Error(`NocoDB POST ${tableId} → ${res.status}: ${await res.text()}`);
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
// NocoDB ne supporte pas bien le filtre datetime — on récupère tout et filtre en JS
try {
const data = await nocodbGet(`${NOCODB_TABLE_STATS}?limit=1000&sort=-timestamp`);
const total = (data.list || []).reduce((sum, row) => {
const ts = new Date(row.timestamp || row.CreatedAt || 0);
if (ts.getFullYear() === year && ts.getMonth() === month) {
return sum + (parseFloat(row.cout_eur) || 0);
}
return sum;
}, 0);
return total;
} catch (e) {
log('Erreur lecture budget:', e.message);
return 0;
}
}
async function logUsage(usage, model, endpoint, orgaId) {
const tokensIn = usage?.prompt_tokens || 0;
const tokensOut = usage?.completion_tokens || 0;
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,
endpoint,
tokens_in: tokensIn,
tokens_out: tokensOut,
cout_eur: parseFloat(coutEur.toFixed(6)),
timestamp: new Date().toISOString(),
orga_id: orgaId || null
});
log(`Usage log: ${tokensIn}in + ${tokensOut}out = €${coutEur.toFixed(6)} (${model})`);
return coutEur;
}
// ─── FETCH FICHES PENDING ────────────────────────────────────────────────────
async function fetchPendingRows() {
const data = await nocodbGet(
`${NOCODB_TABLE_ORGAS}?where=(moderation_status,eq,pending)~and(ai_processed,eq,false)&limit=${WORKER_LIMIT}&sort=submitted_at`
);
return data.list || [];
}
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];
}
// ─── 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 res = await fetch(url, {
signal: controller.signal,
redirect: 'follow',
headers: { 'User-Agent': SCRAPE_USER_AGENT },
});
if (!res.ok) throw new Error(`HTTP ${res.status}`);
const reader = res.body?.getReader();
if (!reader) return await res.text();
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 {
clearTimeout(timer);
}
}
function extractMeta(html, key) {
const re1 = new RegExp(`]*(?:name|property)=["']${key}["'][^>]*content=["']([^"']*)["']`, 'i');
const re2 = new RegExp(`]*content=["']([^"']*)["'][^>]*(?:name|property)=["']${key}["']`, 'i');
return (html.match(re1) || html.match(re2))?.[1]?.trim() || null;
}
function extractOgTags(html) {
const og = {};
const re1 = /]*property=["']og:([a-zA-Z:_-]+)["'][^>]*content=["']([^"']*)["']/gi;
const re2 = /]*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(//g, ' ')
.replace(/