201 lines
9.9 KiB
Python
201 lines
9.9 KiB
Python
# Клиент согласованных снимков GAS-регистратора и проверка карты сигналов. Скачивание
|
||
# выполняется в рабочем потоке транспорта; записи с истёкшим ожиданием нельзя автоматически
|
||
# повторять, поскольку устройство уже могло применить команду.
|
||
|
||
"""Experimental GAS recorder map and transport-independent snapshot client.
|
||
|
||
read(address, count) -> words; write(address, value) must await a confirmed
|
||
success or raise. Run download in the transport's worker, never a Qt UI slot.
|
||
No retries of writes: a timed-out snapshot may already exist on the device.
|
||
"""
|
||
from __future__ import annotations
|
||
import json
|
||
import re
|
||
import zlib
|
||
from pathlib import Path
|
||
|
||
MAGIC, VERSION = 0x474C, 1
|
||
DATA_OFFSET, HEADER_WORDS = 0x100, 5
|
||
|
||
|
||
def validate_map(data):
|
||
if data.get('format') != 'set-gas-logger' or data.get('version') != VERSION:
|
||
raise ValueError('Unsupported GAS logger map')
|
||
def integer(value, low, high):
|
||
if type(value) is not int or not low <= value <= high:
|
||
raise ValueError('Invalid map integer')
|
||
return value
|
||
base = integer(data['service']['base'], 0, 65535)
|
||
if data['service']['data_offset'] != DATA_OFFSET:
|
||
raise ValueError('Invalid snapshot data offset')
|
||
capacity = integer(data['capacity'], 1, 65535)
|
||
channels = data['channels']
|
||
if not 1 <= len(channels) <= 64:
|
||
raise ValueError('Invalid channel count')
|
||
end = base + DATA_OFFSET + capacity * (len(channels) + HEADER_WORDS)
|
||
if end > 65536:
|
||
raise ValueError('Logger window exceeds GAS')
|
||
seen = {key: set() for key in ('key', 'gas', 'modbus', 'can')}
|
||
for channel in channels:
|
||
if not re.fullmatch(r'[a-z][a-z0-9_]*', channel['key']):
|
||
raise ValueError('Invalid channel key')
|
||
if not isinstance(channel['name'], str) or not channel['name'].strip():
|
||
raise ValueError('Missing channel name')
|
||
if channel['type'] not in ('u16', 'i16'):
|
||
raise ValueError('Unsupported channel type')
|
||
if channel['can']['protocol'] != 'protocan_gas':
|
||
raise ValueError('CAN mapping must use ProtoCAN GAS')
|
||
if channel['modbus']['function'] != 3:
|
||
raise ValueError('Modbus mapping must use FC03')
|
||
for key in seen:
|
||
value = channel[key]
|
||
if key in ('modbus', 'can'):
|
||
value = value['address']
|
||
if key != 'key':
|
||
integer(value, 0, 65535)
|
||
if base <= value < end:
|
||
raise ValueError('Signal overlaps logger service')
|
||
if value in seen[key]:
|
||
raise ValueError('Duplicate ' + key)
|
||
seen[key].add(value)
|
||
source = channel['source']
|
||
if source['space'] not in ('pm35_modbus', 'pm67_cache'):
|
||
raise ValueError('Unknown PM35 source')
|
||
integer(source['address'], 0, 127 if source['space'] == 'pm35_modbus' else 1)
|
||
expected = ['magic', 'version', 'flags', 'channels', 'record_words', 'capacity',
|
||
'live_count', 'snapshot_count', 'generation_lo', 'generation_hi',
|
||
'schema_lo', 'schema_hi', 'sequence_lo', 'sequence_hi',
|
||
'missed_lo', 'missed_hi', 'source_errors_lo', 'source_errors_hi',
|
||
'command', 'release_generation_lo', 'release_generation_hi']
|
||
registers = data['service']['registers']
|
||
if [(r['offset'], r['key'], r['access']) for r in registers] != [
|
||
(i, key, 'r' if i < 18 else 'w') for i, key in enumerate(expected)]:
|
||
raise ValueError('Service map does not match protocol version 1')
|
||
return data
|
||
|
||
|
||
def load_map(path=None):
|
||
path = Path(path) if path else Path(__file__).with_name('gas_logger_maps') / 'pm35.json'
|
||
return validate_map(json.loads(path.read_text(encoding='utf-8')))
|
||
|
||
|
||
def schema_id(mapping):
|
||
validate_map(mapping)
|
||
raw = json.dumps(mapping, ensure_ascii=True, sort_keys=True, separators=(',', ':')).encode('ascii')
|
||
return zlib.crc32(raw) & 0xffffffff
|
||
|
||
|
||
def c_header(mapping):
|
||
"""Generate the checked-in C mirror; --check in CI catches drift."""
|
||
validate_map(mapping)
|
||
lines = ['/*\n'
|
||
' * Связь регистратора ПМ35 с картой доступных сигналов. Адреса сигналов согласуются с\n'
|
||
' * JSON-картой клиента; изменение только одной стороны приводит к неверной интерпретации снимка\n'
|
||
' * даже при успешном обмене.\n'
|
||
' */\n'
|
||
'\n'
|
||
'/* Generated from pm35.json. Do not edit. */',
|
||
'#ifndef GL_PM35_MAP_H', '#define GL_PM35_MAP_H', '#include "gas_logger.h"',
|
||
'static const gl_channel gl_pm35_channels[] = {']
|
||
for c in mapping['channels']:
|
||
lines.append(' {%dU, %dU, %dU}, /* %s */' % (
|
||
c['gas'], c['modbus']['address'], c['can']['address'], c['key']))
|
||
lines += ['};', 'static const uint16_t gl_pm35_source_space[] = { ' + ', '.join(
|
||
'0U' if c['source']['space'] == 'pm35_modbus' else '1U' for c in mapping['channels']) + ' };',
|
||
'static const uint16_t gl_pm35_source_address[] = { ' + ', '.join(
|
||
str(c['source']['address'])+'U' for c in mapping['channels'])+' };',
|
||
'static const gl_config gl_pm35_config = {',
|
||
' %dU, %dU, %dU, 0x%08XUL, gl_pm35_channels' % (
|
||
mapping['service']['base'], len(mapping['channels']), mapping['capacity'], schema_id(mapping)),
|
||
'};', '#endif', '']
|
||
return '\n'.join(lines)
|
||
|
||
|
||
def _u32(words, offset):
|
||
return words[offset] | words[offset + 1] << 16
|
||
|
||
|
||
class GasLoggerClient:
|
||
def __init__(self, read, write, mapping=None, *, block_words=120):
|
||
self.mapping = load_map() if mapping is None else validate_map(mapping)
|
||
if not 1 <= block_words <= 125:
|
||
raise ValueError('Block size must be 1..125 words')
|
||
self.read, self.write, self.block_words = read, write, block_words
|
||
self.base = self.mapping['service']['base']
|
||
|
||
def _read(self, address, count):
|
||
result = tuple(self.read(address, count))
|
||
if len(result) != count or any(type(w) is not int or not 0 <= w <= 65535 for w in result):
|
||
raise ValueError('Incomplete or invalid GAS response')
|
||
return result
|
||
|
||
def status(self):
|
||
# Classic CAN may return only four words. Generation and geometry are
|
||
# stable with one control client; the sequence/live counters are diagnostic.
|
||
words = tuple(w for offset in range(0, 18, self.block_words)
|
||
for w in self._read(self.base + offset, min(self.block_words, 18-offset)))
|
||
width = len(self.mapping['channels']) + HEADER_WORDS
|
||
if (words[:2] != (MAGIC, VERSION) or words[3] != len(self.mapping['channels'])
|
||
or words[4] != width or words[5] != self.mapping['capacity']
|
||
or _u32(words, 10) != schema_id(self.mapping)):
|
||
raise ValueError('Device logger/schema does not match JSON map')
|
||
if words[2] & ~3 or words[6] > words[5] or words[7] > words[5]:
|
||
raise ValueError('Invalid logger status')
|
||
return words
|
||
|
||
def start(self):
|
||
self.status() # never write to an unrecognized device
|
||
self.write(self.base+18, 1)
|
||
|
||
def release(self, generation):
|
||
status = self.status()
|
||
if not status[2] & 2 or _u32(status, 8) != generation:
|
||
raise ValueError('Snapshot generation changed')
|
||
self.write(self.base+19, generation & 0xffff)
|
||
self.write(self.base+20, generation >> 16)
|
||
self.write(self.base+18, 3)
|
||
|
||
def download(self, *, resume=False, cancelled=lambda: False, progress=lambda done,total: None):
|
||
"""Pin/copy/verify/release; cancellation/errors retain the pin for resume.
|
||
|
||
resume=True explicitly takes over the retained snapshot. A link failure
|
||
during release may have released it already; inspect status, don't replay.
|
||
"""
|
||
status = self.status()
|
||
if cancelled():
|
||
raise InterruptedError('Snapshot download cancelled')
|
||
if status[2] & 2:
|
||
if not resume:
|
||
raise RuntimeError('Snapshot already pinned; resume or release explicitly')
|
||
else:
|
||
if resume:
|
||
raise RuntimeError('No retained snapshot to resume')
|
||
self.write(self.base+18, 2)
|
||
status = self.status()
|
||
if not status[2] & 2 or not status[7]:
|
||
raise ValueError('Device did not pin a nonempty snapshot')
|
||
generation, count, width = _u32(status, 8), status[7], status[4]
|
||
words = []
|
||
total = count * width
|
||
while len(words) < total:
|
||
if cancelled():
|
||
raise InterruptedError('Snapshot download cancelled; snapshot retained')
|
||
words.extend(self._read(self.base+DATA_OFFSET+len(words), min(self.block_words,total-len(words))))
|
||
progress(len(words),total)
|
||
final = self.status()
|
||
if not final[2] & 2 or _u32(final,8) != generation or final[7] != count:
|
||
raise ValueError('Snapshot changed during download')
|
||
records = []
|
||
for offset in range(0,total,width):
|
||
record = {'time_ms': _u32(words,offset), 'sequence': _u32(words,offset+2), 'event': words[offset+4]}
|
||
if records and record['sequence'] != ((records[-1]['sequence']+1) & 0xffffffff):
|
||
raise ValueError('Snapshot record sequence is not contiguous')
|
||
for i,c in enumerate(self.mapping['channels']):
|
||
v=words[offset+HEADER_WORDS+i]
|
||
record[c['key']] = v-65536 if c['type']=='i16' and v>=32768 else v
|
||
records.append(record)
|
||
if cancelled():
|
||
raise InterruptedError('Snapshot download cancelled; snapshot retained')
|
||
self.release(generation)
|
||
return records
|