Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion ai-docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -351,7 +351,7 @@ store. Epoch IDs are compared as big-endian integers.

### Centralized

A single `RealCentralizedKms` instance holds all key material. No MPC; keys
A single `CentralizedKms` instance holds all key material. No MPC; keys
live in the configured vault backend. Preprocessing / reshare RPCs are not
applicable.

Expand Down
4 changes: 2 additions & 2 deletions core/service/src/bin/kms-server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ use kms_lib::{
signatures::NodeSigningIdentity,
},
engine::{
backup_operator::boot_base_kms, centralized::central_kms::RealCentralizedKms,
backup_operator::boot_base_kms, centralized::central_kms::CentralizedKms,
context::SoftwareVersion, context_manager::create_default_centralized_context_in_storage,
migration::migrate_to_0_15_x, rng_source::RngSource, run_server,
threshold::service::new_real_threshold_kms,
Expand Down Expand Up @@ -680,7 +680,7 @@ async fn main_exec() -> anyhow::Result<()> {
.await?;
// A node without its signing key boots in recovery mode, as a threshold node does,
// so that it can recover its keys from the custodians.
let (kms, (health, health_service)) = RealCentralizedKms::new_from_base_kms(
let (kms, (health, health_service)) = CentralizedKms::new_from_base_kms(
core_config,
public_vault,
private_vault,
Expand Down
40 changes: 19 additions & 21 deletions core/service/src/client/test_tools.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,11 @@ use crate::consts::{DEC_CAPACITY, DEFAULT_PROTOCOL, DEFAULT_URL, MAX_TRIES, MIN_
use crate::cryptography::signatures::PublicSigKey;
use crate::engine::backup_operator::boot_base_kms;
use crate::engine::base::BaseKmsStruct;
use crate::engine::centralized::central_kms::RealCentralizedKms;
use crate::engine::centralized::central_kms::CentralizedKms;
use crate::engine::context_manager::create_default_centralized_context_in_storage;
use crate::engine::rng_source::test_rng_source;
use crate::engine::threshold::service::{RealThresholdKms, new_real_threshold_kms};
use crate::engine::threshold::service::new_real_threshold_kms;
use crate::engine::threshold::threshold_kms::ThresholdKms;
use crate::engine::{Shutdown, run_server};
use crate::grpc::MetaStoreStatusServiceImpl;
use crate::util::rate_limiter::RateLimiterConfig;
Expand Down Expand Up @@ -393,7 +394,7 @@ pub async fn setup_threshold_with_custom_peers<

// Note: explicit some of the types to avoid clippy complaining
let server: anyhow::Result<(
RealThresholdKms<PubS, PrivS>,
ThresholdKms<PubS, PrivS>,
(HealthState, _),
MetaStoreStatusServiceImpl,
)> = new_real_threshold_kms(
Expand Down Expand Up @@ -755,7 +756,7 @@ pub async fn setup_centralized_no_client<
let config_path = format!("{}/config/default_centralized", env!("CARGO_MANIFEST_DIR"));
let mut core_config: CoreConfig = init_conf(&config_path).expect("config must parse");
core_config.rate_limiter_conf = rate_limiter_conf;
let (kms, (health, health_service)) = RealCentralizedKms::new(
let (kms, (health, health_service)) = CentralizedKms::new(
core_config,
pub_storage,
priv_storage,
Expand Down Expand Up @@ -790,9 +791,8 @@ pub async fn setup_centralized_no_client<
.await
.expect("Could not start server");
});
let service_name = <CoreServiceEndpointServer<
RealCentralizedKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let service_name =
<CoreServiceEndpointServer<CentralizedKms<FileStorage, FileStorage>> as NamedService>::NAME;
await_server_ready(service_name, listen_port).await;
ServerHandle::new_centralized(arc_kms_clone, listen_port, tx, handle_health)
}
Expand Down Expand Up @@ -901,17 +901,16 @@ pub async fn setup_recovery_mode<
.unwrap();
let config_path = format!("{}/config/default_centralized", env!("CARGO_MANIFEST_DIR"));
let core_config: CoreConfig = init_conf(&config_path).expect("config must parse");
let (kms, (health, health_service)) =
RealCentralizedKms::<PubS, PrivS>::new_from_base_kms(
core_config,
pub_storage,
priv_storage,
Some(backup_vault),
None,
base_kms,
)
.await
.expect("a server without its signing key must boot in recovery mode");
let (kms, (health, health_service)) = CentralizedKms::<PubS, PrivS>::new_from_base_kms(
core_config,
pub_storage,
priv_storage,
Some(backup_vault),
None,
base_kms,
)
.await
.expect("a server without its signing key must boot in recovery mode");
let kms = Arc::new(kms);
let server = Arc::clone(&kms);
let handle_health = health.clone();
Expand Down Expand Up @@ -999,9 +998,8 @@ pub async fn setup_recovery_mode<
}
};
// The service name does not depend on the type parameters of the server.
let service_name = <CoreServiceEndpointServer<
RealCentralizedKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let service_name =
<CoreServiceEndpointServer<CentralizedKms<FileStorage, FileStorage>> as NamedService>::NAME;
await_server_ready(service_name, service_port).await;
let uri = Uri::from_str(&format!(
"{DEFAULT_PROTOCOL}://{DEFAULT_URL}:{service_port}"
Expand Down
12 changes: 5 additions & 7 deletions core/service/src/client/tests/centralized/misc_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
use crate::client::tests::common::{PollConfig, retrying_poll};
use crate::client::tests::common::{get_pub_dec_resp, send_dec_reqs};
use crate::consts::TEST_CENTRAL_KEY_ID;
use crate::engine::centralized::central_kms::RealCentralizedKms;
use crate::engine::centralized::central_kms::CentralizedKms;
use crate::testing::prelude::*;
use crate::testing::utils::{get_health_client, get_status};
use kms_grpc::kms_service::v1::core_service_endpoint_server::CoreServiceEndpointServer;
Expand Down Expand Up @@ -44,9 +44,8 @@ async fn test_central_health_endpoint_availability() -> Result<()> {
let mut health_client = get_health_client(env.server.service_port)
.await
.expect("Failed to get health client");
let service_name = <CoreServiceEndpointServer<
RealCentralizedKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let service_name =
<CoreServiceEndpointServer<CentralizedKms<FileStorage, FileStorage>> as NamedService>::NAME;
let request = tonic::Request::new(HealthCheckRequest {
service: service_name.to_string(),
});
Expand Down Expand Up @@ -103,9 +102,8 @@ async fn test_central_close_after_drop() -> Result<()> {
let mut health_client = get_health_client(kms_server.service_port)
.await
.expect("Failed to get health client");
let service_name = <CoreServiceEndpointServer<
RealCentralizedKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let service_name =
<CoreServiceEndpointServer<CentralizedKms<FileStorage, FileStorage>> as NamedService>::NAME;
let request = tonic::Request::new(HealthCheckRequest {
service: service_name.to_string(),
});
Expand Down
17 changes: 7 additions & 10 deletions core/service/src/client/tests/threshold/misc_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use crate::client::tests::common::send_dec_reqs;
use crate::client::tests::common::{PollConfig, retrying_poll};
use crate::consts::TEST_THRESHOLD_KEY_ID_4P;
use crate::consts::{DEFAULT_EPOCH_ID, DEFAULT_MPC_CONTEXT};
use crate::engine::threshold::service::RealThresholdKms;
use crate::engine::threshold::threshold_kms::ThresholdKms;
use crate::engine::utils::make_extra_data;
use crate::testing::material::{material_subdir, threshold_material_subdir};
use crate::testing::prelude::*;
Expand Down Expand Up @@ -59,9 +59,8 @@ async fn test_threshold_health_endpoint_availability() -> Result<()> {
let servers = env.servers;

// Wait for all core servers to be ready before sending requests
let core_service_name = <CoreServiceEndpointServer<
RealThresholdKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let core_service_name =
<CoreServiceEndpointServer<ThresholdKms<FileStorage, FileStorage>> as NamedService>::NAME;
for cur_handle in servers.values() {
await_server_ready(core_service_name, cur_handle.service_port).await;
}
Expand Down Expand Up @@ -183,9 +182,8 @@ async fn test_threshold_close_after_drop() -> Result<()> {
let mut core_health_client = get_health_client(servers.get(&1).unwrap().service_port)
.await
.expect("Failed to get core health client");
let core_service_name = <CoreServiceEndpointServer<
RealThresholdKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let core_service_name =
<CoreServiceEndpointServer<ThresholdKms<FileStorage, FileStorage>> as NamedService>::NAME;

// Get health client for MPC threshold service on server 1
let mut threshold_health_client = get_health_client(servers.get(&1).unwrap().mpc_port.unwrap())
Expand Down Expand Up @@ -257,9 +255,8 @@ async fn test_threshold_shutdown() -> Result<()> {
let mut servers = env.servers;

// Ensure that the servers are ready
let core_service_name = <CoreServiceEndpointServer<
RealThresholdKms<FileStorage, FileStorage>,
> as NamedService>::NAME;
let core_service_name =
<CoreServiceEndpointServer<ThresholdKms<FileStorage, FileStorage>> as NamedService>::NAME;
for cur_handle in servers.values() {
await_server_ready(core_service_name, cur_handle.service_port).await;
}
Expand Down
20 changes: 9 additions & 11 deletions core/service/src/engine/backup_operator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,7 @@ use crate::{
signatures::{PrivateSigKey, PublicSigKey},
},
engine::{
base::BaseKmsStruct, threshold::service::ThresholdFheKeys, traits::BackupOperator,
validation::RequestIdParsingErr,
base::BaseKmsStruct, threshold::service::ThresholdFheKeys, validation::RequestIdParsingErr,
},
vault::{
Vault, VaultDataType,
Expand Down Expand Up @@ -68,7 +67,7 @@ use tokio::sync::{Mutex, MutexGuard};
use tonic::{Request, Response};
use zeroize::Zeroizing;

pub struct RealBackupOperator<
pub(crate) struct RealBackupOperator<
PubS: Storage + Sync + Send + 'static,
PrivS: StorageExt + Sync + Send + 'static,
> {
Expand Down Expand Up @@ -214,7 +213,7 @@ where
.map_err(fail)
}

pub fn new(
pub(crate) fn new(
base_kms: BaseKmsStruct,
crypto_storage: CryptoMaterialStorage<PubS, PrivS>,
security_module: Option<Arc<SecurityModuleProxy>>,
Expand Down Expand Up @@ -328,8 +327,7 @@ where
}
}

#[tonic::async_trait]
impl<PubS, PrivS> BackupOperator for RealBackupOperator<PubS, PrivS>
impl<PubS, PrivS> RealBackupOperator<PubS, PrivS>
where
PubS: Storage + Sync + Send + 'static,
PrivS: StorageExt + Sync + Send + 'static,
Expand All @@ -339,7 +337,7 @@ where
/// A digest is attested rather than the key because the composite key exceeds the attestation
/// document's [`crate::cryptography::attestation::NSM_ATTESTATION_FIELD_MAX_BYTES`] `public_key`
/// field.
async fn get_operator_public_key(
pub(crate) async fn get_operator_public_key(
Comment thread
dvdplm marked this conversation as resolved.
&self,
_request: Request<Empty>,
) -> Result<Response<OperatorPublicKey>, MetricedError> {
Expand Down Expand Up @@ -395,7 +393,7 @@ where
}

/// Restores the most recent custodian based backup.
async fn custodian_recovery_init(
pub(crate) async fn custodian_recovery_init(
&self,
request: Request<CustodianRecoveryInitRequest>,
) -> Result<Response<RecoveryRequest>, MetricedError> {
Expand Down Expand Up @@ -493,7 +491,7 @@ where
///
/// Observe that the decryption key is NOT persisted on disc and in fact removed immediately after a call to `restore_from_backup`
/// in order to minimize the possibility of leakage.
async fn custodian_backup_recovery(
pub(crate) async fn custodian_backup_recovery(
&self,
request: Request<CustodianRecoveryRequest>,
) -> Result<Response<Empty>, MetricedError> {
Expand Down Expand Up @@ -693,7 +691,7 @@ where
/// Observe that if secret sharing is used for backup (i.e. with a master key being shared with a set of custodians)
/// then [`custodian_recovery`] _must_ be called first in order to ensure that the master key is restored,
/// which is needed to allow decryption of the backup data.
async fn restore_from_backup(
pub(crate) async fn restore_from_backup(
&self,
_request: Request<Empty>,
) -> Result<Response<Empty>, MetricedError> {
Expand Down Expand Up @@ -728,7 +726,7 @@ where
}
}

async fn get_key_material_availability(
pub(crate) async fn get_key_material_availability(
&self,
_request: Request<Empty>,
) -> Result<Response<KeyMaterialAvailabilityResponse>, MetricedError> {
Expand Down
25 changes: 2 additions & 23 deletions core/service/src/engine/base.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,11 @@
pub use super::signed_payload::UserDecSignedPayload;
use super::signed_payload::user_dec_payload;
use super::traits::BaseKms;
use crate::consts::ID_LENGTH;
use crate::consts::SAFE_SER_SIZE_LIMIT;
use crate::cryptography::decompression;
use crate::cryptography::internal_crypto_types::WrappedDKGParams;
use crate::cryptography::signatures::PublicSigKey;
use crate::cryptography::signatures::internal_sign;
use crate::cryptography::signatures::{PublicSigKey, Signature};
use crate::cryptography::signing::SigningSchemeType;
use crate::cryptography::signing::identity::NodeSigningIdentity;
use crate::cryptography::signing::typed_signature::StoredTypedSignature;
Expand All @@ -19,7 +18,7 @@ use alloy_primitives::U256;
use alloy_primitives::{Address, B256, Bytes, FixedBytes, Uint};
use alloy_sol_types::Eip712Domain;
use alloy_sol_types::SolStruct;
use hashing::{DomainSep, hash_element, hash_versioned, serialize_hash_element};
use hashing::{DomainSep, hash_versioned, serialize_hash_element};
use kms_grpc::RequestId;
use kms_grpc::kms::v1::{
CiphertextFormat, FheParameter, PublicDecryptionResponsePayload, TypedPlaintext,
Expand Down Expand Up @@ -1193,26 +1192,6 @@ impl BaseKmsStruct {
}
}

impl BaseKms for BaseKmsStruct {
/// sign `msg` using the KMS' private signing key
fn sign<T>(&self, dsep: &DomainSep, msg: &T) -> anyhow::Result<Signature>
where
T: Serialize + AsRef<[u8]>,
{
match self.signing_identity.as_ref() {
None => anyhow::bail!("KMS has no signing key"),
Some(identity) => internal_sign(dsep, msg, identity.ecdsa()),
}
}

fn digest<T>(domain_separator: &DomainSep, msg: &T) -> anyhow::Result<Vec<u8>>
where
T: ?Sized + AsRef<[u8]>,
{
Ok(hash_element(domain_separator, msg))
}
}

/// ABI encodes a list of typed plaintexts into a single byte vector for Ethereum compatibility.
/// This follows the encoding pattern used in the JavaScript version for decrypted results and also supports `ebytes`.
/// This function is NOT compatible with fhevm v0.9.0 and is only intended for future use with fhevm supporting `ebytes`.
Expand Down
Loading
Loading