1use 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; struct FileStoreEntry {
46 origin: ImmutableOrigin,
48 file_impl: FileImpl,
50 refs: AtomicUsize,
55 is_valid_url: AtomicBool,
59 outstanding_tokens: FxHashSet<Uuid>,
62}
63
64#[derive(Clone)]
65struct FileMetaData {
66 path: PathBuf,
67 size: u64,
68}
69
70#[derive(Clone)]
72enum FileImpl {
73 MetaDataOnly(FileMetaData),
75 Memory(BlobBuf),
77 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 #[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 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 let remaining_bytes =
235 end as usize - range.start as usize - body.len() + 1;
236 if remaining_bytes <= FILE_CHUNK_SIZE {
237 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 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 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 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
407struct 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 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 let zero_refs = entry.refs.load(Ordering::Acquire) == 0;
457
458 let no_outstanding_tokens = entry.outstanding_tokens.is_empty();
460
461 let valid = entry.is_valid_url.load(Ordering::Acquire);
463
464 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 is_valid_url: AtomicBool::new(true),
541 outstanding_tokens: Default::default(),
542 },
543 );
544
545 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 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 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 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 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 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 (false, None)
763 } else {
764 let is_valid = entry.is_valid_url.load(Ordering::Acquire);
767
768 let no_outstanding_tokens = entry.outstanding_tokens.is_empty();
770
771 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 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 let zero_refs = entry.refs.load(Ordering::Acquire) == 0;
836
837 let no_outstanding_tokens = entry.outstanding_tokens.is_empty();
839
840 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 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 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}