Files
demat-facturation/server/emailImportService.ts
Manus cd61c75c4c Checkpoint: Checkpoint saved: Implémentation complète de l'import automatique depuis dossier et du système de notifications.
Nouvelles fonctionnalités :

**1. Import automatique depuis dossier (folderImportService.ts)**
 Surveillance périodique d'un dossier local configuré
 Détection automatique des nouveaux fichiers PDF
 Traitement identique à l'upload manuel (extraction Mistral AI, détection doublons, stockage)
 Déplacement automatique des fichiers traités vers un sous-dossier "processed"
 Scheduler configurable (fréquence en minutes)
 Routes tRPC pour démarrer/arrêter/vérifier le statut du service
 Boutons de contrôle dans la page Paramètres de réception

**2. Système de notifications (notificationService.ts)**
 Notifications automatiques après chaque import (email ou dossier)
 Résumé détaillé : nombre de factures importées, doublons ignorés, erreurs
 Format adapté selon la source (email avec nom de fichier, dossier avec nombre de fichiers)
 Intégration avec le système de notifications Manus (notifyOwner)
 Notifications envoyées au propriétaire du projet

**Fonctionnement de l'import dossier :**
1. L'utilisateur configure le chemin du dossier source dans "Paramètres de réception"
2. Il active l'import automatique et configure la fréquence de vérification
3. Il clique sur "Démarrer le service" pour lancer la surveillance
4. Le service scanne le dossier selon la fréquence configurée
5. Il détecte les nouveaux fichiers PDF (non encore traités)
6. Il traite chaque PDF avec la même logique que l'upload manuel
7. Les fichiers traités sont déplacés vers un sous-dossier "processed"
8. Une notification est envoyée avec le résumé des résultats

**Fonctionnement des notifications :**
- Après chaque import automatique (email ou dossier), une notification est envoyée
- La notification contient :
  * Source de l'import (📧 email ou 📁 dossier)
  * Nom du fichier (email) ou nombre de fichiers (dossier)
  *  Nombre de factures importées
  * 🔄 Nombre de doublons ignorés
  *  Nombre d'erreurs
  * Total de factures détectées
- Les notifications apparaissent dans l'interface Manus du propriétaire

**Interface utilisateur :**
- Chaque carte (Import dossier, Import email) dispose maintenant de :
  * Indicateur visuel "Service actif" avec point vert animé
  * Bouton "Démarrer le service" (désactivé si l'import n'est pas activé)
  * Bouton "Arrêter le service" (rouge) pour stopper la surveillance
  * Messages toast pour confirmer le démarrage/arrêt

**Dépendances ajoutées :**
- chokidar : Surveillance de fichiers et dossiers pour Node.js

**Architecture :**
- folderImportService.ts : Service d'import automatique depuis dossier
- notificationService.ts : Service de notifications centralisé
- Routes tRPC : folderImportService.{start, stop, status}
- Intégration des notifications dans emailImportService et folderImportService

Les deux services d'import automatique (email et dossier) sont maintenant opérationnels avec notifications intégrées.
2026-02-11 05:16:44 -05:00

455 lines
15 KiB
TypeScript

