/build/source/nativelink-store/src/default_store_factory.rs
Line | Count | Source |
1 | | // Copyright 2024 The NativeLink Authors. All rights reserved. |
2 | | // |
3 | | // Licensed under the Functional Source License, Version 1.1, Apache 2.0 Future License (the "License"); |
4 | | // you may not use this file except in compliance with the License. |
5 | | // You may obtain a copy of the License at |
6 | | // |
7 | | // See LICENSE file for details |
8 | | // |
9 | | // Unless required by applicable law or agreed to in writing, software |
10 | | // distributed under the License is distributed on an "AS IS" BASIS, |
11 | | // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
12 | | // See the License for the specific language governing permissions and |
13 | | // limitations under the License. |
14 | | |
15 | | use core::pin::Pin; |
16 | | use std::sync::Arc; |
17 | | use std::time::SystemTime; |
18 | | |
19 | | use futures::stream::FuturesOrdered; |
20 | | use futures::{Future, TryStreamExt}; |
21 | | use nativelink_config::stores::{ExperimentalCloudObjectSpec, RedisMode, StoreSpec}; |
22 | | use nativelink_error::Error; |
23 | | use nativelink_util::health_utils::HealthRegistryBuilder; |
24 | | use nativelink_util::store_trait::{Store, StoreDriver}; |
25 | | |
26 | | use crate::azure_blob_store::AzureBlobStore; |
27 | | use crate::cache_metrics_store::CacheMetricsStore; |
28 | | use crate::completeness_checking_store::CompletenessCheckingStore; |
29 | | use crate::compression_store::CompressionStore; |
30 | | use crate::dedup_store::DedupStore; |
31 | | use crate::existence_cache_store::ExistenceCacheStore; |
32 | | use crate::fast_slow_store::FastSlowStore; |
33 | | use crate::filesystem_store::FilesystemStore; |
34 | | use crate::gcs_store::GcsStore; |
35 | | use crate::grpc_store::GrpcStore; |
36 | | use crate::memory_store::MemoryStore; |
37 | | use crate::mongo_store::ExperimentalMongoStore; |
38 | | use crate::noop_store::NoopStore; |
39 | | use crate::oci_store::OciStore; |
40 | | use crate::ontap_s3_existence_cache_store::OntapS3ExistenceCache; |
41 | | use crate::ontap_s3_store::OntapS3Store; |
42 | | use crate::r2_store::R2Store; |
43 | | use crate::redis_store::RedisStore; |
44 | | use crate::ref_store::RefStore; |
45 | | use crate::s3_store::S3Store; |
46 | | use crate::shard_store::ShardStore; |
47 | | use crate::size_partitioning_store::SizePartitioningStore; |
48 | | use crate::store_manager::StoreManager; |
49 | | use crate::verify_store::VerifyStore; |
50 | | |
51 | | type FutureMaybeStore<'a> = Box<dyn Future<Output = Result<Store, Error>> + Send + 'a>; |
52 | | |
53 | 103 | pub fn store_factory<'a>( |
54 | 103 | backend: &'a StoreSpec, |
55 | 103 | store_manager: &'a Arc<StoreManager>, |
56 | 103 | maybe_health_registry_builder: Option<&'a mut HealthRegistryBuilder>, |
57 | 103 | ) -> Pin<FutureMaybeStore<'a>> { |
58 | 103 | Box::pin(async move { |
59 | 103 | let store: Arc<dyn StoreDriver> = match backend { |
60 | 2 | StoreSpec::CacheMetrics(spec) => CacheMetricsStore::new( |
61 | 2 | spec, |
62 | 2 | store_factory(&spec.backend, store_manager, None).await?0 , |
63 | | ), |
64 | 94 | StoreSpec::Memory(spec) => MemoryStore::new(spec), |
65 | 0 | StoreSpec::ExperimentalCloudObjectStore(spec) => match spec { |
66 | 0 | ExperimentalCloudObjectSpec::Aws(aws_config) => { |
67 | 0 | S3Store::new(aws_config, SystemTime::now).await? |
68 | | } |
69 | 0 | ExperimentalCloudObjectSpec::Ontap(ontap_config) => { |
70 | 0 | OntapS3Store::new(ontap_config, SystemTime::now).await? |
71 | | } |
72 | 0 | ExperimentalCloudObjectSpec::Gcs(gcs_config) => { |
73 | 0 | GcsStore::new(gcs_config, SystemTime::now).await? |
74 | | } |
75 | 0 | ExperimentalCloudObjectSpec::Azure(azure_config) => { |
76 | 0 | AzureBlobStore::new(azure_config, SystemTime::now).await? |
77 | | } |
78 | 0 | ExperimentalCloudObjectSpec::R2(r2_config) => { |
79 | 0 | R2Store::new(r2_config, SystemTime::now).await? |
80 | | } |
81 | 0 | ExperimentalCloudObjectSpec::Oci(oci_config) => { |
82 | 0 | OciStore::new(oci_config, SystemTime::now).await? |
83 | | } |
84 | | }, |
85 | 0 | StoreSpec::RedisStore(spec) => { |
86 | 0 | if spec.mode == RedisMode::Cluster { |
87 | 0 | RedisStore::new_cluster(spec.clone()).await? |
88 | | } else { |
89 | 0 | RedisStore::new_standard(spec.clone()).await? |
90 | | } |
91 | | } |
92 | 3 | StoreSpec::Verify(spec) => VerifyStore::new( |
93 | 3 | spec, |
94 | 3 | store_factory(&spec.backend, store_manager, None).await?0 , |
95 | | ), |
96 | 0 | StoreSpec::Compression(spec) => CompressionStore::new( |
97 | 0 | &spec.clone(), |
98 | 0 | store_factory(&spec.backend, store_manager, None).await?, |
99 | 0 | )?, |
100 | 0 | StoreSpec::Dedup(spec) => DedupStore::new( |
101 | 0 | spec, |
102 | 0 | store_factory(&spec.index_store, store_manager, None).await?, |
103 | 0 | store_factory(&spec.content_store, store_manager, None).await?, |
104 | 0 | )?, |
105 | 0 | StoreSpec::ExistenceCache(spec) => ExistenceCacheStore::new( |
106 | 0 | spec, |
107 | 0 | store_factory(&spec.backend, store_manager, None).await?, |
108 | | ), |
109 | 3 | StoreSpec::OntapS3ExistenceCache(spec) => { |
110 | 3 | OntapS3ExistenceCache::new(spec, SystemTime::now).await?0 |
111 | | } |
112 | 0 | StoreSpec::CompletenessChecking(spec) => CompletenessCheckingStore::new( |
113 | 0 | store_factory(&spec.backend, store_manager, None).await?, |
114 | 0 | store_factory(&spec.cas_store, store_manager, None).await?, |
115 | | ), |
116 | 0 | StoreSpec::FastSlow(spec) => FastSlowStore::new( |
117 | 0 | spec, |
118 | 0 | store_factory(&spec.fast, store_manager, None).await?, |
119 | 0 | store_factory(&spec.slow, store_manager, None).await?, |
120 | | ), |
121 | 0 | StoreSpec::Filesystem(spec) => <FilesystemStore>::new(spec).await?, |
122 | 0 | StoreSpec::RefStore(spec) => RefStore::new(spec, Arc::downgrade(store_manager)), |
123 | 0 | StoreSpec::SizePartitioning(spec) => SizePartitioningStore::new( |
124 | 0 | spec, |
125 | 0 | store_factory(&spec.lower_store, store_manager, None).await?, |
126 | 0 | store_factory(&spec.upper_store, store_manager, None).await?, |
127 | | ), |
128 | 1 | StoreSpec::Grpc(spec) => GrpcStore::new(spec).await?0 , |
129 | 0 | StoreSpec::Noop(_) => NoopStore::new(), |
130 | 0 | StoreSpec::ExperimentalMongo(spec) => ExperimentalMongoStore::new(spec.clone()).await?, |
131 | 0 | StoreSpec::Shard(spec) => { |
132 | 0 | let stores = spec |
133 | 0 | .stores |
134 | 0 | .iter() |
135 | 0 | .map(|store_spec| store_factory(&store_spec.store, store_manager, None)) |
136 | 0 | .collect::<FuturesOrdered<_>>() |
137 | 0 | .try_collect::<Vec<_>>() |
138 | 0 | .await?; |
139 | 0 | ShardStore::new(spec, stores)? |
140 | | } |
141 | | }; |
142 | | |
143 | 103 | if let Some(health_registry_builder0 ) = maybe_health_registry_builder { |
144 | 0 | store.clone().register_health(health_registry_builder); |
145 | 103 | } |
146 | | |
147 | 103 | Ok(Store::new(store)) |
148 | 103 | }) |
149 | 103 | } |