Skip to main content

net/
resource_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
5//! A thread that takes a URL and streams back the binary data.
6
7use std::borrow::ToOwned;
8use std::collections::HashMap;
9use std::fs::File;
10use std::io::{self, BufReader};
11use std::path::{Path, PathBuf};
12use std::sync::{Arc, Weak};
13use std::thread;
14
15use cookie::Cookie;
16use crossbeam_channel::Sender;
17use devtools_traits::DevtoolsControlMsg;
18use embedder_traits::GenericEmbedderProxy;
19use hyper_serde::Serde;
20use ipc_channel::ipc::IpcSender;
21use log::{debug, trace, warn};
22use malloc_size_of_derive::MallocSizeOf;
23use net_traits::blob_url_store::{BlobTokenCommunicator, parse_blob_url};
24use net_traits::filemanager_thread::FileTokenCheck;
25use net_traits::pub_domains::public_suffix_list_size_of;
26use net_traits::request::{Destination, PreloadEntry, PreloadId, RequestBuilder, RequestId};
27use net_traits::response::{Response, ResponseInit};
28use net_traits::{
29    AsyncRuntime, CookieAsyncResponse, CookieData, CookieSource, CoreResourceMsg,
30    CoreResourceThread, CustomResponseMediator, DiscardFetch, FetchChannels, FetchTaskTarget,
31    NetworkError, ResourceFetchTiming, ResourceThreads, ResourceTimingType, WebSocketDomAction,
32    WebSocketNetworkEvent,
33};
34use parking_lot::{Mutex, RwLock};
35use profile_traits::mem::{
36    ProcessReports, ProfilerChan as MemProfilerChan, Report, ReportKind, ReportsChan,
37    perform_memory_report,
38};
39use profile_traits::path;
40use profile_traits::time::ProfilerChan;
41use rustc_hash::FxHashMap;
42use rustls_pki_types::CertificateDer;
43use rustls_pki_types::pem::PemObject;
44use serde::{Deserialize, Serialize};
45use servo_base::generic_channel::{
46    self, CallbackSetter, GenericCallback, GenericReceiver, GenericReceiverSet,
47    GenericSelectionResult,
48};
49use servo_base::id::CookieStoreId;
50use servo_url::{ImmutableOrigin, ServoUrl};
51use tokio::sync::Mutex as TokioMutex;
52
53use crate::async_runtime::{init_async_runtime, spawn_task};
54use crate::connector::{
55    CACertificates, CertificateErrorOverrideManager, create_http_client, create_tls_config,
56};
57use crate::cookie::ServoCookie;
58use crate::cookie_storage::CookieStorage;
59use crate::embedder::NetToEmbedderMsg;
60use crate::fetch::cors_cache::CorsCache;
61use crate::fetch::fetch_params::{FetchParams, SharedPreloadedResources};
62use crate::fetch::methods::{
63    AutoRequestBodyStreamCloser, CancellationListener, FetchContext,
64    SharedInflightKeepAliveRecords, WebSocketChannel, fetch,
65    transfers_request_body_stream_to_later_manual_redirect,
66};
67use crate::filemanager_thread::FileManager;
68use crate::hsts::{self, HstsList};
69use crate::http_cache::{HttpCache, HttpCacheAssignment};
70use crate::http_loader::{HttpState, http_redirect_fetch};
71use crate::protocols::ProtocolRegistry;
72use crate::request_interceptor::RequestInterceptor;
73use crate::websocket_loader::create_handshake_request;
74
75/// Load a file with CA certificate and produce a RootCertStore with the results.
76fn load_root_cert_store_from_file(file_path: String) -> io::Result<Vec<CertificateDer<'static>>> {
77    let mut pem = BufReader::new(File::open(file_path)?);
78
79    let certs = CertificateDer::pem_reader_iter(&mut pem)
80        .filter_map(|cert| {
81            cert.inspect_err(|e| log::error!("Could not load certificate ({e}). Ignoring it."))
82                .ok()
83        })
84        .collect();
85    Ok(certs)
86}
87
88/// Returns a tuple of (public, private) senders to the new threads.
89#[expect(clippy::too_many_arguments)]
90pub fn new_resource_threads(
91    devtools_sender: Option<Sender<DevtoolsControlMsg>>,
92    time_profiler_chan: ProfilerChan,
93    mem_profiler_chan: MemProfilerChan,
94    embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
95    config_dir: Option<PathBuf>,
96    certificate_path: Option<String>,
97    ignore_certificate_errors: bool,
98    protocols: Arc<ProtocolRegistry>,
99) -> (ResourceThreads, ResourceThreads, Box<dyn AsyncRuntime>) {
100    // Initialize the async runtime, and get a handle to it for use in clean shutdown.
101    let async_runtime = init_async_runtime();
102
103    let ca_certificates = certificate_path
104        .and_then(|path| {
105            Some(CACertificates::Override(
106                load_root_cert_store_from_file(path).ok()?,
107            ))
108        })
109        .unwrap_or_default();
110
111    let (public_core, private_core) = new_core_resource_thread(
112        devtools_sender,
113        time_profiler_chan,
114        mem_profiler_chan,
115        embedder_proxy,
116        config_dir,
117        ca_certificates,
118        ignore_certificate_errors,
119        protocols,
120    );
121    (
122        ResourceThreads::new(public_core),
123        ResourceThreads::new(private_core),
124        async_runtime,
125    )
126}
127
128/// Create a CoreResourceThread
129#[expect(clippy::too_many_arguments)]
130pub fn new_core_resource_thread(
131    devtools_sender: Option<Sender<DevtoolsControlMsg>>,
132    time_profiler_chan: ProfilerChan,
133    mem_profiler_chan: MemProfilerChan,
134    embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
135    config_dir: Option<PathBuf>,
136    ca_certificates: CACertificates<'static>,
137    ignore_certificate_errors: bool,
138    protocols: Arc<ProtocolRegistry>,
139) -> (CoreResourceThread, CoreResourceThread) {
140    let (public_setup_chan, public_setup_port) = generic_channel::channel().unwrap();
141    let (private_setup_chan, private_setup_port) = generic_channel::channel().unwrap();
142    let (report_chan, report_port) = generic_channel::channel().unwrap();
143    let (revoke_sender, revoke_receiver) = generic_channel::channel().unwrap();
144    let (refresh_sender, refresh_receiver) = generic_channel::channel().unwrap();
145
146    let blob_token_communicator = Arc::new(Mutex::new(BlobTokenCommunicator {
147        revoke_sender,
148        refresh_token_sender: refresh_sender,
149    }));
150    thread::Builder::new()
151        .name("ResourceManager".to_owned())
152        .spawn(move || {
153            let resource_manager = CoreResourceManager::new(
154                devtools_sender,
155                time_profiler_chan,
156                embedder_proxy.clone(),
157                ca_certificates.clone(),
158                ignore_certificate_errors,
159                blob_token_communicator,
160            );
161
162            let mut channel_manager = ResourceChannelManager {
163                resource_manager,
164                config_dir,
165                ca_certificates,
166                ignore_certificate_errors,
167                cancellation_listeners: Default::default(),
168                cookie_listeners: Default::default(),
169            };
170
171            mem_profiler_chan.run_with_memory_reporting(
172                || {
173                    channel_manager.start(
174                        public_setup_port,
175                        private_setup_port,
176                        report_port,
177                        revoke_receiver,
178                        refresh_receiver,
179                        protocols,
180                        embedder_proxy,
181                    )
182                },
183                String::from("network-cache-reporter"),
184                report_chan,
185                CoreResourceMsg::CollectMemoryReport,
186            );
187        })
188        .expect("Thread spawning failed");
189    (public_setup_chan, private_setup_chan)
190}
191
192struct ResourceChannelManager {
193    resource_manager: CoreResourceManager,
194    config_dir: Option<PathBuf>,
195    ca_certificates: CACertificates<'static>,
196    ignore_certificate_errors: bool,
197    cancellation_listeners: FxHashMap<RequestId, Weak<CancellationListener>>,
198    cookie_listeners: FxHashMap<CookieStoreId, GenericCallback<CookieAsyncResponse>>,
199}
200
201/// This returns a tuple HttpState and a private HttpState.
202fn create_http_states(
203    config_dir: Option<&Path>,
204    ca_certificates: CACertificates<'static>,
205    ignore_certificate_errors: bool,
206    embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
207) -> (Arc<HttpState>, Arc<HttpState>) {
208    let mut hsts_list = HstsList::default();
209    let mut auth_cache = AuthCache::default();
210    let mut cookie_jar = CookieStorage::new(150);
211    if let Some(config_dir) = config_dir {
212        servo_base::read_json_from_file(&mut auth_cache, config_dir, "auth_cache.json");
213        servo_base::read_json_from_file(&mut hsts_list, config_dir, "hsts_list.json");
214        servo_base::read_json_from_file(&mut cookie_jar, config_dir, "cookie_jar.json");
215    }
216
217    let override_manager = CertificateErrorOverrideManager::new();
218    let http_state = HttpState {
219        hsts_list: RwLock::new(hsts_list),
220        cookie_jar: RwLock::new(cookie_jar),
221        auth_cache: RwLock::new(auth_cache),
222        history_states: RwLock::new(FxHashMap::default()),
223        http_cache: HttpCache::new(HttpCacheAssignment::Public),
224        client: create_http_client(create_tls_config(
225            ca_certificates.clone(),
226            ignore_certificate_errors,
227            override_manager.clone(),
228        )),
229        override_manager,
230        embedder_proxy: embedder_proxy.clone(),
231    };
232
233    let override_manager = CertificateErrorOverrideManager::new();
234    let private_http_state = HttpState {
235        hsts_list: RwLock::new(HstsList::default()),
236        cookie_jar: RwLock::new(CookieStorage::new(150)),
237        auth_cache: RwLock::new(AuthCache::default()),
238        history_states: RwLock::new(FxHashMap::default()),
239        http_cache: HttpCache::new(HttpCacheAssignment::Private),
240        client: create_http_client(create_tls_config(
241            ca_certificates,
242            ignore_certificate_errors,
243            override_manager.clone(),
244        )),
245        override_manager,
246        embedder_proxy,
247    };
248
249    (Arc::new(http_state), Arc::new(private_http_state))
250}
251
252impl ResourceChannelManager {
253    #[expect(clippy::too_many_arguments)]
254    fn start(
255        &mut self,
256        public_receiver: GenericReceiver<CoreResourceMsg>,
257        private_receiver: GenericReceiver<CoreResourceMsg>,
258        memory_reporter: GenericReceiver<CoreResourceMsg>,
259        revoke_receiver: GenericReceiver<CoreResourceMsg>,
260        refresh_receiver: GenericReceiver<CoreResourceMsg>,
261        protocols: Arc<ProtocolRegistry>,
262        embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
263    ) {
264        let (public_http_state, private_http_state) = create_http_states(
265            self.config_dir.as_deref(),
266            self.ca_certificates.clone(),
267            self.ignore_certificate_errors,
268            embedder_proxy,
269        );
270
271        let mut rx_set = GenericReceiverSet::new();
272        let private_id = rx_set.add(private_receiver);
273        let public_id = rx_set.add(public_receiver);
274        let reporter_id = rx_set.add(memory_reporter);
275        let revoker_id = rx_set.add(revoke_receiver);
276        let refresh_id = rx_set.add(refresh_receiver);
277        let mut selector = rx_set.selector();
278
279        loop {
280            for received in selector.select().into_iter() {
281                // Handles case where profiler thread shuts down before resource thread.
282                match received {
283                    GenericSelectionResult::ChannelClosed(_) => continue,
284                    GenericSelectionResult::Error(error) => {
285                        log::error!("Found selection error: {error}")
286                    },
287                    GenericSelectionResult::MessageReceived(id, msg) => {
288                        if id == revoker_id {
289                            let CoreResourceMsg::RevokeTokenForFile(revocation_request) = msg
290                            else {
291                                log::error!("Blob revocation channel received unexpected message");
292                                continue;
293                            };
294                            self.resource_manager.filemanager.invalidate_token(
295                                &FileTokenCheck::Required(revocation_request.token),
296                                &revocation_request.blob_id,
297                            )
298                        } else if id == refresh_id {
299                            let CoreResourceMsg::RefreshTokenForFile(refresh_request) = msg else {
300                                log::error!("Blob revocation channel received unexpected message");
301                                continue;
302                            };
303
304                            let FileTokenCheck::Required(refreshed_token) = self
305                                .resource_manager
306                                .filemanager
307                                .get_token_for_file(&refresh_request.blob_id, true)
308                            else {
309                                unreachable!();
310                            };
311                            let _ = refresh_request.new_token_sender.send(refreshed_token);
312                        } else if id == reporter_id {
313                            if let CoreResourceMsg::CollectMemoryReport(report_chan) = msg {
314                                self.process_report(
315                                    report_chan,
316                                    &public_http_state,
317                                    &private_http_state,
318                                );
319                                continue;
320                            } else {
321                                log::error!("memory reporter should only send CollectMemoryReport");
322                            }
323                        } else {
324                            let group = if id == private_id {
325                                &private_http_state
326                            } else {
327                                assert_eq!(id, public_id);
328                                &public_http_state
329                            };
330                            if !self.process_msg(msg, group, Arc::clone(&protocols)) {
331                                return;
332                            }
333                        }
334                    },
335                }
336            }
337        }
338    }
339
340    fn process_report(
341        &mut self,
342        msg: ReportsChan,
343        public_http_state: &Arc<HttpState>,
344        private_http_state: &Arc<HttpState>,
345    ) {
346        perform_memory_report(|ops| {
347            let mut reports = public_http_state.memory_reports("public", ops);
348            reports.extend(private_http_state.memory_reports("private", ops));
349            reports.extend(vec![
350                Report {
351                    path: path!["hsts-preload-list"],
352                    kind: ReportKind::ExplicitJemallocHeapSize,
353                    size: hsts::hsts_preload_size_of(ops),
354                },
355                Report {
356                    path: path!["public-suffix-list"],
357                    kind: ReportKind::ExplicitJemallocHeapSize,
358                    size: public_suffix_list_size_of(ops),
359                },
360            ]);
361            msg.send(ProcessReports::new(reports));
362        })
363    }
364
365    fn cancellation_listener(&self, request_id: RequestId) -> Option<Arc<CancellationListener>> {
366        self.cancellation_listeners
367            .get(&request_id)
368            .and_then(Weak::upgrade)
369    }
370
371    fn get_or_create_cancellation_listener(
372        &mut self,
373        request_id: RequestId,
374    ) -> Arc<CancellationListener> {
375        if let Some(listener) = self.cancellation_listener(request_id) {
376            return listener;
377        }
378
379        // Clear away any cancellation listeners that are no longer valid.
380        self.cancellation_listeners
381            .retain(|_, listener| listener.strong_count() > 0);
382
383        let cancellation_listener = Arc::new(Default::default());
384        self.cancellation_listeners
385            .insert(request_id, Arc::downgrade(&cancellation_listener));
386        cancellation_listener
387    }
388
389    fn send_cookie_response(&self, store_id: CookieStoreId, data: CookieData) {
390        let Some(sender) = self.cookie_listeners.get(&store_id) else {
391            warn!(
392                "Async cookie request made for store id that is non-existent {:?}",
393                store_id
394            );
395            return;
396        };
397        let res = sender.send(CookieAsyncResponse { data });
398        if res.is_err() {
399            warn!("Unable to send cookie response to script thread");
400        }
401    }
402
403    /// Returns false if the thread should exit.
404    fn process_msg(
405        &mut self,
406        msg: CoreResourceMsg,
407        http_state: &Arc<HttpState>,
408        protocols: Arc<ProtocolRegistry>,
409    ) -> bool {
410        match msg {
411            CoreResourceMsg::Fetch(request_builder, channels) => match channels {
412                FetchChannels::ResponseMsg(sender) => {
413                    let cancellation_listener =
414                        self.get_or_create_cancellation_listener(request_builder.id);
415                    self.resource_manager.fetch(
416                        request_builder,
417                        None,
418                        sender,
419                        http_state,
420                        cancellation_listener,
421                        protocols,
422                    );
423                },
424                FetchChannels::WebSocket {
425                    event_sender,
426                    action_receiver,
427                } => {
428                    let cancellation_listener =
429                        self.get_or_create_cancellation_listener(request_builder.id);
430
431                    self.resource_manager.websocket_connect(
432                        request_builder,
433                        event_sender,
434                        action_receiver,
435                        http_state,
436                        cancellation_listener,
437                        protocols,
438                    )
439                },
440                FetchChannels::Prefetch => self.resource_manager.fetch(
441                    request_builder,
442                    None,
443                    DiscardFetch,
444                    http_state,
445                    Arc::new(Default::default()),
446                    protocols,
447                ),
448            },
449            CoreResourceMsg::Cancel(request_ids) => {
450                for cancellation_listener in request_ids
451                    .into_iter()
452                    .filter_map(|request_id| self.cancellation_listener(request_id))
453                {
454                    cancellation_listener.cancel();
455                }
456            },
457            CoreResourceMsg::DeleteCookiesForSites(sites, sender) => {
458                http_state
459                    .cookie_jar
460                    .write()
461                    .delete_cookies_for_sites(&sites);
462                let _ = sender.send(());
463            },
464            CoreResourceMsg::DeleteSessionCookies(sender) => {
465                http_state.cookie_jar.write().clear_session_cookies();
466                let _ = sender.send(());
467            },
468            CoreResourceMsg::DeleteCookies(request, sender) => {
469                http_state
470                    .cookie_jar
471                    .write()
472                    .clear_storage(request.as_ref());
473                if let Some(sender) = sender {
474                    let _ = sender.send(());
475                }
476                return true;
477            },
478            CoreResourceMsg::DeleteCookie(request, name) => {
479                http_state
480                    .cookie_jar
481                    .write()
482                    .delete_cookie_with_name(&request, name);
483                return true;
484            },
485            CoreResourceMsg::DeleteCookieAsync(cookie_store_id, url, name) => {
486                http_state
487                    .cookie_jar
488                    .write()
489                    .delete_cookie_with_name(&url, name);
490                self.send_cookie_response(cookie_store_id, CookieData::Delete(Ok(())));
491            },
492            CoreResourceMsg::FetchRedirect(request_builder, res_init, sender) => {
493                let cancellation_listener =
494                    self.get_or_create_cancellation_listener(request_builder.id);
495                self.resource_manager.fetch(
496                    request_builder,
497                    Some(res_init),
498                    sender,
499                    http_state,
500                    cancellation_listener,
501                    protocols,
502                )
503            },
504            CoreResourceMsg::SetCookieForUrl(request, cookie, source, sender) => {
505                self.resource_manager.set_cookie_for_url(
506                    &request,
507                    cookie.into_inner().to_owned(),
508                    source,
509                    http_state,
510                );
511                if let Some(sender) = sender {
512                    let _ = sender.send(());
513                }
514            },
515            CoreResourceMsg::SetCookiesForUrl(request, cookies, source) => {
516                for cookie in cookies {
517                    self.resource_manager.set_cookie_for_url(
518                        &request,
519                        cookie.into_inner(),
520                        source,
521                        http_state,
522                    );
523                }
524            },
525            CoreResourceMsg::SetCookieForUrlAsync(cookie_store_id, url, cookie, source) => {
526                self.resource_manager.set_cookie_for_url(
527                    &url,
528                    cookie.into_inner().to_owned(),
529                    source,
530                    http_state,
531                );
532                self.send_cookie_response(cookie_store_id, CookieData::Set(Ok(())));
533            },
534            CoreResourceMsg::GetCookieStringForUrl(url, consumer, source) => {
535                let mut cookie_jar = http_state.cookie_jar.write();
536                cookie_jar.remove_expired_cookies_for_url(&url);
537                consumer.send_or_ignore(cookie_jar.cookies_for_url(&url, source));
538            },
539            CoreResourceMsg::GetCookiesForUrl(url, consumer, source) => {
540                let mut cookie_jar = http_state.cookie_jar.write();
541                cookie_jar.remove_expired_cookies_for_url(&url);
542                let cookies = cookie_jar
543                    .cookies_data_for_url(&url, source)
544                    .map(Serde)
545                    .collect();
546                consumer.send_or_ignore(cookies);
547            },
548            CoreResourceMsg::GetCookieDataForUrlAsync(cookie_store_id, url, name) => {
549                let mut cookie_jar = http_state.cookie_jar.write();
550                cookie_jar.remove_expired_cookies_for_url(&url);
551                let cookie = cookie_jar
552                    .query_cookies(&url, name)
553                    .into_iter()
554                    .map(Serde)
555                    .next();
556                self.send_cookie_response(cookie_store_id, CookieData::Get(cookie));
557            },
558            CoreResourceMsg::GetAllCookieDataForUrlAsync(cookie_store_id, url, name) => {
559                let mut cookie_jar = http_state.cookie_jar.write();
560                cookie_jar.remove_expired_cookies_for_url(&url);
561                let cookies = cookie_jar
562                    .query_cookies(&url, name)
563                    .into_iter()
564                    .map(Serde)
565                    .collect();
566                self.send_cookie_response(cookie_store_id, CookieData::GetAll(cookies));
567            },
568            CoreResourceMsg::EmbedderGetCookiesForUrl(operation_id, url, source) => {
569                let mut cookie_jar = http_state.cookie_jar.write();
570                cookie_jar.remove_expired_cookies_for_url(&url);
571                let cookies: Vec<Cookie<'static>> =
572                    cookie_jar.cookies_data_for_url(&url, source).collect();
573                http_state.embedder_proxy.send(
574                    NetToEmbedderMsg::EmbedderCookieOperationResponseWithCookies(
575                        operation_id,
576                        cookies,
577                    ),
578                );
579            },
580            CoreResourceMsg::EmbedderSetCookieForUrl(operation_id, url, cookie, source) => {
581                self.resource_manager.set_cookie_for_url(
582                    &url,
583                    cookie.into_inner(),
584                    source,
585                    http_state,
586                );
587                http_state
588                    .embedder_proxy
589                    .send(NetToEmbedderMsg::EmbedderCookieOperationResponse(
590                        operation_id,
591                    ));
592            },
593            CoreResourceMsg::EmbedderClearCookies(operation_id) => {
594                http_state.cookie_jar.write().clear_storage(None);
595                http_state
596                    .embedder_proxy
597                    .send(NetToEmbedderMsg::EmbedderCookieOperationResponse(
598                        operation_id,
599                    ));
600            },
601            CoreResourceMsg::EmbedderClearSessionCookies(operation_id) => {
602                http_state.cookie_jar.write().clear_session_cookies();
603                http_state
604                    .embedder_proxy
605                    .send(NetToEmbedderMsg::EmbedderCookieOperationResponse(
606                        operation_id,
607                    ));
608            },
609            CoreResourceMsg::NewCookieListener(cookie_store_id, callback, _url) => {
610                // TODO: Use the URL for setting up the actual monitoring
611                self.cookie_listeners.insert(cookie_store_id, callback);
612            },
613            CoreResourceMsg::RemoveCookieListener(cookie_store_id) => {
614                self.cookie_listeners.remove(&cookie_store_id);
615            },
616            CoreResourceMsg::NetworkMediator(mediator_chan, origin) => {
617                self.resource_manager
618                    .sw_managers
619                    .insert(origin, mediator_chan);
620            },
621            CoreResourceMsg::ListCookies(sender) => {
622                let mut cookie_jar = http_state.cookie_jar.write();
623                cookie_jar.remove_all_expired_cookies();
624                sender.send_or_ignore(cookie_jar.cookie_site_descriptors());
625            },
626            CoreResourceMsg::GetHistoryState(history_state_id, consumer) => {
627                let history_states = http_state.history_states.read();
628                consumer.send_or_ignore(history_states.get(&history_state_id).cloned());
629            },
630            CoreResourceMsg::SetHistoryState(history_state_id, structured_data) => {
631                let mut history_states = http_state.history_states.write();
632                history_states.insert(history_state_id, structured_data);
633            },
634            CoreResourceMsg::RemoveHistoryStates(states_to_remove) => {
635                let mut history_states = http_state.history_states.write();
636                for history_state in states_to_remove {
637                    history_states.remove(&history_state);
638                }
639            },
640            CoreResourceMsg::GetCacheEntries(sender) => {
641                sender.send_or_ignore(http_state.http_cache.cache_entry_descriptors());
642            },
643            CoreResourceMsg::ClearCache(sender) => {
644                http_state.http_cache.clear();
645                if let Some(sender) = sender {
646                    sender.send_or_ignore(());
647                }
648            },
649            CoreResourceMsg::ToFileManager(msg) => self.resource_manager.filemanager.handle(msg),
650            CoreResourceMsg::TotalSizeOfInFlightKeepAliveRecords(pipeline_id, sender) => {
651                let total = self
652                    .resource_manager
653                    .in_flight_keep_alive_records
654                    .lock()
655                    .get(&pipeline_id)
656                    .map(|records| {
657                        records
658                            .iter()
659                            .map(|record| record.keep_alive_body_length)
660                            .sum()
661                    })
662                    .unwrap_or_default();
663                sender.send_or_ignore(total);
664            },
665            CoreResourceMsg::Exit(sender) => {
666                if let Some(ref config_dir) = self.config_dir {
667                    let auth_cache = http_state.auth_cache.read();
668                    servo_base::write_json_to_file(&*auth_cache, config_dir, "auth_cache.json");
669                    let jar = http_state.cookie_jar.read();
670                    servo_base::write_json_to_file(&*jar, config_dir, "cookie_jar.json");
671                    let hsts = http_state.hsts_list.read();
672                    servo_base::write_json_to_file(&*hsts, config_dir, "hsts_list.json");
673                }
674                self.resource_manager.exit();
675
676                let _ = sender.send(());
677                return false;
678            },
679            // Ignore these messages as they are only sent on very specific channels.
680            CoreResourceMsg::CollectMemoryReport(_) |
681            CoreResourceMsg::RevokeTokenForFile(..) |
682            CoreResourceMsg::RefreshTokenForFile(..) => {},
683        }
684        true
685    }
686}
687
688#[derive(Clone, Debug, Deserialize, Serialize, MallocSizeOf)]
689pub struct AuthCacheEntry {
690    pub user_name: String,
691    pub password: String,
692}
693
694impl Default for AuthCache {
695    fn default() -> Self {
696        Self {
697            version: 1,
698            entries: HashMap::new(),
699        }
700    }
701}
702
703#[derive(Clone, Debug, Deserialize, Serialize, MallocSizeOf)]
704pub struct AuthCache {
705    pub version: u32,
706    pub entries: HashMap<String, AuthCacheEntry>,
707}
708
709pub struct CoreResourceManager {
710    devtools_sender: Option<Sender<DevtoolsControlMsg>>,
711    sw_managers: HashMap<ImmutableOrigin, IpcSender<CustomResponseMediator>>,
712    filemanager: FileManager,
713    request_interceptor: RequestInterceptor,
714    ca_certificates: CACertificates<'static>,
715    ignore_certificate_errors: bool,
716    preloaded_resources: SharedPreloadedResources,
717    /// <https://fetch.spec.whatwg.org/#concept-fetch-record>
718    in_flight_keep_alive_records: SharedInflightKeepAliveRecords,
719}
720
721impl CoreResourceManager {
722    pub fn new(
723        devtools_sender: Option<Sender<DevtoolsControlMsg>>,
724        _profiler_chan: ProfilerChan,
725        embedder_proxy: GenericEmbedderProxy<NetToEmbedderMsg>,
726        ca_certificates: CACertificates<'static>,
727        ignore_certificate_errors: bool,
728        blob_token_communicator: Arc<Mutex<BlobTokenCommunicator>>,
729    ) -> CoreResourceManager {
730        CoreResourceManager {
731            devtools_sender,
732            sw_managers: Default::default(),
733            filemanager: FileManager::new(embedder_proxy.clone(), blob_token_communicator),
734            request_interceptor: RequestInterceptor::new(embedder_proxy),
735            ca_certificates,
736            ignore_certificate_errors,
737            preloaded_resources: Default::default(),
738            in_flight_keep_alive_records: Default::default(),
739        }
740    }
741
742    fn handle_preloaded_response(
743        preloaded_resources: SharedPreloadedResources,
744        preload_id: PreloadId,
745        response: Response,
746    ) {
747        // https://html.spec.whatwg.org/multipage/#preload
748        // Step 11.1. If bodyBytes is a byte sequence, then set response's body to bodyBytes as a body.
749        // Step 11.2. Otherwise, set response to a network error.
750        let response = response
751            .get_network_error()
752            .map(|_| {
753                Response::network_error(NetworkError::ResourceLoadError("Failed to preload".into()))
754            })
755            .unwrap_or(response);
756        let mut preloaded_resources = preloaded_resources.lock().unwrap();
757        // Step 11.5. If entry's on response available is null, then set entry's response to response;
758        // otherwise call entry's on response available given response.
759        if let Some(entry) = preloaded_resources.get_mut(&preload_id) {
760            entry.with_response(response);
761        }
762    }
763
764    /// Exit the core resource manager.
765    pub fn exit(&mut self) {
766        debug!("Exited CoreResourceManager");
767    }
768
769    fn set_cookie_for_url(
770        &mut self,
771        request: &ServoUrl,
772        cookie: Cookie<'static>,
773        source: CookieSource,
774        http_state: &Arc<HttpState>,
775    ) {
776        if let Some(cookie) = ServoCookie::new_wrapped(cookie, request, source) {
777            let mut cookie_jar = http_state.cookie_jar.write();
778            cookie_jar.push(cookie, request, source)
779        }
780    }
781
782    fn fetch<Target: 'static + FetchTaskTarget + Send>(
783        &self,
784        request_builder: RequestBuilder,
785        res_init_: Option<ResponseInit>,
786        mut sender: Target,
787        http_state: &Arc<HttpState>,
788        cancellation_listener: Arc<CancellationListener>,
789        protocols: Arc<ProtocolRegistry>,
790    ) {
791        let http_state = http_state.clone();
792        let devtools_chan = self.devtools_sender.clone();
793        let filemanager = self.filemanager.clone();
794        let request_interceptor = self.request_interceptor.clone();
795
796        let timing_type = match request_builder.destination {
797            Destination::Document => ResourceTimingType::Navigation,
798            _ => ResourceTimingType::Resource,
799        };
800
801        let request = request_builder.build();
802        let url = request.current_url();
803
804        // In the case of a valid blob URL, acquiring a token granting access to a file,
805        // regardless if the URL is revoked after token acquisition.
806        //
807        // Ideally all callers should have claimed the blob entry themselves, but we're not there
808        // yet.
809        let (file_token, blob_url_file_id) = match url.scheme() {
810            "blob" => {
811                if let Some(token) = request.current_url_with_blob_claim().token() {
812                    (FileTokenCheck::Required(token.token), Some(token.file_id))
813                } else if let Ok(id) = parse_blob_url(&url) {
814                    // See https://github.com/servo/servo/issues/25226
815                    log::warn!(
816                        "Failed to claim blob URL entry of valid blob URL before passing it to `net`. This causes race conditions."
817                    );
818                    (self.filemanager.get_token_for_file(&id, false), Some(id))
819                } else {
820                    (FileTokenCheck::ShouldFail, None)
821                }
822            },
823            _ => (FileTokenCheck::NotRequired, None),
824        };
825
826        let ca_certificates = self.ca_certificates.clone();
827        let ignore_certificate_errors = self.ignore_certificate_errors;
828        let in_flight_keep_alive_records = self.in_flight_keep_alive_records.clone();
829        let preloaded_resources = self.preloaded_resources.clone();
830        if let Some(ref preload_id) = request.preload_id {
831            let mut preloaded_resources = self.preloaded_resources.lock().unwrap();
832            let entry = PreloadEntry::new(request.integrity_metadata.clone());
833            preloaded_resources.insert(preload_id.clone(), entry);
834        }
835
836        spawn_task(async move {
837            // XXXManishearth: Check origin against pipeline id (also ensure that the mode is allowed)
838            // todo load context / mimesniff in fetch
839            // todo referrer policy?
840            // todo service worker stuff
841            let context = FetchContext {
842                state: http_state,
843                user_agent: servo_config::pref!(user_agent),
844                devtools_chan,
845                filemanager,
846                file_token,
847                request_interceptor: Arc::new(TokioMutex::new(request_interceptor)),
848                cancellation_listener,
849                timing: ResourceFetchTiming::new(request.timing_type()).into(),
850                protocols,
851                websocket_chan: None,
852                ca_certificates,
853                ignore_certificate_errors,
854                preloaded_resources: preloaded_resources.clone(),
855                in_flight_keep_alive_records,
856            };
857
858            match res_init_ {
859                Some(res_init) => {
860                    let response = Response::from_init(res_init, timing_type);
861
862                    let mut fetch_params = FetchParams::new(request);
863                    let mut request_body_stream_closer =
864                        AutoRequestBodyStreamCloser::new(fetch_params.request.body.as_ref());
865                    let response = http_redirect_fetch(
866                        &mut fetch_params,
867                        &mut CorsCache::default(),
868                        response,
869                        true,
870                        &mut sender,
871                        &mut None,
872                        &context,
873                    )
874                    .await;
875                    if transfers_request_body_stream_to_later_manual_redirect(
876                        &fetch_params.request,
877                        &response,
878                    ) {
879                        request_body_stream_closer.disarm();
880                    }
881                },
882                None => {
883                    let preload_id = request.preload_id.clone();
884                    let response = fetch(request, &mut sender, &context).await;
885                    if let Some(preload_id) = preload_id {
886                        Self::handle_preloaded_response(preloaded_resources, preload_id, response);
887                    }
888                },
889            };
890
891            // Remove token after fetch.
892            if let Some(id) = blob_url_file_id.as_ref() {
893                context
894                    .filemanager
895                    .invalidate_token(&context.file_token, id);
896            }
897        });
898    }
899
900    /// <https://websockets.spec.whatwg.org/#concept-websocket-establish>
901    fn websocket_connect(
902        &self,
903        mut request: RequestBuilder,
904        event_sender: IpcSender<WebSocketNetworkEvent>,
905        action_receiver: CallbackSetter<WebSocketDomAction>,
906        http_state: &Arc<HttpState>,
907        cancellation_listener: Arc<CancellationListener>,
908        protocols: Arc<ProtocolRegistry>,
909    ) {
910        let http_state = http_state.clone();
911        let devtools_chan = self.devtools_sender.clone();
912        let filemanager = self.filemanager.clone();
913        let request_interceptor = self.request_interceptor.clone();
914
915        let ca_certificates = self.ca_certificates.clone();
916        let ignore_certificate_errors = self.ignore_certificate_errors;
917        let in_flight_keep_alive_records = self.in_flight_keep_alive_records.clone();
918        let preloaded_resources = self.preloaded_resources.clone();
919
920        spawn_task(async move {
921            let mut event_sender = event_sender;
922
923            // Let requestURL be a copy of url, with its scheme set to "http", if url’s scheme is
924            // "ws"; otherwise to "https"
925            let scheme = match request.url.scheme() {
926                "ws" => "http",
927                _ => "https",
928            };
929            request
930                .url
931                .as_mut_url()
932                .set_scheme(scheme)
933                .unwrap_or_else(|_| panic!("Can't set scheme to {scheme}"));
934
935            match create_handshake_request(request, http_state.clone()) {
936                Ok(request) => {
937                    let context = FetchContext {
938                        state: http_state,
939                        user_agent: servo_config::pref!(user_agent),
940                        devtools_chan,
941                        filemanager,
942                        file_token: FileTokenCheck::NotRequired,
943                        request_interceptor: Arc::new(TokioMutex::new(request_interceptor)),
944                        cancellation_listener,
945                        timing: ResourceFetchTiming::new(request.timing_type()).into(),
946                        protocols: protocols.clone(),
947                        websocket_chan: Some(Arc::new(Mutex::new(WebSocketChannel::new(
948                            event_sender.clone(),
949                            Some(action_receiver),
950                        )))),
951                        ca_certificates,
952                        ignore_certificate_errors,
953                        preloaded_resources,
954                        in_flight_keep_alive_records,
955                    };
956                    fetch(request, &mut event_sender, &context).await;
957                },
958                Err(e) => {
959                    trace!("unable to create websocket handshake request {:?}", e);
960                    let _ = event_sender.send(WebSocketNetworkEvent::Fail);
961                },
962            }
963        });
964    }
965}