use std::sync::Arc;
use kameo::{
actor::ActorRef,
message::{Context, Message},
};
use tokio::sync::watch;
use crate::{Error, env::Env};
pub type CapGrants = Arc<Vec<ts_packetfilter_state::CapGrant>>;
pub struct PacketfilterUpdater {
env: Env,
pf_state: ts_packetfilter::CheckingFilter<
ts_packetfilter::HashbrownFilter,
ts_bart_packetfilter::BartFilter,
>,
self_addrs: Vec<std::net::IpAddr>,
cap_grants_tx: watch::Sender<CapGrants>,
live_filter_tx: LiveFilterTx,
unlocked_allowed_ips: std::collections::BTreeMap<ts_control::NodeId, Vec<ipnet::IpNet>>,
unlocked_nodes_permitted: bool,
}
#[derive(Clone)]
pub struct PacketFilterState(pub Arc<dyn ts_packetfilter::Filter + Send + Sync>);
pub(crate) type LiveFilterRx = watch::Receiver<Option<PacketFilterState>>;
pub(crate) type LiveFilterTx = watch::Sender<Option<PacketFilterState>>;
impl kameo::Actor for PacketfilterUpdater {
type Args = (Env, watch::Sender<CapGrants>, LiveFilterTx);
type Error = Error;
async fn on_start(
(env, cap_grants_tx, live_filter_tx): Self::Args,
slf: ActorRef<Self>,
) -> Result<Self, Self::Error> {
env.subscribe::<Arc<ts_control::StateUpdate>>(&slf).await?;
Ok(Self {
env,
pf_state: Default::default(),
self_addrs: Vec::new(),
cap_grants_tx,
live_filter_tx,
unlocked_allowed_ips: Default::default(),
unlocked_nodes_permitted: false,
})
}
}
impl PacketfilterUpdater {
fn published_filter(&self) -> PacketFilterState {
if self.unlocked_nodes_permitted {
self.wrap_for_publication(ts_packetfilter::HashbrownFilter::default())
} else {
self.wrap_for_publication(self.pf_state.clone())
}
}
fn wrap_for_publication(
&self,
inner: impl ts_packetfilter::Filter + Send + Sync + 'static,
) -> PacketFilterState {
if self.env.block_incoming {
PacketFilterState(Arc::new(ts_packetfilter::ShieldsUpFilter {
inner,
self_addrs: self.self_addrs.clone(),
}))
} else {
PacketFilterState(Arc::new(inner))
}
}
fn track_unlocked_peers(&mut self, state_update: &ts_control::StateUpdate) -> bool {
let Some(peer_update) = state_update.peer_update.as_ref() else {
return false;
};
let before = self.unlocked_allowed_ips.clone();
match peer_update {
ts_control::PeerUpdate::Full(peers) => {
self.unlocked_allowed_ips = peers
.iter()
.filter(|peer| peer.unsigned_peer_api_only)
.map(|peer| (peer.id, peer.accepted_routes.clone()))
.collect();
}
ts_control::PeerUpdate::Delta { upsert, remove } => {
for peer in upsert {
if peer.unsigned_peer_api_only {
self.unlocked_allowed_ips
.insert(peer.id, peer.accepted_routes.clone());
} else {
self.unlocked_allowed_ips.remove(&peer.id);
}
}
for id in remove {
self.unlocked_allowed_ips.remove(id);
}
}
}
self.unlocked_allowed_ips != before
}
fn reassess_unlocked_nodes(&mut self) -> bool {
let allowed: Vec<ipnet::IpNet> = self
.unlocked_allowed_ips
.values()
.flatten()
.copied()
.collect();
let permitted = self
.pf_state
.0
.values()
.any(|ruleset| ts_packetfilter::permits_unlocked_nodes(ruleset, &allowed));
if permitted == self.unlocked_nodes_permitted {
return false;
}
if permitted {
tracing::warn!(
unlocked_peers = self.unlocked_allowed_ips.len(),
"control's packet filter grants network access to a peer outside tailnet lock \
(UnsignedPeerAPIOnly); ignoring the whole filter"
);
} else {
tracing::info!("control's packet filter no longer grants an unlocked peer access");
}
self.unlocked_nodes_permitted = permitted;
true
}
}
impl Message<Arc<ts_control::StateUpdate>> for PacketfilterUpdater {
type Reply = ();
async fn handle(
&mut self,
state_update: Arc<ts_control::StateUpdate>,
_ctx: &mut Context<Self, Self::Reply>,
) {
if self.env.block_incoming
&& let Some(self_node) = state_update.node.as_ref()
{
self.self_addrs = self_node_addresses(self_node, self.env.enable_ipv6);
}
if let Some(grants) = &state_update.cap_grants {
self.cap_grants_tx.send_replace(Arc::new(grants.clone()));
}
let peers_changed = self.track_unlocked_peers(&state_update);
let filter_changed = match &state_update.packetfilter {
Some((pf_ruleset, pf_map)) => {
ts_packetfilter_state::apply_update(&mut self.pf_state, pf_ruleset.clone(), pf_map);
tracing::trace!(updated_packet_filter = ?self.pf_state.0);
true
}
None => false,
};
if !filter_changed && !peers_changed {
return;
}
let verdict_changed = self.reassess_unlocked_nodes();
if !filter_changed && !verdict_changed {
return;
}
let filter = self.published_filter();
self.live_filter_tx.send_replace(Some(filter.clone()));
if let Err(e) = self.env.publish(filter).await {
tracing::error!(error = %e, "publishing packet filter state");
}
}
}
fn self_node_addresses(self_node: &ts_control::Node, enable_ipv6: bool) -> Vec<std::net::IpAddr> {
let tailnet_address = &self_node.tailnet_address;
let mut addrs: Vec<std::net::IpAddr> = vec![tailnet_address.ipv4.addr().into()];
if enable_ipv6 {
addrs.push(tailnet_address.ipv6.addr().into());
}
addrs.push(core::net::Ipv4Addr::new(100, 100, 100, 100).into());
for vip in self_node.service_addresses() {
if vip.is_ipv6() && !enable_ipv6 {
continue;
}
if !addrs.contains(&vip) {
addrs.push(vip);
}
}
addrs
}
#[cfg(test)]
mod unlocked_node_tests {
use std::sync::Arc;
use kameo::actor::Spawn;
use tokio::sync::watch;
use ts_control::{Node, StableNodeId, TailnetAddress};
use ts_packetfilter::FilterExt;
use super::*;
fn forwarder_cfg() -> crate::env::ForwarderConfig {
crate::env::ForwarderConfig {
accept_routes: false,
accept_dns: true,
exit_node: None,
forward_routes: vec![],
forward_tcp_ports: vec![],
forward_udp_ports: vec![],
forward_all_ports: false,
forward_exit_egress: false,
block_incoming: false,
exit_proxy: None,
peerapi_port: None,
taildrop_dir: None,
enable_ipv6: false,
wireguard_listen_port: None,
network_monitor: false,
persistent_keepalive_interval: None,
ingress_active: Arc::new(std::sync::atomic::AtomicBool::new(false)),
}
}
fn updater() -> (ActorRef<PacketfilterUpdater>, LiveFilterRx) {
let (_shutdown_tx, shutdown_rx) = watch::channel(false);
let env =
crate::env::Env::new(ts_keys::NodeState::generate(), shutdown_rx, forwarder_cfg());
let (cap_grants_tx, _cap_grants_rx) = watch::channel(Default::default());
let (filter_tx, filter_rx) = watch::channel(None);
let updater = PacketfilterUpdater::spawn((env, cap_grants_tx, filter_tx));
(updater, filter_rx)
}
fn peer(id: i64, unsigned: bool) -> Node {
let v4: ipnet::IpNet = format!("100.64.0.{id}/32").parse().unwrap();
Node {
id,
stable_id: StableNodeId(format!("n{id}")),
hostname: format!("peer{id}"),
user_id: 0,
tailnet: None,
tags: vec![],
addresses: vec![v4],
tailnet_address: TailnetAddress {
ipv4: v4.to_string().parse().unwrap(),
ipv6: "fd7a::1/128".parse().unwrap(),
},
node_key: [0u8; 32].into(),
node_key_expiry: None,
expired: false,
online: None,
last_seen: None,
key_signature: vec![],
machine_key: None,
disco_key: None,
accepted_routes: vec![v4],
underlay_addresses: vec![],
derp_region: None,
cap: Default::default(),
cap_map: Default::default(),
peerapi_port: None,
peerapi_dns_proxy: false,
is_wireguard_only: false,
exit_node_dns_resolvers: vec![],
peer_relay: false,
ssh_host_keys: vec![],
service_vips: Default::default(),
unsigned_peer_api_only: unsigned,
}
}
fn tcp_rule(src: &str) -> ts_packetfilter::Rule {
ts_packetfilter::Rule {
src: ts_packetfilter::SrcMatch {
pfxs: vec![src.parse().unwrap()],
caps: vec![],
},
protos: vec![ts_packetfilter::IpProto::TCP],
dst: vec![ts_packetfilter::DstMatch {
ports: 0..=u16::MAX,
ips: vec!["100.64.0.1/32".parse().unwrap()],
}],
}
}
fn netmap(
peer_update: Option<ts_control::PeerUpdate>,
rules: Option<Vec<ts_packetfilter::Rule>>,
) -> Arc<ts_control::StateUpdate> {
Arc::new(ts_control::StateUpdate {
session_handle: None,
seq: 0,
keep_alive: false,
derp: None,
node: None,
peer_update,
peer_patches: Vec::new(),
user_profiles: Vec::new(),
ping: None,
packetfilter: rules.map(|r| (Some(r), Default::default())),
cap_grants: None,
pop_browser_url: None,
dial_plan: None,
dns_config: None,
ssh_policy: None,
tka: None,
online_change: Default::default(),
peer_seen_change: Default::default(),
control_time: None,
})
}
fn admits(filter_rx: &LiveFilterRx, src: &str) -> bool {
let filter = filter_rx.borrow().clone().expect("a filter was published");
filter.0.can_access(
&ts_packetfilter::PacketInfo {
src: src.parse().unwrap(),
dst: "100.64.0.1".parse().unwrap(),
ip_proto: ts_packetfilter::IpProto::TCP,
port: 22,
l4: ts_packetfilter::L4Header::Unknown,
},
core::iter::empty(),
)
}
async fn deliver(
updater: &ActorRef<PacketfilterUpdater>,
filter_rx: &mut LiveFilterRx,
update: Arc<ts_control::StateUpdate>,
) {
filter_rx.mark_unchanged();
updater.tell(update).await.expect("netmap delivered");
tokio::time::timeout(std::time::Duration::from_secs(5), filter_rx.changed())
.await
.expect("the updater republished within five seconds")
.expect("the live-filter cell is still open");
}
#[tokio::test]
async fn an_acl_granting_an_unlocked_peer_access_is_ignored_wholesale() {
let (updater, mut filter_rx) = updater();
deliver(
&updater,
&mut filter_rx,
netmap(
Some(ts_control::PeerUpdate::Full(vec![
peer(9, true),
peer(10, false),
])),
Some(vec![tcp_rule("100.64.0.9/32"), tcp_rule("100.64.0.10/32")]),
),
)
.await;
assert!(
!admits(&filter_rx, "100.64.0.9"),
"an ACL naming a peer outside tailnet lock must not be installed"
);
assert!(
!admits(&filter_rx, "100.64.0.10"),
"the rest of that filter must go with it, not be kept as a partial ACL"
);
}
#[tokio::test]
async fn the_same_acl_is_installed_when_no_peer_is_unlocked() {
let (updater, mut filter_rx) = updater();
deliver(
&updater,
&mut filter_rx,
netmap(
Some(ts_control::PeerUpdate::Full(vec![
peer(9, false),
peer(10, false),
])),
Some(vec![tcp_rule("100.64.0.9/32"), tcp_rule("100.64.0.10/32")]),
),
)
.await;
assert!(admits(&filter_rx, "100.64.0.9"));
assert!(admits(&filter_rx, "100.64.0.10"));
}
#[tokio::test]
async fn an_unlocked_peer_arriving_later_invalidates_the_installed_filter() {
let (updater, mut filter_rx) = updater();
deliver(
&updater,
&mut filter_rx,
netmap(
Some(ts_control::PeerUpdate::Full(vec![peer(9, false)])),
Some(vec![tcp_rule("100.64.0.0/24")]),
),
)
.await;
assert!(
admits(&filter_rx, "100.64.0.9"),
"a filter with no unlocked peer in the netmap is installed"
);
deliver(
&updater,
&mut filter_rx,
netmap(
Some(ts_control::PeerUpdate::Delta {
upsert: vec![peer(9, true)],
remove: vec![],
}),
None,
),
)
.await;
assert!(
!admits(&filter_rx, "100.64.0.9"),
"an unlocked peer appearing under an existing ACL must invalidate that ACL"
);
}
#[tokio::test]
async fn removing_the_unlocked_peer_restores_the_filter() {
let (updater, mut filter_rx) = updater();
deliver(
&updater,
&mut filter_rx,
netmap(
Some(ts_control::PeerUpdate::Full(vec![peer(9, true)])),
Some(vec![tcp_rule("100.64.0.0/24")]),
),
)
.await;
assert!(!admits(&filter_rx, "100.64.0.9"));
deliver(
&updater,
&mut filter_rx,
netmap(
Some(ts_control::PeerUpdate::Delta {
upsert: vec![],
remove: vec![9],
}),
None,
),
)
.await;
assert!(
admits(&filter_rx, "100.64.0.9"),
"the compiled filter is set aside while an unlocked peer is present, not dropped"
);
}
}