System Design Problem

Design an Online Judge System (LeetCode)

Commonly Asked By:LeetCodeHackerRankGoogleMeta

Interview Setup

Interview Prompt

Design an Online Judge (LeetCode) for 1M DAU with ~1,600 submissions/sec peak during contests, ~16K problem reads/sec, and strict sandbox isolation for 20+ programming languages.

Clarifying Questions (ask before designing)

QuestionWhy it matters
What is the peak load?Contests drive massive traffic spikes in short windows. For example, 50k submissions in 5 minutes is about 167 submissions/sec, while this design reserves capacity for a 1,600 submissions/sec contest peak.
What programming languages must we support?Different languages need different compilation and sandboxing strategies (e.g., JVM vs C++).
Is this for real time interviews or asynchronous contests?Interactive interviews may require sub-second feedback, whereas contests can tolerate queued evaluation as long as the execution queue stays within its SLO.

Scope

In scope

  • Code submission and queuing workflow
  • Secure code execution environment (Sandboxing)
  • Result aggregation and ranking
  • Handling massive traffic spikes during contests

Out of scope (state explicitly)

  • Discussion forums and comments
  • Payment processing for premium users
  • IDE UI implementation

Functional Requirements

Start by asking your interviewer which programming languages must be supported and whether real time contest spikes or steady practice mode traffic dominates. Confirm whether result streaming via SSE is required before designing the sandbox fleet.

  • View Problems: Users can browse and view coding problems, including descriptions, constraints, and public test cases.
  • Code Submission: Users can write and submit code in multiple programming languages (e.g., Python, Java, C++).
  • Code Execution: The system compiles when required and evaluates submitted code against hidden test cases.
  • Results Reporting: Return execution results such as Accepted, Wrong Answer, Time Limit Exceeded (TLE), Memory Limit Exceeded (MLE), Runtime Error, and System Error.
  • Leaderboard: Rank users based on problems solved and competition ratings, with contest identity preserved on each submission.
  • Reproducible Evaluation: Bind every queued submission to immutable problem, test suite, and language runtime versions, and reject idempotency key reuse with a different normalized request.

Non-Functional Requirements

Your interviewer will stress test sandbox isolation because every submission is malicious until proven otherwise. They will also test whether the system survives 100x contest spikes without polling storms overwhelming Redis. Lead with Firecracker or gVisor, push delivery, and then leaderboard scaling.

  • High Availability: The platform must be highly available for browsing and submissions.
  • Scalability: Must handle sudden spikes during contests (e.g., 100x normal load).
  • Security & Isolation: User-submitted code must be executed in a strictly isolated environment to prevent malicious actions (e.g., infinite loops, network access, file system tampering).
  • Low Latency Execution: Sandbox startup overhead should be minimal. The normal execution target after a worker begins the submission is < 2-3 seconds, while queued contest work remains bounded by the separate queue SLO.
  • Fairness: Resource allocation for CPU, RAM, process count, disk, and runtime must be strictly enforced so executions have bounded and comparable resource budgets. Exact wall clock time can still vary with shared host scheduling, so the judge enforces explicit resource limits rather than assuming perfect physical determinism.

Capacity Estimations

Size the sandbox fleet for peak contest load rather than average. A peak rate of 1,600 submissions/sec with a 2-second runtime requires approximately 3,200 concurrent isolated environments. Kafka absorbs bursts while workers autoscale on consumer lag and available execution capacity.

MetricCalculationValue
Daily Active Users (DAU)Given1M
Peak concurrent usersDAU x 10%100K
Submissions / sec (peak)100K x (1 submission / min ÷ 60)~1,600 submissions/s
Problem reads / sec10x submissions~16,000 reads/s
Storage per submissionMetadata + S3 pointer~2 KB
Storage per year (subs)1M DAU x 2 subs/day x 365 x 2KB~1.5 TB/year
Concurrent sandbox slots1,600 submissions/s x 2s execution~3,200 concurrent slots at 1 CPU core/slot

Compute Constraints

If a peak submission spike during a contest reaches 1,600 submissions/sec, and each submission takes ~2 seconds to compile and run against 50 test cases, we need 1600 x 2 = 3200 concurrent sandbox environments available. This requires a horizontally scalable fleet of execution workers.

Architecture Diagram

An online judge splits cleanly into a read-heavy web tier for problem browsing and leaderboards and a isolated execution pipeline for writes for untrusted code. The API never runs user submissions synchronously. It validates the request, binds the submission to immutable problem and test suite versions plus a pinned runtime image digest, stores source code privately in S3, commits the submission record and a dispatch outbox row in one PostgreSQL transaction, and returns a 202 Accepted response. A relay publishes the durable execution metadata from the outbox to Kafka, while a horizontally scaled sandbox fleet compiles, executes, and streams results back via SSE.

