1use super::{artifacts, evidence::FrozenEvidence, protocol::*};
5use crate::command_exec::evaluator_io::{self, Output};
6use crate::pack::evaluator::PinnedRegistration;
7use cap_fs_ext::DirExt;
8use cap_std::{ambient_authority, fs::Dir};
9use serde::Serialize;
10use std::collections::HashMap;
11use std::os::unix::fs::{DirBuilderExt, OpenOptionsExt, PermissionsExt};
12use std::path::{Path, PathBuf};
13use std::sync::atomic::AtomicBool;
14use std::time::Duration;
15
16const DRIVER: &str = "set -eu\nexpected=$1\nshift\nfor root in /gate/inputs /checker; do\n test \"$(cat \"$root/engine-mount-proof\")\" = \"$expected\" || exit 125\ndone\nfor root in /gate/outputs /gate/build /gate/home; do\n printf '%s' \"$expected\" > \"$root/engine-mount-proof\"\ndone\nexec /usr/bin/env -i HOME=/gate/home TMPDIR=/gate/build PATH=/usr/local/bin:/usr/bin:/bin \"$@\"\n";
17const CONTROL_LIMIT: usize = 1_048_576;
18
19#[derive(Clone)]
22pub struct DockerEvaluator {
23 program: PathBuf,
24 env: HashMap<String, String>,
25}
26
27pub struct RunOptions<'a> {
28 pub attempt_parent: &'a Path,
31 pub retain_private_inputs: bool,
33}
34
35#[derive(Debug, Serialize)]
36#[serde(rename_all = "camelCase")]
37pub struct AcceptedEvaluation {
38 pub raw_stdout_digest: Digest,
40 pub raw_stdout_retained: bool,
41 pub retained_result_digest: Digest,
42 pub transformation: String,
43 pub result: EvaluationResult,
44 pub artifacts: Vec<artifacts::ImportedArtifact>,
45 pub diagnostics: String,
46 pub exit_code: i32,
47 pub container_removed: bool,
48}
49
50#[derive(Debug)]
51pub struct AttemptOutcome {
52 pub directory: PathBuf,
55 pub cleanup_confirmed: bool,
56 pub evaluation: Result<AcceptedEvaluation, String>,
57}
58
59impl DockerEvaluator {
60 pub fn new(program: &Path) -> Result<Self, String> {
62 if !program.is_absolute() {
63 return Err("Docker CLI must be an absolute trusted host path".into());
64 }
65 let program = program.canonicalize().map_err(|e| e.to_string())?;
66 if !program.is_file() {
67 return Err("Docker CLI is not a regular file".into());
68 }
69 Ok(Self {
70 program,
71 env: crate::sandbox_container::ContainerRuntime::Docker.client_env(),
72 })
73 }
74 async fn command(
75 &self,
76 args: &[String],
77 input: &[u8],
78 wall: Duration,
79 write: Duration,
80 caps: (usize, usize),
81 cancelled: &AtomicBool,
82 ) -> Result<Output, String> {
83 evaluator_io::run(
84 &self.program,
85 args,
86 &self.env,
87 input,
88 wall,
89 write,
90 caps.0,
91 caps.1,
92 cancelled,
93 )
94 .await
95 }
96 pub(crate) fn attached_command(&self, args: &[String]) -> tokio::process::Command {
97 let mut command = tokio::process::Command::new(&self.program);
98 command.args(args).env_clear().envs(&self.env);
99 command
100 }
101
102 pub(crate) async fn control(&self, args: &[String]) -> Result<Output, String> {
103 self.bounded_control(args, Duration::from_secs(5), &AtomicBool::new(false))
104 .await
105 }
106
107 pub(crate) async fn bounded_control(
108 &self,
109 args: &[String],
110 wall: Duration,
111 cancelled: &AtomicBool,
112 ) -> Result<Output, String> {
113 self.command(
114 args,
115 b"",
116 wall,
117 Duration::from_secs(1),
118 (CONTROL_LIMIT, CONTROL_LIMIT),
119 cancelled,
120 )
121 .await
122 .map_err(|error| format!("Docker {} control failed: {error}", args[0]))
123 }
124
125 pub async fn evaluate(
126 &self,
127 registration: &PinnedRegistration,
128 evidence: &FrozenEvidence,
129 options: RunOptions<'_>,
130 cancelled: &AtomicBool,
131 ) -> Result<AttemptOutcome, String> {
132 if evidence.request.params.binding.registration_digest != registration.digest()
133 || evidence.request.params.gate_id != registration.declaration.name
134 {
135 return Err("registration changed before execution".into());
136 }
137 let deadline = chrono::DateTime::parse_from_rfc3339(&evidence.request.params.deadline)
138 .map_err(|_| "invalid deadline")?;
139 let remaining = (deadline.with_timezone(&chrono::Utc) - chrono::Utc::now())
140 .to_std()
141 .map_err(|_| "evaluation deadline expired")?;
142 let wall = remaining.min(Duration::from_millis(
143 evidence.request.params.limits.wall_time_ms,
144 ));
145 let deadline = tokio::time::Instant::now() + wall;
146 let parent = options
147 .attempt_parent
148 .canonicalize()
149 .map_err(|e| e.to_string())?;
150 let name = format!("kranz-evaluator-{}", uuid::Uuid::new_v4().simple());
151 let root = parent.join(&name);
152 std::fs::DirBuilder::new()
153 .mode(0o700)
154 .create(&root)
155 .map_err(|e| e.to_string())?;
156 let mount_nonce = uuid::Uuid::new_v4().simple().to_string();
157 let mut guard = ContainerGuard {
158 client: self.clone(),
159 name: name.clone(),
160 armed: false,
161 creation_finished: false,
162 };
163 let execution = async {
164 prepare(&root, &mount_nonce, registration, evidence)?;
165 let pin = ®istration.declaration.image;
166 let image = self
167 .control(&["image".into(), "inspect".into(), pin.clone()])
168 .await?;
169 if image.code != Some(0) {
170 return Err("pinned evaluator image is not installed locally (automatic pulls are disabled)".into());
171 }
172 let metadata: serde_json::Value =
173 crate::strict_json::parse(&image.stdout).map_err(|_| "invalid image inspection")?;
174 let config = metadata.get(0).ok_or("missing image inspection")?;
175 if !config["RepoDigests"]
176 .as_array()
177 .is_some_and(|values| values.iter().any(|d| d.as_str() == Some(pin)))
178 || config["Config"]["Volumes"]
179 .as_object()
180 .is_some_and(|v| !v.is_empty())
181 {
182 return Err(
183 "image digest is unresolved or the image declares extra writable volumes"
184 .into(),
185 );
186 }
187 let args = create_args(&name, &root, &mount_nonce, registration)?;
188 private_write(
191 &root.join("container.json"),
192 &serde_json::to_vec(
193 &serde_json::json!({"name":name,"image":pin,"docker":self.program}),
194 )
195 .map_err(|e| e.to_string())?,
196 )?;
197 guard.armed = true;
198 let create = self
202 .bounded_control(
203 &args,
204 deadline.saturating_duration_since(tokio::time::Instant::now()),
205 cancelled,
206 )
207 .await?;
208 let id = std::str::from_utf8(&create.stdout)
209 .map_err(|_| "invalid container ID")?
210 .trim();
211 if create.code != Some(0)
212 || id.len() != 64
213 || !id.bytes().all(|b| b.is_ascii_hexdigit())
214 {
215 return Err("could not create the contained evaluator".into());
216 }
217 guard.creation_finished = true;
218 let mut input = serde_json::to_vec(&evidence.request).map_err(|e| e.to_string())?;
219 input.push(b'\n');
220 private_write(&root.join("request.ndjson"), &input)?;
221 let limits = &evidence.request.params.limits;
222 let output = self
223 .command(
224 &[
225 "start".into(),
226 "--attach".into(),
227 "--interactive".into(),
228 name.clone(),
229 ],
230 &input,
231 deadline.saturating_duration_since(tokio::time::Instant::now()),
232 Duration::from_millis(limits.write_time_ms),
233 (
234 limits.max_stdout_bytes as usize,
235 limits.max_stderr_bytes as usize,
236 ),
237 cancelled,
238 )
239 .await
240 .map_err(|error| format!("Docker evaluator start failed: {error}"))?;
241 if output.code != Some(0) {
242 return Err(format!(
243 "evaluator or its control process exited unsuccessfully ({:?}): {}",
244 output.code,
245 crate::scrub::scrub_and_truncate(
246 &String::from_utf8_lossy(&output.stderr),
247 4096
248 )
249 ));
250 }
251 let inspect = self
252 .control(&[
253 "inspect".into(),
254 "--format".into(),
255 "{{json .State}}".into(),
256 name.clone(),
257 ])
258 .await?;
259 if inspect.code != Some(0) {
260 return Err("cannot verify evaluator exit state".into());
261 }
262 let state = crate::strict_json::parse(&inspect.stdout)
263 .map_err(|_| "invalid container state")?;
264 if state["Status"] != "exited"
265 || state["Running"] != false
266 || state["OOMKilled"] != false
267 || state["ExitCode"] != 0
268 || state["Pid"] != 0
269 || state["Error"] != ""
270 {
271 return Err("evaluator container did not exit cleanly".into());
272 }
273 Ok(output)
274 };
275 let cancellation = async {
276 while !cancelled.load(std::sync::atomic::Ordering::Acquire) {
277 tokio::time::sleep(Duration::from_millis(10)).await;
278 }
279 };
280 let execution = async {
281 tokio::select! { result = execution => result, () = cancellation => Err("evaluation cancelled".to_string()) }
282 };
283 let execution = tokio::time::timeout_at(deadline, execution)
284 .await
285 .map_err(|_| "evaluation deadline expired".to_string())
286 .and_then(|r| r);
287 let cleanup = if guard.armed {
288 guard.remove().await
289 } else {
290 Ok(())
291 };
292 let evaluation = match cleanup {
293 Err(error) => Err(format!(
294 "{error}; no result accepted; recovery ledger: {}",
295 root.join("container.json").display()
296 )),
297 Ok(()) => execution.and_then(|output| {
298 if cancelled.load(std::sync::atomic::Ordering::Acquire)
299 || tokio::time::Instant::now() >= deadline
300 {
301 return Err("evaluation cancelled or deadline expired before acceptance".into());
302 }
303 if options.retain_private_inputs {
304 private_write(&root.join("raw-stdout.ndjson"), &output.stdout)?;
305 }
306 accept(
307 &root,
308 &mount_nonce,
309 evidence,
310 output,
311 options.retain_private_inputs,
312 )
313 }),
314 };
315 let receipt = match &evaluation {
318 Ok(accepted) => serde_json::to_vec(accepted),
319 Err(error) => serde_json::to_vec(&serde_json::json!({"error":crate::scrub::scrub(error),"containerRemoved":!guard.armed})),
320 }.map_err(|e| e.to_string())?;
321 private_write(&root.join("receipt.json"), &receipt)?;
322 if !options.retain_private_inputs && !guard.armed {
323 for dir in ["inputs", "checker", "driver", "outputs", "build", "home"] {
324 std::fs::remove_dir_all(root.join(dir))
325 .map_err(|e| format!("private attempt cleanup failed: {e}"))?;
326 }
327 for file in ["request.ndjson", "container.json"] {
328 match std::fs::remove_file(root.join(file)) {
329 Ok(()) => {}
330 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
331 Err(e) => return Err(e.to_string()),
332 }
333 }
334 }
335 Ok(AttemptOutcome {
336 directory: root,
337 cleanup_confirmed: !guard.armed,
338 evaluation,
339 })
340 }
341}
342
343struct ContainerGuard {
344 client: DockerEvaluator,
345 name: String,
346 armed: bool,
347 creation_finished: bool,
348}
349impl ContainerGuard {
350 async fn remove(&mut self) -> Result<(), String> {
351 let removed = self
352 .client
353 .control(&["rm".into(), "--force".into(), self.name.clone()])
354 .await?;
355 let listing = self
358 .client
359 .control(&[
360 "container".into(),
361 "ls".into(),
362 "--all".into(),
363 "--filter".into(),
364 format!("name=^/{}$", self.name),
365 "--format".into(),
366 "{{.ID}}".into(),
367 ])
368 .await?;
369 if listing.code != Some(0) || !listing.stdout.iter().all(u8::is_ascii_whitespace) {
370 return Err("container cleanup could not be confirmed".into());
371 }
372 if !self.creation_finished && removed.code != Some(0) {
373 return Err("container creation was interrupted; daemon completion and cleanup remain uncertain".into());
374 }
375 self.armed = false;
376 Ok(())
377 }
378}
379impl Drop for ContainerGuard {
380 fn drop(&mut self) {
381 if !self.armed {
382 return;
383 }
384 let client = self.client.clone();
385 let name = self.name.clone();
386 std::thread::spawn(move || {
389 let Ok(runtime) = tokio::runtime::Builder::new_current_thread()
390 .enable_all()
391 .build()
392 else {
393 return;
394 };
395 runtime.block_on(async {
396 let result = client
397 .control(&["rm".into(), "--force".into(), name.clone()])
398 .await;
399 if !matches!(result, Ok(output) if output.code == Some(0)) {
400 tracing::error!(container = %name, "evaluator cleanup needs operator recovery");
401 }
402 });
403 });
404 }
405}
406
407fn private_write(path: &Path, bytes: &[u8]) -> Result<(), String> {
408 use std::io::Write;
409 let mut file = std::fs::OpenOptions::new()
410 .write(true)
411 .create_new(true)
412 .mode(0o600)
413 .open(path)
414 .map_err(|e| e.to_string())?;
415 file.write_all(bytes).map_err(|e| e.to_string())
416}
417fn prepare(
418 root: &Path,
419 mount_nonce: &str,
420 registration: &PinnedRegistration,
421 evidence: &FrozenEvidence,
422) -> Result<(), String> {
423 for dir in ["inputs", "outputs", "build", "home", "checker", "driver"] {
424 std::fs::DirBuilder::new()
425 .mode(0o700)
426 .create(root.join(dir))
427 .map_err(|e| e.to_string())?;
428 }
429 for artifact in &evidence.manifest.artifacts {
430 let path = root.join(artifact.content.path.as_str());
431 std::fs::create_dir_all(path.parent().ok_or("input has no parent")?)
432 .map_err(|e| e.to_string())?;
433 private_write(&path, &evidence.inputs[&artifact.id])?;
434 }
435 let manifest = root.join(evidence.request.params.evidence.path.as_str());
436 std::fs::create_dir_all(manifest.parent().ok_or("manifest has no parent")?)
437 .map_err(|e| e.to_string())?;
438 private_write(&manifest, &evidence.manifest_bytes)?;
439 for file in ®istration.files {
440 let path = root.join("checker").join(file.path.as_str());
441 std::fs::create_dir_all(path.parent().ok_or("checker file has no parent")?)
442 .map_err(|e| e.to_string())?;
443 private_write(&path, &file.bytes)?;
444 if file.executable {
445 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
446 .map_err(|e| e.to_string())?;
447 }
448 }
449 for part in ["inputs", "checker"] {
450 private_write(
451 &root.join(part).join("engine-mount-proof"),
452 mount_nonce.as_bytes(),
453 )?;
454 }
455 private_write(&root.join("driver/entry.sh"), DRIVER.as_bytes())
456}
457fn create_args(
458 name: &str,
459 root: &Path,
460 mount_nonce: &str,
461 registration: &PinnedRegistration,
462) -> Result<Vec<String>, String> {
463 let mut args: Vec<String> = [
464 "create",
465 "--pull=never",
466 "--interactive",
467 "--name",
468 name,
469 "--network=none",
470 "--read-only",
471 "--cap-drop=ALL",
472 "--security-opt=no-new-privileges",
473 "--pids-limit=32",
474 "--memory=128m",
475 "--memory-swap=128m",
476 "--cpus=1",
477 "--ulimit",
478 "fsize=67108864:67108864",
479 "--no-healthcheck",
480 "--workdir=/gate",
481 "--entrypoint=/usr/bin/env",
482 ]
483 .into_iter()
484 .map(str::to_owned)
485 .collect();
486 args.extend([
487 "--user".into(),
488 format!("{}:{}", unsafe { libc::geteuid() }, unsafe {
489 libc::getegid()
490 }),
491 ]);
492 for part in ["inputs", "outputs", "build", "home", "checker", "driver"] {
493 let source = root.join(part);
494 let source = source
495 .to_str()
496 .filter(|s| !s.contains([',', '\n', '\r']))
497 .ok_or("host mount path is not representable safely")?;
498 let target = if ["checker", "driver"].contains(&part) {
499 format!("/{part}")
500 } else {
501 format!("/gate/{part}")
502 };
503 let readonly = if ["inputs", "checker", "driver"].contains(&part) {
504 ",readonly"
505 } else {
506 ""
507 };
508 args.extend([
509 "--mount".into(),
510 format!("type=bind,src={source},dst={target}{readonly}"),
511 ]);
512 }
513 args.extend([
514 registration.declaration.image.clone(),
515 "-i".into(),
516 "HOME=/gate/home".into(),
517 "TMPDIR=/gate/build".into(),
518 "PATH=/usr/local/bin:/usr/bin:/bin".into(),
519 "/bin/sh".into(),
520 "/driver/entry.sh".into(),
521 mount_nonce.into(),
522 registration.declaration.executable.clone(),
523 ]);
524 args.extend(registration.declaration.args.clone());
525 Ok(args)
526}
527fn accept(
528 root: &Path,
529 mount_nonce: &str,
530 evidence: &FrozenEvidence,
531 output: Output,
532 raw_stdout_retained: bool,
533) -> Result<AcceptedEvaluation, String> {
534 let dir = Dir::open_ambient_dir(root, ambient_authority()).map_err(|e| e.to_string())?;
535 for part in ["outputs", "build", "home"] {
536 let root = dir
537 .open_dir_nofollow(part)
538 .map_err(|_| "scratch root was replaced")?;
539 let reference = ArtifactRef {
540 path: WirePath::try_from("engine-mount-proof".to_string())?,
541 digest: Digest::of(mount_nonce.as_bytes()),
542 bytes: mount_nonce.len() as u64,
543 };
544 artifacts::import(
545 &root,
546 &[reference],
547 &Limits {
548 max_artifacts: 1,
549 max_artifact_bytes: 64,
550 ..evidence.request.params.limits.clone()
551 },
552 )?;
553 }
554 let newline = output
555 .stdout
556 .iter()
557 .position(|b| *b == b'\n')
558 .ok_or("terminal response is missing its NDJSON newline")?;
559 if newline as u64 > evidence.request.params.limits.max_frame_bytes
560 || !output.stdout[newline + 1..]
561 .iter()
562 .all(u8::is_ascii_whitespace)
563 {
564 return Err("duplicate response, trailing stdout or frame overflow".into());
565 }
566 let response = Response::from_bytes(&output.stdout[..newline])?;
567 response.correlate(&evidence.request)?;
568 let Response::Result(mut response) = response else {
569 return Err("checker did not judge".into());
570 };
571 evidence.validate_findings(&response.result)?;
572 for label in response
573 .result
574 .artifacts
575 .iter()
576 .map(|a| a.path.as_str())
577 .chain(
578 response
579 .result
580 .findings
581 .iter()
582 .flatten()
583 .map(|f| f.id.as_str()),
584 )
585 {
586 if crate::scrub::scrub(label) != label {
587 return Err("secret-shaped output identifier is not retained".into());
588 }
589 }
590 let outputs = dir
591 .open_dir_nofollow("outputs")
592 .map_err(|_| "output root was replaced")?;
593 let artifacts = artifacts::import(
594 &outputs,
595 &response.result.artifacts,
596 &evidence.request.params.limits,
597 )?;
598 response.result.rationale = crate::scrub::scrub(&response.result.rationale);
599 for finding in response.result.findings.iter_mut().flatten() {
600 finding.summary = crate::scrub::scrub(&finding.summary);
601 }
602 let retained = serde_json::to_vec(&response.result).map_err(|e| e.to_string())?;
603 Ok(AcceptedEvaluation {
604 raw_stdout_digest: Digest::of(&output.stdout),
605 raw_stdout_retained,
606 retained_result_digest: Digest::of(&retained),
607 transformation: artifacts::transformation(),
608 result: response.result,
609 artifacts,
610 diagnostics: crate::scrub::scrub(&String::from_utf8_lossy(&output.stderr)),
611 exit_code: 0,
612 container_removed: true,
613 })
614}