Diff
checker
文本
文本
圖像
文檔
Excel
文件夾
Legal
Enterprise
桌面版
定價
登入
下載 Diffchecker 桌面版
比較文本
尋找兩個文字檔案之間的差異
工具
歷史
即時編輯器
隱藏空白變更
摺疊未變更行
關閉換行
檢視
拆分
統一
比對精度
智能
單詞
字符
文字樣式
變更外觀
語法突出顯示
選擇語法
忽略
文字轉換
前往第一個差異
編輯輸入
Diffchecker Desktop
執行Diffchecker最安全的方式。取得Diffchecker桌面應用程式:您的差異永遠不會離開您的電腦!
取得桌面版
EvolveDatacenterTCPDiff
建立於
6 個月前
差異永不過期
清除
匯出
分享
解釋
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;
}
}
}
}
已保存差異
原始文本
開啟檔案
#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; } }
尋找差異