hat.drivers.copp

Connection oriented presentation protocol

  1"""Connection oriented presentation protocol"""
  2
  3import asyncio
  4import importlib.resources
  5import logging
  6import typing
  7
  8from hat import aio
  9from hat import asn1
 10from hat import json
 11
 12from hat.drivers import cosp
 13from hat.drivers import net
 14
 15
 16mlog = logging.getLogger(__name__)
 17
 18with importlib.resources.open_text(__package__, 'asn1_repo.json') as _f:
 19    _encoder = asn1.ber.BerEncoder(
 20        asn1.repository_from_json(
 21            json.decode_stream(_f)))
 22
 23
 24class TcpConnectionInfo(typing.NamedTuple):
 25    name: str | None
 26    local_addr: net.TcpAddress
 27    local_tsel: int | None
 28    local_ssel: int | None
 29    local_psel: int | None
 30    remote_addr: net.TcpAddress
 31    remote_tsel: int | None
 32    remote_ssel: int | None
 33    remote_psel: int | None
 34
 35
 36class UnixConnectionInfo(typing.NamedTuple):
 37    name: str | None
 38    addr: net.UnixAddress
 39    local_tsel: int | None
 40    local_ssel: int | None
 41    local_psel: int | None
 42    remote_tsel: int | None
 43    remote_ssel: int | None
 44    remote_psel: int | None
 45
 46
 47ConnectionInfo: typing.TypeAlias = TcpConnectionInfo | UnixConnectionInfo
 48
 49IdentifiedEntity: typing.TypeAlias = tuple[asn1.ObjectIdentifier, asn1.Entity]
 50"""Identified entity"""
 51
 52ValidateCb: typing.TypeAlias = aio.AsyncCallable[['SyntaxNames',
 53                                                  IdentifiedEntity],
 54                                                 IdentifiedEntity | None]
 55"""Validate callback"""
 56
 57ConnectionCb: typing.TypeAlias = aio.AsyncCallable[['Connection'], None]
 58"""Connection callback"""
 59
 60
 61class SyntaxNames:
 62    """Syntax name registry
 63
 64    Args:
 65        syntax_names: list of ASN.1 ObjectIdentifiers representing syntax names
 66
 67    """
 68
 69    def __init__(self, syntax_names: list[asn1.ObjectIdentifier]):
 70        self._syntax_id_names = {(i * 2 + 1): name
 71                                 for i, name in enumerate(syntax_names)}
 72        self._syntax_name_ids = {v: k
 73                                 for k, v in self._syntax_id_names.items()}
 74
 75    def get_name(self, syntax_id: int) -> asn1.ObjectIdentifier:
 76        """Get syntax name associated with id"""
 77        return self._syntax_id_names[syntax_id]
 78
 79    def get_id(self, syntax_name: asn1.ObjectIdentifier) -> int:
 80        """Get syntax id associated with name"""
 81        return self._syntax_name_ids[syntax_name]
 82
 83
 84async def connect(addr: net.StreamAddress,
 85                  syntax_names: SyntaxNames,
 86                  user_data: IdentifiedEntity | None = None,
 87                  *,
 88                  local_psel: int | None = None,
 89                  remote_psel: int | None = None,
 90                  copp_receive_queue_size: int = 1024,
 91                  copp_send_queue_size: int = 1024,
 92                  **kwargs
 93                  ) -> 'Connection':
 94    """Connect to COPP server
 95
 96    Additional arguments are passed directly to `hat.drivers.cosp.connect`.
 97
 98    """
 99    cp_ppdu = _cp_ppdu(syntax_names, local_psel, remote_psel, user_data)
