Skip to main content

pkg_resolver/
main.rs

1// Copyright 2018 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5#![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
62// FIXME: allow for multiple threads and sendable futures once repo updates support it.
63// FIXME(43342): trace durations assume they start and end on the same thread, but since the
64// package resolver's executor is multi-threaded, a trace duration that includes an 'await' may not
65// end on the same thread it starts on, resulting in invalid trace events.
66// const SERVER_THREADS: usize = 2;
67
68const MAX_CONCURRENT_PACKAGE_FETCHES: usize = 5;
69
70// Each fetch_blob call emits an event, and a system update fetches about 1,000 blobs in about a
71// minute.
72const COBALT_CONNECTOR_BUFFER_SIZE: usize = 1000;
73
74const STATIC_REPO_DIR: &str = "/config/data/repositories";
75// Relative to /data.
76const DYNAMIC_REPO_PATH: &str = "repositories.json";
77
78// Relative to /config/data.
79const STATIC_RULES_PATH: &str = "rewrites.json";
80// Relative to /data.
81const 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        // Use anyhow to print the error chain.
92        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    // The list of cache packages from the system image, not to be confused with the PackageCache.
114    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    // `pkg-resolver` is required for an OTA and EagerPackageManager isn't.
224    // Also `EagerPackageManager` depends on /data, which may or may not be available, especially in
225    // tests. Wrapping `EagerPackageManager` in Arc<Option<_>> allows it to be used if available
226    // during package resolve process.
227    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    // report any errors we saw, but don't error out because otherwise we won't be able
416    // to update the system.
417    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                // Given a fresh /data, it's expected the file doesn't exist.
486                LoadRulesError::FileOpen(fuchsia_fs::node::OpenError::OpenError(
487                    zx::Status::NOT_FOUND,
488                )) => {}
489                // Unable to open /data dir proxy.
490                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                // No static rules are configured for this system version.
504                LoadRulesError::FileOpen(fuchsia_fs::node::OpenError::OpenError(
505                    zx::Status::NOT_FOUND,
506                )) => {}
507                // Unable to open /config/data dir proxy.
508                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}