Skip to main content

net/
filemanager_thread.rs

1/* This Source Code Form is subject to the terms of the Mozilla Public
2 * License, v. 2.0. If a copy of the MPL was not distributed with this
3 * file, You can obtain one at https://mozilla.org/MPL/2.0/. */
4
5use std::fs::File;
6use std::io::{BufRead, BufReader, Seek, SeekFrom};
7use std::ops::Index;
8use std::path::{Path, PathBuf};
9use std::sync::Arc;
10use std::sync::atomic::{self, AtomicBool, AtomicUsize, Ordering};
11
12use bytes::Bytes;
13use embedder_traits::{
14    EmbedderControlId, EmbedderControlResponse, FilePickerRequest, GenericEmbedderProxy,
15    SelectedFile,
16};
17use headers::{ContentLength, ContentRange, ContentType, HeaderMap, HeaderMapExt, Range};
18use log::warn;
19use mime::Mime;
20use net_traits::blob_url_store::{BlobBuf, BlobTokenCommunicator, BlobURLStoreError};
21use net_traits::filemanager_thread::{
22    FileManagerResult, FileManagerThreadError, FileManagerThreadMsg, FileTokenCheck,
23    GetTokenForFileReply, ReadFileProgress, RelativePos,
24};
25use net_traits::response::{Response, ResponseBody};
26use parking_lot::{Mutex, RwLock};
27use profile_traits::generic_callback::GenericCallback;
28use rustc_hash::{FxHashMap, FxHashSet};
29use servo_arc::Arc as ServoArc;
30use servo_base::generic_channel::GenericSender;
31use servo_url::ImmutableOrigin;
32use tokio::io::{AsyncReadExt, AsyncSeekExt};
33use tokio::sync::mpsc::UnboundedSender as TokioSender;
34use tokio::task::yield_now;
35use uuid::Uuid;
36
37use crate::async_runtime::spawn_task;
38use crate::embedder::NetToEmbedderMsg;
39use crate::fetch::methods::{CancellationListener, Data, RangeRequestBounds};
40use crate::protocols::get_range_request_bounds;
41
42pub const FILE_CHUNK_SIZE: usize = 32768; // 32 KB
43
44/// FileManagerStore's entry
45struct FileStoreEntry {
46    /// Origin of the entry's "creator"
47    origin: ImmutableOrigin,
48    /// Backend implementation
49    file_impl: FileImpl,
50    /// Number of FileID holders that the ID is used to
51    /// index this entry in `FileManagerStore`.
52    /// Reference holders include a FileStoreEntry or
53    /// a script-side File-based Blob
54    refs: AtomicUsize,
55    /// UUIDs only become valid blob URIs when explicitly requested
56    /// by the user with createObjectURL. Validity can be revoked as well.
57    /// (The UUID is the one that maps to this entry in `FileManagerStore`)
58    is_valid_url: AtomicBool,
59    /// UUIDs of fetch instances that acquired an interest in this file,
60    /// when the url was still valid.
61    outstanding_tokens: FxHashSet<Uuid>,
62}
63
64#[derive(Clone)]
65struct FileMetaData {
66    path: PathBuf,
67    size: u64,
68}
69
70/// File backend implementation
71#[derive(Clone)]
72enum FileImpl {
73    /// Metadata of on-disk file
74    MetaDataOnly(FileMetaData),
75    /// In-memory Blob buffer object
76    Memory(BlobBuf),
77    /// A reference to parent entry in `FileManagerStore`,
78    /// representing a sliced version of the parent entry data
79    Sliced(Uuid, RelativePos),
80}
81
82#[derive(Clone)]
83pub struct FileManager {
84    embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
85    store: Arc<FileManagerStore>,
86    blob_token_communicator: Arc<Mutex<BlobTokenCommunicator>>,
87}
88
89impl FileManager {
90    pub fn new(
91        embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
92        blob_token_communicator: Arc<Mutex<BlobTokenCommunicator>>,
93    ) -> FileManager {
94        FileManager {
95            embedder_proxy,
96            store: Arc::new(FileManagerStore::new()),
97            blob_token_communicator,
98        }
99    }
100
101    fn read_file(
102        &self,
103        sender: GenericCallback<FileManagerResult<ReadFileProgress>>,
104        id: Uuid,
105        origin: ImmutableOrigin,
106    ) {
107        let store = self.store.clone();
108        spawn_task(async move {
109            if let Err(e) = store.try_read_file(&sender, id, origin).await {
110                let _ = sender.send(Err(FileManagerThreadError::BlobURLStoreError(e)));
111            }
112        });
113    }
114
115    pub(crate) fn get_token_for_file(&self, file_id: &Uuid, allow_revoked: bool) -> FileTokenCheck {
116        self.store.get_token_for_file(file_id, allow_revoked)
117    }
118
119    pub(crate) fn invalidate_token(&self, token: &FileTokenCheck, file_id: &Uuid) {
120        self.store.invalidate_token(token, file_id);
121    }
122
123    /// Read a file for the Fetch implementation.
124    /// It gets the required headers synchronously and reads the actual content
125    /// in a separate thread.
126    #[expect(clippy::too_many_arguments)]
127    pub(crate) fn fetch_file(
128        &self,
129        done_sender: &mut TokioSender<Data>,
130        cancellation_listener: Arc<CancellationListener>,
131        id: Uuid,
132        file_token: &FileTokenCheck,
133        origin: ImmutableOrigin,
134        response: &mut Response,
135        range: Option<Range>,
136    ) -> Result<(), BlobURLStoreError> {
137        self.fetch_blob_buf(
138            done_sender,
139            cancellation_listener,
140            &id,
141            file_token,
142            &origin,
143            BlobBounds::Unresolved(range),
144            response,
145        )
146    }
147
148    pub fn promote_memory(
149        &self,
150        id: Uuid,
151        blob_buf: BlobBuf,
152        set_valid: bool,
153        origin: ImmutableOrigin,
154    ) {
155        self.store.promote_memory(id, blob_buf, set_valid, origin);
156    }
157
158    /// Message handler
159    pub fn handle(&self, msg: FileManagerThreadMsg) {
160        match msg {
161            FileManagerThreadMsg::SelectFiles(control_id, file_picker_request, response_sender) => {
162                let store = self.store.clone();
163                let embedder = self.embedder_proxy.clone();
164                spawn_task(async move {
165                    let embedder_control_msg = store
166                        .select_files(control_id, file_picker_request, embedder)
167                        .await;
168                    response_sender.send(embedder_control_msg).unwrap();
169                });
170            },
171            FileManagerThreadMsg::ReadFile(sender, id, origin) => {
172                self.read_file(sender, id, origin);
173            },
174            FileManagerThreadMsg::PromoteMemory(id, blob_buf, set_valid, origin) => {
175                self.promote_memory(id, blob_buf, set_valid, origin);
176            },
177            FileManagerThreadMsg::AddSlicedURLEntry(id, rel_pos, sender, origin) => {
178                self.store.add_sliced_url_entry(id, rel_pos, sender, origin);
179            },
180            FileManagerThreadMsg::DecRef(id, origin, sender) => {
181                let _ = sender.send(self.store.dec_ref(&id, &origin));
182            },
183            FileManagerThreadMsg::RevokeBlobURL(id, origin, sender) => {
184                let _ = sender.send(self.store.set_blob_url_validity(false, &id, &origin));
185            },
186            FileManagerThreadMsg::ActivateBlobURL(id, sender, origin) => {
187                let _ = sender.send(self.store.set_blob_url_validity(true, &id, &origin));
188            },
189            FileManagerThreadMsg::GetTokenForFile(id, sender) => {
190                let token = match self.get_token_for_file(&id, false) {
191                    FileTokenCheck::Required(token) => Some(token),
192                    _ => None,
193                };
194
195                let communicator = self.blob_token_communicator.lock();
196                let _ = sender.send(GetTokenForFileReply {
197                    token,
198                    revoke_sender: communicator.revoke_sender.clone(),
199                    refresh_sender: communicator.refresh_token_sender.clone(),
200                });
201            },
202            FileManagerThreadMsg::RevokeTokenForFile(token, id) => {
203                self.invalidate_token(&FileTokenCheck::Required(token), &id);
204            },
205        }
206    }
207
208    pub fn fetch_file_in_chunks(
209        &self,
210        done_sender: &mut TokioSender<Data>,
211        mut reader: BufReader<File>,
212        res_body: ServoArc<Mutex<ResponseBody>>,
213        cancellation_listener: Arc<CancellationListener>,
214        range: RelativePos,
215    ) {
216        let done_sender = done_sender.clone();
217        spawn_task(async move {
218            loop {
219                if cancellation_listener.cancelled() {
220                    *res_body.lock() = ResponseBody::Done(vec![]);
221                    let _ = done_sender.send(Data::Cancelled);
222                    return;
223                }
224                let length = {
225                    let buffer = reader.fill_buf().unwrap().to_vec();
226                    let mut buffer_len = buffer.len();
227                    if let ResponseBody::Receiving(ref mut body) = *res_body.lock() {
228                        let offset = usize::min(
229                            {
230                                if let Some(end) = range.end {
231                                    // HTTP Range requests are specified with closed ranges,
232                                    // while Rust uses half-open ranges. We add +1 here so
233                                    // we don't skip the last requested byte.
234                                    let remaining_bytes =
235                                        end as usize - range.start as usize - body.len() + 1;
236                                    if remaining_bytes <= FILE_CHUNK_SIZE {
237                                        // This is the last chunk so we set buffer
238                                        // len to 0 to break the reading loop.
239                                        buffer_len = 0;
240                                        remaining_bytes
241                                    } else {
242                                        FILE_CHUNK_SIZE
243                                    }
244                                } else {
245                                    FILE_CHUNK_SIZE
246                                }
247                            },
248                            buffer.len(),
249                        );
250                        let chunk = &buffer[0..offset];
251                        body.extend_from_slice(chunk);
252                        let _ = done_sender.send(Data::Payload(Bytes::copy_from_slice(chunk)));
253                    }
254                    buffer_len
255                };
256                if length == 0 {
257                    let mut body = res_body.lock();
258                    let completed_body = match *body {
259                        ResponseBody::Receiving(ref mut body) => std::mem::take(body),
260                        _ => vec![],
261                    };
262                    *body = ResponseBody::Done(completed_body);
263                    let _ = done_sender.send(Data::Done);
264                    break;
265                }
266                reader.consume(length);
267                yield_now().await
268            }
269        });
270    }
271
272    #[expect(clippy::too_many_arguments)]
273    fn fetch_blob_buf(
274        &self,
275        done_sender: &mut TokioSender<Data>,
276        cancellation_listener: Arc<CancellationListener>,
277        id: &Uuid,
278        file_token: &FileTokenCheck,
279        origin_in: &ImmutableOrigin,
280        bounds: BlobBounds,
281        response: &mut Response,
282    ) -> Result<(), BlobURLStoreError> {
283        let file_impl = self.store.get_impl(id, file_token, origin_in)?;
284        /*
285           Only Fetch Blob Range Request would have unresolved range, and only in that case we care about range header.
286        */
287        let mut is_range_requested = false;
288        match file_impl {
289            FileImpl::Memory(buf) => {
290                let bounds = match bounds {
291                    BlobBounds::Unresolved(range) => {
292                        if range.is_some() {
293                            is_range_requested = true;
294                        }
295                        get_range_request_bounds(range, buf.size)
296                    },
297                    BlobBounds::Resolved(bounds) => bounds,
298                };
299                let range = bounds
300                    .get_final(Some(buf.size))
301                    .map_err(|_| BlobURLStoreError::InvalidRange)?;
302
303                let range = range.to_abs_blob_range(buf.size as usize);
304                let len = range.len() as u64;
305                let content_range = if is_range_requested {
306                    ContentRange::bytes(range.start as u64..range.end as u64, buf.size).ok()
307                } else {
308                    None
309                };
310
311                set_blob_response_headers(
312                    &mut response.headers,
313                    len,
314                    buf.type_string.parse().unwrap_or(mime::TEXT_PLAIN),
315                    content_range,
316                );
317
318                let mut bytes = vec![];
319                bytes.extend_from_slice(buf.bytes.index(range));
320
321                let _ = done_sender.send(Data::Payload(Bytes::copy_from_slice(&bytes)));
322                let _ = done_sender.send(Data::Done);
323
324                Ok(())
325            },
326            FileImpl::MetaDataOnly(metadata) => {
327                /* XXX: Snapshot state check (optional) https://w3c.github.io/FileAPI/#snapshot-state.
328                        Concretely, here we create another file, and this file might not
329                        has the same underlying file state (meta-info plus content) as the time
330                        create_entry is called.
331                */
332
333                let file = File::open(&metadata.path)
334                    .map_err(|e| BlobURLStoreError::External(e.to_string()))?;
335                let mut is_range_requested = false;
336                let bounds = match bounds {
337                    BlobBounds::Unresolved(range) => {
338                        if range.is_some() {
339                            is_range_requested = true;
340                        }
341                        get_range_request_bounds(range, metadata.size)
342                    },
343                    BlobBounds::Resolved(bounds) => bounds,
344                };
345                let range = bounds
346                    .get_final(Some(metadata.size))
347                    .map_err(|_| BlobURLStoreError::InvalidRange)?;
348
349                let mut reader = BufReader::with_capacity(FILE_CHUNK_SIZE, file);
350                if reader.seek(SeekFrom::Start(range.start as u64)).is_err() {
351                    return Err(BlobURLStoreError::External(
352                        "Unexpected method for blob".into(),
353                    ));
354                }
355
356                let content_range = if is_range_requested {
357                    let abs_range = range.to_abs_blob_range(metadata.size as usize);
358                    ContentRange::bytes(abs_range.start as u64..abs_range.end as u64, metadata.size)
359                        .ok()
360                } else {
361                    None
362                };
363                set_blob_response_headers(
364                    &mut response.headers,
365                    metadata.size,
366                    mime_guess::from_path(metadata.path)
367                        .first()
368                        .unwrap_or(mime::TEXT_PLAIN),
369                    content_range,
370                );
371
372                self.fetch_file_in_chunks(
373                    &mut done_sender.clone(),
374                    reader,
375                    response.body.clone(),
376                    cancellation_listener,
377                    range,
378                );
379
380                Ok(())
381            },
382            FileImpl::Sliced(parent_id, inner_rel_pos) => {
383                // Next time we don't need to check validity since
384                // we have already done that for requesting URL if necessary.
385                let bounds = RangeRequestBounds::Final(
386                    RelativePos::full_range().slice_inner(&inner_rel_pos),
387                );
388                self.fetch_blob_buf(
389                    done_sender,
390                    cancellation_listener,
391                    &parent_id,
392                    file_token,
393                    origin_in,
394                    BlobBounds::Resolved(bounds),
395                    response,
396                )
397            },
398        }
399    }
400}
401
402enum BlobBounds {
403    Unresolved(Option<Range>),
404    Resolved(RangeRequestBounds),
405}
406
407/// File manager's data store. It maintains a thread-safe mapping
408/// from FileID to FileStoreEntry which might have different backend implementation.
409/// Access to the content is encapsulated as methods of this struct.
410struct FileManagerStore {
411    entries: RwLock<FxHashMap<Uuid, FileStoreEntry>>,
412}
413
414impl FileManagerStore {
415    fn new() -> Self {
416        FileManagerStore {
417            entries: RwLock::new(FxHashMap::default()),
418        }
419    }
420
421    /// Copy out the file backend implementation content
422    fn get_impl(
423        &self,
424        id: &Uuid,
425        file_token: &FileTokenCheck,
426        origin_in: &ImmutableOrigin,
427    ) -> Result<FileImpl, BlobURLStoreError> {
428        match self.entries.read().get(id) {
429            Some(entry) => {
430                if *origin_in != entry.origin {
431                    Err(BlobURLStoreError::InvalidOrigin)
432                } else {
433                    match file_token {
434                        FileTokenCheck::NotRequired => Ok(entry.file_impl.clone()),
435                        FileTokenCheck::Required(token) => {
436                            if entry.outstanding_tokens.contains(token) {
437                                return Ok(entry.file_impl.clone());
438                            }
439                            Err(BlobURLStoreError::InvalidFileID)
440                        },
441                        FileTokenCheck::ShouldFail => Err(BlobURLStoreError::InvalidFileID),
442                    }
443                }
444            },
445            None => Err(BlobURLStoreError::InvalidFileID),
446        }
447    }
448
449    fn invalidate_token(&self, token: &FileTokenCheck, file_id: &Uuid) {
450        if let FileTokenCheck::Required(token) = token {
451            let mut entries = self.entries.write();
452            if let Some(entry) = entries.get_mut(file_id) {
453                entry.outstanding_tokens.remove(token);
454
455                // Check if there are references left.
456                let zero_refs = entry.refs.load(Ordering::Acquire) == 0;
457
458                // Check if no other fetch has acquired a token for this file.
459                let no_outstanding_tokens = entry.outstanding_tokens.is_empty();
460
461                // Check if there is still a blob URL outstanding.
462                let valid = entry.is_valid_url.load(Ordering::Acquire);
463
464                // Can we remove this file?
465                let do_remove = zero_refs && no_outstanding_tokens && !valid;
466
467                if do_remove {
468                    entries.remove(file_id);
469                }
470            }
471        }
472    }
473
474    pub(crate) fn get_token_for_file(&self, file_id: &Uuid, allow_revoked: bool) -> FileTokenCheck {
475        let mut entries = self.entries.write();
476        let parent_id = match entries.get(file_id) {
477            Some(entry) => {
478                if let FileImpl::Sliced(ref parent_id, _) = entry.file_impl {
479                    Some(*parent_id)
480                } else {
481                    None
482                }
483            },
484            None => return FileTokenCheck::ShouldFail,
485        };
486        let file_id = parent_id.as_ref().unwrap_or(file_id);
487
488        if let Some(entry) = entries.get_mut(file_id) {
489            if !allow_revoked && !entry.is_valid_url.load(Ordering::Acquire) {
490                log::warn!("Refusing to grant token for revoked blob url: {file_id:?}");
491                return FileTokenCheck::ShouldFail;
492            }
493            let token = Uuid::new_v4();
494            entry.outstanding_tokens.insert(token);
495            return FileTokenCheck::Required(token);
496        }
497        FileTokenCheck::ShouldFail
498    }
499
500    fn insert(&self, id: Uuid, entry: FileStoreEntry) {
501        self.entries.write().insert(id, entry);
502    }
503
504    fn remove(&self, id: &Uuid) {
505        self.entries.write().remove(id);
506    }
507
508    fn inc_ref(&self, id: &Uuid, origin_in: &ImmutableOrigin) -> Result<(), BlobURLStoreError> {
509        match self.entries.read().get(id) {
510            Some(entry) => {
511                if entry.origin == *origin_in {
512                    entry.refs.fetch_add(1, Ordering::Relaxed);
513                    Ok(())
514                } else {
515                    Err(BlobURLStoreError::InvalidOrigin)
516                }
517            },
518            None => Err(BlobURLStoreError::InvalidFileID),
519        }
520    }
521
522    fn add_sliced_url_entry(
523        &self,
524        parent_id: Uuid,
525        rel_pos: RelativePos,
526        sender: GenericSender<Result<Uuid, BlobURLStoreError>>,
527        origin_in: ImmutableOrigin,
528    ) {
529        match self.inc_ref(&parent_id, &origin_in) {
530            Ok(_) => {
531                let new_id = Uuid::new_v4();
532                self.insert(
533                    new_id,
534                    FileStoreEntry {
535                        origin: origin_in,
536                        file_impl: FileImpl::Sliced(parent_id, rel_pos),
537                        refs: AtomicUsize::new(1),
538                        // Valid here since AddSlicedURLEntry implies URL creation
539                        // from a BlobImpl::Sliced
540                        is_valid_url: AtomicBool::new(true),
541                        outstanding_tokens: Default::default(),
542                    },
543                );
544
545                // We assume that the returned id will be held by BlobImpl::File
546                let _ = sender.send(Ok(new_id));
547            },
548            Err(e) => {
549                let _ = sender.send(Err(e));
550            },
551        }
552    }
553
554    async fn select_files(
555        &self,
556        control_id: EmbedderControlId,
557        file_picker_request: FilePickerRequest,
558        embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
559    ) -> EmbedderControlResponse {
560        let (sender, receiver) = tokio::sync::oneshot::channel();
561
562        let origin = file_picker_request.origin.clone();
563        embedder_proxy.send(NetToEmbedderMsg::SelectFiles(
564            control_id,
565            file_picker_request,
566            sender,
567        ));
568
569        let paths = match receiver.await {
570            Ok(Some(result)) => result,
571            Ok(None) => {
572                return EmbedderControlResponse::FilePicker(None);
573            },
574            Err(error) => {
575                warn!("Failed to receive files from embedder ({:?}).", error);
576                return EmbedderControlResponse::FilePicker(None);
577            },
578        };
579
580        let mut failed = false;
581        let files: Vec<_> = paths
582            .into_iter()
583            .filter_map(|path| match self.create_entry(&path, origin.clone()) {
584                Ok(entry) => Some(entry),
585                Err(error) => {
586                    failed = true;
587                    warn!("Failed to create entry for selected file: {error:?}");
588                    None
589                },
590            })
591            .collect();
592
593        // From <https://w3c.github.io/webdriver/#dfn-element-send-keys>:
594        //
595        // > Step 8.5: Verify that each file given by the user exists. If any do not,
596        // > return error with error code invalid argument.
597        //
598        // WebDriver expects that if any of the files isn't found we don't select any files.
599        if failed {
600            for file in files.iter() {
601                self.remove(&file.id);
602            }
603            return EmbedderControlResponse::FilePicker(Some(Vec::new()));
604        }
605
606        EmbedderControlResponse::FilePicker(Some(files))
607    }
608
609    fn create_entry(
610        &self,
611        file_path: &Path,
612        origin: ImmutableOrigin,
613    ) -> Result<SelectedFile, FileManagerThreadError> {
614        use net_traits::filemanager_thread::FileManagerThreadError::FileSystemError;
615
616        let file = File::open(file_path).map_err(|e| FileSystemError(e.to_string()))?;
617        let metadata = file
618            .metadata()
619            .map_err(|e| FileSystemError(e.to_string()))?;
620        let modified = metadata
621            .modified()
622            .map_err(|e| FileSystemError(e.to_string()))?;
623        let file_size = metadata.len();
624        let file_name = file_path
625            .file_name()
626            .ok_or(FileSystemError("Invalid filepath".to_string()))?;
627
628        let file_impl = FileImpl::MetaDataOnly(FileMetaData {
629            path: file_path.to_path_buf(),
630            size: file_size,
631        });
632
633        let id = Uuid::new_v4();
634
635        self.insert(
636            id,
637            FileStoreEntry {
638                origin,
639                file_impl,
640                refs: AtomicUsize::new(1),
641                // Invalid here since create_entry is called by file selection
642                is_valid_url: AtomicBool::new(false),
643                outstanding_tokens: Default::default(),
644            },
645        );
646
647        let filename_path = Path::new(file_name);
648        let type_string = match mime_guess::from_path(filename_path).first() {
649            Some(x) => format!("{}", x),
650            None => String::new(),
651        };
652
653        Ok(SelectedFile {
654            id,
655            filename: filename_path.to_path_buf(),
656            modified,
657            size: file_size,
658            type_string,
659        })
660    }
661
662    async fn get_blob_buf(
663        &self,
664        sender: &GenericCallback<FileManagerResult<ReadFileProgress>>,
665        id: &Uuid,
666        file_token: &FileTokenCheck,
667        origin_in: &ImmutableOrigin,
668        rel_pos: RelativePos,
669    ) -> Result<(), BlobURLStoreError> {
670        let file_impl = self.get_impl(id, file_token, origin_in)?;
671        match file_impl {
672            FileImpl::Memory(buf) => {
673                let range = rel_pos.to_abs_range(buf.size as usize);
674                let buf = BlobBuf {
675                    filename: None,
676                    type_string: buf.type_string,
677                    size: range.len() as u64,
678                    bytes: buf.bytes.index(range).to_vec(),
679                };
680
681                let _ = sender.send(Ok(ReadFileProgress::Meta(buf)));
682                let _ = sender.send(Ok(ReadFileProgress::EOF));
683
684                Ok(())
685            },
686            FileImpl::MetaDataOnly(metadata) => {
687                /* XXX: Snapshot state check (optional) https://w3c.github.io/FileAPI/#snapshot-state.
688                        Concretely, here we create another file, and this file might not
689                        has the same underlying file state (meta-info plus content) as the time
690                        create_entry is called.
691                */
692
693                let opt_filename = metadata
694                    .path
695                    .file_name()
696                    .and_then(|osstr| osstr.to_str())
697                    .map(|s| s.to_string());
698
699                let mime = mime_guess::from_path(metadata.path.clone()).first();
700                let range = rel_pos.to_abs_range(metadata.size as usize);
701
702                let mut file = tokio::fs::File::open(&metadata.path)
703                    .await
704                    .map_err(|e| BlobURLStoreError::External(e.to_string()))?;
705                let seeked_start = file
706                    .seek(SeekFrom::Start(range.start as u64))
707                    .await
708                    .map_err(|e| BlobURLStoreError::External(e.to_string()))?;
709
710                if seeked_start == (range.start as u64) {
711                    let type_string = match mime {
712                        Some(x) => format!("{}", x),
713                        None => String::new(),
714                    };
715
716                    read_file_in_chunks(sender, file, range.len(), opt_filename, type_string).await;
717                    Ok(())
718                } else {
719                    Err(BlobURLStoreError::InvalidEntry)
720                }
721            },
722            FileImpl::Sliced(parent_id, inner_rel_pos) => {
723                // Next time we don't need to check validity since
724                // we have already done that for requesting URL if necessary
725                Box::pin(self.get_blob_buf(
726                    sender,
727                    &parent_id,
728                    file_token,
729                    origin_in,
730                    rel_pos.slice_inner(&inner_rel_pos),
731                ))
732                .await
733            },
734        }
735    }
736
737    // Convenient wrapper over get_blob_buf
738    async fn try_read_file(
739        &self,
740        sender: &GenericCallback<FileManagerResult<ReadFileProgress>>,
741        id: Uuid,
742        origin_in: ImmutableOrigin,
743    ) -> Result<(), BlobURLStoreError> {
744        self.get_blob_buf(
745            sender,
746            &id,
747            &FileTokenCheck::NotRequired,
748            &origin_in,
749            RelativePos::full_range(),
750        )
751        .await
752    }
753
754    fn dec_ref(&self, id: &Uuid, origin_in: &ImmutableOrigin) -> Result<(), BlobURLStoreError> {
755        let (do_remove, opt_parent_id) = match self.entries.read().get(id) {
756            Some(entry) => {
757                if entry.origin == *origin_in {
758                    let old_refs = entry.refs.fetch_sub(1, Ordering::Release);
759
760                    if old_refs > 1 {
761                        // not the last reference, no need to touch parent
762                        (false, None)
763                    } else {
764                        // last reference, and if it has a reference to parent id
765                        // dec_ref on parent later if necessary
766                        let is_valid = entry.is_valid_url.load(Ordering::Acquire);
767
768                        // Check if no fetch has acquired a token for this file.
769                        let no_outstanding_tokens = entry.outstanding_tokens.is_empty();
770
771                        // Can we remove this file?
772                        let do_remove = !is_valid && no_outstanding_tokens;
773
774                        if let FileImpl::Sliced(ref parent_id, _) = entry.file_impl {
775                            (do_remove, Some(*parent_id))
776                        } else {
777                            (do_remove, None)
778                        }
779                    }
780                } else {
781                    return Err(BlobURLStoreError::InvalidOrigin);
782                }
783            },
784            None => return Err(BlobURLStoreError::InvalidFileID),
785        };
786
787        // Trigger removing if its last reference is gone and it is
788        // not a part of a valid Blob URL
789        if do_remove {
790            atomic::fence(Ordering::Acquire);
791            self.remove(id);
792
793            if let Some(parent_id) = opt_parent_id {
794                return self.dec_ref(&parent_id, origin_in);
795            }
796        }
797
798        Ok(())
799    }
800
801    fn promote_memory(
802        &self,
803        id: Uuid,
804        blob_buf: BlobBuf,
805        set_valid: bool,
806        origin: ImmutableOrigin,
807    ) {
808        self.insert(
809            id,
810            FileStoreEntry {
811                origin,
812                file_impl: FileImpl::Memory(blob_buf),
813                refs: AtomicUsize::new(1),
814                is_valid_url: AtomicBool::new(set_valid),
815                outstanding_tokens: Default::default(),
816            },
817        );
818    }
819
820    fn set_blob_url_validity(
821        &self,
822        validity: bool,
823        id: &Uuid,
824        origin_in: &ImmutableOrigin,
825    ) -> Result<(), BlobURLStoreError> {
826        let (do_remove, opt_parent_id, res) = match self.entries.read().get(id) {
827            Some(entry) => {
828                if entry.origin == *origin_in {
829                    entry.is_valid_url.store(validity, Ordering::Release);
830
831                    if !validity {
832                        // Check if it is the last possible reference
833                        // since refs only accounts for blob id holders
834                        // and store entry id holders
835                        let zero_refs = entry.refs.load(Ordering::Acquire) == 0;
836
837                        // Check if no fetch has acquired a token for this file.
838                        let no_outstanding_tokens = entry.outstanding_tokens.is_empty();
839
840                        // Can we remove this file?
841                        let do_remove = zero_refs && no_outstanding_tokens;
842
843                        if let FileImpl::Sliced(ref parent_id, _) = entry.file_impl {
844                            (do_remove, Some(*parent_id), Ok(()))
845                        } else {
846                            (do_remove, None, Ok(()))
847                        }
848                    } else {
849                        (false, None, Ok(()))
850                    }
851                } else {
852                    (false, None, Err(BlobURLStoreError::InvalidOrigin))
853                }
854            },
855            None => (false, None, Err(BlobURLStoreError::InvalidFileID)),
856        };
857
858        if do_remove {
859            atomic::fence(Ordering::Acquire);
860            self.remove(id);
861
862            if let Some(parent_id) = opt_parent_id {
863                return self.dec_ref(&parent_id, origin_in);
864            }
865        }
866        res
867    }
868}
869
870async fn read_file_in_chunks(
871    sender: &GenericCallback<FileManagerResult<ReadFileProgress>>,
872    mut file: tokio::fs::File,
873    size: usize,
874    opt_filename: Option<String>,
875    type_string: String,
876) {
877    // First chunk
878    let mut buf = vec![0; FILE_CHUNK_SIZE];
879    match file.read(&mut buf).await {
880        Ok(n) => {
881            buf.truncate(n);
882            let blob_buf = BlobBuf {
883                filename: opt_filename,
884                type_string,
885                size: size as u64,
886                bytes: buf,
887            };
888            let _ = sender.send(Ok(ReadFileProgress::Meta(blob_buf)));
889        },
890        Err(e) => {
891            let _ = sender.send(Err(FileManagerThreadError::FileSystemError(e.to_string())));
892            return;
893        },
894    }
895
896    // Send the remaining chunks
897    loop {
898        let mut buf = vec![0; FILE_CHUNK_SIZE];
899        match file.read(&mut buf).await {
900            Ok(0) => {
901                let _ = sender.send(Ok(ReadFileProgress::EOF));
902                return;
903            },
904            Ok(n) => {
905                buf.truncate(n);
906                let _ = sender.send(Ok(ReadFileProgress::Partial(buf)));
907            },
908            Err(e) => {
909                let _ = sender.send(Err(FileManagerThreadError::FileSystemError(e.to_string())));
910                return;
911            },
912        }
913    }
914}
915
916fn set_blob_response_headers(
917    headers: &mut HeaderMap,
918    content_length: u64,
919    mime: Mime,
920    content_range: Option<ContentRange>,
921) {
922    headers.typed_insert(ContentLength(content_length));
923    if let Some(content_range) = content_range {
924        headers.typed_insert(content_range);
925    }
926    headers.typed_insert(ContentType::from(mime));
927}