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