Contest spikes can push submission rate 100x above steady state, so Kafka absorbs bursts while autoscaling workers drain the queue. S3 holds raw source code and massive hidden test-case files. Kafka carries only execution metadata, submission identifiers, immutable test suite versions, and S3 pointers rather than multi megabyte payloads.

The read path is independently scaled for the ~16K problem reads per second. A CDN serves cacheable problem statements and static assets, Redis can cache hot problem metadata, and PostgreSQL read replicas or a dedicated read pool serve cache misses. Hidden test data never enters this client-facing path.

Loading...

In the room

Lead with the trust boundary: all user code is malicious. Never run submissions synchronously on the API thread, because a single infinite loop would take down the entire web tier. Mention decoupling via Kafka before interviewers probe contest spikes.

Component Deep Dives

The execution path is the security boundary. Every submission is treated as malicious until proven otherwise. Workers pull from Kafka, provision an isolated sandbox, run all test cases in a single sandbox session, and push incremental status to clients through SSE rather than creating polling storms.

1. Submission Queue & Dispatcher (Kafka & Worker Fleet)

The API must never block on code execution. Kafka absorbs contest bursts while a horizontally scaled worker fleet drains the queue at its own pace. See Message Queues Fundamentals.

Synchronous execution is impossible at scale because a slow Python script could occupy an API thread for 10 seconds. The API Gateway therefore commits the submission and a dispatch outbox record transactionally, then a relay publishes execution metadata to a Kafka Topic keyed by submission_id for even distribution. The full source code stays in private S3 storage. A separate per-user rate limiter and fair scheduling policy prevent one participant from monopolizing execution capacity.

  • Why Kafka? It provides extreme throughput and durability. If the execution fleet crashes during a contest, Kafka retains the submissions until workers recover.
  • Dispatcher/Consumer Group: A pool of Go/Rust worker nodes listens to the Kafka topic. Each worker pulls a batch of code submissions, provisions a sandbox, and streams the standard output and execution metrics back to a Result Aggregator.
  • Backpressure: If Kafka lag grows during a coding contest, an Auto-Scaler provisions additional execution capacity. Keep a baseline of on demand workers for predictable capacity and use EC2 Spot Instances for burst capacity. Handle spot interruption signals by draining workers and allowing uncommitted Kafka work to be retried.

2. Execution Engine & Secure Sandboxing

Secure sandboxing is the core security challenge, where cgroups, seccomp, and MicroVM isolation separate a production judge from a simple Docker demo.

Loading...

Executing untrusted code is the most critical vulnerability vector. A user could write code to mine crypto, read host environment variables, or launch a DDoS attack.

  • Container Isolation (Docker/containerd): Each submission runs in an ephemeral container. The filesystem is entirely read-only except for a tiny /tmp scratch space.
  • Resource Quotas (cgroups v2): We strictly enforce CPU (e.g., 1 core max), memory (e.g., 256MB max), PID count, file descriptors, and I/O. The trusted supervisor also enforces output byte limits because cgroups alone do not provide a semantic stdout or stderr quota. If memory exceeds the configured cgroup limit, the process is terminated and the supervisor maps the outcome to MLE (Memory Limit Exceeded).
  • System Call Filtering (seccomp/AppArmor): Drop unnecessary Linux capabilities and apply language specific syscall profiles. Network access is disabled through namespace policy and network isolation. Process creation is bounded with PID limits and restricted with seccomp rather than relying on a blanket fork() or execve() denial that could break the trusted runtime launcher.
  • MicroVMs and hardened sandboxes: Standard Docker shares the host kernel. Firecracker provides hardware virtualized MicroVM isolation, while gVisor interposes a user space kernel to reduce the host kernel attack surface. A representative startup target can be <125ms after prewarming, but actual startup depends on image size, host load, and the fleet configuration.

3. Test Case Execution Strategy

Boot one sandbox per submission and batch all test cases inside it, because 150 sandbox boots per submission would make contest latency unacceptable.

