1#![allow(clippy::unused_unit, deprecated)]
16
17use dioxus::prelude::*;
18
19use crate::api::code_runner::{ExecRequest, ExecTask};
21#[cfg(feature = "server")]
23use crate::api::code_runner::{ExecResult, ExecStatus};
24
25#[cfg(feature = "server")]
27use crate::api::auth::get_current_admin_user;
28#[cfg(feature = "server")]
29use crate::api::code_runner::languages::{is_supported_lang, normalize_lang, LANGUAGES};
30#[cfg(feature = "server")]
31use crate::api::code_runner::progress::{
32 gc_old_tasks, insert_task, update_task_result, update_task_stage, StreamEntry, EXEC_STREAMS,
33 EXEC_TASKS,
34};
35#[cfg(feature = "server")]
36use crate::api::rate_limit::{check_code_exec_limit, get_client_ip};
37#[cfg(feature = "server")]
38use crate::infra::docker::{run_in_container, run_in_container_stream, OutputChunk};
39#[cfg(feature = "server")]
40use crate::infra::runner_config::{clamp_limits, RUNNER_CONFIG};
41#[cfg(feature = "server")]
42use std::sync::{Arc, LazyLock};
43#[cfg(feature = "server")]
44use std::time::Duration;
45#[cfg(feature = "server")]
46use tokio::sync::Semaphore;
47
48#[cfg(feature = "server")]
50pub static RUNNER_SEMAPHORE: LazyLock<Arc<Semaphore>> =
51 LazyLock::new(|| Arc::new(Semaphore::new(RUNNER_CONFIG.max_concurrent)));
52
53#[cfg(feature = "server")]
55async fn client_ip() -> String {
56 match dioxus::fullstack::FullstackContext::current() {
57 Some(ctx) => {
58 let headers = ctx.parts_mut().headers.clone();
59 get_client_ip(&headers).await
60 }
61 None => "unknown".to_string(),
62 }
63}
64
65#[cfg(feature = "server")]
70fn validate_exec_request(req: &ExecRequest) -> Result<(), ServerFnError> {
71 if !is_supported_lang(&req.language) {
73 return Err(ServerFnError::new("不支持该执行语言".to_string()));
74 }
75
76 if req.source.len() > RUNNER_CONFIG.max_source_bytes as usize {
78 return Err(ServerFnError::new("源代码过大".to_string()));
79 }
80
81 Ok(())
82}
83
84#[cfg(feature = "server")]
86async fn check_rate_limit_for_user() -> Result<(), ServerFnError> {
87 let is_admin = get_current_admin_user().await.is_ok();
88 if !is_admin {
89 let ip = client_ip().await;
90 if let Err(msg) = check_code_exec_limit(&ip) {
91 return Err(ServerFnError::new(msg));
92 }
93 }
94 Ok(())
95}
96
97#[cfg(feature = "server")]
107fn spawn_exec_task(
108 task_id: String,
109 req: ExecRequest,
110 stream_tx: Option<tokio::sync::mpsc::Sender<OutputChunk>>,
111) {
112 let lang_key = normalize_lang(&req.language);
115 tokio::spawn(async move {
116 let sem = &*RUNNER_SEMAPHORE;
117
118 let ticket = match tokio::time::timeout(
120 Duration::from_secs(RUNNER_CONFIG.queue_timeout_secs),
121 sem.acquire(),
122 )
123 .await
124 {
125 Ok(Ok(t)) => t,
126 _ => {
127 update_task_stage(&task_id, ExecStatus::Failed, "系统繁忙,排队超时");
128 return;
129 }
130 };
131
132 update_task_stage(&task_id, ExecStatus::Running, "启动容器");
133
134 let lang_def = match LANGUAGES.get(&lang_key) {
135 Some(d) => d,
136 None => {
137 update_task_stage(&task_id, ExecStatus::Failed, "语言未注册");
139 return;
140 }
141 };
142
143 let base_limits = req
145 .overrides
146 .unwrap_or_else(|| lang_def.default_limits.clone());
147 let final_limits = clamp_limits(base_limits, lang_def.allow_network);
148
149 let start_time = chrono::Utc::now();
150 let stream_suffix = if stream_tx.is_some() { " (stream)" } else { "" };
152 let res = match &stream_tx {
153 Some(tx) => run_in_container_stream(
154 &lang_def.image,
155 &lang_def.run_cmd,
156 &req.source,
157 &lang_def.extension,
158 final_limits,
159 lang_def.cache_volume.as_ref(),
160 tx.clone(),
161 )
162 .await
163 .map(|(exit_code, stdout, stderr, oom_killed, _)| {
164 (exit_code, stdout, stderr, oom_killed)
165 }),
166 None => {
167 run_in_container(
168 &lang_def.image,
169 &lang_def.run_cmd,
170 &req.source,
171 &lang_def.extension,
172 final_limits,
173 lang_def.cache_volume.as_ref(),
174 )
175 .await
176 }
177 };
178 let duration_ms = (chrono::Utc::now() - start_time).num_milliseconds().max(0) as u64;
179
180 drop(ticket); match res {
183 Ok((exit_code, stdout, stderr, oom_killed)) => {
184 let status = if oom_killed {
185 ExecStatus::OomKilled
186 } else if exit_code == Some(0) {
187 ExecStatus::Success
188 } else {
189 ExecStatus::Error
190 };
191 let exec_res = ExecResult {
192 status: status.clone(),
193 stdout,
194 stderr,
195 exit_code,
196 duration_ms,
197 language: lang_key,
198 };
199 update_task_result(&task_id, status, exec_res);
200 }
201 Err(e) => {
202 tracing::error!(error = ?e, task_id = %task_id, "container execution failed{}", stream_suffix);
204 let (status, stderr_msg) = classify_runner_error(&e, &lang_def.image);
205
206 if let Some(tx) = &stream_tx {
211 let _ = tx
212 .send(OutputChunk::Done {
213 exit_code: None,
214 oom_killed: false,
215 timed_out: status == ExecStatus::Timeout,
216 duration_ms,
217 error: Some(stderr_msg.clone()),
218 })
219 .await;
220 }
221
222 let exec_res = ExecResult {
223 status: status.clone(),
224 stdout: String::new(),
225 stderr: stderr_msg,
226 exit_code: None,
227 duration_ms,
228 language: lang_key,
229 };
230 update_task_result(&task_id, status, exec_res);
231 }
232 }
233 });
234}
235
236#[cfg(feature = "server")]
250fn classify_runner_error(e: &bollard::errors::Error, image: &str) -> (ExecStatus, String) {
251 if let bollard::errors::Error::IOError { err } = e {
253 if err.kind() == std::io::ErrorKind::TimedOut {
254 return (ExecStatus::Timeout, "执行超时".to_string());
255 }
256 }
257
258 let s = e.to_string();
259 if s.contains("No such image") {
260 (
261 ExecStatus::Failed,
262 format!("运行器镜像未构建:{image}。请在宿主执行:bash docker/build-runners.sh"),
263 )
264 } else if s.contains("Docker daemon 不可用") {
265 (
266 ExecStatus::Failed,
267 "Docker 未运行或 socket 未挂载,代码运行不可用".to_string(),
268 )
269 } else {
270 (ExecStatus::Failed, "系统暂时不可用".to_string())
271 }
272}
273
274#[server(StartExec, "/api")]
282pub async fn start_exec(req: ExecRequest) -> Result<String, ServerFnError> {
283 check_rate_limit_for_user().await?;
284 validate_exec_request(&req)?;
285
286 let task_id = uuid::Uuid::new_v4().to_string();
288 insert_task(task_id.clone());
289
290 gc_old_tasks();
292
293 spawn_exec_task(task_id.clone(), req, None);
296
297 Ok(task_id)
298}
299
300#[server(StartExecStream, "/api")]
310pub async fn start_exec_stream(req: ExecRequest) -> Result<String, ServerFnError> {
311 check_rate_limit_for_user().await?;
312 validate_exec_request(&req)?;
313
314 let task_id = uuid::Uuid::new_v4().to_string();
315 insert_task(task_id.clone());
317
318 let (tx, rx) = tokio::sync::mpsc::channel(64);
320 EXEC_STREAMS.insert(
321 task_id.clone(),
322 StreamEntry {
323 rx,
324 created_at: chrono::Utc::now(),
325 },
326 );
327
328 gc_old_tasks();
329
330 spawn_exec_task(task_id.clone(), req, Some(tx));
332
333 Ok(task_id)
334}
335
336#[server(GetExecResult, "/api")]
338pub async fn get_exec_result(task_id: String) -> Result<ExecTask, ServerFnError> {
339 if let Some(task) = EXEC_TASKS.get(&task_id) {
340 Ok(task.clone())
341 } else {
342 Err(ServerFnError::new("找不到指定的任务".to_string()))
343 }
344}
345
346#[cfg(all(test, feature = "server"))]
347mod tests {
348 use super::*;
349
350 fn req(language: &str, source: &str) -> ExecRequest {
351 ExecRequest {
352 language: language.to_string(),
353 source: source.to_string(),
354 overrides: None,
355 }
356 }
357
358 #[test]
359 fn validate_accepts_registered_language() {
360 for lang in ["python", "node", "go", "rust", "bun"] {
362 assert!(
363 validate_exec_request(&req(lang, "x")).is_ok(),
364 "{lang} 应被支持"
365 );
366 }
367 }
368
369 #[test]
370 fn validate_accepts_aliases() {
371 for lang in ["js", "javascript", "rs", "ts", "typescript"] {
375 assert!(
376 validate_exec_request(&req(lang, "x")).is_ok(),
377 "别名 {lang} 应被支持"
378 );
379 }
380 }
381
382 #[test]
383 fn validate_language_is_case_and_whitespace_insensitive() {
384 assert!(validate_exec_request(&req("Python", "x")).is_ok());
386 assert!(validate_exec_request(&req(" RUST ", "x")).is_ok());
387 assert!(validate_exec_request(&req("Go", "x")).is_ok());
388 }
389
390 #[test]
391 fn validate_rejects_unregistered_language() {
392 assert!(validate_exec_request(&req("c", "x")).is_err());
394 assert!(validate_exec_request(&req("bash", "x")).is_err());
395 assert!(validate_exec_request(&req("python2", "x")).is_err());
396 assert!(validate_exec_request(&req("", "x")).is_err());
397 }
398
399 #[test]
400 fn validate_rejects_multi_token_language() {
401 assert!(validate_exec_request(&req("python; rm -rf /", "x")).is_err());
404 assert!(validate_exec_request(&req("python node", "x")).is_err());
405 assert!(validate_exec_request(&req("python$(whoami)", "x")).is_err());
406 }
407
408 #[test]
409 fn validate_language_tolerates_surrounding_whitespace() {
410 assert!(validate_exec_request(&req("python\n", "x")).is_ok());
413 assert!(validate_exec_request(&req("\tpython\t", "x")).is_ok());
414 }
415
416 #[test]
417 fn validate_rejects_source_exceeding_max_bytes() {
418 let max = RUNNER_CONFIG.max_source_bytes as usize;
419 let exactly = "a".repeat(max);
421 assert!(
422 validate_exec_request(&req("python", &exactly)).is_ok(),
423 "源码恰好等于上限应放行"
424 );
425 let over = "a".repeat(max + 1);
426 let err = validate_exec_request(&req("python", &over)).expect_err("超出上限应拒绝");
427 assert!(
428 err.to_string().contains("过大"),
429 "错误信息应提及大小: {err}"
430 );
431 }
432
433 #[test]
434 fn validate_empty_source_accepted_for_supported_lang() {
435 assert!(validate_exec_request(&req("python", "")).is_ok());
438 }
439
440 #[test]
441 fn validate_checks_language_before_size() {
442 let huge = "a".repeat((RUNNER_CONFIG.max_source_bytes as usize) + 100);
444 let err = validate_exec_request(&req("brainfuck", &huge)).unwrap_err();
445 assert!(err.to_string().contains("语言"), "应先报语言错误: {err}");
446 }
447
448 #[test]
449 fn classify_runner_error_image_missing() {
450 let e = bollard::errors::Error::DockerResponseServerError {
453 status_code: 404,
454 message: "No such image: yggdrasil-runner-python:latest".to_string(),
455 };
456 let (status, stderr) = classify_runner_error(&e, "yggdrasil-runner-python:latest");
457 assert_eq!(status, ExecStatus::Failed);
458 assert_eq!(
459 stderr,
460 "运行器镜像未构建:yggdrasil-runner-python:latest。请在宿主执行:bash docker/build-runners.sh"
461 );
462 }
463
464 #[test]
465 fn classify_runner_error_daemon_unavailable() {
466 let e = bollard::errors::Error::IOError {
468 err: std::io::Error::new(
469 std::io::ErrorKind::NotFound,
470 "Docker daemon 不可用(未安装或未运行)",
471 ),
472 };
473 let (status, stderr) = classify_runner_error(&e, "yggdrasil-runner-python:latest");
474 assert_eq!(status, ExecStatus::Failed);
475 assert_eq!(stderr, "Docker 未运行或 socket 未挂载,代码运行不可用");
476 }
477
478 #[test]
479 fn classify_runner_error_timeout() {
480 let e = bollard::errors::Error::IOError {
484 err: std::io::Error::new(std::io::ErrorKind::TimedOut, "Execution timed out"),
485 };
486 let (status, stderr) = classify_runner_error(&e, "yggdrasil-runner-python:latest");
487 assert_eq!(status, ExecStatus::Timeout);
488 assert_eq!(stderr, "执行超时");
489 }
490
491 #[test]
492 fn classify_runner_error_generic_fallback() {
493 let e = bollard::errors::Error::DockerResponseServerError {
495 status_code: 500,
496 message: "boom".to_string(),
497 };
498 let (status, stderr) = classify_runner_error(&e, "yggdrasil-runner-python:latest");
499 assert_eq!(status, ExecStatus::Failed);
500 assert_eq!(stderr, "系统暂时不可用");
501 }
502}