Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,22 @@ in your migration, hidden ones included. It fails if:
- the starter code does not compile, has no `Solution` class, or passes the
tests without being filled in

### Optional: support complexity analysis

An accepted run can be rerun on bigger inputs to estimate how its time grows.
A problem supports this when its migration also sets two columns:

- `complexity_generator`: Python source defining `SIZES`, the input sizes to
try in order, and `generate(n)`, which returns the stdin for an input of
size `n`. Make the inputs force a full solve, for example by putting the
answer at the end.
- `expected_complexity`: the class your reference solution achieves, one of
`constant`, `linear`, `quadratic` or `cubic`.

The test suite measures your reference solution with your generator and fails
if it does not land in the expected class. `V13__complexity_analysis.sql` has
examples.

### 5. Check it

Run `make test`. To see the problem in the app,
Expand Down
15 changes: 14 additions & 1 deletion apps/backend/app/dao/problems.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,10 @@

# The full test list never leaves the DAO through here: only the public
# samples do, so a client cannot read the hidden tests off the API.
COLUMNS = "id, slug, title, difficulty, time_limit_ms, mem_limit_mb, statement, starter_code, tests"
COLUMNS = (
"id, slug, title, difficulty, time_limit_ms, mem_limit_mb, statement, starter_code, tests, "
"complexity_generator IS NOT NULL AS analyzable, expected_complexity"
)


def _map_problem(row: asyncpg.Record) -> dict:
Expand All @@ -25,6 +28,8 @@ def _map_problem(row: asyncpg.Record) -> dict:
"samples": [
{"input": t["input"], "expected": t["expected"]} for t in tests if not t.get("hidden")
],
"analyzable": row["analyzable"],
"expectedComplexity": row["expected_complexity"],
}


Expand Down Expand Up @@ -84,6 +89,14 @@ async def get_judge_spec(self, problem_id: UUID) -> dict:
"memLimitMb": row["mem_limit_mb"],
}

async def get_analysis_spec(self, problem_id: UUID) -> dict | None:
query = "SELECT complexity_generator, mem_limit_mb FROM problems WHERE id = $1"
async with self.pool.acquire() as conn:
row = await conn.fetchrow(query, problem_id)
if row is None or row["complexity_generator"] is None:
return None
return {"generator": row["complexity_generator"], "memLimitMb": row["mem_limit_mb"]}

async def exists(self, problem_id: UUID) -> bool:
query = "SELECT EXISTS(SELECT 1 FROM problems WHERE id = $1)"
async with self.pool.acquire() as conn:
Expand Down
55 changes: 54 additions & 1 deletion apps/backend/app/dao/submissions.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
SELECT_SUBMISSION = """
SELECT s.id, s.room_id, s.user_id, s.problem_id, s.language, s.code, s.status,
s.time_ms, s.created_at, s.s3_key_stdout, s.s3_key_stderr,
s.s3_key_result_json, s.result, u.name AS user_name
s.s3_key_result_json, s.result, s.analysis_status, s.analysis, u.name AS user_name
FROM submissions s
JOIN users u ON u.id = s.user_id
"""
Expand Down Expand Up @@ -40,6 +40,8 @@ def _map_submission(row: asyncpg.Record) -> dict:
"s3KeyStderr": row["s3_key_stderr"],
"s3KeyResultJson": row["s3_key_result_json"],
"result": json.loads(row["result"]) if row["result"] else None,
"analysisStatus": row["analysis_status"],
"analysis": json.loads(row["analysis"]) if row["analysis"] else None,
}


Expand Down Expand Up @@ -119,3 +121,54 @@ async def requeue_running(self) -> int:
async with self.pool.acquire() as conn:
tag = await conn.execute("UPDATE submissions SET status = 'pending' WHERE status = 'running'")
return int(tag.rsplit(" ", 1)[1])

async def count_analyses_since(self, user_id: str, seconds: int) -> int:
query = """
SELECT COUNT(*) FROM submissions
WHERE analysis_requested_by = $1 AND analysis_requested_at > NOW() - make_interval(secs => $2)
"""
async with self.pool.acquire() as conn:
return await conn.fetchval(query, user_id, seconds)