The capacity baseline assumes about 50 test cases per submission, while a representative problem may have as many as 150 hidden test cases. Booting 150 isolated containers per submission would be disastrously slow.

  • In-Memory Injection: We boot one sandbox per submission. Inside it, a trusted runner uses an immutable language toolchain image to compile the pinned source and then invokes the resulting program or function without exposing host credentials, host mounts, or writable toolchain paths.
  • Batch Execution: The runner script loads the submission's pinned language runtime and iterates through all 150 test cases, capturing output, CPU time, wall time, memory usage, and exit status for each. If a test case fails, the runner immediately halts and reports Wrong Answer or Runtime Error. The test suite version is pinned to the submission so an in-place problem update cannot change an already queued result.
  • Test Data Storage: Huge test cases, such as arrays with 100,000 elements, are not sent over Kafka. They are fetched from private S3 storage or a bounded read through cache available only to trusted workers. Hidden inputs are exposed to the submitted program only as the test data it must process. Expected outputs remain in the trusted runner or comparison service and are never exposed to the submitted program or client. The runner compares captured program output against the expected output after the untrusted process exits.

4. Result Aggregation & Real-time Delivery

Polling at contest scale exhausts Redis and PostgreSQL, making Server-Sent Events push delivery from the result aggregator essential to avoid 200,000 reads per second.

The execution worker publishes incremental and final result events to a Kafka results topic keyed by submission_id. Each event carries a monotonically increasing status sequence, and the Result Aggregator applies only a newer sequence than the last persisted value in PostgreSQL. It performs the durable state update before acknowledging the result event. Redis caches the latest accepted status with a short TTL for reconnects, so a delayed or duplicated status event cannot overwrite a newer terminal result.

JSON
{
  "event_id": "result-evt-21",
  "submission_id": "sub_987654321",
  "sequence": 4,
  "status": "ACCEPTED",
  "runtime_ms": 45,
  "memory_kb": 14200
}
  • Client Delivery (SSE / WebSockets / Polling): The user UI needs the result quickly. Instead of forcing the UI to poll the database heavily, the API uses Server-Sent Events (SSE) as the preferred one-way transport. When a result or incremental status event is available, a routing layer forwards it to the API server holding the user's open SSE connection. WebSockets remain an option when the product also needs bidirectional events.

API Design

Submit Code

Submit returns 202 immediately, and the client opens an SSE stream for incremental status. The submission API carries an idempotency key. The server stores a hash of the normalized request with that key, so a retry returns the original submission while reusing the same key with different code, language, problem, or contest context is rejected as an idempotency conflict.

TYPESCRIPT
export type SubmissionStatus =
  | "QUEUED"
  | "COMPILING"
  | "RUNNING"
  | "ACCEPTED"
  | "WRONG_ANSWER"
  | "TLE"
  | "MLE"
  | "COMPILE_ERROR"
  | "RUNTIME_ERROR"
  | "SYSTEM_ERROR"
  | "CANCELED";

export type ProblemId = number;
export type ContestId = string;
export type LanguageId = string;

export interface SubmitSubmissionRequest {
  problemId: ProblemId;
  contestId?: ContestId;
  language: LanguageId;
  code: string;
  idempotencyKey: string;
}

export interface SubmitSubmissionResponse {
  submissionId: string;
  status: SubmissionStatus;
}

export interface SubmissionStatusEvent {
  eventId: string;
  submissionId: string;
  sequence: number;
  status: SubmissionStatus;
  testcasesPassed?: number;
  runtimeMs?: number;
  memoryKb?: number;
  message?: string;
}
HTTP
POST /api/v1/submissions
Authorization: Bearer <token>
Idempotency-Key: 7f6c8a2d-1
Content-Type: application/json

{
  "problem_id": 123,
  "contest_id": "550e8400-e29b-41d4-a716-446655440000",
  "language": "python3",
  "code": "def twoSum(nums, target):\n    d = {}\n    for i, n in enumerate(nums):\n        if target - n in d:\n            return [d[target - n], i]\n        d[n] = i"
}

HTTP/1.1 202 Accepted
Content-Type: application/json

{
  "submission_id": "sub_987654321",
  "status": "QUEUED"
}

Stream Submission Status (SSE)

HTTP
GET /api/v1/submissions/sub_987654321/stream HTTP/1.1
Authorization: Bearer <token>
Accept: text/event-stream
Last-Event-ID: 0

HTTP/1.1 200 OK
Content-Type: text/event-stream

: heartbeat

id: 1
data: {"sequence": 1, "status": "COMPILING"}

id: 2
data: {"sequence": 2, "status": "RUNNING", "testcases_passed": 10}

id: 3
data: {"sequence": 3, "status": "RUNNING", "testcases_passed": 30}

id: 4
data: {"sequence": 4, "status": "ACCEPTED", "runtime_ms": 45, "memory_kb": 14200}

Common Error Responses

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

We use a relational database (PostgreSQL) for problem metadata, immutable problem and test suite versions, users, contest context, submission state, and durable execution results. Raw source code and hidden test data remain in private S3 object storage.

