Skip to main content

aft/bash_background/
persistence.rs

1#[cfg(unix)]
2use std::ffi::CString;
3use std::ffi::{OsStr, OsString};
4use std::fs::{self, File, OpenOptions};
5use std::io::{self, Read, Seek, SeekFrom, Write};
6#[cfg(unix)]
7use std::os::fd::{AsRawFd, FromRawFd, RawFd};
8#[cfg(unix)]
9use std::os::unix::ffi::{OsStrExt, OsStringExt};
10#[cfg(unix)]
11use std::os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt};
12#[cfg(windows)]
13use std::os::windows::fs::{MetadataExt, OpenOptionsExt};
14#[cfg(windows)]
15use std::os::windows::io::AsRawHandle;
16use std::path::{Path, PathBuf};
17use std::sync::Arc;
18use std::time::{SystemTime, UNIX_EPOCH};
19
20use serde::{Deserialize, Serialize};
21
22use crate::backup::hash_session;
23use crate::bash_permissions::PermissionAsk;
24use crate::db::bash_tasks::BashTaskRow;
25
26use super::BgTaskStatus;
27
28pub const SCHEMA_VERSION: u32 = 6;
29const CONTROL_DIR: &str = "control";
30const IO_DIR: &str = "io";
31const METADATA_FILE: &str = "metadata.json";
32pub const COMMAND_FILE: &str = "command.sh";
33pub const WRAPPER_FILE: &str = "wrapper.sh";
34pub const ENVIRONMENT_FILE: &str = "environment.bin";
35pub const MANIFEST_FILE: &str = "manifest.blake3";
36pub const SANDBOX_PROFILE_FILE: &str = "sandbox-profile.json";
37
38#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub enum TaskLayout {
40    Flat,
41    Directory,
42}
43
44#[derive(Debug, Clone)]
45pub struct TaskPaths {
46    pub layout: TaskLayout,
47    pub task_id: String,
48    pub session_dir: PathBuf,
49    /// Root directory for this task's persisted artifacts; legacy flat-layout tasks
50    /// use the session directory instead of a per-task directory.
51    pub dir: PathBuf,
52    pub control_dir: PathBuf,
53    pub io_dir: PathBuf,
54    pub json: PathBuf,
55    pub stdout: PathBuf,
56    pub stderr: PathBuf,
57    pub exit: PathBuf,
58    pub pty: PathBuf,
59    pub sandbox_unavailable: PathBuf,
60    pub command: PathBuf,
61    pub wrapper: PathBuf,
62    pub environment: PathBuf,
63    pub manifest: PathBuf,
64    pub sandbox_profile: PathBuf,
65}
66
67impl TaskPaths {
68    fn directory(session_dir: PathBuf, task_id: &str) -> Self {
69        let dir = session_dir.join(task_id);
70        let control_dir = dir.join(CONTROL_DIR);
71        let io_dir = dir.join(IO_DIR);
72        Self {
73            layout: TaskLayout::Directory,
74            task_id: task_id.to_string(),
75            session_dir,
76            dir,
77            json: control_dir.join(METADATA_FILE),
78            stdout: io_dir.join(TaskArtifact::Stdout.file_name()),
79            stderr: io_dir.join(TaskArtifact::Stderr.file_name()),
80            exit: io_dir.join(TaskArtifact::Exit.file_name()),
81            pty: io_dir.join(TaskArtifact::Pty.file_name()),
82            sandbox_unavailable: io_dir.join(TaskArtifact::SandboxUnavailable.file_name()),
83            command: control_dir.join(COMMAND_FILE),
84            wrapper: control_dir.join(WRAPPER_FILE),
85            environment: control_dir.join(ENVIRONMENT_FILE),
86            manifest: control_dir.join(MANIFEST_FILE),
87            sandbox_profile: control_dir.join(SANDBOX_PROFILE_FILE),
88            control_dir,
89            io_dir,
90        }
91    }
92
93    fn flat(session_dir: PathBuf, task_id: &str) -> Self {
94        let prefix = |extension: &str| session_dir.join(format!("{task_id}.{extension}"));
95        Self {
96            layout: TaskLayout::Flat,
97            task_id: task_id.to_string(),
98            dir: session_dir.clone(),
99            control_dir: session_dir.clone(),
100            io_dir: session_dir.clone(),
101            json: prefix("json"),
102            stdout: prefix("stdout"),
103            stderr: prefix("stderr"),
104            exit: prefix("exit"),
105            pty: prefix("pty"),
106            sandbox_unavailable: prefix("sandbox-unavailable"),
107            command: prefix("sh"),
108            wrapper: prefix("wrapper.sh"),
109            environment: prefix("env"),
110            manifest: prefix("manifest"),
111            sandbox_profile: prefix("sandbox-profile.json"),
112            session_dir,
113        }
114    }
115
116    pub fn artifact_path(&self, artifact: TaskArtifact) -> &Path {
117        match artifact {
118            TaskArtifact::Stdout => &self.stdout,
119            TaskArtifact::Stderr => &self.stderr,
120            TaskArtifact::Exit => &self.exit,
121            TaskArtifact::Pty => &self.pty,
122            TaskArtifact::SandboxUnavailable => &self.sandbox_unavailable,
123        }
124    }
125
126    fn artifact_name(&self, artifact: TaskArtifact) -> OsString {
127        match self.layout {
128            TaskLayout::Directory => OsString::from(artifact.file_name()),
129            TaskLayout::Flat => {
130                OsString::from(format!("{}.{}", self.task_id, artifact.flat_extension()))
131            }
132        }
133    }
134}
135
136#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
137pub enum TaskArtifact {
138    Stdout,
139    Stderr,
140    Exit,
141    Pty,
142    SandboxUnavailable,
143}
144
145impl TaskArtifact {
146    pub const ALL: [Self; 5] = [
147        Self::Stdout,
148        Self::Stderr,
149        Self::Exit,
150        Self::Pty,
151        Self::SandboxUnavailable,
152    ];
153
154    pub fn file_name(self) -> &'static str {
155        match self {
156            Self::Stdout => "stdout",
157            Self::Stderr => "stderr",
158            Self::Exit => "exit",
159            Self::Pty => "pty",
160            Self::SandboxUnavailable => "sandbox-unavailable",
161        }
162    }
163
164    fn flat_extension(self) -> &'static str {
165        match self {
166            Self::Stdout => "stdout",
167            Self::Stderr => "stderr",
168            Self::Exit => "exit",
169            Self::Pty => "pty",
170            Self::SandboxUnavailable => "sandbox-unavailable",
171        }
172    }
173}
174
175#[derive(Debug)]
176pub struct PinnedDir {
177    file: File,
178    path: PathBuf,
179}
180
181impl PinnedDir {
182    pub fn open(path: &Path) -> io::Result<Self> {
183        #[cfg(unix)]
184        let file = OpenOptions::new()
185            .read(true)
186            .custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC)
187            .open(path)?;
188        #[cfg(windows)]
189        let file = OpenOptions::new()
190            .read(true)
191            .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT | FILE_FLAG_BACKUP_SEMANTICS)
192            .open(path)?;
193        validate_directory_handle(&file)?;
194        Ok(Self {
195            file,
196            path: path.to_path_buf(),
197        })
198    }
199
200    pub fn path(&self) -> &Path {
201        &self.path
202    }
203
204    fn modified(&self) -> io::Result<SystemTime> {
205        self.file.metadata()?.modified()
206    }
207
208    fn same_identity(&self, other: &Self) -> io::Result<bool> {
209        #[cfg(unix)]
210        {
211            let left = self.file.metadata()?;
212            let right = other.file.metadata()?;
213            Ok(left.dev() == right.dev() && left.ino() == right.ino())
214        }
215        #[cfg(windows)]
216        {
217            let left = windows_file_information(&self.file)?;
218            let right = windows_file_information(&other.file)?;
219            Ok(left.volume_serial_number == right.volume_serial_number
220                && left.file_index_high == right.file_index_high
221                && left.file_index_low == right.file_index_low)
222        }
223    }
224
225    #[cfg(unix)]
226    fn open_dir_at(&self, name: &OsStr) -> io::Result<Self> {
227        let file = openat_file(
228            self.file.as_raw_fd(),
229            name,
230            libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC,
231            0,
232        )?;
233        validate_directory_handle(&file)?;
234        Ok(Self {
235            file,
236            path: self.path.join(name),
237        })
238    }
239
240    #[cfg(windows)]
241    fn open_dir_at(&self, name: &OsStr) -> io::Result<Self> {
242        self.ensure_current_identity()?;
243        let child = Self::open(&self.path.join(name))?;
244        self.ensure_current_identity()?;
245        Ok(child)
246    }
247
248    #[cfg(unix)]
249    fn create_dir_at(&self, name: &OsStr) -> io::Result<Self> {
250        let name = os_cstring(name)?;
251        let result = unsafe { libc::mkdirat(self.file.as_raw_fd(), name.as_ptr(), 0o700) };
252        if result != 0 {
253            return Err(io::Error::last_os_error());
254        }
255        self.open_dir_at(OsStr::from_bytes(name.as_bytes()))
256    }
257
258    #[cfg(windows)]
259    fn create_dir_at(&self, name: &OsStr) -> io::Result<Self> {
260        self.ensure_current_identity()?;
261        let path = self.path.join(name);
262        fs::create_dir(&path)?;
263        self.ensure_current_identity()?;
264        Self::open(&path)
265    }
266
267    pub fn open_new_file(&self, name: &OsStr) -> io::Result<File> {
268        #[cfg(unix)]
269        let file = openat_file(
270            self.file.as_raw_fd(),
271            name,
272            libc::O_RDWR | libc::O_CREAT | libc::O_EXCL | libc::O_NOFOLLOW | libc::O_CLOEXEC,
273            0o600,
274        )?;
275        #[cfg(windows)]
276        self.ensure_current_identity()?;
277        #[cfg(windows)]
278        let file = OpenOptions::new()
279            .read(true)
280            .write(true)
281            .create_new(true)
282            .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
283            .open(self.path.join(name))?;
284        validate_regular_handle(&file)?;
285        #[cfg(windows)]
286        self.ensure_current_identity()?;
287        Ok(file)
288    }
289
290    pub fn open_file(&self, name: &OsStr, write: bool) -> io::Result<File> {
291        #[cfg(unix)]
292        let file = openat_file(
293            self.file.as_raw_fd(),
294            name,
295            (if write {
296                libc::O_RDWR
297            } else {
298                libc::O_RDONLY | libc::O_NONBLOCK
299            }) | libc::O_NOFOLLOW
300                | libc::O_CLOEXEC,
301            0,
302        )?;
303        #[cfg(windows)]
304        self.ensure_current_identity()?;
305        #[cfg(windows)]
306        let file = OpenOptions::new()
307            .read(true)
308            .write(write)
309            .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT)
310            .open(self.path.join(name))?;
311        validate_regular_handle(&file)?;
312        #[cfg(unix)]
313        if !write {
314            clear_nonblocking(&file)?;
315        }
316        #[cfg(windows)]
317        self.ensure_current_identity()?;
318        Ok(file)
319    }
320
321    pub fn list_names(&self) -> io::Result<Vec<OsString>> {
322        #[cfg(unix)]
323        {
324            let dot = b".\0";
325            let fresh = unsafe {
326                libc::openat(
327                    self.file.as_raw_fd(),
328                    dot.as_ptr().cast(),
329                    libc::O_RDONLY | libc::O_DIRECTORY | libc::O_CLOEXEC | libc::O_NOFOLLOW,
330                )
331            };
332            if fresh < 0 {
333                return Err(io::Error::last_os_error());
334            }
335            let directory = unsafe { libc::fdopendir(fresh) };
336            if directory.is_null() {
337                let error = io::Error::last_os_error();
338                unsafe { libc::close(fresh) };
339                return Err(error);
340            }
341            let mut names = Vec::new();
342            loop {
343                let entry = unsafe { libc::readdir(directory) };
344                if entry.is_null() {
345                    break;
346                }
347                let bytes = unsafe {
348                    std::ffi::CStr::from_ptr((*entry).d_name.as_ptr())
349                        .to_bytes()
350                        .to_vec()
351                };
352                if bytes != b"." && bytes != b".." {
353                    names.push(OsString::from_vec(bytes));
354                }
355            }
356            unsafe { libc::closedir(directory) };
357            Ok(names)
358        }
359        #[cfg(windows)]
360        {
361            self.ensure_current_identity()?;
362            let names = fs::read_dir(&self.path)?
363                .map(|entry| entry.map(|entry| entry.file_name()))
364                .collect::<io::Result<Vec<_>>>()?;
365            self.ensure_current_identity()?;
366            Ok(names)
367        }
368    }
369
370    fn rename(&self, from: &OsStr, to: &OsStr) -> io::Result<()> {
371        self.rename_to(from, self, to)
372    }
373
374    fn rename_to(&self, from: &OsStr, target: &PinnedDir, to: &OsStr) -> io::Result<()> {
375        #[cfg(unix)]
376        {
377            let from = os_cstring(from)?;
378            let to = os_cstring(to)?;
379            let result = unsafe {
380                libc::renameat(
381                    self.file.as_raw_fd(),
382                    from.as_ptr(),
383                    target.file.as_raw_fd(),
384                    to.as_ptr(),
385                )
386            };
387            if result != 0 {
388                return Err(io::Error::last_os_error());
389            }
390            Ok(())
391        }
392        #[cfg(windows)]
393        {
394            self.ensure_current_identity()?;
395            target.ensure_current_identity()?;
396            fs::rename(self.path.join(from), target.path.join(to))?;
397            self.ensure_current_identity()?;
398            target.ensure_current_identity()
399        }
400    }
401
402    #[cfg(windows)]
403    fn ensure_current_identity(&self) -> io::Result<()> {
404        let current = OpenOptions::new()
405            .read(true)
406            .custom_flags(FILE_FLAG_OPEN_REPARSE_POINT | FILE_FLAG_BACKUP_SEMANTICS)
407            .open(&self.path)?;
408        validate_directory_handle(&current)?;
409        let held = windows_file_information(&self.file)?;
410        let observed = windows_file_information(&current)?;
411        if held.volume_serial_number != observed.volume_serial_number
412            || held.file_index_high != observed.file_index_high
413            || held.file_index_low != observed.file_index_low
414        {
415            return Err(io::Error::new(
416                io::ErrorKind::PermissionDenied,
417                "pinned directory path identity changed",
418            ));
419        }
420        Ok(())
421    }
422
423    fn remove_file(&self, name: &OsStr) -> io::Result<()> {
424        #[cfg(unix)]
425        {
426            let name = os_cstring(name)?;
427            let result = unsafe { libc::unlinkat(self.file.as_raw_fd(), name.as_ptr(), 0) };
428            if result != 0 {
429                return Err(io::Error::last_os_error());
430            }
431            Ok(())
432        }
433        #[cfg(windows)]
434        {
435            self.ensure_current_identity()?;
436            fs::remove_file(self.path.join(name))?;
437            self.ensure_current_identity()
438        }
439    }
440}
441
442#[derive(Debug)]
443pub struct TaskDirs {
444    pub session: Arc<PinnedDir>,
445    pub task: Arc<PinnedDir>,
446    pub control: Arc<PinnedDir>,
447    pub io: Arc<PinnedDir>,
448}
449
450impl Clone for TaskDirs {
451    fn clone(&self) -> Self {
452        Self {
453            session: Arc::clone(&self.session),
454            task: Arc::clone(&self.task),
455            control: Arc::clone(&self.control),
456            io: Arc::clone(&self.io),
457        }
458    }
459}
460
461#[derive(Debug)]
462pub struct ResolvedTask {
463    pub paths: TaskPaths,
464    pub dirs: TaskDirs,
465}
466
467#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
468#[serde(rename_all = "lowercase")]
469pub enum BgMode {
470    #[default]
471    Pipes,
472    Pty,
473}
474
475#[derive(Debug, Clone, Serialize, Deserialize)]
476pub struct PersistedTask {
477    pub schema_version: u32,
478    pub task_id: String,
479    pub session_id: String,
480    pub command: String,
481    #[serde(default)]
482    pub mode: BgMode,
483    pub workdir: PathBuf,
484    #[serde(default)]
485    pub project_root: Option<PathBuf>,
486    pub status: BgTaskStatus,
487    pub started_at: u64,
488    pub finished_at: Option<u64>,
489    pub duration_ms: Option<u64>,
490    pub timeout_ms: Option<u64>,
491    pub exit_code: Option<i32>,
492    pub child_pid: Option<u32>,
493    pub pgid: Option<i32>,
494    pub completion_delivered: bool,
495    #[serde(default = "default_notify_on_completion")]
496    pub notify_on_completion: bool,
497    #[serde(default = "default_compressed")]
498    pub compressed: bool,
499    #[serde(default)]
500    pub pty_rows: Option<u16>,
501    #[serde(default)]
502    pub pty_cols: Option<u16>,
503    #[serde(default, skip_serializing_if = "Vec::is_empty")]
504    pub scanner_report: Vec<PermissionAsk>,
505    #[serde(default)]
506    pub sandbox_native: bool,
507    #[serde(default, skip_serializing_if = "Option::is_none")]
508    pub sandbox_temp_dir: Option<PathBuf>,
509    pub status_reason: Option<String>,
510}
511
512fn default_notify_on_completion() -> bool {
513    true
514}
515
516fn default_compressed() -> bool {
517    true
518}
519
520#[derive(Debug, Clone, PartialEq, Eq)]
521pub enum ExitMarker {
522    Code(i32),
523    Killed,
524}
525
526impl PersistedTask {
527    #[allow(clippy::too_many_arguments)]
528    pub fn starting(
529        task_id: String,
530        session_id: String,
531        command: String,
532        workdir: PathBuf,
533        project_root: Option<PathBuf>,
534        timeout_ms: Option<u64>,
535        notify_on_completion: bool,
536        compressed: bool,
537    ) -> Self {
538        Self {
539            schema_version: SCHEMA_VERSION,
540            task_id,
541            session_id,
542            command,
543            mode: BgMode::Pipes,
544            workdir,
545            project_root,
546            status: BgTaskStatus::Starting,
547            started_at: unix_millis(),
548            finished_at: None,
549            duration_ms: None,
550            timeout_ms,
551            exit_code: None,
552            child_pid: None,
553            pgid: None,
554            completion_delivered: !notify_on_completion,
555            notify_on_completion,
556            compressed,
557            pty_rows: None,
558            pty_cols: None,
559            scanner_report: Vec::new(),
560            sandbox_native: false,
561            sandbox_temp_dir: None,
562            status_reason: None,
563        }
564    }
565
566    pub fn is_terminal(&self) -> bool {
567        self.status.is_terminal()
568    }
569
570    pub fn mark_running(&mut self, child_pid: u32, pgid: i32) {
571        self.status = BgTaskStatus::Running;
572        self.child_pid = Some(child_pid);
573        self.pgid = Some(pgid);
574    }
575
576    pub fn mark_terminal(
577        &mut self,
578        status: BgTaskStatus,
579        exit_code: Option<i32>,
580        reason: Option<String>,
581    ) {
582        let finished_at = unix_millis();
583        self.status = status;
584        self.exit_code = exit_code;
585        self.finished_at = Some(finished_at);
586        self.duration_ms = Some(finished_at.saturating_sub(self.started_at));
587        self.child_pid = None;
588        self.status_reason = reason;
589        self.completion_delivered = !self.notify_on_completion;
590    }
591
592    pub fn to_bash_task_row(
593        &self,
594        harness: &str,
595        paths: &TaskPaths,
596    ) -> Result<BashTaskRow, serde_json::Error> {
597        let project_root = self.project_root.as_deref().unwrap_or(&self.workdir);
598        let output_bytes = capture_output_bytes(&self.mode, paths);
599        let stdout_path = match self.mode {
600            BgMode::Pipes => Some(paths.stdout.display().to_string()),
601            BgMode::Pty => Some(paths.pty.display().to_string()),
602        };
603        let stderr_path = match self.mode {
604            BgMode::Pipes => Some(paths.stderr.display().to_string()),
605            BgMode::Pty => None,
606        };
607        let mut metadata = self.clone();
608        metadata.schema_version = SCHEMA_VERSION;
609        Ok(BashTaskRow {
610            harness: harness.to_string(),
611            session_id: self.session_id.clone(),
612            task_id: self.task_id.clone(),
613            project_key: crate::path_identity::project_scope_key(project_root),
614            command: self.command.clone(),
615            cwd: self.workdir.display().to_string(),
616            status: status_name(&self.status).to_string(),
617            exit_code: self.exit_code,
618            pid: self.child_pid.map(i64::from),
619            pgid: self.pgid.map(i64::from),
620            started_at: self.started_at as i64,
621            completed_at: self.finished_at.map(|value| value as i64),
622            stdout_path,
623            stderr_path,
624            compressed: self.compressed,
625            timeout_ms: self.timeout_ms.map(|value| value as i64),
626            completion_delivered: self.completion_delivered,
627            output_bytes,
628            metadata: serde_json::to_string(&metadata)?,
629        })
630    }
631}
632
633impl From<BashTaskRow> for PersistedTask {
634    fn from(row: BashTaskRow) -> Self {
635        if let Ok(task) = serde_json::from_str::<PersistedTask>(&row.metadata) {
636            return task;
637        }
638        let status = match row.status.as_str() {
639            "starting" => BgTaskStatus::Starting,
640            "running" => BgTaskStatus::Running,
641            "killing" => BgTaskStatus::Killing,
642            "completed" => BgTaskStatus::Completed,
643            "failed" => BgTaskStatus::Failed,
644            "killed" => BgTaskStatus::Killed,
645            "timed_out" => BgTaskStatus::TimedOut,
646            _ => BgTaskStatus::Failed,
647        };
648        let started_at = u64::try_from(row.started_at).unwrap_or_default();
649        let finished_at = row.completed_at.and_then(|value| u64::try_from(value).ok());
650        Self {
651            schema_version: SCHEMA_VERSION,
652            task_id: row.task_id,
653            session_id: row.session_id,
654            command: row.command,
655            mode: BgMode::Pipes,
656            workdir: PathBuf::from(row.cwd),
657            project_root: None,
658            status,
659            started_at,
660            finished_at,
661            duration_ms: finished_at.map(|finished_at| finished_at.saturating_sub(started_at)),
662            timeout_ms: row.timeout_ms.and_then(|value| u64::try_from(value).ok()),
663            exit_code: row.exit_code,
664            child_pid: row.pid.and_then(|value| u32::try_from(value).ok()),
665            pgid: row.pgid.and_then(|value| i32::try_from(value).ok()),
666            completion_delivered: row.completion_delivered,
667            notify_on_completion: !row.completion_delivered,
668            compressed: row.compressed,
669            pty_rows: None,
670            pty_cols: None,
671            scanner_report: Vec::new(),
672            sandbox_native: false,
673            sandbox_temp_dir: None,
674            status_reason: None,
675        }
676    }
677}
678
679fn status_name(status: &BgTaskStatus) -> &'static str {
680    match status {
681        BgTaskStatus::Starting => "starting",
682        BgTaskStatus::Running => "running",
683        BgTaskStatus::Killing => "killing",
684        BgTaskStatus::Completed => "completed",
685        BgTaskStatus::Failed => "failed",
686        BgTaskStatus::Killed => "killed",
687        BgTaskStatus::TimedOut => "timed_out",
688    }
689}
690
691fn capture_output_bytes(mode: &BgMode, paths: &TaskPaths) -> Option<i64> {
692    let len = |artifact| {
693        open_task_artifact(paths, artifact)
694            .ok()
695            .and_then(|file| file.len().ok())
696    };
697    match mode {
698        BgMode::Pipes => match (len(TaskArtifact::Stdout), len(TaskArtifact::Stderr)) {
699            (Some(stdout), Some(stderr)) => Some(stdout.saturating_add(stderr) as i64),
700            (Some(bytes), None) | (None, Some(bytes)) => Some(bytes as i64),
701            (None, None) => None,
702        },
703        BgMode::Pty => len(TaskArtifact::Pty).map(|bytes| bytes as i64),
704    }
705}
706
707pub fn validate_task_id(task_id: &str) -> io::Result<()> {
708    let bytes = task_id.as_bytes();
709    if bytes.len() == 21
710        && bytes.starts_with(b"bash-")
711        && bytes[5..]
712            .iter()
713            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(byte))
714    {
715        Ok(())
716    } else {
717        Err(io::Error::new(
718            io::ErrorKind::InvalidInput,
719            "background task id must match ^bash-[0-9a-f]{16}$",
720        ))
721    }
722}
723
724pub fn session_tasks_dir(storage_dir: &Path, session_id: &str) -> PathBuf {
725    let session_hash = hash_session(session_id);
726    let direct = storage_dir.join("bash-tasks").join(&session_hash);
727    if direct.exists() {
728        return direct;
729    }
730    let mut harness_matches = ["opencode", "pi"]
731        .into_iter()
732        .map(|harness| {
733            storage_dir
734                .join(harness)
735                .join("bash-tasks")
736                .join(&session_hash)
737        })
738        .filter(|path| path.exists())
739        .collect::<Vec<_>>();
740    if harness_matches.len() == 1 {
741        return harness_matches.remove(0);
742    }
743    direct
744}
745
746pub fn task_paths(storage_dir: &Path, session_id: &str, task_id: &str) -> io::Result<TaskPaths> {
747    validate_task_id(task_id)?;
748    Ok(TaskPaths::flat(
749        session_tasks_dir(storage_dir, session_id),
750        task_id,
751    ))
752}
753
754pub fn allocate_task_layout(storage_dir: &Path, session_id: &str) -> io::Result<ResolvedTask> {
755    let session_dir = session_tasks_dir(storage_dir, session_id);
756    create_private_task_store(&session_dir)?;
757    let session = Arc::new(PinnedDir::open(&session_dir)?);
758    for _ in 0..32 {
759        let task_id = random_task_id()?;
760        match create_task_layout_from_session(Arc::clone(&session), &task_id) {
761            Ok(task) => return Ok(task),
762            Err(error) if error.kind() == io::ErrorKind::AlreadyExists => continue,
763            Err(error) => return Err(error),
764        }
765    }
766    Err(io::Error::new(
767        io::ErrorKind::AlreadyExists,
768        "failed to allocate unique background task id after 32 attempts",
769    ))
770}
771
772pub fn create_task_layout(
773    storage_dir: &Path,
774    session_id: &str,
775    task_id: &str,
776) -> io::Result<ResolvedTask> {
777    validate_task_id(task_id)?;
778    let session_dir = session_tasks_dir(storage_dir, session_id);
779    create_private_task_store(&session_dir)?;
780    create_task_layout_from_session(Arc::new(PinnedDir::open(&session_dir)?), task_id)
781}
782
783fn create_private_task_store(session_dir: &Path) -> io::Result<()> {
784    fs::create_dir_all(session_dir)?;
785    #[cfg(unix)]
786    {
787        let parent = session_dir.parent().ok_or_else(|| {
788            io::Error::new(
789                io::ErrorKind::InvalidInput,
790                "session task directory has no bash-tasks parent",
791            )
792        })?;
793        fs::set_permissions(parent, fs::Permissions::from_mode(0o700))?;
794        fs::set_permissions(session_dir, fs::Permissions::from_mode(0o700))?;
795    }
796    Ok(())
797}
798
799fn create_task_layout_from_session(
800    session: Arc<PinnedDir>,
801    task_id: &str,
802) -> io::Result<ResolvedTask> {
803    validate_task_id(task_id)?;
804    let task = session.create_dir_at(OsStr::new(task_id))?;
805    let control = task.create_dir_at(OsStr::new(CONTROL_DIR))?;
806    let io_dir = task.create_dir_at(OsStr::new(IO_DIR))?;
807    let paths = TaskPaths::directory(session.path.clone(), task_id);
808    Ok(ResolvedTask {
809        paths,
810        dirs: TaskDirs {
811            session,
812            task: Arc::new(task),
813            control: Arc::new(control),
814            io: Arc::new(io_dir),
815        },
816    })
817}
818
819pub fn resolve_task_layout(session_dir: &Path, task_id: &str) -> io::Result<ResolvedTask> {
820    let task = resolve_uninitialized_task_layout(session_dir, task_id)?;
821    let metadata = read_task_at(&task)?;
822    if metadata.task_id != task_id {
823        return Err(io::Error::new(
824            io::ErrorKind::InvalidData,
825            "background task metadata identity mismatch",
826        ));
827    }
828    Ok(task)
829}
830
831pub fn resolve_uninitialized_task_layout(
832    session_dir: &Path,
833    task_id: &str,
834) -> io::Result<ResolvedTask> {
835    validate_task_id(task_id)?;
836    let session = Arc::new(PinnedDir::open(session_dir)?);
837    let directory = session.open_dir_at(OsStr::new(task_id));
838    let flat_name = OsString::from(format!("{task_id}.json"));
839    let flat = session.open_file(&flat_name, false);
840    let has_directory = match &directory {
841        Ok(_) => true,
842        Err(error) if error.kind() == io::ErrorKind::NotFound => false,
843        Err(error) => {
844            return Err(io::Error::new(
845                error.kind(),
846                format!("invalid task directory: {error}"),
847            ))
848        }
849    };
850    let has_flat = match &flat {
851        Ok(_) => true,
852        Err(error) if error.kind() == io::ErrorKind::NotFound => false,
853        Err(error) => {
854            return Err(io::Error::new(
855                error.kind(),
856                format!("invalid flat task metadata: {error}"),
857            ))
858        }
859    };
860    if has_directory && has_flat {
861        return Err(io::Error::new(
862            io::ErrorKind::InvalidData,
863            "duplicate flat and directory background task layouts",
864        ));
865    }
866    if has_directory {
867        let task = directory.expect("directory result checked");
868        let control = task.open_dir_at(OsStr::new(CONTROL_DIR))?;
869        let io_dir = task.open_dir_at(OsStr::new(IO_DIR))?;
870        let paths = TaskPaths::directory(session_dir.to_path_buf(), task_id);
871        return Ok(ResolvedTask {
872            paths,
873            dirs: TaskDirs {
874                session,
875                task: Arc::new(task),
876                control: Arc::new(control),
877                io: Arc::new(io_dir),
878            },
879        });
880    }
881    if has_flat {
882        let paths = TaskPaths::flat(session_dir.to_path_buf(), task_id);
883        return Ok(ResolvedTask {
884            paths,
885            dirs: TaskDirs {
886                session: Arc::clone(&session),
887                task: Arc::clone(&session),
888                control: Arc::clone(&session),
889                io: session,
890            },
891        });
892    }
893    Err(io::Error::new(
894        io::ErrorKind::NotFound,
895        "background task layout not found",
896    ))
897}
898
899pub fn resolve_task(
900    storage_dir: &Path,
901    session_id: &str,
902    task_id: &str,
903) -> io::Result<ResolvedTask> {
904    resolve_task_layout(&session_tasks_dir(storage_dir, session_id), task_id)
905}
906
907pub fn discover_task_ids(session_dir: &Path) -> io::Result<(Vec<String>, Vec<OsString>)> {
908    let session = PinnedDir::open(session_dir)?;
909    let mut ids = std::collections::BTreeSet::new();
910    let mut invalid = Vec::new();
911    for name in session.list_names()? {
912        let Some(text) = name.to_str() else {
913            invalid.push(name);
914            continue;
915        };
916        if validate_task_id(text).is_ok() {
917            ids.insert(text.to_string());
918            continue;
919        }
920        if let Some((task_id, _suffix)) = text.split_once('.') {
921            if validate_task_id(task_id).is_ok() {
922                ids.insert(task_id.to_string());
923            } else if task_id.starts_with("bash-") {
924                invalid.push(name);
925            }
926        } else if text.starts_with("bash-") {
927            invalid.push(name);
928        }
929    }
930    Ok((ids.into_iter().collect(), invalid))
931}
932
933pub fn uninitialized_layout_is_recent(
934    session_dir: &Path,
935    task_id: &str,
936    grace: std::time::Duration,
937) -> io::Result<bool> {
938    let task = resolve_uninitialized_task_layout(session_dir, task_id)?;
939    let modified = match task.paths.layout {
940        TaskLayout::Directory => task.dirs.control.modified()?,
941        TaskLayout::Flat => task
942            .dirs
943            .session
944            .open_file(&task.paths.artifact_name(TaskArtifact::Exit), false)
945            .and_then(|file| file.metadata()?.modified())
946            .or_else(|_| {
947                task.dirs
948                    .session
949                    .open_file(OsStr::new(&format!("{task_id}.json")), false)
950                    .and_then(|file| file.metadata()?.modified())
951            })?,
952    };
953    Ok(SystemTime::now()
954        .duration_since(modified)
955        .unwrap_or_default()
956        < grace)
957}
958
959pub fn quarantine_task_layout(
960    storage_dir: &Path,
961    session_dir: &Path,
962    task_id: &str,
963    reason: &str,
964) -> io::Result<()> {
965    validate_task_id(task_id)?;
966    let session = PinnedDir::open(session_dir)?;
967    let names = session.list_names()?;
968    let flat_prefix = format!("{task_id}.");
969    let selected = names
970        .into_iter()
971        .filter(|name| {
972            name == OsStr::new(task_id)
973                || name
974                    .to_str()
975                    .is_some_and(|name| name.starts_with(&flat_prefix))
976        })
977        .collect::<Vec<_>>();
978    quarantine_names(storage_dir, session_dir, &session, selected, reason)
979}
980
981pub fn quarantine_invalid_entry(
982    storage_dir: &Path,
983    session_dir: &Path,
984    entry: &OsStr,
985) -> io::Result<()> {
986    let session = PinnedDir::open(session_dir)?;
987    quarantine_names(
988        storage_dir,
989        session_dir,
990        &session,
991        vec![entry.to_os_string()],
992        "invalid",
993    )
994}
995
996fn quarantine_names(
997    storage_dir: &Path,
998    session_dir: &Path,
999    session: &PinnedDir,
1000    names: Vec<OsString>,
1001    reason: &str,
1002) -> io::Result<()> {
1003    if names.is_empty() {
1004        return Ok(());
1005    }
1006    let session_hash = session_dir.file_name().ok_or_else(|| {
1007        io::Error::new(io::ErrorKind::InvalidInput, "session dir has no identity")
1008    })?;
1009    let quarantine_path = storage_dir.join("bash-tasks-quarantine").join(session_hash);
1010    fs::create_dir_all(&quarantine_path)?;
1011    let quarantine = PinnedDir::open(&quarantine_path)?;
1012    for name in names {
1013        let mut random = [0_u8; 8];
1014        getrandom::fill(&mut random).map_err(io::Error::other)?;
1015        let target = OsString::from(format!(
1016            "{}.{}-{}",
1017            name.to_string_lossy(),
1018            reason,
1019            hex_lower(&random)
1020        ));
1021        session.rename_to(&name, &quarantine, &target)?;
1022    }
1023    Ok(())
1024}
1025
1026fn hex_lower(bytes: &[u8]) -> String {
1027    bytes.iter().map(|byte| format!("{byte:02x}")).collect()
1028}
1029
1030pub fn read_task(path: &Path) -> io::Result<PersistedTask> {
1031    let mut file = open_validated_path(path, false)?;
1032    read_task_file(&mut file)
1033}
1034
1035pub fn read_task_at(task: &ResolvedTask) -> io::Result<PersistedTask> {
1036    let name = match task.paths.layout {
1037        TaskLayout::Directory => OsString::from(METADATA_FILE),
1038        TaskLayout::Flat => OsString::from(format!("{}.json", task.paths.task_id)),
1039    };
1040    let mut file = task.dirs.control.open_file(&name, false)?;
1041    let metadata = read_task_file(&mut file)?;
1042    if metadata.task_id != task.paths.task_id {
1043        return Err(io::Error::new(
1044            io::ErrorKind::InvalidData,
1045            "background task metadata identity does not match its layout name",
1046        ));
1047    }
1048    Ok(metadata)
1049}
1050
1051fn read_task_file(file: &mut File) -> io::Result<PersistedTask> {
1052    file.seek(SeekFrom::Start(0))?;
1053    let mut content = String::new();
1054    file.read_to_string(&mut content)?;
1055    let task: PersistedTask = serde_json::from_str(&content).map_err(io::Error::other)?;
1056    if !matches!(task.schema_version, 2 | 3 | 4 | 5 | SCHEMA_VERSION) {
1057        return Err(io::Error::new(
1058            io::ErrorKind::InvalidData,
1059            format!(
1060                "unsupported background task schema_version {} (expected 2, 3, 4, 5, or {SCHEMA_VERSION})",
1061                task.schema_version
1062            ),
1063        ));
1064    }
1065    validate_task_id(&task.task_id)?;
1066    Ok(task)
1067}
1068
1069pub fn write_task(path: &Path, task: &PersistedTask) -> io::Result<()> {
1070    validate_task_id(&task.task_id)?;
1071    if let Some(parent) = path.parent() {
1072        fs::create_dir_all(parent)?;
1073    }
1074    let parent = path.parent().unwrap_or_else(|| Path::new("."));
1075    let dir = PinnedDir::open(parent)?;
1076    let name = path
1077        .file_name()
1078        .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "metadata path has no name"))?;
1079    write_task_in_dir(&dir, name, task)
1080}
1081
1082pub fn write_task_at(task: &ResolvedTask, metadata: &PersistedTask) -> io::Result<()> {
1083    if metadata.task_id != task.paths.task_id {
1084        return Err(io::Error::new(
1085            io::ErrorKind::InvalidInput,
1086            "refusing to write metadata under a different task identity",
1087        ));
1088    }
1089    let name = match task.paths.layout {
1090        TaskLayout::Directory => OsString::from(METADATA_FILE),
1091        TaskLayout::Flat => OsString::from(format!("{}.json", task.paths.task_id)),
1092    };
1093    write_task_in_dir(&task.dirs.control, &name, metadata)
1094}
1095
1096fn write_task_in_dir(dir: &PinnedDir, name: &OsStr, task: &PersistedTask) -> io::Result<()> {
1097    let mut upgraded = task.clone();
1098    upgraded.schema_version = SCHEMA_VERSION;
1099    let content = serde_json::to_vec_pretty(&upgraded).map_err(io::Error::other)?;
1100    randomized_atomic_replace(dir, name, &content)
1101}
1102
1103pub fn update_task_at<F>(task: &ResolvedTask, update: F) -> io::Result<PersistedTask>
1104where
1105    F: FnOnce(&mut PersistedTask),
1106{
1107    let mut metadata = read_task_at(task)?;
1108    let original_terminal = metadata.is_terminal();
1109    let original = metadata.clone();
1110    update(&mut metadata);
1111    metadata.schema_version = SCHEMA_VERSION;
1112    if original_terminal {
1113        let completion_delivered = metadata.completion_delivered;
1114        metadata = original;
1115        metadata.completion_delivered = completion_delivered;
1116        metadata.schema_version = SCHEMA_VERSION;
1117    }
1118    write_task_at(task, &metadata)?;
1119    Ok(metadata)
1120}
1121
1122pub fn delete_task_bundle(paths: &TaskPaths) -> io::Result<()> {
1123    validate_task_id(&paths.task_id)?;
1124    let resolved = resolve_task_layout(&paths.session_dir, &paths.task_id)?;
1125    if resolved.paths.layout != paths.layout {
1126        return Err(io::Error::new(
1127            io::ErrorKind::InvalidData,
1128            "background task layout changed before deletion",
1129        ));
1130    }
1131    delete_resolved_task(&resolved)
1132}
1133
1134pub fn delete_resolved_task(task: &ResolvedTask) -> io::Result<()> {
1135    validate_task_id(&task.paths.task_id)?;
1136    match task.paths.layout {
1137        TaskLayout::Flat => {
1138            for path in task_bundle_files(&task.paths) {
1139                let Some(name) = path.file_name() else {
1140                    continue;
1141                };
1142                match task.dirs.session.remove_file(name) {
1143                    Ok(()) => {}
1144                    Err(error) if error.kind() == io::ErrorKind::NotFound => {}
1145                    Err(error) => return Err(error),
1146                }
1147            }
1148            Ok(())
1149        }
1150        TaskLayout::Directory => remove_directory_task(task),
1151    }
1152}
1153
1154fn remove_directory_task(task: &ResolvedTask) -> io::Result<()> {
1155    let current = task
1156        .dirs
1157        .session
1158        .open_dir_at(OsStr::new(&task.paths.task_id))?;
1159    if !current.same_identity(&task.dirs.task)? {
1160        return Err(io::Error::new(
1161            io::ErrorKind::PermissionDenied,
1162            "task directory identity changed before deletion",
1163        ));
1164    }
1165    #[cfg(unix)]
1166    let tombstone = rename_task_to_tombstone(task)?;
1167    for name in task.dirs.control.list_names()? {
1168        task.dirs.control.remove_file(&name)?;
1169    }
1170    remove_tree_contents(&task.dirs.io)?;
1171    #[cfg(unix)]
1172    {
1173        remove_dir_entry(&task.dirs.task, IO_DIR)?;
1174        remove_dir_entry(&task.dirs.task, CONTROL_DIR)?;
1175        let name = os_cstring(&tombstone)?;
1176        let result = unsafe {
1177            libc::unlinkat(
1178                task.dirs.session.file.as_raw_fd(),
1179                name.as_ptr(),
1180                libc::AT_REMOVEDIR,
1181            )
1182        };
1183        if result != 0 {
1184            return Err(io::Error::last_os_error());
1185        }
1186        Ok(())
1187    }
1188    #[cfg(windows)]
1189    {
1190        task.dirs.io.ensure_current_identity()?;
1191        task.dirs.control.ensure_current_identity()?;
1192        task.dirs.session.ensure_current_identity()?;
1193        fs::remove_dir(task.dirs.io.path())?;
1194        fs::remove_dir(task.dirs.control.path())?;
1195        fs::remove_dir(&task.paths.dir)?;
1196        task.dirs.session.ensure_current_identity()
1197    }
1198}
1199
1200fn remove_tree_contents(dir: &PinnedDir) -> io::Result<()> {
1201    for name in dir.list_names()? {
1202        match dir.open_dir_at(&name) {
1203            Ok(child) => {
1204                remove_tree_contents(&child)?;
1205                #[cfg(unix)]
1206                {
1207                    let name = os_cstring(&name)?;
1208                    let result = unsafe {
1209                        libc::unlinkat(dir.file.as_raw_fd(), name.as_ptr(), libc::AT_REMOVEDIR)
1210                    };
1211                    if result != 0 {
1212                        return Err(io::Error::last_os_error());
1213                    }
1214                }
1215                #[cfg(windows)]
1216                fs::remove_dir(child.path())?;
1217            }
1218            Err(error) if error.kind() == io::ErrorKind::NotADirectory => {
1219                dir.remove_file(&name)?;
1220            }
1221            Err(error) => return Err(error),
1222        }
1223    }
1224    Ok(())
1225}
1226
1227#[cfg(unix)]
1228fn rename_task_to_tombstone(task: &ResolvedTask) -> io::Result<OsString> {
1229    for _ in 0..32 {
1230        let tombstone = random_temp_name()?;
1231        match task.dirs.session.open_dir_at(&tombstone) {
1232            Ok(_) => continue,
1233            Err(error) if error.kind() == io::ErrorKind::NotFound => {}
1234            Err(error) => return Err(error),
1235        }
1236        task.dirs
1237            .session
1238            .rename(OsStr::new(&task.paths.task_id), &tombstone)?;
1239        let moved = task.dirs.session.open_dir_at(&tombstone)?;
1240        if !moved.same_identity(&task.dirs.task)? {
1241            return Err(io::Error::new(
1242                io::ErrorKind::PermissionDenied,
1243                "task directory identity changed during deletion",
1244            ));
1245        }
1246        return Ok(tombstone);
1247    }
1248    Err(io::Error::new(
1249        io::ErrorKind::AlreadyExists,
1250        "failed to allocate randomized task deletion name",
1251    ))
1252}
1253
1254#[cfg(unix)]
1255fn remove_dir_entry(task: &PinnedDir, child: &str) -> io::Result<()> {
1256    let child = os_cstring(OsStr::new(child))?;
1257    let result =
1258        unsafe { libc::unlinkat(task.file.as_raw_fd(), child.as_ptr(), libc::AT_REMOVEDIR) };
1259    if result != 0 {
1260        return Err(io::Error::last_os_error());
1261    }
1262    Ok(())
1263}
1264
1265pub fn task_bundle_files(paths: &TaskPaths) -> Vec<PathBuf> {
1266    if paths.layout == TaskLayout::Directory {
1267        return vec![paths.dir.clone()];
1268    }
1269    vec![
1270        paths.json.clone(),
1271        paths.stdout.clone(),
1272        paths.stderr.clone(),
1273        paths.exit.clone(),
1274        paths.pty.clone(),
1275        paths.sandbox_unavailable.clone(),
1276        paths.command.clone(),
1277        paths.wrapper.clone(),
1278        paths.environment.clone(),
1279        paths.manifest.clone(),
1280        paths.sandbox_profile.clone(),
1281        paths.dir.join(format!("{}.ps1", paths.task_id)),
1282        paths.dir.join(format!("{}.bat", paths.task_id)),
1283    ]
1284}
1285
1286pub fn write_kill_marker_if_absent(paths: &TaskPaths) -> io::Result<()> {
1287    match open_task_artifact(paths, TaskArtifact::Exit) {
1288        Ok(file) if file.len()? > 0 => Ok(()),
1289        Ok(mut file) => file.replace_contents(b"killed"),
1290        Err(error) if error.kind() == io::ErrorKind::NotFound => {
1291            let resolved = resolve_task_layout(&paths.session_dir, &paths.task_id)?;
1292            randomized_atomic_replace(
1293                &resolved.dirs.io,
1294                &resolved.paths.artifact_name(TaskArtifact::Exit),
1295                b"killed",
1296            )
1297        }
1298        Err(error) => Err(error),
1299    }
1300}
1301
1302pub fn read_exit_marker(paths: &TaskPaths) -> io::Result<Option<ExitMarker>> {
1303    let mut file = match open_task_artifact(paths, TaskArtifact::Exit) {
1304        Ok(file) => file,
1305        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
1306        Err(error) => return Err(error),
1307    };
1308    let mut content = String::new();
1309    file.read_to_string(&mut content)?;
1310    let content = content.trim();
1311    if content.is_empty() {
1312        return Ok(None);
1313    }
1314    if content == "killed" {
1315        return Ok(Some(ExitMarker::Killed));
1316    }
1317    Ok(content.parse::<i32>().ok().map(ExitMarker::Code))
1318}
1319
1320pub fn randomized_atomic_replace(dir: &PinnedDir, name: &OsStr, content: &[u8]) -> io::Result<()> {
1321    for _ in 0..32 {
1322        let temporary = random_temp_name()?;
1323        let mut file = match dir.open_new_file(&temporary) {
1324            Ok(file) => file,
1325            Err(error) if error.kind() == io::ErrorKind::AlreadyExists => continue,
1326            Err(error) => return Err(error),
1327        };
1328        let result = (|| {
1329            file.write_all(content)?;
1330            file.sync_all()?;
1331            validate_regular_handle(&file)?;
1332            dir.rename(&temporary, name)
1333        })();
1334        if result.is_err() {
1335            let _ = dir.remove_file(&temporary);
1336        }
1337        return result;
1338    }
1339    Err(io::Error::new(
1340        io::ErrorKind::AlreadyExists,
1341        "failed to allocate a randomized atomic-write name",
1342    ))
1343}
1344
1345pub fn create_control_file(dirs: &TaskDirs, name: &str, content: &[u8]) -> io::Result<File> {
1346    let mut file = dirs.control.open_new_file(OsStr::new(name))?;
1347    file.write_all(content)?;
1348    file.sync_all()?;
1349    file.seek(SeekFrom::Start(0))?;
1350    validate_regular_handle(&file)?;
1351    Ok(file)
1352}
1353
1354pub fn open_control_file(task: &ResolvedTask, name: &str) -> io::Result<File> {
1355    if name.is_empty() || name.contains('/') || name.contains('\\') || name == "." || name == ".." {
1356        return Err(io::Error::new(
1357            io::ErrorKind::InvalidInput,
1358            "invalid control file name",
1359        ));
1360    }
1361    task.dirs.control.open_file(OsStr::new(name), false)
1362}
1363
1364#[derive(Debug)]
1365pub struct ValidatedArtifact {
1366    file: File,
1367}
1368
1369impl ValidatedArtifact {
1370    fn new(file: File) -> io::Result<Self> {
1371        validate_regular_handle(&file)?;
1372        Ok(Self { file })
1373    }
1374
1375    pub fn len(&self) -> io::Result<u64> {
1376        validate_regular_handle(&self.file)?;
1377        Ok(self.file.metadata()?.len())
1378    }
1379
1380    pub fn rewind(&mut self) -> io::Result<()> {
1381        self.file.seek(SeekFrom::Start(0)).map(|_| ())
1382    }
1383
1384    pub fn tail(&mut self, max_bytes: usize) -> io::Result<(Vec<u8>, bool)> {
1385        let len = self.len()?;
1386        let read_len = len.min(max_bytes as u64);
1387        self.file
1388            .seek(SeekFrom::Start(len.saturating_sub(read_len)))?;
1389        let mut bytes = Vec::with_capacity(read_len as usize);
1390        Read::by_ref(&mut self.file)
1391            .take(read_len)
1392            .read_to_end(&mut bytes)?;
1393        Ok((bytes, len > max_bytes as u64))
1394    }
1395
1396    pub fn read_range(&mut self, start: u64, len: u64) -> io::Result<Vec<u8>> {
1397        self.file.seek(SeekFrom::Start(start))?;
1398        let mut bytes = Vec::with_capacity(len.min(usize::MAX as u64) as usize);
1399        Read::by_ref(&mut self.file)
1400            .take(len)
1401            .read_to_end(&mut bytes)?;
1402        Ok(bytes)
1403    }
1404
1405    pub fn read_all(&mut self) -> io::Result<Vec<u8>> {
1406        self.rewind()?;
1407        let mut bytes = Vec::new();
1408        self.file.read_to_end(&mut bytes)?;
1409        Ok(bytes)
1410    }
1411
1412    pub fn replace_contents(&mut self, content: &[u8]) -> io::Result<()> {
1413        validate_regular_handle(&self.file)?;
1414        self.file.set_len(0)?;
1415        self.file.seek(SeekFrom::Start(0))?;
1416        self.file.write_all(content)?;
1417        self.file.sync_all()
1418    }
1419
1420    pub fn try_clone_file(&self) -> io::Result<File> {
1421        validate_regular_handle(&self.file)?;
1422        self.file.try_clone()
1423    }
1424}
1425
1426impl Read for ValidatedArtifact {
1427    fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
1428        self.file.read(buffer)
1429    }
1430}
1431
1432impl Seek for ValidatedArtifact {
1433    fn seek(&mut self, position: SeekFrom) -> io::Result<u64> {
1434        self.file.seek(position)
1435    }
1436}
1437
1438pub fn open_task_artifact(
1439    paths: &TaskPaths,
1440    artifact: TaskArtifact,
1441) -> io::Result<ValidatedArtifact> {
1442    validate_task_id(&paths.task_id)?;
1443    let resolved = resolve_task_layout(&paths.session_dir, &paths.task_id)?;
1444    if resolved.paths.layout != paths.layout {
1445        return Err(io::Error::new(
1446            io::ErrorKind::InvalidData,
1447            "background task layout identity changed",
1448        ));
1449    }
1450    let metadata = read_task_at(&resolved)?;
1451    if metadata.task_id != paths.task_id {
1452        return Err(io::Error::new(
1453            io::ErrorKind::InvalidData,
1454            "background task metadata identity mismatch",
1455        ));
1456    }
1457    let file = resolved
1458        .dirs
1459        .io
1460        .open_file(&resolved.paths.artifact_name(artifact), false)?;
1461    ValidatedArtifact::new(file)
1462}
1463
1464pub fn replace_artifact_with_tail(
1465    paths: &TaskPaths,
1466    artifact: TaskArtifact,
1467    retain_bytes: u64,
1468) -> io::Result<u64> {
1469    let mut source = open_task_artifact(paths, artifact)?;
1470    let len = source.len()?;
1471    if len <= retain_bytes {
1472        return Ok(0);
1473    }
1474    let mut tail = source.read_range(len.saturating_sub(retain_bytes), retain_bytes)?;
1475    align_tail_start(&mut tail);
1476    let resolved = resolve_task_layout(&paths.session_dir, &paths.task_id)?;
1477    randomized_atomic_replace(
1478        &resolved.dirs.io,
1479        &resolved.paths.artifact_name(artifact),
1480        &tail,
1481    )?;
1482    Ok(len.saturating_sub(tail.len() as u64))
1483}
1484
1485#[derive(Debug)]
1486pub struct TaskIoHandles {
1487    pub dirs: TaskDirs,
1488    stdout: Option<File>,
1489    stderr: Option<File>,
1490    exit: File,
1491    pty: Option<File>,
1492    sandbox_unavailable: File,
1493}
1494
1495impl TaskIoHandles {
1496    pub fn create(task: &ResolvedTask, mode: BgMode) -> io::Result<Self> {
1497        if task.paths.layout != TaskLayout::Directory {
1498            return Err(io::Error::new(
1499                io::ErrorKind::InvalidInput,
1500                "new task output handles require the directory layout",
1501            ));
1502        }
1503        let (stdout, stderr, pty) = match mode {
1504            BgMode::Pipes => (
1505                Some(
1506                    task.dirs
1507                        .io
1508                        .open_new_file(OsStr::new(TaskArtifact::Stdout.file_name()))?,
1509                ),
1510                Some(
1511                    task.dirs
1512                        .io
1513                        .open_new_file(OsStr::new(TaskArtifact::Stderr.file_name()))?,
1514                ),
1515                None,
1516            ),
1517            BgMode::Pty => (
1518                None,
1519                None,
1520                Some(
1521                    task.dirs
1522                        .io
1523                        .open_new_file(OsStr::new(TaskArtifact::Pty.file_name()))?,
1524                ),
1525            ),
1526        };
1527        Ok(Self {
1528            dirs: task.dirs.clone(),
1529            stdout,
1530            stderr,
1531            exit: task
1532                .dirs
1533                .io
1534                .open_new_file(OsStr::new(TaskArtifact::Exit.file_name()))?,
1535            pty,
1536            sandbox_unavailable: task
1537                .dirs
1538                .io
1539                .open_new_file(OsStr::new(TaskArtifact::SandboxUnavailable.file_name()))?,
1540        })
1541    }
1542
1543    pub fn clone_file(&self, artifact: TaskArtifact) -> io::Result<File> {
1544        let file = match artifact {
1545            TaskArtifact::Stdout => self.stdout.as_ref(),
1546            TaskArtifact::Stderr => self.stderr.as_ref(),
1547            TaskArtifact::Exit => Some(&self.exit),
1548            TaskArtifact::Pty => self.pty.as_ref(),
1549            TaskArtifact::SandboxUnavailable => Some(&self.sandbox_unavailable),
1550        }
1551        .ok_or_else(|| {
1552            io::Error::new(io::ErrorKind::NotFound, "task artifact is not pre-opened")
1553        })?;
1554        validate_regular_handle(file)?;
1555        file.try_clone()
1556    }
1557
1558    #[cfg(unix)]
1559    pub fn inheritable_file(&self, artifact: TaskArtifact) -> io::Result<File> {
1560        let file = self.clone_file(artifact)?;
1561        set_close_on_exec(file.as_raw_fd(), false)?;
1562        Ok(file)
1563    }
1564
1565    pub fn write(&mut self, artifact: TaskArtifact, content: &[u8]) -> io::Result<()> {
1566        let file = match artifact {
1567            TaskArtifact::Stdout => self.stdout.as_mut(),
1568            TaskArtifact::Stderr => self.stderr.as_mut(),
1569            TaskArtifact::Exit => Some(&mut self.exit),
1570            TaskArtifact::Pty => self.pty.as_mut(),
1571            TaskArtifact::SandboxUnavailable => Some(&mut self.sandbox_unavailable),
1572        }
1573        .ok_or_else(|| {
1574            io::Error::new(io::ErrorKind::NotFound, "task artifact is not pre-opened")
1575        })?;
1576        validate_regular_handle(file)?;
1577        file.set_len(0)?;
1578        file.seek(SeekFrom::Start(0))?;
1579        file.write_all(content)?;
1580        file.sync_all()
1581    }
1582
1583    pub fn artifact_len(&self, artifact: TaskArtifact) -> io::Result<u64> {
1584        let file = match artifact {
1585            TaskArtifact::Stdout => self.stdout.as_ref(),
1586            TaskArtifact::Stderr => self.stderr.as_ref(),
1587            TaskArtifact::Exit => Some(&self.exit),
1588            TaskArtifact::Pty => self.pty.as_ref(),
1589            TaskArtifact::SandboxUnavailable => Some(&self.sandbox_unavailable),
1590        }
1591        .ok_or_else(|| {
1592            io::Error::new(io::ErrorKind::NotFound, "task artifact is not pre-opened")
1593        })?;
1594        validate_regular_handle(file)?;
1595        Ok(file.metadata()?.len())
1596    }
1597}
1598
1599pub fn repin_task_io(paths: &TaskPaths) -> io::Result<TaskDirs> {
1600    let resolved = resolve_task_layout(&paths.session_dir, &paths.task_id)?;
1601    let metadata = read_task_at(&resolved)?;
1602    if metadata.task_id != paths.task_id {
1603        return Err(io::Error::new(
1604            io::ErrorKind::InvalidData,
1605            "background task metadata identity mismatch",
1606        ));
1607    }
1608    Ok(resolved.dirs)
1609}
1610
1611pub fn unix_millis() -> u64 {
1612    SystemTime::now()
1613        .duration_since(UNIX_EPOCH)
1614        .map(|duration| duration.as_millis() as u64)
1615        .unwrap_or(0)
1616}
1617
1618fn random_task_id() -> io::Result<String> {
1619    let mut bytes = [0_u8; 8];
1620    getrandom::fill(&mut bytes).map_err(io::Error::other)?;
1621    Ok(format!(
1622        "bash-{}",
1623        bytes
1624            .iter()
1625            .map(|byte| format!("{byte:02x}"))
1626            .collect::<String>()
1627    ))
1628}
1629
1630fn random_temp_name() -> io::Result<OsString> {
1631    let mut bytes = [0_u8; 16];
1632    getrandom::fill(&mut bytes).map_err(io::Error::other)?;
1633    Ok(OsString::from(format!(
1634        ".aft-tmp-{}",
1635        bytes
1636            .iter()
1637            .map(|byte| format!("{byte:02x}"))
1638            .collect::<String>()
1639    )))
1640}
1641
1642#[cfg(test)]
1643pub(crate) fn open_unregistered_artifact(path: &Path) -> io::Result<ValidatedArtifact> {
1644    ValidatedArtifact::new(open_validated_path(path, false)?)
1645}
1646
1647#[cfg(test)]
1648pub(crate) fn replace_unregistered_with_tail(path: &Path, retain_bytes: u64) -> io::Result<u64> {
1649    let mut source = open_unregistered_artifact(path)?;
1650    let len = source.len()?;
1651    if len <= retain_bytes {
1652        return Ok(0);
1653    }
1654    let mut tail = source.read_range(len.saturating_sub(retain_bytes), retain_bytes)?;
1655    align_tail_start(&mut tail);
1656    let parent = PinnedDir::open(path.parent().unwrap_or_else(|| Path::new(".")))?;
1657    let name = path
1658        .file_name()
1659        .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "path has no file name"))?;
1660    randomized_atomic_replace(&parent, name, &tail)?;
1661    Ok(len.saturating_sub(tail.len() as u64))
1662}
1663
1664fn align_tail_start(bytes: &mut Vec<u8>) {
1665    let prefix = bytes
1666        .iter()
1667        .take_while(|byte| **byte & 0xc0 == 0x80)
1668        .count();
1669    if prefix > 0 {
1670        bytes.drain(..prefix);
1671    }
1672}
1673
1674#[cfg(unix)]
1675fn clear_nonblocking(file: &File) -> io::Result<()> {
1676    let flags = unsafe { libc::fcntl(file.as_raw_fd(), libc::F_GETFL) };
1677    if flags == -1 {
1678        return Err(io::Error::last_os_error());
1679    }
1680    if flags & libc::O_NONBLOCK != 0 {
1681        let result =
1682            unsafe { libc::fcntl(file.as_raw_fd(), libc::F_SETFL, flags & !libc::O_NONBLOCK) };
1683        if result == -1 {
1684            return Err(io::Error::last_os_error());
1685        }
1686    }
1687    Ok(())
1688}
1689
1690fn open_validated_path(path: &Path, write: bool) -> io::Result<File> {
1691    let parent = path.parent().unwrap_or_else(|| Path::new("."));
1692    let name = path
1693        .file_name()
1694        .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "path has no file name"))?;
1695    PinnedDir::open(parent)?.open_file(name, write)
1696}
1697
1698#[cfg(unix)]
1699fn openat_file(dirfd: RawFd, name: &OsStr, flags: i32, mode: libc::mode_t) -> io::Result<File> {
1700    let name = os_cstring(name)?;
1701    let fd = unsafe { libc::openat(dirfd, name.as_ptr(), flags, libc::c_uint::from(mode)) };
1702    if fd < 0 {
1703        return Err(io::Error::last_os_error());
1704    }
1705    Ok(unsafe { File::from_raw_fd(fd) })
1706}
1707
1708#[cfg(unix)]
1709fn os_cstring(value: &OsStr) -> io::Result<CString> {
1710    CString::new(value.as_bytes())
1711        .map_err(|_| io::Error::new(io::ErrorKind::InvalidInput, "path contains a NUL byte"))
1712}
1713
1714fn validate_directory_handle(file: &File) -> io::Result<()> {
1715    let metadata = file.metadata()?;
1716    if !metadata.is_dir() {
1717        return Err(io::Error::new(
1718            io::ErrorKind::InvalidData,
1719            "expected a non-reparse directory handle",
1720        ));
1721    }
1722    #[cfg(windows)]
1723    validate_windows_handle(file, true)?;
1724    Ok(())
1725}
1726
1727fn validate_regular_handle(file: &File) -> io::Result<()> {
1728    let metadata = file.metadata()?;
1729    if !metadata.is_file() {
1730        return Err(io::Error::new(
1731            io::ErrorKind::InvalidData,
1732            "task artifact is not a regular file",
1733        ));
1734    }
1735    #[cfg(unix)]
1736    if metadata.nlink() != 1 {
1737        return Err(io::Error::new(
1738            io::ErrorKind::InvalidData,
1739            "task artifact has multiple hard links",
1740        ));
1741    }
1742    #[cfg(windows)]
1743    validate_windows_handle(file, false)?;
1744    Ok(())
1745}
1746
1747#[cfg(unix)]
1748pub fn set_close_on_exec(fd: RawFd, enabled: bool) -> io::Result<()> {
1749    let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
1750    if flags < 0 {
1751        return Err(io::Error::last_os_error());
1752    }
1753    let flags = if enabled {
1754        flags | libc::FD_CLOEXEC
1755    } else {
1756        flags & !libc::FD_CLOEXEC
1757    };
1758    if unsafe { libc::fcntl(fd, libc::F_SETFD, flags) } < 0 {
1759        return Err(io::Error::last_os_error());
1760    }
1761    Ok(())
1762}
1763
1764#[cfg(windows)]
1765const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
1766#[cfg(windows)]
1767const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
1768#[cfg(windows)]
1769const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x0000_0400;
1770#[cfg(windows)]
1771const FILE_TYPE_DISK: u32 = 0x0001;
1772#[cfg(windows)]
1773const HANDLE_FLAG_INHERIT: u32 = 0x0000_0001;
1774
1775#[cfg(windows)]
1776fn windows_file_information(file: &File) -> io::Result<ByHandleFileInformation> {
1777    let mut information = std::mem::MaybeUninit::<ByHandleFileInformation>::zeroed();
1778    if unsafe { GetFileInformationByHandle(file.as_raw_handle(), information.as_mut_ptr()) } == 0 {
1779        return Err(io::Error::last_os_error());
1780    }
1781    Ok(unsafe { information.assume_init() })
1782}
1783
1784#[cfg(windows)]
1785fn validate_windows_handle(file: &File, directory: bool) -> io::Result<()> {
1786    if file.metadata()?.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0 {
1787        return Err(io::Error::new(
1788            io::ErrorKind::InvalidData,
1789            "task path is a reparse point",
1790        ));
1791    }
1792    let handle = file.as_raw_handle();
1793    let file_type = unsafe { GetFileType(handle) };
1794    if file_type != FILE_TYPE_DISK {
1795        return Err(io::Error::new(
1796            io::ErrorKind::InvalidData,
1797            "task artifact is not a regular disk file",
1798        ));
1799    }
1800    let information = windows_file_information(file)?;
1801    if !directory && information.number_of_links != 1 {
1802        return Err(io::Error::new(
1803            io::ErrorKind::InvalidData,
1804            "task artifact has multiple hard links",
1805        ));
1806    }
1807    if unsafe { SetHandleInformation(handle, HANDLE_FLAG_INHERIT, 0) } == 0 {
1808        return Err(io::Error::last_os_error());
1809    }
1810    let mut flags = 0_u32;
1811    if unsafe { GetHandleInformation(handle, &mut flags) } == 0 {
1812        return Err(io::Error::last_os_error());
1813    }
1814    if flags & HANDLE_FLAG_INHERIT != 0 {
1815        return Err(io::Error::new(
1816            io::ErrorKind::PermissionDenied,
1817            "validated task handles must not be inherited",
1818        ));
1819    }
1820    Ok(())
1821}
1822
1823#[cfg(windows)]
1824#[repr(C)]
1825struct ByHandleFileInformation {
1826    file_attributes: u32,
1827    creation_time: [u32; 2],
1828    last_access_time: [u32; 2],
1829    last_write_time: [u32; 2],
1830    volume_serial_number: u32,
1831    file_size_high: u32,
1832    file_size_low: u32,
1833    number_of_links: u32,
1834    file_index_high: u32,
1835    file_index_low: u32,
1836}
1837
1838#[cfg(windows)]
1839#[link(name = "kernel32")]
1840extern "system" {
1841    fn GetFileType(file: std::os::windows::io::RawHandle) -> u32;
1842    fn GetFileInformationByHandle(
1843        file: std::os::windows::io::RawHandle,
1844        information: *mut ByHandleFileInformation,
1845    ) -> i32;
1846    fn SetHandleInformation(object: std::os::windows::io::RawHandle, mask: u32, flags: u32) -> i32;
1847    fn GetHandleInformation(object: std::os::windows::io::RawHandle, flags: *mut u32) -> i32;
1848}
1849
1850#[cfg(test)]
1851mod tests {
1852    use super::*;
1853
1854    fn valid_id(suffix: u64) -> String {
1855        format!("bash-{suffix:016x}")
1856    }
1857
1858    #[test]
1859    fn task_id_validation_is_exact() {
1860        assert!(validate_task_id("bash-0123456789abcdef").is_ok());
1861        for invalid in [
1862            "bash-0123456789abcde",
1863            "bash-0123456789abcdef0",
1864            "bash-0123456789ABCDEf",
1865            "bash-0123456789abcdeg",
1866            "../bash-0123456789abcdef",
1867        ] {
1868            assert!(validate_task_id(invalid).is_err(), "accepted {invalid}");
1869        }
1870    }
1871
1872    #[test]
1873    fn new_layout_separates_control_and_io() {
1874        let storage = tempfile::tempdir().unwrap();
1875        let task = create_task_layout(storage.path(), "session", &valid_id(1)).unwrap();
1876        assert_eq!(
1877            task.paths.json.parent(),
1878            Some(task.paths.control_dir.as_path())
1879        );
1880        assert_eq!(
1881            task.paths.stdout.parent(),
1882            Some(task.paths.io_dir.as_path())
1883        );
1884        assert_ne!(task.paths.control_dir, task.paths.io_dir);
1885    }
1886
1887    #[cfg(unix)]
1888    #[test]
1889    fn task_layout_directories_are_private() {
1890        use std::os::unix::fs::PermissionsExt;
1891
1892        let storage = tempfile::tempdir().unwrap();
1893        let task = create_task_layout(storage.path(), "session", &valid_id(5)).unwrap();
1894        for path in [
1895            &task.paths.session_dir,
1896            &task.paths.dir,
1897            &task.paths.control_dir,
1898            &task.paths.io_dir,
1899        ] {
1900            let mode = fs::metadata(path).unwrap().permissions().mode() & 0o777;
1901            assert_eq!(mode, 0o700, "unexpected permissions for {}", path.display());
1902        }
1903        let bash_tasks = task.paths.session_dir.parent().unwrap();
1904        let mode = fs::metadata(bash_tasks).unwrap().permissions().mode() & 0o777;
1905        assert_eq!(
1906            mode,
1907            0o700,
1908            "unexpected permissions for {}",
1909            bash_tasks.display()
1910        );
1911    }
1912
1913    #[test]
1914    fn resolver_refuses_duplicate_layouts() {
1915        let storage = tempfile::tempdir().unwrap();
1916        let task = create_task_layout(storage.path(), "session", &valid_id(2)).unwrap();
1917        let flat = task
1918            .paths
1919            .session_dir
1920            .join(format!("{}.json", task.paths.task_id));
1921        fs::write(flat, b"{}").unwrap();
1922        let error = resolve_task_layout(&task.paths.session_dir, &task.paths.task_id).unwrap_err();
1923        assert_eq!(error.kind(), io::ErrorKind::InvalidData);
1924    }
1925
1926    #[test]
1927    fn resolver_rejects_metadata_identity_mismatch() {
1928        let storage = tempfile::tempdir().unwrap();
1929        let task = create_task_layout(storage.path(), "session", &valid_id(30)).unwrap();
1930        let metadata = PersistedTask::starting(
1931            valid_id(31),
1932            "session".into(),
1933            "true".into(),
1934            storage.path().into(),
1935            None,
1936            None,
1937            true,
1938            false,
1939        );
1940        fs::write(&task.paths.json, serde_json::to_vec(&metadata).unwrap()).unwrap();
1941        let error = resolve_task_layout(&task.paths.session_dir, &task.paths.task_id).unwrap_err();
1942        assert_eq!(error.kind(), io::ErrorKind::InvalidData);
1943    }
1944
1945    // A task directory swapped underneath the daemon must never let deletion
1946    // touch the impostor's control content. The two platforms enforce this the
1947    // same guarantee through different mechanisms, so each is asserted against
1948    // its real mechanism rather than a shared code path.
1949    //
1950    // Unix: POSIX permits renaming a directory while a fd is held open on it, so
1951    // the swap succeeds on disk and `remove_directory_task`'s `same_identity`
1952    // check is what refuses the deletion.
1953    #[cfg(unix)]
1954    #[test]
1955    fn deletion_refuses_replaced_task_directory_without_touching_victim() {
1956        let storage = tempfile::tempdir().unwrap();
1957        let first = create_task_layout(storage.path(), "session", &valid_id(40)).unwrap();
1958        let second = create_task_layout(storage.path(), "session", &valid_id(41)).unwrap();
1959        let victim = second.paths.control_dir.join("victim");
1960        fs::write(&victim, b"victim-bytes").unwrap();
1961        let moved_first = first.paths.session_dir.join("moved-first");
1962        fs::rename(&first.paths.dir, &moved_first).unwrap();
1963        fs::rename(&second.paths.dir, &first.paths.dir).unwrap();
1964
1965        assert!(delete_resolved_task(&first).is_err());
1966        assert_eq!(
1967            fs::read(first.paths.control_dir.join("victim")).unwrap(),
1968            b"victim-bytes"
1969        );
1970    }
1971
1972    // Windows: the daemon's retained `PinnedDir` handle on the task directory
1973    // makes the OS refuse to rename it (Access denied), so the swap cannot occur
1974    // at all while the daemon is live — the impostor's content is never reachable
1975    // for deletion. Assert that structural refusal directly.
1976    #[cfg(windows)]
1977    #[test]
1978    fn deletion_refuses_replaced_task_directory_without_touching_victim() {
1979        let storage = tempfile::tempdir().unwrap();
1980        let first = create_task_layout(storage.path(), "session", &valid_id(40)).unwrap();
1981        let second = create_task_layout(storage.path(), "session", &valid_id(41)).unwrap();
1982        let victim = second.paths.control_dir.join("victim");
1983        fs::write(&victim, b"victim-bytes").unwrap();
1984
1985        // The daemon still holds `first`'s pinned directory handles, so moving
1986        // its task directory out of the way is refused by the OS.
1987        let moved_first = first.paths.session_dir.join("moved-first");
1988        let refusal = fs::rename(&first.paths.dir, &moved_first)
1989            .expect_err("open pinned-dir handle must block the task-dir rename on Windows");
1990        assert_eq!(refusal.kind(), io::ErrorKind::PermissionDenied);
1991
1992        // The victim's control content is untouched because the swap never happened.
1993        assert_eq!(fs::read(&victim).unwrap(), b"victim-bytes");
1994    }
1995
1996    #[test]
1997    fn legacy_flat_layout_is_readable_and_deleted_as_a_bundle() {
1998        let storage = tempfile::tempdir().unwrap();
1999        let task_id = valid_id(32);
2000        let paths = task_paths(storage.path(), "session", &task_id).unwrap();
2001        fs::create_dir_all(&paths.session_dir).unwrap();
2002        let metadata = PersistedTask::starting(
2003            task_id.clone(),
2004            "session".into(),
2005            "true".into(),
2006            storage.path().into(),
2007            None,
2008            None,
2009            true,
2010            false,
2011        );
2012        write_task(&paths.json, &metadata).unwrap();
2013        fs::write(&paths.stdout, b"legacy").unwrap();
2014        fs::write(&paths.stderr, b"").unwrap();
2015        assert_eq!(
2016            resolve_task_layout(&paths.session_dir, &task_id)
2017                .unwrap()
2018                .paths
2019                .layout,
2020            TaskLayout::Flat
2021        );
2022        assert_eq!(
2023            open_task_artifact(&paths, TaskArtifact::Stdout)
2024                .unwrap()
2025                .read_all()
2026                .unwrap(),
2027            b"legacy"
2028        );
2029        delete_task_bundle(&paths).unwrap();
2030        assert!(!paths.json.exists());
2031        assert!(!paths.stdout.exists());
2032    }
2033
2034    #[cfg(unix)]
2035    #[test]
2036    fn live_output_creation_and_later_writes_refuse_link_attacks() {
2037        use std::os::unix::fs::symlink;
2038
2039        let storage = tempfile::tempdir().unwrap();
2040        let task = create_task_layout(storage.path(), "session", &valid_id(33)).unwrap();
2041        let metadata = PersistedTask::starting(
2042            task.paths.task_id.clone(),
2043            "session".into(),
2044            "true".into(),
2045            storage.path().into(),
2046            None,
2047            None,
2048            true,
2049            false,
2050        );
2051        write_task_at(&task, &metadata).unwrap();
2052        let victim = storage.path().join("victim");
2053        fs::write(&victim, b"victim-bytes").unwrap();
2054
2055        symlink(&victim, &task.paths.stdout).unwrap();
2056        assert!(TaskIoHandles::create(&task, BgMode::Pipes).is_err());
2057        assert_eq!(fs::read(&victim).unwrap(), b"victim-bytes");
2058        fs::remove_file(&task.paths.stdout).unwrap();
2059
2060        let mut handles = TaskIoHandles::create(&task, BgMode::Pipes).unwrap();
2061        fs::hard_link(&task.paths.stdout, task.paths.io_dir.join("linked-stdout")).unwrap();
2062        assert!(handles
2063            .write(TaskArtifact::Stdout, b"daemon-write")
2064            .is_err());
2065        assert_eq!(fs::read(&victim).unwrap(), b"victim-bytes");
2066
2067        fs::remove_file(&task.paths.stdout).unwrap();
2068        symlink(&victim, &task.paths.stdout).unwrap();
2069        assert!(replace_artifact_with_tail(&task.paths, TaskArtifact::Stdout, 1).is_err());
2070        assert_eq!(fs::read(&victim).unwrap(), b"victim-bytes");
2071    }
2072
2073    #[test]
2074    fn registered_artifact_consumers_do_not_reopen_paths_directly() {
2075        let rust_sources = [
2076            include_str!("buffer.rs"),
2077            include_str!("registry.rs"),
2078            include_str!("process.rs"),
2079            include_str!("pty_process.rs"),
2080            include_str!("watches.rs"),
2081            include_str!("watchdog.rs"),
2082            include_str!("../commands/bash_status.rs"),
2083        ];
2084        for source in rust_sources {
2085            let production = source
2086                .split("#[cfg(test)]\nmod tests")
2087                .next()
2088                .unwrap_or(source);
2089            for forbidden in [
2090                "File::open(&task.paths",
2091                "fs::read(&task.paths",
2092                "fs::read_to_string(&task.paths",
2093                "File::open(path)?",
2094            ] {
2095                assert!(
2096                    !production.contains(forbidden),
2097                    "registered artifact consumer contains raw path read: {forbidden}"
2098                );
2099            }
2100        }
2101        for source in [
2102            include_str!("../../../../packages/opencode-plugin/src/tools/bash.ts"),
2103            include_str!("../../../../packages/opencode-plugin/src/tools/bash_watch.ts"),
2104            include_str!("../../../../packages/pi-plugin/src/tools/bash.ts"),
2105        ] {
2106            for forbidden in [
2107                "fs.readFile(outputPath)",
2108                "fs.readFile(details.output_path)",
2109                "fs.open(outputPath",
2110            ] {
2111                assert!(
2112                    !source.contains(forbidden),
2113                    "plugin artifact consumer contains raw path read: {forbidden}"
2114                );
2115            }
2116        }
2117    }
2118
2119    #[cfg(unix)]
2120    #[test]
2121    fn validated_artifact_refuses_symlink_hardlink_and_fifo() {
2122        use std::os::unix::fs::symlink;
2123
2124        let storage = tempfile::tempdir().unwrap();
2125        let task = create_task_layout(storage.path(), "session", &valid_id(3)).unwrap();
2126        let metadata = PersistedTask::starting(
2127            task.paths.task_id.clone(),
2128            "session".into(),
2129            "true".into(),
2130            storage.path().into(),
2131            None,
2132            None,
2133            true,
2134            false,
2135        );
2136        write_task_at(&task, &metadata).unwrap();
2137        let canary = storage.path().join("canary");
2138        fs::write(&canary, b"secret").unwrap();
2139
2140        symlink(&canary, &task.paths.stdout).unwrap();
2141        assert!(open_task_artifact(&task.paths, TaskArtifact::Stdout).is_err());
2142        fs::remove_file(&task.paths.stdout).unwrap();
2143
2144        fs::hard_link(&canary, &task.paths.stdout).unwrap();
2145        assert!(open_task_artifact(&task.paths, TaskArtifact::Stdout).is_err());
2146        fs::remove_file(&task.paths.stdout).unwrap();
2147
2148        let path = CString::new(task.paths.stdout.as_os_str().as_bytes()).unwrap();
2149        assert_eq!(unsafe { libc::mkfifo(path.as_ptr(), 0o600) }, 0);
2150        assert!(open_task_artifact(&task.paths, TaskArtifact::Stdout).is_err());
2151    }
2152}