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::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
75fn 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#[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 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#[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
201fn 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 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 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 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 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 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 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 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 if let Some(entry) = preloaded_resources.get_mut(&preload_id) {
760 entry.with_response(response);
761 }
762 }
763
764 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 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 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 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 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 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 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}