SQL
CREATE TABLE problems (
    id SERIAL PRIMARY KEY,
    title VARCHAR(255) NOT NULL,
    difficulty VARCHAR(20) NOT NULL CHECK (difficulty IN ('EASY', 'MEDIUM', 'HARD')),
    current_version INT NOT NULL DEFAULT 1
);

CREATE TABLE problem_versions (
    problem_id INT NOT NULL REFERENCES problems(id),
    version INT NOT NULL,
    statement_s3_key TEXT NOT NULL,
    time_limit_ms INT NOT NULL DEFAULT 2000,
    memory_limit_kb INT NOT NULL DEFAULT 256000,
    created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (problem_id, version)
);

CREATE TABLE test_suites (
    problem_id INT NOT NULL REFERENCES problems(id),
    version INT NOT NULL,
    manifest_sha256 CHAR(64) NOT NULL,
    created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (problem_id, version),
    FOREIGN KEY (problem_id, version) REFERENCES problem_versions(problem_id, version)
);

CREATE TABLE submissions (
    id VARCHAR(50) PRIMARY KEY,
    user_id UUID NOT NULL REFERENCES users(id),
    contest_id UUID,
    problem_id INT NOT NULL REFERENCES problems(id),
    problem_version INT NOT NULL,
    test_suite_version INT NOT NULL,
    language VARCHAR(20) NOT NULL,
    language_version VARCHAR(50) NOT NULL,
    runtime_image_digest VARCHAR(128) NOT NULL,
    code_s3_url TEXT NOT NULL, -- Raw source is private in S3
    code_sha256 CHAR(64) NOT NULL,
    status VARCHAR(20) NOT NULL CHECK (status IN ('QUEUED', 'COMPILING', 'RUNNING', 'ACCEPTED', 'WRONG_ANSWER', 'TLE', 'MLE', 'COMPILE_ERROR', 'RUNTIME_ERROR', 'SYSTEM_ERROR', 'CANCELED')),
    result_sequence INT NOT NULL DEFAULT 0,
    runtime_ms INT,
    memory_kb INT,
    idempotency_key VARCHAR(128) NOT NULL,
    request_sha256 CHAR(64) NOT NULL, -- Hash of the normalized request for idempotency conflict detection
    created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
    started_at TIMESTAMPTZ,
    finished_at TIMESTAMPTZ,
    UNIQUE (user_id, idempotency_key),
    FOREIGN KEY (problem_id, problem_version) REFERENCES problem_versions(problem_id, version),
    FOREIGN KEY (problem_id, test_suite_version) REFERENCES test_suites(problem_id, version)
);

CREATE INDEX idx_submissions_user_created ON submissions(user_id, created_at DESC);
CREATE INDEX idx_submissions_problem_created ON submissions(problem_id, created_at DESC);
CREATE INDEX idx_submissions_contest_created ON submissions(contest_id, created_at DESC);

CREATE TABLE submission_dispatch_outbox (
    id UUID PRIMARY KEY,
    submission_id VARCHAR(50) NOT NULL UNIQUE REFERENCES submissions(id),
    payload JSONB NOT NULL,
    created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
    published_at TIMESTAMPTZ
);

CREATE INDEX idx_submission_dispatch_pending
    ON submission_dispatch_outbox(created_at)
    WHERE published_at IS NULL;

CREATE TABLE test_cases (
    problem_id INT NOT NULL,
    suite_version INT NOT NULL,
    test_case_id INT NOT NULL,
    input_s3_key TEXT NOT NULL, -- Hidden test data stays private
    output_s3_key TEXT NOT NULL, -- Expected output stays outside the untrusted process
    is_hidden BOOLEAN NOT NULL,
    PRIMARY KEY (problem_id, suite_version, test_case_id),
    FOREIGN KEY (problem_id, suite_version) REFERENCES test_suites(problem_id, version)
);

Note on Code and Test Storage: The actual submitted source and massive hidden test cases, which can contain gigabytes of integers, are offloaded to AWS S3. Buckets remain private, encryption is enabled, and trusted workers access objects through scoped IAM permissions or short-lived signed access. PostgreSQL stores only object references and metadata, keeping transactional tables small and index efficient.

For submission ingestion, the platform verifies the source object and its SHA-256 hash before publishing the execution job. The database stores the object reference, code hash, request hash, exact version bindings, and idempotency key before the job becomes visible to workers. The stored content hash makes retries and integrity checks deterministic. Failed submissions can leave temporary S3 objects, so a lifecycle policy and asynchronous garbage collector remove unreferenced uploads without touching objects still referenced by durable submission records.

Fault Tolerance

