110 lines
3.7 KiB
Python
Raw Normal View History

2016-11-21 15:47:29 +01:00
"""Asyncio protocol implementation for handling telegrams."""
from functools import partial
2016-11-21 15:47:29 +01:00
import asyncio
import logging
from serial_asyncio import create_serial_connection
2017-01-10 20:09:33 +01:00
from dsmr_parser import telegram_specifications
from dsmr_parser.clients.telegram_buffer import TelegramBuffer
from dsmr_parser.exceptions import ParseError
from dsmr_parser.parsers import TelegramParser
2017-01-10 20:09:33 +01:00
from dsmr_parser.clients.settings import SERIAL_SETTINGS_V2_2, \
SERIAL_SETTINGS_V4
2016-11-21 15:47:29 +01:00
2017-01-03 22:27:39 +01:00
def create_dsmr_protocol(dsmr_version, telegram_callback, loop=None):
2016-12-08 22:14:43 +01:00
"""Creates a DSMR asyncio protocol."""
2016-11-21 15:47:29 +01:00
if dsmr_version == '2.2':
specification = telegram_specifications.V2_2
2016-11-21 15:47:29 +01:00
serial_settings = SERIAL_SETTINGS_V2_2
elif dsmr_version == '4':
specification = telegram_specifications.V4
2016-11-21 15:47:29 +01:00
serial_settings = SERIAL_SETTINGS_V4
else:
raise NotImplementedError("No telegram parser found for version: %s",
dsmr_version)
2016-11-21 15:47:29 +01:00
protocol = partial(DSMRProtocol, loop, TelegramParser(specification),
2016-11-21 15:47:29 +01:00
telegram_callback=telegram_callback)
2016-12-08 22:14:43 +01:00
return protocol, serial_settings
def create_dsmr_reader(port, dsmr_version, telegram_callback, loop=None):
"""Creates a DSMR asyncio protocol coroutine using serial port."""
2017-01-03 22:27:39 +01:00
protocol, serial_settings = create_dsmr_protocol(
2016-12-08 22:14:43 +01:00
dsmr_version, telegram_callback, loop=None)
serial_settings['url'] = port
2016-11-21 15:47:29 +01:00
conn = create_serial_connection(loop, protocol, **serial_settings)
2016-12-08 22:14:43 +01:00
return conn
2016-11-21 15:47:29 +01:00
def create_tcp_dsmr_reader(host, port, dsmr_version,
telegram_callback, loop=None):
2016-12-08 22:14:43 +01:00
"""Creates a DSMR asyncio protocol coroutine using TCP connection."""
2017-01-03 22:27:39 +01:00
protocol, _ = create_dsmr_protocol(
2016-12-08 22:14:43 +01:00
dsmr_version, telegram_callback, loop=None)
conn = loop.create_connection(protocol, host, port)
2016-11-21 15:47:29 +01:00
return conn
class DSMRProtocol(asyncio.Protocol):
"""Assemble and handle incoming data into complete DSM telegrams."""
transport = None
telegram_callback = None
def __init__(self, loop, telegram_parser, telegram_callback=None):
"""Initialize class."""
self.loop = loop
self.log = logging.getLogger(__name__)
self.telegram_parser = telegram_parser
# callback to call on complete telegram
self.telegram_callback = telegram_callback
# buffer to keep incomplete incoming data
2017-01-07 22:29:02 +01:00
self.telegram_buffer = TelegramBuffer()
# keep a lock until the connection is closed
self._closed = asyncio.Event()
2016-11-21 15:47:29 +01:00
def connection_made(self, transport):
"""Just logging for now."""
self.transport = transport
self.log.debug('connected')
def data_received(self, data):
"""Add incoming data to buffer."""
data = data.decode('ascii')
2017-01-07 21:26:21 +01:00
self.log.debug('received data: %s', data)
self.telegram_buffer.append(data)
2017-01-08 11:28:15 +01:00
for telegram in self.telegram_buffer.get_all():
self.handle_telegram(telegram)
2016-11-21 15:47:29 +01:00
def connection_lost(self, exc):
"""Stop when connection is lost."""
if exc:
self.log.exception('disconnected due to exception')
else:
self.log.info('disconnected because of close/abort.')
self._closed.set()
2016-11-21 15:47:29 +01:00
def handle_telegram(self, telegram):
"""Send off parsed telegram to handling callback."""
self.log.debug('got telegram: %s', telegram)
2017-01-07 22:29:02 +01:00
try:
parsed_telegram = self.telegram_parser.parse(telegram)
except ParseError:
self.log.exception("failed to parse telegram")
else:
self.telegram_callback(parsed_telegram)
@asyncio.coroutine
def wait_closed(self):
"""Wait until connection is closed."""
yield from self._closed.wait()