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 D1
Add Architectural Feedback / Question

Planning a Similar Distributed System?

Available for architecture advisory and technical design reviews.

Consult with Khaled →