100    cp_ppdu_data = _encode('CP-type', cp_ppdu)
101    conn = await cosp.connect(addr, cp_ppdu_data, **kwargs)
102
103    try:
104        cpa_ppdu = _decode('CPA-PPDU', conn.conn_res_user_data)
105        _validate_connect_response(cp_ppdu, cpa_ppdu)
106
107        calling_psel, called_psel = _get_psels(cp_ppdu)
108        return Connection(conn, syntax_names, cp_ppdu, cpa_ppdu,
109                          calling_psel, called_psel,
110                          copp_receive_queue_size, copp_send_queue_size)
111
112    except Exception:
113        await aio.uncancellable(_close_cosp(conn, _arp_ppdu(), mlog))
114        raise
115
116
117async def listen(validate_cb: ValidateCb,
118                 connection_cb: ConnectionCb,
119                 addr: net.StreamAddress = net.TcpAddress('0.0.0.0', 102),
120                 *,
121                 bind_connections: bool = False,
122                 copp_receive_queue_size: int = 1024,
123                 copp_send_queue_size: int = 1024,
124                 **kwargs
125                 ) -> 'Server':
126    """Create COPP listening server
127
128    Additional arguments are passed directly to `hat.drivers.cosp.listen`.
129
130    Args:
131        validate_cb: callback function or coroutine called on new
132            incomming connection request prior to creating connection object
133        connection_cb: new connection callback
134        addr: local listening address
135
136    """
137    server = Server()
138    server._validate_cb = validate_cb
139    server._connection_cb = connection_cb
140    server._bind_connections = bind_connections
141    server._receive_queue_size = copp_receive_queue_size
142    server._send_queue_size = copp_send_queue_size
143    server._log = mlog
144
145    server._srv = await cosp.listen(server._on_validate,
146                                    server._on_connection,
147                                    addr,
148                                    bind_connections=False,
149                                    **kwargs)
150
151    server._log = _create_server_logger(server._srv.info)
152
153    return server
154
155
156class Server(aio.Resource):
157    """COPP listening server
158
159    For creating new server see `listen`.
160
161    """
162
163    @property
164    def async_group(self) -> aio.Group:
165        """Async group"""
166        return self._srv.async_group
167
168    @property
169    def info(self) -> net.ServerInfo:
170        """Server info"""
171        return self._srv.info
172
173    async def _on_validate(self, user_data):
174        cp_ppdu = _decode('CP-type', user_data)
175        cp_params = cp_ppdu['normal-mode-parameters']
176        called_psel_data = cp_params.get('called-presentation-selector')
177        called_psel = (int.from_bytes(called_psel_data, 'big')
178                       if called_psel_data else None)
179        cp_pdv_list = cp_params['user-data'][1][0]
180        syntax_names = _sytax_names_from_cp_ppdu(cp_ppdu)
181        cp_user_data = (
182            syntax_names.get_name(
183                cp_pdv_list['presentation-context-identifier']),
184            cp_pdv_list['presentation-data-values'][1])
185
186        cpa_user_data = await aio.call(self._validate_cb, syntax_names,
187                                       cp_user_data)
188
189        cpa_ppdu = _cpa_ppdu(syntax_names, called_psel, cpa_user_data)
190        cpa_ppdu_data = _encode('CPA-PPDU', cpa_ppdu)
191        return cpa_ppdu_data
192
193    async def _on_connection(self, cosp_conn):
194        try:
195            try:
196                cp_ppdu = _decode('CP-type', cosp_conn.conn_req_user_data)
197                cpa_ppdu = _decode('CPA-PPDU', cosp_conn.conn_res_user_data)
198
199                syntax_names = _sytax_names_from_cp_ppdu(cp_ppdu)
200                calling_psel, called_psel = _get_psels(cp_ppdu)
201
202                conn = Connection(cosp_conn, syntax_names, cp_ppdu, cpa_ppdu,
203                                  called_psel, calling_psel,
204                                  self._receive_queue_size,
205                                  self._send_queue_size)
206
207            except Exception:
208                await aio.uncancellable(
209                    _close_cosp(cosp_conn, _arp_ppdu(), self._log))
210                raise
211
212            try:
213                await aio.call(self._connection_cb, conn)
214
215            except BaseException:
216                await aio.uncancellable(conn.async_close())
217                raise
218
219        except Exception as e:
220            self._log.error("error creating new incomming connection: %s",
221                            e, exc_info=e)
222            return
223
224        if not self._bind_connections:
225            return
226
227        try:
228            await conn.wait_closed()
229
230        except BaseException:
231            await aio.uncancellable(conn.async_close())
232            raise
233
234
235class Connection(aio.Resource):
236    """COPP connection
237
238    For creating new connection see `connect` or `listen`.
239
240    """
241
242    def __init__(self,
243                 conn: cosp.Connection,
244                 syntax_names: SyntaxNames,
245                 cp_ppdu: asn1.Value,
246                 cpa_ppdu: asn1.Value,
247                 local_psel: int | None,
248                 remote_psel: int | None,
249                 receive_queue_size: int,
250                 send_queue_size: int):
251        cp_user_data = cp_ppdu['normal-mode-parameters']['user-data']
252        cpa_user_data = cpa_ppdu['normal-mode-parameters']['user-data']
253
254        conn_req_user_data = (
255            syntax_names.get_name(
256                cp_user_data[1][0]['presentation-context-identifier']),
257            cp_user_data[1][0]['presentation-data-values'][1])
258        conn_res_user_data = (
259            syntax_names.get_name(
260                cpa_user_data[1][0]['presentation-context-identifier']),
261            cpa_user_data[1][0]['presentation-data-values'][1])
262
263        self._conn = conn
264        self._syntax_names = syntax_names
265        self._conn_req_user_data = conn_req_user_data
266        self._conn_res_user_data = conn_res_user_data
267        self._loop = asyncio.get_running_loop()
268        self._info = _connection_info_from_cosp(info=conn.info,
269                                                local_psel=local_psel,
270                                                remote_psel=remote_psel)
271        self._close_ppdu = _arp_ppdu()
272        self._receive_queue = aio.Queue(receive_queue_size)
273        self._send_queue = aio.Queue(send_queue_size)
274        self._async_group = aio.Group()
275        self._log = _create_connection_logger(self._info)
276
277        self.async_group.spawn(aio.call_on_cancel, self._on_close)
278        self.async_group.spawn(self._receive_loop)
279        self.async_group.spawn(self._send_loop)
280        self.async_group.spawn(aio.call_on_done, conn.wait_closing(),
281                               self.close)
282
283    @property
284    def async_group(self) -> aio.Group:
285        """Async group"""
286        return self._async_group
287
288    @property
289    def info(self) -> ConnectionInfo:
290        """Connection info"""
291        return self._info
292
293    @property
294    def syntax_names(self) -> SyntaxNames:
295        """Syntax names"""
296        return self._syntax_names
297
298    @property
299    def conn_req_user_data(self) -> IdentifiedEntity:
300        """Connect request's user data"""
301        return self._conn_req_user_data
302
303    @property
304    def conn_res_user_data(self) -> IdentifiedEntity:
305        """Connect response's user data"""
306        return self._conn_res_user_data
307
308    def close(self, user_data: IdentifiedEntity | None = None):
309        """Close connection"""
310        self._close(_aru_ppdu(self._syntax_names, user_data))
311
312    async def async_close(self, user_data: IdentifiedEntity | None = None):
313        """Async close"""
314        self.close(user_data)
315        await self.wait_closed()
316
317    async def receive(self) -> IdentifiedEntity:
318        """Receive data"""
319        try:
320            return await self._receive_queue.get()
321
322        except aio.QueueClosedError:
323            raise ConnectionError()
324
325    async def send(self, data: IdentifiedEntity):
326        """Send data"""
327        try:
328            await self._send_queue.put((data, None))
329
330        except aio.QueueClosedError:
331            raise ConnectionError()
332
333    async def drain(self):
334        """Drain output buffer"""
335        try:
336            future = self._loop.create_future()
337            await self._send_queue.put((None, future))
338            await future
339
340        except aio.QueueClosedError:
341            raise ConnectionError()
342
343    async def _on_close(self):
344        await _close_cosp(self._conn, self._close_ppdu, self._log)
345
346    def _close(self, ppdu):
347        if not self.is_open:
348            return
349
350        self._close_ppdu = ppdu
351        self._async_group.close()
352
353    async def _receive_loop(self):
354        try:
355            while True:
356                cosp_data = await self._conn.receive()
357
358                user_data = _decode('User-data', cosp_data)
359
360                pdv_list = user_data[1][0]
361                syntax_name = self._syntax_names.get_name(
362                    pdv_list['presentation-context-identifier'])
363                data = pdv_list['presentation-data-values'][1]
364
365                await self._receive_queue.put((syntax_name, data))
366
367        except ConnectionError:
368            pass
369
370        except Exception as e:
371            self._log.error("receive loop error: %s", e, exc_info=e)
372
373        finally:
374            self._close(_arp_ppdu())
375            self._receive_queue.close()
376
377    async def _send_loop(self):
378        future = None
379        try:
380            while True:
381                data, future = await self._send_queue.get()
382
383                if data is None:
384                    await self._conn.drain()
385
386                else:
387                    user_data = _user_data(self._syntax_names, data)
388                    ppdu_data = _encode('User-data', user_data)
389
390                    await self._conn.send(ppdu_data)
391
392                if future and not future.done():
393                    future.set_result(None)
394
395        except ConnectionError:
396            pass
397
398        except Exception as e:
399            self._log.error("send loop error: %s", e, exc_info=e)
400
401        finally:
402            self._close(_arp_ppdu())
403            self._send_queue.close()
404
405            while True:
406                if future and not future.done():
407                    future.set_result(None)
408                if self._send_queue.empty():
409                    break
410                _, future = self._send_queue.get_nowait()
411
412
413async def _close_cosp(cosp_conn, ppdu, log):
414    try:
415        data = _encode('Abort-type', ppdu)
416
417    except Exception as e:
418        log.error("error encoding abort ppdu: %s", e, exc_info=e)
419        data = None
420
421    finally:
422        await cosp_conn.async_close(data)
423
424
425def _get_psels(cp_ppdu):
426    cp_params = cp_ppdu['normal-mode-parameters']
427    calling_psel_data = cp_params.get('calling-presentation-selector')
428    calling_psel = (int.from_bytes(calling_psel_data, 'big')
429                    if calling_psel_data else None)
430    called_psel_data = cp_params.get('called-presentation-selector')
431    called_psel = (int.from_bytes(called_psel_data, 'big')
432                   if called_psel_data else None)
433    return calling_psel, called_psel
434
435
436def _validate_connect_response(cp_ppdu, cpa_ppdu):
437    cp_params = cp_ppdu['normal-mode-parameters']
438    cpa_params = cpa_ppdu['normal-mode-parameters']
439    called_psel_data = cp_params.get('called-presentation-selector')
440    responding_psel_data = cpa_params.get('responding-presentation-selector')
441
442    if called_psel_data and responding_psel_data:
443        called_psel = int.from_bytes(called_psel_data, 'big')
444        responding_psel = int.from_bytes(responding_psel_data, 'big')
445
446        if called_psel != responding_psel:
447            raise Exception('presentation selectors not matching')
448
449    result_list = cpa_params['presentation-context-definition-result-list']
450    if any(i['result'] != 0 for i in result_list):
451        raise Exception('presentation context not accepted')
452
453
454def _cp_ppdu(syntax_names, calling_psel, called_psel, user_data):
455    cp_params = {
456        'presentation-context-definition-list': [
457            {'presentation-context-identifier': i,
458             'abstract-syntax-name': name,
459             'transfer-syntax-name-list': [_encoder.syntax_name]}
460            for i, name in syntax_names._syntax_id_names.items()]}
461
462    if calling_psel is not None:
463        cp_params['calling-presentation-selector'] = \
464            calling_psel.to_bytes(4, 'big')
465
466    if called_psel is not None:
467        cp_params['called-presentation-selector'] = \
468            called_psel.to_bytes(4, 'big')
469
470    if user_data:
471        cp_params['user-data'] = _user_data(syntax_names, user_data)
472
473    return {
474        'mode-selector': {
475            'mode-value': 1},
476        'normal-mode-parameters': cp_params}
477
478
479def _cpa_ppdu(syntax_names, responding_psel, user_data):
480    cpa_params = {
481        'presentation-context-definition-result-list': [
482            {'result': 0,
483             'transfer-syntax-name': _encoder.syntax_name}
484            for _ in syntax_names._syntax_id_names.keys()]}
485
486    if responding_psel is not None:
487        cpa_params['responding-presentation-selector'] = \
488            responding_psel.to_bytes(4, 'big')
489
490    if user_data:
491        cpa_params['user-data'] = _user_data(syntax_names, user_data)
492
493    return {
494        'mode-selector': {
495            'mode-value': 1},
496        'normal-mode-parameters': cpa_params}
497
498
499def _aru_ppdu(syntax_names, user_data):
500    aru_params = {}
501
502    if user_data:
503        aru_params['user-data'] = _user_data(syntax_names, user_data)
504
505    return 'aru-ppdu', ('normal-mode-parameters', aru_params)
506
507
508def _arp_ppdu():
509    return 'arp-ppdu', {}
510
511
512def _user_data(syntax_names, user_data):
513    return 'fully-encoded-data', [{
514        'presentation-context-identifier': syntax_names.get_id(user_data[0]),
515        'presentation-data-values': (
516            'single-ASN1-type', user_data[1])}]
517
518
519def _sytax_names_from_cp_ppdu(cp_ppdu):
520    cp_params = cp_ppdu['normal-mode-parameters']
521    syntax_names = SyntaxNames([])
522    syntax_names._syntax_id_names = {
523        i['presentation-context-identifier']: i['abstract-syntax-name']
524        for i in cp_params['presentation-context-definition-list']}
525    syntax_names._syntax_name_ids = {
526        v: k for k, v in syntax_names._syntax_id_names.items()}
527    return syntax_names
528
529
530def _encode(name, value):
531    return _encoder.encode(asn1.TypeRef('ISO8823-PRESENTATION', name), value)
532
533
534def _decode(name, data):
535    res, _ = _encoder.decode(asn1.TypeRef('ISO8823-PRESENTATION', name),
536                             memoryview(data))
537    return res
538
539
540def _connection_info_to_net(info):
541    if isinstance(info, TcpConnectionInfo):
542        return net.TcpConnectionInfo(name=info.name,
543                                     local_addr=info.local_addr,
544                                     remote_addr=info.remote_addr)
545
546    if isinstance(info, UnixConnectionInfo):
547        return net.UnixConnectionInfo(name=info.name,
548                                      addr=info.addr)
549
550    raise TypeError('unsupported info type')
551
552
553def _connection_info_from_cosp(info, local_psel, remote_psel):
554    if isinstance(info, cosp.TcpConnectionInfo):
555        return TcpConnectionInfo(name=info.name,
556                                 local_addr=info.local_addr,
557                                 local_tsel=info.local_tsel,
558                                 local_ssel=info.local_ssel,
559                                 local_psel=local_psel,
560                                 remote_addr=info.remote_addr,
561                                 remote_tsel=info.remote_tsel,
562                                 remote_ssel=info.remote_ssel,
563                                 remote_psel=remote_psel)
564
565    if isinstance(info, cosp.UnixConnectionInfo):
566        return UnixConnectionInfo(name=info.name,
567                                  addr=info.addr,
568                                  local_tsel=info.local_tsel,
569                                  local_ssel=info.local_ssel,
570                                  local_psel=local_psel,
571                                  remote_tsel=info.remote_tsel,
572                                  remote_ssel=info.remote_ssel,
573                                  remote_psel=remote_psel)
574
575    raise TypeError('unsupported info type')
576
577
578def _create_server_logger(info):
579    extra = {'meta': {'type': 'CoppServer',
580                      **net.server_info_to_json(info)}}
581
582    return logging.LoggerAdapter(mlog, extra)
583
584
585def _create_connection_logger(info):
586    net_info = _connection_info_to_net(info)
587
588    extra = {'meta': {'type': 'CoppConnection',
589                      **net.connection_info_to_json(net_info)}}
590
591    return logging.LoggerAdapter(mlog, extra)
mlog = <Logger hat.drivers.copp (WARNING)>
class TcpConnectionInfo(typing.NamedTuple):
25class TcpConnectionInfo(typing.NamedTuple):
26    name: str | None
27    local_addr: net.TcpAddress
28    local_tsel: int | None
29    local_ssel: int | None
30    local_psel: int | None
31    remote_addr: net.TcpAddress
32    remote_tsel: int | None
33    remote_ssel: int | None
34    remote_psel: int | None

