// 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). #define _GNU_SOURCE // pthread_setaffinity_np, cpu_set_t #include "common.h" #include #include #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 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 uint32_t last_seq = 0; static int seq_init = 0; // Разделяемых на запись счётчиков нет. main читает их раз в секунду для // телеметрии без синхронизации: на 32-бит ARM возможен косметический «рваный» // показ 64-бит значения — на итоговую строку (печатается после join) не влияет. static fec decoder = NULL; static int fec_on = 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: seq-детектор + RS + CRC-поверх-RS + запись. // Порядок сохранён (fq FIFO, единственный писатель в stdout — этот поток). 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; // Детектор пропусков seq (техдолг #6). Кадр, выброшенный переполнением fq, // тоже проявится здесь как дыра seq → в `lost`; счётчик FQdrop разделяет // причины (эфирная потеря против переполнения очереди). uint32_t seq = ((uint32_t)hdr[2] << 24) | ((uint32_t)hdr[3] << 16) | ((uint32_t)hdr[4] << 8) | hdr[5]; if (seq_init && seq > last_seq + 1) { lost += seq - last_seq - 1; fprintf(stderr, "\n[ПОТЕРЯ] seq %u→%u, дыра %u кадров\n", last_seq, seq, seq - last_seq - 1); } last_seq = seq; seq_init = 1; uint16_t original_len = (hdr[6] << 8) | hdr[7]; uint8_t last_block_bytes = hdr[10]; 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) return; uint8_t recovered[4096]; 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++; return; } 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 не справился, кадр битый — не пишем return; } if (!pay_valid) rs_saved++; // liquid CRC по payload не прошёл, а наш // CRC поверх RS — да: значит RS реально спас fwrite(recovered, 1, original_len, stdout); fflush(stdout); total_bytes += original_len; crc_final++; } else if (pay_valid && original_len > 0 && original_len <= pay_len) { fwrite(pay, 1, original_len, stdout); fflush(stdout); total_bytes += original_len; crc_final++; } } // Привязка потока к ядру: ядро 0 — только sync, ядро 1 — захват+FEC+IO // (доктрина README §2). Ошибка не фатальна — планировщик разведёт сам. static void pin_to_cpu(int cpu) { cpu_set_t cs; CPU_ZERO(&cs); CPU_SET(cpu, &cs); if (pthread_setaffinity_np(pthread_self(), sizeof(cs), &cs) != 0) fprintf(stderr, "[warn] не удалось привязать поток к ядру %d\n", cpu); } // Ждать данных в 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 + запись из очереди кадров. static void *fec_thread(void *arg) { (void)arg; pin_to_cpu(1); for (;;) { if (!wait_slot(&fq.mtx, &fq.not_empty, &fq.count, &fq.eof)) break; 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); } 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"; int opt; while ((opt = getopt(argc, argv, "f:r:b:u:c: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 'h': print_usage(argv[0], "ПРИЁМНИК"); fprintf(stderr, "\nСПЕЦИФИЧНЫЕ ОПЦИИ RX:\n"); fprintf(stderr, " Усиление фиксированное: %.0f дБ\n", RX_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(); print_config("ПРИЁМНИК", freq, rate, bw, fec_on); fprintf(stderr, "Усиление: %.0f дБ (ручное)\n", RX_GAIN); fprintf(stderr, "CFO фильтр: α=%.2f\n", CFO_FILTER_ALPHA); fprintf(stderr, "URI: %s\n", uri); fprintf(stderr, "Конвейер: 3 стадии (A захват / B sync / C FEC+IO)\n"); fprintf(stderr, "═══════════════════════════════════════════════════\n"); struct iio_context *ctx = NULL; struct iio_device *phy = NULL, *rx = NULL; 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); fprintf(stderr, "Приёмник готов. Ожидание сигнала...\n\n"); fprintf(stderr, "Кадры | Загол. | CRC | RS испр | RS бит | Потери | Overr | FQdrp | Вывод | 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 | %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)lost, (unsigned long long)ring.overruns, (unsigned long long)fq.drops, (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 | Потери seq: %llu | Overrun: %llu | FQ drops: %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)lost, (unsigned long long)ring.overruns, (unsigned long long)fq.drops, (unsigned long long)crc_final, (unsigned long long)total_bytes); if (decoder) fec_destroy(decoder); ofdmflexframesync_destroy(fs); iio_buffer_destroy(buf); iio_context_destroy(ctx); return 0; }