async def request_analysis(self, submission_id: UUID, user_id: str) -> dict | None:
"""Queue an analysis of an accepted run that has never had one."""
query = """
UPDATE submissions
SET analysis_status = 'pending', analysis_requested_by = $2, analysis_requested_at = NOW()
WHERE id = $1 AND status = 'accepted' AND analysis_status IS NULL
RETURNING id
"""
async with self.pool.acquire() as conn:
updated = await conn.fetchval(query, submission_id, user_id)
return await self.get_by_id(updated) if updated else None
Comment thread
coderabbitai[bot] marked this conversation as resolved.

async def claim_pending_analysis(self) -> dict | None:
query = """
UPDATE submissions SET analysis_status = 'running'
WHERE id = (
SELECT id FROM submissions
WHERE analysis_status = 'pending'
ORDER BY analysis_requested_at
LIMIT 1
FOR UPDATE SKIP LOCKED
)
RETURNING id
"""
async with self.pool.acquire() as conn:
submission_id = await conn.fetchval(query)
return await self.get_by_id(submission_id) if submission_id else None

async def complete_analysis(self, submission_id: UUID, status: str, analysis: dict) -> None:
query = "UPDATE submissions SET analysis_status = $2, analysis = $3::jsonb WHERE id = $1"
async with self.pool.acquire() as conn:
async with conn.transaction():
await conn.execute(query, submission_id, status, json.dumps(analysis))
# The room hears about it the way it hears about a verdict.
await conn.execute("SELECT pg_notify($1, $2)", JUDGED_CHANNEL, str(submission_id))