TcpConnectionInfo(name, local_addr, local_tsel, local_ssel, local_psel, remote_addr, remote_tsel, remote_ssel, remote_psel)

TcpConnectionInfo( name: str | None, local_addr: hat.drivers.net.TcpAddress, local_tsel: int | None, local_ssel: int | None, local_psel: int | None, remote_addr: hat.drivers.net.TcpAddress, remote_tsel: int | None, remote_ssel: int | None, remote_psel: int | None)

Create new instance of TcpConnectionInfo(name, local_addr, local_tsel, local_ssel, local_psel, remote_addr, remote_tsel, remote_ssel, remote_psel)

name: str | None

Alias for field number 0

Alias for field number 1

local_tsel: int | None

Alias for field number 2

local_ssel: int | None

Alias for field number 3

local_psel: int | None

Alias for field number 4

Alias for field number 5

remote_tsel: int | None

Alias for field number 6

remote_ssel: int | None

Alias for field number 7

remote_psel: int | None

Alias for field number 8

class UnixConnectionInfo(typing.NamedTuple):
37class UnixConnectionInfo(typing.NamedTuple):
38    name: str | None
39    addr: net.UnixAddress
40    local_tsel: int | None
41    local_ssel: int | None
42    local_psel: int | None
43    remote_tsel: int | None
44    remote_ssel: int | None
45    remote_psel: int | None

