issue-ai/src/app/api/monitor/scan-emails/route.ts

381 lines
12 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// src/app/api/monitor/scan-emails/route.ts
import { NextRequest, NextResponse } from 'next/server'
import { initDatabase } from '@/lib/db-schema'
import { getDb } from '@/lib/db'
import { getCurrentUser } from '@/lib/auth'
import { hasPermission } from '@/lib/permissions'
import { getMonitorConfig } from '@/lib/monitor/settings-manager'
import { ImapFlow } from 'imapflow'
import { simpleParser } from 'mailparser'
import { formatBeijingTime } from '@/lib/monitor/types'
import { writeAuditLog, getClientIP } from '@/lib/audit'
import { getScanState, setScanState, isCancelRequested, resetCancelFlag } from '@/lib/monitor/scan-state'
interface ScanResult {
status: 'running' | 'completed' | 'error' | 'cancelling'
startedAt: string
completedAt?: string
timeRange: { value: number | null; unit: string }
stats: {
total: number
matched: number
imported: number
skipped: number
errors: number
}
details: {
msg_id: string
subject: string
date: string
order_number: string | null
status: 'imported' | 'skipped' | 'error'
ticket_no?: string
error?: string
}[]
detailsTruncated: boolean
error?: string
}
const VALID_UNITS = ['minute', 'hour', 'day', 'week', 'month', 'all']
const MAX_VALUE = 365
const MAX_DETAILS = 500
const SCAN_TIMEOUT_MS = 5 * 60 * 1000 // 5 分钟超时
// 辅助函数:添加 detail 并追踪是否被截断
function addDetail(
scanState: ScanResult,
detail: ScanResult['details'][0],
counter: { total: number }
): void {
counter.total++
if (scanState.details.length < MAX_DETAILS) {
scanState.details.push(detail)
} else {
scanState.detailsTruncated = true
}
}
function extractOrderNumber(subject: string): string | null {
const match = subject.match(/【服务器故障单】(\d+),/)
return match?.[1] ?? null
}
function getDateRange(value: number | null, unit: string): Date {
const now = new Date()
if (!value || unit === 'all') {
// 全部:搜索最近 30 天
now.setDate(now.getDate() - 30)
return now
}
switch (unit) {
case 'minute':
now.setMinutes(now.getMinutes() - value)
break
case 'hour':
now.setHours(now.getHours() - value)
break
case 'day':
now.setDate(now.getDate() - value)
break
case 'week':
now.setDate(now.getDate() - value * 7)
break
case 'month':
now.setMonth(now.getMonth() - value)
break
default:
now.setDate(now.getDate() - 7)
}
return now
}
export async function POST(request: NextRequest) {
initDatabase()
const user = await getCurrentUser()
if (!user || !hasPermission(user, 'monitor:write')) {
return NextResponse.json({ error: '权限不足' }, { status: 403 })
}
const clientIP = getClientIP(request)
// 检查是否有正在运行的扫描
const currentScan = getScanState()
if (currentScan?.status === 'running') {
return NextResponse.json({ error: '扫描正在进行中,请稍后再试' }, { status: 409 })
}
const body = await request.json()
const { value, unit } = body.timeRange || { value: 7, unit: 'day' }
// 输入验证
if (unit !== 'all') {
if (typeof value !== 'number' || value < 1 || value > MAX_VALUE) {
return NextResponse.json({ error: `时间范围值必须在 1-${MAX_VALUE} 之间` }, { status: 400 })
}
}
if (!VALID_UNITS.includes(unit)) {
return NextResponse.json({ error: `无效的时间单位: ${unit}` }, { status: 400 })
}
const config = getMonitorConfig()
const db = getDb()
const since = getDateRange(value, unit)
// 初始化扫描结果
const scanState: ScanResult = {
status: 'running',
startedAt: formatBeijingTime(),
timeRange: { value, unit },
stats: { total: 0, matched: 0, imported: 0, skipped: 0, errors: 0 },
details: [],
detailsTruncated: false,
}
setScanState(scanState)
// 重置取消标志
resetCancelFlag()
// 异步执行扫描
scanEmails(config, db, since, scanState, clientIP, user.id).catch(err => {
// 内层已设置错误状态,此处不再重复设置
console.error('[Scan] Unexpected error:', err)
})
return NextResponse.json({ success: true, message: '扫描已开始' })
}
async function scanEmails(
config: ReturnType<typeof getMonitorConfig>,
db: ReturnType<typeof getDb>,
since: Date,
scanState: ScanResult,
clientIP: string | null,
userId: number
): Promise<void> {
const client = new ImapFlow({
host: config.mail.imap_server,
port: config.mail.imap_port,
secure: true,
auth: { user: config.mail.address, pass: config.mail.password },
logger: false,
connectionTimeout: 15_000,
})
// 整体超时控制
const startTime = Date.now()
const checkTimeout = (): boolean => {
if (Date.now() - startTime > SCAN_TIMEOUT_MS) {
scanState.status = 'error'
scanState.completedAt = formatBeijingTime()
scanState.error = `扫描超时(超过 ${SCAN_TIMEOUT_MS / 60000} 分钟)`
return true
}
return false
}
try {
await client.connect()
const lock = await client.getMailboxLock('INBOX')
// 详情计数器
const detailsCounter = { total: 0 }
try {
// 搜索邮件
const uids = await client.search({ since }, { uid: true })
if (!uids || uids.length === 0) {
scanState.status = 'completed'
scanState.completedAt = formatBeijingTime()
return
}
scanState.stats.total = uids.length
// 处理每封邮件
for (const uid of uids) {
// 检查取消标志
if (isCancelRequested()) {
scanState.status = 'completed'
scanState.completedAt = formatBeijingTime()
scanState.error = '用户取消'
break
}
// 检查超时
if (checkTimeout()) {
break
}
try {
// 使用 body 获取邮件内容source 在某些 IMAP 服务器上可能不返回内容
const msg = await client.fetchOne(uid, { uid: true, bodyStructure: true, body: '.' }, { uid: true })
if (!msg) {
console.log(`[Scan] 邮件 ${uid} 获取失败,跳过`)
continue
}
// 优先使用 source其次使用 body
const rawSource = (msg as any).source || (msg as any).body
if (!rawSource) {
console.log(`[Scan] 邮件 ${uid} 无内容,跳过 (keys: ${Object.keys(msg).join(',')})`)
continue
}
const parsed = await simpleParser(rawSource as Buffer)
const subject = parsed.subject || ''
const date = parsed.date ? formatBeijingTime(parsed.date) : formatBeijingTime()
const msgId = parsed.messageId || String(uid)
// 提取工单号
const orderNumber = extractOrderNumber(subject)
if (!orderNumber) {
scanState.stats.matched++
addDetail(scanState, {
msg_id: msgId,
subject,
date,
order_number: null,
status: 'skipped',
error: '无法提取工单号',
}, detailsCounter)
continue
}
// 检查工单是否已存在tickets 表以 id 作为工单号)
const existing = db.prepare('SELECT id FROM tickets WHERE id = ?').get(parseInt(orderNumber))
if (existing) {
scanState.stats.skipped++
scanState.stats.matched++
addDetail(scanState, {
msg_id: msgId,
subject,
date,
order_number: orderNumber,
status: 'skipped',
ticket_no: orderNumber,
error: '工单已存在',
}, detailsCounter)
continue
}
// 检查是否已处理过
const processed = db.prepare('SELECT msg_id FROM processed_emails WHERE msg_id = ?').get(msgId)
if (processed) {
scanState.stats.skipped++
scanState.stats.matched++
addDetail(scanState, {
msg_id: msgId,
subject,
date,
order_number: orderNumber,
status: 'skipped',
error: '邮件已处理过',
}, detailsCounter)
continue
}
// 创建工单(使用 id 作为工单号)
const ticketId = parseInt(orderNumber)
db.prepare(`
INSERT INTO tickets (id, content, current_status, created_by, created_at, updated_at)
VALUES (?, ?, 'open', ?, datetime('now', '+8 hours'), datetime('now', '+8 hours'))
`).run(ticketId, `从邮件导入:${subject}`, userId)
// 记录已处理邮件
db.prepare("INSERT OR IGNORE INTO processed_emails (msg_id, subject) VALUES (?, ?)").run(msgId, subject)
// 审计日志
writeAuditLog({
userId: userId,
apiKeyId: null,
action: 'import',
entityType: 'ticket',
entityId: ticketId,
details: { created: { ticket_no: orderNumber, source: 'email_scan' } },
ipAddress: clientIP,
})
scanState.stats.imported++
scanState.stats.matched++
addDetail(scanState, {
msg_id: msgId,
subject,
date,
order_number: orderNumber,
status: 'imported',
ticket_no: orderNumber,
}, detailsCounter)
} catch (err) {
scanState.stats.errors++
addDetail(scanState, {
msg_id: String(uid),
subject: '',
date: '',
order_number: null,
status: 'error',
error: err instanceof Error ? err.message : '处理失败',
}, detailsCounter)
}
}
} finally {
lock.release()
}
// 记录截断信息
if (scanState.detailsTruncated) {
console.log(`[Scan] 详情列表已截断: 共 ${detailsCounter.total} 条,仅保存前 ${MAX_DETAILS}`)
}
// 只有状态仍然是 running 时才设置为 completed避免覆盖超时/取消状态)
if (scanState.status === 'running') {
scanState.status = 'completed'
scanState.completedAt = formatBeijingTime()
}
} catch (err) {
scanState.status = 'error'
scanState.completedAt = formatBeijingTime()
scanState.error = err instanceof Error ? err.message : '连接失败'
// 不再 re-throw避免外层 catch 重复处理
} finally {
try { await client.logout() } catch {}
}
// 保存扫描历史到数据库
saveScanHistory(scanState, userId)
}
// 保存扫描历史到数据库
function saveScanHistory(scanState: ScanResult, userId: number): void {
try {
const db = getDb()
db.prepare(`
INSERT INTO scan_history (
status, started_at, completed_at,
time_range_value, time_range_unit,
total_count, matched_count, imported_count, skipped_count, error_count,
details_truncated, error_message, details_json, created_by
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`).run(
scanState.status,
scanState.startedAt,
scanState.completedAt || null,
scanState.timeRange.value,
scanState.timeRange.unit,
scanState.stats.total,
scanState.stats.matched,
scanState.stats.imported,
scanState.stats.skipped,
scanState.stats.errors,
scanState.detailsTruncated ? 1 : 0,
scanState.error || null,
JSON.stringify(scanState.details),
userId
)
console.log(`[Scan] 扫描历史已保存: 状态=${scanState.status}, 导入=${scanState.stats.imported}`)
} catch (err) {
console.error('[Scan] 保存扫描历史失败:', err)
}
}