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; } }
查找差异