UnixConnectionInfo(name, addr, local_tsel, local_ssel, local_psel, remote_tsel, remote_ssel, remote_psel)

UnixConnectionInfo( name: str | None, addr: pathlib.Path, local_tsel: int | None, local_ssel: int | None, local_psel: int | None, remote_tsel: int | None, remote_ssel: int | None, remote_psel: int | None)

Create new instance of UnixConnectionInfo(name, addr, local_tsel, local_ssel, local_psel, remote_tsel, remote_ssel, remote_psel)

name: str | None

Alias for field number 0

addr: pathlib.Path

Alias for field number 1

local_tsel: int | None

Alias for field number 2

local_ssel: int | None

Alias for field number 3

local_psel: int | None

Alias for field number 4

remote_tsel: int | None

Alias for field number 5

remote_ssel: int | None

Alias for field number 6

remote_psel: int | None

Alias for field number 7

ConnectionInfo: TypeAlias = TcpConnectionInfo | UnixConnectionInfo
IdentifiedEntity: TypeAlias = tuple[tuple[int, ...], hat.asn1.common.Entity]

Identified entity

ValidateCb: TypeAlias = Callable[[ForwardRef('SyntaxNames'), tuple[tuple[int, ...], hat.asn1.common.Entity]], tuple[tuple[int, ...], hat.asn1.common.Entity] | None | Awaitable[tuple[tuple[int, ...], hat.asn1.common.Entity] | None]]