ScenarioHandling Strategy
Submission DB commit succeeds but Kafka is unavailableThe transactional outbox row remains unpublished. A relay retries until Kafka accepts the job, so a broker outage cannot silently drop an accepted submission.
Contest Spikes (100x traffic)Auto-scale the Execution Worker fleet based on the Kafka topic lag. Use prewarmed container pools to avoid startup latency.
Infinite Loops (TLE)The container runtime monitors CPU time using cgroups. A supervisor process kills the container if execution exceeds the problem's time limit.
Unbounded Program OutputCap stdout and stderr bytes inside the sandbox. Truncate diagnostic output for the user and terminate the process if it exceeds the configured output budget so log flooding cannot exhaust worker memory, network, or the result pipeline.
Worker Node CrashIf a worker dies before committing the Kafka offset, the submission is redelivered and may execute again. Results are keyed by submission_id and written idempotently so duplicate execution cannot create duplicate durable results.
Malicious Code EscapeNetwork access is disabled inside the sandbox. seccomp blocks dangerous syscalls. Firecracker microVMs provide a hardware-level isolation barrier.

Additional Considerations

Contest Fairness, Versions, and Leaderboards

Correct judging depends on more than sandbox isolation. Every submission captures the exact immutable problem version, test suite version, and language runtime image digest used for evaluation. A contest leaderboard can maintain live ranks in a Redis Sorted Set keyed by contest ID, while accepted submission results remain authoritative in PostgreSQL and durable score updates are persisted asynchronously so a cache failure does not lose the contest result. The leaderboard projection can be rebuilt from durable submission results after a cache loss. Fair scheduling uses per-user rate limits and queue isolation so one participant cannot consume the entire execution fleet.

Interview Walkthrough

  • 25-minute cut

    Prioritize the trust boundary and execution pipeline. Add staff depth only when the interview requires it.

    • Trust boundary and sandbox isolation for malicious user code (5 min)
    • Decouple ingestion from execution with Kafka and explain worker backpressure (6 min)
    • Capacity math: 1,600 submissions/sec x 2s = ~3,200 concurrent sandbox slots (5 min)
    • One sandbox per submission with batched test cases and pinned runtime versions (5 min)
    • Deliver incremental results through SSE and explain the polling failure mode (4 min)
  • Lead with the trust boundary. Treat all user code as malicious and enforce sandbox isolation with Firecracker or gVisor, seccomp, namespaces, resource limits, and disabled networking.
  • Decouple ingestion from execution. The API validates and enqueues metadata, while workers drain Kafka. Never block HTTP threads on compilation or execution.
  • Capacity math: peak submissions/sec x average execution seconds gives the required concurrent sandbox capacity, which is ~3,200 at 1,600 submissions/sec x 2 seconds.
  • Use one sandbox per submission and batch all test cases inside it to avoid 150 sandbox boots per submission.
  • Deliver results through SSE or WebSocket push. Polling at contest scale can exhaust Redis and the database. For leaderboard ranking patterns, see the Real-Time Gaming Leaderboard.
  • Store source code and large test cases in S3. PostgreSQL holds transactional metadata, version bindings, and durable submission status.
  • Common pitfall: running user code synchronously on the API server, where one infinite loop can consume an HTTP worker and degrade the entire web tier.

Engineering Trade-offs

1. Container Isolation vs. MicroVMs

Standard Docker containers share the host kernel. A kernel vulnerability (like a zero-day in eBPF) could allow a user to break out of the sandbox and compromise the worker node. MicroVMs (like AWS Firecracker) provide hardware virtualized isolation, which is much safer, but have slightly higher startup times (~120ms vs ~50ms). Modern OJ systems trade the slight overhead of MicroVMs for the absolute security they provide against kernel exploits.

2. Real-time Delivery: SSE vs. WebSockets vs. Polling

For returning execution results, we need real time UI updates.

  • Polling: Simplest to implement, entirely stateless. But polling every 500ms creates massive read intensive load on the DB/Redis layer.
  • WebSockets: Full duplex, low latency. But stateful, requiring sticky sessions or complex Pub/Sub routing, overkill for a unidirectional status update.
  • Server-Sent Events (SSE): The perfect middle ground. Unidirectional (Server to Client), runs over standard HTTP, natively supported by browsers, and much easier to scale than WebSockets.

3. Code Storage: DB vs. Object Storage

Storing submitted code strings directly in PostgreSQL is convenient but degrades DB performance as the table grows to Terabytes of text. A better architecture offloads the actual code string to AWS S3 (or a similar object store) and only stores the S3 URI path in PostgreSQL. This keeps the relational database lean and index efficient, trading slight read latency for massive scalability.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...