Skip to content

Commit 52ee454

Browse files
committed
use mg-common macros
1 parent 8b282cf commit 52ee454

9 files changed

Lines changed: 33 additions & 23 deletions

File tree

Cargo.lock

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

bfd-async/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ edition = "2024"
55

66
[dependencies]
77
bfd.workspace = true
8+
mg-common.workspace = true
89
mg-api-types.workspace = true
910
rdb.workspace = true
1011
slog.workspace = true

bfd-async/src/dispatcher.rs

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55
use crate::AddPeerError;
66
use bfd::SessionCounters;
77
use bfd::packet;
8+
use mg_common::read_lock;
9+
use mg_common::write_lock;
810
use slog::Logger;
911
use slog::warn;
1012
use slog_error_chain::InlineErrorChain;
@@ -228,7 +230,7 @@ impl Listener {
228230
}
229231

230232
fn remove_peer(&self, peer: IpAddr) -> ListenerRemovePeerResult {
231-
let mut sessions = self.sessions.write().unwrap();
233+
let mut sessions = write_lock!(self.sessions);
232234
if sessions.remove(&peer).is_some() {
233235
if sessions.is_empty() {
234236
ListenerRemovePeerResult::RemovedNowEmpty
@@ -245,7 +247,7 @@ impl Listener {
245247
peer: IpAddr,
246248
session: SessionHandle,
247249
) -> Result<(), AddPeerError> {
248-
let mut sessions = self.sessions.write().unwrap();
250+
let mut sessions = write_lock!(self.sessions);
249251
match sessions.entry(peer) {
250252
hash_map::Entry::Occupied(_) => Err(AddPeerError::PeerExists(peer)),
251253
hash_map::Entry::Vacant(entry) => {
@@ -346,7 +348,7 @@ impl ListenerTask {
346348

347349
// Do we expect traffic from this peer address?
348350
let Some(session) =
349-
self.sessions.read().unwrap().get(&peer.ip()).cloned()
351+
read_lock!(self.sessions).get(&peer.ip()).cloned()
350352
else {
351353
warn!(
352354
self.log, "unknown peer; dropping";

bfd-async/src/dispatcher/end_to_end.rs

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ use crate::wait_for_condition;
2323
use bfd::DEFAULT_DETECT_MULTIPLIER;
2424
use mg_api_types::bfd::BfdPeerState;
2525
use mg_api_types::bfd::SessionMode;
26+
use mg_common::lock;
2627
use slog::Logger;
2728
use std::net::IpAddr;
2829
use std::net::Ipv4Addr;
@@ -78,10 +79,7 @@ impl ListenerBackend for PreBoundSocketBackend {
7879
sessions: SharedSessions,
7980
log: Logger,
8081
) -> Result<JoinHandle<()>, AddPeerError> {
81-
let socket = self
82-
.socket
83-
.lock()
84-
.unwrap()
82+
let socket = lock!(self.socket)
8583
.take()
8684
.expect("PreBoundSocketBackend::spawn called more than once");
8785
let listen_task = ListenerTask::new(socket, sessions, log);

bfd-async/src/dispatcher/proptests.rs

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,9 @@ use super::Dispatcher;
1313
use super::ListenerBackend;
1414
use super::SharedSessions;
1515
use crate::AddPeerError;
16+
use mg_common::parse;
17+
use mg_common::read_lock;
18+
use mg_common::sockaddr;
1619
use proptest::collection::vec;
1720
use proptest::prelude::*;
1821
use proptest::sample::select;
@@ -63,9 +66,9 @@ impl ListenerBackend for FakeBackend {
6366
// IPs/ports are arbitrary.
6467
fn arb_listen_addr() -> impl Strategy<Value = SocketAddr> {
6568
let addrs = vec![
66-
"127.0.0.1:3784".parse().unwrap(),
67-
"127.0.0.1:4784".parse().unwrap(),
68-
"192.0.2.1:3784".parse().unwrap(),
69+
sockaddr!("127.0.0.1:3784"),
70+
sockaddr!("127.0.0.1:4784"),
71+
sockaddr!("192.0.2.1:3784"),
6972
];
7073
select(addrs)
7174
}
@@ -146,7 +149,7 @@ fn check_invariants(dispatcher: &Dispatcher, model: &Model) {
146149
// that address, and no listener is empty.
147150
for (addr, listener) in &dispatcher.listeners {
148151
let actual: BTreeSet<IpAddr> =
149-
listener.sessions.read().unwrap().keys().copied().collect();
152+
read_lock!(listener.sessions).keys().copied().collect();
150153
let expected: BTreeSet<IpAddr> = model
151154
.peer_to_addr
152155
.iter()

bfd-async/src/dispatcher/tests.rs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,9 @@ use super::ListenerBackend;
1010
use super::ListenerTask;
1111
use bfd::SessionCounters;
1212
use bfd::packet::Control;
13+
use mg_common::ip;
14+
use mg_common::lock;
15+
use mg_common::parse;
1316
use slog::Discard;
1417
use slog::Logger;
1518
use slog::o;
@@ -51,10 +54,7 @@ impl TestBackend {
5154
}
5255

5356
fn expect_bound_address(&self) -> SocketAddr {
54-
self.bound_address
55-
.lock()
56-
.unwrap()
57-
.expect("socket has been bound")
57+
lock!(self.bound_address).expect("socket has been bound")
5858
}
5959

6060
fn was_listener_task_dropped(&self) -> bool {
@@ -71,7 +71,7 @@ impl ListenerBackend for TestBackend {
7171
) -> Result<tokio::task::JoinHandle<()>, crate::AddPeerError> {
7272
assert_eq!(listen_addr, LISTEN_ADDR);
7373

74-
let mut bound_address = self.bound_address.lock().unwrap();
74+
let mut bound_address = lock!(self.bound_address);
7575
assert!(
7676
bound_address.is_none(),
7777
"TestBackend can only spawn one listening socket"
@@ -154,7 +154,7 @@ async fn drops_packet_from_unregistered_peer() {
154154

155155
// Register a peer that is *not* the loopback source address our client
156156
// sends from, so the listener sees a mismatched source IP.
157-
let peer: IpAddr = "192.0.2.1".parse().unwrap();
157+
let peer = ip!("192.0.2.1");
158158

159159
let mut rx = dispatcher
160160
.ensure(LISTEN_ADDR, peer, Arc::default(), &log)

bfd-async/src/egress/tests.rs

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
// file, You can obtain one at https://mozilla.org/MPL/2.0/.
44

55
use super::*;
6+
use mg_common::parse;
7+
use mg_common::sockaddr;
68
use slog::Discard;
79
use slog::o;
810
use std::net::Ipv4Addr;
@@ -12,7 +14,7 @@ use tokio::time::timeout;
1214

1315
#[tokio::test]
1416
async fn egress_socket_sets_ttl_v4() {
15-
let sk = bind_egress_socket("127.0.0.1:0".parse().unwrap())
17+
let sk = bind_egress_socket(sockaddr!("127.0.0.1:0"))
1618
.await
1719
.expect("created v4 egress socket");
1820
let sock = Socket::from(sk.into_std().expect("converted socket"));
@@ -21,7 +23,7 @@ async fn egress_socket_sets_ttl_v4() {
2123

2224
#[tokio::test]
2325
async fn egress_socket_sets_hop_limit_v6() {
24-
let sk = bind_egress_socket("[::1]:0".parse().unwrap())
26+
let sk = bind_egress_socket(sockaddr!("[::1]:0"))
2527
.await
2628
.expect("created v6 egress socket");
2729
let sock = Socket::from(sk.into_std().expect("converted socket"));
@@ -145,7 +147,7 @@ async fn binds_ports_in_expected_range() {
145147
#[tokio::test(flavor = "multi_thread")]
146148
async fn exits_when_channel_closed() {
147149
let (tx, _counters, handle) =
148-
spawn_egress("127.0.0.1:1".parse().unwrap(), Arc::default());
150+
spawn_egress(sockaddr!("127.0.0.1:1"), Arc::default());
149151
drop(tx);
150152
timeout(Duration::from_secs(5), handle)
151153
.await

bfd-async/src/rib.rs

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -100,6 +100,7 @@ impl RibTask {
100100
mod tests {
101101
use super::*;
102102
use crate::wait_for_condition;
103+
use mg_common::lock;
103104
use std::net::Ipv4Addr;
104105
use std::sync::Arc;
105106
use std::sync::Mutex;
@@ -115,13 +116,13 @@ mod tests {
115116

116117
impl NexthopSink for RecordingSink {
117118
fn set_nexthop_shutdown(&self, nexthop: IpAddr, shutdown: bool) {
118-
self.0.lock().unwrap().push((nexthop, shutdown));
119+
lock!(self.0).push((nexthop, shutdown));
119120
}
120121
}
121122

122123
impl RecordingSink {
123124
fn writes(&self) -> Vec<(IpAddr, bool)> {
124-
self.0.lock().unwrap().clone()
125+
lock!(self.0).clone()
125126
}
126127
}
127128

bfd-async/src/session/tests.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@ use super::*;
1414
use crate::wait_for_condition;
1515
use bfd::DEFAULT_DETECT_MULTIPLIER;
1616
use bfd::packet::Control;
17+
use mg_common::parse;
18+
use mg_common::sockaddr;
1719
use slog::Discard;
1820
use slog::o;
1921
use tokio::time::timeout;
@@ -49,7 +51,7 @@ fn spawn_driver(required_rx: Duration) -> Harness {
4951
DriverTask {
5052
sm,
5153
counters: Arc::clone(&counters),
52-
remote_addr: "127.0.0.1:3784".parse().unwrap(),
54+
remote_addr: sockaddr!("127.0.0.1:3784"),
5355
listener_rx,
5456
egress_tx,
5557
state_tx,

0 commit comments

Comments
 (0)