Validate callback

ConnectionCb: TypeAlias = Callable[[ForwardRef('Connection')], None | Awaitable[None]]

Connection callback

class SyntaxNames:
62class SyntaxNames:
63    """Syntax name registry
64
65    Args:
66        syntax_names: list of ASN.1 ObjectIdentifiers representing syntax names
67
68    """
69
70    def __init__(self, syntax_names: list[asn1.ObjectIdentifier]):
71        self._syntax_id_names = {(i * 2 + 1): name
72                                 for i, name in enumerate(syntax_names)}
73        self._syntax_name_ids = {v: k
74                                 for k, v in self._syntax_id_names.items()}
75
76    def get_name(self, syntax_id: int) -> asn1.ObjectIdentifier:
77        """Get syntax name associated with id"""
78        return self._syntax_id_names[syntax_id]
79
80    def get_id(self, syntax_name: asn1.ObjectIdentifier) -> int:
81        """Get syntax id associated with name"""
82        return self._syntax_name_ids[syntax_name]

Syntax name registry

Arguments:
  • syntax_names: list of ASN.1 ObjectIdentifiers representing syntax names
SyntaxNames(syntax_names: list[tuple[int, ...]])
70    def __init__(self, syntax_names: list[asn1.ObjectIdentifier]):
71        self._syntax_id_names = {(i * 2 + 1): name
72                                 for i, name in enumerate(syntax_names)}
73        self._syntax_name_ids = {v: k
74                                 for k, v in self._syntax_id_names.items()}
def get_name(self, syntax_id: int) -> tuple[int, ...]:
76    def get_name(self, syntax_id: int) -> asn1.ObjectIdentifier:
77        """Get syntax name associated with id"""
78        return self._syntax_id_names[syntax_id]

Get syntax name associated with id

def get_id(self, syntax_name: tuple[int, ...]) -> int:
80    def get_id(self, syntax_name: asn1.ObjectIdentifier) -> int:
81        """Get syntax id associated with name"""
82        return self._syntax_name_ids[syntax_name]

Get syntax id associated with name

async def connect( addr: hat.drivers.net.TcpAddress | pathlib.Path, syntax_names: SyntaxNames, user_data: tuple[tuple[int, ...], hat.asn1.common.Entity] | None = None, *, local_psel: int | None = None, remote_psel: int | None = None, copp_receive_queue_size: int = 1024, copp_send_queue_size: int = 1024, **kwargs) -> Connection:
 85async def connect(addr: net.StreamAddress,
 86                  syntax_names: SyntaxNames,
 87                  user_data: IdentifiedEntity | None = None,
 88                  *,
 89                  local_psel: int | None = None,
 90                  remote_psel: int | None = None,
 91                  copp_receive_queue_size: int = 1024,
 92                  copp_send_queue_size: int = 1024,
 93                  **kwargs
 94                  ) -> 'Connection':
 95    """Connect to COPP server
 96
 97    Additional arguments are passed directly to `hat.drivers.cosp.connect`.
 98
 99    """
100    cp_ppdu = _cp_ppdu(syntax_names, local_psel, remote_psel, user_data)
101    cp_ppdu_data = _encode('CP-type', cp_ppdu)
102    conn = await cosp.connect(addr, cp_ppdu_data, **kwargs)
103
104    try:
105        cpa_ppdu = _decode('CPA-PPDU', conn.conn_res_user_data)
106        _validate_connect_response(cp_ppdu, cpa_ppdu)
107
108        calling_psel, called_psel = _get_psels(cp_ppdu)
109        return Connection(conn, syntax_names, cp_ppdu, cpa_ppdu,
110                          calling_psel, called_psel,
111                          copp_receive_queue_size, copp_send_queue_size)
112
113    except Exception:
114        await aio.uncancellable(_close_cosp(conn, _arp_ppdu(), mlog))
115        raise

