Skip to content
Merged
54 changes: 54 additions & 0 deletions crates/e2e-tests/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2236,6 +2236,26 @@ mod tests {
proxy_allowlist().await?.members,
tls_keys(&test_networks.hashi_network().nodes()[..INITIAL_NODES])?
);
// The proxy forwards a committee handoff only once the chain stores
// it, which is when the rotation completes.
use hashi_guardian_proxy::node::handoffs::HandoffGate;
use hashi_guardian_proxy::node::members::ChainSource;
use hashi_types::proto;
let guardian = test_networks
.guardian_harness
.as_ref()
.context("no guardian harness")?
.endpoint()
.to_string();
let handoff_metrics =
std::sync::Arc::new(hashi_guardian_proxy::metrics::ProxyMetrics::new());
let handoff_gate = HandoffGate::new(
ChainSource::new(
tonic::transport::Endpoint::from_shared(guardian)?.connect_lazy(),
&test_networks.sui_network.rpc_url,
)?,
handoff_metrics.clone(),
);

// Force epoch change → key rotation 19→20.
test_networks.sui_network.force_close_epoch().await?;
Expand All @@ -2259,6 +2279,28 @@ mod tests {
proxy_allowlist().await?.members,
tls_keys(test_networks.hashi_network().nodes())?
);
// A handoff out of the current epoch is not stored while the rotation
// is pending, so the proxy refuses it.
let early = proto::SignedCommitteeTransition {
data: Some(proto::CommitteeTransition {
new_committee: Some(proto::Committee {
epoch: Some(initial_epoch + 1),
..Default::default()
}),
}),
committee_signature: Some(proto::CommitteeSignature {
epoch: Some(initial_epoch),
..Default::default()
}),
};
handoff_gate.admit(&[early]).await.unwrap_err();
assert_eq!(
handoff_metrics
.handoff_refused
.with_label_values(&["not_on_chain"])
.get(),
1
);

wait_for_rotation(test_networks.hashi_network().nodes(), initial_epoch + 1).await;
assert_nodes_agree_on_mpc_key(test_networks.hashi_network().nodes()).await;
Expand All @@ -2267,6 +2309,18 @@ mod tests {
proxy_allowlist().await?.members,
tls_keys(test_networks.hashi_network().nodes())?
);
// The handoff a node pushes for the completed rotation is admitted.
let pushed = test_networks.hashi_network().nodes()[0]
.hashi()
.onchain_state()
.committee_handoff(initial_epoch)
.context("node 0 holds no handoff out of the initial epoch")?;
let pushed =
hashi_types::guardian::proto_conversions::signed_committee_transition_to_pb(&pushed);
handoff_gate
.admit(&[pushed])
.await
.map_err(|status| anyhow::anyhow!("the proxy refused a node's handoff: {status}"))?;
crate::test_helpers::assert_no_member_refusals(&test_networks);

