33#include < type_traits>
44
55#include " device/interfaces/i_host.hpp"
6- #include " event/generate.hpp"
76#include " flow/i_flow.hpp"
87#include " i_tcp_cc.hpp"
98#include " metrics/metrics_collector.hpp"
@@ -20,14 +19,12 @@ class TcpFlow : public IFlow,
2019public:
2120 TcpFlow (Id a_id, std::shared_ptr<IHost> a_src,
2221 std::shared_ptr<IHost> a_dest, TTcpCC a_cc, Size a_packet_size,
23- Time a_delay_between_packets, std::uint32_t a_packets_to_send,
24- bool a_ecn_capable = true )
22+ std::uint32_t a_packets_to_send, bool a_ecn_capable = true )
2523 : m_id(std::move(a_id)),
2624 m_src (a_src),
2725 m_dest(a_dest),
2826 m_cc(std::move(a_cc)),
2927 m_packet_size(a_packet_size),
30- m_delay_between_packets(a_delay_between_packets),
3128 m_packets_to_send(a_packets_to_send),
3229 m_ecn_capable(a_ecn_capable),
3330 m_packets_in_flight(0 ),
@@ -42,22 +39,7 @@ class TcpFlow : public IFlow,
4239 initialize_flag_manager ();
4340 }
4441
45- void start () final {
46- Time curr_time = Scheduler::get_instance ().get_current_time ();
47- Scheduler::get_instance ().add <Generate>(
48- curr_time, this ->shared_from_this (), m_packet_size);
49- }
50-
51- Time create_new_data_packet () final {
52- if (m_packets_to_send == 0 ) {
53- return 0 ;
54- }
55- if (try_to_put_data_to_device ()) {
56- --m_packets_to_send;
57- }
58-
59- return m_delay_between_packets;
60- }
42+ void start () final { send_packets (); }
6143
6244 void update (Packet packet, DeviceType type) final {
6345 (void )type;
@@ -84,15 +66,11 @@ class TcpFlow : public IFlow,
8466
8567 double old_cwnd = m_cc.get_cwnd ();
8668
87- if (m_cc.on_ack (rtt, packet.congestion_experienced )) {
88- // Trigger congestion
89- if (m_packets_in_flight > 0 ) {
90- m_packets_in_flight--;
91- }
92- } else {
93- if (m_packets_in_flight > 0 ) {
94- m_packets_in_flight--;
95- }
69+ if (m_packets_in_flight > 0 ) {
70+ m_packets_in_flight--;
71+ }
72+ if (!m_cc.on_ack (rtt, packet.congestion_experienced )) {
73+ // No congestion
9674 m_packets_acked++;
9775 }
9876
@@ -113,6 +91,7 @@ class TcpFlow : public IFlow,
11391 m_flag_manager.set_flag (packet, packet_type_label, PacketType::ACK );
11492 m_dest.lock ()->enqueue_packet (ack);
11593 }
94+ send_packets ();
11695 }
11796
11897 std::shared_ptr<IHost> get_sender () const final { return m_src.lock (); }
@@ -130,7 +109,6 @@ class TcpFlow : public IFlow,
130109 oss << " , CC module: " << m_cc.to_string ();
131110 oss << " , packet size: " << m_packet_size;
132111 oss << " , to send packets: " << m_packets_to_send;
133- oss << " , delay: " << m_delay_between_packets;
134112 oss << " , packets_in_flight: " << m_packets_in_flight;
135113 oss << " , acked packets: " << m_packets_acked;
136114 oss << " ]" ;
@@ -152,7 +130,7 @@ class TcpFlow : public IFlow,
152130 return ;
153131 }
154132 auto flow = m_flow.lock ();
155- flow->send_packet_now (m_packet);
133+ flow->send_packet_now (std::move ( m_packet) );
156134 }
157135
158136 private:
@@ -174,26 +152,34 @@ class TcpFlow : public IFlow,
174152 }
175153
176154 void send_packet_now (Packet packet) {
177- m_packets_in_flight++;
155+ // TODO: think about this place(should be here or in send_packets)
178156 m_sent_bytes += packet.size_byte ;
179157 packet.sent_time = Scheduler::get_instance ().get_current_time ();
180- m_src.lock ()->enqueue_packet (packet);
158+ m_src.lock ()->enqueue_packet (std::move ( packet) );
181159 }
182160
183- bool try_to_put_data_to_device () {
184- if (m_packets_in_flight < m_cc.get_cwnd ()) {
161+ // Send (ot plan sending) as many packets as possible
162+ void send_packets () {
163+ constexpr double EPS = 1e-6 ;
164+
165+ Time total_delay = 0 ;
166+ Time pacing_delay = m_cc.get_pacing_delay ();
167+ Time curr_time = Scheduler::get_instance ().get_current_time ();
168+
169+ while (m_packets_to_send > 0 &&
170+ m_packets_in_flight < m_cc.get_cwnd () + EPS ) {
185171 Packet packet = generate_packet ();
186- Time pacing_delay = m_cc. get_pacing_delay () ;
172+ total_delay += pacing_delay ;
187173 if (pacing_delay == 0 ) {
188- send_packet_now (packet);
174+ send_packet_now (std::move ( packet) );
189175 } else {
190- Time curr_time = Scheduler::get_instance ().get_current_time ();
191176 Scheduler::get_instance ().add <SendAtTime>(
192- curr_time + pacing_delay, this ->shared_from_this (), packet);
177+ curr_time + total_delay, this ->shared_from_this (),
178+ std::move (packet));
193179 }
194- return true ;
180+ m_packets_in_flight++;
181+ m_packets_to_send--;
195182 }
196- return false ;
197183 }
198184
199185 static void initialize_flag_manager () {
@@ -217,7 +203,6 @@ class TcpFlow : public IFlow,
217203 TTcpCC m_cc;
218204
219205 Size m_packet_size;
220- Time m_delay_between_packets;
221206 std::uint32_t m_packets_to_send;
222207 bool m_ecn_capable;
223208
0 commit comments