Diff
checker
टेक्स्ट
टेक्स्ट
छवियां
दस्तावेज़
Excel
फ़ोल्डर्स
Legal
Enterprise
डेस्कटॉप
मूल्य
साइन इन करें
Diffchecker डेस्कटॉप डाउनलोड करें
टेक्स्ट की तुलना करें
दो टेक्स्ट फ़ाइलों के बीच अंतर ढूंढें
उपकरण
इतिहास
रियल-टाइम एडिटर
रिक्त स्थान छिपाएँ
अपरिवर्तित संक्षिप्त करें
लाइन रैप बंद
लेआउट
विभाजित
संयुक्त
परिवर्तन हाइलाइट करें
स्मार्ट
शब्द
अक्षर
टेक्स्ट शैलियां
दिखावट बदलें
सिंटैक्स हाइलाइटिंग
सिंटैक्स चुनें
अनदेखा करें
टेक्स्ट बदलें
पहले अंतर पर जाएँ
इनपुट संपादित करें
Diffchecker Desktop
Diffchecker चलाने का सबसे सुरक्षित तरीका। Diffchecker Desktop ऐप पाएं: आपके diffs कभी आपके कंप्यूटर से बाहर नहीं जाते!
Desktop पाएं
EvolveDatacenterTCPDiff
बनाया गया
6 माह पहले
Diff कभी समाप्त नहीं होता
साफ़
निर्यात करें
शेयर करें
समझाएं
8 हटाए गए
लाइनें
कुल
हटाया गया
अक्षर
कुल
हटाया गया
इस सुविधा का उपयोग जारी रखने के लिए, अपग्रेड करें
Diff
checker
Pro
मूल्य देखें
383 लाइनें
सभी को कॉपी करें
19 जोड़े गए
लाइनें
कुल
जोड़ा गया
अक्षर
कुल
जोड़ा गया
इस सुविधा का उपयोग जारी रखने के लिए, अपग्रेड करें
Diff
checker
Pro
मूल्य देखें
389 लाइनें
सभी को कॉपी करें
#include "tcp-datacenter.h"
#include "tcp-datacenter.h"
namespace ns3 {
namespace ns3 {
void TcpDataCenter::AddQueuePair(uint64_t size, uint16_t pg, Ipv4Address sip, Ipv4Address dip, uint16_t sport, uint16_t dport, uint32_t win, uint64_t baseRtt, Callback<void> notifyAppFinish, Time stopTime) {
void TcpDataCenter::AddQueuePair(uint64_t size, uint16_t pg, Ipv4Address sip, Ipv4Address dip, uint16_t sport, uint16_t dport, uint32_t win, uint64_t baseRtt, Callback<void> notifyAppFinish, Time stopTime) {
Ptr<RdmaQueuePair> qp = CreateObject<RdmaQueuePair>(pg, sip, dip, sport, dport);
Ptr<RdmaQueuePair> qp = CreateObject<RdmaQueuePair>(pg, sip, dip, sport, dport);
qp->SetSize(size);
qp->SetSize(size);
qp->SetWin(win);
qp->SetWin(win);
qp->SetBaseRtt(baseRtt);
qp->SetBaseRtt(baseRtt);
qp->SetVarWin(m_var_win);
qp->SetVarWin(m_var_win);
qp->SetAppNotifyCallback(notifyAppFinish);
qp->SetAppNotifyCallback(notifyAppFinish);
qp->stopTime = stopTime;
qp->stopTime = stopTime;
uint32_t nic_idx = GetNicIdxOfQp(qp);
uint32_t nic_idx = GetNicIdxOfQp(qp);
m_nic[nic_idx].qpGrp->AddQp(qp);
m_nic[nic_idx].qpGrp->AddQp(qp);
uint64_t key = GetQpKey(dip.Get(), sport, pg);
uint64_t key = GetQpKey(dip.Get(), sport, pg);
m_qpMap[key] = qp;
m_qpMap[key] = qp;
qp->powerEnabled = true;
qp->powerEnabled = true;
if (m_nic[nic_idx].dev == NULL) {
if (m_nic[nic_idx].dev == NULL) {
std::cout << "sip " << sip << " dip " << dip << " sport " << sport << " dport " << dport << std::endl;
std::cout << "sip " << sip << " dip " << dip << " sport " << sport << " dport " << dport << std::endl;
}
}
DataRate m_bps = m_nic[nic_idx].dev->GetDataRate();
DataRate m_bps = m_nic[nic_idx].dev->GetDataRate();
if(win)
if(win)
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
qp->SetWin(m_bps.GetBitRate() * 1
* baseRtt * 1e-9 / 8);
qp->SetWin(m_bps.GetBitRate() * 1
.0
* baseRtt * 1e-9 / 8);
qp->m_rate =
m_bps
;
qp->m_rate =
DataRate(
m_bps
.GetBitRate() / 10)
;
qp->m_max_rate = m_bps;
qp->m_max_rate = m_bps;
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
qp->hp.m_curRate =
m_bps
;
qp->hp.m_curRate =
DataRate(
m_bps
.GetBitRate() / 10)
;
m_nic[nic_idx].dev->NewQp(qp);
m_nic[nic_idx].dev->NewQp(qp);
}
}
int TcpDataCenter::ReceiveCnp(Ptr<Packet> p, CustomHeader &ch) {
int TcpDataCenter::ReceiveCnp(Ptr<Packet> p, CustomHeader &ch) {
uint32_t qIndex = ch.cnp.qIndex;
uint32_t qIndex = ch.cnp.qIndex;
if (qIndex == 1) {
if (qIndex == 1) {
std::cout << "TCP--ignore\n";
std::cout << "TCP--ignore\n";
return 0;
return 0;
}
}
uint16_t udpport = ch.cnp.fid;
uint16_t udpport = ch.cnp.fid;
Ptr<RdmaQueuePair> qp = GetQp(ch.sip, udpport, qIndex);
Ptr<RdmaQueuePair> qp = GetQp(ch.sip, udpport, qIndex);
if (qp == NULL)
if (qp == NULL)
std::cout << "ERROR: QCN NIC cannot find the flow\n";
std::cout << "ERROR: QCN NIC cannot find the flow\n";
uint32_t nic_idx = GetNicIdxOfQp(qp);
uint32_t nic_idx = GetNicIdxOfQp(qp);
Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev;
Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev;
if (qp->m_rate == 0)
if (qp->m_rate == 0)
{
{
qp->m_rate = dev->GetDataRate();
qp->m_rate = dev->GetDataRate();
qp->hp.m_curRate = dev->GetDataRate();
qp->hp.m_curRate = dev->GetDataRate();
}
}
return 0;
return 0;
}
}
int TcpDataCenter::ReceiveAck(Ptr<Packet> p, CustomHeader &ch) {
int TcpDataCenter::ReceiveAck(Ptr<Packet> p, CustomHeader &ch) {
uint16_t qIndex = ch.ack.pg;
uint16_t qIndex = ch.ack.pg;
uint16_t port = ch.ack.dport;
uint16_t port = ch.ack.dport;
uint32_t seq = ch.ack.seq;
uint32_t seq = ch.ack.seq;
uint8_t cnp = (ch.ack.flags >> qbbHeader::FLAG_CNP) & 1;
uint8_t cnp = (ch.ack.flags >> qbbHeader::FLAG_CNP) & 1;
Ptr<RdmaQueuePair> qp = GetQp(ch.sip, port, qIndex);
Ptr<RdmaQueuePair> qp = GetQp(ch.sip, port, qIndex);
if (qp == NULL) {
if (qp == NULL) {
std::cout << "ERROR: " << "node:" << m_node->GetId() << ' ' << (ch.l3Prot == 0xFC ? "ACK" : "NACK") << " NIC cannot find the flow\n";
std::cout << "ERROR: " << "node:" << m_node->GetId() << ' ' << (ch.l3Prot == 0xFC ? "ACK" : "NACK") << " NIC cannot find the flow\n";
return 0;
return 0;
}
}
uint32_t nic_idx = GetNicIdxOfQp(qp);
uint32_t nic_idx = GetNicIdxOfQp(qp);
Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev;
Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev;
if (m_ack_interval == 0)
if (m_ack_interval == 0)
std::cout << "ERROR: shouldn't receive ack\n";
std::cout << "ERROR: shouldn't receive ack\n";
else {
else {
if (!m_backto0) {
if (!m_backto0) {
qp->Acknowledge(seq);
qp->Acknowledge(seq);
} else {
} else {
uint32_t goback_seq = seq / m_chunk * m_chunk;
uint32_t goback_seq = seq / m_chunk * m_chunk;
qp->Acknowledge(goback_seq);
qp->Acknowledge(goback_seq);
}
}
if (qp->IsFinished()) {
if (qp->IsFinished()) {
QpComplete(qp);
QpComplete(qp);
}
}
}
}
if (ch.l3Prot == 0xFD)
if (ch.l3Prot == 0xFD)
RecoverQueue(qp);
RecoverQueue(qp);
HandleAckHp(qp, p, ch);
HandleAckHp(qp, p, ch);
dev->TriggerTransmit();
dev->TriggerTransmit();
return 0;
return 0;
}
}
#define PRINT_LOG 0
#define PRINT_LOG 0
void TcpDataCenter::HandleAckHp(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) {
void TcpDataCenter::HandleAckHp(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) {
uint32_t ack_seq = ch.ack.seq;
uint32_t ack_seq = ch.ack.seq;
if (ack_seq > qp->hp.m_lastUpdateSeq) {
if (ack_seq > qp->hp.m_lastUpdateSeq) {
UpdateRatePower(qp, p, ch, false);
UpdateRatePower(qp, p, ch, false);
} else {
} else {
FastReactPower(qp, p, ch);
FastReactPower(qp, p, ch);
}
}
}
}
void TcpDataCenter::UpdateRatePower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch, bool fast_react) {
void TcpDataCenter::UpdateRatePower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch, bool fast_react) {
uint32_t next_seq = qp->snd_nxt;
uint32_t next_seq = qp->snd_nxt;
double prevRtt = qp->m_baseRtt;
double prevRtt = qp->m_baseRtt;
double prevCompletion = Simulator::Now().GetNanoSeconds();
double prevCompletion = Simulator::Now().GetNanoSeconds();
std::map<uint32_t, double>::iterator it = qp->rates.find(ch.ack.seq);
std::map<uint32_t, double>::iterator it = qp->rates.find(ch.ack.seq);
if (it != qp->rates.end()) {
if (it != qp->rates.end()) {
qp->rates.erase(it);
qp->rates.erase(it);
prevRtt = Simulator::Now().GetNanoSeconds() - it->second;
prevRtt = Simulator::Now().GetNanoSeconds() - it->second;
qp->m_baseRtt = std::min(uint64_t(Simulator::Now().GetNanoSeconds() - it->second), qp->m_baseRtt);
qp->m_baseRtt = std::min(uint64_t(Simulator::Now().GetNanoSeconds() - it->second), qp->m_baseRtt);
prevCompletion = Simulator::Now().GetNanoSeconds();
prevCompletion = Simulator::Now().GetNanoSeconds();
}
}
if (qp->hp.m_lastUpdateSeq == 0) {
if (qp->hp.m_lastUpdateSeq == 0) {
qp->prevRtt = prevRtt;
qp->prevRtt = prevRtt;
qp->prevCompletion = Simulator::Now().GetNanoSeconds();
qp->prevCompletion = Simulator::Now().GetNanoSeconds();
qp->hp.m_lastUpdateSeq = next_seq;
qp->hp.m_lastUpdateSeq = next_seq;
}else {
}else {
double max_c = 0;
double max_c = 0;
double U = 0;
double U = 0;
uint64_t dt = 0;
uint64_t dt = 0;
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
double A = ( double(prevRtt - qp->prevRtt) /
(prevCompletion - qp->prevCompletion)
+ 1 );
double dt_sample = prevCompletion - qp->prevCompletion;
if (A < 0.5)
if (dt_sample < 1.0) dt_sample = 1.0;
A = 0.5;
double A = ( double(prevRtt - qp->prevRtt) /
dt_sample
+ 1 );
if (A < 0.5)
A = 0.5;
double power = ( A ) * (prevRtt);
double power = ( A ) * (prevRtt);
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
double powerx = (power) / (1.
05
* qp->m_baseRtt);
double powerx = (power) / (1.
03
* qp->m_baseRtt);
double u = powerx;
double u = powerx;
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
uint8_t cnp = (ch.ack.flags >> qbbHeader::FLAG_CNP) & 1;
if (cnp) {
if (u < 2.0) u = 2.0;
}
if (u > U) {
if (u > U) {
U = u;
U = u;
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
dt =
prevCompletion - qp->prevCompletion
;
dt =
(uint64_t)dt_sample
;
}
}
DataRate new_rate;
DataRate new_rate;
if (dt > 1.0 * qp->m_baseRtt)
if (dt > 1.0 * qp->m_baseRtt)
dt = 1.0 * qp->m_baseRtt;
dt = 1.0 * qp->m_baseRtt;
if (U < 0) {
if (U < 0) {
U = qp->hp.u;
U = qp->hp.u;
}
}
qp->hp.u = (qp->hp.u * (1.0 * qp->m_baseRtt - dt) + U * dt) / double(1.0 * qp->m_baseRtt);
qp->hp.u = (qp->hp.u * (1.0 * qp->m_baseRtt - dt) + U * dt) / double(1.0 * qp->m_baseRtt);
max_c = qp->hp.u;
max_c = qp->hp.u;
कॉपी
कॉपी हुआ
कॉपी
कॉपी हुआ
new_rate = (0.
7
* ( qp->hp.m_curRate / max_c + DataRate("
1
50Mbps") ) + 0.
3
* qp->hp.m_curRate);
new_rate = (0.
8
* ( qp->hp.m_curRate / max_c + DataRate("
2
50Mbps") ) + 0.
2
* qp->hp.m_curRate);
if (new_rate < m_minRate)
if (new_rate < m_minRate)
new_rate = m_minRate;
new_rate = m_minRate;
if (new_rate > qp->m_max_rate)
if (new_rate > qp->m_max_rate)
new_rate = qp->m_max_rate;
new_rate = qp->m_max_rate;
qp->prevRtt = prevRtt;
qp->prevRtt = prevRtt;
qp->prevCompletion = Simulator::Now().GetNanoSeconds();
qp->prevCompletion = Simulator::Now().GetNanoSeconds();
ChangeRate(qp, new_rate);
ChangeRate(qp, new_rate);
if (!fast_react) {
if (!fast_react) {
qp->hp.m_curRate = new_rate;
qp->hp.m_curRate = new_rate;
}
}
if (!fast_react) {
if (!fast_react) {
if (next_seq > qp->hp.m_lastUpdateSeq)
if (next_seq > qp->hp.m_lastUpdateSeq)
qp->hp.m_lastUpdateSeq = next_seq;
qp->hp.m_lastUpdateSeq = next_seq;
}
}
}
}
}
}
void TcpDataCenter::FastReactPower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) {
void TcpDataCenter::FastReactPower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) {
if (m_fast_react)
if (m_fast_react)
UpdateRatePower(qp, p, ch, true);
UpdateRatePower(qp, p, ch, true);
}
}
void TcpDataCenter::Setup(QpCompleteCallback cb) {
void TcpDataCenter::Setup(QpCompleteCallback cb) {
for (uint32_t i = 0; i < m_nic.size(); i++) {
for (uint32_t i = 0; i < m_nic.size(); i++) {
Ptr<QbbNetDevice> dev = m_nic[i].dev;
Ptr<QbbNetDevice> dev = m_nic[i].dev;
if (dev == NULL)
if (dev == NULL)
continue;
continue;
dev->m_rdmaEQ->m_qpGrp = m_nic[i].qpGrp;
dev->m_rdmaEQ->m_qpGrp = m_nic[i].qpGrp;
dev->m_rdmaReceiveCb = MakeCallback(&TcpDataCenter::Receive, this);
dev->m_rdmaReceiveCb = MakeCallback(&TcpDataCenter::Receive, this);
dev->m_rdmaLinkDownCb = MakeCallback(&TcpDataCenter::SetLinkDown, this);
dev->m_rdmaLinkDownCb = MakeCallback(&TcpDataCenter::SetLinkDown, this);
dev->m_rdmaPktSent = MakeCallback(&TcpDataCenter::PktSent, this);
dev->m_rdmaPktSent = MakeCallback(&TcpDataCenter::PktSent, this);
dev->m_rdmaEQ->m_rdmaGetNxtPkt = MakeCallback(&TcpDataCenter::GetNxtPacket, this);
dev->m_rdmaEQ->m_rdmaGetNxtPkt = MakeCallback(&TcpDataCenter::GetNxtPacket, this);
}
}
m_qpCompleteCallback = cb;
m_qpCompleteCallback = cb;
}
}
void TcpDataCenter::DeleteQueuePair(Ptr<RdmaQueuePair> qp) {
void TcpDataCenter::DeleteQueuePair(Ptr<RdmaQueuePair> qp) {
uint64_t key = GetQpKey(qp->dip.Get(), qp->sport, qp->m_pg);
uint64_t key = GetQpKey(qp->dip.Get(), qp->sport, qp->m_pg);
m_qpMap.erase(key);
m_qpMap.erase(key);
}
}
Ptr<RdmaRxQueuePair> TcpDataCenter::GetRxQp(uint32_t sip, uint32_t dip, uint16_t sport, uint16_t dport, uint16_t pg, bool create) {
Ptr<RdmaRxQueuePair> TcpDataCenter::GetRxQp(uint32_t sip, uint32_t dip, uint16_t sport, uint16_t dport, uint16_t pg, bool create) {
uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport;
uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport;
auto it = m_rxQpMap.find(key);
auto it = m_rxQpMap.find(key);
if (it != m_rxQpMap.end())
if (it != m_rxQpMap.end())
return it->second;
return it->second;
if (create) {
if (create) {
Ptr<RdmaRxQueuePair> q = CreateObject<RdmaRxQueuePair>();
Ptr<RdmaRxQueuePair> q = CreateObject<RdmaRxQueuePair>();
q->sip = sip;
q->sip = sip;
q->dip = dip;
q->dip = dip;
q->sport = sport;
q->sport = sport;
q->dport = dport;
q->dport = dport;
q->m_ecn_source.qIndex = pg;
q->m_ecn_source.qIndex = pg;
m_rxQpMap[key] = q;
m_rxQpMap[key] = q;
return q;
return q;
}
}
return NULL;
return NULL;
}
}
void TcpDataCenter::DeleteRxQp(uint32_t dip, uint16_t pg, uint16_t dport) {
void TcpDataCenter::DeleteRxQp(uint32_t dip, uint16_t pg, uint16_t dport) {
uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport;
uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport;
m_rxQpMap.erase(key);
m_rxQpMap.erase(key);
}
}
int TcpDataCenter::ReceiveUdp(Ptr<Packet> p, CustomHeader &ch) {
int TcpDataCenter::ReceiveUdp(Ptr<Packet> p, CustomHeader &ch) {
uint8_t ecnbits = ch.GetIpv4EcnBits();
uint8_t ecnbits = ch.GetIpv4EcnBits();
uint32_t payload_size = p->GetSize() - ch.GetSerializedSize();
uint32_t payload_size = p->GetSize() - ch.GetSerializedSize();
Ptr<RdmaRxQueuePair> rxQp = GetRxQp(ch.dip, ch.sip, ch.udp.dport, ch.udp.sport, ch.udp.pg, true);
Ptr<RdmaRxQueuePair> rxQp = GetRxQp(ch.dip, ch.sip, ch.udp.dport, ch.udp.sport, ch.udp.pg, true);
if (ecnbits != 0) {
if (ecnbits != 0) {
rxQp->m_ecn_source.ecnbits |= ecnbits;
rxQp->m_ecn_source.ecnbits |= ecnbits;
rxQp->m_ecn_source.qfb++;
rxQp->m_ecn_source.qfb++;
}
}
rxQp->m_ecn_source.total++;
rxQp->m_ecn_source.total++;
rxQp->m_milestone_rx = m_ack_interval;
rxQp->m_milestone_rx = m_ack_interval;
int x = ReceiverCheckSeq(ch.udp.seq, rxQp, payload_size);
int x = ReceiverCheckSeq(ch.udp.seq, rxQp, payload_size);
if (x == 1 || x == 2) {
if (x == 1 || x == 2) {
qbbHeader seqh;
qbbHeader seqh;
seqh.SetSeq(rxQp->ReceiverNextExpectedSeq);
seqh.SetSeq(rxQp->ReceiverNextExpectedSeq);
seqh.SetPG(ch.udp.pg);
seqh.SetPG(ch.udp.pg);
seqh.SetSport(ch.udp.dport);
seqh.SetSport(ch.udp.dport);
seqh.SetDport(ch.udp.sport);
seqh.SetDport(ch.udp.sport);
if (ecnbits)
if (ecnbits)
seqh.SetCnp();
seqh.SetCnp();
Ptr<Packet> newp = Create<Packet>(std::max(60 - 14 - 20 - (int)seqh.GetSerializedSize(), 0));
Ptr<Packet> newp = Create<Packet>(std::max(60 - 14 - 20 - (int)seqh.GetSerializedSize(), 0));
newp->AddHeader(seqh);
newp->AddHeader(seqh);
Ipv4Header head;
Ipv4Header head;
head.SetDestination(Ipv4Address(ch.sip));
head.SetDestination(Ipv4Address(ch.sip));
head.SetSource(Ipv4Address(ch.dip));
head.SetSource(Ipv4Address(ch.dip));
head.SetProtocol(x == 1 ? 0xFC : 0xFD);
head.SetProtocol(x == 1 ? 0xFC : 0xFD);
head.SetTtl(64);
head.SetTtl(64);
head.SetPayloadSize(newp->GetSize());
head.SetPayloadSize(newp->GetSize());
head.SetIdentification(rxQp->m_ipid++);
head.SetIdentification(rxQp->m_ipid++);
newp->AddHeader(head);
newp->AddHeader(head);
AddHeader(newp, 0x800);
AddHeader(newp, 0x800);
uint32_t nic_idx = GetNicIdxOfRxQp(rxQp);
uint32_t nic_idx = GetNicIdxOfRxQp(rxQp);
m_nic[nic_idx].dev->RdmaEnqueueHighPrioQ(newp);
m_nic[nic_idx].dev->RdmaEnqueueHighPrioQ(newp);
m_nic[nic_idx].dev->TriggerTransmit();
m_nic[nic_idx].dev->TriggerTransmit();
}
}
return 0;
return 0;
}
}
int TcpDataCenter::Receive(Ptr<Packet> p, CustomHeader &ch) {
int TcpDataCenter::Receive(Ptr<Packet> p, CustomHeader &ch) {
if (ch.l3Prot == 0x11) {
if (ch.l3Prot == 0x11) {
ReceiveUdp(p, ch);
ReceiveUdp(p, ch);
} else if (ch.l3Prot == 0xFF) {
} else if (ch.l3Prot == 0xFF) {
ReceiveCnp(p, ch);
ReceiveCnp(p, ch);
} else if (ch.l3Prot == 0xFD) {
} else if (ch.l3Prot == 0xFD) {
ReceiveAck(p, ch);
ReceiveAck(p, ch);
} else if (ch.l3Prot == 0xFC) {
} else if (ch.l3Prot == 0xFC) {
ReceiveAck(p, ch);
ReceiveAck(p, ch);
}
}
return 0;
return 0;
}
}
void TcpDataCenter::RecoverQueue(Ptr<RdmaQueuePair> qp) {
void TcpDataCenter::RecoverQueue(Ptr<RdmaQueuePair> qp) {
qp->snd_nxt = qp->snd_una;
qp->snd_nxt = qp->snd_una;
}
}
void TcpDataCenter::QpComplete(Ptr<RdmaQueuePair> qp) {
void TcpDataCenter::QpComplete(Ptr<RdmaQueuePair> qp) {
NS_ASSERT(!m_qpCompleteCallback.IsNull());
NS_ASSERT(!m_qpCompleteCallback.IsNull());
if (m_cc_mode == 1) {
if (m_cc_mode == 1) {
Simulator::Cancel(qp->mlx.m_eventUpdateAlpha);
Simulator::Cancel(qp->mlx.m_eventUpdateAlpha);
Simulator::Cancel(qp->mlx.m_eventDecreaseRate);
Simulator::Cancel(qp->mlx.m_eventDecreaseRate);
Simulator::Cancel(qp->mlx.m_rpTimer);
Simulator::Cancel(qp->mlx.m_rpTimer);
}
}
m_qpCompleteCallback(qp);
m_qpCompleteCallback(qp);
qp->m_notifyAppFinish();
qp->m_notifyAppFinish();
DeleteQueuePair(qp);
DeleteQueuePair(qp);
}
}
void TcpDataCenter::SetLinkDown(Ptr<QbbNetDevice> dev) {
void TcpDataCenter::SetLinkDown(Ptr<QbbNetDevice> dev) {
printf("RdmaHw: node:%u a link down\n", m_node->GetId());
printf("RdmaHw: node:%u a link down\n", m_node->GetId());
}
}
void TcpDataCenter::AddTableEntry(Ipv4Address &dstAddr, uint32_t intf_idx) {
void TcpDataCenter::AddTableEntry(Ipv4Address &dstAddr, uint32_t intf_idx) {
uint32_t dip = dstAddr.Get();
uint32_t dip = dstAddr.Get();
m_rtTable[dip].push_back(intf_idx);
m_rtTable[dip].push_back(intf_idx);
}
}
void TcpDataCenter::ClearTable() {
void TcpDataCenter::ClearTable() {
m_rtTable.clear();
m_rtTable.clear();
}
}
void TcpDataCenter::RedistributeQp() {
void TcpDataCenter::RedistributeQp() {
for (uint32_t i = 0; i < m_nic.size(); i++) {
for (uint32_t i = 0; i < m_nic.size(); i++) {
if (m_nic[i].dev == NULL)
if (m_nic[i].dev == NULL)
continue;
continue;
m_nic[i].qpGrp->Clear();
m_nic[i].qpGrp->Clear();
}
}
for (auto &it : m_qpMap) {
for (auto &it : m_qpMap) {
Ptr<RdmaQueuePair> qp = it.second;
Ptr<RdmaQueuePair> qp = it.second;
uint32_t nic_idx = GetNicIdxOfQp(qp);
uint32_t nic_idx = GetNicIdxOfQp(qp);
m_nic[nic_idx].qpGrp->AddQp(qp);
m_nic[nic_idx].qpGrp->AddQp(qp);
m_nic[nic_idx].dev->ReassignedQp(qp);
m_nic[nic_idx].dev->ReassignedQp(qp);
}
}
}
}
Ptr<Packet> TcpDataCenter::GetNxtPacket(Ptr<RdmaQueuePair> qp) {
Ptr<Packet> TcpDataCenter::GetNxtPacket(Ptr<RdmaQueuePair> qp) {
uint32_t payload_size = qp->GetBytesLeft();
uint32_t payload_size = qp->GetBytesLeft();
if (m_mtu < payload_size)
if (m_mtu < payload_size)
payload_size = m_mtu;
payload_size = m_mtu;
Ptr<Packet> p = Create<Packet> (payload_size);
Ptr<Packet> p = Create<Packet> (payload_size);
SeqTsHeader seqTs;
SeqTsHeader seqTs;
seqTs.SetSeq (qp->snd_nxt);
seqTs.SetSeq (qp->snd_nxt);
seqTs.SetPG (qp->m_pg);
seqTs.SetPG (qp->m_pg);
p->AddHeader (seqTs);
p->AddHeader (seqTs);
UdpHeader udpHeader;
UdpHeader udpHeader;
udpHeader.SetDestinationPort (qp->dport);
udpHeader.SetDestinationPort (qp->dport);
udpHeader.SetSourcePort (qp->sport);
udpHeader.SetSourcePort (qp->sport);
p->AddHeader (udpHeader);
p->AddHeader (udpHeader);
Ipv4Header ipHeader;
Ipv4Header ipHeader;
ipHeader.SetSource (qp->sip);
ipHeader.SetSource (qp->sip);
ipHeader.SetDestination (qp->dip);
ipHeader.SetDestination (qp->dip);
ipHeader.SetProtocol (0x11);
ipHeader.SetProtocol (0x11);
ipHeader.SetPayloadSize (p->GetSize());
ipHeader.SetPayloadSize (p->GetSize());
ipHeader.SetTtl (64);
ipHeader.SetTtl (64);
ipHeader.SetTos (0);
ipHeader.SetTos (0);
ipHeader.SetIdentification (qp->m_ipid);
ipHeader.SetIdentification (qp->m_ipid);
p->AddHeader(ipHeader);
p->AddHeader(ipHeader);
PppHeader ppp;
PppHeader ppp;
ppp.SetProtocol (0x0021);
ppp.SetProtocol (0x0021);
p->AddHeader (ppp);
p->AddHeader (ppp);
qp->snd_nxt += payload_size;
qp->snd_nxt += payload_size;
qp->m_ipid++;
qp->m_ipid++;
return p;
return p;
}
}
void TcpDataCenter::PktSent(Ptr<RdmaQueuePair> qp, Ptr<Packet> pkt, Time interframeGap) {
void TcpDataCenter::PktSent(Ptr<RdmaQueuePair> qp, Ptr<Packet> pkt, Time interframeGap) {
qp->lastPktSize = pkt->GetSize();
qp->lastPktSize = pkt->GetSize();
uint32_t seq = qp->snd_nxt;
uint32_t seq = qp->snd_nxt;
qp->rates[qp->snd_nxt] = Simulator::Now().GetNanoSeconds();
qp->rates[qp->snd_nxt] = Simulator::Now().GetNanoSeconds();
UpdateNextAvail(qp, interframeGap, pkt->GetSize());
UpdateNextAvail(qp, interframeGap, pkt->GetSize());
}
}
void TcpDataCenter::UpdateNextAvail(Ptr<RdmaQueuePair> qp, Time interframeGap, uint32_t pkt_size) {
void TcpDataCenter::UpdateNextAvail(Ptr<RdmaQueuePair> qp, Time interframeGap, uint32_t pkt_size) {
Time sendingTime;
Time sendingTime;
if (m_rateBound)
if (m_rateBound)
sendingTime = interframeGap + qp->m_rate.CalculateBytesTxTime(pkt_size);
sendingTime = interframeGap + qp->m_rate.CalculateBytesTxTime(pkt_size);
else
else
sendingTime = interframeGap + qp->m_max_rate.CalculateBytesTxTime(pkt_size);
sendingTime = interframeGap + qp->m_max_rate.CalculateBytesTxTime(pkt_size);
qp->m_nextAvail = Simulator::Now() + sendingTime;
qp->m_nextAvail = Simulator::Now() + sendingTime;
}
}
void TcpDataCenter::ChangeRate(Ptr<RdmaQueuePair> qp, DataRate new_rate) {
void TcpDataCenter::ChangeRate(Ptr<RdmaQueuePair> qp, DataRate new_rate) {
#if 1
#if 1
Time sendingTime = qp->m_rate.CalculateBytesTxTime(qp->lastPktSize);
Time sendingTime = qp->m_rate.CalculateBytesTxTime(qp->lastPktSize);
Time new_sendintTime = new_rate.CalculateBytesTxTime(qp->lastPktSize);
Time new_sendintTime = new_rate.CalculateBytesTxTime(qp->lastPktSize);
qp->m_nextAvail = qp->m_nextAvail + new_sendintTime - sendingTime;
qp->m_nextAvail = qp->m_nextAvail + new_sendintTime - sendingTime;
uint32_t nic_idx = GetNicIdxOfQp(qp);
uint32_t nic_idx = GetNicIdxOfQp(qp);
m_nic[nic_idx].dev->UpdateNextAvail(qp->m_nextAvail);
m_nic[nic_idx].dev->UpdateNextAvail(qp->m_nextAvail);
#endif
#endif
qp->m_rate = new_rate;
qp->m_rate = new_rate;
}
}
}
}
सेव किए गए Diffs
ऑरिजनल टेक्स्ट
फ़ाइल खोलें
#include "tcp-datacenter.h" namespace ns3 { void TcpDataCenter::AddQueuePair(uint64_t size, uint16_t pg, Ipv4Address sip, Ipv4Address dip, uint16_t sport, uint16_t dport, uint32_t win, uint64_t baseRtt, Callback<void> notifyAppFinish, Time stopTime) { Ptr<RdmaQueuePair> qp = CreateObject<RdmaQueuePair>(pg, sip, dip, sport, dport); qp->SetSize(size); qp->SetWin(win); qp->SetBaseRtt(baseRtt); qp->SetVarWin(m_var_win); qp->SetAppNotifyCallback(notifyAppFinish); qp->stopTime = stopTime; uint32_t nic_idx = GetNicIdxOfQp(qp); m_nic[nic_idx].qpGrp->AddQp(qp); uint64_t key = GetQpKey(dip.Get(), sport, pg); m_qpMap[key] = qp; qp->powerEnabled = true; if (m_nic[nic_idx].dev == NULL) { std::cout << "sip " << sip << " dip " << dip << " sport " << sport << " dport " << dport << std::endl; } DataRate m_bps = m_nic[nic_idx].dev->GetDataRate(); if(win) qp->SetWin(m_bps.GetBitRate() * 1 * baseRtt * 1e-9 / 8); qp->m_rate = m_bps; qp->m_max_rate = m_bps; qp->hp.m_curRate = m_bps; m_nic[nic_idx].dev->NewQp(qp); } int TcpDataCenter::ReceiveCnp(Ptr<Packet> p, CustomHeader &ch) { uint32_t qIndex = ch.cnp.qIndex; if (qIndex == 1) { std::cout << "TCP--ignore\n"; return 0; } uint16_t udpport = ch.cnp.fid; Ptr<RdmaQueuePair> qp = GetQp(ch.sip, udpport, qIndex); if (qp == NULL) std::cout << "ERROR: QCN NIC cannot find the flow\n"; uint32_t nic_idx = GetNicIdxOfQp(qp); Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev; if (qp->m_rate == 0) { qp->m_rate = dev->GetDataRate(); qp->hp.m_curRate = dev->GetDataRate(); } return 0; } int TcpDataCenter::ReceiveAck(Ptr<Packet> p, CustomHeader &ch) { uint16_t qIndex = ch.ack.pg; uint16_t port = ch.ack.dport; uint32_t seq = ch.ack.seq; uint8_t cnp = (ch.ack.flags >> qbbHeader::FLAG_CNP) & 1; Ptr<RdmaQueuePair> qp = GetQp(ch.sip, port, qIndex); if (qp == NULL) { std::cout << "ERROR: " << "node:" << m_node->GetId() << ' ' << (ch.l3Prot == 0xFC ? "ACK" : "NACK") << " NIC cannot find the flow\n"; return 0; } uint32_t nic_idx = GetNicIdxOfQp(qp); Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev; if (m_ack_interval == 0) std::cout << "ERROR: shouldn't receive ack\n"; else { if (!m_backto0) { qp->Acknowledge(seq); } else { uint32_t goback_seq = seq / m_chunk * m_chunk; qp->Acknowledge(goback_seq); } if (qp->IsFinished()) { QpComplete(qp); } } if (ch.l3Prot == 0xFD) RecoverQueue(qp); HandleAckHp(qp, p, ch); dev->TriggerTransmit(); return 0; } #define PRINT_LOG 0 void TcpDataCenter::HandleAckHp(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) { uint32_t ack_seq = ch.ack.seq; if (ack_seq > qp->hp.m_lastUpdateSeq) { UpdateRatePower(qp, p, ch, false); } else { FastReactPower(qp, p, ch); } } void TcpDataCenter::UpdateRatePower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch, bool fast_react) { uint32_t next_seq = qp->snd_nxt; double prevRtt = qp->m_baseRtt; double prevCompletion = Simulator::Now().GetNanoSeconds(); std::map<uint32_t, double>::iterator it = qp->rates.find(ch.ack.seq); if (it != qp->rates.end()) { qp->rates.erase(it); prevRtt = Simulator::Now().GetNanoSeconds() - it->second; qp->m_baseRtt = std::min(uint64_t(Simulator::Now().GetNanoSeconds() - it->second), qp->m_baseRtt); prevCompletion = Simulator::Now().GetNanoSeconds(); } if (qp->hp.m_lastUpdateSeq == 0) { qp->prevRtt = prevRtt; qp->prevCompletion = Simulator::Now().GetNanoSeconds(); qp->hp.m_lastUpdateSeq = next_seq; }else { double max_c = 0; double U = 0; uint64_t dt = 0; double A = ( double(prevRtt - qp->prevRtt) / (prevCompletion - qp->prevCompletion) + 1 ); if (A < 0.5) A = 0.5; double power = ( A ) * (prevRtt); double powerx = (power) / (1.05 * qp->m_baseRtt); double u = powerx; if (u > U) { U = u; dt = prevCompletion - qp->prevCompletion; } DataRate new_rate; if (dt > 1.0 * qp->m_baseRtt) dt = 1.0 * qp->m_baseRtt; if (U < 0) { U = qp->hp.u; } qp->hp.u = (qp->hp.u * (1.0 * qp->m_baseRtt - dt) + U * dt) / double(1.0 * qp->m_baseRtt); max_c = qp->hp.u; new_rate = (0.7 * ( qp->hp.m_curRate / max_c + DataRate("150Mbps") ) + 0.3 * qp->hp.m_curRate); if (new_rate < m_minRate) new_rate = m_minRate; if (new_rate > qp->m_max_rate) new_rate = qp->m_max_rate; qp->prevRtt = prevRtt; qp->prevCompletion = Simulator::Now().GetNanoSeconds(); ChangeRate(qp, new_rate); if (!fast_react) { qp->hp.m_curRate = new_rate; } if (!fast_react) { if (next_seq > qp->hp.m_lastUpdateSeq) qp->hp.m_lastUpdateSeq = next_seq; } } } void TcpDataCenter::FastReactPower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) { if (m_fast_react) UpdateRatePower(qp, p, ch, true); } void TcpDataCenter::Setup(QpCompleteCallback cb) { for (uint32_t i = 0; i < m_nic.size(); i++) { Ptr<QbbNetDevice> dev = m_nic[i].dev; if (dev == NULL) continue; dev->m_rdmaEQ->m_qpGrp = m_nic[i].qpGrp; dev->m_rdmaReceiveCb = MakeCallback(&TcpDataCenter::Receive, this); dev->m_rdmaLinkDownCb = MakeCallback(&TcpDataCenter::SetLinkDown, this); dev->m_rdmaPktSent = MakeCallback(&TcpDataCenter::PktSent, this); dev->m_rdmaEQ->m_rdmaGetNxtPkt = MakeCallback(&TcpDataCenter::GetNxtPacket, this); } m_qpCompleteCallback = cb; } void TcpDataCenter::DeleteQueuePair(Ptr<RdmaQueuePair> qp) { uint64_t key = GetQpKey(qp->dip.Get(), qp->sport, qp->m_pg); m_qpMap.erase(key); } Ptr<RdmaRxQueuePair> TcpDataCenter::GetRxQp(uint32_t sip, uint32_t dip, uint16_t sport, uint16_t dport, uint16_t pg, bool create) { uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport; auto it = m_rxQpMap.find(key); if (it != m_rxQpMap.end()) return it->second; if (create) { Ptr<RdmaRxQueuePair> q = CreateObject<RdmaRxQueuePair>(); q->sip = sip; q->dip = dip; q->sport = sport; q->dport = dport; q->m_ecn_source.qIndex = pg; m_rxQpMap[key] = q; return q; } return NULL; } void TcpDataCenter::DeleteRxQp(uint32_t dip, uint16_t pg, uint16_t dport) { uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport; m_rxQpMap.erase(key); } int TcpDataCenter::ReceiveUdp(Ptr<Packet> p, CustomHeader &ch) { uint8_t ecnbits = ch.GetIpv4EcnBits(); uint32_t payload_size = p->GetSize() - ch.GetSerializedSize(); Ptr<RdmaRxQueuePair> rxQp = GetRxQp(ch.dip, ch.sip, ch.udp.dport, ch.udp.sport, ch.udp.pg, true); if (ecnbits != 0) { rxQp->m_ecn_source.ecnbits |= ecnbits; rxQp->m_ecn_source.qfb++; } rxQp->m_ecn_source.total++; rxQp->m_milestone_rx = m_ack_interval; int x = ReceiverCheckSeq(ch.udp.seq, rxQp, payload_size); if (x == 1 || x == 2) { qbbHeader seqh; seqh.SetSeq(rxQp->ReceiverNextExpectedSeq); seqh.SetPG(ch.udp.pg); seqh.SetSport(ch.udp.dport); seqh.SetDport(ch.udp.sport); if (ecnbits) seqh.SetCnp(); Ptr<Packet> newp = Create<Packet>(std::max(60 - 14 - 20 - (int)seqh.GetSerializedSize(), 0)); newp->AddHeader(seqh); Ipv4Header head; head.SetDestination(Ipv4Address(ch.sip)); head.SetSource(Ipv4Address(ch.dip)); head.SetProtocol(x == 1 ? 0xFC : 0xFD); head.SetTtl(64); head.SetPayloadSize(newp->GetSize()); head.SetIdentification(rxQp->m_ipid++); newp->AddHeader(head); AddHeader(newp, 0x800); uint32_t nic_idx = GetNicIdxOfRxQp(rxQp); m_nic[nic_idx].dev->RdmaEnqueueHighPrioQ(newp); m_nic[nic_idx].dev->TriggerTransmit(); } return 0; } int TcpDataCenter::Receive(Ptr<Packet> p, CustomHeader &ch) { if (ch.l3Prot == 0x11) { ReceiveUdp(p, ch); } else if (ch.l3Prot == 0xFF) { ReceiveCnp(p, ch); } else if (ch.l3Prot == 0xFD) { ReceiveAck(p, ch); } else if (ch.l3Prot == 0xFC) { ReceiveAck(p, ch); } return 0; } void TcpDataCenter::RecoverQueue(Ptr<RdmaQueuePair> qp) { qp->snd_nxt = qp->snd_una; } void TcpDataCenter::QpComplete(Ptr<RdmaQueuePair> qp) { NS_ASSERT(!m_qpCompleteCallback.IsNull()); if (m_cc_mode == 1) { Simulator::Cancel(qp->mlx.m_eventUpdateAlpha); Simulator::Cancel(qp->mlx.m_eventDecreaseRate); Simulator::Cancel(qp->mlx.m_rpTimer); } m_qpCompleteCallback(qp); qp->m_notifyAppFinish(); DeleteQueuePair(qp); } void TcpDataCenter::SetLinkDown(Ptr<QbbNetDevice> dev) { printf("RdmaHw: node:%u a link down\n", m_node->GetId()); } void TcpDataCenter::AddTableEntry(Ipv4Address &dstAddr, uint32_t intf_idx) { uint32_t dip = dstAddr.Get(); m_rtTable[dip].push_back(intf_idx); } void TcpDataCenter::ClearTable() { m_rtTable.clear(); } void TcpDataCenter::RedistributeQp() { for (uint32_t i = 0; i < m_nic.size(); i++) { if (m_nic[i].dev == NULL) continue; m_nic[i].qpGrp->Clear(); } for (auto &it : m_qpMap) { Ptr<RdmaQueuePair> qp = it.second; uint32_t nic_idx = GetNicIdxOfQp(qp); m_nic[nic_idx].qpGrp->AddQp(qp); m_nic[nic_idx].dev->ReassignedQp(qp); } } Ptr<Packet> TcpDataCenter::GetNxtPacket(Ptr<RdmaQueuePair> qp) { uint32_t payload_size = qp->GetBytesLeft(); if (m_mtu < payload_size) payload_size = m_mtu; Ptr<Packet> p = Create<Packet> (payload_size); SeqTsHeader seqTs; seqTs.SetSeq (qp->snd_nxt); seqTs.SetPG (qp->m_pg); p->AddHeader (seqTs); UdpHeader udpHeader; udpHeader.SetDestinationPort (qp->dport); udpHeader.SetSourcePort (qp->sport); p->AddHeader (udpHeader); Ipv4Header ipHeader; ipHeader.SetSource (qp->sip); ipHeader.SetDestination (qp->dip); ipHeader.SetProtocol (0x11); ipHeader.SetPayloadSize (p->GetSize()); ipHeader.SetTtl (64); ipHeader.SetTos (0); ipHeader.SetIdentification (qp->m_ipid); p->AddHeader(ipHeader); PppHeader ppp; ppp.SetProtocol (0x0021); p->AddHeader (ppp); qp->snd_nxt += payload_size; qp->m_ipid++; return p; } void TcpDataCenter::PktSent(Ptr<RdmaQueuePair> qp, Ptr<Packet> pkt, Time interframeGap) { qp->lastPktSize = pkt->GetSize(); uint32_t seq = qp->snd_nxt; qp->rates[qp->snd_nxt] = Simulator::Now().GetNanoSeconds(); UpdateNextAvail(qp, interframeGap, pkt->GetSize()); } void TcpDataCenter::UpdateNextAvail(Ptr<RdmaQueuePair> qp, Time interframeGap, uint32_t pkt_size) { Time sendingTime; if (m_rateBound) sendingTime = interframeGap + qp->m_rate.CalculateBytesTxTime(pkt_size); else sendingTime = interframeGap + qp->m_max_rate.CalculateBytesTxTime(pkt_size); qp->m_nextAvail = Simulator::Now() + sendingTime; } void TcpDataCenter::ChangeRate(Ptr<RdmaQueuePair> qp, DataRate new_rate) { #if 1 Time sendingTime = qp->m_rate.CalculateBytesTxTime(qp->lastPktSize); Time new_sendintTime = new_rate.CalculateBytesTxTime(qp->lastPktSize); qp->m_nextAvail = qp->m_nextAvail + new_sendintTime - sendingTime; uint32_t nic_idx = GetNicIdxOfQp(qp); m_nic[nic_idx].dev->UpdateNextAvail(qp->m_nextAvail); #endif qp->m_rate = new_rate; } }
परिवर्तित टेक्स्ट
फ़ाइल खोलें
#include "tcp-datacenter.h" namespace ns3 { void TcpDataCenter::AddQueuePair(uint64_t size, uint16_t pg, Ipv4Address sip, Ipv4Address dip, uint16_t sport, uint16_t dport, uint32_t win, uint64_t baseRtt, Callback<void> notifyAppFinish, Time stopTime) { Ptr<RdmaQueuePair> qp = CreateObject<RdmaQueuePair>(pg, sip, dip, sport, dport); qp->SetSize(size); qp->SetWin(win); qp->SetBaseRtt(baseRtt); qp->SetVarWin(m_var_win); qp->SetAppNotifyCallback(notifyAppFinish); qp->stopTime = stopTime; uint32_t nic_idx = GetNicIdxOfQp(qp); m_nic[nic_idx].qpGrp->AddQp(qp); uint64_t key = GetQpKey(dip.Get(), sport, pg); m_qpMap[key] = qp; qp->powerEnabled = true; if (m_nic[nic_idx].dev == NULL) { std::cout << "sip " << sip << " dip " << dip << " sport " << sport << " dport " << dport << std::endl; } DataRate m_bps = m_nic[nic_idx].dev->GetDataRate(); if(win) qp->SetWin(m_bps.GetBitRate() * 1.0 * baseRtt * 1e-9 / 8); qp->m_rate = DataRate(m_bps.GetBitRate() / 10); qp->m_max_rate = m_bps; qp->hp.m_curRate = DataRate(m_bps.GetBitRate() / 10); m_nic[nic_idx].dev->NewQp(qp); } int TcpDataCenter::ReceiveCnp(Ptr<Packet> p, CustomHeader &ch) { uint32_t qIndex = ch.cnp.qIndex; if (qIndex == 1) { std::cout << "TCP--ignore\n"; return 0; } uint16_t udpport = ch.cnp.fid; Ptr<RdmaQueuePair> qp = GetQp(ch.sip, udpport, qIndex); if (qp == NULL) std::cout << "ERROR: QCN NIC cannot find the flow\n"; uint32_t nic_idx = GetNicIdxOfQp(qp); Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev; if (qp->m_rate == 0) { qp->m_rate = dev->GetDataRate(); qp->hp.m_curRate = dev->GetDataRate(); } return 0; } int TcpDataCenter::ReceiveAck(Ptr<Packet> p, CustomHeader &ch) { uint16_t qIndex = ch.ack.pg; uint16_t port = ch.ack.dport; uint32_t seq = ch.ack.seq; uint8_t cnp = (ch.ack.flags >> qbbHeader::FLAG_CNP) & 1; Ptr<RdmaQueuePair> qp = GetQp(ch.sip, port, qIndex); if (qp == NULL) { std::cout << "ERROR: " << "node:" << m_node->GetId() << ' ' << (ch.l3Prot == 0xFC ? "ACK" : "NACK") << " NIC cannot find the flow\n"; return 0; } uint32_t nic_idx = GetNicIdxOfQp(qp); Ptr<QbbNetDevice> dev = m_nic[nic_idx].dev; if (m_ack_interval == 0) std::cout << "ERROR: shouldn't receive ack\n"; else { if (!m_backto0) { qp->Acknowledge(seq); } else { uint32_t goback_seq = seq / m_chunk * m_chunk; qp->Acknowledge(goback_seq); } if (qp->IsFinished()) { QpComplete(qp); } } if (ch.l3Prot == 0xFD) RecoverQueue(qp); HandleAckHp(qp, p, ch); dev->TriggerTransmit(); return 0; } #define PRINT_LOG 0 void TcpDataCenter::HandleAckHp(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) { uint32_t ack_seq = ch.ack.seq; if (ack_seq > qp->hp.m_lastUpdateSeq) { UpdateRatePower(qp, p, ch, false); } else { FastReactPower(qp, p, ch); } } void TcpDataCenter::UpdateRatePower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch, bool fast_react) { uint32_t next_seq = qp->snd_nxt; double prevRtt = qp->m_baseRtt; double prevCompletion = Simulator::Now().GetNanoSeconds(); std::map<uint32_t, double>::iterator it = qp->rates.find(ch.ack.seq); if (it != qp->rates.end()) { qp->rates.erase(it); prevRtt = Simulator::Now().GetNanoSeconds() - it->second; qp->m_baseRtt = std::min(uint64_t(Simulator::Now().GetNanoSeconds() - it->second), qp->m_baseRtt); prevCompletion = Simulator::Now().GetNanoSeconds(); } if (qp->hp.m_lastUpdateSeq == 0) { qp->prevRtt = prevRtt; qp->prevCompletion = Simulator::Now().GetNanoSeconds(); qp->hp.m_lastUpdateSeq = next_seq; }else { double max_c = 0; double U = 0; uint64_t dt = 0; double dt_sample = prevCompletion - qp->prevCompletion; if (dt_sample < 1.0) dt_sample = 1.0; double A = ( double(prevRtt - qp->prevRtt) / dt_sample + 1 ); if (A < 0.5) A = 0.5; double power = ( A ) * (prevRtt); double powerx = (power) / (1.03 * qp->m_baseRtt); double u = powerx; uint8_t cnp = (ch.ack.flags >> qbbHeader::FLAG_CNP) & 1; if (cnp) { if (u < 2.0) u = 2.0; } if (u > U) { U = u; dt = (uint64_t)dt_sample; } DataRate new_rate; if (dt > 1.0 * qp->m_baseRtt) dt = 1.0 * qp->m_baseRtt; if (U < 0) { U = qp->hp.u; } qp->hp.u = (qp->hp.u * (1.0 * qp->m_baseRtt - dt) + U * dt) / double(1.0 * qp->m_baseRtt); max_c = qp->hp.u; new_rate = (0.8 * ( qp->hp.m_curRate / max_c + DataRate("250Mbps") ) + 0.2 * qp->hp.m_curRate); if (new_rate < m_minRate) new_rate = m_minRate; if (new_rate > qp->m_max_rate) new_rate = qp->m_max_rate; qp->prevRtt = prevRtt; qp->prevCompletion = Simulator::Now().GetNanoSeconds(); ChangeRate(qp, new_rate); if (!fast_react) { qp->hp.m_curRate = new_rate; } if (!fast_react) { if (next_seq > qp->hp.m_lastUpdateSeq) qp->hp.m_lastUpdateSeq = next_seq; } } } void TcpDataCenter::FastReactPower(Ptr<RdmaQueuePair> qp, Ptr<Packet> p, CustomHeader &ch) { if (m_fast_react) UpdateRatePower(qp, p, ch, true); } void TcpDataCenter::Setup(QpCompleteCallback cb) { for (uint32_t i = 0; i < m_nic.size(); i++) { Ptr<QbbNetDevice> dev = m_nic[i].dev; if (dev == NULL) continue; dev->m_rdmaEQ->m_qpGrp = m_nic[i].qpGrp; dev->m_rdmaReceiveCb = MakeCallback(&TcpDataCenter::Receive, this); dev->m_rdmaLinkDownCb = MakeCallback(&TcpDataCenter::SetLinkDown, this); dev->m_rdmaPktSent = MakeCallback(&TcpDataCenter::PktSent, this); dev->m_rdmaEQ->m_rdmaGetNxtPkt = MakeCallback(&TcpDataCenter::GetNxtPacket, this); } m_qpCompleteCallback = cb; } void TcpDataCenter::DeleteQueuePair(Ptr<RdmaQueuePair> qp) { uint64_t key = GetQpKey(qp->dip.Get(), qp->sport, qp->m_pg); m_qpMap.erase(key); } Ptr<RdmaRxQueuePair> TcpDataCenter::GetRxQp(uint32_t sip, uint32_t dip, uint16_t sport, uint16_t dport, uint16_t pg, bool create) { uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport; auto it = m_rxQpMap.find(key); if (it != m_rxQpMap.end()) return it->second; if (create) { Ptr<RdmaRxQueuePair> q = CreateObject<RdmaRxQueuePair>(); q->sip = sip; q->dip = dip; q->sport = sport; q->dport = dport; q->m_ecn_source.qIndex = pg; m_rxQpMap[key] = q; return q; } return NULL; } void TcpDataCenter::DeleteRxQp(uint32_t dip, uint16_t pg, uint16_t dport) { uint64_t key = ((uint64_t)dip << 32) | ((uint64_t)pg << 16) | (uint64_t)dport; m_rxQpMap.erase(key); } int TcpDataCenter::ReceiveUdp(Ptr<Packet> p, CustomHeader &ch) { uint8_t ecnbits = ch.GetIpv4EcnBits(); uint32_t payload_size = p->GetSize() - ch.GetSerializedSize(); Ptr<RdmaRxQueuePair> rxQp = GetRxQp(ch.dip, ch.sip, ch.udp.dport, ch.udp.sport, ch.udp.pg, true); if (ecnbits != 0) { rxQp->m_ecn_source.ecnbits |= ecnbits; rxQp->m_ecn_source.qfb++; } rxQp->m_ecn_source.total++; rxQp->m_milestone_rx = m_ack_interval; int x = ReceiverCheckSeq(ch.udp.seq, rxQp, payload_size); if (x == 1 || x == 2) { qbbHeader seqh; seqh.SetSeq(rxQp->ReceiverNextExpectedSeq); seqh.SetPG(ch.udp.pg); seqh.SetSport(ch.udp.dport); seqh.SetDport(ch.udp.sport); if (ecnbits) seqh.SetCnp(); Ptr<Packet> newp = Create<Packet>(std::max(60 - 14 - 20 - (int)seqh.GetSerializedSize(), 0)); newp->AddHeader(seqh); Ipv4Header head; head.SetDestination(Ipv4Address(ch.sip)); head.SetSource(Ipv4Address(ch.dip)); head.SetProtocol(x == 1 ? 0xFC : 0xFD); head.SetTtl(64); head.SetPayloadSize(newp->GetSize()); head.SetIdentification(rxQp->m_ipid++); newp->AddHeader(head); AddHeader(newp, 0x800); uint32_t nic_idx = GetNicIdxOfRxQp(rxQp); m_nic[nic_idx].dev->RdmaEnqueueHighPrioQ(newp); m_nic[nic_idx].dev->TriggerTransmit(); } return 0; } int TcpDataCenter::Receive(Ptr<Packet> p, CustomHeader &ch) { if (ch.l3Prot == 0x11) { ReceiveUdp(p, ch); } else if (ch.l3Prot == 0xFF) { ReceiveCnp(p, ch); } else if (ch.l3Prot == 0xFD) { ReceiveAck(p, ch); } else if (ch.l3Prot == 0xFC) { ReceiveAck(p, ch); } return 0; } void TcpDataCenter::RecoverQueue(Ptr<RdmaQueuePair> qp) { qp->snd_nxt = qp->snd_una; } void TcpDataCenter::QpComplete(Ptr<RdmaQueuePair> qp) { NS_ASSERT(!m_qpCompleteCallback.IsNull()); if (m_cc_mode == 1) { Simulator::Cancel(qp->mlx.m_eventUpdateAlpha); Simulator::Cancel(qp->mlx.m_eventDecreaseRate); Simulator::Cancel(qp->mlx.m_rpTimer); } m_qpCompleteCallback(qp); qp->m_notifyAppFinish(); DeleteQueuePair(qp); } void TcpDataCenter::SetLinkDown(Ptr<QbbNetDevice> dev) { printf("RdmaHw: node:%u a link down\n", m_node->GetId()); } void TcpDataCenter::AddTableEntry(Ipv4Address &dstAddr, uint32_t intf_idx) { uint32_t dip = dstAddr.Get(); m_rtTable[dip].push_back(intf_idx); } void TcpDataCenter::ClearTable() { m_rtTable.clear(); } void TcpDataCenter::RedistributeQp() { for (uint32_t i = 0; i < m_nic.size(); i++) { if (m_nic[i].dev == NULL) continue; m_nic[i].qpGrp->Clear(); } for (auto &it : m_qpMap) { Ptr<RdmaQueuePair> qp = it.second; uint32_t nic_idx = GetNicIdxOfQp(qp); m_nic[nic_idx].qpGrp->AddQp(qp); m_nic[nic_idx].dev->ReassignedQp(qp); } } Ptr<Packet> TcpDataCenter::GetNxtPacket(Ptr<RdmaQueuePair> qp) { uint32_t payload_size = qp->GetBytesLeft(); if (m_mtu < payload_size) payload_size = m_mtu; Ptr<Packet> p = Create<Packet> (payload_size); SeqTsHeader seqTs; seqTs.SetSeq (qp->snd_nxt); seqTs.SetPG (qp->m_pg); p->AddHeader (seqTs); UdpHeader udpHeader; udpHeader.SetDestinationPort (qp->dport); udpHeader.SetSourcePort (qp->sport); p->AddHeader (udpHeader); Ipv4Header ipHeader; ipHeader.SetSource (qp->sip); ipHeader.SetDestination (qp->dip); ipHeader.SetProtocol (0x11); ipHeader.SetPayloadSize (p->GetSize()); ipHeader.SetTtl (64); ipHeader.SetTos (0); ipHeader.SetIdentification (qp->m_ipid); p->AddHeader(ipHeader); PppHeader ppp; ppp.SetProtocol (0x0021); p->AddHeader (ppp); qp->snd_nxt += payload_size; qp->m_ipid++; return p; } void TcpDataCenter::PktSent(Ptr<RdmaQueuePair> qp, Ptr<Packet> pkt, Time interframeGap) { qp->lastPktSize = pkt->GetSize(); uint32_t seq = qp->snd_nxt; qp->rates[qp->snd_nxt] = Simulator::Now().GetNanoSeconds(); UpdateNextAvail(qp, interframeGap, pkt->GetSize()); } void TcpDataCenter::UpdateNextAvail(Ptr<RdmaQueuePair> qp, Time interframeGap, uint32_t pkt_size) { Time sendingTime; if (m_rateBound) sendingTime = interframeGap + qp->m_rate.CalculateBytesTxTime(pkt_size); else sendingTime = interframeGap + qp->m_max_rate.CalculateBytesTxTime(pkt_size); qp->m_nextAvail = Simulator::Now() + sendingTime; } void TcpDataCenter::ChangeRate(Ptr<RdmaQueuePair> qp, DataRate new_rate) { #if 1 Time sendingTime = qp->m_rate.CalculateBytesTxTime(qp->lastPktSize); Time new_sendintTime = new_rate.CalculateBytesTxTime(qp->lastPktSize); qp->m_nextAvail = qp->m_nextAvail + new_sendintTime - sendingTime; uint32_t nic_idx = GetNicIdxOfQp(qp); m_nic[nic_idx].dev->UpdateNextAvail(qp->m_nextAvail); #endif qp->m_rate = new_rate; } }
अंतर खोजें