Skip to main content

microsandbox_agentd/
fs.rs

1//! Guest-side filesystem operation handlers.
2//!
3//! Handles `core.fs.*` protocol messages by performing filesystem operations
4//! using `std::fs` and `tokio::fs`, then sending responses back to the host.
5
6use 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
38//--------------------------------------------------------------------------------------------------
39// Constants
40//--------------------------------------------------------------------------------------------------
41
42/// Default maximum number of entries returned by one `ReadDir` request.
43const DEFAULT_READ_DIR_LIMIT: u32 = 128;
44
45/// Maximum number of open filesystem handles owned by one relay client.
46const MAX_OPEN_HANDLES_PER_OWNER: usize = 1024;
47
48//--------------------------------------------------------------------------------------------------
49// Types
50//--------------------------------------------------------------------------------------------------
51
52/// Mutable filesystem protocol state held by agentd.
53#[derive(Default)]
54pub struct FsState {
55    next_handle: u64,
56    handles: HashMap<u64, FsHandleEntry>,
57}
58
59/// Tracks an in-progress streaming write operation.
60pub 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
71/// Tracks an in-progress streaming read operation.
72pub struct FsReadSession {
73    owner_id: u32,
74    handle: u64,
75    task: JoinHandle<()>,
76    credit_tx: Option<watch::Sender<Option<BulkCredit>>>,
77}
78
79/// A filesystem stream session started by a request.
80pub enum FsStreamSession {
81    /// Read stream task.
82    Read(FsReadSession),
83
84    /// Write stream awaiting `FsData` chunks.
85    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
104//--------------------------------------------------------------------------------------------------
105// Methods
106//--------------------------------------------------------------------------------------------------
107
108impl FsState {
109    /// Close handles opened by a disconnected relay client.
110    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    /// Close all handles.
118    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    /// Correlation ID whose relay client owns this read stream.
269    pub fn owner_id(&self) -> u32 {
270        self.owner_id
271    }
272
273    /// Filesystem handle being read.
274    pub fn handle(&self) -> u64 {
275        self.handle
276    }
277
278    /// Abort the background read task.
279    pub fn abort(self) {
280        self.task.abort();
281    }
282
283    /// Delivers the latest absolute credit update to a generation-8 read task.
284    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    /// Whether this read session negotiated generation-8 raw bulk.
296    pub fn is_bulk(&self) -> bool {
297        self.credit_tx.is_some()
298    }
299}
300
301impl FsWriteSession {
302    /// Correlation ID whose relay client owns this write stream.
303    pub fn owner_id(&self) -> u32 {
304        self.owner_id
305    }
306
307    /// Filesystem handle being written.
308    pub fn handle(&self) -> u64 {
309        self.handle
310    }
311
312    /// Whether this write session negotiated generation-8 raw bulk.
313    pub fn is_bulk(&self) -> bool {
314        self.bulk.is_some()
315    }
316}
317
318//--------------------------------------------------------------------------------------------------
319// Functions
320//--------------------------------------------------------------------------------------------------
321
322fn 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
326/// Handles an incoming `FsRequest` message.
327pub 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
538/// Handles an incoming `FsData` message for a streaming write session.
539///
540/// If `data` is empty, the file is flushed and a terminal `FsResponse` is sent.
541/// Returns `true` if the session should be removed (EOF received).
542pub 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
605/// Handles one generation-8 raw record for a filesystem write correlation.
606pub 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
615/// Handles a bounded contiguous batch of generation-8 filesystem records.
616///
617/// The protocol still validates and accounts for each wire record independently. Coalescing only
618/// changes the local filesystem effect: one lock, one seek and a vectored write replace multiple
619/// blocking-file tasks when a peer negotiates records below the current filesystem default.
620pub 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
687/// Completes a vectored Tokio file write even when the kernel accepts only a prefix of the batch.
688async 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
736/// Handles the exact end marker for a generation-8 filesystem write.
737pub 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
773//--------------------------------------------------------------------------------------------------
774// Functions: Handlers
775//--------------------------------------------------------------------------------------------------
776
777async 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
990/// Runs a blocking filesystem call with the filesystem identity of `user`.
991///
992/// The identity is thread-local (`setfsuid`/`setfsgid` and the raw `setgroups` syscall), so it never
993/// leaks into other requests or exec sessions, and the kernel enforces the user's permissions and
994/// assigns ownership of anything the call creates. `None` runs the call unchanged, as root.
995async 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
1016/// Restores the thread's filesystem identity on drop, including when the call panics.
1017struct 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        // Each call returns the previous id even when it fails, so a second call confirms the change.
1028        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
1065/// Sets the calling thread's supplementary groups. `libc::setgroups` would change every thread in
1066/// the process, including concurrent root requests, so this calls the syscall directly.
1067fn 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        // Reserve before the file read and CBOR materialization so concurrent streams cannot
1275        // each create an uncharged maximum-sized frame while aggregate output is saturated.
1276        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
1333//--------------------------------------------------------------------------------------------------
1334// Functions: Attribute Helpers
1335//--------------------------------------------------------------------------------------------------
1336
1337async 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
1490//--------------------------------------------------------------------------------------------------
1491// Functions: Helpers
1492//--------------------------------------------------------------------------------------------------
1493
1494fn 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//--------------------------------------------------------------------------------------------------
1674// Tests
1675//--------------------------------------------------------------------------------------------------
1676
1677#[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        // Switching the filesystem identity needs root.
1827        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        // Entries created without a user stay root-owned, and the user cannot open them.
1865        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        // The identity is confined to the request.
1880        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        // The rest edits /etc/group, so it runs only where the caller opted in (a throwaway container).
1886        if std::env::var_os("MSB_TEST_EDIT_ETC_GROUP").is_none() {
1887            std::fs::remove_dir_all(dir).unwrap();
1888            return;
1889        }
1890        // A directory writable only through a supplementary group of the user.
1891        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        // Only the calling thread takes the user's groups: a concurrent thread keeps root's.
1921        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        // The user's groups do not outlive the call, on this thread or on reused blocking threads.
1936        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}