254 lines
11 KiB
Python
254 lines
11 KiB
Python
"""Bounded host traffic recording and synthetic DSView v3 logic captures.
|
||
|
||
No hardware timing is inferred: timestamps are host observations and overlapping
|
||
packets on one wire are serialized. Original events remain in the archive.
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import os
|
||
from pathlib import Path
|
||
import tempfile
|
||
from time import monotonic_ns
|
||
import zipfile
|
||
|
||
|
||
def uart_event(source: str, direction: str, data: bytes, rate: int) -> dict:
|
||
return dict(source=source, kind="uart", direction=direction, data=bytes(data),
|
||
rate=rate, time_ns=monotonic_ns())
|
||
|
||
|
||
def can_event(source: str, direction: str, identifier: int, data: bytes,
|
||
rate: int, extended: bool = True, remote: bool = False) -> dict:
|
||
return dict(source=source, kind="can", direction=direction, data=bytes(data),
|
||
rate=rate, time_ns=monotonic_ns(), identifier=identifier,
|
||
extended=extended, remote=remote)
|
||
|
||
|
||
def validate_event(event: dict) -> None:
|
||
if event['kind'] not in ('uart', 'can') or event['direction'] not in ('RX', 'TX'):
|
||
raise ValueError('Unknown traffic type/direction')
|
||
if not isinstance(event['rate'], int) or not 1 <= event['rate'] <= 10_000_000:
|
||
raise ValueError('Invalid line rate')
|
||
if not isinstance(event['time_ns'], int) or event['time_ns'] < 0:
|
||
raise ValueError('Invalid event timestamp')
|
||
if not event['source'] or any(c in event['source'] for c in '\r\n=[]'):
|
||
raise ValueError('Invalid source name')
|
||
if event['kind'] == 'can':
|
||
if len(event['data']) > 8 or not 0 <= event['identifier'] <= (0x1fffffff if event['extended'] else 0x7ff):
|
||
raise ValueError('Only classic CAN IDs and DLC 0..8 are supported')
|
||
|
||
|
||
class TrafficRecording:
|
||
MAX_EVENTS = 100_000
|
||
MAX_BYTES = 8 * 1024 * 1024
|
||
|
||
def __init__(self):
|
||
self.events: list[dict] = []
|
||
self.byte_count = 0
|
||
|
||
def append(self, event: dict) -> None:
|
||
validate_event(event)
|
||
data = bytes(event['data'])
|
||
if event['kind'] == 'uart' and not data:
|
||
return
|
||
if len(self.events) >= self.MAX_EVENTS or self.byte_count + len(data) > self.MAX_BYTES:
|
||
raise ValueError('Достигнут лимит записи: 100 000 событий / 8 МиБ данных')
|
||
self.events.append(dict(event, data=data))
|
||
self.byte_count += len(data)
|
||
|
||
|
||
def bits(value: int, count: int):
|
||
return [(value >> i) & 1 for i in range(count-1, -1, -1)]
|
||
|
||
|
||
def uart_bits(data: bytes):
|
||
for byte in data:
|
||
yield 0
|
||
for i in range(8):
|
||
yield (byte >> i) & 1
|
||
yield 1
|
||
|
||
|
||
def can_bits(event: dict):
|
||
"""CAN 2.0 data/remote frame, CRC15, stuffing, synthetic successful ACK."""
|
||
ident, remote, payload = event['identifier'], event['remote'], event['data']
|
||
if event['extended']:
|
||
raw = [0] + bits(ident >> 18, 11) + [1, 1] + bits(ident & 0x3ffff, 18) + [int(remote), 0, 0]
|
||
else:
|
||
raw = [0] + bits(ident, 11) + [int(remote), 0, 0]
|
||
raw += bits(len(payload), 4)
|
||
if not remote:
|
||
for byte in payload:
|
||
raw += bits(byte, 8)
|
||
crc = 0
|
||
for bit in raw:
|
||
feedback = ((crc >> 14) & 1) ^ bit
|
||
crc = (crc << 1) & 0x7fff
|
||
if feedback:
|
||
crc ^= 0x4599
|
||
raw += bits(crc, 15)
|
||
last, run = None, 0
|
||
for bit in raw:
|
||
yield bit
|
||
run = run + 1 if bit == last else 1
|
||
last = bit
|
||
if run == 5:
|
||
last = 1 - bit
|
||
yield last
|
||
run = 1
|
||
# CRC delimiter, reconstructed ACK, ACK delimiter, EOF, intermission.
|
||
yield from [1, 0, 1] + [1] * 10
|
||
|
||
|
||
def channel_name(event):
|
||
# A reconnect at another line rate must remain independently decodable.
|
||
return event['source'] + ('_' + event['direction'] if event['kind'] == 'uart' else '_CAN') + '_' + str(event['rate'])
|
||
|
||
|
||
class _PackedWriter:
|
||
"""LSB-first samples streamed in bounded chunks, including long idle runs."""
|
||
def __init__(self, stream, cancelled):
|
||
self.stream, self.cancelled = stream, cancelled
|
||
self.partial = self.used = 0
|
||
self.buffer = bytearray()
|
||
|
||
def _flush(self):
|
||
if self.cancelled():
|
||
raise InterruptedError('Экспорт отменён')
|
||
self.stream.write(self.buffer)
|
||
self.buffer.clear()
|
||
|
||
def run(self, level: int, count: int):
|
||
if count < 0:
|
||
raise ValueError('Non-monotonic waveform')
|
||
if self.used:
|
||
take = min(count, 8 - self.used)
|
||
if level:
|
||
self.partial |= ((1 << take)-1) << self.used
|
||
self.used += take
|
||
count -= take
|
||
if self.used == 8:
|
||
self.buffer.append(self.partial)
|
||
self.partial = self.used = 0
|
||
full, tail = divmod(count, 8)
|
||
while full:
|
||
take = min(full, 65536)
|
||
self.buffer.extend(bytes([255 if level else 0]) * take)
|
||
full -= take
|
||
if len(self.buffer) >= 65536:
|
||
self._flush()
|
||
if tail:
|
||
self.partial = (1 << tail)-1 if level else 0
|
||
self.used = tail
|
||
|
||
def finish(self):
|
||
if self.used:
|
||
raise ValueError('Unaligned sample count')
|
||
self._flush()
|
||
|
||
|
||
def export_dsl(path, events, samplerate=20_000_000, *, overwrite=False,
|
||
cancelled=lambda: False, progress=lambda value: None):
|
||
"""Atomic export. Existing destination survives failed/cancelled exports."""
|
||
path = Path(path)
|
||
if not events:
|
||
raise ValueError('Нет записанного трафика')
|
||
if path.exists() and not overwrite:
|
||
raise FileExistsError(path)
|
||
if not isinstance(samplerate, int) or not 1 <= samplerate <= 100_000_000:
|
||
raise ValueError('Частота дискретизации должна быть от 1 до 100 МГц')
|
||
groups = {}
|
||
origin = min(event['time_ns'] for event in events)
|
||
shifted = 0
|
||
for event in events:
|
||
validate_event(event)
|
||
if samplerate < event['rate'] * 8:
|
||
raise ValueError('Частота дискретизации должна быть не меньше 8× скорости линии')
|
||
groups.setdefault(channel_name(event), []).append(event)
|
||
if len(groups) > 32:
|
||
raise ValueError('DSView: максимум 32 цифровых канала')
|
||
schedule, maximum = {}, 0
|
||
for name, channel_events in groups.items():
|
||
cursor = 16
|
||
scheduled = []
|
||
for event in sorted(channel_events, key=lambda item: item['time_ns']):
|
||
if cancelled():
|
||
raise InterruptedError('Экспорт отменён')
|
||
data_bits = list(can_bits(event)) if event['kind'] == 'can' else None
|
||
bit_count = len(data_bits) if data_bits is not None else len(event['data']) * 10
|
||
wanted = 16 + (event['time_ns'] - origin) * samplerate // 1_000_000_000
|
||
start = max(cursor, wanted)
|
||
shifted += int(start > wanted)
|
||
count = (bit_count * samplerate + event['rate'] - 1) // event['rate']
|
||
scheduled.append((start, event, data_bits, bit_count))
|
||
cursor = start + count
|
||
maximum = max(maximum, cursor + 16)
|
||
schedule[name] = scheduled
|
||
padded = (maximum + 63) // 64 * 64
|
||
if padded > 2_000_000_000:
|
||
raise ValueError('Запись превышает 2 млрд отсчётов: сократите длительность или частоту')
|
||
names = list(groups)
|
||
header = ('[version]\nversion = 3\n[header]\ndriver = virtual-session\n'
|
||
'device mode = 0\ncapturefile = data\n'
|
||
f'total samples = {padded}\ntotal probes = {len(names)}\ntotal blocks = 1\n'
|
||
f'samplerate = {samplerate} Hz\ntrigger time = 0\ntrigger pos = 0\n'
|
||
+ ''.join(f'probe{i} = {name}\n' for i, name in enumerate(names)))
|
||
session = {'Device': 'virtual-session', 'DeviceMode': 0, 'Version': 3,
|
||
'Title': 'DSView v1.3.2', 'Sample count': str(padded),
|
||
'Sample rate': str(samplerate), 'Max Height': '1X', 'decoder': [],
|
||
'channel': [{'colour': 'default', 'enabled': True, 'index': i,
|
||
'name': name, 'strigger': 0, 'type': 10000, 'view_index': i}
|
||
for i, name in enumerate(names)]}
|
||
fd, temporary = tempfile.mkstemp(prefix='.setgui-dsl-', suffix='.tmp', dir=path.parent)
|
||
os.close(fd)
|
||
try:
|
||
with zipfile.ZipFile(temporary, 'w', zipfile.ZIP_DEFLATED) as archive:
|
||
archive.writestr('header', header)
|
||
archive.writestr('session', json.dumps(session, ensure_ascii=False))
|
||
archive.writestr('decoders', '[]')
|
||
metadata = dict(schema='setgui.protocol-capture.v1', synthetic=True,
|
||
timing='host observation; overlapping events serialized per wire',
|
||
uart='8N1, non-inverted', can_ack='synthetic dominant ACK',
|
||
shifted_events=shifted, samplerate=samplerate,
|
||
origin_monotonic_ns=origin,
|
||
events=[dict(event, data=event['data'].hex()) for event in events])
|
||
archive.writestr('protocol-log.json', json.dumps(metadata, ensure_ascii=False))
|
||
done, last_progress = 0, -1
|
||
for index, name in enumerate(names):
|
||
with archive.open(f'L-{index}/0', 'w', force_zip64=True) as output:
|
||
writer = _PackedWriter(output, cancelled)
|
||
cursor = 0
|
||
for start, event, data_bits, bit_count in schedule[name]:
|
||
writer.run(1, start - cursor)
|
||
bitstream = data_bits if data_bits is not None else uart_bits(event['data'])
|
||
previous = 0
|
||
for i, level in enumerate(bitstream, 1):
|
||
boundary = (i * samplerate + event['rate'] - 1) // event['rate']
|
||
writer.run(level, boundary - previous)
|
||
previous = boundary
|
||
cursor = start + previous
|
||
done += 1
|
||
value = done * 95 // len(events)
|
||
if value != last_progress:
|
||
progress(value)
|
||
last_progress = value
|
||
writer.run(1, padded - cursor)
|
||
writer.finish()
|
||
if cancelled():
|
||
raise InterruptedError('Экспорт отменён')
|
||
with zipfile.ZipFile(temporary) as archive:
|
||
if archive.testzip() is not None:
|
||
raise ValueError('Ошибка проверки архива DSL')
|
||
if cancelled():
|
||
raise InterruptedError('Экспорт отменён')
|
||
if overwrite:
|
||
os.replace(temporary, path)
|
||
else:
|
||
# Atomic no-clobber publication on the same filesystem.
|
||
os.link(temporary, path)
|
||
progress(100)
|
||
finally:
|
||
Path(temporary).unlink(missing_ok=True)
|
||
return dict(samples=padded, samplerate=samplerate, channels=names, shifted_events=shifted)
|