Diff
checker
Testo
Testo
Immagini
Documenti
Excel
Cartelle
Legal
Enterprise
Applicazione per desktop
Prezzi
Accedi
Scarica Diffchecker Desktop
Confronta il testo
Trova la differenza tra due file di testo
Strumenti
Cronologia
Editor live
Nascondi spazi bianchi
Comprimi invariate
Senza a capo
Layout
Diviso
Unificato
Livello di dettaglio
Intelligente
Parola
Carattere
Stili testo
Modifica aspetto
Evidenziazione sintassi
Scegli sintassi
Ignora
Trasforma testo
Vai alla prima modifica
Modifica input
Diffchecker Desktop
Il modo più sicuro per usare Diffchecker. Ottieni l'app Diffchecker Desktop: i tuoi diff non lasciano mai il tuo computer!
Ottieni Desktop
EvolveDatacenterTCPDiff
Creato
6 mesi fa
Il diff non scade mai
Eliminare
Esporta
Condividere
Spiegare
8 rimozioni
Linee
Totale
Rimosso
Caratteri
Totale
Rimosso
Per continuare a utilizzare questa funzione, aggiorna a
Diff
checker
Pro
Visualizza prezzi
383 linee
Copia tutti
19 aggiunte
Linee
Totale
Aggiunto
Caratteri
Totale
Aggiunto
Per continuare a utilizzare questa funzione, aggiorna a
Diff
checker
Pro
Visualizza prezzi
389 linee
Copia tutti
#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)
Copia
Copiato
Copia
Copiato
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;
Copia
Copiato
Copia
Copiato
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;
Copia
Copiato
Copia
Copiato
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);
Copia
Copiato
Copia
Copiato
double powerx = (power) / (1.
05
* qp->m_baseRtt);
double powerx = (power) / (1.
03
* qp->m_baseRtt);
double u = powerx;
double u = powerx;
Copia
Copiato
Copia
Copiato
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;
Copia
Copiato
Copia
Copiato
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;
Copia
Copiato
Copia
Copiato
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 salvati
Testo originale
Apri file
#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; } }
Testo modificato
Apri file
#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; } }
Trovare la differenza