From b5d8f643209879b8a7955cb6863399d9fedb4155 Mon Sep 17 00:00:00 2001 From: dreamhunter2333 Date: Tue, 25 Aug 2026 21:13:34 +0800 Subject: [PATCH] refactor: centralize mail flag handling --- db/schema.sql | 2 +- worker/src/admin_api/admin_mail_api.ts | 7 +- worker/src/admin_api/db_api.ts | 2 +- worker/src/common.ts | 6 +- worker/src/email/index.ts | 71 +++++++++--------- worker/src/gzip.ts | 15 ++-- worker/src/mail_flags.ts | 96 ++++++++++++++++--------- worker/src/mails_api/mails_crud.ts | 58 ++++++--------- worker/src/mails_api/parsed_mail_api.ts | 9 +-- worker/src/models/index.ts | 1 + worker/src/user_api/user_mail_api.ts | 70 ++++++------------ worker/src/utils.ts | 63 +++++++--------- 12 files changed, 194 insertions(+), 206 deletions(-) diff --git a/db/schema.sql b/db/schema.sql index b7e7d9c..a9a8523 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -6,7 +6,7 @@ CREATE TABLE IF NOT EXISTS raw_mails ( raw TEXT, raw_blob BLOB, metadata TEXT, - flags INTEGER NOT NULL DEFAULT 0, + flags INTEGER, created_at DATETIME DEFAULT CURRENT_TIMESTAMP ); diff --git a/worker/src/admin_api/admin_mail_api.ts b/worker/src/admin_api/admin_mail_api.ts index 98e6645..3868dc9 100644 --- a/worker/src/admin_api/admin_mail_api.ts +++ b/worker/src/admin_api/admin_mail_api.ts @@ -1,7 +1,6 @@ import { Context } from "hono"; import { handleMailListQuery } from "../common"; import { resolveRawEmailRow } from "../gzip"; -import { serializeMailState } from "../mail_flags"; import { getBooleanValue } from "../utils"; export default { @@ -33,8 +32,10 @@ export default { `SELECT * FROM raw_mails WHERE id = ?` ).bind(id).first(); if (!result) return c.json(null); - const resolved = await resolveRawEmailRow(result); - return c.json(serializeMailState(resolved, getBooleanValue(c.env.ENABLE_MAIL_FLAGS))); + return c.json(await resolveRawEmailRow( + result, + getBooleanValue(c.env.ENABLE_MAIL_FLAGS), + )); }, deleteMail: async (c: Context) => { const { id } = c.req.param(); diff --git a/worker/src/admin_api/db_api.ts b/worker/src/admin_api/db_api.ts index 38ff3cb..38e09fe 100644 --- a/worker/src/admin_api/db_api.ts +++ b/worker/src/admin_api/db_api.ts @@ -11,7 +11,7 @@ CREATE TABLE IF NOT EXISTS raw_mails ( raw TEXT, raw_blob BLOB, metadata TEXT, - flags INTEGER NOT NULL DEFAULT 0, + flags INTEGER, created_at DATETIME DEFAULT CURRENT_TIMESTAMP ); diff --git a/worker/src/common.ts b/worker/src/common.ts index 744b2f4..e61995a 100644 --- a/worker/src/common.ts +++ b/worker/src/common.ts @@ -7,7 +7,6 @@ import { unbindTelegramByAddress } from './telegram_api/common'; import { CONSTANTS } from './constants'; import { AddressCreationSettings, AdminWebhookSettings, ExtractResult, WebhookMail, WebhookSettings } from './models'; import i18n from './i18n'; -import { serializeMailState } from './mail_flags'; const DEFAULT_NAME_REGEX = /[^a-z0-9]/g; const DEFAULT_RANDOM_SUBDOMAIN_LENGTH = 8; @@ -721,8 +720,9 @@ export const handleMailListQuery = async ( const { results } = await c.env.DB.prepare(resultsQuery).bind( ...params, limit, offset ).all(); - const resolvedResults = (await resolveRawEmailList(results)).map(row => - serializeMailState(row, getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) + const resolvedResults = await resolveRawEmailList( + results, + getBooleanValue(c.env.ENABLE_MAIL_FLAGS), ); const count = offset == 0 ? await c.env.DB.prepare( countQuery diff --git a/worker/src/email/index.ts b/worker/src/email/index.ts index e3b4156..2145250 100644 --- a/worker/src/email/index.ts +++ b/worker/src/email/index.ts @@ -12,7 +12,7 @@ import { forwardEmail } from "./forward"; import { EmailRuleSettings } from "../models"; import { CONSTANTS } from "../constants"; import { compressText } from "../gzip"; -import { insertRawMail, resolveInitialMailFlags } from "../mail_flags"; +import { updateInitialMailFlags } from "../mail_flags"; async function email(message: ForwardableEmailMessage, env: Bindings, ctx: ExecutionContext) { @@ -68,10 +68,8 @@ async function email(message: ForwardableEmailMessage, env: Bindings, ctx: Execu const message_id = message.headers.get("Message-ID"); // save email try { - const initialFlags = await resolveInitialMailFlags( - getBooleanValue(env.ENABLE_MAIL_FLAGS), env, toAddress, parsedEmailContext - ); let success = false; + let insertResult: D1Result | null = null; if (getBooleanValue(env.ENABLE_MAIL_GZIP)) { let compressed: ArrayBuffer | null = null; try { @@ -81,54 +79,55 @@ async function email(message: ForwardableEmailMessage, env: Bindings, ctx: Execu } if (compressed) { try { - ({ success } = await insertRawMail(env.DB, { - source: message.from, - address: toAddress, - content: compressed, - contentColumn: 'raw_blob', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw_blob, message_id) VALUES (?, ?, ?, ?)` + ).bind( + message.from, toAddress, compressed, message_id + ).run(); + ({ success } = insertResult); } catch (dbError) { // Fallback to plaintext only if raw_blob column is missing (migration not applied) const errMsg = String(dbError); if (errMsg.includes('raw_blob') || errMsg.includes('no such column')) { console.error("raw_blob column missing, falling back to plaintext", dbError); - ({ success } = await insertRawMail(env.DB, { - source: message.from, - address: toAddress, - content: parsedEmailContext.rawEmail, - contentColumn: 'raw', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw, message_id) VALUES (?, ?, ?, ?)` + ).bind( + message.from, toAddress, parsedEmailContext.rawEmail, message_id + ).run(); + ({ success } = insertResult); } else { throw dbError; } } } else { - ({ success } = await insertRawMail(env.DB, { - source: message.from, - address: toAddress, - content: parsedEmailContext.rawEmail, - contentColumn: 'raw', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw, message_id) VALUES (?, ?, ?, ?)` + ).bind( + message.from, toAddress, parsedEmailContext.rawEmail, message_id + ).run(); + ({ success } = insertResult); } } else { - ({ success } = await insertRawMail(env.DB, { - source: message.from, - address: toAddress, - content: parsedEmailContext.rawEmail, - contentColumn: 'raw', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw, message_id) VALUES (?, ?, ?, ?)` + ).bind( + message.from, toAddress, parsedEmailContext.rawEmail, message_id + ).run(); + ({ success } = insertResult); } if (!success) { message.setReject(`Failed save message to ${toAddress}`); console.error(`Failed save message from ${message.from} to ${toAddress}`); + } else { + await updateInitialMailFlags( + env.DB, + getBooleanValue(env.ENABLE_MAIL_FLAGS), + insertResult?.meta.last_row_id ?? 0, + env, + toAddress, + parsedEmailContext, + ); } } catch (error) { diff --git a/worker/src/gzip.ts b/worker/src/gzip.ts index 8d5e993..bd60e8b 100644 --- a/worker/src/gzip.ts +++ b/worker/src/gzip.ts @@ -4,6 +4,7 @@ */ import { RawMailRow } from "./models"; +import { serializeMailState } from "./mail_flags"; export async function compressText(text: string): Promise { const stream = new Blob([text]).stream().pipeThrough(new CompressionStream('gzip')); @@ -34,15 +35,21 @@ export async function resolveRawEmail(row: RawMailRow): Promise { /** * Resolve a single row: decompress raw_blob if present, strip raw_blob from result. */ -export async function resolveRawEmailRow(row: RawMailRow): Promise { +export async function resolveRawEmailRow( + row: RawMailRow, + enableReadStatus = false, +): Promise { const raw = await resolveRawEmail(row); const { raw_blob: _, ...rest } = row; - return { ...rest, raw }; + return serializeMailState({ ...rest, raw }, enableReadStatus); } /** * Batch resolve raw emails for list queries using Promise.all. */ -export async function resolveRawEmailList(rows: RawMailRow[]): Promise { - return Promise.all(rows.map(row => resolveRawEmailRow(row))); +export async function resolveRawEmailList( + rows: RawMailRow[], + enableReadStatus = false, +): Promise { + return Promise.all(rows.map(row => resolveRawEmailRow(row, enableReadStatus))); } diff --git a/worker/src/mail_flags.ts b/worker/src/mail_flags.ts index 653df9f..a31666a 100644 --- a/worker/src/mail_flags.ts +++ b/worker/src/mail_flags.ts @@ -33,38 +33,25 @@ export const serializeMailState = >( return result; }; -export const resolveInitialMailFlags = async ( - enabled: boolean, +const resolveInitialMailFlags = async ( _env: Bindings, _address: string, _parsedEmailContext: ParsedEmailContext, -): Promise => { - if (!enabled) return null; +): Promise => { return MAIL_FLAGS.UNREAD; }; -type InsertRawMailParams = { - source: string; - address: string; - content: string | ArrayBuffer; - contentColumn: 'raw' | 'raw_blob'; - messageId: string | null; - flags: number | null; -}; - -export const insertRawMail = async ( +export const updateInitialMailFlags = async ( db: D1Database, - params: InsertRawMailParams, + enabled: boolean, + mailId: number, + env: Bindings, + address: string, + parsedEmailContext: ParsedEmailContext, ) => { - const { source, address, content, contentColumn, messageId, flags } = params; - if (flags === null) { - return db.prepare( - `INSERT INTO raw_mails (source, address, ${contentColumn}, message_id) VALUES (?, ?, ?, ?)` - ).bind(source, address, content, messageId).run(); - } - return db.prepare( - `INSERT INTO raw_mails (source, address, ${contentColumn}, message_id, flags) VALUES (?, ?, ?, ?, ?)` - ).bind(source, address, content, messageId, flags).run(); + if (!enabled || !Number.isInteger(mailId) || mailId <= 0) return; + const flags = await resolveInitialMailFlags(env, address, parsedEmailContext); + await db.prepare(`UPDATE raw_mails SET flags = ? WHERE id = ?`).bind(flags, mailId).run(); }; export type MailReadStatusUpdate = { @@ -73,21 +60,26 @@ export type MailReadStatusUpdate = { action: MailReadStatusAction; }; -export type MailReadStatusFilter = { - mask: number; - state: 'set' | 'unset'; +export type MailReadStatusQuery = { + clause: string; + params: string[]; }; -export const parseReadStatusFilter = ( +export const getReadStatusQuery = ( value: string | undefined, -): MailReadStatusFilter | undefined | null => { + column: 'flags' | 'rm.flags', +): MailReadStatusQuery | undefined | null => { if (value === undefined || value === 'all') return undefined; - if (value === 'unread') return { mask: MAIL_FLAGS.UNREAD, state: 'set' }; - if (value === 'read') return { mask: MAIL_FLAGS.UNREAD, state: 'unset' }; + if (value === 'unread') { + return { clause: `(COALESCE(${column}, 0) & ?) != 0`, params: [String(MAIL_FLAGS.UNREAD)] }; + } + if (value === 'read') { + return { clause: `(COALESCE(${column}, 0) & ?) = 0`, params: [String(MAIL_FLAGS.UNREAD)] }; + } return null; }; -export const parseMailReadStatusUpdate = (value: unknown): MailReadStatusUpdate | null => { +const parseMailReadStatusUpdate = (value: unknown): MailReadStatusUpdate | null => { if (!value || typeof value !== 'object') return null; const body = value as Record; if (!Array.isArray(body.ids) || body.ids.length === 0 || body.ids.length > 100) return null; @@ -101,7 +93,7 @@ export const parseMailReadStatusUpdate = (value: unknown): MailReadStatusUpdate return { ids, mask: MAIL_FLAGS.UNREAD, action: body.action }; }; -export const getMailReadStatusUpdateExpression = ( +const getMailReadStatusUpdateExpression = ( update: MailReadStatusUpdate, column = 'flags', ): { expression: string; params: number[]; condition?: string; conditionParams?: number[] } => { @@ -126,3 +118,41 @@ export const getMailReadStatusUpdateExpression = ( params: [update.mask, update.mask], }; }; + +type MailScope = { + clause: string; + params: (string | number)[]; +}; + +export const applyMailReadStatusUpdate = async ( + db: D1Database, + scope: MailScope, + value: unknown, +) => { + const update = parseMailReadStatusUpdate(value); + if (!update) return null; + + const placeholders = update.ids.map(() => '?').join(','); + const statusUpdate = getMailReadStatusUpdateExpression(update); + const condition = statusUpdate.condition ? ` AND ${statusUpdate.condition}` : ''; + const result = await db.prepare( + `UPDATE raw_mails SET flags = ${statusUpdate.expression}` + + ` WHERE id IN (${placeholders}) AND (${scope.clause})${condition}` + ).bind( + ...statusUpdate.params, + ...update.ids, + ...scope.params, + ...(statusUpdate.conditionParams ?? []), + ).run(); + if (!result.success) return { success: false, changes: 0, results: [] }; + + const { results } = await db.prepare( + `SELECT id, flags FROM raw_mails` + + ` WHERE id IN (${placeholders}) AND (${scope.clause})` + ).bind(...update.ids, ...scope.params).all(); + return { + success: true, + changes: result.meta.changes ?? 0, + results: results.map(row => serializeMailState(row, true)), + }; +}; diff --git a/worker/src/mails_api/mails_crud.ts b/worker/src/mails_api/mails_crud.ts index f350cca..f2eb07e 100644 --- a/worker/src/mails_api/mails_crud.ts +++ b/worker/src/mails_api/mails_crud.ts @@ -6,10 +6,8 @@ import { handleMailListQuery, deleteAddressWithData, updateAddressUpdatedAt } fr import { resolveRawEmailRow } from '../gzip' import { getSendBalanceState } from './send_balance'; import { - getMailReadStatusUpdateExpression, - parseMailReadStatusUpdate, - parseReadStatusFilter, - serializeMailState, + getReadStatusQuery, + applyMailReadStatusUpdate, } from '../mail_flags'; const listMails = async (c: Context) => { @@ -19,17 +17,17 @@ const listMails = async (c: Context) => { } const { limit, offset, read_status } = c.req.query(); if (Number.parseInt(offset) <= 0) updateAddressUpdatedAt(c, address); - const readStatusFilter = parseReadStatusFilter(read_status); - if (readStatusFilter === null) return c.json({ error: "Invalid mail read status filter" }, 400); - if (readStatusFilter && !getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) { + const readStatusQuery = getReadStatusQuery(read_status, 'flags'); + if (readStatusQuery === null) return c.json({ error: "Invalid mail read status filter" }, 400); + if (readStatusQuery && !getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) { return c.json({ error: "Mail read status is disabled" }, 403); } const filters = [`address = ?`]; const params = [address]; - if (readStatusFilter) { - filters.push(`(COALESCE(flags, 0) & ?) ${readStatusFilter.state === 'set' ? '!=' : '='} 0`); - params.push(String(readStatusFilter.mask)); + if (readStatusQuery) { + filters.push(readStatusQuery.clause); + params.push(...readStatusQuery.params); } const whereClause = filters.join(' AND '); return await handleMailListQuery(c, @@ -46,8 +44,10 @@ const getMail = async (c: Context) => { `SELECT * FROM raw_mails where id = ? and address = ?` ).bind(mail_id, address).first(); if (!result) return c.json(null); - const resolved = await resolveRawEmailRow(result); - return c.json(serializeMailState(resolved, getBooleanValue(c.env.ENABLE_MAIL_FLAGS))); + return c.json(await resolveRawEmailRow( + result, + getBooleanValue(c.env.ENABLE_MAIL_FLAGS), + )); }; const deleteMail = async (c: Context) => { @@ -68,33 +68,15 @@ const updateMailReadStatus = async (c: Context) => { if (!getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) { return c.json({ error: "Mail read status is disabled" }, 403); } - const update = parseMailReadStatusUpdate(await c.req.json().catch(() => null)); - if (!update) return c.json({ error: "Invalid mail read status request" }, 400); - const { address } = c.get("jwtPayload"); - const placeholders = update.ids.map(() => '?').join(','); - const statusUpdate = getMailReadStatusUpdateExpression(update); - const condition = statusUpdate.condition ? ` AND ${statusUpdate.condition}` : ''; - const result = await c.env.DB.prepare( - `UPDATE raw_mails` - + ` SET flags = ${statusUpdate.expression}` - + ` WHERE address = ? AND id IN (${placeholders})${condition}` - ).bind( - ...statusUpdate.params, - address, - ...update.ids, - ...(statusUpdate.conditionParams ?? []), - ).run(); - if (!result.success) return c.json({ success: false, changes: 0, results: [] }, 500); - - const { results } = await c.env.DB.prepare( - `SELECT id, flags FROM raw_mails WHERE address = ? AND id IN (${placeholders})` - ).bind(address, ...update.ids).all(); - return c.json({ - success: true, - changes: result.meta.changes ?? 0, - results: results.map(row => serializeMailState(row, true)), - }); + const result = await applyMailReadStatusUpdate( + c.env.DB, + { clause: 'address = ?', params: [address] }, + await c.req.json().catch(() => null), + ); + if (!result) return c.json({ error: "Invalid mail read status request" }, 400); + if (!result.success) return c.json(result, 500); + return c.json(result); }; const getSettings = async (c: Context) => { diff --git a/worker/src/mails_api/parsed_mail_api.ts b/worker/src/mails_api/parsed_mail_api.ts index eb03d17..e7d30a6 100644 --- a/worker/src/mails_api/parsed_mail_api.ts +++ b/worker/src/mails_api/parsed_mail_api.ts @@ -2,7 +2,6 @@ import { Context } from 'hono' import { commonParseMail, handleMailListQuery, updateAddressUpdatedAt } from '../common' import { resolveRawEmailRow } from '../gzip' -import { serializeMailState } from '../mail_flags'; import { getBooleanValue } from '../utils'; const toParsedMailRow = async (row: Record): Promise> => { @@ -47,9 +46,11 @@ const getParsedMail = async (c: Context) => { `SELECT * FROM raw_mails where id = ? and address = ?` ).bind(mail_id, address).first(); if (!row) return c.json(null); - const resolved = await resolveRawEmailRow(row); - const serialized = serializeMailState(resolved, getBooleanValue(c.env.ENABLE_MAIL_FLAGS)); - return c.json(await toParsedMailRow(serialized)); + const resolved = await resolveRawEmailRow( + row, + getBooleanValue(c.env.ENABLE_MAIL_FLAGS), + ); + return c.json(await toParsedMailRow(resolved)); }; export default { listParsedMails, getParsedMail }; diff --git a/worker/src/models/index.ts b/worker/src/models/index.ts index 463c499..e951db6 100644 --- a/worker/src/models/index.ts +++ b/worker/src/models/index.ts @@ -214,6 +214,7 @@ export type RawMailRow = { raw_blob?: unknown; metadata?: string; flags?: number | null; + unread?: boolean; created_at?: string; } diff --git a/worker/src/user_api/user_mail_api.ts b/worker/src/user_api/user_mail_api.ts index b8598fe..c1dde10 100644 --- a/worker/src/user_api/user_mail_api.ts +++ b/worker/src/user_api/user_mail_api.ts @@ -3,10 +3,8 @@ import i18n from "../i18n"; import { handleMailListQuery } from "../common"; import { getBooleanValue } from "../utils"; import { - getMailReadStatusUpdateExpression, - parseMailReadStatusUpdate, - parseReadStatusFilter, - serializeMailState, + getReadStatusQuery, + applyMailReadStatusUpdate, } from "../mail_flags"; export default { @@ -19,14 +17,14 @@ export default { filterQuerys.push(`rm.address = ?`); filterParams.push(address); } - const readStatusFilter = parseReadStatusFilter(read_status); - if (readStatusFilter === null) return c.json({ error: "Invalid mail read status filter" }, 400); - if (readStatusFilter && !getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) { + const readStatusQuery = getReadStatusQuery(read_status, 'rm.flags'); + if (readStatusQuery === null) return c.json({ error: "Invalid mail read status filter" }, 400); + if (readStatusQuery && !getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) { return c.json({ error: "Mail read status is disabled" }, 403); } - if (readStatusFilter) { - filterQuerys.push(`(COALESCE(rm.flags, 0) & ?) ${readStatusFilter.state === 'set' ? '!=' : '='} 0`); - filterParams.push(String(readStatusFilter.mask)); + if (readStatusQuery) { + filterQuerys.push(readStatusQuery.clause); + filterParams.push(...readStatusQuery.params); } const fromQuery = ` FROM users_address ua` + ` JOIN address a ON a.id = ua.address_id` @@ -61,43 +59,21 @@ export default { if (!getBooleanValue(c.env.ENABLE_MAIL_FLAGS)) { return c.json({ error: "Mail read status is disabled" }, 403); } - const update = parseMailReadStatusUpdate(await c.req.json().catch(() => null)); - if (!update) return c.json({ error: "Invalid mail read status request" }, 400); - const { user_id } = c.get("userPayload"); - const placeholders = update.ids.map(() => '?').join(','); - const statusUpdate = getMailReadStatusUpdateExpression(update); - const condition = statusUpdate.condition ? ` AND ${statusUpdate.condition}` : ''; - const result = await c.env.DB.prepare( - `UPDATE raw_mails` - + ` SET flags = ${statusUpdate.expression}` - + ` WHERE id IN (${placeholders})` - + ` AND EXISTS (` - + `SELECT 1 FROM users_address ua` - + ` JOIN address a ON a.id = ua.address_id` - + ` WHERE ua.user_id = ? AND a.name = raw_mails.address` - + `)${condition}` - ).bind( - ...statusUpdate.params, - ...update.ids, - user_id, - ...(statusUpdate.conditionParams ?? []), - ).run(); - if (!result.success) return c.json({ success: false, changes: 0, results: [] }, 500); - - const { results } = await c.env.DB.prepare( - `SELECT id, flags FROM raw_mails` - + ` WHERE id IN (${placeholders})` - + ` AND EXISTS (` - + `SELECT 1 FROM users_address ua` - + ` JOIN address a ON a.id = ua.address_id` - + ` WHERE ua.user_id = ? AND a.name = raw_mails.address` - + `)` - ).bind(...update.ids, user_id).all(); - return c.json({ - success: true, - changes: result.meta.changes ?? 0, - results: results.map(row => serializeMailState(row, true)), - }); + const result = await applyMailReadStatusUpdate( + c.env.DB, + { + clause: `EXISTS (` + + `SELECT 1 FROM users_address ua` + + ` JOIN address a ON a.id = ua.address_id` + + ` WHERE ua.user_id = ? AND a.name = raw_mails.address` + + `)`, + params: [user_id], + }, + await c.req.json().catch(() => null), + ); + if (!result) return c.json({ error: "Invalid mail read status request" }, 400); + if (!result.success) return c.json(result, 500); + return c.json(result); } } diff --git a/worker/src/utils.ts b/worker/src/utils.ts index a3444fd..2106657 100644 --- a/worker/src/utils.ts +++ b/worker/src/utils.ts @@ -3,7 +3,7 @@ import { createMimeMessage } from "mimetext"; import { UserSettings, RoleAddressConfig } from "./models"; import { CONSTANTS } from "./constants"; import { compressText } from "./gzip"; -import { insertRawMail, resolveInitialMailFlags } from "./mail_flags"; +import { updateInitialMailFlags } from "./mail_flags"; export const getJsonObjectValue = ( value: string | any @@ -373,10 +373,8 @@ export const sendAdminInternalMail = async ( const message_id = Math.random().toString(36).substring(2, 15); const rawText = msg.asRaw(); const parsedEmailContext: ParsedEmailContext = { rawEmail: rawText }; - const initialFlags = await resolveInitialMailFlags( - getBooleanValue(c.env.ENABLE_MAIL_FLAGS), c.env, toMail, parsedEmailContext - ); let success = false; + let insertResult: D1Result | null = null; if (getBooleanValue(c.env.ENABLE_MAIL_GZIP)) { let compressed: ArrayBuffer | null = null; try { @@ -386,52 +384,45 @@ export const sendAdminInternalMail = async ( } if (compressed) { try { - ({ success } = await insertRawMail(c.env.DB, { - source: "admin@internal", - address: toMail, - content: compressed, - contentColumn: 'raw_blob', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await c.env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw_blob, message_id) VALUES (?, ?, ?, ?)` + ).bind("admin@internal", toMail, compressed, message_id).run(); + ({ success } = insertResult); } catch (dbError) { const errMsg = String(dbError); if (errMsg.includes('raw_blob') || errMsg.includes('no such column')) { console.error("raw_blob column missing, falling back to plaintext", dbError); - ({ success } = await insertRawMail(c.env.DB, { - source: "admin@internal", - address: toMail, - content: rawText, - contentColumn: 'raw', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await c.env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw, message_id) VALUES (?, ?, ?, ?)` + ).bind("admin@internal", toMail, rawText, message_id).run(); + ({ success } = insertResult); } else { throw dbError; } } } else { - ({ success } = await insertRawMail(c.env.DB, { - source: "admin@internal", - address: toMail, - content: rawText, - contentColumn: 'raw', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await c.env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw, message_id) VALUES (?, ?, ?, ?)` + ).bind("admin@internal", toMail, rawText, message_id).run(); + ({ success } = insertResult); } } else { - ({ success } = await insertRawMail(c.env.DB, { - source: "admin@internal", - address: toMail, - content: rawText, - contentColumn: 'raw', - messageId: message_id, - flags: initialFlags, - })); + insertResult = await c.env.DB.prepare( + `INSERT INTO raw_mails (source, address, raw, message_id) VALUES (?, ?, ?, ?)` + ).bind("admin@internal", toMail, rawText, message_id).run(); + ({ success } = insertResult); } if (!success) { console.log(`Failed save message from admin@internal to ${toMail}`); + } else { + await updateInitialMailFlags( + c.env.DB, + getBooleanValue(c.env.ENABLE_MAIL_FLAGS), + insertResult?.meta.last_row_id ?? 0, + c.env, + toMail, + parsedEmailContext, + ); } return success; } catch (error) {