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)
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)
Create new instance of TcpConnectionInfo(name, local_addr, local_tsel, local_ssel, local_psel, remote_addr, remote_tsel, remote_ssel, remote_psel)
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)
Create new instance of UnixConnectionInfo(name, addr, local_tsel, local_ssel, local_psel, remote_tsel, remote_ssel, remote_psel)
Identified entity
Validate callback
Connection callback
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
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.
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
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.
164 @property 165 def async_group(self) -> aio.Group: 166 """Async group""" 167 return self._srv.async_group
Async group
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()
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)
284 @property 285 def async_group(self) -> aio.Group: 286 """Async group""" 287 return self._async_group
Async group
294 @property 295 def syntax_names(self) -> SyntaxNames: 296 """Syntax names""" 297 return self._syntax_names
Syntax names
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
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
309 def close(self, user_data: IdentifiedEntity | None = None): 310 """Close connection""" 311 self._close(_aru_ppdu(self._syntax_names, user_data))
Close connection
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
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