Ok(())
Expand Down
5 changes: 3 additions & 2 deletions crates/hashi-guardian-proxy/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,8 +53,9 @@ pub struct Config {
/// bitcoin|testnet|signet|regtest). Must match the guardian's config; used
/// to recompute sighashes when verifying a log replay.
pub btc_network: Network,
/// Sui fullnode gRPC endpoint the committee member allowlist is read from
/// (`SUI_RPC_URL`, required), on the chain of the guardian's Hashi object.
/// Sui fullnode gRPC endpoint the committee member allowlist and stored
/// committee handoffs are read from (`SUI_RPC_URL`, required), on the chain
/// of the guardian's Hashi object.
pub sui_rpc_url: String,
/// Push metrics to a Prometheus remote-write endpoint; `None` leaves them
/// on `/metrics`, which nothing can scrape (`MIMIR_URL`, `MIMIR_USERNAME`
Expand Down
72 changes: 64 additions & 8 deletions crates/hashi-guardian-proxy/src/forward.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,11 @@
//! internet-facing and `OperatorInit` is one-shot and unauthenticated, so
//! exposing it would let anyone wedge the guardian. KP-signed RPCs are
//! forwarded after a signature and roster check; `ConfirmCeremony` goes to the
//! ceremony guardian, which is the relay's backend. Wrapped by
//! [`crate::node::cache::CachingGuardianGrpc`] to cache `StandardWithdrawal` and
//! `GetGuardianInfo`.
//! ceremony guardian, which is the relay's backend. A committee handoff is
//! forwarded only once the chain stores one between the same two epochs
//! ([`crate::node::handoffs`]). Wrapped by
//! [`crate::node::cache::CachingGuardianGrpc`] to cache `StandardWithdrawal`
//! and `GetGuardianInfo`.

use std::sync::Arc;

Expand All @@ -25,6 +27,7 @@ use tonic::Status;
use crate::kp;
use crate::kp::roster::RosterCache;
use crate::log_store::LogStore;
use crate::node::handoffs::HandoffGate;

/// Holds a plain [`Channel`] rather than the node's boxed transport: the generated
/// server trait requires `Send + Sync + 'static`, and `BoxCloneService` is not `Sync`.
Expand All @@ -37,14 +40,21 @@ pub struct Forwarding<L> {
/// Shared with the relay: one gate admits every KP-signed RPC, and a cert
/// rotation drops the cached roster for both.
roster: Arc<RosterCache<L>>,
handoffs: Arc<HandoffGate>,
}

impl<L: LogStore> Forwarding<L> {
pub fn new(channel: Channel, ceremony_channel: Channel, roster: Arc<RosterCache<L>>) -> Self {
pub fn new(
channel: Channel,
ceremony_channel: Channel,
roster: Arc<RosterCache<L>>,
handoffs: Arc<HandoffGate>,
) -> Self {
Self {
client: GuardianServiceClient::new(channel),
ceremony_client: GuardianServiceClient::new(ceremony_channel),
roster,
handoffs,
}
}
}
Expand Down Expand Up @@ -88,13 +98,17 @@ impl<L: LogStore> GuardianService for Forwarding<L> {
&self,
request: Request<proto::SignedCommitteeTransition>,
) -> Result<Response<proto::UpdateCommitteeResponse>, Status> {
self.handoffs
.admit(std::slice::from_ref(request.get_ref()))
.await?;
self.client.clone().update_committee(request).await
}

async fn update_committee_chain(
&self,
request: Request<proto::UpdateCommitteeChainRequest>,
) -> Result<Response<proto::UpdateCommitteeResponse>, Status> {
self.handoffs.admit(&request.get_ref().transitions).await?;
self.client.clone().update_committee_chain(request).await
}

Expand Down Expand Up @@ -179,6 +193,7 @@ pub(crate) mod test_utils {
pub(crate) get_guardian_info_calls: Arc<AtomicUsize>,
pub(crate) get_attested_guardian_info_calls: Arc<AtomicUsize>,
pub(crate) confirm_ceremony_calls: Arc<AtomicUsize>,
pub(crate) update_committee_calls: Arc<AtomicUsize>,
/// Served by `GetGuardianInfo`; the default response when unset.
pub(crate) info: Arc<std::sync::Mutex<Option<proto::GetGuardianInfoResponse>>>,
}
Expand Down Expand Up @@ -271,13 +286,15 @@ pub(crate) mod test_utils {
&self,
_: Request<proto::SignedCommitteeTransition>,
) -> Result<Response<proto::UpdateCommitteeResponse>, Status> {
unimplemented!("not exercised by tests")
self.update_committee_calls.fetch_add(1, Ordering::SeqCst);
Ok(Response::new(proto::UpdateCommitteeResponse::default()))
}
async fn update_committee_chain(
&self,
_: Request<proto::UpdateCommitteeChainRequest>,
) -> Result<Response<proto::UpdateCommitteeResponse>, Status> {
unimplemented!("not exercised by tests")
self.update_committee_calls.fetch_add(1, Ordering::SeqCst);
Ok(Response::new(proto::UpdateCommitteeResponse::default()))
}
async fn rotate_kp_set(
&self,
Expand Down Expand Up @@ -329,6 +346,8 @@ mod tests {
use super::test_utils::*;
use super::*;
use crate::node::cache::CachingGuardianGrpc;
use crate::node::handoffs::test_utils::gate_over;
use crate::node::handoffs::test_utils::transition;
use crate::node::widlog::test_utils::withdrawal_log_json;
use crate::node::widlog::WidLogIndex;
use hashi_types::guardian::now_timestamp_ms;
Expand All @@ -339,14 +358,21 @@ mod tests {

type StubStore = crate::log_store::test_store::MemStore;

/// A proxy whose chain stores `handoffs`, each as `(from_epoch, next_epoch)`.
async fn proxy_over(
active: tonic::transport::Channel,
ceremony: tonic::transport::Channel,
store: StubStore,
handoffs: &[(u64, u64)],
) -> CachingGuardianGrpc<Forwarding<StubStore>, StubStore> {
let metrics = Arc::new(crate::metrics::ProxyMetrics::new());
CachingGuardianGrpc::new(
Forwarding::new(active, ceremony, Arc::new(RosterCache::new(store))),
Forwarding::new(
active,
ceremony,
Arc::new(RosterCache::new(store)),
gate_over(handoffs),
),
WidLogIndex::ready_for_tests(StubStore::default(), metrics.clone()).await,
metrics,
)
Expand All @@ -360,7 +386,7 @@ mod tests {
CachingGuardianGrpc<Forwarding<StubStore>, StubStore>,
) {
let (stub, channel) = spawn_stub().await;
(stub, proxy_over(channel.clone(), channel, store).await)
(stub, proxy_over(channel.clone(), channel, store, &[]).await)
}

#[tokio::test]
Expand Down Expand Up @@ -439,6 +465,36 @@ mod tests {
assert_eq!(stub.confirm_ceremony_calls.load(Ordering::SeqCst), 0);
}

#[tokio::test]
async fn forwards_only_committee_handoffs_the_chain_stores() {
let (stub, channel) = spawn_stub().await;
let proxy = proxy_over(channel.clone(), channel, StubStore::default(), &[(5, 7)]).await;
let chain = |transitions| Request::new(proto::UpdateCommitteeChainRequest { transitions });

proxy
.update_committee(Request::new(transition(5, 7)))
.await
.unwrap();
proxy
.update_committee_chain(chain(vec![transition(5, 7)]))
.await
.unwrap();
assert_eq!(stub.update_committee_calls.load(Ordering::SeqCst), 2);

// The reconfig out of epoch 7 has not completed.
let early = proxy
.update_committee(Request::new(transition(7, 9)))
.await
.unwrap_err();
assert_eq!(early.code(), tonic::Code::Unavailable);
let early = proxy
.update_committee_chain(chain(vec![transition(5, 7), transition(7, 9)]))
.await
.unwrap_err();
assert_eq!(early.code(), tonic::Code::Unavailable);
assert_eq!(stub.update_committee_calls.load(Ordering::SeqCst), 2);
}

// The stub `unimplemented!()`s the rejected RPCs, so a forwarded call would panic
// the server rather than return `PERMISSION_DENIED` — proof the proxy short-circuits.
#[tokio::test]
Expand Down
12 changes: 10 additions & 2 deletions crates/hashi-guardian-proxy/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,9 @@
//! log ([`node::widlog`]), which the proxy only reads.
//! [`node::member_auth`] gates every route: node RPCs are served only to
//! current or pending committee members ([`node::members`]), who present
//! their registered TLS key as a client certificate.
//! their registered TLS key as a client certificate. [`node::handoffs`]
//! forwards a committee handoff only once the chain stores one between the
//! same two epochs.
//! - [`kp`]: [`kp::relay`] serves `GuardianRelayService`: key provisioners
//! submit one share each — authenticated against the ceremony's committed
//! roster read from the S3 share log ([`kp::roster`]) — and the relay batches
Expand Down Expand Up @@ -121,6 +123,7 @@ mod tests {
use crate::kp::roster::RosterCache;
use crate::log_store::test_store::MemStore;
use crate::metrics::ProxyMetrics;
use crate::node::handoffs::test_utils::gate_over;
use crate::node::members::test_utils::snapshot;
use crate::node::members::MemberAllowlist;
use crate::node::widlog::WidLogIndex;
Expand Down Expand Up @@ -201,7 +204,12 @@ mod tests {
let metrics = Arc::new(ProxyMetrics::new());
let roster = Arc::new(RosterCache::new(MemStore::default()));
let guardian = CachingGuardianGrpc::new(
Forwarding::new(backend.clone(), backend.clone(), roster.clone()),
Forwarding::new(
backend.clone(),
backend.clone(),
roster.clone(),
gate_over(&[]),
),
WidLogIndex::ready_for_tests(MemStore::default(), metrics.clone()).await,
metrics.clone(),
);
Expand Down
20 changes: 14 additions & 6 deletions crates/hashi-guardian-proxy/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,9 @@ use hashi_guardian_proxy::kp::roster::RosterCache;
use hashi_guardian_proxy::log_store::S3LogStore;
use hashi_guardian_proxy::metrics::ProxyMetrics;
use hashi_guardian_proxy::node::cache::CachingGuardianGrpc;
use hashi_guardian_proxy::node::handoffs::HandoffGate;
use hashi_guardian_proxy::node::member_auth::MemberGate;
use hashi_guardian_proxy::node::members::ChainMemberSource;
use hashi_guardian_proxy::node::members::ChainSource;
use hashi_guardian_proxy::node::members::MemberAllowlist;
use hashi_guardian_proxy::node::widlog::WidLogIndex;
use hashi_guardian_proxy::public::info;
Expand Down Expand Up @@ -85,11 +86,18 @@ async fn main() -> Result<()> {
// The member gate's allowlist follows the committee of the Hashi object the
// active guardian serves.
let allowlist = Arc::new(MemberAllowlist::new(metrics.clone()));
tokio::spawn(allowlist.clone().refresh_forever(ChainMemberSource::new(
channel.clone(),
&config.sui_rpc_url,
)?));
tokio::spawn(
allowlist
.clone()
.refresh_forever(ChainSource::new(channel.clone(), &config.sui_rpc_url)?),
);
let gate = Arc::new(MemberGate::new(allowlist, metrics.clone()));
// Its own Sui connection: lookups are request-driven, and on a shared one
// they could crowd out the allowlist refresh.
let handoffs = Arc::new(HandoffGate::new(
ChainSource::new(channel.clone(), &config.sui_rpc_url)?,
metrics.clone(),
));

// One roster cache, shared: the relay authorizes submissions against it and
// a cert rotation through the forwarder invalidates it.
Expand All @@ -115,7 +123,7 @@ async fn main() -> Result<()> {
// KPs confirm a ceremony to the guardian they are provisioning, so that
// RPC follows the relay's backend.
let guardian_svc = CachingGuardianGrpc::new(
Forwarding::new(channel, relay_channel, roster),
Forwarding::new(channel, relay_channel, roster, handoffs),
widlog,
metrics.clone(),
);
Expand Down
14 changes: 14 additions & 0 deletions crates/hashi-guardian-proxy/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ pub struct ProxyMetrics {
pub member_snapshot_timestamp_seconds: IntGauge,
/// Failed allowlist reads.
pub member_refresh_failures: IntCounter,
/// Committee updates the handoff gate refused, by reason.
pub handoff_refused: IntCounterVec,
}

impl ProxyMetrics {
Expand Down Expand Up @@ -104,6 +106,14 @@ impl ProxyMetrics {
"Failed committee member allowlist reads",
)
.expect("valid metric");
let handoff_refused = IntCounterVec::new(
Opts::new(
"guardian_proxy_handoff_refused_total",
"Committee updates refused by the handoff gate, by reason",
),
&["reason"],
)
.expect("valid metric");

registry
.register(Box::new(requests.clone()))
Expand Down Expand Up @@ -135,6 +145,9 @@ impl ProxyMetrics {
registry
.register(Box::new(member_refresh_failures.clone()))
.expect("register");
registry
.register(Box::new(handoff_refused.clone()))
.expect("register");

Self {
registry,
Expand All @@ -148,6 +161,7 @@ impl ProxyMetrics {
member_allowlist_size,
member_snapshot_timestamp_seconds,
member_refresh_failures,
handoff_refused,
}
}

Expand Down
Loading
Loading