@@ -20,13 +20,8 @@ void NewTcpFlow::send(std::vector<PacketInfo> packets_info) {
2020 if (packets_info.empty ()) {
2121 return ;
2222 }
23- auto exp_endpoints = get_endpoints ();
24- if (!exp_endpoints.has_value ()) {
25- LOG_ERROR (fmt::format (" Flow {}; error while sending packets: {}" , m_id,
26- exp_endpoints.error ()));
27- return ;
28- }
29- auto [sender, receiver] = exp_endpoints.value ();
23+ auto sender = m_context.sender ;
24+ auto receiver = m_context.receiver ;
3025 for (auto info : packets_info) {
3126 Packet data = create_data_packet (std::move (info), sender, receiver);
3227
@@ -42,31 +37,19 @@ NewTcpFlow::NewTcpFlow(Id a_id, std::shared_ptr<IHost> a_sender,
4237 std::shared_ptr<IHost> a_receiver, bool a_ecn_capable,
4338 RTO a_rto)
4439 : m_id(std::move(a_id)),
45- m_context ({SizeByte (0 ), SizeByte (0 ), SizeByte (0 ),
40+ m_context ({Endpoints{a_sender, a_receiver}, SizeByte (0 ), SizeByte (0 ),
41+ SizeByte (0 ),
4642
47- std::nullopt , std::nullopt , utils::Statistics<TimeNs>(),
48-
49- a_sender, a_receiver}),
43+ std::nullopt , std::nullopt , utils::Statistics<TimeNs>()}),
5044 m_ecn_capable(a_ecn_capable),
5145 m_rto(std::move(a_rto)) {
5246 if (!m_flag_manager.register_flag_by_amount (m_packet_type_label,
5347 PacketType::ENUM_SIZE )) {
54- throw std::runtime_error (" Can not registrate packet type label" );
48+ throw std::runtime_error (" Can not register packet type label" );
5549 }
5650 if (!register_packet_avg_rtt_flag (m_flag_manager)) {
57- throw std::runtime_error (" Can not registrate packet avg rtt label" );
58- }
59- }
60-
61- utils::StrExpected<NewTcpFlow::Endpoints> NewTcpFlow::get_endpoints () const {
62- if (m_context.sender .expired ()) {
63- return std::unexpected (" sender pointer expired" );
64- }
65-
66- if (m_context.receiver .expired ()) {
67- return std::unexpected (" receiver pointer expired" );
51+ throw std::runtime_error (" Can not register packet avg rtt label" );
6852 }
69- return Endpoints{m_context.sender .lock (), m_context.receiver .lock ()};
7053}
7154
7255Packet NewTcpFlow::create_data_packet (PacketInfo info,
@@ -120,38 +103,21 @@ void NewTcpFlow::set_avg_rtt_if_present(Packet& packet) {
120103}
121104
122105void NewTcpFlow::send_data_packet (Packet data) {
123- if (m_context.sender .expired ()) {
124- LOG_ERROR (fmt::format (
125- " Flow {}: sender pointer expired; could not send data packet {}" ,
126- m_id, data.to_string ()));
127- return ;
128- }
129- auto sender = m_context.sender .lock ();
130-
131106 TimeNs now = Scheduler::get_instance ().get_current_time ();
132107
133108 Scheduler::get_instance ().add <Timeout>(now + m_rto.current ,
134109 shared_from_this (), data);
135110 m_context.sent_size += data.size ;
136111
137112 data.sent_time = now;
138- sender->enqueue_packet (std::move (data));
113+ m_context. sender ->enqueue_packet (std::move (data));
139114}
140115
141116void NewTcpFlow::process_data_packet (const Packet& data,
142117 PacketCallback callback) {
143- auto exp_endpoints = get_endpoints ();
144- if (!exp_endpoints.has_value ()) {
145- LOG_ERROR (
146- fmt::format (" Flow {}: error while processing data packet {}: {}" ,
147- m_id, data.to_string (), exp_endpoints.error ()));
148- return ;
149- }
150- auto [sender, receiver] = exp_endpoints.value ();
151-
152118 Packet ack = data;
153- ack.source_id = receiver->get_id ();
154- ack.dest_id = sender->get_id ();
119+ ack.source_id = m_context. receiver ->get_id ();
120+ ack.dest_id = m_context. sender ->get_id ();
155121 ack.size = SizeByte (1 );
156122 ack.ttl = M_MAX_TTL ;
157123 ack.flags .set_flag (m_packet_type_label, PacketType::ACK )
@@ -173,7 +139,7 @@ void NewTcpFlow::process_data_packet(const Packet& data,
173139 callback);
174140 };
175141
176- receiver->enqueue_packet (ack);
142+ m_context. receiver ->enqueue_packet (ack);
177143}
178144
179145void NewTcpFlow::process_ack (const Packet& ack, SizeByte data_packet_size,
@@ -189,8 +155,9 @@ void NewTcpFlow::process_ack(const Packet& ack, SizeByte data_packet_size,
189155 update_rto_on_ack ();
190156
191157 if (!m_ack_monitor.confirm_one (ack.packet_num )) {
192- LOG_WARN (fmt::format (" Flow {} got ack {} that confirms nothing; ignored" , m_id,
193- ack.to_string ()));
158+ LOG_WARN (
159+ fmt::format (" Flow {} got ack {} that confirms nothing; ignored" ,
160+ m_id, ack.to_string ()));
194161 return ;
195162 }
196163
@@ -217,7 +184,7 @@ void NewTcpFlow::update_rto_on_ack() {
217184void NewTcpFlow::on_timeout (Packet data) {
218185 if (m_ack_monitor.is_confirmed (data.packet_num )) {
219186 LOG_INFO (fmt::format (
220- " Flow {}: packet {} is confirmed when timeout reashed ; no "
187+ " Flow {}: packet {} is confirmed when timeout reached ; no "
221188 " retransmit" ,
222189 m_id, data.packet_num ));
223190 return ;
@@ -247,7 +214,7 @@ class NewTcpFlow::Timeout : public Event {
247214 void operator ()() {
248215 if (m_flow.expired ()) {
249216 LOG_ERROR (
250- " Could not run TCP flow timout event: flow pointer expired" );
217+ " Could not run TCP flow timeout event: flow pointer expired" );
251218 return ;
252219 }
253220 m_flow.lock ()->on_timeout (std::move (m_packet));
0 commit comments