-
Notifications
You must be signed in to change notification settings - Fork 1
/
reliable_SR.cpp
218 lines (167 loc) · 5.88 KB
/
reliable_SR.cpp
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
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
#include <thread>
#include <condition_variable>
#include <utility>
#include <deque>
#include <set>
#include <functional>
#include "log.h"
#include "unreliable.h"
#include "reliable_SR.h"
const static auto waitTime = std::chrono::milliseconds(50);
const static uint32_t N = 3;
struct Task {
std::mutex m;
std::condition_variable cv;
bool ackReceived = false;
std::function<void(std::shared_ptr<Task>)> sender;
};
class WindowSR {
const uint32_t N;
uint32_t base;
uint32_t end;
std::deque<std::shared_ptr<Task>> queue;
std::mutex m;
std::condition_variable cvQueue;
public:
WindowSR(uint32_t base, uint32_t end, uint32_t N)
: base(base), end(end), N(N) {}
void push(std::shared_ptr<Task> task) {
std::unique_lock lock(m);
cvQueue.wait(lock, [this] { return queue.size() < N; });
queue.push_back(task);
LOG << "after push, queue size = " << queue.size() << std::endl;
std::thread senderThread(task->sender, task);
senderThread.detach();
}
void recvAck(uint32_t ack) {
std::lock_guard lock(m);
// invalid ack
if (ack < base || ack >= base + queue.size()) {
return;
}
// notify sender thread, so the sender thread stop sending packet and exit
const auto &currTask = queue.at(ack - base);
{
std::lock_guard taskLock(currTask->m);
currTask->ackReceived = true;
currTask->cv.notify_all();
}
// try to move window
uint32_t moving = 0;
for (const auto &task: queue) {
if (task == currTask) {
moving++;
break;
}
std::lock_guard taskLock(task->m);
if (task->ackReceived) {
moving++;
} else {
break;
}
}
if (moving > 0) {
for (uint32_t i = 0; i < moving; i++) {
LOG << "move window" << std::endl;
base++;
queue.pop_front();
}
LOG << "after move, queue size = " << queue.size() << std::endl;
cvQueue.notify_all();
}
}
};
ReliableSR::ReliableSR(Unreliable unreliable)
: unreliable(std::move(unreliable)) {}
bool ReliableSR::send(uint8_t *buf, int len) {
const int dataSize = MAX_PACKET_SIZE - sizeof(Packet);
uint32_t seq = 0;
uint32_t end = ROUND_UP(len, dataSize) / dataSize;
WindowSR window(seq, end, N);
std::set<uint32_t> allSeqs;
for (uint32_t i = seq; i < end; i++) {
allSeqs.insert(i);
}
std::thread ackReceiver([this, &window, &allSeqs] {
while (true) {
auto packet = unreliable.recv();
if (packet &&
PacketHelper::isValidPacket(packet) &&
packet->type == PacketType::ACK) {
LOG << "recveive ACK " << packet->num << std::endl;
window.recvAck(packet->num);
allSeqs.erase(packet->num);
if (allSeqs.empty()) {
break;
}
} else {
LOG << "invalid ACK" << std::endl;
}
}
LOG << "receive ACK thread exit" << std::endl;
});
for (uint8_t *sliceBuf = buf;
sliceBuf < buf + len;
sliceBuf += dataSize, seq++) {
int sliceLen = (std::min)(static_cast<int>(len - (sliceBuf - buf)), dataSize);
auto task = std::make_shared<Task>();
task->sender = [this, seq, sliceBuf, sliceLen](std::shared_ptr<Task> task) {
std::unique_lock lock(task->m);
do {
LOG << "sending slice " << seq << std::endl;
unreliable.send(PacketHelper::makePacket(
PacketType::DATA,
seq,
sliceBuf,
sliceLen
));
} while (!task->cv.wait_for(lock, waitTime,
[&] { return task->ackReceived; }));
LOG << "slice " << seq << " sent successfully" << std::endl;
};
window.push(task);
}
// waiting for received all ACKs
ackReceiver.join();
LOG << "sending FIN" << std::endl;
if (!unreliable.send(PacketHelper::makePacket(PacketType::FIN))) {
LOG << "failed to send FIN" << std::endl;
return false;
}
LOG << "receiving FIN_ACK" << std::endl;
std::unique_ptr<Packet> packet = unreliable.recv();
if (packet == nullptr ||
!PacketHelper::isValidPacket(packet) ||
packet->type != PacketType::FIN_ACK) {
return false;
}
LOG << "sent all slices successfully" << std::endl;
return true;
}
int ReliableSR::recv(uint8_t *buf, int len) {
const int dataSize = MAX_PACKET_SIZE - sizeof(Packet);
uint32_t recvSize = 0;
while (true) {
auto packet = unreliable.recv();
if (packet &&
PacketHelper::isValidPacket(packet) &&
packet->type == PacketType::DATA) {
LOG << "received slice " << packet->num << std::endl;
int sliceLen = packet->len - sizeof(Packet);
memcpy(buf + packet->num * dataSize, packet->data, sliceLen);
recvSize = (std::max)(recvSize, packet->num * dataSize + sliceLen);
LOG << "sending ACK " << packet->num << std::endl;
unreliable.send(PacketHelper::makePacket(PacketType::ACK, packet->num));
} else if (packet &&
PacketHelper::isValidPacket(packet) &&
packet->type == PacketType::FIN) {
LOG << "received FIN" << std::endl;
LOG << "sending FIN_ACK" << std::endl;
unreliable.send(PacketHelper::makePacket(PacketType::FIN_ACK));
break;
} else {
LOG << "received invalid packet" << std::endl;
}
}
return recvSize;
}