import Imap from "imap";
import { simpleParser, ParsedMail, Attachment } from "mailparser";
import {
getImportSettingsByUser,
createSourceFile,
updateSourceFile,
getUserSettings,
findDuplicateInvoice,
createInvoice,
createImportLog,
} from "./db";
import { extractInvoicesWithMistral, generateMetadataJSON } from "./invoiceExtractor";
import { localStoragePut, generateStorageKey } from "./localStorage";
import { sendImportNotification } from "./notificationService";
interface EmailImportConfig {
userId: number;
emailAddress: string;
password: string;
host: string;
port: number;
}
// Store active intervals for each user
const activeIntervals = new Map<number, NodeJS.Timeout>();
/**
* Process a single email attachment (PDF)
* Replicates the same logic as manual upload
*/
async function processEmailAttachment(
userId: number,
attachment: Attachment,
emailSubject: string
): Promise<{ success: boolean; totalInvoices: number; imported: number; duplicates: number; errors: number }> {
const fileName = attachment.filename || `email-attachment-${Date.now()}.pdf`;
console.log(`[EmailImport] Processing attachment: ${fileName} from email: ${emailSubject}`);
try {
// Convert attachment content to Buffer
const fileBuffer = attachment.content;
console.log(`[EmailImport] File size: ${fileBuffer.length} bytes`);
// Store source file
const sourceFileKey = generateStorageKey(userId, fileName);
console.log(`[EmailImport] Generated storage key: ${sourceFileKey}`);
let sourceFileUrl: string;
try {
const result = await localStoragePut(sourceFileKey, fileBuffer, "application/pdf");
sourceFileUrl = result.url;
console.log(`[EmailImport] File stored successfully at: ${sourceFileUrl}`);
} catch (error) {
console.error(`[EmailImport] FAILED to store file:`, error);
throw new Error("Failed to store PDF file");
}
// Create source file record
const sourceFile = await createSourceFile({
userId,
fileName,
fileKey: sourceFileKey,
fileUrl: sourceFileUrl,
processingStatus: "processing",
});
console.log(`[EmailImport] Source file record created with ID: ${sourceFile.id}`);
// Get user settings for custom keywords
const settings = await getUserSettings(userId);
const customKeywords = settings ? {
invoiceNumber: settings.invoiceNumberKeywords,
deliveryNote: settings.deliveryNoteKeywords,
orderNumber: settings.orderNumberKeywords,
supplier: settings.supplierKeywords,
totalAmount: settings.totalAmountKeywords,
} : undefined;
const model = settings?.llmModel || "mistral-large-latest";
// Extract invoices
console.log(`[EmailImport] Starting invoice extraction...`);
const result = await extractInvoicesWithMistral(
fileBuffer,
userId,
sourceFile.id,
model,
customKeywords
);
console.log(`[EmailImport] Extraction complete: ${result.invoiceCount} invoice(s) detected`);
// Update source file with total count
await updateSourceFile(sourceFile.id, {
totalInvoicesDetected: result.invoiceCount,
processingProgress: `Extraction ${result.invoiceCount} facture(s) détectée(s)`,
});
// Process each invoice
let importedCount = 0;
let duplicatesCount = 0;
let errorsCount = 0;
const duplicateDetails: any[] = [];
const errorDetails: any[] = [];
for (let i = 0; i < result.invoices.length; i++) {
const invoiceData = result.invoices[i]!;
try {
// Update progress
await updateSourceFile(sourceFile.id, {
processingProgress: `Extraction ${i + 1}/${result.invoiceCount} factures...`,
});
// Check for duplicates
const duplicate = await findDuplicateInvoice(
invoiceData.supplierName,
invoiceData.invoiceNumber,
invoiceData.invoiceDate
);
if (duplicate) {
duplicatesCount++;
duplicateDetails.push({
supplierName: invoiceData.supplierName,
invoiceNumber: invoiceData.invoiceNumber,
invoiceDate: invoiceData.invoiceDate,
});
continue;
}
// Generate metadata JSON
const metadataJson = generateMetadataJSON(invoiceData);
const metadataKey = generateStorageKey(userId, `${fileName}-${i + 1}-metadata.json`);
const { url: metadataUrl } = await localStoragePut(
metadataKey,
Buffer.from(metadataJson),
"application/json"
);
// Create invoice record
await createInvoice({
userId,
sourceFileId: sourceFile.id,
invoiceIndexInFile: i + 1,
fileName: `${fileName} - Facture ${i + 1}`,
fileKey: sourceFileKey,
fileUrl: sourceFileUrl,
supplierName: invoiceData.supplierName,
invoiceNumber: invoiceData.invoiceNumber,
invoiceDate: invoiceData.invoiceDate,
deliveryNoteNumber: invoiceData.deliveryNoteNumber,
orderNumber: invoiceData.orderNumber,
totalAmount: invoiceData.totalAmount?.toString(),
pageRange: invoiceData.pageRange,
qualityScore: invoiceData.qualityScore,
metadataFileKey: metadataKey,
metadataFileUrl: metadataUrl,
status: "completed",
});
importedCount++;
} catch (error: any) {
errorsCount++;
errorDetails.push({
invoiceIndex: i + 1,
error: error.message,
});
}
}
// Update source file status
await updateSourceFile(sourceFile.id, {
processingStatus: "completed",
processingProgress: `Terminé: ${importedCount} importée(s), ${duplicatesCount} doublon(s)`,
});
// Create import log
await createImportLog({
userId,
sourceFileId: sourceFile.id,
fileName,
totalInvoicesDetected: result.invoiceCount,
invoicesImported: importedCount,
duplicatesIgnored: duplicatesCount,
errors: errorsCount,
duplicateDetails: duplicateDetails.length > 0 ? JSON.stringify(duplicateDetails) : null,
errorDetails: errorDetails.length > 0 ? JSON.stringify(errorDetails) : null,
});
console.log(`[EmailImport] Successfully processed attachment: ${fileName}`);
console.log(`[EmailImport] Results: ${importedCount} imported, ${duplicatesCount} duplicates, ${errorsCount} errors`);
return {
success: true,
totalInvoices: result.invoiceCount,
imported: importedCount,
duplicates: duplicatesCount,
errors: errorsCount,
};
} catch (error) {
console.error(`[EmailImport] Error processing attachment ${attachment.filename}:`, error);
return {
success: false,
totalInvoices: 0,
imported: 0,
duplicates: 0,
errors: 1,
};
}
}
/**
* Connect to IMAP and process unread emails with PDF attachments
*/
async function checkEmailsForPDFs(config: EmailImportConfig): Promise<void> {
return new Promise((resolve, reject) => {
const imap = new Imap({
user: config.emailAddress,
password: config.password,
host: config.host,
port: config.port,
tls: true,
tlsOptions: { rejectUnauthorized: false },
});
function openInbox(cb: (err: Error | null, box?: any) => void) {
imap.openBox("INBOX", false, cb);
}
imap.once("ready", () => {
console.log(`[EmailImport] Connected to IMAP server for user ${config.userId}`);
openInbox((err, box) => {
if (err) {
console.error("[EmailImport] Error opening inbox:", err);
imap.end();
reject(err);
return;
}
// Search for unread emails
imap.search(["UNSEEN"], (err, results) => {
if (err) {
console.error("[EmailImport] Error searching emails:", err);
imap.end();
reject(err);
return;
}
if (!results || results.length === 0) {
console.log(`[EmailImport] No unread emails found for user ${config.userId}`);
imap.end();
resolve();
return;
}
console.log(`[EmailImport] Found ${results.length} unread emails for user ${config.userId}`);
const fetch = imap.fetch(results, {
bodies: "",
markSeen: false, // Don't mark as seen yet
});
const processedEmails: number[] = [];
fetch.on("message", (msg, seqno) => {
msg.on("body", (stream) => {
simpleParser(stream as any, async (err, parsed: ParsedMail) => {
if (err) {
console.error("[EmailImport] Error parsing email:", err);
return;
}
// Check if email has PDF attachments
const pdfAttachments = parsed.attachments.filter(
(att) =>
att.contentType === "application/pdf" ||
att.filename?.toLowerCase().endsWith(".pdf")
);
if (pdfAttachments.length === 0) {
console.log(`[EmailImport] Email ${seqno} has no PDF attachments, skipping`);
return;
}
console.log(
`[EmailImport] Email ${seqno} has ${pdfAttachments.length} PDF attachment(s)`
);
// Process each PDF attachment
for (const attachment of pdfAttachments) {
try {
const result = await processEmailAttachment(
config.userId,
attachment,
parsed.subject || "No subject"
);
// Mark this email as successfully processed
if (!processedEmails.includes(seqno)) {
processedEmails.push(seqno);
}
// Send notification after successful processing
if (result.success) {
await sendImportNotification(config.userId, {
source: "email",
fileName: attachment.filename || "email-attachment.pdf",
totalInvoices: result.totalInvoices,
imported: result.imported,
duplicates: result.duplicates,
errors: result.errors,
});
}
} catch (error) {
console.error(
`[EmailImport] Failed to process attachment from email ${seqno}:`,
error
);
}
}
});
});
});
fetch.once("error", (err) => {
console.error("[EmailImport] Fetch error:", err);
imap.end();
reject(err);
});
fetch.once("end", () => {
console.log(`[EmailImport] Finished fetching emails for user ${config.userId}`);
// Mark successfully processed emails as seen
if (processedEmails.length > 0) {
imap.addFlags(processedEmails, ["\\Seen"], (err) => {
if (err) {
console.error("[EmailImport] Error marking emails as seen:", err);
} else {
console.log(`[EmailImport] Marked ${processedEmails.length} emails as seen`);
}
imap.end();
resolve();
});
} else {
imap.end();
resolve();
}
});
});
});
});
imap.once("error", (err) => {
console.error("[EmailImport] IMAP connection error:", err);
reject(err);
});
imap.once("end", () => {
console.log(`[EmailImport] IMAP connection ended for user ${config.userId}`);
});
imap.connect();
});
}
/**
* Start email import service for a user
*/
export async function startEmailImportService(userId: number): Promise<boolean> {
try {
// Get user's import settings
const settings = await getImportSettingsByUser(userId);
if (!settings || settings.emailImportEnabled !== 1) {
console.log(`[EmailImport] Email import not enabled for user ${userId}`);
return false;
}
if (!settings.emailImportAddress || !settings.emailImportPassword || !settings.emailImportHost) {
console.log(`[EmailImport] Email import configuration incomplete for user ${userId}`);
return false;
}
// Stop existing service if running
stopEmailImportService(userId);
const config: EmailImportConfig = {
userId,
emailAddress: settings.emailImportAddress,
password: settings.emailImportPassword,
host: settings.emailImportHost,
port: settings.emailImportPort || 993,
};
const frequencyMs = (settings.emailImportFrequency || 30) * 60 * 1000; // Convert minutes to milliseconds
console.log(
`[EmailImport] Starting email import service for user ${userId} with frequency ${settings.emailImportFrequency} minutes`
);
// Run immediately on start
checkEmailsForPDFs(config).catch((error) => {
console.error(`[EmailImport] Error checking emails for user ${userId}:`, error);
});
// Set up interval for periodic checks
const interval = setInterval(() => {
checkEmailsForPDFs(config).catch((error) => {
console.error(`[EmailImport] Error checking emails for user ${userId}:`, error);
});
}, frequencyMs);
activeIntervals.set(userId, interval);
console.log(`[EmailImport] Email import service started for user ${userId}`);
return true;
} catch (error) {
console.error(`[EmailImport] Error starting email import service for user ${userId}:`, error);
return false;
}
}
/**
* Stop email import service for a user
*/
export function stopEmailImportService(userId: number): void {
const interval = activeIntervals.get(userId);
if (interval) {
clearInterval(interval);
activeIntervals.delete(userId);
console.log(`[EmailImport] Email import service stopped for user ${userId}`);
}
}
/**
* Check if email import service is running for a user
*/
export function isEmailImportServiceRunning(userId: number): boolean {
return activeIntervals.has(userId);
}
/**
* Stop all email import services
*/
export function stopAllEmailImportServices(): void {
activeIntervals.forEach((interval, userId) => {
clearInterval(interval);
console.log(`[EmailImport] Stopped email import service for user ${userId}`);
});
activeIntervals.clear();
}