Add a way of sending groups of packets as one

This commit is contained in:
xssfox 2023-12-27 21:29:22 +11:00
parent deefd79402
commit a596cabd2d
4 changed files with 154 additions and 80 deletions

View file

@ -87,7 +87,7 @@ if __name__ == '__main__':
def tx(data):
try:
logging.debug(f"Sending {str(data)}")
output_device.write(modem_tx.write(data))
output_device.write(Packet(data))
except:
logging.critical(
traceback.format_exc()
@ -151,6 +151,7 @@ if __name__ == '__main__':
input_device = audio.InputDevice(modem_rx.write, modem_rx.sample_rate, name_or_id=input_device_name_or_id)
output_device = audio.OutputDevice(
modem_rx.sample_rate,
modem = modem_tx,
name_or_id=output_device_name_or_id,
ptt_release=ptt_release,
ptt_trigger=ptt_trigger,

View file

@ -4,11 +4,12 @@ from tabulate import tabulate
import logging
import audioop as pyaudioop
import time
from threading import Lock
from threading import Lock, Thread
from typing import Callable
#from pydub import pyaudioop
import pydub
import math
from .modem import FreeDVTX, Packet
p = pyaudio.PyAudio()
@ -161,8 +162,10 @@ class OutputDevice():
buffer = bytearray()
output_buffer_lock = Lock()
send_queue_lock = Lock()
inhibit = False
output_buffer_thread = None
@property
def queue_ms(self):
@ -170,6 +173,7 @@ class OutputDevice():
def __init__(self,
sample_rate: int,
modem: FreeDVTX,
name_or_id:int|str|None=None,
ptt_trigger:Callable[[],None]=None,
ptt_release:Callable[[],None]=None,
@ -182,6 +186,8 @@ class OutputDevice():
self.ptt_on_delay_ms = ptt_on_delay_ms
self.ptt_off_delay_ms = ptt_off_delay_ms
self.db = db
self.send_queue = []
self.modem = modem
if name_or_id != None:
try:
@ -219,13 +225,46 @@ class OutputDevice():
self.ptt_release = ptt_release
self.ptt = False
def write_raw(self,data:bytes):
if self.device.sample_rate != self.sample_rate:
(data, self.rate_state) = pyaudioop.ratecv(
data,
pyaudio.get_sample_size(FORMAT),
1,
self.sample_rate,
self.device.sample_rate,
self.rate_state,
)
if self.db:
data = pyaudioop.mul(data, 2, 10**(self.db/20.0))
def write(self, data: bytes):
if self.device.output_channels == 2:
data = pyaudioop.tostereo(
data,
pyaudio.get_sample_size(FORMAT),
1,
1
)
with self.output_buffer_lock:
self.buffer += data
def write(self, data: Packet):
with self.send_queue_lock:
self.send_queue.append(data)
def audio_buffer(self):
logging.debug("Populating audio buffer")
# ptt delay
pre_silence = pydub.AudioSegment.silent(duration=self.ptt_on_delay_ms, frame_rate=self.device.sample_rate)
pre_silence = pre_silence.set_channels(self.device.output_channels)
write_buffer = pre_silence.raw_data
with self.send_queue_lock:
send_queue = self.send_queue
self.send_queue = []
data = self.modem.write(send_queue)
if self.device.sample_rate != self.sample_rate:
(data, self.rate_state) = pyaudioop.ratecv(
data,
@ -265,6 +304,7 @@ class OutputDevice():
ptt = False
# if we aren't transmitting and we have inhibited tx then skip
if self.inhibit == True and self.ptt == False:
return (bytes(output), pyaudio.paContinue)
@ -274,6 +314,10 @@ class OutputDevice():
output[:chunk_size] = self.buffer[:chunk_size]
if self.buffer:
ptt = True
elif self.send_queue and (not self.output_buffer_thread or not self.output_buffer_thread.is_alive()):
# if we have no output buffer and queued messages we should start a thread to generate an output buffer
self.output_buffer_thread = Thread(target=self.audio_buffer)
self.output_buffer_thread.start()
del self.buffer[:chunk_size]
if self.ptt != ptt:

View file

@ -23,7 +23,11 @@ class FreeDVFrame:
snr: float
modem: Modems
@dataclass
class Packet():
data: bytes
header: int|bytes = b"\xff"
mode: str = None
class Modem():
def __init__(self, modem: Modems, callback: Callable[[FreeDVFrame],None]|None=None):
@ -113,7 +117,7 @@ class Modem():
data_in = ffi.from_buffer(f"unsigned char[{self.bytes_per_frame - 2}]", data)
return lib.freedv_gen_crc16(data_in, self.bytes_per_frame - 2).to_bytes(2, byteorder="big")
def modulate(self, data: bytes, header_byte=b'\xff') -> bytes:
def modulate(self, queue: list[Packet]) -> bytes:
"""
Modulates bytes into audio samples (also bytes)
"""
@ -131,46 +135,62 @@ class Modem():
"""
# Convert to byte array as it will be easier to slice
data = bytearray(data)
chunks = []
pop_packet_length = self.bytes_per_frame - 2 - 3 # first iteration we use 3 bytes for the header
while data:
chunks.append(data[:pop_packet_length])
del data[:pop_packet_length]
pop_packet_length = self.bytes_per_frame - 2 - 1 # next iterations only use 1 byte for sequence
frames = []
# first frame includes header
pop_packet_length = self.bytes_per_frame - 2 - 3 # first iteration we use 3 bytes for the header
frame=bytearray(self.bytes_per_frame)
used_bytes = 0
while queue:
packet = queue.pop(0)
data = bytearray(packet.data)
chunks = []
header_byte = packet.header
while data:
chunks.append(data[:pop_packet_length])
del data[:pop_packet_length]
pop_packet_length = self.bytes_per_frame - 2 - 1 # next iterations only use 1 byte for sequence
# header
frame[0:3] = header_byte + sum([len(x) for x in chunks]).to_bytes(2)
# data
frame[3:3+len(chunks[0])] = chunks[0]
# crc
frame[-2:] = self.crc(bytes(frame)[:-2])
frames.append(frame)
for seq, next_chunk in enumerate(chunks[1:]):
frame=bytearray(self.bytes_per_frame)
# header
frame[0] = seq
frame[1:1+len(next_chunk)] = next_chunk
header = header_byte + sum([len(x) for x in chunks]).to_bytes(2)
frame[0+used_bytes:3+used_bytes] = header
# crc
frame[-2:] = self.crc(bytes(frame)[:-2])
chunk = chunks[0]
frames.append(frame)
# data
frame[3+used_bytes:3+len(chunk)+used_bytes] = chunk
used_bytes = used_bytes + len(chunk) + len(header)
for header, chunk in enumerate(chunks[1:]):
used_bytes = 0
frames.append(frame)
frame=bytearray(self.bytes_per_frame)
# header
header = header.to_bytes(1)
frame[0:len(header)] = header
frame[1:1+len(chunk)] = chunk
used_bytes = used_bytes + len(chunk) + len(header)
# can we fit a little more data in?
# 2 for crc
if used_bytes + 2 <= self.bytes_per_frame - 3: # we need three bytes to start the next payload
pop_packet_length = self.bytes_per_frame - used_bytes - 2 - 3
else:
if queue:
frames.append(frame)
frame=bytearray(self.bytes_per_frame)
used_bytes = 0
pop_packet_length = self.bytes_per_frame - 2 - 3
frames.append(frame)
output = bytes()
for frame in frames:
# calculate CRCs
frame[-2:] = self.crc(bytes(frame)[:-2])
#logging.debug(f"modulating {str(bytes(frame))}")
from_modem = ffi.new(f"short mod_out[{lib.freedv_get_n_tx_modem_samples(self.modem)}]")
@ -187,18 +207,14 @@ class Modem():
#postamble
samples=lib.freedv_rawdatapostambletx(self.modem, from_modem)
output += ffi.buffer(from_modem)[:(samples*ffi.sizeof("short"))]
# add an extra bit of silence to clear out buffers
output += bytes(lib.freedv_get_n_nom_modem_samples(self.modem)*ffi.sizeof("short")*2)
return output
@dataclass
class Packet():
data: bytes
header: int
mode: str
class FreeDVRX():
def __init__(self, callback: Callable[[bytes],None], progress: Callable[[int,int],None], inhibit: Callable[[bool],None]):
@ -234,45 +250,55 @@ class FreeDVRX():
def rx(self, data_frame: FreeDVFrame):
logging.debug(f"Received data. snr:{data_frame.snr}")
data = bytearray(data_frame.data)
header = data.pop(0)
if header > 200: # start of packet
self.remaining_bytes = int.from_bytes(data[0:2])
self.total_bytes = self.remaining_bytes
del data[0:2]
logging.debug(f"Found packet start - Expecting {self.remaining_bytes} bytes")
self.next_seq_number = 0
self.partial_data=b''
self.header = header
elif self.next_seq_number != None: # should be a seq number
if self.next_seq_number != header:
logging.debug(f"Missing data - header seq expected {self.next_seq_number}, got {header}")
while data:
#logging.debug(f"Received data: {str(data)}")
header = data.pop(0)
if header > 200: # start of packet
self.remaining_bytes = int.from_bytes(data[0:2])
self.total_bytes = self.remaining_bytes
del data[0:2]
logging.debug(f"Found packet start - Expecting {self.remaining_bytes} bytes")
self.next_seq_number = 0
self.partial_data=b''
self.header = header
elif self.next_seq_number != None: # should be a seq number
if self.next_seq_number != header:
logging.debug(f"Missing data - header seq expected {self.next_seq_number}, got {header}")
logging.debug(f"Full data frame: {str(data_frame.data)}")
self.next_seq_number = None
self.remaining_bytes = None
return
else:
logging.debug(f"Received frame {header}")
self.next_seq_number += 1
else:
if header != 0:
logging.debug(f"Not expecting data - got {header}")
logging.debug(f"Full data frame: {str(data_frame.data)}")
self.next_seq_number = None
self.remaining_bytes = None
return
self.partial_data += data[:self.remaining_bytes]
received_bytes = len(data[:self.remaining_bytes])
self.remaining_bytes -= received_bytes
logging.debug(f"Seq: {header} Remaining data: {self.remaining_bytes}")
self.progress(self.total_bytes, self.remaining_bytes, data_frame.modem)
if self.remaining_bytes == 0:
self.next_seq_number = None
del data[:received_bytes]
self.remaining_bytes = None
self.callback(Packet(header=self.header, data=self.partial_data, mode=data_frame.modem))
else:
logging.debug(f"Received frame {header}")
self.next_seq_number += 1
else:
logging.debug(f"Not expecting data - got {header}")
self.next_seq_number = None
self.remaining_bytes = None
return
self.partial_data += data[:self.remaining_bytes]
self.remaining_bytes -= len(data[:self.remaining_bytes])
return
logging.debug(f"Seq: {header} Remaining data: {self.remaining_bytes}")
self.progress(self.total_bytes, self.remaining_bytes, data_frame.modem)
if self.remaining_bytes == 0:
self.next_seq_number = None
self.remaining_bytes = None
self.callback(Packet(header=self.header, data=self.partial_data, mode=data_frame.modem))
class FreeDVTX():
def __init__(self, modem: str = Modems.DATAC1.name):
self.modem = Modem(modem={x.name:x for x in Modems}[modem])
def set_mode(self, modem: str):
self.modem = Modem(modem={x.name:x for x in Modems}[modem])
def write(self, data: bytes, header_byte=b'\xff'):
return self.modem.modulate(data, header_byte)
def write(self, data: list[Packet]):
return self.modem.modulate(data)

View file

@ -18,7 +18,7 @@ import readline
import code
import rlcompleter
import pydub.generators
from .modem import Modems, FreeDVRX, FreeDVTX
from .modem import Modems, FreeDVRX, FreeDVTX, Packet
import traceback
from pathlib import Path
import argparse
@ -84,7 +84,8 @@ class FreeDVShellCommands():
bit_depth=16,
).to_audio_segment(2000, volume=-6)
sin_wave.set_channels(1)
self.output_device.write(sin_wave.raw_data)
self.output_device.write_raw(sin_wave.raw_data)
def help_mode(self):
return f"Change TX Mode: mode [{', '.join([x.name for x in Modems])}]"
@ -110,6 +111,8 @@ class FreeDVShellCommands():
def do_clear(self, arg):
"Clears TX queues"
self.output_device.clear()
with self.output_device.send_queue_lock:
self.output_device.send_queue = []
return "TX buffer cleared"
def do_list_audio_devices(self, arg):
@ -124,7 +127,7 @@ class FreeDVShellCommands():
def do_send_string(self, arg):
"Sends string over the modem"
self.output_device.write(self.modem_tx.write(arg.encode()))
self.output_device.write(Packet(arg.encode()))
return "Queued for sending"
def do_volume(self,arg):
@ -149,8 +152,7 @@ class FreeDVShellCommands():
return "Set callsign with the callsign command first\n"
data = self.options.callsign.encode() + b"\xff" + arg.encode()
self.output_device.write(self.modem_tx.write(data, header_byte=b"\xfe"))
self.output_device.write(Packet(data, header=b"\xfe"))
def do_exit(self, arg):
"Exits FreeDVTNC2"
@ -229,7 +231,7 @@ class FreeDVShell():
def accept(buff):
try:
new_text = self.log.text + f"\n> {input_field.text}\n"
new_text = self.log.text + f"> {input_field.text}\n"
self.log.buffer.document = Document(
text=new_text, cursor_position=(len(new_text))
)
@ -295,7 +297,8 @@ class FreeDVShell():
(f"class:status.{ 'red' if self.output_device.ptt else 'green' }", f"{ ' on' if self.output_device.ptt else 'off' }"),
("class:status", f" | "),
("class:status", f"TX Queue: { (self.output_device.queue_ms / 1000) :5.1f}s | "),
("class:status", f"Audio Queue: { (self.output_device.queue_ms / 1000) :5.1f}s | "),
("class:status", f"TX Queue: { len(self.output_device.send_queue) :3.0f} | "),
("class:status", f"Channel: "),
(f"class:status.{'red' if self.output_device.inhibit else 'green'}", f"{'busy' if self.output_device.inhibit else 'clear'}"),
("class:status", " |\n"),