OmniMail Engine
Real-Time Autonomous Multi-Inbox Sync & Webhook Delivery Platform
Role: Principal Systems Architect•Timeline: 2024 - 2025•Category: Distributed Systems & Messaging
Webhook Ingestion SLA
<300ms
Sub-second real-time delivery
Storage Optimization
4.8 TB
Deduplicated via R2 Content Hash
Worker Cold Starts
0ms
V8 Isolate instant spawn
Inbox Reliability
99.99%
Self-healing connection re-auth
System Specifications
Architecture Pattern
Event-Driven Microservices & Webhook Mesh
Target Throughput
50,000 events/min
Latency Profile
280ms end-to-end delivery
Availability SLA
99.98%
Storage Subsystem
Cloudflare D1 (Indexed Messages) + R2 (MIME Blobs)
Compute Runtime
Cloudflare Workers + Queues
The Engineering Challenge
Synchronizing disparate legacy IMAP mailboxes with modern webhook architectures requires persistent socket pools, memory-efficient MIME stream parsing, and strict resilience against provider rate limits.
The Architectural Solution
Engineered a serverless polling and socket worker harness on Cloudflare Workers and Queues that streams binary MIME payloads directly into Cloudflare R2 while indexing metadata into D1 with transactional safety.
Implementation Details
Built custom streaming MIME parser with chunked SHA-256 deduplication.
Implemented exponential backoff with jitter and automated IMAP session self-healing.
Constructed end-to-end encrypted payload vault using AES-256-GCM before persistent storage.
Delivered webhook retry queue with Dead Letter Queues (DLQ) and WhatsApp operator alerts.
Technologies Used
Node.jsTypeScriptCloudflare QueuesCloudflare D1Cloudflare R2DockerEvolution API
mime-stream-ingest.tstypescript
export async function ingestMimeMessage(rawStream: ReadableStream, env: Env): Promise<IngestResult> {
const hashTransform = new HashStream('sha-256');
const [streamA, streamB] = rawStream.tee();
// Stream to R2 with zero in-memory buffer overhead
const hashPromise = streamA.pipeTo(hashTransform);
const contentHash = await hashPromise;
await env.STORAGE.put(`emails/${contentHash}.eml`, streamB, {
httpMetadata: { contentType: 'message/rfc822' }
});
await env.DB.prepare(
'INSERT INTO indexed_messages (hash, status, received_at) VALUES (?, ?, CURRENT_TIMESTAMP)'
).bind(contentHash, 'indexed').run();
return { success: true, hash: contentHash };
}Technical Discussion (0)
Live on D1Planning a Similar Distributed System?
Available for architecture advisory and technical design reviews.