Skip to content

Commit 4622836

Browse files
committed
Replace all Packet with PacketPtr
1 parent ebda36c commit 4622836

51 files changed

Lines changed: 205 additions & 260 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

source/network/connection/flow/packet.hpp

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
#pragma once
22

3+
#include <memory>
34
#include <string>
45

56
#include "four_tuple.hpp"
@@ -12,7 +13,9 @@ using PathHash = std::uint32_t;
1213

1314
struct Packet;
1415

15-
using OnPacketDeliveryCallback = std::function<void(const Packet&)>;
16+
using PacketPtr = std::shared_ptr<Packet>;
17+
18+
using OnPacketDeliveryCallback = std::function<void(PacketPtr)>;
1619

1720
struct Packet : FourTuple {
1821
Packet(SizeByte a_size = SizeByte(0ul), Id a_source_id = "",
@@ -33,7 +36,7 @@ struct Packet : FourTuple {
3336
// Note: callback takes packet itself because it might be changed while
3437
// travelling along the net
3538
OnPacketDeliveryCallback callback =
36-
[]([[maybe_unused]] const Packet& packet) {};
39+
[]([[maybe_unused]] PacketPtr packet) {};
3740
TimeNs generated_time; // Note: ACK's generated time is the data packet
3841
// generated time
3942
TimeNs sent_time; // Note: ACK's sent time is the data packet sent time

source/network/connection/flow/packet_ack_info.hpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ namespace sim {
88
struct PacketAckInfo {
99
TimeNs rtt;
1010
const utils::Statistics<TimeNs>& rtt_stat;
11-
const Packet& ack;
11+
PacketPtr ack;
1212
};
1313

1414
} // namespace sim

source/network/connection/flow/tcp/tcp_flow.cpp

Lines changed: 19 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ TcpFlow::TcpFlow(Id a_id, FlowFourTuple a_four_tuple, bool a_ecn_capable,
7272
}
7373
}
7474

75-
std::shared_ptr<Packet> TcpFlow::create_data_packet(
75+
PacketPtr TcpFlow::create_data_packet(
7676
PacketInfo info, std::shared_ptr<IHost> sender,
7777
std::shared_ptr<IHost> receiver) {
7878
auto packet = std::make_shared<Packet>();
@@ -82,7 +82,7 @@ std::shared_ptr<Packet> TcpFlow::create_data_packet(
8282
packet->flags.set_flag(m_packet_type_label, PacketType::DATA)
8383
.log_err_if_not_present("Failed to set packet type (DATA)");
8484

85-
set_avg_rtt_if_present(*packet);
85+
set_avg_rtt_if_present(packet);
8686

8787
packet->data_id = std::move(info.id);
8888
packet->sender_id = sender->get_id();
@@ -94,7 +94,7 @@ std::shared_ptr<Packet> TcpFlow::create_data_packet(
9494
std::shared_ptr<TcpFlow> flow = shared_from_this();
9595

9696
packet->callback = [flow, delivery_callback = info.callback](
97-
const Packet& delivered_packet) {
97+
PacketPtr delivered_packet) {
9898
flow->process_data_packet(delivered_packet, delivery_callback);
9999
};
100100
packet->generated_time = info.generated_time;
@@ -107,15 +107,15 @@ std::shared_ptr<Packet> TcpFlow::create_data_packet(
107107
return packet;
108108
}
109109

110-
void TcpFlow::set_avg_rtt_if_present(Packet& packet) {
110+
void TcpFlow::set_avg_rtt_if_present(PacketPtr packet) {
111111
std::optional<TimeNs> avg_rtt = m_context.rtt_statistics.get_mean();
112112
if (avg_rtt.has_value()) {
113-
set_avg_rtt_flag(packet.flags, avg_rtt.value())
113+
set_avg_rtt_flag(packet->flags, avg_rtt.value())
114114
.log_err_if_not_present("Failed to set average RTT");
115115
}
116116
}
117117

118-
void TcpFlow::send_data_packet(std::shared_ptr<Packet> data) {
118+
void TcpFlow::send_data_packet(PacketPtr data) {
119119
TimeNs now = Scheduler::get_instance().get_current_time();
120120

121121
Scheduler::get_instance().add(
@@ -127,13 +127,13 @@ void TcpFlow::send_data_packet(std::shared_ptr<Packet> data) {
127127
m_context.sender->enqueue_packet(data);
128128
}
129129

130-
void TcpFlow::process_data_packet(const Packet& data,
130+
void TcpFlow::process_data_packet(PacketPtr data,
131131
const PacketCallback& callback) {
132-
m_packet_reordering.add_record(data.packet_num);
132+
m_packet_reordering.add_record(data->packet_num);
133133
m_metrics.packet_reordering->add_record(
134134
Scheduler::get_instance().get_current_time(),
135135
m_packet_reordering.value());
136-
auto ack = std::make_shared<Packet>(data);
136+
auto ack = std::make_shared<Packet>(*data);
137137
ack->sender_id = IdWithHash(m_context.receiver->get_id());
138138
ack->sender_port = m_context.receiver_port;
139139
ack->receiver_id = IdWithHash(m_context.sender->get_id());
@@ -146,42 +146,42 @@ void TcpFlow::process_data_packet(const Packet& data,
146146
fmt::format("Flow {}: could not set type label to ack packet {}",
147147
m_id, ack->to_string()));
148148
}
149-
SizeByte data_packet_size = data.size;
149+
SizeByte data_packet_size = data->size;
150150

151151
std::shared_ptr<TcpFlow> flow = shared_from_this();
152152
ack->callback = [flow, callback,
153-
data_packet_size](const Packet& delivered_ack) {
153+
data_packet_size](PacketPtr delivered_ack) {
154154
flow->process_ack(delivered_ack, data_packet_size, callback);
155155
};
156156

157157
m_context.receiver->enqueue_packet(std::move(ack));
158158
}
159159

160-
void TcpFlow::process_ack(const Packet& ack, SizeByte data_packet_size,
160+
void TcpFlow::process_ack(PacketPtr ack, SizeByte data_packet_size,
161161
PacketCallback callback) {
162162
TimeNs now = Scheduler::get_instance().get_current_time();
163163

164-
TimeNs rtt = now - ack.sent_time;
164+
TimeNs rtt = now - ack->sent_time;
165165

166166
m_context.rtt_statistics.add_record(rtt);
167167

168168
m_metrics.rtt->add_record(now, rtt.value());
169169

170170
update_rto_on_ack();
171171

172-
if (!m_ack_monitor.confirm_one(ack.packet_num)) {
172+
if (!m_ack_monitor.confirm_one(ack->packet_num)) {
173173
LOG_WARN(
174174
fmt::format("Flow {} got ack {} that confirms nothing; ignored",
175-
m_id, ack.to_string()));
175+
m_id, ack->to_string()));
176176
return;
177177
}
178178
m_context.last_ack_receive_time = now;
179179

180180
m_context.delivered_size += data_packet_size;
181181

182182
SpeedGbps delivery_rate =
183-
(m_context.delivered_size - ack.delivered_data_size_at_origin) /
184-
(now - ack.generated_time);
183+
(m_context.delivered_size - ack->delivered_data_size_at_origin) /
184+
(now - ack->generated_time);
185185

186186
m_metrics.delivery_rate->add_record(now, delivery_rate.value());
187187
m_context.delivery_rate_statistics.add_record(delivery_rate);
@@ -196,7 +196,7 @@ void TcpFlow::update_rto_on_ack() {
196196
m_rto.is_steady = true;
197197
}
198198

199-
void TcpFlow::on_timeout(std::shared_ptr<Packet> data) {
199+
void TcpFlow::on_timeout(PacketPtr data) {
200200
if (m_ack_monitor.is_confirmed(data->packet_num)) {
201201
LOG_INFO(fmt::format(
202202
"Flow {}: packet {} is confirmed when timeout reached; no retransmit",
@@ -213,7 +213,7 @@ void TcpFlow::update_rto_on_timeout() {
213213
}
214214
}
215215

216-
void TcpFlow::retransmit_packet(std::shared_ptr<Packet> data) {
216+
void TcpFlow::retransmit_packet(PacketPtr data) {
217217
m_context.retransmit_size += data->size;
218218
send_data_packet(data);
219219
}

source/network/connection/flow/tcp/tcp_flow.hpp

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -43,25 +43,25 @@ class TcpFlow : public IFlow, public std::enable_shared_from_this<TcpFlow> {
4343
TcpFlow(Id a_id, FlowFourTuple a_four_tuple, bool a_ecn_capable, RTO a_rto,
4444
TcpFlowMetricsFilters a_metrics_flags);
4545

46-
std::shared_ptr<Packet> create_data_packet(PacketInfo info,
46+
PacketPtr create_data_packet(PacketInfo info,
4747
std::shared_ptr<IHost> sender,
4848
std::shared_ptr<IHost> receiver);
4949

50-
void set_avg_rtt_if_present(Packet& packet);
50+
void set_avg_rtt_if_present(PacketPtr packet);
5151

52-
void send_data_packet(std::shared_ptr<Packet> data);
52+
void send_data_packet(PacketPtr data);
5353

54-
void process_data_packet(const Packet& data_packet,
54+
void process_data_packet(PacketPtr data_packet,
5555
const PacketCallback& callback);
5656

57-
void process_ack(const Packet& ack, SizeByte data_packet_size,
57+
void process_ack(PacketPtr ack, SizeByte data_packet_size,
5858
PacketCallback callback);
5959

6060
void update_rto_on_ack();
6161

62-
void on_timeout(std::shared_ptr<Packet> data);
62+
void on_timeout(PacketPtr data);
6363
void update_rto_on_timeout();
64-
void retransmit_packet(std::shared_ptr<Packet> data);
64+
void retransmit_packet(PacketPtr data);
6565

6666
private:
6767
// flag labels

source/network/connection/mplb/cc/tahoe/tcp_tahoe_cc.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ TcpTahoeCC::TcpTahoeCC(double a_start_cwnd, double a_sstresh)
1313

1414
void TcpTahoeCC::on_ack(const PacketAckInfo& info) {
1515
m_last_avg_rtt = info.rtt_stat.get_mean().value();
16-
if (info.ack.congestion_experienced) {
16+
if (info.ack->congestion_experienced) {
1717
on_timeout();
1818
} else if (m_cwnd < m_ssthresh) {
1919
// Slow start

source/network/connection/rdma/rdma_connection.cpp

Lines changed: 19 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -154,7 +154,7 @@ void RdmaConnection::send_next_data_packet() {
154154
schedule_data_send();
155155
}
156156

157-
std::shared_ptr<Packet> RdmaConnection::create_data_packet(const Data& data) {
157+
PacketPtr RdmaConnection::create_data_packet(const Data& data) {
158158
auto packet = std::make_shared<Packet>();
159159
packet->packet_num = m_next_packet_num++;
160160

@@ -167,7 +167,7 @@ std::shared_ptr<Packet> RdmaConnection::create_data_packet(const Data& data) {
167167

168168
RdmaConnectionPtr conn = shared_from_this();
169169

170-
packet->callback = [conn](const Packet& delivered_packet) {
170+
packet->callback = [conn](PacketPtr delivered_packet) {
171171
conn->process_data_packet(delivered_packet);
172172
};
173173
TimeNs now = Scheduler::get_instance().get_current_time();
@@ -178,27 +178,27 @@ std::shared_ptr<Packet> RdmaConnection::create_data_packet(const Data& data) {
178178
return packet;
179179
}
180180

181-
void RdmaConnection::process_data_packet(const Packet& packet) {
181+
void RdmaConnection::process_data_packet(PacketPtr packet) {
182182
if (!m_receiver_started) {
183183
m_receiver_started = true;
184184
schedule_ack_timer();
185185
}
186-
if (packet.packet_num < m_next_expected_packet_num) {
186+
if (packet->packet_num < m_next_expected_packet_num) {
187187
LOG_ERROR(fmt::format(
188188
"RDMA {}: receiver got data packet {} with number smaller "
189189
"than expected {}; ignored",
190-
m_id, packet.to_string(), m_next_expected_packet_num));
191-
} else if (packet.packet_num == m_next_expected_packet_num) {
190+
m_id, packet->to_string(), m_next_expected_packet_num));
191+
} else if (packet->packet_num == m_next_expected_packet_num) {
192192
process_expected_data_packet();
193193
} else {
194194
LOG_ERROR(fmt::format(
195195
"RDMA {}: receiver got data packet {} with number greater "
196196
"than expected {}; could not put it to reorder buffer; ignored",
197-
m_id, packet.to_string(), m_next_expected_packet_num));
197+
m_id, packet->to_string(), m_next_expected_packet_num));
198198
send_nak();
199199
return;
200200
}
201-
if (packet.congestion_experienced) {
201+
if (packet->congestion_experienced) {
202202
send_cnp();
203203
}
204204
}
@@ -230,7 +230,7 @@ void RdmaConnection::send_nak() {
230230

231231
RdmaConnectionPtr conn = shared_from_this();
232232

233-
nak->callback = [conn](const Packet& nak) { conn->process_nak(nak); };
233+
nak->callback = [conn](PacketPtr nak) { conn->process_nak(nak); };
234234
TimeNs now = Scheduler::get_instance().get_current_time();
235235
nak->generated_time = now;
236236
nak->sent_time = now;
@@ -241,8 +241,8 @@ void RdmaConnection::send_nak() {
241241
m_receiver->enqueue_packet(std::move(nak));
242242
}
243243

244-
void RdmaConnection::process_nak(const Packet& nak) {
245-
while (m_last_acked_pcn < nak.packet_num) {
244+
void RdmaConnection::process_nak(PacketPtr nak) {
245+
while (m_last_acked_pcn < nak->packet_num) {
246246
m_last_acked_pcn++;
247247
confirm_first_unconfirmed_packet();
248248
}
@@ -269,7 +269,7 @@ void RdmaConnection::send_ack() {
269269

270270
RdmaConnectionPtr conn = shared_from_this();
271271

272-
ack->callback = [conn](const Packet& delivered_packet) {
272+
ack->callback = [conn](PacketPtr delivered_packet) {
273273
conn->process_ack(delivered_packet);
274274
};
275275
TimeNs now = Scheduler::get_instance().get_current_time();
@@ -282,8 +282,8 @@ void RdmaConnection::send_ack() {
282282
m_receiver->enqueue_packet(std::move(ack));
283283
}
284284

285-
void RdmaConnection::process_ack(const Packet& ack) {
286-
PacketNum ack_num = ack.packet_num;
285+
void RdmaConnection::process_ack(PacketPtr ack) {
286+
PacketNum ack_num = ack->packet_num;
287287
if (ack_num < m_last_acked_pcn) {
288288
LOG_ERROR(
289289
fmt::format("RDMA connection {}: sender got ack with number "
@@ -322,7 +322,7 @@ void RdmaConnection::send_cnp() {
322322

323323
RdmaConnectionPtr conn = shared_from_this();
324324

325-
cnp->callback = [conn](const Packet& cnp) { conn->process_cnp(cnp); };
325+
cnp->callback = [conn](PacketPtr cnp) { conn->process_cnp(cnp); };
326326
TimeNs now = Scheduler::get_instance().get_current_time();
327327
cnp->generated_time = now;
328328
cnp->sent_time = now;
@@ -333,7 +333,7 @@ void RdmaConnection::send_cnp() {
333333
m_receiver->enqueue_packet(std::move(cnp));
334334
}
335335

336-
void RdmaConnection::process_cnp([[maybe_unused]] const Packet& cnp) {
336+
void RdmaConnection::process_cnp([[maybe_unused]] PacketPtr cnp) {
337337
m_dcqcn.on_cnp();
338338
}
339339

@@ -346,15 +346,15 @@ void RdmaConnection::confirm_first_unconfirmed_packet() {
346346
m_id));
347347
return;
348348
}
349-
const Packet& confirmed = *m_send_queue.front();
349+
PacketPtr confirmed = m_send_queue.front();
350350
utils::Defer defer([this]() {
351351
m_send_queue.pop_front();
352352
if (m_send_queue.empty()) {
353353
m_dcqcn.stop();
354354
}
355355
});
356356

357-
const DataId& data_id = confirmed.data_id;
357+
const DataId& data_id = confirmed->data_id;
358358
auto it = m_data_context_table.find(data_id);
359359
if (it == m_data_context_table.end()) {
360360
LOG_ERROR(
@@ -386,7 +386,7 @@ void RdmaConnection::send_ack_request() {
386386
RdmaConnectionPtr conn = shared_from_this();
387387

388388
ack_request->callback =
389-
[conn]([[maybe_unused]] const Packet& delivered_packet) {
389+
[conn]([[maybe_unused]] PacketPtr delivered_packet) {
390390
conn->process_ack_request();
391391
};
392392
TimeNs now = Scheduler::get_instance().get_current_time();

source/network/connection/rdma/rdma_connectrion.hpp

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -51,23 +51,23 @@ class RdmaConnection : public IConnection,
5151

5252
void send_next_data_packet();
5353

54-
std::shared_ptr<Packet> create_data_packet(const Data& data);
54+
PacketPtr create_data_packet(const Data& data);
5555

56-
void process_data_packet(const Packet& data);
56+
void process_data_packet(PacketPtr data);
5757

5858
void process_expected_data_packet();
5959

6060
void send_nak();
6161

62-
void process_nak(const Packet& nak);
62+
void process_nak(PacketPtr nak);
6363

6464
void send_ack();
6565

66-
void process_ack(const Packet& ack);
66+
void process_ack(PacketPtr ack);
6767

6868
void send_cnp();
6969

70-
void process_cnp(const Packet& cnp);
70+
void process_cnp(PacketPtr cnp);
7171

7272
void confirm_first_unconfirmed_packet();
7373

@@ -89,7 +89,7 @@ class RdmaConnection : public IConnection,
8989

9090
// Invariant: packet_num of first packet in this queue is equal to
9191
// m_last_acked_pcn
92-
std::deque<std::shared_ptr<Packet>> m_send_queue;
92+
std::deque<PacketPtr> m_send_queue;
9393

9494
uint32_t m_send_window;
9595

0 commit comments

Comments
 (0)