1#![allow(clippy::let_unit_value)]
6#![allow(clippy::too_many_arguments)]
7#![allow(clippy::enum_variant_names)]
8
9use anyhow::{Context as _, Error, anyhow};
10use async_lock::RwLock as AsyncRwLock;
11use cobalt_sw_delivery_registry as metrics;
12use delivery_blob::DeliveryBlobType;
13use fdio::Namespace;
14use fidl::endpoints::DiscoverableProtocolMarker as _;
15use fidl_contrib::ProtocolConnector;
16use fidl_contrib::protocol_connector::ProtocolSender;
17use fidl_fuchsia_io as fio;
18use fidl_fuchsia_metrics as fmetrics;
19use fidl_fuchsia_pkg as fpkg;
20use fidl_fuchsia_pkg_http as fpkg_http;
21use fuchsia_async as fasync;
22use fuchsia_cobalt_builders::MetricEventExt as _;
23use fuchsia_component::server::ServiceFs;
24use fuchsia_inspect as inspect;
25use fuchsia_trace as ftrace;
26use futures::prelude::*;
27use futures::stream::FuturesUnordered;
28use log::{error, info, warn};
29use std::sync::Arc;
30use std::time::{Duration, Instant};
31
32mod authority_service;
33mod cache;
34mod cache_package_index;
35mod clock;
36mod component_resolver;
37mod config;
38mod eager_package_manager;
39mod error;
40mod inspect_util;
41mod metrics_util;
42mod ota_downloader;
43mod repository;
44mod repository_manager;
45mod repository_service;
46mod resolver_service;
47mod rewrite_manager;
48mod rewrite_service;
49mod util;
50
51#[cfg(test)]
52mod test_util;
53
54use crate::cache::BasePackageIndex;
55use crate::config::Config;
56use crate::repository_manager::{RepositoryManager, RepositoryManagerBuilder};
57use crate::repository_service::RepositoryService;
58use crate::resolver_service::ResolverServiceInspectState;
59use crate::rewrite_manager::{LoadRulesError, RewriteManager, RewriteManagerBuilder};
60use crate::rewrite_service::RewriteService;
61
62const MAX_CONCURRENT_PACKAGE_FETCHES: usize = 5;
69
70const COBALT_CONNECTOR_BUFFER_SIZE: usize = 1000;
73
74const STATIC_REPO_DIR: &str = "/config/data/repositories";
75const DYNAMIC_REPO_PATH: &str = "repositories.json";
77
78const STATIC_RULES_PATH: &str = "rewrites.json";
80const DYNAMIC_RULES_PATH: &str = "rewrites.json";
82
83#[fuchsia::main(logging_tags = ["pkg-resolver"])]
84pub fn main() -> Result<(), Error> {
85 let startup_time = Instant::now();
86 fuchsia_trace_provider::trace_provider_create_with_fdio();
87 info!("starting package resolver");
88
89 let mut executor = fasync::LocalExecutorBuilder::new().build();
90 executor.run_singlethreaded(main_inner_async(startup_time)).map_err(|err| {
91 let err = anyhow!(err);
93 error!("error running pkg-resolver: {:#}", err);
94 err
95 })
96}
97
98async fn main_inner_async(startup_time: Instant) -> Result<(), Error> {
99 let config = Config::load_from_config_data_or_default();
100 let structured_config = pkg_resolver_config::Config::take_from_startup_handle();
101
102 let pkg_cache_proxy =
103 fuchsia_component::client::connect_to_protocol::<fpkg::PackageCacheMarker>()
104 .context("error connecting to package cache")?;
105 let pkg_cache = fidl_fuchsia_pkg_ext::cache::Client::from_proxy(pkg_cache_proxy);
106
107 let base_package_index = Arc::new(
108 BasePackageIndex::from_proxy(pkg_cache.proxy())
109 .await
110 .context("failed to load base package index")?,
111 );
112
113 let system_cache_list = Arc::new(cache_package_index::from_proxy(pkg_cache.proxy()).await);
115
116 let inspector = fuchsia_inspect::Inspector::default();
117 inspector
118 .root()
119 .record_child("structured_config", |node| structured_config.record_inspect(node));
120
121 let futures = FuturesUnordered::new();
122
123 let (mut cobalt_sender, cobalt_fut) = ProtocolConnector::new_with_buffer_size(
124 metrics_util::CobaltConnectedService,
125 COBALT_CONNECTOR_BUFFER_SIZE,
126 )
127 .serve_and_log_errors();
128 futures.push(cobalt_fut.boxed_local());
129
130 let data_proxy = match fuchsia_fs::directory::open_in_namespace(
131 "/data",
132 fio::PERM_READABLE | fio::PERM_WRITABLE,
133 ) {
134 Ok(proxy) => Some(proxy),
135 Err(e) => {
136 warn!("failed to open /data: {:#}", anyhow!(e));
137 None
138 }
139 };
140
141 if data_proxy.is_some() {
142 let namespace = Namespace::installed().context("failed to get installed namespace")?;
143 namespace.unbind("/data").context("failed to unbind /data from default namespace")?;
144 }
145
146 let config_proxy =
147 match fuchsia_fs::directory::open_in_namespace("/config/data", fio::PERM_READABLE) {
148 Ok(proxy) => Some(proxy),
149 Err(e) => {
150 warn!("failed to open /config/data: {:#}", anyhow!(e));
151 None
152 }
153 };
154
155 let delivery_blob_type: DeliveryBlobType =
156 structured_config.delivery_blob_type.try_into().with_context(|| {
157 format!("invalid delivery blob type {}", structured_config.delivery_blob_type)
158 })?;
159
160 let repo_manager = Arc::new(AsyncRwLock::new(
161 load_repo_manager(
162 inspector.root().create_child("repository_manager"),
163 &config,
164 cobalt_sender.clone(),
165 TufTimeouts {
166 metadata: std::time::Duration::from_secs(
167 structured_config.tuf_metadata_timeout_seconds.into(),
168 ),
169 network_header: zx::BootDuration::from_seconds(
170 structured_config.tuf_network_header_timeout_seconds.into(),
171 ),
172 },
173 data_proxy.clone(),
174 delivery_blob_type,
175 )
176 .await,
177 ));
178 let rewrite_manager = Arc::new(AsyncRwLock::new(
179 load_rewrite_manager(
180 inspector.root().create_child("rewrite_manager"),
181 &config,
182 data_proxy.clone(),
183 config_proxy,
184 )
185 .await,
186 ));
187
188 let (blob_fetch_queue, blob_fetcher) = crate::cache::BlobFetcher::new(
189 fuchsia_component::client::connect_to_protocol::<fpkg_http::ClientMarker>()
190 .context("error connecting to fuchsia.pkg.http/Client")?,
191 inspector.root().create_child("blob_fetcher"),
192 structured_config.blob_download_concurrency_limit.into(),
193 repo_manager.read().await.stats(),
194 cache::BlobFetchParams::builder()
195 .header_network_timeout(zx::BootDuration::from_seconds(
196 structured_config.blob_network_header_timeout_seconds.into(),
197 ))
198 .body_network_timeout(zx::BootDuration::from_seconds(
199 structured_config.blob_network_body_timeout_seconds.into(),
200 ))
201 .download_resumption_attempts_limit(
202 structured_config.blob_download_resumption_attempts_limit,
203 )
204 .build(),
205 );
206 futures.push(blob_fetch_queue.boxed_local());
207
208 let resolver_service_inspect_state = Arc::new(ResolverServiceInspectState::from_node(
209 inspector.root().create_child("resolver_service"),
210 ));
211 let (package_fetch_queue, package_resolver) = resolver_service::QueuedResolver::new(
212 pkg_cache.clone(),
213 Arc::clone(&base_package_index),
214 Arc::clone(&system_cache_list),
215 Arc::clone(&repo_manager),
216 Arc::clone(&rewrite_manager),
217 blob_fetcher.clone(),
218 MAX_CONCURRENT_PACKAGE_FETCHES,
219 Arc::clone(&resolver_service_inspect_state),
220 );
221 futures.push(package_fetch_queue.boxed_local());
222
223 let eager_package_manager = Arc::new(
228 crate::eager_package_manager::EagerPackageManager::from_namespace(
229 package_resolver.clone(),
230 pkg_cache.clone(),
231 data_proxy,
232 &system_cache_list,
233 cobalt_sender.clone(),
234 )
235 .await
236 .map_err(|e| {
237 error!("failed to create EagerPackageManager: {:#}", e);
238 })
239 .ok()
240 .map(AsyncRwLock::new),
241 );
242
243 let make_resolver_cb = {
244 let repo_manager = Arc::clone(&repo_manager);
245 let rewrite_manager = Arc::clone(&rewrite_manager);
246 let package_resolver = package_resolver.clone();
247 let pkg_cache = pkg_cache.clone();
248 let cobalt_sender = cobalt_sender.clone();
249 let eager_package_manager = Arc::clone(&eager_package_manager);
250 move |gc_protection| {
251 let repo_manager = Arc::clone(&repo_manager);
252 let rewrite_manager = Arc::clone(&rewrite_manager);
253 let package_resolver = package_resolver.clone();
254 let pkg_cache = pkg_cache.clone();
255 let base_package_index = Arc::clone(&base_package_index);
256 let system_cache_list = Arc::clone(&system_cache_list);
257 let cobalt_sender = cobalt_sender.clone();
258 let resolver_service_inspect_state = Arc::clone(&resolver_service_inspect_state);
259 let eager_package_manager = Arc::clone(&eager_package_manager);
260 move |stream| {
261 fasync::Task::local(
262 resolver_service::run_resolver_service(
263 Arc::clone(&repo_manager),
264 Arc::clone(&rewrite_manager),
265 package_resolver.clone(),
266 pkg_cache.clone(),
267 Arc::clone(&base_package_index),
268 Arc::clone(&system_cache_list),
269 stream,
270 gc_protection,
271 cobalt_sender.clone(),
272 Arc::clone(&resolver_service_inspect_state),
273 Arc::clone(&eager_package_manager),
274 )
275 .unwrap_or_else(|e| error!("run_resolver_service failed: {:#}", anyhow!(e))),
276 )
277 .detach()
278 }
279 }
280 };
281
282 let resolver_toolbox_cb = {
283 let package_resolver = package_resolver.clone();
284 let cobalt_sender = cobalt_sender.clone();
285 let eager_package_manager = Arc::clone(&eager_package_manager);
286 move |stream| {
287 fasync::Task::local(
288 resolver_service::run_resolver_toolbox_service(
289 package_resolver.clone(),
290 stream,
291 fpkg::GcProtection::OpenPackageTracking,
292 cobalt_sender.clone(),
293 Arc::clone(&eager_package_manager),
294 )
295 .unwrap_or_else(|e| {
296 error!("run_resolver_toolbox_service failed: {:#}", anyhow!(e))
297 }),
298 )
299 .detach()
300 }
301 };
302
303 let authority_cb = {
304 let rewrite_manager = Arc::clone(&rewrite_manager);
305 let repo_manager = Arc::clone(&repo_manager);
306 move |stream| {
307 fasync::Task::local(
308 authority_service::serve(
309 stream,
310 Arc::clone(&rewrite_manager),
311 Arc::clone(&repo_manager),
312 )
313 .unwrap_or_else(|e: anyhow::Error| error!("serving authority service: {e:#}")),
314 )
315 .detach()
316 }
317 };
318
319 let repo_cb = move |stream| {
320 let repo_manager = Arc::clone(&repo_manager);
321
322 fasync::Task::local(
323 async move {
324 let mut repo_service = RepositoryService::new(repo_manager);
325 repo_service.run(stream).await
326 }
327 .unwrap_or_else(|e| error!("error encountered: {:#}", anyhow!(e))),
328 )
329 .detach()
330 };
331
332 let rewrite_cb = move |stream| {
333 let mut rewrite_service = RewriteService::new(Arc::clone(&rewrite_manager));
334
335 fasync::Task::local(
336 async move { rewrite_service.handle_client(stream).await }
337 .unwrap_or_else(|e| error!("while handling rewrite client {:#}", anyhow!(e))),
338 )
339 .detach()
340 };
341
342 let cup_cb = {
343 let cobalt_sender = cobalt_sender.clone();
344 move |stream| {
345 fasync::Task::local(
346 eager_package_manager::run_cup_service(
347 Arc::clone(&eager_package_manager),
348 stream,
349 cobalt_sender.clone(),
350 )
351 .unwrap_or_else(|e| error!("run_cup_service failed: {:#}", anyhow!(e))),
352 )
353 .detach()
354 }
355 };
356
357 let component_resolver_cb = move |stream| {
358 fasync::Task::local(
359 component_resolver::serve(stream)
360 .unwrap_or_else(|e| error!("serve_component_resolver_failed: {:#}", e)),
361 )
362 .detach()
363 };
364
365 let ota_downloader_cb = move |stream| {
366 fasync::Task::local(
367 ota_downloader::serve(stream, blob_fetcher.clone(), pkg_cache.clone())
368 .unwrap_or_else(|e| error!("run_cup_service failed: {:#}", anyhow!(e))),
369 )
370 .detach()
371 };
372
373 let mut fs = ServiceFs::new();
374 fs.dir("svc")
375 .add_fidl_service(make_resolver_cb(fpkg::GcProtection::OpenPackageTracking))
376 .add_fidl_service_at(
377 format!("{}-ota", fpkg::PackageResolverMarker::PROTOCOL_NAME),
378 make_resolver_cb(fpkg::GcProtection::Retained),
379 )
380 .add_fidl_service(resolver_toolbox_cb)
381 .add_fidl_service(authority_cb)
382 .add_fidl_service(repo_cb)
383 .add_fidl_service(rewrite_cb)
384 .add_fidl_service(cup_cb)
385 .add_fidl_service(component_resolver_cb)
386 .add_fidl_service(ota_downloader_cb);
387
388 fs.take_and_serve_directory_handle().context("while serving directory handle")?;
389
390 futures.push(fs.collect().boxed_local());
391
392 cobalt_sender.send(
393 fmetrics::MetricEvent::builder(metrics::PKG_RESOLVER_STARTUP_DURATION_MIGRATED_METRIC_ID)
394 .as_integer(Instant::now().duration_since(startup_time).as_micros() as i64),
395 );
396
397 ftrace::instant!("app", "startup", ftrace::Scope::Process);
398
399 let _inspect_server_task =
400 inspect_runtime::publish(&inspector, inspect_runtime::PublishOptions::default());
401
402 futures.collect::<()>().await;
403
404 Ok(())
405}
406
407async fn load_repo_manager(
408 node: inspect::Node,
409 config: &Config,
410 mut cobalt_sender: ProtocolSender<fmetrics::MetricEvent>,
411 tuf_timeouts: TufTimeouts,
412 data_proxy: Option<fio::DirectoryProxy>,
413 delivery_blob_type: DeliveryBlobType,
414) -> RepositoryManager {
415 let dynamic_repo_path =
418 if config.enable_dynamic_configuration() { Some(DYNAMIC_REPO_PATH) } else { None };
419 let builder = match RepositoryManagerBuilder::new(data_proxy, dynamic_repo_path, tuf_timeouts)
420 .await
421 .unwrap_or_else(|(builder, err)| {
422 error!("error loading dynamic repo config: {:#}", anyhow!(err));
423 builder
424 })
425 .delivery_blob_type(delivery_blob_type)
426 .inspect_node(node)
427 .load_static_configs_dir(STATIC_REPO_DIR)
428 {
429 Ok(builder) => {
430 cobalt_sender.send(
431 fmetrics::MetricEvent::builder(
432 metrics::REPOSITORY_MANAGER_LOAD_STATIC_CONFIGS_MIGRATED_METRIC_ID,
433 )
434 .with_event_codes(
435 metrics::RepositoryManagerLoadStaticConfigsMigratedMetricDimensionResult::Success,
436 )
437 .as_occurrence(1),
438 );
439 builder
440 }
441 Err((builder, errs)) => {
442 for err in errs {
443 let dimension_result: metrics::RepositoryManagerLoadStaticConfigsMigratedMetricDimensionResult
444 = (&err).into();
445 cobalt_sender.send(
446 fmetrics::MetricEvent::builder(
447 metrics::REPOSITORY_MANAGER_LOAD_STATIC_CONFIGS_MIGRATED_METRIC_ID,
448 )
449 .with_event_codes(dimension_result)
450 .as_occurrence(1),
451 );
452 match &err {
453 crate::repository_manager::LoadError::Io { path: _, error }
454 if error.kind() == std::io::ErrorKind::NotFound =>
455 {
456 info!("no statically configured repositories present");
457 }
458 _ => error!("error loading static repo config: {:#}", anyhow!(err)),
459 };
460 }
461 builder
462 }
463 };
464
465 match config.persisted_repos_dir() {
466 Some(repo) => builder.with_persisted_repos_dir(repo),
467 None => builder,
468 }
469 .cobalt_sender(cobalt_sender)
470 .build()
471}
472
473async fn load_rewrite_manager(
474 node: inspect::Node,
475 config: &Config,
476 data_proxy: Option<fio::DirectoryProxy>,
477 config_proxy: Option<fio::DirectoryProxy>,
478) -> RewriteManager {
479 let dynamic_rules_path =
480 if config.enable_dynamic_configuration() { Some(DYNAMIC_RULES_PATH) } else { None };
481 let builder = RewriteManagerBuilder::new(data_proxy, dynamic_rules_path)
482 .await
483 .unwrap_or_else(|(builder, err)| {
484 match err {
485 LoadRulesError::FileOpen(fuchsia_fs::node::OpenError::OpenError(
487 zx::Status::NOT_FOUND,
488 )) => {}
489 LoadRulesError::DirOpen(_) => {}
491 err => error!(
492 "unable to load dynamic rewrite rules from disk, using defaults: {:#}",
493 anyhow!(err)
494 ),
495 };
496 builder
497 })
498 .inspect_node(node)
499 .static_rules_path(config_proxy, STATIC_RULES_PATH)
500 .await
501 .unwrap_or_else(|(builder, err)| {
502 match err {
503 LoadRulesError::FileOpen(fuchsia_fs::node::OpenError::OpenError(
505 zx::Status::NOT_FOUND,
506 )) => {}
507 LoadRulesError::DirOpen(_) => {}
509 err => {
510 error!("unable to load static rewrite rules from disk: {:#}", anyhow!(err))
511 }
512 };
513 builder
514 });
515
516 builder.build()
517}
518
519#[derive(Clone, Copy, Debug)]
520struct TufTimeouts {
521 metadata: Duration,
522 network_header: zx::BootDuration,
523}