Skip to main content

kranz_engine/gate_evaluation/
subprocess.rs

1//! One contained process per evaluation. Docker is the trusted host control
2//! plane; only the explicitly assembled inputs/checker/scratch mounts cross
3//! into the untrusted process. Mission stages share this contained runner.
4use 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/// Host-owned Docker executable and a frozen client environment. Neither is
20/// selected from checker output or forwarded to the container.
21#[derive(Clone)]
22pub struct DockerEvaluator {
23    program: PathBuf,
24    env: HashMap<String, String>,
25}
26
27pub struct RunOptions<'a> {
28    /// An existing trusted host directory shared with the Docker daemon.
29    /// A fresh mode-0700 child is created for every attempt.
30    pub attempt_parent: &'a Path,
31    /// Explicit audit policy. False removes raw inputs/output after cleanup.
32    pub retain_private_inputs: bool,
33}
34
35#[derive(Debug, Serialize)]
36#[serde(rename_all = "camelCase")]
37pub struct AcceptedEvaluation {
38    /// Digest of all raw stdout; `result` below is scrubbed for retention.
39    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    /// Contains a scrubbed receipt; raw files remain only when explicitly
53    /// retained or when container cleanup was not confirmed.
54    pub directory: PathBuf,
55    pub cleanup_confirmed: bool,
56    pub evaluation: Result<AcceptedEvaluation, String>,
57}
58
59impl DockerEvaluator {
60    /// `program` must be a trusted installed Docker CLI, not a repo script.
61    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 = &registration.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            // A ledger and guard exist before create, so cancellation cannot
189            // forget a daemon-side object while the CLI is in flight.
190            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            // Creation can materialize a cold image's root filesystem. Charge
199            // it to the evaluation's existing deadline, not the short control
200            // timeout used for inspection/cleanup. Never retry an uncertain create.
201            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        // Only scrubbed data goes into the default retained receipt. Raw
316        // stdout is never copied from a failing process into diagnostics.
317        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        // Removal may report absent after an interrupted create. A successful
356        // daemon listing, not a failed inspect, proves absence in both cases.
357        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        // Dropping the async driver is cancellation too. This worker owns a
387        // bounded cleanup runtime independent of the cancelled caller.
388        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 &registration.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}