1use std::collections::HashMap;
7use std::ffi::CString;
8use std::io::IoSlice;
9use std::os::fd::AsRawFd;
10use std::os::unix::ffi::OsStrExt;
11use std::os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt};
12use std::path::Path;
13use std::sync::Arc;
14
15use bytes::Bytes;
16use microsandbox_protocol::bulk::{
17 BULK_FLOW_MASK_GUEST_TO_HOST, BULK_FLOW_MASK_HOST_TO_GUEST, BULK_PROTOCOL_VERSION,
18 BulkAccepted, BulkCredit, BulkFinish, BulkFlow, BulkKind, BulkOffer, BulkReceiveState,
19 BulkRecord, BulkSendState, DEFAULT_BULK_WINDOW, DEFAULT_FILESYSTEM_BULK_RECORD_PAYLOAD,
20};
21use microsandbox_protocol::codec;
22use microsandbox_protocol::fs::{
23 FS_CHUNK_SIZE, FsData, FsEntryInfo, FsOp, FsOpenOptions, FsRequest, FsResponse, FsResponseData,
24 FsSetAttrs,
25};
26use microsandbox_protocol::message::{Message, MessageType};
27use microsandbox_protocol::transport::relay_client_slot;
28use serde::Serialize;
29use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
30use tokio::sync::{Mutex, watch};
31use tokio::task::JoinHandle;
32
33use crate::session::{
34 BulkSessionOutput, RawActivity, RawSessionCompletion, RawSessionOutput, SessionOutput,
35 SessionOutputSender,
36};
37
38const DEFAULT_READ_DIR_LIMIT: u32 = 128;
44
45const MAX_OPEN_HANDLES_PER_OWNER: usize = 1024;
47
48#[derive(Default)]
54pub struct FsState {
55 next_handle: u64,
56 handles: HashMap<u64, FsHandleEntry>,
57}
58
59pub struct FsWriteSession {
61 owner_id: u32,
62 handle: u64,
63 file: Arc<Mutex<tokio::fs::File>>,
64 offset: u64,
65 append: bool,
66 expected_len: Option<u64>,
67 written: u64,
68 bulk: Option<BulkReceiveState>,
69}
70
71pub struct FsReadSession {
73 owner_id: u32,
74 handle: u64,
75 task: JoinHandle<()>,
76 credit_tx: Option<watch::Sender<Option<BulkCredit>>>,
77}
78
79pub enum FsStreamSession {
81 Read(FsReadSession),
83
84 Write(FsWriteSession),
86}
87
88enum FsHandleEntry {
89 File {
90 owner_id: u32,
91 file: Arc<Mutex<tokio::fs::File>>,
92 read: bool,
93 write: bool,
94 append: bool,
95 path: String,
96 },
97 Dir {
98 owner_id: u32,
99 dir: Arc<Mutex<tokio::fs::ReadDir>>,
100 path: String,
101 },
102}
103
104impl FsState {
109 pub fn close_owner_range(&mut self, id_start: u32, id_end_exclusive: u32) {
111 self.handles.retain(|_, handle| {
112 let owner_id = handle.owner_id();
113 owner_id < id_start || owner_id >= id_end_exclusive
114 });
115 }
116
117 pub fn clear(&mut self) {
119 self.handles.clear();
120 }
121
122 fn insert_file(
123 &mut self,
124 owner_id: u32,
125 file: tokio::fs::File,
126 read: bool,
127 write: bool,
128 append: bool,
129 path: String,
130 ) -> Result<u64, String> {
131 self.enforce_owner_limit(owner_id)?;
132 let handle = self.alloc_handle();
133 self.handles.insert(
134 handle,
135 FsHandleEntry::File {
136 owner_id,
137 file: Arc::new(Mutex::new(file)),
138 read,
139 write,
140 append,
141 path,
142 },
143 );
144 Ok(handle)
145 }
146
147 fn insert_dir(
148 &mut self,
149 owner_id: u32,
150 dir: tokio::fs::ReadDir,
151 path: String,
152 ) -> Result<u64, String> {
153 self.enforce_owner_limit(owner_id)?;
154 let handle = self.alloc_handle();
155 self.handles.insert(
156 handle,
157 FsHandleEntry::Dir {
158 owner_id,
159 dir: Arc::new(Mutex::new(dir)),
160 path,
161 },
162 );
163 Ok(handle)
164 }
165
166 fn close_handle(&mut self, caller_id: u32, handle: u64) -> Result<FsHandleEntry, String> {
167 if let Some(entry) = self.handles.get(&handle) {
168 entry.ensure_owner(handle, caller_id)?;
169 }
170 self.handles
171 .remove(&handle)
172 .ok_or_else(|| format!("invalid handle: {handle}"))
173 }
174
175 fn file(
176 &self,
177 caller_id: u32,
178 handle: u64,
179 need_read: bool,
180 need_write: bool,
181 ) -> Result<(Arc<Mutex<tokio::fs::File>>, bool, String), String> {
182 match self.handles.get(&handle) {
183 Some(FsHandleEntry::File {
184 file,
185 read,
186 write,
187 append,
188 path,
189 ..
190 }) => {
191 self.handles
192 .get(&handle)
193 .expect("entry just matched")
194 .ensure_owner(handle, caller_id)?;
195 if need_read && !read {
196 return Err(format!("handle {handle} is not open for reading"));
197 }
198 if need_write && !write && !append {
199 return Err(format!("handle {handle} is not open for writing"));
200 }
201 Ok((Arc::clone(file), *append, path.clone()))
202 }
203 Some(FsHandleEntry::Dir { .. }) => Err(format!("handle {handle} is a directory")),
204 None => Err(format!("invalid handle: {handle}")),
205 }
206 }
207
208 fn dir(
209 &self,
210 caller_id: u32,
211 handle: u64,
212 ) -> Result<(Arc<Mutex<tokio::fs::ReadDir>>, String), String> {
213 match self.handles.get(&handle) {
214 Some(FsHandleEntry::Dir { dir, path, .. }) => {
215 self.handles
216 .get(&handle)
217 .expect("entry just matched")
218 .ensure_owner(handle, caller_id)?;
219 Ok((Arc::clone(dir), path.clone()))
220 }
221 Some(FsHandleEntry::File { .. }) => Err(format!("handle {handle} is a file")),
222 None => Err(format!("invalid handle: {handle}")),
223 }
224 }
225
226 fn alloc_handle(&mut self) -> u64 {
227 self.next_handle = self.next_handle.wrapping_add(1).max(1);
228 while self.handles.contains_key(&self.next_handle) {
229 self.next_handle = self.next_handle.wrapping_add(1).max(1);
230 }
231 self.next_handle
232 }
233
234 fn enforce_owner_limit(&self, owner_id: u32) -> Result<(), String> {
235 let count = self
236 .handles
237 .values()
238 .filter(|entry| same_relay_client(entry.owner_id(), owner_id))
239 .count();
240 if count >= MAX_OPEN_HANDLES_PER_OWNER {
241 return Err(format!(
242 "too many open filesystem handles for relay client: {count}"
243 ));
244 }
245 Ok(())
246 }
247}
248
249impl FsHandleEntry {
250 fn owner_id(&self) -> u32 {
251 match self {
252 Self::File { owner_id, .. } | Self::Dir { owner_id, .. } => *owner_id,
253 }
254 }
255
256 fn ensure_owner(&self, handle: u64, caller_id: u32) -> Result<(), String> {
257 if same_relay_client(self.owner_id(), caller_id) {
258 Ok(())
259 } else {
260 Err(format!(
261 "handle {handle} is owned by a different relay client"
262 ))
263 }
264 }
265}
266
267impl FsReadSession {
268 pub fn owner_id(&self) -> u32 {
270 self.owner_id
271 }
272
273 pub fn handle(&self) -> u64 {
275 self.handle
276 }
277
278 pub fn abort(self) {
280 self.task.abort();
281 }
282
283 pub fn apply_credit(&self, credit: BulkCredit) -> Result<(), String> {
285 let Some(tx) = &self.credit_tx else {
286 return Err("filesystem read session is not using raw bulk".into());
287 };
288 if credit.kind != BulkKind::Filesystem || credit.flow != BulkFlow::GuestToHost {
289 return Err("filesystem read received credit for another kind or flow".into());
290 }
291 tx.send_replace(Some(credit));
292 Ok(())
293 }
294
295 pub fn is_bulk(&self) -> bool {
297 self.credit_tx.is_some()
298 }
299}
300
301impl FsWriteSession {
302 pub fn owner_id(&self) -> u32 {
304 self.owner_id
305 }
306
307 pub fn handle(&self) -> u64 {
309 self.handle
310 }
311
312 pub fn is_bulk(&self) -> bool {
314 self.bulk.is_some()
315 }
316}
317
318fn same_relay_client(left: u32, right: u32) -> bool {
323 relay_client_slot(left).is_some_and(|left| Some(left) == relay_client_slot(right))
324}
325
326pub async fn handle_fs_request(
328 id: u32,
329 protocol_version: u8,
330 req: FsRequest,
331 state: &mut FsState,
332 out_buf: &mut Vec<u8>,
333 session_tx: &SessionOutputSender,
334) -> Result<Option<FsStreamSession>, String> {
335 let FsRequest { op, bulk } = req;
336 if bulk.is_some() && protocol_version < BULK_PROTOCOL_VERSION {
337 encode_response(
338 id,
339 error_response(format!(
340 "raw bulk offer requires protocol generation {BULK_PROTOCOL_VERSION}"
341 )),
342 out_buf,
343 )?;
344 return Ok(None);
345 }
346 if bulk.is_some() && !matches!(&op, FsOp::Read { .. } | FsOp::Write { .. }) {
347 encode_response(
348 id,
349 error_response("raw bulk is only valid for streaming reads and writes".into()),
350 out_buf,
351 )?;
352 return Ok(None);
353 }
354
355 match op {
356 FsOp::RealPath { path } => {
357 let resp = handle_realpath(&path).await;
358 encode_response(id, resp, out_buf)?;
359 Ok(None)
360 }
361 FsOp::Stat {
362 path,
363 follow_symlink,
364 } => {
365 let resp = handle_stat(&path, follow_symlink).await;
366 encode_response(id, resp, out_buf)?;
367 Ok(None)
368 }
369 FsOp::SetStat {
370 path,
371 follow_symlink,
372 attrs,
373 } => {
374 let resp = handle_setstat(&path, follow_symlink, attrs).await;
375 encode_response(id, resp, out_buf)?;
376 Ok(None)
377 }
378 FsOp::List { path } => {
379 let resp = handle_list(&path).await;
380 encode_response(id, resp, out_buf)?;
381 Ok(None)
382 }
383 FsOp::ReadLink { path } => {
384 let resp = handle_readlink(&path).await;
385 encode_response(id, resp, out_buf)?;
386 Ok(None)
387 }
388 FsOp::Symlink {
389 target,
390 link_path,
391 user,
392 } => {
393 let resp = handle_symlink(&target, &link_path, user).await;
394 encode_response(id, resp, out_buf)?;
395 Ok(None)
396 }
397 FsOp::Mkdir { path, mode, user } => {
398 let resp = handle_mkdir(&path, mode, user).await;
399 encode_response(id, resp, out_buf)?;
400 Ok(None)
401 }
402 FsOp::Remove { path } => {
403 let resp = handle_remove(&path).await;
404 encode_response(id, resp, out_buf)?;
405 Ok(None)
406 }
407 FsOp::RemoveDir { path, recursive } => {
408 let resp = handle_remove_dir(&path, recursive).await;
409 encode_response(id, resp, out_buf)?;
410 Ok(None)
411 }
412 FsOp::Copy { src, dst } => {
413 let resp = handle_copy(&src, &dst).await;
414 encode_response(id, resp, out_buf)?;
415 Ok(None)
416 }
417 FsOp::Rename { src, dst } => {
418 let resp = handle_rename(&src, &dst).await;
419 encode_response(id, resp, out_buf)?;
420 Ok(None)
421 }
422 FsOp::OpenFile { path, options } => {
423 let resp = handle_open_file(id, state, &path, options).await;
424 encode_response(id, resp, out_buf)?;
425 Ok(None)
426 }
427 FsOp::OpenDir { path } => {
428 let resp = handle_open_dir(id, state, &path).await;
429 encode_response(id, resp, out_buf)?;
430 Ok(None)
431 }
432 FsOp::CloseHandle { handle } => {
433 let resp = handle_close_handle(id, state, handle).await;
434 encode_response(id, resp, out_buf)?;
435 Ok(None)
436 }
437 FsOp::Read {
438 handle,
439 offset,
440 len,
441 } => match state.file(id, handle, true, false) {
442 Ok((file, _, _)) => {
443 let tx = session_tx.clone();
444 let (task, credit_tx) = match bulk {
445 Some(offer) => {
446 let accepted = accept_fs_read_offer(offer)?;
447 encode_control(MessageType::BulkAccepted, id, &accepted, out_buf)?;
448 let sender = BulkSendState::new(
449 BulkKind::Filesystem,
450 BulkFlow::GuestToHost,
451 accepted.max_record_payload,
452 accepted.guest_to_host_credit_limit,
453 )
454 .map_err(|error| format!("accept read bulk state: {error}"))?;
455 let (credit_tx, credit_rx) = watch::channel(None);
456 let task = tokio::spawn(async move {
457 handle_bulk_read_stream(id, file, offset, len, sender, credit_rx, &tx)
458 .await;
459 });
460 (task, Some(credit_tx))
461 }
462 None => {
463 let task = tokio::spawn(async move {
464 handle_read_stream(id, file, offset, len, &tx).await;
465 });
466 (task, None)
467 }
468 };
469 Ok(Some(FsStreamSession::Read(FsReadSession {
470 owner_id: id,
471 handle,
472 task,
473 credit_tx,
474 })))
475 }
476 Err(e) => {
477 encode_response(id, error_response(format!("read: {e}")), out_buf)?;
478 Ok(None)
479 }
480 },
481 FsOp::Write {
482 handle,
483 offset,
484 len,
485 } => match state.file(id, handle, false, true) {
486 Ok((file, append, _)) => {
487 let bulk = match bulk {
488 Some(offer) => {
489 let accepted = accept_fs_write_offer(offer)?;
490 encode_control(MessageType::BulkAccepted, id, &accepted, out_buf)?;
491 Some(
492 BulkReceiveState::new(
493 BulkKind::Filesystem,
494 BulkFlow::HostToGuest,
495 accepted.max_record_payload,
496 accepted.host_to_guest_credit_limit,
497 DEFAULT_BULK_WINDOW,
498 )
499 .map_err(|error| format!("accept write bulk state: {error}"))?,
500 )
501 }
502 None => None,
503 };
504 Ok(Some(FsStreamSession::Write(FsWriteSession {
505 owner_id: id,
506 handle,
507 file,
508 offset,
509 append,
510 expected_len: len,
511 written: 0,
512 bulk,
513 })))
514 }
515 Err(e) => {
516 encode_response(id, error_response(format!("write: {e}")), out_buf)?;
517 Ok(None)
518 }
519 },
520 FsOp::ReadDir { handle, limit } => {
521 let resp = handle_read_dir(id, state, handle, limit).await;
522 encode_response(id, resp, out_buf)?;
523 Ok(None)
524 }
525 FsOp::FStat { handle } => {
526 let resp = handle_fstat(id, state, handle).await;
527 encode_response(id, resp, out_buf)?;
528 Ok(None)
529 }
530 FsOp::FSetStat { handle, attrs } => {
531 let resp = handle_fsetstat(id, state, handle, attrs).await;
532 encode_response(id, resp, out_buf)?;
533 Ok(None)
534 }
535 }
536}
537
538pub async fn handle_fs_data(
543 id: u32,
544 data: FsData,
545 session: &mut FsWriteSession,
546 out_buf: &mut Vec<u8>,
547) -> Result<bool, String> {
548 if session.bulk.is_some() {
549 encode_response(
550 id,
551 error_response("CBOR filesystem data is invalid after raw bulk acceptance".into()),
552 out_buf,
553 )?;
554 return Ok(true);
555 }
556 if data.data.is_empty() {
557 if let Some(expected) = session.expected_len
558 && session.written != expected
559 {
560 let resp = error_response(format!(
561 "write length mismatch: expected {expected}, wrote {}",
562 session.written
563 ));
564 encode_response(id, resp, out_buf)?;
565 return Ok(true);
566 }
567
568 if let Some(expected) = session.expected_len {
569 let next_written = session.written.saturating_add(data.data.len() as u64);
570 if next_written > expected {
571 let resp = error_response(format!(
572 "write length mismatch: expected {expected}, received at least {next_written}"
573 ));
574 encode_response(id, resp, out_buf)?;
575 return Ok(true);
576 }
577 }
578
579 let mut file = session.file.lock().await;
580 if let Err(e) = file.flush().await {
581 encode_response(id, error_response(format!("flush: {e}")), out_buf)?;
582 return Ok(true);
583 }
584
585 encode_response(id, ok_response(None), out_buf)?;
586 Ok(true)
587 } else {
588 let mut file = session.file.lock().await;
589 if !session.append
590 && let Err(e) = file.seek(std::io::SeekFrom::Start(session.offset)).await
591 {
592 encode_response(id, error_response(format!("seek: {e}")), out_buf)?;
593 return Ok(true);
594 }
595 if let Err(e) = file.write_all(&data.data).await {
596 encode_response(id, error_response(format!("write: {e}")), out_buf)?;
597 return Ok(true);
598 }
599 session.offset = session.offset.saturating_add(data.data.len() as u64);
600 session.written = session.written.saturating_add(data.data.len() as u64);
601 Ok(false)
602 }
603}
604
605pub async fn handle_fs_bulk_record(
607 id: u32,
608 record: &BulkRecord,
609 session: &mut FsWriteSession,
610 out_buf: &mut Vec<u8>,
611) -> Result<bool, String> {
612 handle_fs_bulk_records(id, std::slice::from_ref(record), session, out_buf).await
613}
614
615pub async fn handle_fs_bulk_records(
621 id: u32,
622 records: &[BulkRecord],
623 session: &mut FsWriteSession,
624 out_buf: &mut Vec<u8>,
625) -> Result<bool, String> {
626 if records.is_empty() {
627 return Err("filesystem bulk record batch is empty".into());
628 }
629
630 let mut final_end = session.written;
631 let mut payload_bytes = 0usize;
632 {
633 let Some(receiver) = session.bulk.as_mut() else {
634 return Err("raw bulk record sent to a generation-6 filesystem write".into());
635 };
636 for record in records {
637 final_end = receiver
638 .accept_record(record)
639 .map_err(|error| format!("invalid filesystem bulk record: {error}"))?;
640 payload_bytes = payload_bytes
641 .checked_add(record.payload.len())
642 .ok_or_else(|| "filesystem bulk batch byte count overflowed usize".to_string())?;
643
644 if let Some(expected) = session.expected_len
645 && final_end > expected
646 {
647 encode_response(
648 id,
649 error_response(format!(
650 "write length mismatch: expected {expected}, received at least {final_end}"
651 )),
652 out_buf,
653 )?;
654 return Ok(true);
655 }
656 }
657 }
658
659 let mut file = session.file.lock().await;
660 if !session.append
661 && let Err(error) = file.seek(std::io::SeekFrom::Start(session.offset)).await
662 {
663 encode_response(id, error_response(format!("seek: {error}")), out_buf)?;
664 return Ok(true);
665 }
666 if let Err(error) = write_bulk_payloads_vectored(&mut file, records).await {
667 encode_response(id, error_response(format!("write: {error}")), out_buf)?;
668 return Ok(true);
669 }
670 drop(file);
671
672 session.offset = session.offset.saturating_add(payload_bytes as u64);
673 session.written = final_end;
674 let receiver = session
675 .bulk
676 .as_mut()
677 .expect("bulk receiver was validated before the filesystem write");
678 if let Some(credit) = receiver
679 .consume(final_end)
680 .map_err(|error| format!("advance filesystem bulk credit: {error}"))?
681 {
682 encode_control(MessageType::BulkCredit, id, &credit, out_buf)?;
683 }
684 Ok(false)
685}
686
687async fn write_bulk_payloads_vectored(
689 file: &mut tokio::fs::File,
690 records: &[BulkRecord],
691) -> std::io::Result<()> {
692 if records.len() == 1 {
693 return file.write_all(&records[0].payload).await;
694 }
695
696 let mut record_index = 0usize;
697 let mut record_offset = 0usize;
698
699 while record_index < records.len() {
700 let mut slices = Vec::with_capacity(records.len() - record_index);
701 slices.push(IoSlice::new(
702 &records[record_index].payload[record_offset..],
703 ));
704 slices.extend(
705 records[record_index + 1..]
706 .iter()
707 .map(|record| IoSlice::new(&record.payload)),
708 );
709
710 let written = file.write_vectored(&slices).await?;
711 if written == 0 {
712 return Err(std::io::Error::new(
713 std::io::ErrorKind::WriteZero,
714 "failed to write filesystem bulk batch",
715 ));
716 }
717
718 let mut remaining = written;
719 while record_index < records.len() {
720 let record_remaining = records[record_index].payload.len() - record_offset;
721 if remaining < record_remaining {
722 record_offset += remaining;
723 break;
724 }
725 remaining -= record_remaining;
726 record_index += 1;
727 record_offset = 0;
728 if remaining == 0 {
729 break;
730 }
731 }
732 }
733 Ok(())
734}
735
736pub async fn handle_fs_bulk_finish(
738 id: u32,
739 finish: BulkFinish,
740 session: &mut FsWriteSession,
741 out_buf: &mut Vec<u8>,
742) -> Result<bool, String> {
743 let Some(receiver) = session.bulk.as_mut() else {
744 return Err("bulk finish sent to a generation-6 filesystem write".into());
745 };
746 receiver
747 .accept_finish(finish)
748 .map_err(|error| format!("invalid filesystem bulk finish: {error}"))?;
749
750 if let Some(expected) = session.expected_len
751 && session.written != expected
752 {
753 encode_response(
754 id,
755 error_response(format!(
756 "write length mismatch: expected {expected}, wrote {}",
757 session.written
758 )),
759 out_buf,
760 )?;
761 return Ok(true);
762 }
763
764 let mut file = session.file.lock().await;
765 if let Err(error) = file.flush().await {
766 encode_response(id, error_response(format!("flush: {error}")), out_buf)?;
767 return Ok(true);
768 }
769 encode_response(id, ok_response(None), out_buf)?;
770 Ok(true)
771}
772
773async fn handle_realpath(path: &str) -> FsResponse {
778 match realpath(path).await {
779 Ok(path) => ok_response(Some(FsResponseData::Path(path))),
780 Err(e) => error_response(format!("realpath: {e}")),
781 }
782}
783
784async fn handle_stat(path: &str, follow_symlink: bool) -> FsResponse {
785 let result = if follow_symlink {
786 tokio::fs::metadata(path).await
787 } else {
788 tokio::fs::symlink_metadata(path).await
789 };
790
791 match result {
792 Ok(meta) => ok_response(Some(FsResponseData::Stat(metadata_to_entry_info(
793 path, &meta,
794 )))),
795 Err(e) => error_response(format!("stat: {e}")),
796 }
797}
798
799async fn handle_setstat(path: &str, follow_symlink: bool, attrs: FsSetAttrs) -> FsResponse {
800 match apply_path_attrs(path, follow_symlink, attrs).await {
801 Ok(()) => ok_response(None),
802 Err(e) => error_response(format!("setstat: {e}")),
803 }
804}
805
806async fn handle_list(path: &str) -> FsResponse {
807 match read_all_dir(path).await {
808 Ok(entries) => ok_response(Some(FsResponseData::List(entries))),
809 Err(e) => error_response(format!("readdir: {e}")),
810 }
811}
812
813async fn handle_readlink(path: &str) -> FsResponse {
814 match tokio::fs::read_link(path).await {
815 Ok(target) => ok_response(Some(FsResponseData::Path(
816 target.to_string_lossy().to_string(),
817 ))),
818 Err(e) => error_response(format!("readlink: {e}")),
819 }
820}
821
822async fn handle_symlink(target: &str, link_path: &str, user: Option<String>) -> FsResponse {
823 let target = target.to_string();
824 let link_path = link_path.to_string();
825 match run_as_user(user, move || std::os::unix::fs::symlink(target, link_path)).await {
826 Ok(Ok(())) => ok_response(None),
827 Ok(Err(e)) => error_response(format!("symlink: {e}")),
828 Err(e) => error_response(e),
829 }
830}
831
832async fn handle_open_file(
833 id: u32,
834 state: &mut FsState,
835 path: &str,
836 options: FsOpenOptions,
837) -> FsResponse {
838 let mut open_options = std::fs::OpenOptions::new();
839 open_options
840 .read(options.read)
841 .write(options.write)
842 .append(options.append)
843 .create(options.create)
844 .truncate(options.truncate)
845 .create_new(options.create_new);
846 if let Some(mode) = options.mode {
847 open_options.mode(mode);
848 }
849
850 let open_path = path.to_string();
851 match run_as_user(options.user, move || open_options.open(open_path)).await {
852 Ok(Ok(file)) => match state.insert_file(
853 id,
854 tokio::fs::File::from_std(file),
855 options.read,
856 options.write,
857 options.append,
858 path.to_string(),
859 ) {
860 Ok(handle) => ok_response(Some(FsResponseData::Handle(handle))),
861 Err(e) => error_response(format!("open: {e}")),
862 },
863 Ok(Err(e)) => error_response(format!("open: {e}")),
864 Err(e) => error_response(e),
865 }
866}
867
868async fn handle_open_dir(id: u32, state: &mut FsState, path: &str) -> FsResponse {
869 match tokio::fs::read_dir(path).await {
870 Ok(dir) => match state.insert_dir(id, dir, path.to_string()) {
871 Ok(handle) => ok_response(Some(FsResponseData::Handle(handle))),
872 Err(e) => error_response(format!("opendir: {e}")),
873 },
874 Err(e) => error_response(format!("opendir: {e}")),
875 }
876}
877
878async fn handle_close_handle(id: u32, state: &mut FsState, handle: u64) -> FsResponse {
879 match state.close_handle(id, handle) {
880 Ok(FsHandleEntry::File { file, .. }) => {
881 let mut file = file.lock().await;
882 match file.flush().await {
883 Ok(()) => ok_response(None),
884 Err(e) => error_response(format!("close: {e}")),
885 }
886 }
887 Ok(FsHandleEntry::Dir { .. }) => ok_response(None),
888 Err(e) => error_response(format!("close: {e}")),
889 }
890}
891
892async fn handle_read_dir(id: u32, state: &FsState, handle: u64, limit: Option<u32>) -> FsResponse {
893 let (dir, path) = match state.dir(id, handle) {
894 Ok(v) => v,
895 Err(e) => return error_response(format!("readdir: {e}")),
896 };
897
898 let limit = limit.unwrap_or(DEFAULT_READ_DIR_LIMIT).max(1);
899 let mut dir = dir.lock().await;
900 let mut entries = Vec::new();
901
902 for _ in 0..limit {
903 match dir.next_entry().await {
904 Ok(Some(entry)) => {
905 let entry_path = entry.path();
906 let path_str = entry_path.to_string_lossy().to_string();
907 match tokio::fs::symlink_metadata(&entry_path).await {
908 Ok(meta) => entries.push(metadata_to_entry_info(&path_str, &meta)),
909 Err(_) => entries.push(unknown_entry_info(&path_str)),
910 }
911 }
912 Ok(None) => break,
913 Err(e) => return error_response(format!("readdir {path}: {e}")),
914 }
915 }
916
917 ok_response(Some(FsResponseData::List(entries)))
918}
919
920async fn handle_fstat(id: u32, state: &FsState, handle: u64) -> FsResponse {
921 match state.handles.get(&handle) {
922 Some(FsHandleEntry::File { file, path, .. }) => {
923 if let Err(e) = state
924 .handles
925 .get(&handle)
926 .expect("entry just matched")
927 .ensure_owner(handle, id)
928 {
929 return error_response(format!("fstat: {e}"));
930 }
931 let file = file.lock().await;
932 match file.metadata().await {
933 Ok(meta) => ok_response(Some(FsResponseData::Stat(metadata_to_entry_info(
934 path, &meta,
935 )))),
936 Err(e) => error_response(format!("fstat: {e}")),
937 }
938 }
939 Some(FsHandleEntry::Dir { path, .. }) => {
940 if let Err(e) = state
941 .handles
942 .get(&handle)
943 .expect("entry just matched")
944 .ensure_owner(handle, id)
945 {
946 return error_response(format!("fstat: {e}"));
947 }
948 match tokio::fs::metadata(path).await {
949 Ok(meta) => ok_response(Some(FsResponseData::Stat(metadata_to_entry_info(
950 path, &meta,
951 )))),
952 Err(e) => error_response(format!("fstat: {e}")),
953 }
954 }
955 None => error_response(format!("fstat: invalid handle: {handle}")),
956 }
957}
958
959async fn handle_fsetstat(id: u32, state: &FsState, handle: u64, attrs: FsSetAttrs) -> FsResponse {
960 let (file, _, path) = match state.file(id, handle, false, false) {
961 Ok(v) => v,
962 Err(e) => return error_response(format!("fsetstat: {e}")),
963 };
964
965 let mut file = file.lock().await;
966 match apply_file_attrs(&mut file, &path, attrs).await {
967 Ok(()) => ok_response(None),
968 Err(e) => error_response(format!("fsetstat: {e}")),
969 }
970}
971
972async fn handle_mkdir(path: &str, mode: Option<u32>, user: Option<String>) -> FsResponse {
973 let path = path.to_string();
974 match run_as_user(user, move || {
975 std::fs::create_dir_all(&path).map_err(|e| format!("mkdir: {e}"))?;
976 if let Some(mode) = mode {
977 std::fs::set_permissions(&path, std::fs::Permissions::from_mode(mode))
978 .map_err(|e| format!("chmod: {e}"))?;
979 }
980 Ok(())
981 })
982 .await
983 {
984 Ok(Ok(())) => ok_response(None),
985 Ok(Err(e)) => error_response(e),
986 Err(e) => error_response(e),
987 }
988}
989
990async fn run_as_user<T, F>(user: Option<String>, f: F) -> Result<T, String>
996where
997 T: Send + 'static,
998 F: FnOnce() -> T + Send + 'static,
999{
1000 let Some(user) = user else {
1001 return tokio::task::spawn_blocking(f)
1002 .await
1003 .map_err(|e| format!("filesystem task: {e}"));
1004 };
1005 let (uid, gid, groups) =
1006 crate::session::resolve_user_groups(&user).map_err(|e| e.to_string())?;
1007 tokio::task::spawn_blocking(move || {
1008 let _identity =
1009 FsIdentity::switch(uid, gid, &groups).map_err(|e| format!("switch to {user}: {e}"))?;
1010 Ok(f())
1011 })
1012 .await
1013 .map_err(|e| format!("filesystem task: {e}"))?
1014}
1015
1016struct FsIdentity {
1018 uid: libc::c_int,
1019 gid: libc::c_int,
1020 groups: Vec<libc::gid_t>,
1021}
1022
1023impl FsIdentity {
1024 fn switch(uid: u32, gid: u32, groups: &[libc::gid_t]) -> std::io::Result<Self> {
1025 let previous = current_groups()?;
1026 set_groups(groups)?;
1027 let identity = unsafe {
1029 Self {
1030 gid: libc::setfsgid(gid),
1031 uid: libc::setfsuid(uid),
1032 groups: previous,
1033 }
1034 };
1035 if unsafe { libc::setfsuid(uid) } != uid as libc::c_int
1036 || unsafe { libc::setfsgid(gid) } != gid as libc::c_int
1037 {
1038 return Err(std::io::Error::from(std::io::ErrorKind::PermissionDenied));
1039 }
1040 Ok(identity)
1041 }
1042}
1043
1044impl Drop for FsIdentity {
1045 fn drop(&mut self) {
1046 unsafe {
1047 libc::setfsuid(self.uid as libc::uid_t);
1048 libc::setfsgid(self.gid as libc::gid_t);
1049 }
1050 let _ = set_groups(&self.groups);
1051 }
1052}
1053
1054fn current_groups() -> std::io::Result<Vec<libc::gid_t>> {
1055 let count = unsafe { libc::getgroups(0, std::ptr::null_mut()) };
1056 let mut groups = vec![0; count.max(0) as usize];
1057 let count = unsafe { libc::getgroups(groups.len() as libc::c_int, groups.as_mut_ptr()) };
1058 if count < 0 {
1059 return Err(std::io::Error::last_os_error());
1060 }
1061 groups.truncate(count as usize);
1062 Ok(groups)
1063}
1064
1065fn set_groups(groups: &[libc::gid_t]) -> std::io::Result<()> {
1068 if unsafe { libc::syscall(libc::SYS_setgroups, groups.len(), groups.as_ptr()) } != 0 {
1069 return Err(std::io::Error::last_os_error());
1070 }
1071 Ok(())
1072}
1073
1074async fn handle_remove(path: &str) -> FsResponse {
1075 match tokio::fs::remove_file(path).await {
1076 Ok(()) => ok_response(None),
1077 Err(e) => error_response(format!("remove: {e}")),
1078 }
1079}
1080
1081async fn handle_remove_dir(path: &str, recursive: bool) -> FsResponse {
1082 let result = if recursive {
1083 tokio::fs::remove_dir_all(path).await
1084 } else {
1085 tokio::fs::remove_dir(path).await
1086 };
1087 match result {
1088 Ok(()) => ok_response(None),
1089 Err(e) => error_response(format!("remove_dir: {e}")),
1090 }
1091}
1092
1093async fn handle_copy(src: &str, dst: &str) -> FsResponse {
1094 match tokio::fs::copy(src, dst).await {
1095 Ok(_) => ok_response(None),
1096 Err(e) => error_response(format!("copy: {e}")),
1097 }
1098}
1099
1100async fn handle_rename(src: &str, dst: &str) -> FsResponse {
1101 match tokio::fs::rename(src, dst).await {
1102 Ok(()) => ok_response(None),
1103 Err(e) => error_response(format!("rename: {e}")),
1104 }
1105}
1106
1107fn accept_fs_read_offer(offer: BulkOffer) -> Result<BulkAccepted, String> {
1108 let offer = offer
1109 .validate()
1110 .map_err(|error| format!("invalid filesystem read bulk offer: {error}"))?;
1111 if offer.guest_to_host_credit_limit == 0 {
1112 return Err("filesystem read bulk offer must grant guest-to-host credit".into());
1113 }
1114 Ok(BulkAccepted {
1115 kind: BulkKind::Filesystem,
1116 flows: BULK_FLOW_MASK_GUEST_TO_HOST,
1117 format: offer.format,
1118 max_record_payload: offer
1119 .max_record_payload
1120 .min(DEFAULT_FILESYSTEM_BULK_RECORD_PAYLOAD),
1121 host_to_guest_credit_limit: 0,
1122 guest_to_host_credit_limit: offer.guest_to_host_credit_limit,
1123 })
1124}
1125
1126fn accept_fs_write_offer(offer: BulkOffer) -> Result<BulkAccepted, String> {
1127 let offer = offer
1128 .validate()
1129 .map_err(|error| format!("invalid filesystem write bulk offer: {error}"))?;
1130 if offer.guest_to_host_credit_limit != 0 {
1131 return Err("filesystem write bulk offer must not grant guest-to-host credit".into());
1132 }
1133 Ok(BulkAccepted {
1134 kind: BulkKind::Filesystem,
1135 flows: BULK_FLOW_MASK_HOST_TO_GUEST,
1136 format: offer.format,
1137 max_record_payload: offer
1138 .max_record_payload
1139 .min(DEFAULT_FILESYSTEM_BULK_RECORD_PAYLOAD),
1140 host_to_guest_credit_limit: DEFAULT_BULK_WINDOW,
1141 guest_to_host_credit_limit: 0,
1142 })
1143}
1144
1145async fn handle_bulk_read_stream(
1146 id: u32,
1147 file: Arc<Mutex<tokio::fs::File>>,
1148 offset: u64,
1149 len: Option<u64>,
1150 mut sender: BulkSendState,
1151 mut credit_rx: watch::Receiver<Option<BulkCredit>>,
1152 tx: &SessionOutputSender,
1153) {
1154 let mut file = file.lock().await;
1155 if let Err(error) = file.seek(std::io::SeekFrom::Start(offset)).await {
1156 send_raw_response(id, false, Some(format!("seek: {error}")), None, tx).await;
1157 return;
1158 }
1159
1160 let mut remaining = len;
1161 loop {
1162 if remaining == Some(0) {
1163 break;
1164 }
1165 while sender.available_credit() == 0 {
1166 if credit_rx.changed().await.is_err() {
1167 return;
1168 }
1169 let Some(credit) = *credit_rx.borrow_and_update() else {
1170 continue;
1171 };
1172 if let Err(error) = sender.apply_credit(credit) {
1173 send_raw_response(
1174 id,
1175 false,
1176 Some(format!("invalid filesystem bulk credit: {error}")),
1177 None,
1178 tx,
1179 )
1180 .await;
1181 return;
1182 }
1183 }
1184
1185 let read_len = sender
1186 .available_credit()
1187 .min(sender.max_record_payload() as u64)
1188 .min(remaining.unwrap_or(u64::MAX)) as usize;
1189 let Some(permit) = tx.reserve_bulk(read_len).await else {
1190 return;
1191 };
1192 let mut payload = vec![0u8; read_len];
1193 match file.read(&mut payload).await {
1194 Ok(0) => break,
1195 Ok(read) => {
1196 payload.truncate(read);
1197 if let Some(remaining) = &mut remaining {
1198 *remaining = remaining.saturating_sub(read as u64);
1199 }
1200 let record_offset = match sender.admit(read) {
1201 Ok(offset) => offset,
1202 Err(error) => {
1203 send_raw_response(
1204 id,
1205 false,
1206 Some(format!("admit filesystem bulk record: {error}")),
1207 None,
1208 tx,
1209 )
1210 .await;
1211 return;
1212 }
1213 };
1214 let record = BulkRecord {
1215 id,
1216 kind: BulkKind::Filesystem,
1217 flow: BulkFlow::GuestToHost,
1218 offset: record_offset,
1219 payload: Bytes::from(payload),
1220 };
1221 let output = BulkSessionOutput::new(record, RawActivity::fs_bytes(read));
1222 if !tx
1223 .send_reserved(id, SessionOutput::Bulk(output), permit)
1224 .await
1225 {
1226 return;
1227 }
1228 }
1229 Err(error) => {
1230 send_raw_response(id, false, Some(format!("read: {error}")), None, tx).await;
1231 return;
1232 }
1233 }
1234 }
1235
1236 let finish = match sender.finish() {
1237 Ok(finish) => finish,
1238 Err(error) => {
1239 send_raw_response(
1240 id,
1241 false,
1242 Some(format!("finish filesystem bulk read: {error}")),
1243 None,
1244 tx,
1245 )
1246 .await;
1247 return;
1248 }
1249 };
1250 if !send_raw_control(id, MessageType::BulkFinish, &finish, None, tx).await {
1251 return;
1252 }
1253 send_raw_response(id, true, None, None, tx).await;
1254}
1255
1256async fn handle_read_stream(
1257 id: u32,
1258 file: Arc<Mutex<tokio::fs::File>>,
1259 offset: u64,
1260 len: Option<u64>,
1261 tx: &SessionOutputSender,
1262) {
1263 let mut file = file.lock().await;
1264 if let Err(e) = file.seek(std::io::SeekFrom::Start(offset)).await {
1265 send_raw_response(id, false, Some(format!("seek: {e}")), None, tx).await;
1266 return;
1267 }
1268
1269 let mut remaining = len;
1270 let mut chunk = vec![0u8; FS_CHUNK_SIZE];
1271 let mut buf = Vec::new();
1272
1273 loop {
1274 let Some(permit) = tx.reserve(codec::MAX_FRAME_SIZE as usize + 4).await else {
1277 return;
1278 };
1279 let read_len = match remaining {
1280 Some(0) => break,
1281 Some(n) => chunk.len().min(n as usize),
1282 None => chunk.len(),
1283 };
1284
1285 match file.read(&mut chunk[..read_len]).await {
1286 Ok(0) => break,
1287 Ok(n) => {
1288 if let Some(ref mut remaining) = remaining {
1289 *remaining = remaining.saturating_sub(n as u64);
1290 }
1291 let data = FsData {
1292 data: chunk[..n].to_vec(),
1293 };
1294 let msg = match Message::with_payload(MessageType::FsData, id, &data) {
1295 Ok(msg) => msg,
1296 Err(e) => {
1297 send_raw_response(id, false, Some(format!("encode chunk: {e}")), None, tx)
1298 .await;
1299 return;
1300 }
1301 };
1302 buf.clear();
1303 if let Err(e) = codec::encode_to_buf(&msg, &mut buf) {
1304 send_raw_response(
1305 id,
1306 false,
1307 Some(format!("encode chunk frame: {e}")),
1308 None,
1309 tx,
1310 )
1311 .await;
1312 return;
1313 }
1314 let output =
1315 RawSessionOutput::new(std::mem::take(&mut buf), RawActivity::fs_bytes(n), None);
1316 if !tx
1317 .send_reserved(id, SessionOutput::Raw(output), permit)
1318 .await
1319 {
1320 return;
1321 }
1322 }
1323 Err(e) => {
1324 send_raw_response(id, false, Some(format!("read: {e}")), None, tx).await;
1325 return;
1326 }
1327 }
1328 }
1329
1330 send_raw_response(id, true, None, None, tx).await;
1331}
1332
1333async fn apply_path_attrs(
1338 path: &str,
1339 follow_symlink: bool,
1340 attrs: FsSetAttrs,
1341) -> Result<(), String> {
1342 if let Some(size) = attrs.size {
1343 let file = tokio::fs::OpenOptions::new()
1344 .write(true)
1345 .open(path)
1346 .await
1347 .map_err(|e| format!("open for truncate: {e}"))?;
1348 file.set_len(size)
1349 .await
1350 .map_err(|e| format!("set_len: {e}"))?;
1351 }
1352
1353 if let Some(mode) = attrs.mode {
1354 if !follow_symlink
1355 && tokio::fs::symlink_metadata(path)
1356 .await
1357 .map_err(|e| format!("lstat before chmod: {e}"))?
1358 .file_type()
1359 .is_symlink()
1360 {
1361 return Err("chmod on symlink without following is not supported".into());
1362 }
1363 tokio::fs::set_permissions(path, std::fs::Permissions::from_mode(mode))
1364 .await
1365 .map_err(|e| format!("chmod: {e}"))?;
1366 }
1367
1368 if attrs.uid.is_some() || attrs.gid.is_some() {
1369 chown_path(path, follow_symlink, attrs.uid, attrs.gid)?;
1370 }
1371
1372 if attrs.atime.is_some() || attrs.mtime.is_some() {
1373 set_times_path(path, follow_symlink, attrs.atime, attrs.mtime).await?;
1374 }
1375
1376 Ok(())
1377}
1378
1379async fn apply_file_attrs(
1380 file: &mut tokio::fs::File,
1381 path: &str,
1382 attrs: FsSetAttrs,
1383) -> Result<(), String> {
1384 if let Some(size) = attrs.size {
1385 file.set_len(size)
1386 .await
1387 .map_err(|e| format!("set_len: {e}"))?;
1388 }
1389
1390 if let Some(mode) = attrs.mode {
1391 file.set_permissions(std::fs::Permissions::from_mode(mode))
1392 .await
1393 .map_err(|e| format!("chmod: {e}"))?;
1394 }
1395
1396 if attrs.uid.is_some() || attrs.gid.is_some() {
1397 let uid = attrs.uid.map(|v| v as libc::uid_t).unwrap_or(!0);
1398 let gid = attrs.gid.map(|v| v as libc::gid_t).unwrap_or(!0);
1399 let rc = unsafe { libc::fchown(file.as_raw_fd(), uid, gid) };
1400 if rc != 0 {
1401 return Err(format!("fchown: {}", std::io::Error::last_os_error()));
1402 }
1403 }
1404
1405 if attrs.atime.is_some() || attrs.mtime.is_some() {
1406 set_times_fd(file.as_raw_fd(), path, attrs.atime, attrs.mtime).await?;
1407 }
1408
1409 Ok(())
1410}
1411
1412fn chown_path(
1413 path: &str,
1414 follow_symlink: bool,
1415 uid: Option<u32>,
1416 gid: Option<u32>,
1417) -> Result<(), String> {
1418 let c_path = cstring_path(path)?;
1419 let uid = uid.map(|v| v as libc::uid_t).unwrap_or(!0);
1420 let gid = gid.map(|v| v as libc::gid_t).unwrap_or(!0);
1421 let rc = unsafe {
1422 if follow_symlink {
1423 libc::chown(c_path.as_ptr(), uid, gid)
1424 } else {
1425 libc::lchown(c_path.as_ptr(), uid, gid)
1426 }
1427 };
1428 if rc != 0 {
1429 return Err(format!("chown: {}", std::io::Error::last_os_error()));
1430 }
1431 Ok(())
1432}
1433
1434async fn set_times_path(
1435 path: &str,
1436 follow_symlink: bool,
1437 atime: Option<i64>,
1438 mtime: Option<i64>,
1439) -> Result<(), String> {
1440 let meta = if follow_symlink {
1441 tokio::fs::metadata(path).await
1442 } else {
1443 tokio::fs::symlink_metadata(path).await
1444 }
1445 .map_err(|e| format!("stat before utimensat: {e}"))?;
1446 let times = timespecs(atime.unwrap_or(meta.atime()), mtime.unwrap_or(meta.mtime()));
1447 let c_path = cstring_path(path)?;
1448 let flags = if follow_symlink {
1449 0
1450 } else {
1451 libc::AT_SYMLINK_NOFOLLOW
1452 };
1453 let rc = unsafe { libc::utimensat(libc::AT_FDCWD, c_path.as_ptr(), times.as_ptr(), flags) };
1454 if rc != 0 {
1455 return Err(format!("utimensat: {}", std::io::Error::last_os_error()));
1456 }
1457 Ok(())
1458}
1459
1460async fn set_times_fd(
1461 fd: std::os::fd::RawFd,
1462 path: &str,
1463 atime: Option<i64>,
1464 mtime: Option<i64>,
1465) -> Result<(), String> {
1466 let meta = tokio::fs::metadata(path)
1467 .await
1468 .map_err(|e| format!("stat before futimens: {e}"))?;
1469 let times = timespecs(atime.unwrap_or(meta.atime()), mtime.unwrap_or(meta.mtime()));
1470 let rc = unsafe { libc::futimens(fd, times.as_ptr()) };
1471 if rc != 0 {
1472 return Err(format!("futimens: {}", std::io::Error::last_os_error()));
1473 }
1474 Ok(())
1475}
1476
1477fn timespecs(atime: i64, mtime: i64) -> [libc::timespec; 2] {
1478 [
1479 libc::timespec {
1480 tv_sec: atime as _,
1481 tv_nsec: 0,
1482 },
1483 libc::timespec {
1484 tv_sec: mtime as _,
1485 tv_nsec: 0,
1486 },
1487 ]
1488}
1489
1490fn encode_response(id: u32, resp: FsResponse, out_buf: &mut Vec<u8>) -> Result<(), String> {
1495 let msg = Message::with_payload(MessageType::FsResponse, id, &resp)
1496 .map_err(|e| format!("encode fs response: {e}"))?;
1497 codec::encode_to_buf(&msg, out_buf).map_err(|e| format!("encode fs response frame: {e}"))?;
1498 Ok(())
1499}
1500
1501fn encode_control<T: Serialize>(
1502 message_type: MessageType,
1503 id: u32,
1504 payload: &T,
1505 out_buf: &mut Vec<u8>,
1506) -> Result<(), String> {
1507 let message = Message::with_payload(message_type, id, payload)
1508 .map_err(|error| format!("encode {}: {error}", message_type.as_str()))?;
1509 codec::encode_to_buf(&message, out_buf)
1510 .map_err(|error| format!("encode {} frame: {error}", message_type.as_str()))
1511}
1512
1513async fn send_raw_control<T: Serialize>(
1514 id: u32,
1515 message_type: MessageType,
1516 payload: &T,
1517 completion: Option<RawSessionCompletion>,
1518 tx: &SessionOutputSender,
1519) -> bool {
1520 let mut frame = Vec::new();
1521 if let Err(error) = encode_control(message_type, id, payload, &mut frame) {
1522 eprintln!("failed to {error}");
1523 return false;
1524 }
1525 tx.send(
1526 id,
1527 SessionOutput::Raw(RawSessionOutput::new(
1528 frame,
1529 RawActivity::guest_message(),
1530 completion,
1531 )),
1532 )
1533 .await
1534}
1535
1536async fn send_raw_response(
1537 id: u32,
1538 ok: bool,
1539 error: Option<String>,
1540 data: Option<FsResponseData>,
1541 tx: &SessionOutputSender,
1542) {
1543 let resp = FsResponse { ok, error, data };
1544 match Message::with_payload(MessageType::FsResponse, id, &resp) {
1545 Ok(msg) => {
1546 let mut buf = Vec::new();
1547 match codec::encode_to_buf(&msg, &mut buf) {
1548 Ok(()) => {
1549 let output = RawSessionOutput::new(
1550 buf,
1551 RawActivity::guest_message(),
1552 Some(RawSessionCompletion::FsRead),
1553 );
1554 let _ = tx.send(id, SessionOutput::Raw(output)).await;
1555 }
1556 Err(e) => {
1557 eprintln!("failed to encode fs response frame for {id}: {e}");
1558 }
1559 }
1560 }
1561 Err(e) => {
1562 eprintln!("failed to encode fs response for {id}: {e}");
1563 }
1564 }
1565}
1566
1567async fn realpath(path: &str) -> Result<String, String> {
1568 match tokio::fs::canonicalize(path).await {
1569 Ok(path) => Ok(path.to_string_lossy().to_string()),
1570 Err(original_error) => {
1571 let path = Path::new(path);
1572 let Some(parent) = path.parent() else {
1573 return Err(original_error.to_string());
1574 };
1575 let parent = tokio::fs::canonicalize(parent)
1576 .await
1577 .map_err(|_| original_error.to_string())?;
1578 let resolved = match path.file_name() {
1579 Some(name) => parent.join(name),
1580 None => parent,
1581 };
1582 Ok(resolved.to_string_lossy().to_string())
1583 }
1584 }
1585}
1586
1587async fn read_all_dir(path: &str) -> Result<Vec<FsEntryInfo>, String> {
1588 let mut dir = tokio::fs::read_dir(path)
1589 .await
1590 .map_err(|e| format!("opendir: {e}"))?;
1591 let mut entries = Vec::new();
1592
1593 loop {
1594 match dir.next_entry().await {
1595 Ok(Some(entry)) => {
1596 let entry_path = entry.path();
1597 let path_str = entry_path.to_string_lossy().to_string();
1598 match tokio::fs::symlink_metadata(&entry_path).await {
1599 Ok(meta) => entries.push(metadata_to_entry_info(&path_str, &meta)),
1600 Err(_) => entries.push(unknown_entry_info(&path_str)),
1601 }
1602 }
1603 Ok(None) => break,
1604 Err(e) => return Err(e.to_string()),
1605 }
1606 }
1607
1608 Ok(entries)
1609}
1610
1611fn ok_response(data: Option<FsResponseData>) -> FsResponse {
1612 FsResponse {
1613 ok: true,
1614 error: None,
1615 data,
1616 }
1617}
1618
1619fn error_response(error: String) -> FsResponse {
1620 FsResponse {
1621 ok: false,
1622 error: Some(error),
1623 data: None,
1624 }
1625}
1626
1627fn metadata_to_entry_info(path: &str, meta: &std::fs::Metadata) -> FsEntryInfo {
1628 let kind = if meta.is_file() {
1629 "file"
1630 } else if meta.is_dir() {
1631 "dir"
1632 } else if meta.is_symlink() {
1633 "symlink"
1634 } else {
1635 "other"
1636 };
1637
1638 let mtime = Some(meta.mtime());
1639 let atime = Some(meta.atime());
1640
1641 FsEntryInfo {
1642 path: path.to_string(),
1643 kind: kind.to_string(),
1644 size: meta.len(),
1645 mode: meta.mode(),
1646 modified: mtime,
1647 uid: meta.uid(),
1648 gid: meta.gid(),
1649 atime,
1650 mtime,
1651 }
1652}
1653
1654fn unknown_entry_info(path: &str) -> FsEntryInfo {
1655 FsEntryInfo {
1656 path: path.to_string(),
1657 kind: "other".to_string(),
1658 size: 0,
1659 mode: 0,
1660 modified: None,
1661 uid: 0,
1662 gid: 0,
1663 atime: None,
1664 mtime: None,
1665 }
1666}
1667
1668fn cstring_path(path: impl AsRef<Path>) -> Result<CString, String> {
1669 CString::new(path.as_ref().as_os_str().as_bytes())
1670 .map_err(|e| format!("path contains NUL: {e}"))
1671}
1672
1673#[cfg(test)]
1678mod tests {
1679 use std::time::{SystemTime, UNIX_EPOCH};
1680
1681 use super::*;
1682
1683 #[test]
1684 fn filesystem_offer_preserves_an_older_hosts_smaller_record_limit() {
1685 assert_eq!(
1686 DEFAULT_FILESYSTEM_BULK_RECORD_PAYLOAD as usize,
1687 FS_CHUNK_SIZE
1688 );
1689 let old_offer = BulkOffer {
1690 max_record_payload: microsandbox_protocol::bulk::DEFAULT_BULK_RECORD_PAYLOAD,
1691 ..BulkOffer::filesystem_write()
1692 };
1693
1694 let accepted = accept_fs_write_offer(old_offer).unwrap();
1695 assert_eq!(
1696 accepted.max_record_payload,
1697 microsandbox_protocol::bulk::DEFAULT_BULK_RECORD_PAYLOAD
1698 );
1699 }
1700
1701 #[tokio::test]
1702 async fn raw_bulk_write_requires_exact_offsets_and_finish_length() {
1703 let path = test_path("bulk-write");
1704 let file = tokio::fs::File::create(&path).await.unwrap();
1705 let mut session = FsWriteSession {
1706 owner_id: 1,
1707 handle: 1,
1708 file: Arc::new(Mutex::new(file)),
1709 offset: 0,
1710 append: false,
1711 expected_len: Some(7),
1712 written: 0,
1713 bulk: Some(
1714 BulkReceiveState::new(
1715 BulkKind::Filesystem,
1716 BulkFlow::HostToGuest,
1717 microsandbox_protocol::bulk::DEFAULT_BULK_RECORD_PAYLOAD,
1718 DEFAULT_BULK_WINDOW,
1719 DEFAULT_BULK_WINDOW,
1720 )
1721 .unwrap(),
1722 ),
1723 };
1724 let mut out = Vec::new();
1725
1726 let wrong_offset = BulkRecord {
1727 id: 1,
1728 kind: BulkKind::Filesystem,
1729 flow: BulkFlow::HostToGuest,
1730 offset: 1,
1731 payload: Bytes::from_static(b"ignored"),
1732 };
1733 assert!(
1734 handle_fs_bulk_record(1, &wrong_offset, &mut session, &mut out)
1735 .await
1736 .unwrap_err()
1737 .contains("does not match expected")
1738 );
1739
1740 let first = BulkRecord {
1741 offset: 0,
1742 payload: Bytes::from_static(b"ign"),
1743 ..wrong_offset
1744 };
1745 let second = BulkRecord {
1746 offset: 3,
1747 payload: Bytes::from_static(b"ored"),
1748 ..first.clone()
1749 };
1750 assert!(
1751 !handle_fs_bulk_records(1, &[first, second], &mut session, &mut out)
1752 .await
1753 .unwrap()
1754 );
1755 assert!(
1756 handle_fs_bulk_finish(
1757 1,
1758 BulkFinish {
1759 kind: BulkKind::Filesystem,
1760 flow: BulkFlow::HostToGuest,
1761 final_offset: 7,
1762 },
1763 &mut session,
1764 &mut out,
1765 )
1766 .await
1767 .unwrap()
1768 );
1769 drop(session);
1770
1771 let response = codec::try_decode_from_buf(&mut out).unwrap().unwrap();
1772 assert_eq!(response.t, MessageType::FsResponse);
1773 assert!(response.payload::<FsResponse>().unwrap().ok);
1774 assert_eq!(tokio::fs::read(&path).await.unwrap(), b"ignored");
1775 tokio::fs::remove_file(path).await.unwrap();
1776 }
1777
1778 #[tokio::test]
1779 async fn raw_bulk_read_emits_payload_exact_finish_then_terminal_response() {
1780 let path = test_path("bulk-read");
1781 tokio::fs::write(&path, b"raw-read-payload").await.unwrap();
1782 let file = tokio::fs::File::open(&path).await.unwrap();
1783 let sender = BulkSendState::new(
1784 BulkKind::Filesystem,
1785 BulkFlow::GuestToHost,
1786 microsandbox_protocol::bulk::DEFAULT_BULK_RECORD_PAYLOAD,
1787 DEFAULT_BULK_WINDOW,
1788 )
1789 .unwrap();
1790 let (_credit_tx, credit_rx) = watch::channel(None);
1791 let (session_tx, mut session_rx) = SessionOutputSender::channel();
1792
1793 handle_bulk_read_stream(
1794 2,
1795 Arc::new(Mutex::new(file)),
1796 0,
1797 None,
1798 sender,
1799 credit_rx,
1800 &session_tx,
1801 )
1802 .await;
1803
1804 let first = session_rx.recv().await.unwrap();
1805 let SessionOutput::Bulk(first) = first.output else {
1806 panic!("expected raw filesystem record");
1807 };
1808 assert_eq!(first.record.offset, 0);
1809 assert_eq!(
1810 first.record.payload,
1811 Bytes::from_static(b"raw-read-payload")
1812 );
1813
1814 let finish = decode_raw_output(session_rx.recv().await.unwrap().output);
1815 assert_eq!(finish.t, MessageType::BulkFinish);
1816 let finish: BulkFinish = finish.payload().unwrap();
1817 assert_eq!(finish.final_offset, b"raw-read-payload".len() as u64);
1818 let response = decode_raw_output(session_rx.recv().await.unwrap().output);
1819 assert_eq!(response.t, MessageType::FsResponse);
1820 assert!(response.payload::<FsResponse>().unwrap().ok);
1821 tokio::fs::remove_file(path).await.unwrap();
1822 }
1823
1824 #[tokio::test]
1825 async fn requested_user_owns_created_entries_and_is_bound_by_permissions() {
1826 if unsafe { libc::geteuid() } != 0 {
1828 return;
1829 }
1830 let user = Some("nobody".to_string());
1831 let nobody = crate::session::resolve_default_user(user.as_deref()).unwrap();
1832 let owner = |path: &std::path::Path| {
1833 let meta = std::fs::symlink_metadata(path).unwrap();
1834 (meta.uid(), meta.gid())
1835 };
1836 let dir = test_path("as-user");
1837 std::fs::create_dir(&dir).unwrap();
1838 std::fs::set_permissions(&dir, std::fs::Permissions::from_mode(0o777)).unwrap();
1839 let mut state = FsState::default();
1840 let open = |user: Option<String>| FsOpenOptions {
1841 write: true,
1842 create: true,
1843 user,
1844 ..Default::default()
1845 };
1846
1847 let file = dir.join("file");
1848 let resp =
1849 handle_open_file(1, &mut state, file.to_str().unwrap(), open(user.clone())).await;
1850 assert!(resp.ok, "{:?}", resp.error);
1851 assert_eq!(owner(&file), nobody);
1852
1853 let nested = dir.join("a/b");
1854 let resp = handle_mkdir(nested.to_str().unwrap(), Some(0o755), user.clone()).await;
1855 assert!(resp.ok, "{:?}", resp.error);
1856 assert_eq!(owner(&dir.join("a")), nobody);
1857 assert_eq!(owner(&nested), nobody);
1858
1859 let link = dir.join("link");
1860 let resp = handle_symlink("file", link.to_str().unwrap(), user.clone()).await;
1861 assert!(resp.ok, "{:?}", resp.error);
1862 assert_eq!(owner(&link), nobody);
1863
1864 let root_file = dir.join("root-file");
1866 let resp = handle_open_file(2, &mut state, root_file.to_str().unwrap(), open(None)).await;
1867 assert!(resp.ok, "{:?}", resp.error);
1868 std::fs::set_permissions(&root_file, std::fs::Permissions::from_mode(0o600)).unwrap();
1869 assert_eq!(owner(&root_file), (0, 0));
1870 let resp = handle_open_file(
1871 3,
1872 &mut state,
1873 root_file.to_str().unwrap(),
1874 open(user.clone()),
1875 )
1876 .await;
1877 assert!(!resp.ok);
1878
1879 let after = dir.join("after");
1881 let resp = handle_open_file(4, &mut state, after.to_str().unwrap(), open(None)).await;
1882 assert!(resp.ok, "{:?}", resp.error);
1883 assert_eq!(owner(&after), (0, 0));
1884
1885 if std::env::var_os("MSB_TEST_EDIT_ETC_GROUP").is_none() {
1887 std::fs::remove_dir_all(dir).unwrap();
1888 return;
1889 }
1890 struct RestoreGroupFile(Vec<u8>);
1892 impl Drop for RestoreGroupFile {
1893 fn drop(&mut self) {
1894 std::fs::write("/etc/group", &self.0).unwrap();
1895 }
1896 }
1897 let group_gid = 54321;
1898 let original = std::fs::read("/etc/group").unwrap();
1899 let _restore = RestoreGroupFile(original.clone());
1900 let mut entries = original;
1901 entries.extend_from_slice(format!("msb-shared:x:{group_gid}:nobody\n").as_bytes());
1902 std::fs::write("/etc/group", entries).unwrap();
1903 let shared = dir.join("shared");
1904 std::fs::create_dir(&shared).unwrap();
1905 std::os::unix::fs::chown(&shared, None, Some(group_gid)).unwrap();
1906 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o770)).unwrap();
1907 let root_groups = current_groups().unwrap();
1908 let in_shared = shared.join("file");
1909 let resp = handle_open_file(
1910 5,
1911 &mut state,
1912 in_shared.to_str().unwrap(),
1913 open(user.clone()),
1914 )
1915 .await;
1916 assert!(resp.ok, "{:?}", resp.error);
1917 assert_eq!(owner(&in_shared), nobody);
1918 let resp = handle_mkdir(shared.join("sub").to_str().unwrap(), None, user.clone()).await;
1919 assert!(resp.ok, "{:?}", resp.error);
1920 let (ping, pinged) = std::sync::mpsc::channel::<()>();
1922 let (pong, ponged) = std::sync::mpsc::channel();
1923 std::thread::spawn(move || {
1924 while pinged.recv().is_ok() {
1925 pong.send(current_groups().unwrap()).unwrap();
1926 }
1927 });
1928 let concurrent = run_as_user(user.clone(), move || {
1929 ping.send(()).unwrap();
1930 ponged.recv().unwrap()
1931 })
1932 .await
1933 .unwrap();
1934 assert_eq!(concurrent, root_groups);
1935 assert_eq!(current_groups().unwrap(), root_groups);
1937 let checks: Vec<_> = (0..64)
1938 .map(|_| tokio::task::spawn_blocking(|| current_groups().unwrap()))
1939 .collect();
1940 for check in checks {
1941 assert_eq!(check.await.unwrap(), root_groups);
1942 }
1943
1944 std::fs::remove_dir_all(dir).unwrap();
1945 }
1946
1947 fn decode_raw_output(output: SessionOutput) -> Message {
1948 let SessionOutput::Raw(mut output) = output else {
1949 panic!("expected raw control frame");
1950 };
1951 codec::try_decode_from_buf(&mut output.frame)
1952 .unwrap()
1953 .unwrap()
1954 }
1955
1956 fn test_path(name: &str) -> std::path::PathBuf {
1957 let unique = SystemTime::now()
1958 .duration_since(UNIX_EPOCH)
1959 .unwrap()
1960 .as_nanos();
1961 std::env::temp_dir().join(format!("msb-agentd-{name}-{}-{unique}", std::process::id()))
1962 }
1963}