Connect to COPP server

Additional arguments are passed directly to hat.drivers.cosp.connect.

async def listen( validate_cb: Callable[[SyntaxNames, tuple[tuple[int, ...], hat.asn1.common.Entity]], tuple[tuple[int, ...], hat.asn1.common.Entity] | None | Awaitable[tuple[tuple[int, ...], hat.asn1.common.Entity] | None]], connection_cb: Callable[[Connection], None | Awaitable[None]], addr: hat.drivers.net.TcpAddress | pathlib.Path = TcpAddress(host='0.0.0.0', port=102), *, bind_connections: bool = False, copp_receive_queue_size: int = 1024, copp_send_queue_size: int = 1024, **kwargs) -> Server:
118async def listen(validate_cb: ValidateCb,
119                 connection_cb: ConnectionCb,
120                 addr: net.StreamAddress = net.TcpAddress('0.0.0.0', 102),
121                 *,
122                 bind_connections: bool = False,
123                 copp_receive_queue_size: int = 1024,
124                 copp_send_queue_size: int = 1024,
125                 **kwargs
126                 ) -> 'Server':
127    """Create COPP listening server
128
129    Additional arguments are passed directly to `hat.drivers.cosp.listen`.
130
131    Args:
132        validate_cb: callback function or coroutine called on new
133            incomming connection request prior to creating connection object
134        connection_cb: new connection callback
135        addr: local listening address
136
137    """
138    server = Server()
139    server._validate_cb = validate_cb
140    server._connection_cb = connection_cb
141    server._bind_connections = bind_connections
142    server._receive_queue_size = copp_receive_queue_size
143    server._send_queue_size = copp_send_queue_size
144    server._log = mlog
145
146    server._srv = await cosp.listen(server._on_validate,
147                                    server._on_connection,
148                                    addr,
149                                    bind_connections=False,
150                                    **kwargs)
151
152    server._log = _create_server_logger(server._srv.info)
153
154    return server

Create COPP listening server

Additional arguments are passed directly to hat.drivers.cosp.listen.

Arguments:
  • validate_cb: callback function or coroutine called on new incomming connection request prior to creating connection object
  • connection_cb: new connection callback
  • addr: local listening address
class Server(hat.aio.group.Resource):
157class Server(aio.Resource):
158    """COPP listening server
159
160    For creating new server see `listen`.
161
162    """
163
164    @property
165    def async_group(self) -> aio.Group:
166        """Async group"""
167        return self._srv.async_group
168
169    @property
170    def info(self) -> net.ServerInfo:
171        """Server info"""
172        return self._srv.info
173
174    async def _on_validate(self, user_data):
175        cp_ppdu = _decode('CP-type', user_data)
176        cp_params = cp_ppdu['normal-mode-parameters']
177        called_psel_data = cp_params.get('called-presentation-selector')
178        called_psel = (int.from_bytes(called_psel_data, 'big')
179                       if called_psel_data else None)
180        cp_pdv_list = cp_params['user-data'][1][0]
181        syntax_names = _sytax_names_from_cp_ppdu(cp_ppdu)
182        cp_user_data = (
183            syntax_names.get_name(
184                cp_pdv_list['presentation-context-identifier']),
185            cp_pdv_list['presentation-data-values'][1])
186
187        cpa_user_data = await aio.call(self._validate_cb, syntax_names,
188                                       cp_user_data)
189
190        cpa_ppdu = _cpa_ppdu(syntax_names, called_psel, cpa_user_data)
191        cpa_ppdu_data = _encode('CPA-PPDU', cpa_ppdu)
192        return cpa_ppdu_data
193
194    async def _on_connection(self, cosp_conn):
195        try:
196            try:
197                cp_ppdu = _decode('CP-type', cosp_conn.conn_req_user_data)
198                cpa_ppdu = _decode('CPA-PPDU', cosp_conn.conn_res_user_data)
199
200                syntax_names = _sytax_names_from_cp_ppdu(cp_ppdu)
201                calling_psel, called_psel = _get_psels(cp_ppdu)
202
203                conn = Connection(cosp_conn, syntax_names, cp_ppdu, cpa_ppdu,
204                                  called_psel, calling_psel,
205                                  self._receive_queue_size,
206                                  self._send_queue_size)
207
208            except Exception:
209                await aio.uncancellable(
210                    _close_cosp(cosp_conn, _arp_ppdu(), self._log))
211                raise
212
213            try:
214                await aio.call(self._connection_cb, conn)
215
216            except BaseException:
217                await aio.uncancellable(conn.async_close())
218                raise
219
220        except Exception as e:
221            self._log.error("error creating new incomming connection: %s",
222                            e, exc_info=e)
223            return
224
225        if not self._bind_connections:
226            return
227
228        try:
229            await conn.wait_closed()
230
231        except BaseException:
232            await aio.uncancellable(conn.async_close())
233            raise

COPP listening server

For creating new server see listen.

