Checkpoint: Checkpoint saved: Implémentation complète de l'import automatique par email.

Nouvelles fonctionnalités :
 Service d'import automatique par email (emailImportService.ts)
 Connexion IMAP avec support SSL/TLS
 Détection automatique des nouveaux emails non lus avec pièces jointes PDF
 Téléchargement et traitement des pièces jointes (même logique que l'upload manuel)
 Marquage des emails comme lus après traitement réussi
 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

Fonctionnement :
1. L'utilisateur configure ses identifiants email dans "Paramètres de réception"
2. Il active l'import par email et configure la fréquence de vérification
3. Il clique sur "Démarrer le service" pour lancer la surveillance
4. Le service se connecte au serveur IMAP selon la fréquence configurée
5. Il détecte les nouveaux emails non lus avec pièces jointes PDF
6. Il télécharge chaque PDF et applique le même traitement que l'upload manuel :
   - Stockage local du fichier
   - Extraction avec Mistral AI
   - Détection de doublons
   - Création des factures en base de données
   - Génération des métadonnées JSON
   - Création du log d'import
7. Les emails traités sont marqués comme lus
8. Le service continue de tourner en arrière-plan selon la fréquence configurée

Interface utilisateur :
- Indicateur visuel "Service actif" avec point vert animé
- Bouton "Démarrer le service" (désactivé si l'import email 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 :
- imap : Client IMAP pour Node.js
- mailparser : Parser d'emails avec support des pièces jointes
- @types/imap et @types/mailparser : Types TypeScript

Le service est maintenant prêt à être testé avec un vrai compte email IMAP.
This commit is contained in:
Manus
2026-02-11 05:08:45 -05:00
parent 1d79959293
commit ff916d6a79
6 changed files with 796 additions and 1 deletions

View File

@@ -0,0 +1,427 @@
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";
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<void> {
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`);
} catch (error) {
console.error(`[EmailImport] Error processing attachment ${attachment.filename}:`, error);
throw error;
}
}
/**
* 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 {
await processEmailAttachment(
config.userId,
attachment,
parsed.subject || "No subject"
);
// Mark this email as successfully processed
if (!processedEmails.includes(seqno)) {
processedEmails.push(seqno);
}
} 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();
}

View File

@@ -33,6 +33,7 @@ import { loginLocal, hashPassword, getAzureAuthUrl, isAzureAdConfigured, generat
import { extractInvoicesWithMistral, generateMetadataJSON } from "./invoiceExtractor";
import { localStoragePut, generateStorageKey } from "./localStorage";
import { testSftpConnection, exportInvoiceToSftp, getUserSftpConfig } from "./sftpExport";
import { startEmailImportService, stopEmailImportService, isEmailImportServiceRunning } from "./emailImportService";
import { TRPCError } from "@trpc/server";
// Admin-only procedure
@@ -575,6 +576,24 @@ export const appRouter = router({
}),
}),
// ============= EMAIL IMPORT SERVICE ROUTES =============
emailImportService: router({
start: protectedProcedure.mutation(async ({ ctx }) => {
const started = await startEmailImportService(ctx.user.id);
return { success: started };
}),
stop: protectedProcedure.mutation(async ({ ctx }) => {
stopEmailImportService(ctx.user.id);
return { success: true };
}),
status: protectedProcedure.query(async ({ ctx }) => {
const isRunning = isEmailImportServiceRunning(ctx.user.id);
return { isRunning };
}),
}),
// ============= IMPORT SETTINGS ROUTES =============
importSettings: router({
get: protectedProcedure.query(async ({ ctx }) => {