Commit 669be5b4 authored by Bartosz Podrygajlo's avatar Bartosz Podrygajlo

Rfsim: Process received packets on demand

- Received packets are dynamically allocated and kept in a
queue for as long as they are needed.

- The layout of the packets was changed from interleaved antennas
to non-interleaved antennas.

- Generate a contiguous buffer from saved packets as needed.

There is extra work caused by this change due to recv() call no
longer saving data directly into the circular buffer, however further
processing is easier due to contiguous nature of the tx antenna buffers
prepared. The impact on performance was not measured.
parent 4ec69cad
...@@ -41,14 +41,12 @@ ...@@ -41,14 +41,12 @@
either we regenerate the channel (call again random_channel(desc,0)), or we keep it over subframes either we regenerate the channel (call again random_channel(desc,0)), or we keep it over subframes
legacy: we regenerate each sub frame in UL, and each frame only in DL legacy: we regenerate each sub frame in UL, and each frame only in DL
*/ */
void rxAddInput(const c16_t *input_sig, void rxAddInput(const c16_t **input_sig,
cf_t *after_channel_sig, cf_t *after_channel_sig,
int rxAnt, int rxAnt,
channel_desc_t *channelDesc, channel_desc_t *channelDesc,
int nbSamples, int nbSamples,
uint64_t TS, uint64_t TS)
uint32_t CirSize,
bool add_noise)
{ {
static uint64_t last_TS = 0; static uint64_t last_TS = 0;
...@@ -249,8 +247,7 @@ void rxAddInput(const c16_t *input_sig, ...@@ -249,8 +247,7 @@ void rxAddInput(const c16_t *input_sig,
const double pathLossLinear = pow(10, channelDesc->path_loss_dB / 20.0); const double pathLossLinear = pow(10, channelDesc->path_loss_dB / 20.0);
// Energy in one sample to calibrate input noise // Energy in one sample to calibrate input noise
// the normalized OAI value seems to be 256 as average amplitude (numerical amplification = 1) // the normalized OAI value seems to be 256 as average amplitude (numerical amplification = 1)
const double noise_per_sample = add_noise ? pow(10, channelDesc->noise_power_dB / 10.0) * 256 : 0; const double noise_per_sample = pow(10, channelDesc->noise_power_dB / 10.0) * 256;
const uint64_t dd = channelDesc->channel_offset;
const int nbTx = channelDesc->nb_tx; const int nbTx = channelDesc->nb_tx;
double Doppler_phase_cur = channelDesc->Doppler_phase_cur[rxAnt]; double Doppler_phase_cur = channelDesc->Doppler_phase_cur[rxAnt];
Doppler_phase_cur -= 2 * M_PI * round(Doppler_phase_cur / (2 * M_PI)); Doppler_phase_cur -= 2 * M_PI * round(Doppler_phase_cur / (2 * M_PI));
...@@ -264,15 +261,8 @@ void rxAddInput(const c16_t *input_sig, ...@@ -264,15 +261,8 @@ void rxAddInput(const c16_t *input_sig,
// const struct complex *channelModelEnd=channelModel+channelDesc->channel_length; // const struct complex *channelModelEnd=channelModel+channelDesc->channel_length;
for (int l = 0; l < (int)channelDesc->channel_length; l++) { for (int l = 0; l < (int)channelDesc->channel_length; l++) {
// let's assume TS+i >= l const int idx = i - l + channelDesc->channel_length - 1;
// fixme: the rfsimulator current structure is interleaved antennas const struct complex16 tx16 = input_sig[txAnt][idx];
// this has been designed to not have to wait a full block transmission
// but it is not very usefull
// it would be better to split out each antenna in a separate flow
// that will allow to mix ru antennas freely
// (X + cirSize) % cirSize to ensure that index is positive
const int idx = ((TS + i - l - dd) * nbTx + txAnt + CirSize) % CirSize;
const struct complex16 tx16 = input_sig[idx];
rx_tmp.r += tx16.r * channelModel[l].r - tx16.i * channelModel[l].i; rx_tmp.r += tx16.r * channelModel[l].r - tx16.i * channelModel[l].i;
rx_tmp.i += tx16.i * channelModel[l].r + tx16.r * channelModel[l].i; rx_tmp.i += tx16.i * channelModel[l].r + tx16.r * channelModel[l].i;
} // l } // l
...@@ -297,11 +287,10 @@ void rxAddInput(const c16_t *input_sig, ...@@ -297,11 +287,10 @@ void rxAddInput(const c16_t *input_sig,
channelDesc->Doppler_phase_cur[rxAnt] = Doppler_phase_cur; channelDesc->Doppler_phase_cur[rxAnt] = Doppler_phase_cur;
if ((TS * nbTx) % CirSize + nbSamples <= CirSize)
// Cast to a wrong type for compatibility ! // Cast to a wrong type for compatibility !
LOG_D(HW, LOG_D(HW,
"Input power %f, output power: %f, channel path loss %f, noise coeff: %f \n", "Input power %f, output power: %f, channel path loss %f, noise coeff: %f \n",
10 * log10((double)signal_energy((int32_t *)&input_sig[(TS * nbTx) % CirSize], nbSamples)), 10 * log10((double)signal_energy((int32_t *)&input_sig[0], nbSamples)),
10 * log10((double)signal_energy((int32_t *)after_channel_sig, nbSamples)), 10 * log10((double)signal_energy((int32_t *)after_channel_sig, nbSamples)),
channelDesc->path_loss_dB, channelDesc->path_loss_dB,
10 * log10(noise_per_sample)); 10 * log10(noise_per_sample));
......
...@@ -24,13 +24,11 @@ ...@@ -24,13 +24,11 @@
#ifndef __RFSIMULATOR_H #ifndef __RFSIMULATOR_H
#define __RFSIMULATOR_H #define __RFSIMULATOR_H
#include <stdbool.h> #include <stdbool.h>
void rxAddInput(const c16_t *input_sig, void rxAddInput(const c16_t **input_sig,
cf_t *after_channel_sig, cf_t *after_channel_sig,
int rxAnt, int rxAnt,
channel_desc_t *channelDesc, channel_desc_t *channelDesc,
int nbSamples, int nbSamples,
uint64_t TS, uint64_t TS);
uint32_t CirSize,
bool apply_noise);
#endif #endif
...@@ -27,6 +27,7 @@ ...@@ -27,6 +27,7 @@
* When the opposite side switch from passive reading to active R+Write, the synchro is not fully deterministic * When the opposite side switch from passive reading to active R+Write, the synchro is not fully deterministic
*/ */
#include "PHY/defs_common.h"
#include "utils.h" #include "utils.h"
#include <sys/socket.h> #include <sys/socket.h>
#include <sys/types.h> #include <sys/types.h>
...@@ -50,27 +51,15 @@ ...@@ -50,27 +51,15 @@
extern "C" { extern "C" {
#include <common/utils/load_module_shlib.h> #include <common/utils/load_module_shlib.h>
#include <openair1/SIMULATION/TOOLS/sim.h> #include <openair1/SIMULATION/TOOLS/sim.h>
}
#include "rfsimulator.h" #include "rfsimulator.h"
}
#include <queue> #include <queue>
#include <mutex> #include <mutex>
#include <vector> #include <vector>
#include <sstream> #include <sstream>
#include <algorithm>
#define PORT 4043 // default TCP port for this simulator #define PORT 4043 // default TCP port for this simulator
//
// CirSize defines the number of samples inquired for a read cycle
// It is bounded by a slot read capability (which depends on bandwidth and numerology)
// up to multiple slots read to allow I/Q buffering of the I/Q TCP stream
//
// As a rule of thumb:
// -it can't be less than the number of samples for a slot
// -it can range up to multiple slots
//
// The default value is chosen for 10ms buffering which makes 23040*20 = 460800 samples
// The previous value is kept below in comment it was computed for 100ms 1x 20MHz
// #define CirSize 6144000 // 100ms SiSo 20MHz LTE
#define minCirSize 460800 // 10ms SiSo 40Mhz 3/4 sampling NR78 FR1
#define sampleToByte(a, b) ((a) * (b) * sizeof(sample_t)) #define sampleToByte(a, b) ((a) * (b) * sizeof(sample_t))
#define byteToSample(a, b) ((a) / (sizeof(sample_t) * (b))) #define byteToSample(a, b) ((a) / (sizeof(sample_t) * (b)))
...@@ -141,7 +130,6 @@ static telnetshell_cmddef_t rfsimu_cmdarray[] = { ...@@ -141,7 +130,6 @@ static telnetshell_cmddef_t rfsimu_cmdarray[] = {
static telnetshell_cmddef_t *setmodel_cmddef = &(rfsimu_cmdarray[1]); static telnetshell_cmddef_t *setmodel_cmddef = &(rfsimu_cmdarray[1]);
static telnetshell_vardef_t rfsimu_vardef[] = {{"", 0, 0, NULL}}; static telnetshell_vardef_t rfsimu_vardef[] = {{"", 0, 0, NULL}};
static uint64_t CirSize = minCirSize;
typedef c16_t sample_t; // 2*16 bits complex number typedef c16_t sample_t; // 2*16 bits complex number
typedef struct beam_switch_command_t { typedef struct beam_switch_command_t {
...@@ -155,20 +143,25 @@ typedef struct { ...@@ -155,20 +143,25 @@ typedef struct {
std::mutex mutex; std::mutex mutex;
} beam_state_t; } beam_state_t;
typedef struct {
samplesBlockHeader_t header;
char payload[1];
} rfsim_packet_t;
typedef struct buffer_s { typedef struct buffer_s {
int conn_sock; int conn_sock;
openair0_timestamp lastReceivedTS; openair0_timestamp lastReceivedTS;
bool headerMode; bool headerMode;
bool trashingPacket; bool trashingPacket;
samplesBlockHeader_t th; samplesBlockHeader_t th;
uint nbAnt;
char *transferPtr; char *transferPtr;
uint64_t remainToTransfer; uint64_t remainToTransfer;
char *circularBufEnd;
sample_t *circularBuf;
channel_desc_t *channel_model; channel_desc_t *channel_model;
char *packet_ptr; rfsim_packet_t *packet_ptr;
size_t payload_sz; size_t payload_sz;
size_t remainToTransferBeam; size_t remainToTransferBeam;
std::queue<rfsim_packet_t *> received_packets;
} buffer_t; } buffer_t;
typedef struct { typedef struct {
...@@ -273,6 +266,29 @@ static void clear_beam_queue(beam_state_t *beam_state, openair0_timestamp timest ...@@ -273,6 +266,29 @@ static void clear_beam_queue(beam_state_t *beam_state, openair0_timestamp timest
} }
} }
/**
* @brief Clears old packets from the received_packets queue based on a threshold timestamp.
*
* This function iterates through the received_packets queue and removes packets whose
* timestamp plus size is less than or equal to the provided threshold_timestamp.
* It frees the memory allocated for each removed packet.
*
* @param received_packets Reference to the queue of rfsim_packet_t pointers representing received packets.
* @param threshold_timestamp The timestamp threshold used to determine which packets to remove.
*/
static void clear_old_packets(std::queue<rfsim_packet_t *> &received_packets, uint64_t threshold_timestamp)
{
while (!received_packets.empty()) {
rfsim_packet_t *pkt = received_packets.front();
if (pkt->header.timestamp + pkt->header.size <= threshold_timestamp) {
free(pkt);
received_packets.pop();
} else {
break;
}
}
}
static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time); static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time);
static buffer_t *allocCirBuf(rfsimulator_state_t *bridge, int sock) static buffer_t *allocCirBuf(rfsimulator_state_t *bridge, int sock)
...@@ -280,18 +296,13 @@ static buffer_t *allocCirBuf(rfsimulator_state_t *bridge, int sock) ...@@ -280,18 +296,13 @@ static buffer_t *allocCirBuf(rfsimulator_state_t *bridge, int sock)
uint64_t buff_index = bridge->next_buf++ % MAX_FD_RFSIMU; uint64_t buff_index = bridge->next_buf++ % MAX_FD_RFSIMU;
buffer_t *ptr = &bridge->buf[buff_index]; buffer_t *ptr = &bridge->buf[buff_index];
bridge->nb_cnx++; bridge->nb_cnx++;
ptr->circularBuf = static_cast<sample_t *>(calloc(1, sampleToByte(CirSize, 1)));
if (ptr->circularBuf == NULL) {
LOG_E(HW, "malloc(%lu) failed\n", sampleToByte(CirSize, 1));
return NULL;
}
ptr->circularBufEnd = ((char *)ptr->circularBuf) + sampleToByte(CirSize, 1);
ptr->conn_sock = sock; ptr->conn_sock = sock;
ptr->lastReceivedTS = 0; ptr->lastReceivedTS = 0;
ptr->headerMode = true; ptr->headerMode = true;
ptr->trashingPacket = true; ptr->trashingPacket = true;
ptr->transferPtr = (char *)&ptr->th; ptr->transferPtr = (char *)&ptr->th;
ptr->remainToTransfer = sizeof(samplesBlockHeader_t); ptr->remainToTransfer = sizeof(samplesBlockHeader_t);
ptr->received_packets = std::queue<rfsim_packet_t *>();
int sendbuff = SEND_BUFF_SIZE; int sendbuff = SEND_BUFF_SIZE;
if (setsockopt(sock, SOL_SOCKET, SO_SNDBUF, &sendbuff, sizeof(sendbuff)) != 0) { if (setsockopt(sock, SOL_SOCKET, SO_SNDBUF, &sendbuff, sizeof(sendbuff)) != 0) {
LOG_E(HW, "setsockopt(SO_SNDBUF) failed\n"); LOG_E(HW, "setsockopt(SO_SNDBUF) failed\n");
...@@ -338,11 +349,11 @@ static void removeCirBuf(rfsimulator_state_t *bridge, buffer_t *buf) ...@@ -338,11 +349,11 @@ static void removeCirBuf(rfsimulator_state_t *bridge, buffer_t *buf)
LOG_E(HW, "epoll_ctl(EPOLL_CTL_DEL) failed\n"); LOG_E(HW, "epoll_ctl(EPOLL_CTL_DEL) failed\n");
} }
close(buf->conn_sock); close(buf->conn_sock);
free(buf->circularBuf);
// Fixme: no free_channel_desc_scm(bridge->buf[sock].channel_model) implemented // Fixme: no free_channel_desc_scm(bridge->buf[sock].channel_model) implemented
// a lot of mem leaks // a lot of mem leaks
// free(bridge->buf[sock].channel_model); // free(bridge->buf[sock].channel_model);
memset(buf, 0, sizeof(buffer_t)); clear_old_packets(buf->received_packets, INT64_MAX);
*buf = buffer_t{};
buf->conn_sock = -1; buf->conn_sock = -1;
bridge->nb_cnx--; bridge->nb_cnx--;
} }
...@@ -856,37 +867,17 @@ static int rfsimulator_write_internal(rfsimulator_state_t *t, ...@@ -856,37 +867,17 @@ static int rfsimulator_write_internal(rfsimulator_state_t *t,
LOG_E(HW, "rfsimulator sending 0 tx antennas\n"); LOG_E(HW, "rfsimulator sending 0 tx antennas\n");
samplesBlockHeader_t header = {(uint32_t)nsamps, (uint32_t)nbAnt, (uint64_t)timestamp, 0, 0, beam_map}; samplesBlockHeader_t header = {(uint32_t)nsamps, (uint32_t)nbAnt, (uint64_t)timestamp, 0, 0, beam_map};
fullwrite(b->conn_sock, &header, sizeof(header), t); fullwrite(b->conn_sock, &header, sizeof(header), t);
if (t->beam_ctrl->enable_beams) {
int num_beams = __builtin_popcountll(beam_map); int num_beams = __builtin_popcountll(beam_map);
AssertFatal(num_beams > 0, "Must set at least one bit in beam_map\n"); AssertFatal(num_beams > 0, "Must set at least one bit in beam_map\n");
sample_t tmpSamples[nsamps][num_beams][nbAnt];
for (int beam = 0; beam < num_beams; beam++) { for (int beam = 0; beam < num_beams; beam++) {
for (int a = 0; a < nbAnt; a++) { for (int a = 0; a < nbAnt; a++) {
sample_t *in = (sample_t *)samplesVoid[beam][a]; sample_t *in = (sample_t *)samplesVoid[beam][a];
for (int s = 0; s < nsamps; s++) fullwrite(b->conn_sock, (void *)in, sampleToByte(nsamps, 1), t);
tmpSamples[s][beam][a] = in[s];
}
}
fullwrite(b->conn_sock, (void *)tmpSamples, sampleToByte(nsamps, nbAnt * num_beams), t);
} else {
if (nbAnt == 1) {
fullwrite(b->conn_sock, samplesVoid[0][0], sampleToByte(nsamps, nbAnt), t);
} else {
sample_t tmpSamples[nsamps][nbAnt];
for (int a = 0; a < nbAnt; a++) {
sample_t *in = (sample_t *)samplesVoid[0][a];
for (int s = 0; s < nsamps; s++)
tmpSamples[s][a] = in[s];
}
fullwrite(b->conn_sock, (void *)tmpSamples, sampleToByte(nsamps, nbAnt), t);
} }
} }
} }
} }
if (t->lastWroteTS != 0 && fabs((double)t->lastWroteTS - timestamp) > (double)CirSize)
LOG_W(HW, "Discontinuous TX gap too large Tx:%lu, %lu\n", t->lastWroteTS, timestamp);
if (t->lastWroteTS > timestamp) if (t->lastWroteTS > timestamp)
LOG_W(HW, "Not supported to send Tx out of order %lu, %lu\n", t->lastWroteTS, timestamp); LOG_W(HW, "Not supported to send Tx out of order %lu, %lu\n", t->lastWroteTS, timestamp);
...@@ -995,99 +986,110 @@ static bool add_client(rfsimulator_state_t *t) ...@@ -995,99 +986,110 @@ static bool add_client(rfsimulator_state_t *t)
static void process_recv_header(rfsimulator_state_t *t, buffer_t *b, bool first_time) static void process_recv_header(rfsimulator_state_t *t, buffer_t *b, bool first_time)
{ {
b->headerMode = false; // We got the header b->headerMode = false; // We got the header
AssertFatal(b->th.nbAnt != 0, "Number of antennas not set\n");
if (b->nbAnt != b->th.nbAnt) {
LOG_A(HW, "RFsim: Number of antennas changed from %d to %d\n", b->nbAnt, b->th.nbAnt);
b->nbAnt = b->th.nbAnt;
}
if (first_time) { if (first_time) {
b->lastReceivedTS = b->th.timestamp; b->lastReceivedTS = b->th.timestamp;
b->trashingPacket = true; b->trashingPacket = true;
} else { } else {
if (b->lastReceivedTS < (int64_t)b->th.timestamp) { if (b->lastReceivedTS < (int64_t)b->th.timestamp) {
// We have a transmission hole to fill, like TDD
// we create no signal samples up to the beginning of this reception
int nbAnt = b->th.nbAnt; int nbAnt = b->th.nbAnt;
if (!nbAnt) if (!nbAnt)
LOG_E(HW, "rfsimulator receive 0 rx antennas\n"); LOG_E(HW, "rfsimulator receive 0 rx antennas\n");
if (b->th.timestamp - b->lastReceivedTS < CirSize) {
// case we wrap at circular buffer end
for (uint64_t index = b->lastReceivedTS; index < b->th.timestamp; index++) {
for (int a = 0; a < nbAnt; a++) {
b->circularBuf[(index * nbAnt + a) % CirSize] = (c16_t){0};
}
}
} else {
// case no circular buffer wrap
memset(b->circularBuf, 0, sampleToByte(CirSize, 1));
}
b->lastReceivedTS = b->th.timestamp; b->lastReceivedTS = b->th.timestamp;
} else if (b->lastReceivedTS > (int64_t)b->th.timestamp) { } else if (b->lastReceivedTS > (int64_t)b->th.timestamp) {
LOG_W(HW, "Received data in past: current is %lu, new reception: %lu!\n", b->lastReceivedTS, b->th.timestamp); LOG_W(HW, "Received data in past: current is %lu, new reception: %lu!\n", b->lastReceivedTS, b->th.timestamp);
b->trashingPacket = true; b->trashingPacket = true;
} }
// verify the time gap is not too large
mutexlock(t->Sockmutex);
if (t->lastWroteTS != 0 && (fabs((double)t->lastWroteTS - b->lastReceivedTS) > (double)CirSize))
LOG_W(HW, "UEsock(%d) Tx/Rx shift too large Tx:%lu, Rx:%lu\n", b->conn_sock, t->lastWroteTS, b->lastReceivedTS);
mutexunlock(t->Sockmutex);
} }
int num_beams = __builtin_popcountll(b->th.beam_map); int num_beams = __builtin_popcountll(b->th.beam_map);
AssertFatal(b->th.beam_map == 1ULL || t->beam_ctrl->enable_beams == 1, AssertFatal(b->th.beam_map == 1ULL || t->beam_ctrl->enable_beams == 1,
"The transmitter has enabled beam simulation while this receiver has not\n"); "The transmitter has enabled beam simulation while this receiver has not\n");
size_t payload_sz = sampleToByte(b->th.size, b->th.nbAnt) * num_beams; size_t payload_sz = sampleToByte(b->th.size, b->th.nbAnt) * num_beams;
if (b->packet_ptr == NULL || payload_sz > b->payload_sz) { b->packet_ptr = static_cast<rfsim_packet_t *>(calloc_or_fail(1, payload_sz + sizeof(samplesBlockHeader_t)));
free(b->packet_ptr); b->packet_ptr->header = b->th;
b->packet_ptr = static_cast<char *>(calloc_or_fail(1, payload_sz)); b->transferPtr = b->packet_ptr->payload;
b->payload_sz = payload_sz;
}
b->transferPtr = b->packet_ptr;
b->remainToTransfer = payload_sz; b->remainToTransfer = payload_sz;
return; return;
} }
static void process_recv(rfsimulator_state_t *t, buffer_t *b, bool first_time) static std::vector<int> get_beam_ids(uint64_t beam_map)
{ {
int num_beams = __builtin_popcountll(b->th.beam_map); int num_beams = __builtin_popcountll(beam_map);
AssertFatal(num_beams > 0, "Needs at least one beam\n"); AssertFatal(num_beams > 0, "Needs at least one beam\n");
int num_ant = b->th.nbAnt; std::vector<int> beam_ids;
int num_samples = b->th.size;
int tx_beam_ids[MAX_BEAMS];
int *tx_beam_ids_p = tx_beam_ids;
for (int i = 0; i < MAX_BEAMS; i++) {
if (b->th.beam_map & (1ULL << i)) {
*tx_beam_ids_p++ = i;
}
}
c16_t *in = (c16_t *)b->packet_ptr;
int num_samples_to_process = num_samples;
openair0_timestamp first_sample_timestamp = b->th.timestamp;
while (num_samples_to_process > 0) {
uint32_t nsamps_out;
uint64_t rx_beam_map = get_beam_map(&t->beam_ctrl->rx, first_sample_timestamp, num_samples_to_process, &nsamps_out);
AssertFatal(__builtin_popcountll(rx_beam_map) == 1, "Only supports 1 beam\n");
int rx_beam_id = 0;
for (int i = 0; i < MAX_BEAMS; i++) { for (int i = 0; i < MAX_BEAMS; i++) {
if (rx_beam_map & (1ULL << i)) { if (beam_map & (1ULL << i)) {
rx_beam_id = i; beam_ids.push_back(i);
break;
} }
} }
for (uint i = 0; i < nsamps_out; i++) { return beam_ids;
for (int aatx = 0; aatx < num_ant; aatx++) { }
c16_t *out = &b->circularBuf[((first_sample_timestamp + i) * num_ant + aatx) % CirSize];
cf_t sample = {.r = 0, .i = 0}; /**
for (int beam = 0; beam < num_beams; beam++) { * @brief Combines samples from received packets into a single buffer for the rx beam
float gain_dB = get_rx_gain_db(t, rx_beam_id, tx_beam_ids[beam]); *
* This function processes the received_packets queue and combines the transmitted beams
* into a single buffer for each antenna. It applies the appropriate gain and handles
* overlapping timestamps.
*
*
* @param t Pointer to the rfsimulator_state_t structure.
* @param received_packets Reference to the queue of rfsim_packet_t pointers representing received packets.
* @param start_timestamp The start timestamp for the combination process.
* @param num_aatx The number of antennas to process.
* @param num_samples The number of samples to process.
* @param rx_beam_id Receiver configured beam id for the timestamp
* @return A vector of vectors containing the combined samples for each antenna.
*/
static std::vector<std::vector<c16_t>> combine_received_beams(rfsimulator_state_t *t,
std::queue<rfsim_packet_t *> &received_packets,
uint64_t start_timestamp,
int num_aatx,
size_t num_samples,
int rx_beam_id)
{
// Assume received_packets is ordered by timestamp
std::vector<std::vector<c16_t>> ant_buffers(num_aatx, std::vector<c16_t>(num_samples, {0, 0}));
std::queue<rfsim_packet_t *> packets_copy = received_packets;
while (!packets_copy.empty()) {
rfsim_packet_t *pkt = packets_copy.front();
if (pkt->header.timestamp + pkt->header.size <= start_timestamp) {
// This packet is before the start timestamp, discard it
packets_copy.pop();
continue;
}
if (pkt->header.timestamp > start_timestamp + num_samples) {
// This packet is after the end of the buffer, stop processing
break;
}
std::vector<int> tx_beams = get_beam_ids(pkt->header.beam_map);
for (uint beam = 0; beam < tx_beams.size(); beam++) {
float gain_dB = get_rx_gain_db(t, rx_beam_id, tx_beams[beam]);
float gain_linear = powf(10, gain_dB / 20.0); float gain_linear = powf(10, gain_dB / 20.0);
sample.r += in->r * gain_linear; uint64_t overlap_start = std::max(start_timestamp, pkt->header.timestamp);
sample.i += in->i * gain_linear; uint64_t overlap_end = std::min(start_timestamp + num_samples, pkt->header.timestamp + pkt->header.size);
in++; int write_start_idx = overlap_start - start_timestamp;
int write_end_idx = overlap_end - start_timestamp;
int read_start_idx = overlap_start - pkt->header.timestamp;
for (int aatx = 0; aatx < num_aatx; aatx++) {
c16_t *buffer = (c16_t *)pkt->payload;
c16_t *tx_ant_buffer_in = &buffer[(num_aatx * beam + aatx) * pkt->header.size + read_start_idx];
for (int s = write_start_idx; s < write_end_idx; s++) {
ant_buffers[aatx][s].r += tx_ant_buffer_in->r * gain_linear;
ant_buffers[aatx][s].i += tx_ant_buffer_in->i * gain_linear;
tx_ant_buffer_in++;
} }
out->r = sample.r;
out->i = sample.i;
} }
} }
first_sample_timestamp += nsamps_out; packets_copy.pop();
num_samples_to_process -= nsamps_out;
} }
return ant_buffers;
} }
static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time) static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time)
...@@ -1118,7 +1120,7 @@ static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time) ...@@ -1118,7 +1120,7 @@ static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time)
} }
} }
if (b->circularBuf == NULL) { if (b->conn_sock == -1) {
LOG_E(HW, "Received data on not connected socket %d\n", events[nbEv].data.fd); LOG_E(HW, "Received data on not connected socket %d\n", events[nbEv].data.fd);
continue; continue;
} }
...@@ -1142,10 +1144,13 @@ static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time) ...@@ -1142,10 +1144,13 @@ static bool flushInput(rfsimulator_state_t *t, int timeout, bool first_time)
b->remainToTransfer = sizeof(samplesBlockHeader_t); b->remainToTransfer = sizeof(samplesBlockHeader_t);
if (!b->trashingPacket) { if (!b->trashingPacket) {
process_recv(t, b, first_time);
b->lastReceivedTS = b->th.timestamp + b->th.size; b->lastReceivedTS = b->th.timestamp + b->th.size;
LOG_D(HW, "UEsock: %d Set b->lastReceivedTS %ld\n", b->conn_sock, b->lastReceivedTS); LOG_D(HW, "UEsock: %d Set b->lastReceivedTS %ld\n", b->conn_sock, b->lastReceivedTS);
b->received_packets.emplace(b->packet_ptr);
} else {
free(b->packet_ptr);
} }
b->packet_ptr = NULL;
b->trashingPacket = false; b->trashingPacket = false;
} }
} }
...@@ -1167,7 +1172,7 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest ...@@ -1167,7 +1172,7 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest
int first_sock; int first_sock;
for (first_sock = 0; first_sock < MAX_FD_RFSIMU; first_sock++) for (first_sock = 0; first_sock < MAX_FD_RFSIMU; first_sock++)
if (t->buf[first_sock].circularBuf != NULL) if (t->buf[first_sock].conn_sock != -1)
break; break;
if (first_sock == MAX_FD_RFSIMU) { if (first_sock == MAX_FD_RFSIMU) {
...@@ -1194,7 +1199,7 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest ...@@ -1194,7 +1199,7 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest
buffer_t *b = NULL; buffer_t *b = NULL;
for (int sock = 0; sock < MAX_FD_RFSIMU; sock++) { for (int sock = 0; sock < MAX_FD_RFSIMU; sock++) {
b = &t->buf[sock]; b = &t->buf[sock];
if (b->circularBuf && (t->nextRxTstamp + nsamps) > b->lastReceivedTS) { if (b->conn_sock != -1 && (t->nextRxTstamp + nsamps) > b->lastReceivedTS) {
have_to_wait = true; have_to_wait = true;
break; break;
} }
...@@ -1221,13 +1226,15 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest ...@@ -1221,13 +1226,15 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest
for (int a = 0; a < nbAnt; a++) for (int a = 0; a < nbAnt; a++)
memset(samplesVoid[a], 0, sampleToByte(nsamps, 1)); memset(samplesVoid[a], 0, sampleToByte(nsamps, 1));
cf_t temp_array[nbAnt][nsamps]; cf_t temp_array[nbAnt][nsamps];
bool apply_noise_per_channel = get_noise_power_dBFS() == INVALID_DBFS_VALUE; memset(temp_array, 0, sizeof(temp_array));
int num_chanmod_channels = 0;
// Add all input nodes signal in the output buffer // Add all input nodes signal in the output buffer
for (int sock = 0; sock < MAX_FD_RFSIMU; sock++) { for (int sock = 0; sock < MAX_FD_RFSIMU; sock++) {
buffer_t *ptr = &t->buf[sock]; buffer_t *ptr = &t->buf[sock];
if (ptr->circularBuf && !ptr->trashingPacket) { uint64_t timestamp_processed = t->nextRxTstamp;
if (ptr->conn_sock != -1 && !ptr->received_packets.empty()) {
AssertFatal(ptr->nbAnt != 0, "Number of antennas not set\n");
bool reGenerateChannel = false; bool reGenerateChannel = false;
// fixme: when do we regenerate // fixme: when do we regenerate
...@@ -1235,70 +1242,69 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest ...@@ -1235,70 +1242,69 @@ static int rfsimulator_read(openair0_device *device, openair0_timestamp *ptimest
if (reGenerateChannel) if (reGenerateChannel)
random_channel(ptr->channel_model, 0); random_channel(ptr->channel_model, 0);
for (int a_rx = 0; a_rx < nbAnt; a_rx++) { // loop over number of Rx antennas std::vector<int> rx_beams_ids = get_beam_ids(t->beam_ctrl->rx.beam_map);
if (ptr->channel_model != NULL) { // apply a channel model if (ptr->channel_model != NULL) { // apply a channel model
if (num_chanmod_channels == 0) { const uint64_t dd = ptr->channel_model->channel_offset;
memset(temp_array, 0, sizeof(temp_array)); const uint64_t channel_length = ptr->channel_model->channel_length;
timestamp_processed = t->nextRxTstamp - dd;
for (uint beam = 0; beam < rx_beams_ids.size(); beam++) {
std::vector<std::vector<c16_t>> ant_buffers = combine_received_beams(t,
ptr->received_packets,
t->nextRxTstamp - dd,
ptr->nbAnt,
nsamps + channel_length - 1,
rx_beams_ids[beam]);
const c16_t *input[ant_buffers.size()];
for (uint aatx = 0; aatx < ant_buffers.size(); aatx++) {
input[aatx] = ant_buffers[aatx].data();
}
for (int aarx = 0; aarx < nbAnt; aarx++) {
rxAddInput(input, temp_array[aarx], aarx, ptr->channel_model, nsamps, t->nextRxTstamp);
} }
num_chanmod_channels++;
rxAddInput(ptr->circularBuf,
temp_array[a_rx],
a_rx,
ptr->channel_model,
nsamps,
t->nextRxTstamp,
CirSize,
apply_noise_per_channel);
} else { // no channel modeling
int nbAnt_tx = ptr->th.nbAnt; // number of Tx antennas
int firstIndex = (CirSize + t->nextRxTstamp - t->chan_offset) % CirSize;
sample_t *out = (sample_t *)samplesVoid[a_rx];
if (nbAnt_tx == 1 && t->nb_cnx == 1) { // optimized for 1 Tx and 1 UE
sample_t *firstSample = (sample_t *)&(ptr->circularBuf[firstIndex]);
if ((uint64_t)firstIndex + (uint64_t)nsamps > CirSize) {
int tailSz = CirSize - firstIndex;
memcpy(out, firstSample, sampleToByte(tailSz, 1));
memcpy(out + tailSz, &ptr->circularBuf[0], sampleToByte(nsamps - tailSz, 1));
} else {
memcpy(out, firstSample, sampleToByte(nsamps, 1));
} }
} else { } else {
// SIMD (with simde) optimization might be added here later for (uint beam = 0; beam < rx_beams_ids.size(); beam++) {
LOG_D(HW, "nbAnt_tx %d nbAnt %d\n", nbAnt_tx, nbAnt); std::vector<std::vector<c16_t>> ant_buffers =
double H_awgn_mimo_coeff[nbAnt_tx]; combine_received_beams(t, ptr->received_packets, t->nextRxTstamp, ptr->th.nbAnt, nsamps, rx_beams_ids[beam]);
for (int a_tx = 0; a_tx < nbAnt_tx; a_tx++) { for (int aarx = 0; aarx < nbAnt; aarx++) {
uint32_t ant_diff = abs(a_tx - a_rx); double H_awgn_mimo_coeff[ant_buffers.size()];
H_awgn_mimo_coeff[a_tx] = ant_diff ? (0.2 / ant_diff) : 1.0; for (int aatx = 0; aatx < (int)ant_buffers.size(); aatx++) {
} uint32_t ant_diff = std::abs(aatx - aarx);
H_awgn_mimo_coeff[aatx] = ant_diff ? (0.2 / ant_diff) : 1.0;
}
for (uint aatx = 0; aatx < ant_buffers.size(); aatx++) { // sum up signals from nbAnt_tx antennas
c16_t *out = (sample_t *)samplesVoid[aarx];
for (int i = 0; i < nsamps; i++) { // loop over nsamps for (int i = 0; i < nsamps; i++) { // loop over nsamps
for (int a_tx = 0; a_tx < nbAnt_tx; a_tx++) { // sum up signals from nbAnt_tx antennas out[i].r += ant_buffers[aatx][i].r * H_awgn_mimo_coeff[aatx];
out[i].r += (short)(ptr->circularBuf[((firstIndex + i) * nbAnt_tx + a_tx) % CirSize].r * H_awgn_mimo_coeff[a_tx]); out[i].i += ant_buffers[aatx][i].i * H_awgn_mimo_coeff[aatx];
out[i].i += (short)(ptr->circularBuf[((firstIndex + i) * nbAnt_tx + a_tx) % CirSize].i * H_awgn_mimo_coeff[a_tx]);
} // end for a_tx } // end for a_tx
} // end for i (number of samps) } // end for i (number of samps)
} // end of 1 tx antenna optimization } // end for a_rx
} // end of no channel modeling
} // end for a (number of rx antennas)
} }
} }
if (apply_noise_per_channel && num_chanmod_channels > 0) { }
// Noise is already applied through the channel model
clear_old_packets(ptr->received_packets, timestamp_processed);
}
bool apply_global_noise = get_noise_power_dBFS() != INVALID_DBFS_VALUE;
if (apply_global_noise) {
for (int a = 0; a < nbAnt; a++) { for (int a = 0; a < nbAnt; a++) {
sample_t *out = (sample_t *)samplesVoid[a];
for (int i = 0; i < nsamps; i++) { for (int i = 0; i < nsamps; i++) {
out[i].r += lroundf(temp_array[a][i].r); int16_t noise_power = (int16_t)(32767.0 / powf(10.0, .05 * -get_noise_power_dBFS()));
out[i].i += lroundf(temp_array[a][i].i); temp_array[a][i].r += noise_power + gaussZiggurat(0.0, 1.0);
temp_array[a][i].i += noise_power * gaussZiggurat(0.0, 1.0);
} }
} }
} else if (num_chanmod_channels > 0) { }
// Apply noise from global setting
int16_t noise_power = (int16_t)(32767.0 / powf(10.0, .05 * -get_noise_power_dBFS()));
for (int a = 0; a < nbAnt; a++) { for (int a = 0; a < nbAnt; a++) {
sample_t *out = (sample_t *)samplesVoid[a]; sample_t *out = (sample_t *)samplesVoid[a];
for (int i = 0; i < nsamps; i++) { for (int i = 0; i < nsamps; i++) {
out[i].r += lroundf(temp_array[a][i].r + noise_power * gaussZiggurat(0.0, 1.0)); out[i].r += lroundf(temp_array[a][i].r);
out[i].i += lroundf(temp_array[a][i].i + noise_power * gaussZiggurat(0.0, 1.0)); out[i].i += lroundf(temp_array[a][i].i);
}
} }
} }
...@@ -1345,6 +1351,8 @@ static void rfsimulator_end(openair0_device *device) ...@@ -1345,6 +1351,8 @@ static void rfsimulator_end(openair0_device *device)
if (b->conn_sock >= 0) if (b->conn_sock >= 0)
removeCirBuf(s, b); removeCirBuf(s, b);
} }
clear_beam_queue(&s->beam_ctrl->tx, INT64_MAX);
clear_beam_queue(&s->beam_ctrl->rx, INT64_MAX);
close(s->epollfd); close(s->epollfd);
free(s); free(s);
} }
...@@ -1391,10 +1399,6 @@ extern "C" __attribute__((__visibility__("default"))) int device_init(openair0_d ...@@ -1391,10 +1399,6 @@ extern "C" __attribute__((__visibility__("default"))) int device_init(openair0_d
if (rfsimulator->prop_delay_ms > 0.0) if (rfsimulator->prop_delay_ms > 0.0)
rfsimulator->chan_offset = ceil(rfsimulator->sample_rate * rfsimulator->prop_delay_ms / 1000); rfsimulator->chan_offset = ceil(rfsimulator->sample_rate * rfsimulator->prop_delay_ms / 1000);
if (rfsimulator->chan_offset != 0) { if (rfsimulator->chan_offset != 0) {
if (CirSize < minCirSize + rfsimulator->chan_offset) {
CirSize = minCirSize + rfsimulator->chan_offset;
LOG_I(HW, "CirSize = %lu\n", CirSize);
}
rfsimulator->prop_delay_ms = rfsimulator->chan_offset * 1000 / rfsimulator->sample_rate; rfsimulator->prop_delay_ms = rfsimulator->chan_offset * 1000 / rfsimulator->sample_rate;
LOG_I(HW, "propagation delay %f ms, %lu samples\n", rfsimulator->prop_delay_ms, rfsimulator->chan_offset); LOG_I(HW, "propagation delay %f ms, %lu samples\n", rfsimulator->prop_delay_ms, rfsimulator->chan_offset);
} }
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment