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 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(¤t)?;
409 let held = windows_file_information(&self.file)?;
410 let observed = windows_file_information(¤t)?;
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 #[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 #[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 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 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}