async_group: hat.aio.group.Group
164    @property
165    def async_group(self) -> aio.Group:
166        """Async group"""
167        return self._srv.async_group

Async group

169    @property
170    def info(self) -> net.ServerInfo:
171        """Server info"""
172        return self._srv.info

Server info

class Connection(hat.aio.group.Resource):
236class Connection(aio.Resource):
237    """COPP connection
238
239    For creating new connection see `connect` or `listen`.
240
241    """
242
243    def __init__(self,
244                 conn: cosp.Connection,
245                 syntax_names: SyntaxNames,
246                 cp_ppdu: asn1.Value,
247                 cpa_ppdu: asn1.Value,
248                 local_psel: int | None,
249                 remote_psel: int | None,
250                 receive_queue_size: int,
251                 send_queue_size: int):
252        cp_user_data = cp_ppdu['normal-mode-parameters']['user-data']
253        cpa_user_data = cpa_ppdu['normal-mode-parameters']['user-data']
254
255        conn_req_user_data = (
256            syntax_names.get_name(
257                cp_user_data[1][0]['presentation-context-identifier']),
258            cp_user_data[1][0]['presentation-data-values'][1])
259        conn_res_user_data = (
260            syntax_names.get_name(
261                cpa_user_data[1][0]['presentation-context-identifier']),
262            cpa_user_data[1][0]['presentation-data-values'][1])
263
264        self._conn = conn
265        self._syntax_names = syntax_names
266        self._conn_req_user_data = conn_req_user_data
267        self._conn_res_user_data = conn_res_user_data
268        self._loop = asyncio.get_running_loop()
269        self._info = _connection_info_from_cosp(info=conn.info,
270                                                local_psel=local_psel,
271                                                remote_psel=remote_psel)
272        self._close_ppdu = _arp_ppdu()
273        self._receive_queue = aio.Queue(receive_queue_size)
274        self._send_queue = aio.Queue(send_queue_size)
275        self._async_group = aio.Group()
276        self._log = _create_connection_logger(self._info)
277
278        self.async_group.spawn(aio.call_on_cancel, self._on_close)
279        self.async_group.spawn(self._receive_loop)
280        self.async_group.spawn(self._send_loop)
281        self.async_group.spawn(aio.call_on_done, conn.wait_closing(),
282                               self.close)
283
284    @property
285    def async_group(self) -> aio.Group:
286        """Async group"""
287        return self._async_group
288
289    @property
290    def info(self) -> ConnectionInfo:
291        """Connection info"""
292        return self._info
293
294    @property
295    def syntax_names(self) -> SyntaxNames:
296        """Syntax names"""
297        return self._syntax_names
298
299    @property
300    def conn_req_user_data(self) -> IdentifiedEntity:
301        """Connect request's user data"""
302        return self._conn_req_user_data
303
304    @property
305    def conn_res_user_data(self) -> IdentifiedEntity:
306        """Connect response's user data"""
307        return self._conn_res_user_data
308
309    def close(self, user_data: IdentifiedEntity | None = None):
310        """Close connection"""
311        self._close(_aru_ppdu(self._syntax_names, user_data))
312
313    async def async_close(self, user_data: IdentifiedEntity | None = None):
314        """Async close"""
315        self.close(user_data)
316        await self.wait_closed()
317
318    async def receive(self) -> IdentifiedEntity:
319        """Receive data"""
320        try:
321            return await self._receive_queue.get()
322
323        except aio.QueueClosedError:
324            raise ConnectionError()
325
326    async def send(self, data: IdentifiedEntity):
327        """Send data"""
328        try:
329            await self._send_queue.put((data, None))
330
331        except aio.QueueClosedError:
332            raise ConnectionError()
333
334    async def drain(self):
335        """Drain output buffer"""
336        try:
337            future = self._loop.create_future()
338            await self._send_queue.put((None, future))
339            await future
340
341        except aio.QueueClosedError:
342            raise ConnectionError()
343
344    async def _on_close(self):
345        await _close_cosp(self._conn, self._close_ppdu, self._log)
346
347    def _close(self, ppdu):
348        if not self.is_open:
349            return
350
351        self._close_ppdu = ppdu
352        self._async_group.close()
353
354    async def _receive_loop(self):
355        try:
356            while True:
357                cosp_data = await self._conn.receive()
358
359                user_data = _decode('User-data', cosp_data)
360
361                pdv_list = user_data[1][0]
362                syntax_name = self._syntax_names.get_name(
363                    pdv_list['presentation-context-identifier'])
364                data = pdv_list['presentation-data-values'][1]
365
366                await self._receive_queue.put((syntax_name, data))
367
368        except ConnectionError:
369            pass
370
371        except Exception as e:
372            self._log.error("receive loop error: %s", e, exc_info=e)
373
374        finally:
375            self._close(_arp_ppdu())
376            self._receive_queue.close()
377
378    async def _send_loop(self):
379        future = None
380        try:
381            while True:
382                data, future = await self._send_queue.get()
383
384                if data is None:
385                    await self._conn.drain()
386
387                else:
388                    user_data = _user_data(self._syntax_names, data)
389                    ppdu_data = _encode('User-data', user_data)
390
391                    await self._conn.send(ppdu_data)
392
393                if future and not future.done():
394                    future.set_result(None)
395
396        except ConnectionError:
397            pass
398
399        except Exception as e:
400            self._log.error("send loop error: %s", e, exc_info=e)
401
402        finally:
403            self._close(_arp_ppdu())
404            self._send_queue.close()
405
406            while True:
407                if future and not future.done():
408                    future.set_result(None)
409                if self._send_queue.empty():
410                    break
411                _, future = self._send_queue.get_nowait()

