mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-09-03 01:25:37 +00:00
feat: support allocating public IPv6 addresses from a provider (#2162)
* feat: support allocating public IPv6 addresses from a provider Add a provider/leaser architecture for public IPv6 address allocation between nodes in the same network: - A node with `--ipv6-public-addr-provider` advertises a delegable public IPv6 prefix (auto-detected from kernel routes or manually configured via `--ipv6-public-addr-prefix`). - Other nodes with `--ipv6-public-addr-auto` request a /128 lease from the selected provider via a new RPC service (PublicIpv6AddrRpc). - Leases have a 30s TTL, renewed every 10s by the client routine. - The provider allocates addresses deterministically from its prefix using instance-UUID-based hashing to prefer stable assignments. - Routes to peer leases are installed on the TUN device, and each client's own /128 is assigned as its IPv6 address. Also includes netlink IPv6 route table inspection, integration tests, and event-driven route/address reconciliation. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
b20075e3dc
commit
8f862997eb
@@ -292,13 +292,33 @@ impl AclFilter {
|
||||
processor.increment_stat(AclStatKey::PacketsTotal);
|
||||
}
|
||||
|
||||
fn classify_chain_type(
|
||||
is_in: bool,
|
||||
packet_info: &PacketInfo,
|
||||
my_ipv4: Option<Ipv4Addr>,
|
||||
is_local_ipv6: impl Fn(Ipv6Addr) -> bool,
|
||||
) -> ChainType {
|
||||
if !is_in {
|
||||
return ChainType::Outbound;
|
||||
}
|
||||
|
||||
let is_local_dst = packet_info.dst_ip == my_ipv4.unwrap_or(Ipv4Addr::UNSPECIFIED)
|
||||
|| matches!(packet_info.dst_ip, IpAddr::V6(dst) if is_local_ipv6(dst));
|
||||
|
||||
if is_local_dst {
|
||||
ChainType::Inbound
|
||||
} else {
|
||||
ChainType::Forward
|
||||
}
|
||||
}
|
||||
|
||||
/// Common ACL processing logic
|
||||
pub fn process_packet_with_acl(
|
||||
&self,
|
||||
packet: &ZCPacket,
|
||||
is_in: bool,
|
||||
my_ipv4: Option<Ipv4Addr>,
|
||||
my_ipv6: Option<Ipv6Addr>,
|
||||
is_local_ipv6: impl Fn(Ipv6Addr) -> bool,
|
||||
route: &(dyn super::route_trait::Route + Send + Sync + 'static),
|
||||
) -> bool {
|
||||
if !self.acl_enabled.load(Ordering::Relaxed) {
|
||||
@@ -323,17 +343,7 @@ impl AclFilter {
|
||||
}
|
||||
};
|
||||
|
||||
let chain_type = if is_in {
|
||||
if packet_info.dst_ip == my_ipv4.unwrap_or(Ipv4Addr::UNSPECIFIED)
|
||||
|| packet_info.dst_ip == my_ipv6.unwrap_or(Ipv6Addr::UNSPECIFIED)
|
||||
{
|
||||
ChainType::Inbound
|
||||
} else {
|
||||
ChainType::Forward
|
||||
}
|
||||
} else {
|
||||
ChainType::Outbound
|
||||
};
|
||||
let chain_type = Self::classify_chain_type(is_in, &packet_info, my_ipv4, is_local_ipv6);
|
||||
|
||||
// Get current processor atomically
|
||||
let processor = self.get_processor();
|
||||
@@ -384,3 +394,55 @@ impl AclFilter {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::{
|
||||
net::{IpAddr, Ipv4Addr, Ipv6Addr},
|
||||
sync::Arc,
|
||||
};
|
||||
|
||||
use crate::{
|
||||
common::acl_processor::PacketInfo,
|
||||
proto::acl::{ChainType, Protocol},
|
||||
};
|
||||
|
||||
use super::AclFilter;
|
||||
|
||||
fn packet_info(dst_ip: IpAddr) -> PacketInfo {
|
||||
PacketInfo {
|
||||
src_ip: IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)),
|
||||
dst_ip,
|
||||
src_port: Some(1234),
|
||||
dst_port: Some(80),
|
||||
protocol: Protocol::Tcp,
|
||||
packet_size: 64,
|
||||
src_groups: Arc::new(Vec::new()),
|
||||
dst_groups: Arc::new(Vec::new()),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_chain_type_treats_public_ipv6_lease_as_inbound() {
|
||||
let leased_ipv6 = Ipv6Addr::new(0x2001, 0xdb8, 0x100, 0, 0, 0, 0, 0x123);
|
||||
let packet_info = packet_info(IpAddr::V6(leased_ipv6));
|
||||
|
||||
let chain =
|
||||
AclFilter::classify_chain_type(true, &packet_info, None, |ip| ip == leased_ipv6);
|
||||
|
||||
assert_eq!(chain, ChainType::Inbound);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_chain_type_keeps_non_local_ipv6_as_forward() {
|
||||
let leased_ipv6 = Ipv6Addr::new(0x2001, 0xdb8, 0x100, 0, 0, 0, 0, 0x123);
|
||||
let packet_info = packet_info(IpAddr::V6(Ipv6Addr::new(
|
||||
0x2001, 0xdb8, 0xffff, 2, 0, 0, 0, 0x100,
|
||||
)));
|
||||
|
||||
let chain =
|
||||
AclFilter::classify_chain_type(true, &packet_info, None, |ip| ip == leased_ipv6);
|
||||
|
||||
assert_eq!(chain, ChainType::Forward);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ pub mod peer_ospf_route;
|
||||
pub mod peer_rpc;
|
||||
pub mod peer_rpc_service;
|
||||
pub mod peer_session;
|
||||
pub(crate) mod public_ipv6;
|
||||
pub mod relay_peer_map;
|
||||
pub mod route_trait;
|
||||
pub mod rpc_service;
|
||||
|
||||
@@ -1062,7 +1062,7 @@ impl PeerManager {
|
||||
&ret,
|
||||
true,
|
||||
global_ctx.get_ipv4().map(|x| x.address()),
|
||||
global_ctx.get_ipv6().map(|x| x.address()),
|
||||
|dst| global_ctx.is_ip_local_ipv6(&dst),
|
||||
&route,
|
||||
) {
|
||||
continue;
|
||||
@@ -1291,6 +1291,18 @@ impl PeerManager {
|
||||
self.get_route().list_proxy_cidrs_v6().await
|
||||
}
|
||||
|
||||
pub async fn list_public_ipv6_routes(&self) -> BTreeSet<cidr::Ipv6Inet> {
|
||||
self.get_route().list_public_ipv6_routes().await
|
||||
}
|
||||
|
||||
pub async fn get_my_public_ipv6_addr(&self) -> Option<cidr::Ipv6Inet> {
|
||||
self.get_route().get_my_public_ipv6_addr().await
|
||||
}
|
||||
|
||||
pub async fn get_local_public_ipv6_info(&self) -> instance::ListPublicIpv6InfoResponse {
|
||||
self.get_route().get_local_public_ipv6_info().await
|
||||
}
|
||||
|
||||
pub async fn dump_route(&self) -> String {
|
||||
self.get_route().dump().await
|
||||
}
|
||||
@@ -1330,7 +1342,7 @@ impl PeerManager {
|
||||
data,
|
||||
false,
|
||||
None,
|
||||
None,
|
||||
|_| false,
|
||||
&self.get_route(),
|
||||
) {
|
||||
return false;
|
||||
@@ -1532,6 +1544,10 @@ impl PeerManager {
|
||||
dst_peers.extend(self.peers.list_routes().await.iter().map(|x| *x.key()));
|
||||
} else if let Some(peer_id) = self.peers.get_peer_id_by_ipv6(ipv6_addr).await {
|
||||
dst_peers.push(peer_id);
|
||||
} else if !ipv6_addr.is_unicast_link_local()
|
||||
&& let Some(peer_id) = self.get_route().get_public_ipv6_gateway_peer_id().await
|
||||
{
|
||||
dst_peers.push(peer_id);
|
||||
} else if !ipv6_addr.is_unicast_link_local() {
|
||||
// NOTE: never route link local address to exit node.
|
||||
for exit_node in self.exit_nodes.read().await.iter() {
|
||||
@@ -1662,7 +1678,7 @@ impl PeerManager {
|
||||
&& !self.global_ctx.is_ip_local_virtual_ip(&ip_addr)
|
||||
{
|
||||
// Keep the loop-prevention flags for proxy-induced self-delivery where
|
||||
// the destination is not this node's own virtual IP.
|
||||
// the destination is not this node's own EasyTier-managed IP.
|
||||
hdr.set_not_send_to_tun(true);
|
||||
hdr.set_no_proxy(true);
|
||||
}
|
||||
@@ -1879,6 +1895,15 @@ impl PeerManager {
|
||||
version: EASYTIER_VERSION.to_string(),
|
||||
feature_flag: Some(self.global_ctx.get_feature_flags()),
|
||||
ip_list: Some(self.global_ctx.get_ip_collector().collect_ip_addrs().await),
|
||||
public_ipv6_addr: self.get_my_public_ipv6_addr().await.map(Into::into),
|
||||
ipv6_public_addr_prefix: self
|
||||
.global_ctx
|
||||
.get_advertised_ipv6_public_addr_prefix()
|
||||
.map(|prefix| {
|
||||
cidr::Ipv6Inet::new(prefix.first_address(), prefix.network_length())
|
||||
.unwrap()
|
||||
.into()
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ use std::{
|
||||
};
|
||||
|
||||
use arc_swap::ArcSwap;
|
||||
use cidr::{IpCidr, Ipv4Cidr, Ipv6Cidr};
|
||||
use cidr::{IpCidr, Ipv4Cidr, Ipv6Cidr, Ipv6Inet};
|
||||
use crossbeam::atomic::AtomicCell;
|
||||
use dashmap::DashMap;
|
||||
use ordered_hash_map::OrderedHashMap;
|
||||
@@ -46,9 +46,10 @@ use crate::{
|
||||
peer_rpc::{
|
||||
ForeignNetworkRouteInfoEntry, ForeignNetworkRouteInfoKey, OspfRouteRpc,
|
||||
OspfRouteRpcClientFactory, OspfRouteRpcServer, PeerGroupInfo, PeerIdVersion,
|
||||
PeerIdentityType, RouteForeignNetworkInfos, RouteForeignNetworkSummary, RoutePeerInfo,
|
||||
RoutePeerInfos, SyncRouteInfoError, SyncRouteInfoRequest, SyncRouteInfoResponse,
|
||||
TrustedCredentialPubkey, TrustedCredentialPubkeyProof, route_foreign_network_infos,
|
||||
PeerIdentityType, PublicIpv6AddrRpcServer, RouteForeignNetworkInfos,
|
||||
RouteForeignNetworkSummary, RoutePeerInfo, RoutePeerInfos, SyncRouteInfoError,
|
||||
SyncRouteInfoRequest, SyncRouteInfoResponse, TrustedCredentialPubkey,
|
||||
TrustedCredentialPubkeyProof, route_foreign_network_infos,
|
||||
route_foreign_network_summary, sync_route_info_request::ConnInfo,
|
||||
},
|
||||
rpc_types::{
|
||||
@@ -63,6 +64,9 @@ use super::{
|
||||
PeerPacketFilter,
|
||||
graph_algo::dijkstra_with_first_hop,
|
||||
peer_rpc::PeerRpcManager,
|
||||
public_ipv6::{
|
||||
PublicIpv6PeerRouteInfo, PublicIpv6RouteControl, PublicIpv6Service, PublicIpv6SyncTrigger,
|
||||
},
|
||||
route_trait::{
|
||||
DefaultRouteCostCalculator, ForeignNetworkRouteInfoMap, NextHopPolicy, RouteCostCalculator,
|
||||
RouteCostCalculatorInterface,
|
||||
@@ -137,6 +141,10 @@ fn raw_credential_bytes_from_route_info(
|
||||
.map(|credential| credential.encode_to_vec())
|
||||
}
|
||||
|
||||
fn route_peer_inst_id(info: &RoutePeerInfo) -> Option<uuid::Uuid> {
|
||||
info.inst_id.map(Into::into)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
struct AtomicVersion(Arc<AtomicU32>);
|
||||
|
||||
@@ -205,6 +213,8 @@ impl RoutePeerInfo {
|
||||
quic_port: None,
|
||||
noise_static_pubkey: Vec::new(),
|
||||
trusted_credential_pubkeys: Vec::new(),
|
||||
ipv6_public_addr_prefix: None,
|
||||
ipv6_public_addr_lease: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -221,6 +231,7 @@ impl RoutePeerInfo {
|
||||
my_peer_id: PeerId,
|
||||
peer_route_id: u64,
|
||||
global_ctx: &ArcGlobalCtx,
|
||||
public_ipv6_addr_lease: Option<Ipv6Inet>,
|
||||
) -> Self {
|
||||
let stun_info = global_ctx.get_stun_info_collector().get_stun_info();
|
||||
let noise_static_pubkey = global_ctx
|
||||
@@ -259,6 +270,14 @@ impl RoutePeerInfo {
|
||||
.unwrap_or(24),
|
||||
|
||||
ipv6_addr: global_ctx.get_ipv6().map(|x| x.into()),
|
||||
ipv6_public_addr_prefix: global_ctx.get_advertised_ipv6_public_addr_prefix().map(
|
||||
|prefix| {
|
||||
Ipv6Inet::new(prefix.first_address(), prefix.network_length())
|
||||
.unwrap()
|
||||
.into()
|
||||
},
|
||||
),
|
||||
ipv6_public_addr_lease: public_ipv6_addr_lease.map(Into::into),
|
||||
|
||||
groups: global_ctx.get_acl_groups(my_peer_id),
|
||||
|
||||
@@ -349,6 +368,8 @@ impl From<RoutePeerInfo> for crate::proto::api::instance::Route {
|
||||
path_latency_latency_first: None,
|
||||
|
||||
ipv6_addr: val.ipv6_addr,
|
||||
public_ipv6_addr: val.ipv6_public_addr_lease,
|
||||
ipv6_public_addr_prefix: val.ipv6_public_addr_prefix,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -964,8 +985,14 @@ impl SyncedRouteInfo {
|
||||
my_peer_id: PeerId,
|
||||
my_peer_route_id: u64,
|
||||
global_ctx: &ArcGlobalCtx,
|
||||
public_ipv6_addr_lease: Option<Ipv6Inet>,
|
||||
) -> bool {
|
||||
let mut new = RoutePeerInfo::new_updated_self(my_peer_id, my_peer_route_id, global_ctx);
|
||||
let mut new = RoutePeerInfo::new_updated_self(
|
||||
my_peer_id,
|
||||
my_peer_route_id,
|
||||
global_ctx,
|
||||
public_ipv6_addr_lease,
|
||||
);
|
||||
let mut guard = self.peer_infos.upgradable_read();
|
||||
let old = guard.get(&my_peer_id);
|
||||
let new_version = old.map(|x| x.version).unwrap_or(0) + 1;
|
||||
@@ -1588,6 +1615,21 @@ impl RouteTable {
|
||||
.or_insert(peer_id_and_version);
|
||||
}
|
||||
|
||||
if let Some(ipv6_addr) = info
|
||||
.ipv6_public_addr_lease
|
||||
.as_ref()
|
||||
.and_then(|addr| addr.address)
|
||||
{
|
||||
self.ipv6_peer_id_map
|
||||
.entry(ipv6_addr.into())
|
||||
.and_modify(|v| {
|
||||
if is_new_peer_better(v) {
|
||||
*v = peer_id_and_version;
|
||||
}
|
||||
})
|
||||
.or_insert(peer_id_and_version);
|
||||
}
|
||||
|
||||
for cidr in info.proxy_cidrs.iter() {
|
||||
let Ok(cidr) = cidr.parse::<IpCidr>() else {
|
||||
tracing::warn!("invalid proxy cidr: {:?}, from peer: {:?}", cidr, peer_id);
|
||||
@@ -2019,6 +2061,8 @@ struct PeerRouteServiceImpl {
|
||||
foreign_network_owner_map: DashMap<NetworkIdentity, Vec<PeerId>>,
|
||||
foreign_network_my_peer_id_map: DashMap<(String, PeerId), PeerId>,
|
||||
synced_route_info: SyncedRouteInfo,
|
||||
public_ipv6_service: std::sync::Mutex<Weak<PublicIpv6Service>>,
|
||||
self_public_ipv6_addr_lease: std::sync::Mutex<Option<Ipv6Inet>>,
|
||||
cached_local_conn_map: std::sync::Mutex<RouteConnBitmap>,
|
||||
cached_local_conn_map_version: AtomicVersion,
|
||||
cached_interface_peer_snapshot: std::sync::Mutex<Arc<InterfacePeerSnapshot>>,
|
||||
@@ -2081,6 +2125,8 @@ impl PeerRouteServiceImpl {
|
||||
non_reusable_credential_owners: DashMap::new(),
|
||||
version: AtomicVersion::new(),
|
||||
},
|
||||
public_ipv6_service: std::sync::Mutex::new(Weak::new()),
|
||||
self_public_ipv6_addr_lease: std::sync::Mutex::new(None),
|
||||
cached_local_conn_map: std::sync::Mutex::new(RouteConnBitmap::default()),
|
||||
cached_local_conn_map_version: AtomicVersion::new(),
|
||||
cached_interface_peer_snapshot: std::sync::Mutex::new(Arc::new(
|
||||
@@ -2119,6 +2165,20 @@ impl PeerRouteServiceImpl {
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn set_public_ipv6_service(&self, service: Weak<PublicIpv6Service>) {
|
||||
*self.public_ipv6_service.lock().unwrap() = service;
|
||||
}
|
||||
|
||||
fn public_ipv6_service(&self) -> Option<Arc<PublicIpv6Service>> {
|
||||
self.public_ipv6_service.lock().unwrap().upgrade()
|
||||
}
|
||||
|
||||
fn notify_public_ipv6_route_change(&self) -> bool {
|
||||
self.public_ipv6_service()
|
||||
.map(|service| service.handle_route_change())
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn get_or_create_session(&self, dst_peer_id: PeerId) -> Arc<SyncRouteSession> {
|
||||
self.sessions
|
||||
.entry(dst_peer_id)
|
||||
@@ -2230,6 +2290,7 @@ impl PeerRouteServiceImpl {
|
||||
self.my_peer_id,
|
||||
self.my_peer_route_id,
|
||||
&self.global_ctx,
|
||||
*self.self_public_ipv6_addr_lease.lock().unwrap(),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -2618,14 +2679,19 @@ impl PeerRouteServiceImpl {
|
||||
untrusted_changed = self.refresh_credential_trusts_and_disconnect().await;
|
||||
}
|
||||
|
||||
let mut public_ipv6_state_updated = false;
|
||||
if my_peer_info_updated || my_conn_info_updated || untrusted_changed {
|
||||
self.update_route_table_and_cached_local_conn_bitmap();
|
||||
self.update_foreign_network_owner_map();
|
||||
public_ipv6_state_updated = self.notify_public_ipv6_route_change();
|
||||
}
|
||||
if my_peer_info_updated {
|
||||
self.update_peer_info_last_update();
|
||||
}
|
||||
my_peer_info_updated || my_conn_info_updated || my_foreign_network_updated
|
||||
my_peer_info_updated
|
||||
|| my_conn_info_updated
|
||||
|| my_foreign_network_updated
|
||||
|| public_ipv6_state_updated
|
||||
}
|
||||
|
||||
async fn refresh_acl_groups(&self) -> bool {
|
||||
@@ -2652,15 +2718,17 @@ impl PeerRouteServiceImpl {
|
||||
let untrusted = self.refresh_credential_trusts_with_current_topology();
|
||||
self.disconnect_untrusted_peers(&untrusted).await;
|
||||
|
||||
let mut public_ipv6_state_updated = false;
|
||||
if my_peer_info_updated || !untrusted.is_empty() {
|
||||
self.update_route_table_and_cached_local_conn_bitmap();
|
||||
self.update_foreign_network_owner_map();
|
||||
public_ipv6_state_updated = self.notify_public_ipv6_route_change();
|
||||
}
|
||||
if my_peer_info_updated {
|
||||
self.update_peer_info_last_update();
|
||||
}
|
||||
|
||||
my_peer_info_updated || !untrusted.is_empty()
|
||||
my_peer_info_updated || !untrusted.is_empty() || public_ipv6_state_updated
|
||||
}
|
||||
|
||||
fn refresh_credential_trusts(&self) -> Vec<PeerId> {
|
||||
@@ -2968,7 +3036,6 @@ impl PeerRouteServiceImpl {
|
||||
session
|
||||
.update_dst_saved_foreign_network_version(foreign_network, dst_peer_id);
|
||||
}
|
||||
|
||||
session.update_last_sync_succ_timestamp(next_last_sync_succ_timestamp);
|
||||
}
|
||||
}
|
||||
@@ -3493,7 +3560,13 @@ impl RouteSessionManager {
|
||||
}
|
||||
|
||||
if need_update_route_table || foreign_network_changed {
|
||||
service_impl.update_route_table_and_cached_local_conn_bitmap();
|
||||
service_impl.update_foreign_network_owner_map();
|
||||
if need_update_route_table
|
||||
&& let Some(public_ipv6_service) = service_impl.public_ipv6_service()
|
||||
{
|
||||
public_ipv6_service.handle_route_change();
|
||||
}
|
||||
}
|
||||
|
||||
tracing::debug!(
|
||||
@@ -3534,12 +3607,86 @@ impl RouteSessionManager {
|
||||
}
|
||||
}
|
||||
|
||||
struct OspfPublicIpv6RouteHandle {
|
||||
service_impl: Weak<PeerRouteServiceImpl>,
|
||||
}
|
||||
|
||||
impl PublicIpv6RouteControl for OspfPublicIpv6RouteHandle {
|
||||
fn my_peer_id(&self) -> PeerId {
|
||||
self.service_impl
|
||||
.upgrade()
|
||||
.map(|service_impl| service_impl.my_peer_id)
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
fn peer_route_snapshot(&self) -> Vec<PublicIpv6PeerRouteInfo> {
|
||||
let Some(service_impl) = self.service_impl.upgrade() else {
|
||||
return Vec::new();
|
||||
};
|
||||
|
||||
service_impl
|
||||
.synced_route_info
|
||||
.peer_infos
|
||||
.read()
|
||||
.iter()
|
||||
.map(|(peer_id, info)| PublicIpv6PeerRouteInfo {
|
||||
peer_id: *peer_id,
|
||||
inst_id: route_peer_inst_id(info),
|
||||
is_provider: info
|
||||
.feature_flag
|
||||
.as_ref()
|
||||
.map(|flags| flags.ipv6_public_addr_provider)
|
||||
.unwrap_or(false),
|
||||
prefix: info
|
||||
.ipv6_public_addr_prefix
|
||||
.map(Into::into)
|
||||
.map(|prefix: Ipv6Inet| prefix.network()),
|
||||
lease: info.ipv6_public_addr_lease.map(Into::into),
|
||||
reachable: *peer_id == service_impl.my_peer_id
|
||||
|| service_impl.route_table.peer_reachable(*peer_id),
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn publish_self_public_ipv6_lease(&self, lease: Option<Ipv6Inet>) -> bool {
|
||||
let Some(service_impl) = self.service_impl.upgrade() else {
|
||||
return false;
|
||||
};
|
||||
|
||||
let mut current = service_impl.self_public_ipv6_addr_lease.lock().unwrap();
|
||||
if *current == lease {
|
||||
return false;
|
||||
}
|
||||
*current = lease;
|
||||
drop(current);
|
||||
|
||||
let changed = service_impl.update_my_peer_info();
|
||||
if changed {
|
||||
service_impl.update_route_table_and_cached_local_conn_bitmap();
|
||||
service_impl.update_foreign_network_owner_map();
|
||||
}
|
||||
changed
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct OspfPublicIpv6SyncTrigger {
|
||||
session_mgr: RouteSessionManager,
|
||||
}
|
||||
|
||||
impl PublicIpv6SyncTrigger for OspfPublicIpv6SyncTrigger {
|
||||
fn sync_now(&self, reason: &str) {
|
||||
self.session_mgr.sync_now(reason);
|
||||
}
|
||||
}
|
||||
|
||||
pub struct PeerRoute {
|
||||
my_peer_id: PeerId,
|
||||
global_ctx: ArcGlobalCtx,
|
||||
peer_rpc: Weak<PeerRpcManager>,
|
||||
|
||||
service_impl: Arc<PeerRouteServiceImpl>,
|
||||
public_ipv6_service: Arc<PublicIpv6Service>,
|
||||
session_mgr: RouteSessionManager,
|
||||
|
||||
tasks: std::sync::Mutex<JoinSet<()>>,
|
||||
@@ -3563,6 +3710,17 @@ impl PeerRoute {
|
||||
) -> Arc<Self> {
|
||||
let service_impl = Arc::new(PeerRouteServiceImpl::new(my_peer_id, global_ctx.clone()));
|
||||
let session_mgr = RouteSessionManager::new(service_impl.clone(), peer_rpc.clone());
|
||||
let public_ipv6_service = Arc::new(PublicIpv6Service::new(
|
||||
global_ctx.clone(),
|
||||
Arc::downgrade(&peer_rpc),
|
||||
Arc::new(OspfPublicIpv6RouteHandle {
|
||||
service_impl: Arc::downgrade(&service_impl),
|
||||
}),
|
||||
Arc::new(OspfPublicIpv6SyncTrigger {
|
||||
session_mgr: session_mgr.clone(),
|
||||
}),
|
||||
));
|
||||
service_impl.set_public_ipv6_service(Arc::downgrade(&public_ipv6_service));
|
||||
|
||||
Arc::new(PeerRoute {
|
||||
my_peer_id,
|
||||
@@ -3570,6 +3728,7 @@ impl PeerRoute {
|
||||
peer_rpc: Arc::downgrade(&peer_rpc),
|
||||
|
||||
service_impl,
|
||||
public_ipv6_service,
|
||||
session_mgr,
|
||||
|
||||
tasks: std::sync::Mutex::new(JoinSet::new()),
|
||||
@@ -3607,6 +3766,9 @@ impl PeerRoute {
|
||||
tracing::debug!("cost_calculator_need_update");
|
||||
service_impl.synced_route_info.version.inc();
|
||||
service_impl.update_route_table();
|
||||
if let Some(public_ipv6_service) = service_impl.public_ipv6_service() {
|
||||
public_ipv6_service.handle_route_change();
|
||||
}
|
||||
}
|
||||
|
||||
select! {
|
||||
@@ -3631,11 +3793,16 @@ impl PeerRoute {
|
||||
|
||||
// make sure my_peer_id is in the peer_infos.
|
||||
self.service_impl.update_my_infos().await;
|
||||
self.public_ipv6_service.handle_route_change();
|
||||
|
||||
peer_rpc.rpc_server().registry().register(
|
||||
OspfRouteRpcServer::new(self.session_mgr.clone()),
|
||||
&self.global_ctx.get_network_name(),
|
||||
);
|
||||
peer_rpc.rpc_server().registry().register(
|
||||
PublicIpv6AddrRpcServer::new(self.public_ipv6_service.rpc_server()),
|
||||
&self.global_ctx.get_network_name(),
|
||||
);
|
||||
|
||||
self.tasks
|
||||
.lock()
|
||||
@@ -3657,6 +3824,16 @@ impl PeerRoute {
|
||||
.lock()
|
||||
.unwrap()
|
||||
.spawn(Self::clear_expired_peer(self.service_impl.clone()));
|
||||
|
||||
self.tasks
|
||||
.lock()
|
||||
.unwrap()
|
||||
.spawn(self.public_ipv6_service.clone().provider_gc_routine());
|
||||
|
||||
self.tasks
|
||||
.lock()
|
||||
.unwrap()
|
||||
.spawn(self.public_ipv6_service.clone().client_routine());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3677,6 +3854,10 @@ impl Drop for PeerRoute {
|
||||
OspfRouteRpcServer::new(self.session_mgr.clone()),
|
||||
&self.global_ctx.get_network_name(),
|
||||
);
|
||||
peer_rpc.rpc_server().registry().unregister(
|
||||
PublicIpv6AddrRpcServer::new(self.public_ipv6_service.rpc_server()),
|
||||
&self.global_ctx.get_network_name(),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3765,6 +3946,51 @@ impl Route for PeerRoute {
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn list_public_ipv6_routes(&self) -> BTreeSet<Ipv6Inet> {
|
||||
self.public_ipv6_service.list_routes()
|
||||
}
|
||||
|
||||
async fn get_my_public_ipv6_addr(&self) -> Option<Ipv6Inet> {
|
||||
self.public_ipv6_service.my_addr()
|
||||
}
|
||||
|
||||
async fn get_public_ipv6_gateway_peer_id(&self) -> Option<PeerId> {
|
||||
self.public_ipv6_service.provider_peer_id_for_client()
|
||||
}
|
||||
|
||||
async fn get_local_public_ipv6_info(
|
||||
&self,
|
||||
) -> crate::proto::api::instance::ListPublicIpv6InfoResponse {
|
||||
let Some((provider, leases)) = self.public_ipv6_service.local_provider_state() else {
|
||||
return crate::proto::api::instance::ListPublicIpv6InfoResponse::default();
|
||||
};
|
||||
|
||||
crate::proto::api::instance::ListPublicIpv6InfoResponse {
|
||||
provider_prefix: Some(
|
||||
Ipv6Inet::new(
|
||||
provider.prefix.first_address(),
|
||||
provider.prefix.network_length(),
|
||||
)
|
||||
.unwrap()
|
||||
.into(),
|
||||
),
|
||||
provider_leases: leases
|
||||
.into_iter()
|
||||
.map(|lease| crate::proto::api::instance::PublicIpv6LeaseInfo {
|
||||
peer_id: lease.peer_id,
|
||||
inst_id: lease.inst_id.to_string(),
|
||||
leased_addr: Some(lease.addr.into()),
|
||||
valid_until_unix_seconds: lease
|
||||
.valid_until
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs() as i64,
|
||||
reused: lease.reused,
|
||||
})
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
|
||||
async fn get_peer_id_by_ipv4(&self, ipv4_addr: &Ipv4Addr) -> Option<PeerId> {
|
||||
let route_table = &self.service_impl.route_table;
|
||||
if let Some(p) = route_table.ipv4_peer_id_map.get(ipv4_addr) {
|
||||
@@ -5180,6 +5406,7 @@ mod tests {
|
||||
service_impl.my_peer_id,
|
||||
service_impl.my_peer_route_id,
|
||||
&service_impl.global_ctx,
|
||||
None,
|
||||
);
|
||||
let mut self_info = self_info;
|
||||
self_info.version = 1;
|
||||
|
||||
@@ -41,6 +41,10 @@ impl DirectConnectorRpc for DirectConnectorManagerRpcServer {
|
||||
let et_ipv6: crate::proto::common::Ipv6Addr = et_ipv6.address().into();
|
||||
ret.interface_ipv6s.retain(|x| *x != et_ipv6);
|
||||
}
|
||||
if let Some(public_ipv6) = self.global_ctx.get_public_ipv6_lease() {
|
||||
let public_ipv6: crate::proto::common::Ipv6Addr = public_ipv6.address().into();
|
||||
ret.interface_ipv6s.retain(|x| *x != public_ipv6);
|
||||
}
|
||||
tracing::trace!(
|
||||
"get_ip_list: public_ipv4: {:?}, public_ipv6: {:?}, listeners: {:?}",
|
||||
ret.public_ipv4,
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,3 +1,4 @@
|
||||
use cidr::Ipv6Inet;
|
||||
use cidr::{Ipv4Cidr, Ipv6Cidr};
|
||||
use dashmap::DashMap;
|
||||
use std::{
|
||||
@@ -8,9 +9,12 @@ use std::{
|
||||
|
||||
use crate::{
|
||||
common::{PeerId, global_ctx::NetworkIdentity},
|
||||
proto::peer_rpc::{
|
||||
ForeignNetworkRouteInfoEntry, ForeignNetworkRouteInfoKey, PeerIdentityType,
|
||||
RouteForeignNetworkInfos, RouteForeignNetworkSummary, RoutePeerInfo,
|
||||
proto::{
|
||||
api::instance::ListPublicIpv6InfoResponse,
|
||||
peer_rpc::{
|
||||
ForeignNetworkRouteInfoEntry, ForeignNetworkRouteInfoKey, PeerIdentityType,
|
||||
RouteForeignNetworkInfos, RouteForeignNetworkSummary, RoutePeerInfo,
|
||||
},
|
||||
},
|
||||
};
|
||||
|
||||
@@ -93,6 +97,22 @@ pub trait Route {
|
||||
// TODO: rewrite route management, remove this
|
||||
async fn list_proxy_cidrs_v6(&self) -> BTreeSet<Ipv6Cidr>;
|
||||
|
||||
async fn list_public_ipv6_routes(&self) -> BTreeSet<Ipv6Inet> {
|
||||
BTreeSet::new()
|
||||
}
|
||||
|
||||
async fn get_my_public_ipv6_addr(&self) -> Option<Ipv6Inet> {
|
||||
None
|
||||
}
|
||||
|
||||
async fn get_public_ipv6_gateway_peer_id(&self) -> Option<PeerId> {
|
||||
None
|
||||
}
|
||||
|
||||
async fn get_local_public_ipv6_info(&self) -> ListPublicIpv6InfoResponse {
|
||||
ListPublicIpv6InfoResponse::default()
|
||||
}
|
||||
|
||||
async fn get_peer_id_by_ipv4(&self, _ipv4: &Ipv4Addr) -> Option<PeerId> {
|
||||
None
|
||||
}
|
||||
@@ -194,6 +214,14 @@ impl Route for MockRoute {
|
||||
unimplemented!()
|
||||
}
|
||||
|
||||
async fn list_public_ipv6_routes(&self) -> BTreeSet<Ipv6Inet> {
|
||||
unimplemented!()
|
||||
}
|
||||
|
||||
async fn get_my_public_ipv6_addr(&self) -> Option<Ipv6Inet> {
|
||||
panic!("mock route")
|
||||
}
|
||||
|
||||
async fn get_peer_info(&self, _peer_id: PeerId) -> Option<RoutePeerInfo> {
|
||||
panic!("mock route")
|
||||
}
|
||||
|
||||
@@ -13,9 +13,9 @@ use crate::{
|
||||
GetWhitelistRequest, GetWhitelistResponse, ListCredentialsRequest,
|
||||
ListCredentialsResponse, ListForeignNetworkRequest, ListForeignNetworkResponse,
|
||||
ListGlobalForeignNetworkRequest, ListGlobalForeignNetworkResponse, ListPeerRequest,
|
||||
ListPeerResponse, ListRouteRequest, ListRouteResponse, PeerInfo, PeerManageRpc,
|
||||
RevokeCredentialRequest, RevokeCredentialResponse, ShowNodeInfoRequest,
|
||||
ShowNodeInfoResponse,
|
||||
ListPeerResponse, ListPublicIpv6InfoRequest, ListPublicIpv6InfoResponse,
|
||||
ListRouteRequest, ListRouteResponse, PeerInfo, PeerManageRpc, RevokeCredentialRequest,
|
||||
RevokeCredentialResponse, ShowNodeInfoRequest, ShowNodeInfoResponse,
|
||||
},
|
||||
rpc_types::{self, controller::BaseController},
|
||||
},
|
||||
@@ -99,6 +99,16 @@ impl PeerManageRpc for PeerManagerRpcService {
|
||||
Ok(reply)
|
||||
}
|
||||
|
||||
async fn list_public_ipv6_info(
|
||||
&self,
|
||||
_: BaseController,
|
||||
_request: ListPublicIpv6InfoRequest,
|
||||
) -> Result<ListPublicIpv6InfoResponse, rpc_types::error::Error> {
|
||||
Ok(weak_upgrade(&self.peer_manager)?
|
||||
.get_local_public_ipv6_info()
|
||||
.await)
|
||||
}
|
||||
|
||||
async fn list_route(
|
||||
&self,
|
||||
_: BaseController,
|
||||
|
||||
Reference in New Issue
Block a user