Files
Pluto-SDR/src/receiver.c
2026-07-16 14:30:17 +03:00

908 lines
47 KiB
C
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// receiver.c — OFDM-приёмник с RS(255,223) FEC и частотной коррекцией
// Трёхстадийный конвейер (README §12.5, BENCHMARK «Разбивка по стадиям»):
// поток A (ядро 1): iio_buffer_refill + int16→float → кольцо сэмплов
// поток B (ядро 0): ТОЛЬКО ofdmflexframesync_execute; callback кладёт сырой
// кадр в очередь fq (RS с него снят — он насыщал ядро 0)
// поток C (ядро 1): из fq → seq-детектор + RS-декод + CRC + запись
// Замер: sync сам 1.65 Msps < 1.92, RS 26% эфира. Вынос RS на ядро 1 роняет
// p_min ~8 мс → ~1.52 мс. -p 0 на этом PHY недостижим (осознанно, §12.5).
#include "common.h" // _GNU_SOURCE/pthread/sched — уже подключены здесь
#include "group_fec.h"
#include "frame_tx.h" // phy_tx_frame — отправка STATUS обратного канала
#include "feedback.h" // fb_status_t / pack — телеметрия B→A (этап 2, §10 п.7)
#define BUF_SAMPLES 32768
#define RING_SLOTS 16 // 16 × 17 мс = 273 мс поглощения всплесков
#define FQ_SLOTS 32 // очередь кадров B→C, ~147 КБ .bss
#define FQ_PAY_MAX (18 * RS_ENC) // 4590: ≤18 блоков RS (техдолг #2)
#define RX_GAIN 50.0
#define FB_TX_GAIN -20.0 // Усиление feedback-TX по умолч. (§10 п.7)
#define CFO_FILTER_ALPHA 0.95f
#define CFO_MAX_VALID 0.5f
extern volatile sig_atomic_t stop;
// --- Кольцо сэмплов A→B (SPSC) ------------------------------------------
// tail пишет только A, head — только B; count/eof/overruns под мьютексом.
typedef struct {
liquid_float_complex data[BUF_SAMPLES];
int ns; // фактическое число отсчётов (частичный refill)
} slot_t;
static slot_t ring_slots[RING_SLOTS]; // 4 МБ .bss
static struct {
int head, tail, count, eof;
uint64_t overruns; // refill'ы, выброшенные из-за полного кольца
pthread_mutex_t mtx;
pthread_cond_t not_empty;
} ring = {
.mtx = PTHREAD_MUTEX_INITIALIZER,
.not_empty = PTHREAD_COND_INITIALIZER,
};
// --- Очередь кадров B→C (SPSC) ------------------------------------------
// B (producer) копирует сюда сырой кадр и НИКОГДА не блокируется (drop-newest),
// поэтому C, выйдя раньше, не заклинит B. tail пишет только B, head — только C.
typedef struct {
uint8_t hdr[HDR_SIZE]; // 12 Б, парсит C
uint8_t pay[FQ_PAY_MAX]; // сырой RS-payload
unsigned pay_len;
int pay_valid; // liquid-CRC по payload
} fq_slot_t;
static fq_slot_t fq_slots[FQ_SLOTS];
static struct {
int head, tail, count, eof;
uint64_t drops; // переполнение очереди + негабаритные кадры
pthread_mutex_t mtx;
pthread_cond_t not_empty;
} fq = {
.mtx = PTHREAD_MUTEX_INITIALIZER,
.not_empty = PTHREAD_COND_INITIALIZER,
};
static struct iio_buffer *g_buf = NULL;
static ofdmflexframesync g_fs = NULL;
// Счётчики B (пишет callback): calls/hdr_ok/crc_raw + телеметрия.
static uint64_t calls = 0, hdr_ok = 0, crc_raw = 0;
static float rssi = -100, evm = 0, cfo_est = 0.0f;
// Счётчики C (пишет process_frame): целостность и вывод.
static uint64_t rs_saved = 0, crc_final = 0, total_bytes = 0;
static uint64_t rs_fail = 0, lost = 0;
static uint64_t grp_recovered = 0, grp_failed = 0; // group-FEC: восст. кадров / провал. групп
static uint64_t sig_bad = 0; // hdr[0..1] не совпал с известной сигнатурой
static uint32_t highest_seq = 0; // максимальный увиденный seq (не откатывается назад)
static int seq_init = 0;
static uint32_t next_out = 0; // курсор вывода: следующий seq в stdout (этап 3b)
static uint64_t rb_over = 0; // кадров отброшено переполнением окна reorder
static int gf_seen = 0; // в эфире group-FEC (автодетект по hdr[11])
// Разделяемых на запись счётчиков нет. main читает их раз в секунду для
// телеметрии без синхронизации: на 32-бит ARM возможен косметический «рваный»
// показ 64-бит значения — на итоговую строку (печатается после join) не влияет.
static fec decoder = NULL;
static int fec_on = 0;
// --- Missing-tracker обратного канала (этап 2, §10 п.7) -----------------
// Отсортированный по возрастанию набор seq отсутствующих кадров данных. Пишется
// ТОЛЬКО потоком C (в process_frame и maybe_send_status) — синхронизации нет,
// как у прочего C-состояния. Дыры seq → miss_add (по возрастанию), приход кадра
// (в т.ч. будущий ретрансмит) → miss_ack. В STATUS уходят FB_MAX_NACK старейших.
#define MISS_CAP 512
static uint32_t miss[MISS_CAP];
static int miss_cnt = 0;
static uint64_t miss_over = 0; // дыр не влезло в MISS_CAP
static void miss_add(uint32_t seq) {
if (miss_cnt >= MISS_CAP) { miss_over++; return; } // старейшие дыры важнее
miss[miss_cnt++] = seq; // highest_seq только растёт → массив сортирован
}
// Вставка с сохранением порядка и дедупом — для seq, который МЕНЬШЕ фронта:
// битый кадр (rs_fail) пришёл, дырой в seq не выглядит, но дослать его надо.
// miss_add для такого seq сломал бы сортировку (в STATUS уходят FB_MAX_NACK
// первых элементов — они обязаны быть старейшими).
static void miss_add_sorted(uint32_t seq) {
int i;
for (i = 0; i < miss_cnt; i++) {
if (miss[i] == seq) return; // уже в списке
if (miss[i] > seq) break; // точка вставки
}
if (miss_cnt >= MISS_CAP) { miss_over++; return; }
memmove(&miss[i + 1], &miss[i], (size_t)(miss_cnt - i) * sizeof(miss[0]));
miss[i] = seq;
miss_cnt++;
}
static void miss_ack(uint32_t seq) {
for (int i = 0; i < miss_cnt; i++) {
if (miss[i] == seq) {
memmove(&miss[i], &miss[i + 1], (size_t)(miss_cnt - i - 1) * sizeof(miss[0]));
miss_cnt--;
return;
}
}
}
// --- EOF-хендшейк (этап 3, поток C) -------------------------------------
// END (F0 E7) несёт итог передачи: сколько кадров и байт ушло в эфир. Без него
// хвостовые потери молчаливы — у последних кадров нет следующего seq, по которому
// детектор видит дыру. Приняв END, RX добивает недостающие seq в NACK, а собрав
// всё — поднимает FB_FLAG_COMPLETE в STATUS и завершается сам.
#define ARQ_COMPLETE_HOLD_MS 600 // держать COMPLETE ≥6 периодов STATUS: в
// тихом эфире END-фазы доставка ~100%, но
// повтор дешевле, чем таймаут A на 10 с
static int end_seen = 0;
static uint32_t end_total_seq = 0;
static uint64_t end_total_bytes = 0;
static uint64_t end_rx = 0; // принято валидных END (диагностика)
static int64_t complete_since_ms = 0; // 0 = файл ещё не собран
// --- Отправитель STATUS на feedback-TX (только при -F, поток C) ----------
// Свой framegen + буфер на txfb (868), свой RS-энкодер. Кадр STATUS шлётся не
// чаще FB_PERIOD_MS. Данные RX и обратный TX — независимые cf-девайсы (FDD),
// поэтому push сюда не конфликтует с приёмом.
#define FB_BUF_SAMPLES 16384
#define FB_PERIOD_MS 100
#define FB_AMP 0.2f
static ofdmflexframegen g_fb_fg = NULL;
static struct iio_buffer *g_fb_buf = NULL;
static fec g_fb_enc = NULL;
static uint32_t g_fb_seq = 0;
static int64_t g_fb_last_ms = 0;
static const uint8_t SIG_STATUS[2] = { FB_SIG_0, FB_SIG_1 };
static int64_t now_ms(void) {
struct timespec t;
clock_gettime(CLOCK_MONOTONIC, &t);
return (int64_t)t.tv_sec * 1000 + t.tv_nsec / 1000000;
}
// Файл собран? Legacy: курсор вывода дошёл до конца — это сильнее, чем «дыр
// нет»: next_out двигается только по непрерывному префиксу реально записанных
// кадров. Group-FEC: ARQ выключен (этап 3c), вывод идёт мимо кольца — критерий
// прежний, по фронту; ветка уйдёт вместе с гибридом GF+ARQ.
// Держим флаг ARQ_COMPLETE_HOLD_MS, чтобы A успел услышать хотя бы один STATUS,
// затем завершаемся штатным путём через stop (main → join A→B→C).
static void arq_check_complete(void) {
if (!g_fb_fg || !end_seen) return; // симплекс либо передача ещё идёт
int done = (end_total_seq == 0) ||
(miss_cnt == 0 &&
(gf_seen ? (seq_init && highest_seq == end_total_seq - 1)
: next_out == end_total_seq));
if (!done) { complete_since_ms = 0; return; }
int64_t t = now_ms();
if (!complete_since_ms) {
complete_since_ms = t; // с этого момента STATUS несёт COMPLETE
fprintf(stderr, "\n[ARQ] COMPLETE: %u кадров, %llu байт — держу флаг %d мс\n",
end_total_seq, (unsigned long long)end_total_bytes,
ARQ_COMPLETE_HOLD_MS);
}
if (t - complete_since_ms >= ARQ_COMPLETE_HOLD_MS) {
fprintf(stderr, "[ARQ] COMPLETE отправлен — завершаю приём\n");
stop = 1;
}
}
static void maybe_send_status(void) {
if (!g_fb_fg) return; // симплекс — обратного канала нет
arq_check_complete(); // до гейта периода: решение о завершении
// не должно ждать своего слота STATUS
int64_t t = now_ms();
if (g_fb_last_ms && t - g_fb_last_ms < FB_PERIOD_MS) return;
g_fb_last_ms = t;
fb_status_t st;
memset(&st, 0, sizeof(st));
st.type = FB_TYPE_STATUS;
st.flags = (end_seen ? FB_FLAG_END_SEEN : 0) |
(complete_since_ms ? FB_FLAG_COMPLETE : 0);
st.highest_seq = seq_init ? highest_seq : 0xFFFFFFFFu;
st.rx_frames = (uint32_t)crc_final;
st.overruns = (uint32_t)ring.overruns; // телеметрия, гонка косметическая
int cnt = miss_cnt < FB_MAX_NACK ? miss_cnt : FB_MAX_NACK;
st.nack_cnt = (uint8_t)cnt;
for (int i = 0; i < cnt; i++) st.nack[i] = miss[i];
uint8_t pay[FB_STATUS_MAXLEN];
size_t n = fb_status_pack(&st, pay);
phy_tx_frame(g_fb_fg, g_fb_buf, FB_BUF_SAMPLES, g_fb_enc,
SIG_STATUS, g_fb_seq++, 0, pay, n, FB_AMP, 0);
}
// callback потока B — лёгкий: телеметрия + копия сырого кадра в очередь C.
static int callback(unsigned char *hdr, int hdr_valid,
unsigned char *pay, unsigned int pay_len,
int pay_valid, framesyncstats_s stats,
void *user) {
(void)user;
calls++;
rssi = stats.rssi;
evm = stats.evm;
// CFO — только телеметрия: ofdmflexframesync сам компенсирует смещение
// по преамбуле. Внешний NCO-контур убран (техдолг #4).
if (fabsf(stats.cfo) < CFO_MAX_VALID)
cfo_est = CFO_FILTER_ALPHA * cfo_est + (1.0f - CFO_FILTER_ALPHA) * stats.cfo;
if (!hdr_valid) return 0;
hdr_ok++;
if (pay_valid) crc_raw++;
// Копия кадра в очередь C. Негабарит/полная очередь → дроп со счётчиком
// (drop-newest: B не блокируется). Разбор seq/RS/CRC — уже в потоке C.
int oversize = (pay_len > FQ_PAY_MAX);
pthread_mutex_lock(&fq.mtx);
int reject = oversize || (fq.count == FQ_SLOTS);
if (reject) fq.drops++;
pthread_mutex_unlock(&fq.mtx);
if (reject) return 0;
fq_slot_t *fs = &fq_slots[fq.tail]; // tail трогает только producer B
memcpy(fs->hdr, hdr, HDR_SIZE);
memcpy(fs->pay, pay, pay_len);
fs->pay_len = pay_len;
fs->pay_valid = pay_valid;
pthread_mutex_lock(&fq.mtx);
fq.tail = (fq.tail + 1) % FQ_SLOTS;
fq.count++;
pthread_cond_signal(&fq.not_empty);
pthread_mutex_unlock(&fq.mtx);
return 0;
}
// --- Кольцо переупорядочивания вывода (поток C, этап 3b) ----------------
// До ARQ кадры приходили строго по возрастанию seq и писались в stdout сразу.
// Ретрансмит ломает это: кадр 5 приедет после кадра 40 — прямая запись сдвинула
// бы данные и убила md5. Кольцо по seq % RB_CAP + курсор next_out: в stdout
// уходит только непрерывный префикс. Дыра держит вывод, пока ARQ её не закроет.
// len == 0 — «слот занят, но в файл ничего»: паритетный кадр GF тратит seq, но
// данных не несёт (нужно для этапа 3c, чтобы next_out доходил до total_seq).
#define RB_CAP 1024 // ~4 с эфира при 9.3 мс/кадр; TXWIN у A вдвое
// больше — досылка для любого seq в окне есть
typedef struct {
uint32_t seq;
uint16_t len;
uint8_t filled;
uint8_t data[GF_CHUNK];
} rbslot_t;
static rbslot_t rb[RB_CAP]; // ~1 МБ .bss (next_out/rb_over — выше, к ним
// обращается arq_check_complete до этой секции)
// Выдать в stdout непрерывный префикс, начиная с next_out.
static void rb_drain(void) {
int wrote = 0;
for (;;) {
rbslot_t *s = &rb[next_out % RB_CAP];
if (!s->filled || s->seq != next_out) break; // дыра — дальше нельзя
if (s->len) {
fwrite(s->data, 1, s->len, stdout);
total_bytes += s->len;
crc_final++;
}
s->filled = 0;
next_out++;
wrote = 1;
}
if (wrote) fflush(stdout);
}
// Принудительно продвинуть курсор до seq `to`, схлопывая дыры: то, что есть —
// в stdout, чего нет — пропало навсегда (прежняя семантика «потеря = сдвиг»).
static void rb_force_to(uint32_t to) {
while (next_out < to) {
rbslot_t *s = &rb[next_out % RB_CAP];
if (s->filled && s->seq == next_out && s->len) {
fwrite(s->data, 1, s->len, stdout);
total_bytes += s->len;
crc_final++;
}
s->filled = 0;
next_out++;
}
fflush(stdout);
}
static void rb_insert(uint32_t seq, const uint8_t *d, size_t len) {
if (len > GF_CHUNK) len = GF_CHUNK; // защита в глубину: len берётся из
// hdr, слот кольца — ровно GF_CHUNK
if (seq < next_out) return; // уже выведен: дубль ретрансмита
if (seq - next_out >= RB_CAP) {
// Дыра в голове старше окна. С ARQ это не тупик: дропаем новый кадр и
// возвращаем его seq в NACK — A дошлёт, когда голова разъедется. Без
// обратного канала досылать некому → схлопываем окно вперёд.
if (g_fb_fg) {
rb_over++;
miss_add_sorted(seq);
return;
}
rb_force_to(seq - RB_CAP + 1);
}
rbslot_t *s = &rb[seq % RB_CAP];
s->seq = seq;
s->len = (uint16_t)len;
s->filled = 1;
if (len) memcpy(s->data, d, len);
}
// Досдать всё, что осталось в кольце (поток C на выходе): приём кончился, ARQ
// уже не поможет — дыры схлопываются, хвост уходит в файл.
static void rb_flush_all(void) {
uint32_t last = next_out;
for (int i = 0; i < RB_CAP; i++)
if (rb[i].filled && rb[i].seq >= last) last = rb[i].seq + 1;
rb_force_to(last);
}
// --- Групповой стирающий FEC (поток C, README §12.7) --------------------
// Собирает группу из 16 кадров данных + 2 паритетных (P,Q). Пропавший на эфире
// или битый по CRC кадр = стирание; ≤2 стираний на группу закрываются gf_recover.
// Всё состояние — только поток C (как остальные C-счётчики), синхронизации нет.
static struct {
int active;
uint32_t base; // seq первого кадра группы (= seq idx)
int k; // число кадров данных (из меты паритета)
int k_known; // k получен из паритета
int maxidx; // макс. полученный idx данных — оценка k без паритета
uint8_t data[GF_K][GF_CHUNK];
uint16_t dlen[GF_K]; // фактическая длина каждого кадра данных
uint8_t have[GF_K]; // 1 = кадр idx получен и валиден
uint8_t P[GF_CHUNK], Q[GF_CHUNK];
int haveP, haveQ;
uint16_t tail_len; // длина последнего кадра данных (из меты)
} g;
static void group_reset(uint32_t base) {
memset(g.have, 0, sizeof(g.have));
memset(g.dlen, 0, sizeof(g.dlen));
g.active = 1; g.base = base; g.k = 0; g.k_known = 0;
g.maxidx = -1; g.haveP = g.haveQ = 0; g.tail_len = 0;
}
// Финализировать текущую группу: восстановить стирания и записать данные по idx.
static void group_flush(void) {
if (!g.active) return;
int k = g.k_known ? g.k : (g.maxidx + 1);
if (k <= 0 || k > GF_K) { g.active = 0; return; }
int erasures = 0;
for (int i = 0; i < k; i++) if (!g.have[i]) erasures++;
if (erasures > 0) {
if (gf_recover(g.data, g.have, k, g.P, g.haveP, g.Q, g.haveQ)) {
grp_recovered += erasures;
for (int i = 0; i < k; i++) if (!g.have[i]) {
// прочие кадры полные 1024, последний — tail_len
g.dlen[i] = (i == k - 1 && g.tail_len) ? g.tail_len : GF_CHUNK;
g.have[i] = 1;
}
} else {
grp_failed++; // >2 стираний либо не хватило паритета
}
}
for (int i = 0; i < k; i++) {
if (!g.have[i]) continue; // невосстановимая дыра
uint16_t len = g.dlen[i] > GF_CHUNK ? GF_CHUNK : g.dlen[i];
fwrite(g.data[i], 1, len, stdout);
total_bytes += len;
crc_final++;
}
fflush(stdout);
g.active = 0;
}
// Разложить один декодированный group-FEC кадр по группе.
static void group_route(uint8_t flags, uint32_t seq,
const uint8_t *dec, size_t dec_len, int ok) {
int idx = flags & GF_IDX_MASK;
int is_par = flags & GF_FLAG_PAR;
uint32_t base = seq - idx; // idx=k для P, k+1 для Q → та же база
if (!g.active || base != g.base) { // новая группа → финализировать прежнюю
group_flush();
group_reset(base);
}
if (!ok) return; // битый кадр остаётся стиранием
if (!is_par) {
if (idx < GF_K && dec_len <= GF_CHUNK) {
memcpy(g.data[idx], dec, dec_len);
if (dec_len < GF_CHUNK) memset(g.data[idx] + dec_len, 0, GF_CHUNK - dec_len);
g.dlen[idx] = (uint16_t)dec_len;
g.have[idx] = 1;
if (idx > g.maxidx) g.maxidx = idx;
}
} else if (dec_len >= GF_PAR_PAY) { // паритет: [0..1] мета, [2..] P/Q
uint16_t meta = ((uint16_t)dec[0] << 8) | dec[1];
g.k = GF_META_K(meta);
g.tail_len = GF_META_TAIL(meta);
g.k_known = 1;
if (flags & GF_FLAG_Q) { memcpy(g.Q, dec + GF_META, GF_CHUNK); g.haveQ = 1; }
else { memcpy(g.P, dec + GF_META, GF_CHUNK); g.haveP = 1; }
}
}
// Разбор END-кадра (поток C, этап 3). Зафиксировать итог передачи и добить
// хвостовые дыры в NACK. Повторы END (A шлёт их каждые 100 мс до COMPLETE)
// идемпотентны: состояние берём с первого валидного.
static void process_end(const uint8_t *hdr, uint8_t *pay, unsigned pay_len,
int pay_valid) {
static uint8_t edec[RS_ENC]; // END = 12 Б → ровно 1 RS-блок (223 Б)
fb_end_t e;
if (fec_on && decoder) {
size_t elen = 0;
if (!phy_rx_decode(hdr, pay, pay_len, decoder, edec, sizeof(edec), &elen))
return; // RS/CRC не прошли — ждём следующий END
if (fb_end_unpack(edec, elen, &e) != 0) return;
} else { // -c none: payload как есть
uint16_t olen = (uint16_t)((hdr[6] << 8) | hdr[7]);
if (!pay_valid || olen < FB_END_LEN || olen > pay_len) return;
if (fb_end_unpack(pay, olen, &e) != 0) return;
}
end_rx++;
if (end_seen) return;
end_seen = 1;
end_total_seq = e.total_seq;
end_total_bytes = e.total_bytes;
// Backfill: всё от текущего фронта до total_seq-1 в эфире было, но до нас не
// дошло. Массив miss остаётся отсортированным — добавляем seq строго больше
// всех имеющихся. Фронт двигаем на конец файла: дальше приходить могут только
// ретрансмиты (seq ≤ highest_seq), ложных дыр они не породят.
if (e.total_seq > 0 && (!seq_init || highest_seq < e.total_seq - 1)) {
uint32_t from = seq_init ? highest_seq + 1 : 0;
for (uint32_t m = from; m < e.total_seq; m++) miss_add(m);
lost += e.total_seq - from;
seq_init = 1;
highest_seq = e.total_seq - 1;
}
fprintf(stderr, "\n[ARQ] END принят: всего %u кадров, %llu байт; не хватает %d\n",
e.total_seq, (unsigned long long)e.total_bytes, miss_cnt);
}
// Разбор одного кадра в потоке C: seq-детектор + RS + CRC-поверх-RS, затем
// маршрутизация (прямая запись для legacy, сборка группы для group-FEC).
// Единственный писатель в stdout — этот поток, порядок сохранён (fq FIFO).
static void process_frame(fq_slot_t *fs) {
const uint8_t *hdr = fs->hdr;
uint8_t *pay = fs->pay; // не const: fec_decode() liquid не const-correct
unsigned pay_len = fs->pay_len;
int pay_valid = fs->pay_valid;
// END (F0 E7) — маркер конца передачи. Ветвимся ДО seq-детектора: у END свой
// счётчик seq, в пространстве seq данных его быть не должно.
if (hdr[0] == HDR_SIGNATURE_0 && hdr[1] == HDR_SIG_END_1) {
process_end(hdr, pay, pay_len, pay_valid);
return;
}
// Проверка сигнатуры hdr[0..1]: liquid-CRC (hdr_valid) страхует только
// целостность байт, но не то, что это наш протокол. Отсеивает чужие типы
// кадров (roadmap §10 п.7) и мусор, прошедший CRC заголовка.
if (hdr[0] != HDR_SIGNATURE_0 || hdr[1] != HDR_SIGNATURE_1) {
sig_bad++;
return;
}
// Детектор пропусков seq (техдолг #6) на основе МАКСИМАЛЬНОГО увиденного
// seq, а не последнего: повтор/переупорядоченный кадр с seq ≤ highest_seq
// (ретрансмит ARQ, §10 п.7) не должен откатывать состояние назад и рождать
// ложную дыру на следующем свежем кадре. Кадр, выброшенный переполнением fq,
// тоже проявится здесь как дыра seq → в `lost`; счётчик FQdrop разделяет
// причины (эфирная потеря против переполнения очереди). В group-FEC потеря
// одновременно фиксируется как стирание в group_route по idx.
uint32_t seq = ((uint32_t)hdr[2] << 24) | ((uint32_t)hdr[3] << 16) |
((uint32_t)hdr[4] << 8) | hdr[5];
if (!seq_init) {
seq_init = 1;
highest_seq = seq;
if (seq > 0) { // приём начался не с нуля: кадры до
for (uint32_t m = 0; m < seq; m++) miss_add(m); // первого пойманного
lost += seq; // потеряны, а дырой не выглядят —
fprintf(stderr, "\n[ПОТЕРЯ] старт с seq %u, дыра %u кадров в начале\n",
seq, seq); // им не предшествовал никакой seq
}
} else if (seq > highest_seq + 1) {
// Пропавшие seq → NACK-набор обратного канала (этап 2) + счётчик потерь.
for (uint32_t m = highest_seq + 1; m < seq; m++) miss_add(m);
lost += seq - highest_seq - 1;
fprintf(stderr, "\n[ПОТЕРЯ] seq %u→%u, дыра %u кадров\n",
highest_seq, seq, seq - highest_seq - 1);
highest_seq = seq;
} else if (seq > highest_seq) {
highest_seq = seq;
}
// seq ≤ highest_seq: старый/повторный кадр — highest_seq не трогаем.
// Снятие seq с NACK — НЕ здесь, а после декода (см. метку route): факт
// приёма кадра ещё не значит, что он целый.
uint16_t original_len = (hdr[6] << 8) | hdr[7];
uint8_t last_block_bytes = hdr[10];
uint8_t flags = hdr[11];
// Декод payload → dec/dec_len, ok=1 при пройденном CRC-поверх-RS.
static uint8_t recovered[4096]; // поток C один → static безопасен
const uint8_t *dec = NULL;
size_t dec_len = 0;
int ok = 0;
if (fec_on && decoder && pay_len >= RS_ENC) {
size_t nblocks = pay_len / RS_ENC;
// Защита в глубину: B уже отсёк pay_len>FQ_PAY_MAX, но guard оставлен
// (recovered[4096] вмещает 4096/RS_DATA = 18 блоков, техдолг #2).
if (nblocks == 0 || nblocks > 4096 / RS_DATA) goto route;
for (size_t b = 0; b < nblocks; b++) {
// fec_decode() liquid всегда "успешен" — код возврата бесполезен.
fec_decode(decoder, RS_DATA, pay + b * RS_ENC,
recovered + b * RS_DATA);
}
// Настоящий критерий целостности — CRC поверх декода (техдолг #3):
// fec_decode отдаёт мусор молча, поэтому проверяем сами.
size_t framed_len = (nblocks - 1) * RS_DATA +
(last_block_bytes > 0 ? last_block_bytes : RS_DATA);
if (framed_len > nblocks * RS_DATA) framed_len = nblocks * RS_DATA;
if (original_len == 0 || (size_t)original_len + 4 > framed_len) {
rs_fail++;
goto route;
}
unsigned int rx_crc = crc_generate_key(LIQUID_CRC_32, recovered, original_len);
unsigned int tx_crc = ((uint32_t)recovered[original_len] << 24) |
((uint32_t)recovered[original_len + 1] << 16) |
((uint32_t)recovered[original_len + 2] << 8) |
recovered[original_len + 3];
if (rx_crc != tx_crc) {
rs_fail++; // RS не справился, кадр битый — не пишем
goto route;
}
if (!pay_valid) rs_saved++; // liquid CRC по payload не прошёл, а наш
// CRC поверх RS — да: значит RS реально спас
dec = recovered; dec_len = original_len; ok = 1;
} else if (pay_valid && original_len > 0 && original_len <= pay_len) {
dec = pay; dec_len = original_len; ok = 1;
}
route:
// Учёт NACK — ТОЛЬКО после декода. Раньше seq снимался с NACK по факту
// приёма кадра, и битый кадр (rs_fail) выпадал из списка навсегда: ARQ его
// уже не дослал бы, а RX считал бы файл собранным. Битый кадр — такая же
// дыра, как непришедший, просто обнаруженная не детектором seq.
if (ok) miss_ack(seq);
else miss_add_sorted(seq);
if (flags & GF_FLAG_MODE) {
gf_seen = 1;
group_route(flags, seq, dec, dec_len, ok);
} else {
if (g.active) group_flush(); // TX не мешает режимы, но подстрахуемся
if (ok) {
rb_insert(seq, dec, dec_len); // порядок восстанавливает кольцо:
rb_drain(); // ретрансмит приходит вне очереди
}
}
}
// pin_to_cpu — теперь в common.{h,c} (используется также будущим потоком
// обратного канала на TX, roadmap §10 п.7).
// Ждать данных в SPSC-кольце (100 мс страховка от lost wakeup). Возврат:
// 1 — есть слот (мьютекс m ЗАБЛОКИРОВАН на выходе), 0 — eof/stop и пусто (m разблокирован).
static int wait_slot(pthread_mutex_t *m, pthread_cond_t *cv, int *count, int *eof) {
pthread_mutex_lock(m);
while (*count == 0 && !*eof && !stop) {
struct timespec ts;
clock_gettime(CLOCK_REALTIME, &ts);
ts.tv_nsec += 100000000L;
if (ts.tv_nsec >= 1000000000L) { ts.tv_sec++; ts.tv_nsec -= 1000000000L; }
pthread_cond_timedwait(cv, m, &ts);
}
if (*count == 0) { pthread_mutex_unlock(m); return 0; }
return 1;
}
// Поток A (ядро 1): блокирующий refill + конвертация int16→float в слот кольца.
static void *acq_thread(void *arg) {
(void)arg;
pin_to_cpu(1);
while (!stop) {
ssize_t n = iio_buffer_refill(g_buf);
if (n < 0) {
if (stop) break; // EINTR при сигнале
usleep(500);
continue;
}
pthread_mutex_lock(&ring.mtx);
int full = (ring.count == RING_SLOTS);
if (full) ring.overruns++;
pthread_mutex_unlock(&ring.mtx);
if (full) continue; // drop-newest: кольцо полно — выбрасываем refill
slot_t *s = &ring_slots[ring.tail]; // tail трогает только producer A
int16_t *bb = (int16_t *)iio_buffer_start(g_buf);
int ns = (iio_buffer_end(g_buf) - iio_buffer_start(g_buf)) / sizeof(int16_t) / 2;
if (ns > BUF_SAMPLES) ns = BUF_SAMPLES;
for (int i = 0; i < ns; i++) {
s->data[i] = ((float)bb[2*i] / 2048.0f) +
((float)bb[2*i+1] / 2048.0f) * _Complex_I;
}
s->ns = ns;
pthread_mutex_lock(&ring.mtx);
ring.tail = (ring.tail + 1) % RING_SLOTS;
ring.count++;
pthread_cond_signal(&ring.not_empty);
pthread_mutex_unlock(&ring.mtx);
}
pthread_mutex_lock(&ring.mtx);
ring.eof = 1;
pthread_cond_broadcast(&ring.not_empty); // гарантированно будим B
pthread_mutex_unlock(&ring.mtx);
return NULL;
}
// Поток B (ядро 0): только синхронизатор. Callback кладёт кадры в fq.
static void *dsp_thread(void *arg) {
(void)arg;
pin_to_cpu(0);
for (;;) {
if (!wait_slot(&ring.mtx, &ring.not_empty, &ring.count, &ring.eof))
break; // eof/stop и пусто
slot_t *s = &ring_slots[ring.head];
pthread_mutex_unlock(&ring.mtx); // execute — без мьютекса
ofdmflexframesync_execute(g_fs, s->data, s->ns);
pthread_mutex_lock(&ring.mtx);
ring.head = (ring.head + 1) % RING_SLOTS;
ring.count--;
pthread_mutex_unlock(&ring.mtx);
}
pthread_mutex_lock(&fq.mtx); // сигналим C: кадров больше нет
fq.eof = 1;
pthread_cond_broadcast(&fq.not_empty);
pthread_mutex_unlock(&fq.mtx);
return NULL;
}
// Поток C (ядро 1): RS-декод + CRC + запись из очереди кадров + периодический
// STATUS обратного канала. Ждём кадр таймингованно (50 мс), чтобы STATUS уходил
// и в простое (именно тогда TX и нужно знать про дыры). maybe_send_status —
// no-op в симплексе (g_fb_fg==NULL), поэтому поведение без -F не меняется.
static void *fec_thread(void *arg) {
(void)arg;
pin_to_cpu(1);
for (;;) {
pthread_mutex_lock(&fq.mtx);
while (fq.count == 0 && !fq.eof && !stop) {
struct timespec ts;
clock_gettime(CLOCK_REALTIME, &ts);
ts.tv_nsec += 50000000L; // 50 мс
if (ts.tv_nsec >= 1000000000L) { ts.tv_sec++; ts.tv_nsec -= 1000000000L; }
pthread_cond_timedwait(&fq.not_empty, &fq.mtx, &ts);
if (fq.count == 0 && !fq.eof && !stop) { // всё ещё пусто → STATUS без кадра
pthread_mutex_unlock(&fq.mtx);
maybe_send_status();
pthread_mutex_lock(&fq.mtx);
}
}
if (fq.count == 0) { pthread_mutex_unlock(&fq.mtx); break; } // eof/stop
fq_slot_t *fs = &fq_slots[fq.head];
pthread_mutex_unlock(&fq.mtx);
process_frame(fs);
pthread_mutex_lock(&fq.mtx);
fq.head = (fq.head + 1) % FQ_SLOTS;
fq.count--;
pthread_mutex_unlock(&fq.mtx);
maybe_send_status(); // под нагрузкой — период гейтит частоту
}
group_flush(); // дописать последнюю группу файла (нет кадра-триггера)
rb_flush_all(); // и хвост кольца: приём кончился, дыры схлопываем —
return NULL; // досылать их уже некому
}
int main(int argc, char *argv[]) {
long long freq = DEFAULT_FREQ, rate = DEFAULT_RATE, bw = DEFAULT_BW;
const char *uri = RX_URI_DEFAULT;
const char *fec_opt = "rs8";
long long fb_freq = 0; // 0 = симплекс; иначе FDD, LO обратного TX
double fb_gain = FB_TX_GAIN; // усиление feedback-TX (-g)
int opt;
while ((opt = getopt(argc, argv, "f:r:b:u:c:F:g:h")) != -1) {
switch (opt) {
case 'f': freq = atoll(optarg); break;
case 'r': rate = atoll(optarg); break;
case 'b': bw = atoll(optarg); break;
case 'u': uri = optarg; break;
case 'c': fec_opt = optarg; break;
case 'F': fb_freq = atoll(optarg); break;
case 'g': fb_gain = atof(optarg); break;
case 'h':
print_usage(argv[0], "ПРИЁМНИК");
fprintf(stderr, "\nСПЕЦИФИЧНЫЕ ОПЦИИ RX:\n");
fprintf(stderr, " Усиление фиксированное: %.0f дБ\n", RX_GAIN);
fprintf(stderr, " -F ЧАСТОТА FDD: частота обратного TX в Гц (§10 п.7). Без -F — симплекс\n");
fprintf(stderr, " -g УСИЛЕНИЕ Усиление feedback-TX в дБ (по умолч. %.0f)\n", FB_TX_GAIN);
fprintf(stderr, " Частотная коррекция: адаптивная (α=%.2f)\n", CFO_FILTER_ALPHA);
fprintf(stderr, "\nОПТИМАЛЬНЫЙ ЗАПУСК:\n");
fprintf(stderr, " %s -f %lld -r %lld -b %lld -c rs8 > output.bin\n",
argv[0], DEFAULT_FREQ, DEFAULT_RATE, DEFAULT_BW);
return 0;
default: return 1;
}
}
fec_on = !strcmp(fec_opt, "rs8");
setup_signal_handlers();
gf_init(); // group-FEC автодетектится по hdr[11] первого кадра
print_config("ПРИЁМНИК", freq, rate, bw, fec_on);
fprintf(stderr, "Усиление: %.0f дБ (ручное)\n", RX_GAIN);
fprintf(stderr, "CFO фильтр: α=%.2f\n", CFO_FILTER_ALPHA);
fprintf(stderr, "Дуплекс: %s\n", fb_freq
? "FDD (обратный канал: передача STATUS на 868)" : "симплекс");
if (fb_freq) fprintf(stderr, "Обр. TX LO: %.1f МГц, усил. %.0f дБ\n", fb_freq / 1e6, fb_gain);
fprintf(stderr, "URI: %s\n", uri);
fprintf(stderr, "Конвейер: 3 стадии (A захват / B sync / C FEC+IO)\n");
fprintf(stderr, "═══════════════════════════════════════════════════\n");
// С -F открываем оба cf-девайса и настраиваем оба LO (обратный TX пока только
// сконфигурирован; передача STATUS/END по нему — этап 2).
struct iio_context *ctx = NULL;
struct iio_device *phy = NULL, *rx = NULL, *txfb = NULL;
if (fb_freq) {
if (pluto_init_fdd(uri, &ctx, &phy, &txfb, &rx) != 0) return 1;
// данные RX на freq (915), обратный TX на fb_freq (868)
if (pluto_configure_fdd(phy, fb_freq, freq, rate, bw, fb_gain, RX_GAIN) != 0) {
LOG_ERROR("Настройка FDD-радиотракта не удалась");
iio_context_destroy(ctx);
return 1;
}
} else {
if (pluto_init(uri, &ctx, &phy, &rx, 0) != 0) return 1;
if (pluto_configure(phy, freq, rate, bw, RX_GAIN, 0) != 0) {
LOG_ERROR("Настройка радиотракта не удалась");
iio_context_destroy(ctx);
return 1;
}
}
iio_channel_enable(iio_device_find_channel(rx, "voltage0", false));
iio_channel_enable(iio_device_find_channel(rx, "voltage1", false));
// Глубже DMA-очередь ядра (по умолчанию 4): 8 × 17 мс = 137 мс запаса,
// кроет задержки A от RS-всплесков C на общем ядре 1.
if (iio_device_set_kernel_buffers_count(rx, 8) != 0)
fprintf(stderr, "[warn] не удалось увеличить kernel_buffers_count\n");
struct iio_buffer *buf = iio_device_create_buffer(rx, BUF_SAMPLES, false);
if (!buf) {
fprintf(stderr, "Ошибка создания буфера приёма\n");
iio_context_destroy(ctx);
return 1;
}
// Инициализация FEC
if (fec_on) {
decoder = fec_create(LIQUID_FEC_RS_M8, NULL);
if (!decoder) {
fprintf(stderr, "Ошибка создания FEC-декодера\n");
}
}
ofdmflexframesync fs = ofdmflexframesync_create(
OFDM_M, OFDM_CP, OFDM_TAPER, NULL, callback, NULL);
if (!fs) {
fprintf(stderr, "Ошибка создания OFDM-синхронизатора\n");
if (decoder) fec_destroy(decoder);
iio_buffer_destroy(buf);
iio_context_destroy(ctx);
return 1;
}
// Тот же 12-байтовый заголовок, что и на TX (иначе байты 8..11 — мусор)
ofdmflexframesync_set_header_len(fs, HDR_SIZE);
// Обратный канал (этап 2, §10 п.7): framegen + буфер на feedback-TX (868).
// Кадры STATUS (F0 55) шлёт поток C через maybe_send_status. set_header_len
// 12 обязателен и здесь — иначе TX-сторона прочтёт заголовок как мусор.
if (fb_freq && txfb) {
iio_channel_enable(iio_device_find_channel(txfb, "voltage0", true));
iio_channel_enable(iio_device_find_channel(txfb, "voltage1", true));
g_fb_buf = iio_device_create_buffer(txfb, FB_BUF_SAMPLES, false);
if (!g_fb_buf) {
fprintf(stderr, "[warn] буфер обратного канала не создан — STATUS отключён\n");
} else {
ofdmflexframegenprops_s fp;
ofdmflexframegenprops_init_default(&fp);
fp.check = LIQUID_CRC_32;
g_fb_fg = ofdmflexframegen_create(OFDM_M, OFDM_CP, OFDM_TAPER, NULL, &fp);
ofdmflexframegen_set_header_len(g_fb_fg, HDR_SIZE);
g_fb_enc = fec_create(LIQUID_FEC_RS_M8, NULL);
fprintf(stderr, "Обратный канал: STATUS каждые %d мс (F0 55, RS+CRC32)\n",
FB_PERIOD_MS);
}
}
fprintf(stderr, "Приёмник готов. Ожидание сигнала...\n\n");
fprintf(stderr, "Кадры | Загол. | CRC | RS испр | RS бит | SigBd | Потери | Overr | FQdrp | Rec | GrpF | Вывод | RSSI | EVM | CFO\n");
fprintf(stderr, "--------+--------+------+---------+--------+-------+--------+--------+--------+--------+--------+--------+-------+-------+----------\n");
// Три стадии на двух ядрах. main крутит статистику 1 Гц и ждёт stop, затем
// join A→B→C: A ставит ring.eof, B дочищает кольцо и ставит fq.eof, C
// дочищает очередь. Producer'ы (A, B) не блокируются → цепочка без дедлоков.
g_buf = buf;
g_fs = fs;
pthread_t ta, tb, tc;
pthread_create(&tc, NULL, fec_thread, NULL);
pthread_create(&tb, NULL, dsp_thread, NULL);
pthread_create(&ta, NULL, acq_thread, NULL);
time_t last = time(NULL);
while (!stop) {
struct timespec nap = { .tv_sec = 0, .tv_nsec = 200000000L }; // 200 мс
nanosleep(&nap, NULL);
time_t now = time(NULL);
if (now > last) {
fprintf(stderr, "\r%6llu | %5llu | %4llu | %6llu | %5llu | %5llu | %6llu | %6llu | %6llu | %6llu | %6llu | %6llu | %+5.0f | %+5.1f | %+8.4f",
(unsigned long long)calls, (unsigned long long)hdr_ok,
(unsigned long long)crc_raw, (unsigned long long)rs_saved,
(unsigned long long)rs_fail, (unsigned long long)sig_bad,
(unsigned long long)lost,
(unsigned long long)ring.overruns, (unsigned long long)fq.drops,
(unsigned long long)grp_recovered, (unsigned long long)grp_failed,
(unsigned long long)crc_final, rssi, evm, cfo_est);
fflush(stderr);
last = now;
}
}
pthread_join(ta, NULL);
pthread_join(tb, NULL);
pthread_join(tc, NULL);
fprintf(stderr, "\n\nПРИЁМ ЗАВЕРШЁН\n");
fprintf(stderr, "Кадров: %llu | Заголовков: %llu | CRC: %llu | RS испр: %llu | RS битых: %llu | Сигн.битых: %llu | Потери seq: %llu | Overrun: %llu | FQ drops: %llu | GF восст: %llu | GF провал: %llu | Вывод: %llu | Байт: %llu\n",
(unsigned long long)calls, (unsigned long long)hdr_ok,
(unsigned long long)crc_raw, (unsigned long long)rs_saved,
(unsigned long long)rs_fail, (unsigned long long)sig_bad,
(unsigned long long)lost,
(unsigned long long)ring.overruns, (unsigned long long)fq.drops,
(unsigned long long)grp_recovered, (unsigned long long)grp_failed,
(unsigned long long)crc_final, (unsigned long long)total_bytes);
if (fb_freq) {
fprintf(stderr, "Обратный канал: отправлено STATUS: %u\n", g_fb_seq);
fprintf(stderr, "ARQ: END принято: %llu | ожидалось кадров: %u | выведено seq: %u | не хватает: %d | COMPLETE: %s\n",
(unsigned long long)end_rx, end_total_seq, next_out, miss_cnt,
complete_since_ms ? "да" : "НЕТ");
fprintf(stderr, "ARQ: окно reorder переполнено: %llu | NACK переполнен: %llu\n",
(unsigned long long)rb_over, (unsigned long long)miss_over);
}
if (decoder) fec_destroy(decoder);
ofdmflexframesync_destroy(fs);
iio_buffer_destroy(buf);
if (g_fb_fg) ofdmflexframegen_destroy(g_fb_fg);
if (g_fb_enc) fec_destroy(g_fb_enc);
if (g_fb_buf) iio_buffer_destroy(g_fb_buf);
iio_context_destroy(ctx);
return 0;
}