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