async def requeue_running_analyses(self) -> int:
async with self.pool.acquire() as conn:
tag = await conn.execute(
"UPDATE submissions SET analysis_status = 'pending' WHERE analysis_status = 'running'"
)
return int(tag.rsplit(" ", 1)[1])
15 changes: 15 additions & 0 deletions apps/backend/app/routes/submissions.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,21 @@ async def list_submissions(
return [SubmissionResponse.model_validate(submission) for submission in submissions]


@router.post("/{submission_id}/analysis", response_model=SubmissionResponse)
async def request_analysis(
submission_id: UUID,
request: Request,
caller_id: str = Depends(current_user_id),
service: SubmissionService = Depends(get_submission_service),
) -> SubmissionResponse:
submission = await service.request_analysis(submission_id, caller_id)
# Everyone in the room sees it start; the result follows from the listener.
await request.app.state.room_chat_manager.broadcast(
submission["roomId"], {"type": "submission", "submission": submission}
)
return SubmissionResponse.model_validate(submission)


@router.get("/{submission_id}", response_model=SubmissionResponse)
async def get_submission(
submission_id: UUID,
Expand Down
42 changes: 41 additions & 1 deletion apps/backend/app/runner/__main__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
"""Claim pending submissions and judge them, one at a time.

Complexity analyses share the runner but only run when no submission is
waiting, so asking for one never delays anyone's verdict by more than the
analysis already in progress.

The submissions table is the queue: a row is claimed by moving it to
``running`` under ``FOR UPDATE SKIP LOCKED``, so several runners can share a
database without judging the same submission twice.
Expand All @@ -16,6 +20,7 @@
from app.dao.submissions import SubmissionDAO
from app.database import create_pool
from app.runner import sandbox
from app.runner.complexity import analyze
from app.runner.judge import RUNTIME_ERROR, Verdict, judge

logger = logging.getLogger(__name__)
Expand All @@ -36,6 +41,40 @@ def run_judge(code: str, tests: list[dict], time_limit_ms: int, mem_limit_mb: in
return judge(code, tests, time_limit_ms, mem_limit_mb)


def run_analysis(code: str, generator: str, mem_limit_mb: int) -> dict:
if SANDBOX_IMAGE:
return sandbox.analyze_in_container(SANDBOX_IMAGE, code, generator, mem_limit_mb)
return analyze(code, generator, mem_limit_mb)


async def analyze_next(submissions: SubmissionDAO, problems: ProblemDAO) -> bool:
"""Analyze one accepted run if one is waiting. Returns whether it did."""
submission = await submissions.claim_pending_analysis()
if submission is None:
return False

try:
spec = await problems.get_analysis_spec(submission["problemId"])
if spec is None:
raise RuntimeError("the problem has no generator")
analysis = await asyncio.to_thread(run_analysis, submission["code"] or "", spec["generator"], spec["memLimitMb"])
except Exception:
logger.exception("Analysis failed for submission %s", submission["id"])
analysis = {"points": [], "complexity": None, "slope": None, "note": "The analysis could not run."}

status = "done" if analysis["complexity"] else "failed"
try:
await submissions.complete_analysis(submission["id"], status, analysis)
except Exception:
logger.exception("Could not store the analysis for submission %s", submission["id"])
status = "failed"
await submissions.complete_analysis(
submission["id"], status, {"points": [], "complexity": None, "slope": None, "note": "The analysis could not be saved."}
)
logger.info("Submission %s analysis: %s", submission["id"], analysis["complexity"] or status)
return True


async def judge_next(submissions: SubmissionDAO, problems: ProblemDAO) -> bool:
"""Judge one submission if there is one waiting. Returns whether it did."""
submission = await submissions.claim_pending()
Expand Down Expand Up @@ -80,9 +119,10 @@ async def main() -> None:
# claimed again. One runner at a time is the local setup, so on boot
# anything still running is ours from before and goes back in the queue.
await submissions.requeue_running()
await submissions.requeue_running_analyses()
try:
while True:
if not await judge_next(submissions, problems):
if not await judge_next(submissions, problems) and not await analyze_next(submissions, problems):
await asyncio.sleep(POLL_SECONDS)
finally:
await pool.close()
Expand Down
161 changes: 161 additions & 0 deletions apps/backend/app/runner/complexity.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,161 @@
"""Estimate how a program's running time grows with the size of its input.

The program runs on inputs of doubling size from the problem's generator, and
the CPU time of each run is fitted to a power of n. CPU time, not a count of
executed lines, because a line count sees `x in some_list` or `sorted(xs)` as
one step and would call a quadratic brute force linear. It is measured from
outside the program, so the program cannot report a time of its own.

Like the judge, this file is shipped into a sandbox container as the program
text, so it imports nothing from the app.
"""

from __future__ import annotations

import math
import resource
import subprocess
import sys
import tempfile
import time
from pathlib import Path

# Runner time one analysis may use, and the most one size may take. A run
# that reaches a second has shown its growth; bigger inputs only cost time.
BUDGET_SECONDS = 6.0
PER_RUN_SECONDS = 2.5
ENOUGH_SECONDS = 1.0

# Below this a run is mostly noise, not the solution.
MIN_MEASURABLE_MS = 20.0

# Fitted exponent of n -> growth class, split halfway between the powers.
# n log n fits at about 1.1 over the sizes used, too close to n to tell
# apart, so they share a class.
CLASSES = [(0.5, "constant"), (1.5, "linear"), (2.5, "quadratic"), (3.5, "cubic")]
WORSE = "worse"

PROGRAM_ENV = {"PATH": "/usr/local/bin:/usr/bin:/bin", "LANG": "C.UTF-8"}


def _limits(mem_limit_mb: int):
# The same limits the judge applies; this file cannot import them.
def apply() -> None:
mem = mem_limit_mb * 1024 * 1024
resource.setrlimit(resource.RLIMIT_AS, (mem, mem))
resource.setrlimit(resource.RLIMIT_NPROC, (0, 0))
resource.setrlimit(resource.RLIMIT_FSIZE, (0, 0))

return apply


def _children_cpu() -> float:
usage = resource.getrusage(resource.RUSAGE_CHILDREN)
return usage.ru_utime + usage.ru_stime


def classify(points: list[tuple[int, float]]) -> tuple[str, float]:
"""Least-squares slope of log(ms) against log(n), and its class."""
xs = [math.log(n) for n, _ in points]
ys = [math.log(ms) for _, ms in points]
mean_x, mean_y = sum(xs) / len(xs), sum(ys) / len(ys)
spread = sum((x - mean_x) ** 2 for x in xs)
slope = sum((x - mean_x) * (y - mean_y) for x, y in zip(xs, ys)) / spread
for bound, name in CLASSES:
if slope < bound:
return name, slope
return WORSE, slope


def _run(program: Path, workdir: str, stdin: str, timeout: float, mem_limit_mb: int):
"""CPU and wall seconds for one run, or None and a reason if it did not finish."""
cpu_before, started = _children_cpu(), time.perf_counter()
try:
completed = subprocess.run(
[sys.executable, "-I", str(program)],
input=stdin.encode(),
# Nothing is checked here, and a large input can mean a large answer.
stdout=subprocess.DEVNULL,
stderr=subprocess.PIPE,
timeout=timeout,
cwd=workdir,
env=PROGRAM_ENV,
preexec_fn=_limits(mem_limit_mb),
)
except subprocess.TimeoutExpired:
return None, f"took over {timeout:.1f} s"
wall = time.perf_counter() - started
if completed.returncode != 0:
# Postgres rejects NUL in jsonb, and a result that cannot be stored
# would leave the runner retrying the same analysis forever.
last = completed.stderr.decode("utf-8", errors="replace").replace("\x00", "").strip().splitlines()
return None, f"crashed: {last[-1][:200]}" if last else "crashed"
# Some sandboxes do not account children's CPU; wall time is the fallback.
cpu = _children_cpu() - cpu_before
return (cpu if cpu > 0 else wall), None


def analyze(code: str, generator: str, mem_limit_mb: int) -> dict:
namespace: dict = {}
exec(generator, namespace)
generate, sizes = namespace["generate"], namespace["SIZES"]

points: list[tuple[int, float]] = []
note = None
spent = 0.0
with tempfile.TemporaryDirectory() as workdir:
program = Path(workdir) / "main.py"
program.write_text(code)
# Every run pays for starting the interpreter, the program's imports
# and reading its input. The program on a tiny input costs about that
# and nothing more; left in, it flattens the growth of every run.
tiny = generate(max(4, sizes[0] // 32))
startup = min(_run(program, workdir, tiny, PER_RUN_SECONDS, mem_limit_mb)[0] or 0.0 for _ in range(3))

for n in sizes:
remaining = BUDGET_SECONDS - spent
if remaining <= 0.1:
note = f"Stopped before n = {n}: out of time for this analysis."
break
seconds, failure = _run(program, workdir, generate(n), min(PER_RUN_SECONDS, remaining), mem_limit_mb)
if seconds is None:
note = f"Stopped at n = {n}: it {failure}."
break
spent += seconds
ms = (seconds - startup) * 1000
if ms >= MIN_MEASURABLE_MS:
points.append((n, ms))
if seconds >= ENOUGH_SECONDS:
break

result = {
"points": [{"n": n, "ms": round(ms, 1)} for n, ms in points],
"complexity": None,
"slope": None,
"note": note,
}
if len(points) >= 2:
result["complexity"], slope = classify(points)
result["slope"] = round(slope, 2)
if len(points) == 2:
result["note"] = (note + " " if note else "") + "Only two sizes were measurable, so this is rough."
elif not points and note is None:
# Even the largest input ran too fast to measure: it barely grows.
result["complexity"], result["slope"] = "constant", 0.0
elif note is None:
result["note"] = "Only one size was measurable, which is not enough to see growth."
return result


if __name__ == "__main__":
# Entry point inside a sandbox container: the spec arrives on stdin and
# the result leaves on stdout, both as JSON.
import ctypes
import json

# As in the judge: the program must not be able to print our result for us.
PR_SET_DUMPABLE = 4
ctypes.CDLL(None).prctl(PR_SET_DUMPABLE, 0)

spec = json.load(sys.stdin)
print(json.dumps(analyze(spec["code"], spec["generator"], spec["memLimitMb"])))
Loading
Loading