Diff
checker
Texte
Texte
Images
Documents
Excel
Dossiers
Legal
Enterprise
Application de bureau
Prix
Se connecter
Télécharger Diffchecker Desktop
Comparer le texte
Trouver la différence entre deux fichiers texte
Outils
Historique
Éditeur live
Masquer les espaces
Cacher identiques
Sans retour à la ligne
Vue
Divisé
Unifié
Niveau de précision
Intelligent
Mot
Caractère
Styles de texte
Modifier l’apparence
Coloration syntaxique
Choisir la syntaxe
Ignorer
Transformer le texte
Aller au premier écart
Modifier l'entrée
Diffchecker Desktop
La façon la plus sécurisée d'utiliser Diffchecker. Obtenez l'application Diffchecker Desktop : vos diffs ne quittent jamais votre ordinateur !
Obtenir Desktop
EvolveDatacenterTCPDiff
Créé
il y a 6 mois
Le diff n'expire jamais
Effacer
Exporter
Partager
Expliquer
8 suppressions
Lignes
Total
Supprimé
Caractères
Total
Supprimé
Pour continuer à utiliser cette fonctionnalité, passez à
Diff
checker
Pro
Voir les prix
383 lignes
Copier tout
19 ajouts
Lignes
Total
Ajouté
Caractères
Total
Ajouté
Pour continuer à utiliser cette fonctionnalité, passez à
Diff
checker
Pro
Voir les prix
389 lignes
Copier tout
#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)
Copier
Copié
Copier
Copié
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;
Copier
Copié
Copier
Copié
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;
Copier
Copié
Copier
Copié
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);
Copier
Copié
Copier
Copié
double powerx = (power) / (1.
05
* qp->m_baseRtt);
double powerx = (power) / (1.
03
* qp->m_baseRtt);
double u = powerx;
double u = powerx;
Copier
Copié
Copier
Copié
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;
Copier
Copié
Copier
Copié
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;
Copier
Copié
Copier
Copié
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;
}
}
}
}
Différences enregistrées
Texte d'origine
Ouvrir un fichier
#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; } }
Texte modifié
Ouvrir un fichier
#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; } }
Trouver la différence