1use 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
76fn 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#[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 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#[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
202fn 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 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 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 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 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 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 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 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 if let Some(entry) = preloaded_resources.get_mut(&preload_id) {
761 entry.with_response(response);
762 }
763 }
764
765 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 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 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 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 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 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 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}