COPP connection

For creating new connection see connect or listen.

Connection( conn: Connection, syntax_names: SyntaxNames, cp_ppdu: bool | int | Collection[bool] | bytes | bytearray | memoryview | None | tuple[int, ...] | str | float | Tuple[str, ForwardRef('Value')] | Dict[str, ForwardRef('Value')] | Collection['Value'] | hat.asn1.common.Entity | hat.asn1.common.External | hat.asn1.common.EmbeddedPDV, cpa_ppdu: bool | int | Collection[bool] | bytes | bytearray | memoryview | None | tuple[int, ...] | str | float | Tuple[str, ForwardRef('Value')] | Dict[str, ForwardRef('Value')] | Collection['Value'] | hat.asn1.common.Entity | hat.asn1.common.External | hat.asn1.common.EmbeddedPDV, local_psel: int | None, remote_psel: int | None, receive_queue_size: int, send_queue_size: int)
243    def __init__(self,
244                 conn: cosp.Connection,
245                 syntax_names: SyntaxNames,
246                 cp_ppdu: asn1.Value,
247                 cpa_ppdu: asn1.Value,
248                 local_psel: int | None,
249                 remote_psel: int | None,
250                 receive_queue_size: int,
251                 send_queue_size: int):
252        cp_user_data = cp_ppdu['normal-mode-parameters']['user-data']
253        cpa_user_data = cpa_ppdu['normal-mode-parameters']['user-data']
254
255        conn_req_user_data = (
256            syntax_names.get_name(
257                cp_user_data[1][0]['presentation-context-identifier']),
258            cp_user_data[1][0]['presentation-data-values'][1])
259        conn_res_user_data = (
260            syntax_names.get_name(
261                cpa_user_data[1][0]['presentation-context-identifier']),
262            cpa_user_data[1][0]['presentation-data-values'][1])
263
264        self._conn = conn
265        self._syntax_names = syntax_names
266        self._conn_req_user_data = conn_req_user_data
267        self._conn_res_user_data = conn_res_user_data
268        self._loop = asyncio.get_running_loop()
269        self._info = _connection_info_from_cosp(info=conn.info,
270                                                local_psel=local_psel,
271                                                remote_psel=remote_psel)
272        self._close_ppdu = _arp_ppdu()
273        self._receive_queue = aio.Queue(receive_queue_size)
274        self._send_queue = aio.Queue(send_queue_size)
275        self._async_group = aio.Group()
276        self._log = _create_connection_logger(self._info)
277
278        self.async_group.spawn(aio.call_on_cancel, self._on_close)
279        self.async_group.spawn(self._receive_loop)
280        self.async_group.spawn(self._send_loop)
281        self.async_group.spawn(aio.call_on_done, conn.wait_closing(),
282                               self.close)
async_group: hat.aio.group.Group
284    @property
285    def async_group(self) -> aio.Group:
286        """Async group"""
287        return self._async_group

Async group

289    @property
290    def info(self) -> ConnectionInfo:
291        """Connection info"""
292        return self._info

Connection info

syntax_names: SyntaxNames
294    @property
295    def syntax_names(self) -> SyntaxNames:
296        """Syntax names"""
297        return self._syntax_names

Syntax names

conn_req_user_data: tuple[tuple[int, ...], hat.asn1.common.Entity]
299    @property
300    def conn_req_user_data(self) -> IdentifiedEntity:
301        """Connect request's user data"""
302        return self._conn_req_user_data

Connect request's user data

conn_res_user_data: tuple[tuple[int, ...], hat.asn1.common.Entity]
304    @property
305    def conn_res_user_data(self) -> IdentifiedEntity:
306        """Connect response's user data"""
307        return self._conn_res_user_data

Connect response's user data

def close( self, user_data: tuple[tuple[int, ...], hat.asn1.common.Entity] | None = None):
309    def close(self, user_data: IdentifiedEntity | None = None):
310        """Close connection"""
311        self._close(_aru_ppdu(self._syntax_names, user_data))

Close connection

async def async_close( self, user_data: tuple[tuple[int, ...], hat.asn1.common.Entity] | None = None):
313    async def async_close(self, user_data: IdentifiedEntity | None = None):
314        """Async close"""
315        self.close(user_data)
316        await self.wait_closed()

Async close

async def receive(self) -> tuple[tuple[int, ...], hat.asn1.common.Entity]:
318    async def receive(self) -> IdentifiedEntity:
319        """Receive data"""
320        try:
321            return await self._receive_queue.get()
322
323        except aio.QueueClosedError:
324            raise ConnectionError()

Receive data

async def send(self, data: tuple[tuple[int, ...], hat.asn1.common.Entity]):
326    async def send(self, data: IdentifiedEntity):
327        """Send data"""
328        try:
329            await self._send_queue.put((data, None))
330
331        except aio.QueueClosedError:
332            raise ConnectionError()

Send data

async def drain(self):
334    async def drain(self):
335        """Drain output buffer"""
336        try:
337            future = self._loop.create_future()
338            await self._send_queue.put((None, future))
339            await future
340
341        except aio.QueueClosedError:
342            raise ConnectionError()

Drain output buffer