// 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.5–2 мс. -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 uint64_t miss_rebuilt = 0; // пересборок NACK из кольца (этап 3c) // Разделяемых на запись счётчиков нет. 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; } // Файл собран? Критерий один для обоих режимов (этап 3c): курсор вывода дошёл до // конца. Это сильнее, чем «дыр нет»: next_out двигается только по непрерывному // префиксу реально записанных кадров. Паритетные кадры group-FEC занимают слот // нулевой длины — они тратят seq, но в файл не идут, и курсор их проходит. // Держим флаг 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 && 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); } // Вызывается ТОЛЬКО на FDD-пути (симплекс пишет в stdout напрямую — там нечего // переупорядочивать), поэтому досылка по NACK всегда доступна. 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) { // Дыра в голове старше окна: дропаем новый кадр и возвращаем его seq в // NACK — A дошлёт, когда голова разъедется. Схлопнуть окно нельзя: это // молча потеряло бы кадр, который ARQ ещё способен закрыть. rb_over++; miss_add_sorted(seq); return; } 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); } // Пересобрать NACK-набор из кольца (этап 3c, только FDD, только по END). // Зачем: miss_add при переполнении MISS_CAP молча теряет seq, и вернуть его туда // уже некому — детектор дыр по нему не сработает второй раз (highest_seq ушёл // вперёд). В тяжёлом канале это тупик: A не знает, что дослать, next_out стоит, // COMPLETE не наступает, B досиживает до сторожа. Кольцо же знает точно, каких // кадров нет: слот либо занят своим seq, либо это дыра. // Скан ограничен окном кольца: seq дальше next_out+RB_CAP всё равно принять // некуда (rb_insert их отбросит), просить их у A рано. static void miss_rebuild(void) { if (!g_fb_fg || !end_seen) return; if (end_total_seq <= next_out) { miss_cnt = 0; return; } // всё выведено; заодно // страховка от переполнения вычитания ниже uint32_t to = end_total_seq; if (to - next_out > RB_CAP) to = next_out + RB_CAP; miss_cnt = 0; for (uint32_t s = next_out; s < to && miss_cnt < MISS_CAP; s++) { rbslot_t *sl = &rb[s % RB_CAP]; if (sl->filled && sl->seq == s) continue; // кадр уже принят miss[miss_cnt++] = s; // возрастает → сортировка цела } miss_rebuilt++; } // --- Групповой стирающий FEC (поток C, README §12.7) -------------------- // Собирает группу из 16 кадров данных + 2 паритетных (P,Q). Пропавший на эфире // или битый по CRC кадр = стирание; ≤2 стираний на группу закрываются gf_recover. // Всё состояние — только поток C (как остальные C-счётчики), синхронизации нет. static struct { int active; int seen; // хоть одна группа была: g.base осмыслен (этап 3c) 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 fixed[GF_K]; // 1 = кадр восстановлен паритетом, а не принят 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.fixed, 0, sizeof(g.fixed)); memset(g.dlen, 0, sizeof(g.dlen)); g.active = 1; g.seen = 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; g.fixed[i] = 1; // закрыт паритетом — досылать по ARQ не нужно } } else { grp_failed++; // >2 стираний либо не хватило паритета } } if (g_fb_fg) { // FDD (гибрид GF+ARQ, этап 3c): порядок держит кольцо, и принятые кадры // уже легли в него на приёме (group_route). Отсюда идут ТОЛЬКО кадры, // поднятые паритетом, — и их seq снимаются с NACK: GF закрыл дыру, // ретрансмит был бы лишним трафиком. Невосстановимые остаются в miss — // их дошлёт ARQ, для того он здесь и включён. for (int i = 0; i < k; i++) { if (!g.fixed[i]) continue; rb_insert(g.base + i, g.data[i], g.dlen[i]); miss_ack(g.base + i); } rb_drain(); } else { // Симплекс: кольца нет (ретрансмитов не будет — переупорядочивать нечего), // группа сама себе буфер. Пишем её целиком в порядке idx, восстановленные // встают на свои места. Невосстановимая дыра = сдвиг, прежняя семантика. 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 → та же база // FDD: кольцо — единственный источник порядка, и оно живёт независимо от // сборки групп. Кладём кадр в слот ДО всякого bypass: даже ретрансмит давно // закрытой группы обязан занять своё место, иначе next_out встанет на нём // навсегда и COMPLETE не наступит. Паритет тратит seq, но данных не несёт — // ему слот нулевой длины, курсор его проходит насквозь. if (ok && g_fb_fg) { rb_insert(seq, is_par ? NULL : dec, is_par ? 0 : dec_len); rb_drain(); } // Ретрансмит кадра уже закрытой группы (этап 3c). Пересобирать нечего: в FDD // кадр только что лёг в кольцо, а группа нужна лишь ради паритета. Ловим два // случая: старая группа (base < g.base) — иначе она преждевременно flush-нула // бы активную и сбросила её стирания; и текущая, уже финализированная по END // (base == g.base && !g.active) — иначе кадр «воскресил» бы группу с одним // элементом, наврав в grp_failed. if (g.seen && (base < g.base || (base == g.base && !g.active))) return; 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) { end_seen = 1; end_total_seq = e.total_seq; end_total_bytes = e.total_bytes; // Хвостовые потери молчаливы: у последних кадров нет следующего seq, по // которому детектор увидел бы разрыв. Двигаем фронт на конец файла — // дальше приходить могут только ретрансмиты (seq ≤ highest_seq), ложных // дыр они не породят. if (e.total_seq > 0 && (!seq_init || highest_seq < e.total_seq - 1)) { lost += e.total_seq - (seq_init ? highest_seq + 1 : 0); seq_init = 1; highest_seq = e.total_seq - 1; } } // Досрочно финализировать последнюю группу (group-FEC + FDD). Штатно группа // закрывается приходом первого кадра СЛЕДУЮЩЕЙ группы — а для последней его // не будет, и без этого её восстановление паритетом случилось бы только при // выходе потока C, то есть уже после COMPLETE. Дыры пришлось бы закрывать // ретрансмитом вместо готового паритета. END значит «новых данных не будет» // (дальше только повторы), поэтому финализировать безопасно; повторные END // попадают в no-op (g.active уже 0). group_flush(); // NACK пересобираем на КАЖДОМ END, а не только на первом (этап 3c). A шлёт // END каждые 100 мс до самого COMPLETE, так что это ещё и штатный способ // вернуть в запрос seq, выпавшие из miss[] по MISS_CAP: пока хвост пуст, // дыры видит только кольцо. Скан ≤1024 при 10 Гц — сотые доли ядра. miss_rebuild(); if (end_rx == 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) { group_route(flags, seq, dec, dec_len, ok); } else if (ok) { if (g.active) group_flush(); // TX не мешает режимы, но подстрахуемся if (g_fb_fg) { rb_insert(seq, dec, dec_len); // FDD: ретрансмит приходит вне rb_drain(); // очереди — порядок держит кольцо } else { // Симплекс: ретрансмитов нет, seq всегда растёт — переупорядочивать // нечего. Кольцо здесь только задержало бы вывод до EOF (дыра стопорит // next_out, а закрыть её некому). Пишем сразу, как до этапа 3b. fwrite(dec, 1, dec_len, stdout); fflush(stdout); total_bytes += dec_len; crc_final++; } } } // 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 | пересборок NACK: %llu\n", (unsigned long long)rb_over, (unsigned long long)miss_over, (unsigned long long)miss_rebuilt); } 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; }