Repository navigation
Expand file tree
/
Copy pathrender.cpp
More file actions
1860 lines (1646 loc) · 71.2 KB
/
Copy pathrender.cpp
File metadata and controls
1860 lines (1646 loc) · 71.2 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
// RiseTogether Bela forwarder.
//
// The Bela is the hardware hub of the experiment. Everything it sees or emits
// lands on one clock -- the audio frame counter (which bela runs at 48 kHz)
// -- so a latency between two parties is a frame difference, with no
// cross-device clock involved and nothing to correct for.
//
// Raspberry Pi ---trigger--> [ Bela ] ---forwarded verbatim---> EEG amp
// ---jittered ~1 Hz timer--> EEG amp
// iPads (Polly, Pia, ...) ---photodiode---> [ Bela ]
//
// The forwarded trigger is a sample-accurate level mirror of the Pi's line,
// delayed by exactly one block: the PRU writes block k's output word during
// block k+1. The timer trigger is deliberately jittered -- a metronomic pulse
// train would beat against the EEG and inject a correlated artefact -- and both
// its schedule and its RNG seed are logged, so the record stays exact.
//
// A device's photodiode edge can be aligned with the final recorded, collated
// data the ipads record when the photodiode trigger was requested, and then
// when the Bela log is aligned with the coordinator device logs (which uses LSL
// to connect to all of the ipads), the entire sequence can be reconstructed
// referenced to the coordinator's LSL monotonic clock.
//
// LSL lives elsewhere in the stack now. The library, the headers and every line
// of inlet code are kept behind ENABLE_LSL, which is 0.
#include "include/font.h"
#include <Bela.h>
#include "include/linux_i2c.c"
#include <algorithm>
#include <atomic>
#include <cctype>
#include <cerrno>
#include <cstdint>
#include <cstdio>
#include <cstring>
#include <ctime>
#include <fstream>
#include <include/ssd1306.c>
#include <iomanip>
#include <mutex>
#include <random>
#include <sstream>
#include <string>
#include <sys/stat.h>
#include <sys/time.h>
#include <thread>
#include <unistd.h>
#include <vector>
// ---------------------------------------------------------------------------
// Compile-time configuration
// ---------------------------------------------------------------------------
#define ENABLE_LSL 0 // 0 = pins and triggers only; 1 restores the inlet
#define USE_OLED_DISPLAY 0
#if ENABLE_LSL
#include <include/lsl_cpp.h>
#endif
#if USE_OLED_DISPLAY
static const int kOledI2cDev = 1;
#endif
// Digital inputs. `device` names the physical source -- an iPad for a
// photodiode, the Pi for the trigger line -- and is what makes a row in
// edges.csv mean something months later. `active_level` is the logic level that
// means "asserted". `refractory_ms` only drives the `accepted` flag in the log;
// no edge is ever discarded because of it (see the burst guard in render()). It
// is a duration, not a frame count, so it means the same thing whatever rate
// the board comes up at -- gRefractoryFrames[] below holds the conversion, done
// once in setup() against the real digital sample rate.
//
// An edge is `accepted` when the pin had been quiet -- no edge of either
// polarity -- for at least `refractory_ms` before it (meta "accept_rule":
// "quiet_before"). For a photodiode that makes an accepted rising edge exactly
// a flash onset. An iPad backlight is PWM-dimmed (~480 Hz on the lab panels,
// dark gaps under 1 ms), so one 100 ms white square arrives as ~50 pulses; the
// window has to be longer than that gap and shorter than the shortest dark run
// the marker channel can produce, a Manchester half-bit of 3 samples at 120 Hz
// (25 ms, ~17 ms after a frame of display jitter). 10 ms sits between the two.
// Until 2026-09-17 the window was 1 ms and counted from the last *accepted*
// edge, which accepted a PWM edge every millisecond of every flash.
struct InputPin {
unsigned int pin;
const char* role; // "photodiode" | "trigger_in"
const char* device;
bool active_level;
double refractory_ms;
};
// Add a device by adding a row. Two things are indexed by table position rather
// than by pin number -- the gPinStatesAtomic bitfield and gRefractoryFrames[]
// -- so there is a hard limit of 16 rows, checked in setup() along with the
// rest of the map.
static const InputPin kInputPins[] = {
{0, "photodiode", "Pia", true, 10.0}, // active-HIGH comparator
{1, "photodiode", "Parsnip", true, 10.0},
{2, "photodiode", "Polly", true, 10.0},
{3, "photodiode", "Peter", true, 10.0},
{4, "photodiode", "Padme", true, 10.0},
{5, "photodiode", "Patrick", true, 10.0},
{11, "trigger_in", "rpi", true, 0.0}, // no refractory: never mask a trigger
};
static const size_t kNumInputPins = sizeof(kInputPins) / sizeof(kInputPins[0]);
static const size_t kMaxInputPins = 16;
// Which row above carries the Raspberry Pi trigger. setup() checks that this
// index really does name a "trigger_in" row.
static const size_t kTriggerInIndex = 6;
// Outputs to the EEG amp. These must not also appear in kInputPins.
static const unsigned int kTimerOutPin = 12; // local jittered trigger
static const unsigned int kFwdOutPin = 13; // level mirror of the Pi trigger
static const bool kOutIdleLevel = false; // idle LOW, pulse HIGH
// Timer trigger. Interval = period + U(-jitter, +jitter), scheduled from the
// previous *scheduled* frame rather than the emitted one, so the mean rate
// stays exactly 1000/kTimerPeriodMs Hz and quantisation never accumulates.
static const double kTimerPeriodMs = 1000.0;
static const double kTimerJitterMs = 200.0;
static const double kTimerPulseMs = 100.0;
// Both outputs are held idle for this long after the first block: the PRU needs
// a moment to settle, and the startup pin scan below wants a quiet bus.
static const double kStartupHoldMs = 250.0;
// If a pin produces more than this many edges without a quiet refractory
// window between them, stop emitting rows for it until it has been quiet for a
// full refractory period. Suppressed edges are counted and carried on the next
// emitted row, so the record stays lossless in aggregate while staying bounded
// in the worst case.
//
// Sized for PWM light, not for chatter: with quiet-before acceptance every PWM
// edge inside a lit stretch is non-accepted. A 100 ms flash is ~100 edges at
// 480 Hz, but the 2026-09-16 recording also has lit stretches of 1.2-1.3 s with
// no 10 ms gap on several pins (~1300 edges), which a smaller limit would have
// cut from the log. 8192 is ~8.5 s of 480 Hz PWM. The queue is not the
// constraint: even a pin toggling every sample is 960 edges per 20 ms drain,
// against 16384 slots.
static const uint32_t kMaxEdgesPerBurst = 8192;
#if ENABLE_LSL
static const char* kStreamPrefixFilter = "LSLTest";
static const double kPullTimeoutSec = 0.05; // blocking pull; see below
static const double kTimeCorrIntervalSec = 1.0;
static const double kTimeCorrTimeoutSec = 0.2; // fast after the first call
static const double kTimeCorrFirstTimeoutSec = 5.0;
static const int32_t kInletBufLenSec = 10;
static const int32_t kInletChunkLen = 1; // no chunking: minimum latency
static const bool kInletRecover = true;
#endif
// Fallback only, for the window in which context->digitalSampleRate has not
// been read yet; setup() overwrites gSampleRate with what the board reports.
// Bela runs at 44100 on most capes and 48000 on the CTAG ones, so nothing here
// may assume either.
static const double kFallbackFs = 44100.0;
static const unsigned int kPeriodSize = 16;
static const size_t kMaxChannels = 8;
static const size_t kMaxIdLen = 64;
// Resolved once in setup() from context->digitalSampleRate, then read-only for
// the rest of the run, so no synchronisation is needed: setup() completes
// before render() or any worker thread starts.
static double gSampleRate = kFallbackFs;
static uint64_t gRefractoryFrames[kMaxInputPins];
static unsigned int gSyncPeriodBlocks = 1;
// Timer-trigger geometry, resolved the same way against the reported rate.
static uint64_t gTimerPeriodFrames = 0;
static uint64_t gTimerJitterFrames = 0;
static uint64_t gTimerPulseFrames = 0;
static uint64_t gStartupHoldFrames = 0;
// xorshift32. std::mt19937 allocates and is far too heavy to call from the RT
// thread; this is three shifts and a modulo. The seed is logged, so a session's
// entire trigger schedule can be regenerated offline.
static uint32_t gRngSeed = 0;
static uint32_t gRngState = 0;
static inline int64_t nextJitterFrames() {
gRngState ^= gRngState << 13;
gRngState ^= gRngState >> 17;
gRngState ^= gRngState << 5;
if (gTimerJitterFrames == 0)
return 0;
const uint32_t span = static_cast<uint32_t>(2 * gTimerJitterFrames + 1);
return static_cast<int64_t>(gRngState % span) -
static_cast<int64_t>(gTimerJitterFrames);
}
// Queue depths. Drained every 20 ms by the logger thread.
static const size_t kEdgeQueueSize = 16384;
static const size_t kTriggerQueueSize = 4096;
static const size_t kStatusQueueSize = 1024;
static const size_t kSyncQueueSize = 1024;
static const size_t kLslQueueSize = 2048;
static const size_t kCorrQueueSize = 1024;
static const unsigned int kLoggerPollUs = 20000; // 20 ms
static const double kFlushIntervalS = 1.0; // bound data loss on hard kill
static const unsigned int kDisplayEveryNPolls = 10; // ~200 ms
// ---------------------------------------------------------------------------
// Event records
// ---------------------------------------------------------------------------
struct EdgeEvent {
uint64_t frame; // absolute audio frame of the digital sample
uint64_t edge_index; // per-pin monotonic counter (counts ALL edges)
uint64_t dt_frames_prev; // frames since the previous edge on this pin
uint32_t
suppressed_count; // edges dropped by the burst guard since last row
uint32_t pin_index; // index into kInputPins
uint8_t state; // raw logic level
uint8_t accepted; // 1 if outside this pin's refractory window
EdgeEvent()
: frame(0), edge_index(0), dt_frames_prev(0), suppressed_count(0),
pin_index(0), state(0), accepted(0) {}
};
// One row per level change on an output pin -- the definitive record of what
// the EEG amp actually saw, on the same frame axis as every input edge. A timer
// rising edge carries the jitter drawn for the *next* interval, and how late
// the pulse was if a dropped block delayed it.
enum TriggerSource {
TRIG_TIMER = 0,
TRIG_FORWARD,
};
static const char* triggerSourceName(uint32_t s) {
return s == TRIG_TIMER ? "timer" : "forward";
}
struct TriggerEvent {
uint64_t frame;
uint64_t seq; // per-source counter; rising and falling both
int64_t jitter_frames; // signed; timer rising edges only
int64_t late_frames; // > 0 only when a block gap delayed the pulse
uint32_t source;
uint8_t level;
TriggerEvent()
: frame(0), seq(0), jitter_frames(0), late_frames(0),
source(TRIG_TIMER), level(0) {}
};
enum StatusCode {
ST_SESSION_START = 0,
ST_SESSION_END,
ST_XRUN,
ST_BLOCK_GAP,
ST_QUEUE_FULL,
ST_STREAM_OPEN,
ST_STREAM_LOST,
ST_CLOCK_RESET,
ST_INFO,
};
static const char* statusName(int code) {
switch (code) {
case ST_SESSION_START:
return "SESSION_START";
case ST_SESSION_END:
return "SESSION_END";
case ST_XRUN:
return "XRUN";
case ST_BLOCK_GAP:
return "BLOCK_GAP";
case ST_QUEUE_FULL:
return "QUEUE_FULL";
case ST_STREAM_OPEN:
return "STREAM_OPEN";
case ST_STREAM_LOST:
return "STREAM_LOST";
case ST_CLOCK_RESET:
return "CLOCK_RESET";
default:
return "INFO";
}
}
struct StatusEvent {
uint64_t frame;
double lsl_clock; // 0.0 when pushed from RT (no clock calls there)
int64_t detail_num;
int32_t code;
char detail_str[96]; // only ever filled from non-RT threads
StatusEvent() : frame(0), lsl_clock(0.0), detail_num(0), code(ST_INFO) {
detail_str[0] = '\0';
}
};
struct SyncEvent {
uint64_t frame_req; // frame stored by render() just before scheduling
uint64_t frame_latest; // latest frame render() has reported, read here
double lsl_clock;
double monotonic_clock;
SyncEvent()
: frame_req(0), frame_latest(0), lsl_clock(0.0), monotonic_clock(0.0) {}
};
struct LslEvent {
uint64_t rx_index;
double lsl_timestamp; // raw sender clock (post_none)
double arrival_lsl_clock; // Bela clock, captured immediately after pull
uint64_t arrival_frame; // Bela frame axis, +/- one block
uint32_t samples_available_before; // 0 => clean, backlog-free arrival
uint32_t n_channels;
double ch[kMaxChannels];
char source_id[kMaxIdLen];
char stream_name[kMaxIdLen];
LslEvent()
: rx_index(0), lsl_timestamp(0.0), arrival_lsl_clock(0.0),
arrival_frame(0), samples_available_before(0), n_channels(0) {
memset(ch, 0, sizeof(ch));
source_id[0] = '\0';
stream_name[0] = '\0';
}
};
struct TimeCorrEvent {
double lsl_clock;
double correction;
double remote_time;
double uncertainty;
uint8_t clock_reset;
uint8_t ok;
char source_id[kMaxIdLen];
TimeCorrEvent()
: lsl_clock(0.0), correction(0.0), remote_time(0.0), uncertainty(0.0),
clock_reset(0), ok(0) {
source_id[0] = '\0';
}
};
// ---------------------------------------------------------------------------
// Lock-free SPSC queue (one producer, one consumer; logger is sole consumer).
// Safe across the Xenomai/Linux boundary: shared memory + C++ atomics only,
// no blocking primitives, and Bela mlockall()s the process.
// ---------------------------------------------------------------------------
template <typename T, size_t Size> class LockFreeSPSCQueue {
private:
alignas(64) std::atomic<size_t> write_pos{0};
alignas(64) std::atomic<size_t> read_pos{0};
T buffer[Size];
public:
bool push(const T& item) {
size_t current_write = write_pos.load(std::memory_order_relaxed);
size_t next_write = (current_write + 1) % Size;
if (next_write == read_pos.load(std::memory_order_acquire))
return false;
buffer[current_write] = item;
write_pos.store(next_write, std::memory_order_release);
return true;
}
bool pop(T& item) {
size_t current_read = read_pos.load(std::memory_order_relaxed);
if (current_read == write_pos.load(std::memory_order_acquire))
return false;
item = buffer[current_read];
read_pos.store((current_read + 1) % Size, std::memory_order_release);
return true;
}
size_t size_approx() const {
size_t w = write_pos.load(std::memory_order_relaxed);
size_t r = read_pos.load(std::memory_order_relaxed);
return (w >= r) ? (w - r) : (Size - r + w);
}
};
static LockFreeSPSCQueue<EdgeEvent, kEdgeQueueSize> gEdgeQueue;
static LockFreeSPSCQueue<TriggerEvent, kTriggerQueueSize> gTriggerQueue;
static LockFreeSPSCQueue<StatusEvent, kStatusQueueSize> gStatusQueueRT;
static LockFreeSPSCQueue<SyncEvent, kSyncQueueSize> gSyncQueue;
#if ENABLE_LSL
static LockFreeSPSCQueue<LslEvent, kLslQueueSize> gLslQueue;
static LockFreeSPSCQueue<TimeCorrEvent, kCorrQueueSize> gCorrQueue;
#endif
// Status records originate from several threads (render, LSL, setup/cleanup),
// which would break the single-producer contract above. The RT path keeps its
// own lock-free queue; every non-RT producer uses this mutex-guarded vector,
// which is safe because none of those threads are real-time.
static std::mutex gStatusMutex;
static std::vector<StatusEvent> gStatusPending;
// Push a status record from a non-RT thread.
static void pushStatus(int code, uint64_t frame, double lsl_clock,
int64_t detail_num, const char* detail_str) {
StatusEvent e;
e.code = code;
e.frame = frame;
e.lsl_clock = lsl_clock;
e.detail_num = detail_num;
if (detail_str) {
strncpy(e.detail_str, detail_str, sizeof(e.detail_str) - 1);
e.detail_str[sizeof(e.detail_str) - 1] = '\0';
}
std::lock_guard<std::mutex> lock(gStatusMutex);
gStatusPending.push_back(e);
}
// ---------------------------------------------------------------------------
// Shared state
// ---------------------------------------------------------------------------
// Latest frame render() has finished processing. This is the *end* of the most
// recent block: the PRU has already sampled digital input for that whole block
// by the time render() runs, so blockStart + audioFrames is the best estimate
// of "hardware now".
static std::atomic<uint64_t> gElapsedFramesAtomic{0};
// Frame stored by render() immediately before scheduling the clock-sync task.
// The gap between this and gElapsedFramesAtomic when the task actually runs
// bounds the task's scheduling latency -- logged as bracket_frames.
static std::atomic<uint64_t> gSyncRequestFrame{0};
static std::atomic<bool> gRunning{true}; // producers (LSL thread)
static std::atomic<bool> gLoggerRunning{true}; // consumer (logger thread)
static std::atomic<uint16_t> gPinStatesAtomic{0};
static std::atomic<uint64_t> gDroppedEvents{0};
static AuxiliaryTask gClockSyncTask;
// Session paths, fixed at startup.
static std::string gSessionDir;
static std::string gStem;
// ---------------------------------------------------------------------------
// Helpers
// ---------------------------------------------------------------------------
static double monotonicNow() {
struct timespec ts;
clock_gettime(CLOCK_MONOTONIC, &ts);
return static_cast<double>(ts.tv_sec) +
1e-9 * static_cast<double>(ts.tv_nsec);
}
static double nowClock() {
#if ENABLE_LSL
return lsl::local_clock();
#else
return monotonicNow();
#endif
}
static std::string csvEscape(const std::string& in) {
bool needs = in.find_first_of(",\"\n\r") != std::string::npos;
if (!needs)
return in;
std::string out = "\"";
for (char c : in) {
if (c == '"')
out += "\"\"";
else if (c == '\n' || c == '\r')
out += ' ';
else
out += c;
}
out += "\"";
return out;
}
#if ENABLE_LSL
static std::string jsonEscape(const std::string& in) {
std::string out = "\"";
for (char c : in) {
if (c == '"' || c == '\\') {
out += '\\';
out += c;
} else if (c == '\n' || c == '\r' || c == '\t')
out += ' ';
else
out += c;
}
out += "\"";
return out;
}
#endif // ENABLE_LSL
// Ten significant digits after the point: LSL clocks are seconds since an
// arbitrary recent epoch, so this preserves sub-nanosecond resolution.
static std::string fmtClock(double v) {
std::ostringstream os;
os << std::fixed << std::setprecision(9) << v;
return os.str();
}
static bool makeSessionDir() {
struct timeval tv;
gettimeofday(&tv, nullptr);
time_t t = tv.tv_sec;
struct tm tmv;
localtime_r(&t, &tmv);
char stamp[32];
strftime(stamp, sizeof(stamp), "%Y%m%d_%H%M%S", &tmv);
std::random_device rd;
std::mt19937 gen(rd());
std::uniform_int_distribution<> dis(1000, 9999);
mkdir("./logs", 0755); // fine if it already exists
gStem = std::string(stamp) + "_" + std::to_string(dis(gen));
gSessionDir = std::string("./logs/") + gStem;
if (mkdir(gSessionDir.c_str(), 0755) != 0 && errno != EEXIST) {
rt_printf("Error: could not create session directory %s\n",
gSessionDir.c_str());
return false;
}
return true;
}
static std::string sessionPath(const char* suffix) {
return gSessionDir + "/" + gStem + suffix;
}
// ---------------------------------------------------------------------------
// Real-time path
// ---------------------------------------------------------------------------
struct PinRT {
bool prev_state;
bool initialised;
uint64_t last_edge_frame;
uint64_t last_accepted_frame;
uint64_t edge_index;
uint32_t burst_count;
uint32_t suppressed_count;
PinRT()
: prev_state(false), initialised(false), last_edge_frame(0),
last_accepted_frame(0), edge_index(0), burst_count(0),
suppressed_count(0) {}
};
static PinRT gPinRT[kNumInputPins];
// RT-safe status push. No string formatting, no clock call.
static void pushStatusRT(int code, uint64_t frame, int64_t detail_num) {
StatusEvent e;
e.code = code;
e.frame = frame;
e.lsl_clock = 0.0;
e.detail_num = detail_num;
if (!gStatusQueueRT.push(e)) {
gDroppedEvents.fetch_add(1, std::memory_order_relaxed);
}
}
// RT-safe trigger push. render() is the sole producer, which is what preserves
// the queue's single-producer contract.
static void pushTriggerRT(uint32_t source, uint64_t frame, uint64_t seq,
bool level, int64_t jitter_frames,
int64_t late_frames) {
TriggerEvent e;
e.source = source;
e.frame = frame;
e.seq = seq;
e.level = level ? 1 : 0;
e.jitter_frames = jitter_frames;
e.late_frames = late_frames;
if (!gTriggerQueue.push(e)) {
gDroppedEvents.fetch_add(1, std::memory_order_relaxed);
}
}
void clockSyncTask(void*) {
SyncEvent e;
e.frame_req = gSyncRequestFrame.load(std::memory_order_relaxed);
// Read the clocks first, then the frame counter, so frame_latest is never
// earlier than the instant the clocks were read -- keeps the bracket a true
// upper bound on staleness.
e.lsl_clock = nowClock();
e.monotonic_clock = monotonicNow();
e.frame_latest = gElapsedFramesAtomic.load(std::memory_order_relaxed);
if (!gSyncQueue.push(e)) {
gDroppedEvents.fetch_add(1, std::memory_order_relaxed);
}
}
// Survey of every digital pin over the first 100 ms of running, reported once.
// digitalRead() is only meaningful once the audio thread has sampled a block,
// so this cannot live in setup(); and one block is 16 frames (0.33 ms), which
// is far too short a window to tell an idle level from a signal, and is also
// the block most likely to be garbage if the PRU has not settled. Every pin is
// reported, watched or not, so a signal on a pin nobody is listening to shows
// up as such instead of looking like a dead sensor -- which is the usual
// symptom of a pin map that does not match the wiring.
static const unsigned int kScanFrames = 4800; // 100 ms at 44.1-48 kHz
struct PinScan {
unsigned int high_frames;
unsigned int edges;
bool first_state;
bool prev_state;
};
static PinScan gScan[16];
static unsigned int gScanFrames = 0;
static unsigned int gScanBlocks = 0;
static bool gScanDone = false;
static void reportDigitalPinScan(BelaContext* context, unsigned int nPins) {
rt_printf("--- digital pin scan: %u frames (%.1f ms) over %u blocks ---\n",
gScanFrames, 1000.0 * gScanFrames / context->digitalSampleRate,
gScanBlocks);
rt_printf("pin first now high%% edges role\n");
unsigned int totalEdges = 0;
unsigned int totalHigh = 0;
for (unsigned int pin = 0; pin < nPins; pin++) {
const PinScan& c = gScan[pin];
totalEdges += c.edges;
totalHigh += c.high_frames;
char role[48];
snprintf(role, sizeof(role), "(unwatched)");
for (size_t p = 0; p < kNumInputPins; p++) {
if (kInputPins[p].pin == pin) {
snprintf(role, sizeof(role), "%s / %s", kInputPins[p].role,
kInputPins[p].device);
break;
}
}
// Our own outputs read back through the same word, so the level shown
// for them is what we are driving, not a measurement of anything.
if (pin == kFwdOutPin)
snprintf(role, sizeof(role), "OUT trigger_fwd (we drive this)");
else if (pin == kTimerOutPin)
snprintf(role, sizeof(role), "OUT trigger_timer (we drive this)");
rt_printf("%3u %-5s %-5s %5u %5u %s\n", pin,
c.first_state ? "HIGH" : "LOW", c.prev_state ? "HIGH" : "LOW",
(100 * c.high_frames) / gScanFrames, c.edges, role);
}
// Every pin reading a flat zero for 100 ms is far more likely to mean the
// PRU never filled the digital buffer than to mean every input really is
// grounded
// -- an unconnected Bela input floats and an idle active-LOW comparator
// sits HIGH, so a true all-LOW reading takes deliberate wiring.
if (totalHigh == 0 && totalEdges == 0) {
rt_printf(
"*** All %u pins read LOW for the whole window with no edges. "
"That is the signature of a\n",
nPins);
rt_printf("*** digital buffer the PRU never wrote, not of a wiring "
"fault. Check for 'PRU interrupt\n");
rt_printf("*** timeout' or 'McASP error' above before reading anything "
"into these levels.\n");
}
rt_printf("--- end pin scan ---\n");
}
// Accumulate one block into the scan; prints and latches off once the window is
// full. Called every block until then.
static void updateDigitalPinScan(BelaContext* context) {
if (gScanDone)
return;
const unsigned int nPins =
context->digitalChannels < 16 ? context->digitalChannels : 16;
if (nPins == 0 || context->digitalFrames == 0) {
gScanDone = true;
rt_printf("*** NO DIGITAL I/O: digitalChannels=%u digitalFrames=%u "
"digitalSampleRate=%.0f\n",
context->digitalChannels, context->digitalFrames,
context->digitalSampleRate);
rt_printf("*** digitalRead() cannot return anything and no edge will "
"ever be logged. Run line needs --use-digital 1.\n");
return;
}
for (unsigned int n = 0; n < context->digitalFrames; n++) {
for (unsigned int pin = 0; pin < nPins; pin++) {
const bool state = digitalRead(context, n, pin);
PinScan& c = gScan[pin];
if (gScanFrames == 0) {
c.first_state = state;
c.prev_state = state;
c.high_frames = 0;
c.edges = 0;
} else if (state != c.prev_state) {
c.edges++;
c.prev_state = state;
}
if (state)
c.high_frames++;
}
gScanFrames++;
if (gScanFrames >= kScanFrames)
break;
}
gScanBlocks++;
if (gScanFrames >= kScanFrames) {
gScanDone = true;
reportDigitalPinScan(context, nPins);
}
}
void render(BelaContext* context, void* userData) {
const uint64_t blockStart = context->audioFramesElapsed;
const uint64_t blockEnd = blockStart + context->audioFrames;
// --- block continuity: a dropped block means digital frames were never
// read, so an input edge could be missing entirely and an output pulse
// could have gone out late. Without this you cannot tell "the iPad never
// flashed" from "the Bela missed it".
static uint64_t sExpectedFrame = 0;
static bool sFirstBlock = true;
static unsigned int sLastUnderrunCount = 0;
if (sFirstBlock) {
sFirstBlock = false;
sLastUnderrunCount = context->underrunCount;
} else {
if (blockStart != sExpectedFrame) {
pushStatusRT(ST_BLOCK_GAP, blockStart,
static_cast<int64_t>(blockStart) -
static_cast<int64_t>(sExpectedFrame));
}
if (context->underrunCount != sLastUnderrunCount) {
pushStatusRT(ST_XRUN, blockStart,
static_cast<int64_t>(context->underrunCount));
sLastUnderrunCount = context->underrunCount;
}
}
sExpectedFrame = blockEnd;
updateDigitalPinScan(context);
// Re-assert pin direction every block. Directions persist, so this is
// normally a no-op -- but it costs a handful of bit operations and
// guarantees that a stray setting cannot silently corrupt a reading or
// leave an output undriven. Direction lives in bits 0-15 and values in bits
// 16-31, so this cannot disturb what digitalRead() sees.
for (size_t p = 0; p < kNumInputPins; p++) {
pinMode(context, 0, kInputPins[p].pin, INPUT);
}
pinMode(context, 0, kFwdOutPin, OUTPUT);
pinMode(context, 0, kTimerOutPin, OUTPUT);
// Output state. Bela's digital word is bidirectional and the PRU refills it
// with input readings every block, so an output level does NOT persist by
// itself -- it has to be re-driven from frame 0 of every block.
// digitalWrite() covers frame n to the end of the block, so writing the
// held level at frame 0 and again on each change reproduces the waveform
// exactly while touching the word only when something happens.
static bool sFwdLevel = kOutIdleLevel;
static bool sTimerLevel = kOutIdleLevel;
static uint64_t sFwdSeq = 0;
static uint64_t sTimerSeq = 0;
// Both outputs stay idle until the startup hold expires: the PRU's first
// blocks are the ones most likely to be garbage, and a garbage transition
// forwarded to the EEG amp is an event in the recording that never
// happened. Once armed, sNextOnFrame advances by period + jitter from its
// own previous value, never from the frame a pulse actually landed on, so
// quantisation cannot accumulate.
static bool sOutputsArmed = false;
static uint64_t sNextOnFrame = 0;
static uint64_t sOffFrame = 0;
if (!sOutputsArmed && blockStart >= gStartupHoldFrames) {
sOutputsArmed = true;
// Adopt the Pi line's current level rather than waiting for its next
// edge, so the mirror is correct from the first armed frame even if the
// trigger line happens to be asserted right now.
sFwdLevel = gPinRT[kTriggerInIndex].initialised
? gPinRT[kTriggerInIndex].prev_state
: kOutIdleLevel;
sNextOnFrame = blockStart + gTimerPeriodFrames;
pushStatusRT(ST_INFO, blockStart, static_cast<int64_t>(sNextOnFrame));
}
digitalWrite(context, 0, kFwdOutPin, sFwdLevel);
digitalWrite(context, 0, kTimerOutPin, sTimerLevel);
uint16_t states = gPinStatesAtomic.load(std::memory_order_relaxed);
bool statesChanged = false;
for (unsigned int n = 0; n < context->digitalFrames; n++) {
const uint64_t frame = blockStart + n;
for (size_t p = 0; p < kNumInputPins; p++) {
const bool state = digitalRead(context, n, kInputPins[p].pin);
PinRT& s = gPinRT[p];
if (!s.initialised) {
s.initialised = true;
s.prev_state = state;
s.last_edge_frame = frame;
s.last_accepted_frame = frame;
if (state)
states |= (1 << p);
else
states &= ~(1 << p);
statesChanged = true;
continue;
}
if (state == s.prev_state)
continue;
const uint64_t dt = frame - s.last_edge_frame;
s.prev_state = state;
s.last_edge_frame = frame;
s.edge_index++;
if (state)
states |= (1 << p);
else
states &= ~(1 << p);
statesChanged = true;
// Forward the Pi's line to the EEG amp first. The mirror is the one
// thing in this loop with a hard latency budget, and it must not
// sit behind the logging bookkeeping -- nor behind the refractory
// guard below, which may decide not to emit a row but must never
// decide not to pass a trigger on.
if (p == kTriggerInIndex && sOutputsArmed) {
sFwdLevel = state;
digitalWrite(context, n, kFwdOutPin, sFwdLevel);
pushTriggerRT(TRIG_FORWARD, frame, sFwdSeq++, state, 0, 0);
}
// Quiet before: `dt` is the gap since the previous edge of either
// polarity on this pin, suppressed or not.
const bool accepted = dt >= gRefractoryFrames[p];
if (accepted) {
s.last_accepted_frame = frame;
s.burst_count = 0;
} else {
s.burst_count++;
}
// Bounded chatter guard: keep counting but stop emitting rows until
// the pin has been quiet for a full refractory period (which is
// what makes the next edge `accepted` and resets burst_count).
// Rising and falling edges both count, so a PWM-lit flash is ~2
// edges per PWM period.
if (!accepted && s.burst_count > kMaxEdgesPerBurst) {
s.suppressed_count++;
continue;
}
EdgeEvent e;
e.frame = frame;
e.edge_index = s.edge_index;
e.dt_frames_prev = dt;
e.suppressed_count = s.suppressed_count;
e.pin_index = static_cast<uint32_t>(p);
e.state = state ? 1 : 0;
e.accepted = accepted ? 1 : 0;
if (gEdgeQueue.push(e)) {
s.suppressed_count = 0;
} else {
gDroppedEvents.fetch_add(1, std::memory_order_relaxed);
}
}
// --- jittered timer trigger.
// `>=` rather than `==` so that a dropped block cannot swallow a pulse;
// the resulting lateness is logged rather than hidden.
if (sOutputsArmed) {
if (!sTimerLevel && frame >= sNextOnFrame) {
const int64_t late = static_cast<int64_t>(frame - sNextOnFrame);
sTimerLevel = true;
digitalWrite(context, n, kTimerOutPin, true);
sOffFrame = frame + gTimerPulseFrames;
const int64_t jit = nextJitterFrames();
sNextOnFrame += gTimerPeriodFrames + jit;
// After a gap long enough to swallow a whole interval, catch up
// to the grid rather than firing a burst to make up the count.
if (sNextOnFrame <= frame)
sNextOnFrame = frame + gTimerPeriodFrames;
pushTriggerRT(TRIG_TIMER, frame, sTimerSeq++, true, jit, late);
} else if (sTimerLevel && frame >= sOffFrame) {
sTimerLevel = false;
digitalWrite(context, n, kTimerOutPin, false);
pushTriggerRT(TRIG_TIMER, frame, sTimerSeq++, false, 0, 0);
}
}
for (unsigned int j = 0; j < context->audioOutChannels; j++) {
audioWrite(context, n, j, 0.0f);
}
}
if (statesChanged) {
gPinStatesAtomic.store(states, std::memory_order_relaxed);
}
gElapsedFramesAtomic.store(blockEnd, std::memory_order_relaxed);
static unsigned int renderCount = 0;
renderCount++;
// Clock-sync pairs every ~200 ms.
if (renderCount % gSyncPeriodBlocks == 0) {
gSyncRequestFrame.store(blockEnd, std::memory_order_relaxed);
Bela_scheduleAuxiliaryTask(gClockSyncTask);
}
}
// ---------------------------------------------------------------------------
// LSL thread
//
// A plain std::thread, not a Bela AuxiliaryTask, for two reasons: it can block
// on the network without stalling anything real-time, and file/socket work here
// causes no Xenomai mode switches.
//
// The pull is *blocking* with a short timeout. pull_sample() then returns
// within microseconds of liblsl handing over the sample, so local_clock() on
// the very next line is a genuine arrival timestamp. Polling non-blocking at
// some interval instead would add that whole interval as quantisation error to
// the number we are trying to measure.
// ---------------------------------------------------------------------------
#if ENABLE_LSL
static std::atomic<bool> gStreamConnected{false};
static std::string gConnectedSourceId;
static std::mutex gConnectedMutex;
static void writeStreamXml(const std::string& xml) {
std::ofstream f(sessionPath("_stream.xml"),
std::ios::out | std::ios::trunc);
if (f.is_open()) {
f << xml;
f.close();
}
}
static void lslThreadFunc() {
lsl::continuous_resolver* resolver = nullptr;
lsl::stream_inlet* inlet = nullptr;
std::string srcId, streamName;
uint32_t nChan = 0;
uint64_t rxIndex = 0;
std::vector<double> buf;
double nextCorrTime = 0.0;
bool firstCorrection = true;
while (gRunning.load(std::memory_order_relaxed)) {
// ---- discovery
// -------------------------------------------------------
if (!inlet) {
if (!resolver)
resolver = new lsl::continuous_resolver(5.0);
std::vector<lsl::stream_info> results = resolver->results();
for (size_t i = 0; i < results.size() && !inlet; i++) {
lsl::stream_info& info = results[i];
if (info.name().empty())
continue;
if (info.name().find(kStreamPrefixFilter) == std::string::npos)
continue;
// Any numeric format is fine: pull_sample(vector<double>) makes
// liblsl convert, so the Bela side is immune to Flutter-side
// channel changes.
if (info.channel_format() == lsl::cf_string ||
info.channel_format() == lsl::cf_undefined) {
pushStatus(ST_INFO, 0, nowClock(), info.channel_format(),
"skipped non-numeric stream");
continue;
}
try {
inlet = new lsl::stream_inlet(
info, kInletBufLenSec, kInletChunkLen, kInletRecover);
// Raw sender timestamps. The correction is logged
// separately as a time series so it can be fitted offline;
// never let liblsl silently smooth or shift the thing being
// measured.
inlet->set_postprocessing(lsl::post_none);