-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathbasic_flow.cpp
More file actions
129 lines (110 loc) · 4.54 KB
/
Copy pathbasic_flow.cpp
File metadata and controls
129 lines (110 loc) · 4.54 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
#include "basic_flow.hpp"
#include <spdlog/fmt/fmt.h>
#include "event/generate.hpp"
#include "logger/logger.hpp"
#include "metrics/metrics_collector.hpp"
#include "scheduler.hpp"
namespace sim {
std::string BasicFlow::packet_type_label = "type";
FlagManager<std::string, PacketFlagsBase> BasicFlow::m_flag_manager;
bool BasicFlow::m_is_flag_manager_initialized = false;
BasicFlow::BasicFlow(Id a_id, std::shared_ptr<IHost> a_src,
std::shared_ptr<IHost> a_dest, Size a_packet_size,
Time a_delay_between_packets,
std::uint32_t a_packets_to_send)
: m_id(a_id),
m_src(a_src),
m_dest(a_dest),
m_packet_size(a_packet_size),
m_delay_between_packets(a_delay_between_packets),
m_updates_number(0),
m_packets_to_send(a_packets_to_send),
m_sent_bytes(0) {
if (m_src.lock() == nullptr) {
throw std::invalid_argument("Sender for Flow is nullptr");
}
if (m_dest.lock() == nullptr) {
throw std::invalid_argument("Receiver for Flow is nullptr");
}
initialize_flag_manager();
}
void BasicFlow::start() {
schedule_packet_generation(Scheduler::get_instance().get_current_time());
}
Packet BasicFlow::generate_packet() {
sim::Packet packet;
m_flag_manager.set_flag(packet, packet_type_label, PacketType::DATA);
packet.size_byte = m_packet_size;
packet.flow = this;
packet.source_id = get_sender()->get_id();
packet.dest_id = get_receiver()->get_id();
packet.sent_time = Scheduler::get_instance().get_current_time();
return packet;
}
void BasicFlow::update(Packet packet, DeviceType type) {
(void)type;
if (packet.dest_id == m_dest.lock()->get_id() &&
m_flag_manager.get_flag(packet, packet_type_label) == PacketType::DATA) {
// data packet arrived to destination device, send ack
Packet ack(1, this, m_dest.lock()->get_id(),
m_src.lock()->get_id(), packet.sent_time,
packet.sent_bytes_at_origin, packet.ecn_capable_transport,
packet.congestion_experienced);
m_flag_manager.set_flag(ack, packet_type_label, PacketType::ACK);
m_dest.lock()->enqueue_packet(ack);
} else if (packet.dest_id == m_src.lock()->get_id() &&
m_flag_manager.get_flag(packet, packet_type_label) == PacketType::ACK) {
// ask arrived to source device, update metrics
++m_updates_number;
Time current_time = Scheduler::get_instance().get_current_time();
double rtt = current_time - packet.sent_time;
MetricsCollector::get_instance().add_RTT(packet.flow->get_id(),
current_time, rtt);
double delivery_bit_rate =
8 * (m_sent_bytes - packet.sent_bytes_at_origin) / rtt;
MetricsCollector::get_instance().add_delivery_rate(
packet.flow->get_id(), current_time, delivery_bit_rate);
} else {
LOG_ERROR(
fmt::format("Called update on flow {} with some foreign packet {}",
get_id(), packet.to_string()));
}
}
std::uint32_t BasicFlow::get_updates_number() const { return m_updates_number; }
Time BasicFlow::create_new_data_packet() {
if (m_packets_to_send == 0) {
return 0;
}
--m_packets_to_send;
Packet data = generate_packet();
// Note: sent_time and m_sent_bytes are evaluated at time of pushing
// the packet to the m_sending_buffer
m_sent_bytes += data.size_byte;
m_sending_buffer.push(data);
return put_data_to_device();
}
std::shared_ptr<IHost> BasicFlow::get_sender() const { return m_src.lock(); }
std::shared_ptr<IHost> BasicFlow::get_receiver() const { return m_dest.lock(); }
Id BasicFlow::get_id() const { return m_id; }
Time BasicFlow::put_data_to_device() {
if (m_src.expired()) {
LOG_ERROR("Flow source was deleted; can not put data to it");
return 0;
}
m_sending_buffer.front().sent_time =
Scheduler::get_instance().get_current_time();
m_src.lock()->enqueue_packet(m_sending_buffer.front());
m_sending_buffer.pop();
return m_delay_between_packets;
}
void BasicFlow::schedule_packet_generation(Time time) {
Scheduler::get_instance().add<Generate>(time, shared_from_this(),
m_packet_size);
}
void BasicFlow::initialize_flag_manager() {
if (!m_is_flag_manager_initialized) {
m_flag_manager.register_flag_by_amount(packet_type_label, PacketType::ENUM_SIZE);
m_is_flag_manager_initialized = true;
}
}
} // namespace sim