-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmodem.py
More file actions
165 lines (143 loc) · 5.74 KB
/
Copy pathmodem.py
File metadata and controls
165 lines (143 loc) · 5.74 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
"""Modem core.
Owns the DSP, framing and link layers and exposes a block-oriented audio
API that works the same for a real sounddevice stream and for the
software-loopback selftest:
out = modem.pull_tx(n) # (n, 2) float32 to the DAC
modem.push_rx(block) # (n, 2) float32 from the ADC
Lane mapping:
bond on : lane 0 = Left, lane 1 = Right (independent frames, striped)
bond off : lane 0 only; TX is duplicated on L and R, RX listens on L.
"""
from __future__ import annotations
import threading
import time
from collections import deque
import numpy as np
from .config import ModemConfig, lane_count, FT_DATA
from .dsp import OFDM, Demod
from .framing import FrameBuilder, Header
from .link import LinkEngine
class _Lane:
def __init__(self, idx: int, ofdm: OFDM, header_cb, frame_cb):
self.idx = idx
self.demod = Demod(ofdm, header_cb, frame_cb, name="LR"[idx])
self.fifo: deque[np.ndarray] = deque()
self.fifo_n = 0
self.pending_hdr: Header | None = None
self.lock = threading.Lock()
def push_tx(self, samples: np.ndarray):
with self.lock:
self.fifo.append(samples)
self.fifo_n += len(samples)
def pop_tx(self, n: int) -> np.ndarray:
out = np.zeros(n, dtype=np.float32)
with self.lock:
i = 0
while i < n and self.fifo:
head = self.fifo[0]
take = min(n - i, len(head))
out[i:i + take] = head[:take]
if take == len(head):
self.fifo.popleft()
else:
self.fifo[0] = head[take:]
self.fifo_n -= take
i += take
return out
class Modem:
def __init__(self, cfg: ModemConfig, log=lambda *_: None,
clock=time.monotonic):
self.cfg = cfg
self.log = log
self.clock = clock
self.ofdm = OFDM(cfg)
self.builder = FrameBuilder(cfg, self.ofdm)
self.link = LinkEngine(
cfg, self.builder.max_payload, self._to_host, self._local_snr,
frame_dur_fn=lambda mode: self.builder.frame_duration(cfg.max_data_syms),
log=log, clock=clock)
self.lanes = [
_Lane(i, self.ofdm,
header_cb=self._mk_header_cb(i),
frame_cb=self._mk_frame_cb(i))
for i in range(lane_count(cfg))
]
self.host_send = None # set by host interface: fn(pkt bytes)
self._low_water = 2 * cfg.blocksize
self._t0 = time.monotonic()
self.audio_level_in = 0.0
# -------------------------------------------------------- host plumbing
def attach_host(self, send_fn):
self.host_send = send_fn
def host_packet_in(self, pkt: bytes):
self.link.host_packet_in(pkt)
def _to_host(self, pkt: bytes):
if self.host_send:
self.host_send(pkt)
def _local_snr(self) -> float:
vals = [ln.demod.metrics["snr"] for ln in self.lanes]
return float(np.mean(vals)) if vals else 0.0
# ----------------------------------------------------------- RX wiring
def _mk_header_cb(self, idx: int):
def cb(hard_bits):
h = self.builder.parse_header_bits(hard_bits)
if h is None or h.nsyms > self.cfg.max_data_syms:
return None
self.lanes[idx].pending_hdr = h
return (h.nsyms, h.mode, h.mpx)
return cb
def _mk_frame_cb(self, idx: int):
def cb(data_bits, metrics):
h = self.lanes[idx].pending_hdr
self.lanes[idx].pending_hdr = None
if h is None:
return
now = self.clock()
if h.nsyms == 0:
self.link.on_frame(h, b"", idx, now)
elif data_bits is not None:
payload = self.builder.decode_payload(data_bits, h)
self.link.on_frame(h, payload, idx, now)
return cb
# ----------------------------------------------------------- audio API
def pull_tx(self, n: int) -> np.ndarray:
now = self.clock()
for ln in self.lanes:
while ln.fifo_n < self._low_water:
spec = self.link.next_frame(ln.idx, now)
if spec is None:
break
ftype, seq, ack, snr, payload, mode, mpx, use_fec = spec
samples, _ = self.builder.build(ftype, seq, ack, snr,
payload, mode, mpx, use_fec)
ln.push_tx(samples)
out = np.zeros((n, 2), dtype=np.float32)
out[:, 0] = self.lanes[0].pop_tx(n)
if len(self.lanes) > 1:
out[:, 1] = self.lanes[1].pop_tx(n)
else:
out[:, 1] = out[:, 0] # duplicate mono lane
return out
def push_rx(self, block: np.ndarray):
if block.ndim == 1:
block = np.stack([block, block], axis=1)
self.audio_level_in = float(np.sqrt(np.mean(block[:, 0] ** 2)) + 1e-12)
self.lanes[0].demod.feed(block[:, 0])
if len(self.lanes) > 1:
self.lanes[1].demod.feed(block[:, 1])
# -------------------------------------------------------------- status
def metrics(self) -> dict:
m = self.link.stats()
m["lanes"] = []
for ln in self.lanes:
d = dict(ln.demod.metrics)
d["sync"] = ln.demod.state != 0 or (self.clock() - m["t_rx"] < 1.5)
m["lanes"].append(d)
m["level_in_db"] = 20 * np.log10(self.audio_level_in + 1e-9)
m["snr"] = self._local_snr()
return m
def const_points(self, lane=0, n=400):
pts = self.lanes[lane].demod.const_points
out = pts[-n:]
del pts[:max(0, len(pts) - n)]
return out