Files
templates/python/set_devices/gas_logger.py

191 lines
9.0 KiB
Python

"""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 = ['/* 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