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