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
53 changes: 53 additions & 0 deletions tasks/submission_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@
from __future__ import annotations

import json
import logging
import threading
import time

from models import Assignment, Submission, SystemLog, TestCase as TC, User, db
from services.demo_database import activate_demo_run, is_active_demo_run
Expand All @@ -12,6 +14,35 @@
from utils.scoring import normalize_evaluation_score, normalize_feedback_text


logger = logging.getLogger(__name__)


def _log_submission_evaluation_event(
event: str,
submission_id: int,
started_at: float,
*,
level: int = logging.INFO,
**fields,
) -> None:
"""Write a bounded lifecycle signal without logging submission content."""

parts = [
"submission_evaluation",
f"event={event}",
f"submission_id={int(submission_id)}",
f"elapsed_ms={int((time.perf_counter() - started_at) * 1000)}",
]
parts.extend(f"{key}={value}" for key, value in fields.items())
try:
from flask import current_app, has_app_context

target_logger = current_app.logger if has_app_context() else logger
except RuntimeError:
target_logger = logger
target_logger.log(level, " ".join(parts))


def _demo_database_is_available(demo_run_id: str | None) -> bool:
"""Return whether this worker may still use its temporary database."""

Expand Down Expand Up @@ -123,9 +154,17 @@ def evaluate_submission_async(
raise

def _evaluate():
started_at = time.perf_counter()
with app.app_context():
_log_submission_evaluation_event("started", submission_id, started_at)
if demo_run_id and not activate_demo_run(demo_run_id):
print("公开体验会话已失效,跳过提交评测")
_log_submission_evaluation_event(
"skipped",
submission_id,
started_at,
reason="demo_run_unavailable",
)
return

try:
Expand Down Expand Up @@ -337,11 +376,25 @@ def _evaluate():
user_id=student_id,
icon="bi bi-check-circle-fill",
)
_log_submission_evaluation_event(
"finished",
submission_id,
started_at,
state=submission.status,
sandbox_status=submission.sandbox_status or "none",
)
print(f"提交 {submission_id} 评测全部完成")
return "evaluated"

except Exception as error:
print(f"评测线程崩溃: {type(error).__name__}")
_log_submission_evaluation_event(
"failed",
submission_id,
started_at,
level=logging.WARNING,
error_type=type(error).__name__,
)
if not _demo_database_is_available(demo_run_id):
return
try:
Expand Down
58 changes: 57 additions & 1 deletion tests/test_submission_worker.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import inspect
import re
import sys
import time
from concurrent.futures import ThreadPoolExecutor
Expand Down Expand Up @@ -91,6 +92,40 @@ def test_worker_entry_disables_web_threads_and_preset_scanner(monkeypatch):
build.assert_called_once_with(app)


def test_expired_demo_run_emits_terminal_skipped_event(caplog):
_require_worker_contract()
caplog.set_level("INFO")
app = Flask(__name__)

class ImmediateThread:
def __init__(self, target, *args, **kwargs):
self.target = target
self.daemon = False

def start(self):
self.target()

with patch.object(worker_tasks.threading, "Thread", ImmediateThread), \
patch.object(worker_tasks, "activate_demo_run", return_value=False):
worker_tasks.evaluate_submission_async(
app,
41,
"循环题",
demo_run_id="expired-demo-run",
)

messages = [record.getMessage() for record in caplog.records]
assert any(
"submission_evaluation event=started submission_id=41" in message
for message in messages
)
assert any(
"submission_evaluation event=skipped submission_id=41" in message
and "reason=demo_run_unavailable" in message
for message in messages
)


def test_rq_backend_enqueues_without_starting_legacy_thread():
_require_worker_contract()
app = Flask(__name__)
Expand Down Expand Up @@ -193,8 +228,11 @@ def slow_completion(submission_id):
assert submission_queue.get_submission_job_status(app, 42) == "failed"


def test_formal_worker_updates_submission_in_isolated_database(tmp_path, monkeypatch):
def test_formal_worker_updates_submission_in_isolated_database(
tmp_path, monkeypatch, caplog
):
_require_worker_contract()
caplog.set_level("INFO")
from app import create_app
from config import TestingConfig as _TestingConfig
from models import Assignment, StudentLearningVector, StudentVectorIndexState, Submission, User, db
Expand Down Expand Up @@ -288,5 +326,23 @@ def test_formal_worker_updates_submission_in_isolated_database(tmp_path, monkeyp
for vector in active_vectors
)

finished_events = [
record.getMessage()
for record in caplog.records
if "submission_evaluation event=finished" in record.getMessage()
]
started_events = [
record.getMessage()
for record in caplog.records
if "submission_evaluation event=started" in record.getMessage()
]
assert len(started_events) == 1
assert len(finished_events) == 1
finished_event = finished_events[0]
assert f"submission_id={submission_id}" in finished_event
assert "state=evaluated" in finished_event
elapsed_match = re.search(r"elapsed_ms=(\d+)", finished_event)
assert elapsed_match is not None

assert submission_queue.get_submission_job_status(app, submission_id) == "completed"
assert state.operation_id == f"submission-evaluation-{submission_id}"
Loading