hat.drivers.acse
Association Controll Service Element
1"""Association Controll Service Element""" 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 copp 13from hat.drivers import net 14 15 16mlog = logging.getLogger(__name__) 17 18# (joint-iso-itu-t, association-control, abstract-syntax, apdus, version1) 19_acse_syntax_name = (2, 2, 1, 0, 1) 20 21with importlib.resources.open_text(__package__, 'asn1_repo.json') as _f: 22 _encoder = asn1.ber.BerEncoder( 23 asn1.repository_from_json( 24 json.decode_stream(_f))) 25 26 27class TcpConnectionInfo(typing.NamedTuple): 28 name: str | None 29 local_addr: net.TcpAddress 30 local_tsel: int | None 31 local_ssel: int | None 32 local_psel: int | None 33 local_ap_title: asn1.ObjectIdentifier | None 34 local_ae_qualifier: int | None 35 remote_addr: net.TcpAddress 36 remote_tsel: int | None 37 remote_ssel: int | None 38 remote_psel: int | None 39 remote_ap_title: asn1.ObjectIdentifier | None 40 remote_ae_qualifier: int | None 41 42 43class UnixConnectionInfo(typing.NamedTuple): 44 name: str | None 45 addr: net.UnixAddress 46 local_tsel: int | None 47 local_ssel: int | None 48 local_psel: int | None 49 local_ap_title: asn1.ObjectIdentifier | None 50 local_ae_qualifier: int | None 51 remote_tsel: int | None 52 remote_ssel: int | None 53 remote_psel: int | None 54 remote_ap_title: asn1.ObjectIdentifier | None 55 remote_ae_qualifier: int | None 56 57 58ConnectionInfo: typing.TypeAlias = TcpConnectionInfo | UnixConnectionInfo 59 60ValidateCb: typing.TypeAlias = aio.AsyncCallable[[copp.SyntaxNames, 61 copp.IdentifiedEntity], 62 copp.IdentifiedEntity | None] 63"""Validate callback""" 64 65ConnectionCb: typing.TypeAlias = aio.AsyncCallable[['Connection'], None] 66"""Connection callback""" 67 68 69async def connect(addr: net.StreamAddress, 70 syntax_name_list: list[asn1.ObjectIdentifier], 71 app_context_name: asn1.ObjectIdentifier, 72 user_data: copp.IdentifiedEntity | None = None, 73 *, 74 local_ap_title: asn1.ObjectIdentifier | None = None, 75 remote_ap_title: asn1.ObjectIdentifier | None = None, 76 local_ae_qualifier: int | None = None, 77 remote_ae_qualifier: int | None = None, 78 acse_receive_queue_size: int = 1024, 79 acse_send_queue_size: int = 1024, 80 **kwargs 81 ) -> 'Connection': 82 """Connect to ACSE server 83 84 Additional arguments are passed directly to `hat.drivers.copp.connect` 85 (`syntax_names` is set by this coroutine). 86 87 """ 88 syntax_names = copp.SyntaxNames([_acse_syntax_name, *syntax_name_list]) 89 aarq_apdu = _aarq_apdu(syntax_names, app_context_name, 90 local_ap_title, remote_ap_title, 91 local_ae_qualifier, remote_ae_qualifier, 92 user_data) 93 copp_user_data = _acse_syntax_name, _encode(aarq_apdu) 94 conn = await copp.connect(addr, syntax_names, copp_user_data, **kwargs) 95 96 try: 97 aare_apdu_syntax_name, aare_apdu_entity = conn.conn_res_user_data 98 if aare_apdu_syntax_name != _acse_syntax_name: 99 raise Exception("invalid syntax name") 100 101 aare_apdu = _decode(aare_apdu_entity) 102 if aare_apdu[0] != 'aare' or aare_apdu[1]['result'] != 0: 103 raise Exception("invalid apdu") 104 105 calling_ap_title, called_ap_title = _get_ap_titles(aarq_apdu) 106 calling_ae_qualifier, called_ae_qualifier = _get_ae_qualifiers( 107 aarq_apdu) 108 return Connection(conn, aarq_apdu, aare_apdu, 109 calling_ap_title, called_ap_title, 110 calling_ae_qualifier, called_ae_qualifier, 111 acse_receive_queue_size, acse_send_queue_size) 112 113 except Exception: 114 await aio.uncancellable(_close_copp(conn, _abrt_apdu(1))) 115 raise 116 117 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 acse_receive_queue_size: int = 1024, 124 acse_send_queue_size: int = 1024, 125 **kwargs 126 ) -> 'Server': 127 """Create ACSE listening server 128 129 Additional arguments are passed directly to `hat.drivers.copp.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 = acse_receive_queue_size 143 server._send_queue_size = acse_send_queue_size 144 server._log = mlog 145 146 server._srv = await copp.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 155 156 157class Server(aio.Resource): 158 """ACSE 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, syntax_names, user_data): 175 aarq_apdu_syntax_name, aarq_apdu_entity = user_data 176 if aarq_apdu_syntax_name != _acse_syntax_name: 177 raise Exception('invalid acse syntax name') 178 179 aarq_apdu = _decode(aarq_apdu_entity) 180 if aarq_apdu[0] != 'aarq': 181 raise Exception('not aarq message') 182 183 aarq_external = aarq_apdu[1]['user-information'][0] 184 if aarq_external.direct_ref is not None: 185 if aarq_external.direct_ref != _encoder.syntax_name: 186 raise Exception('invalid encoder identifier') 187 188 _, called_ap_title = _get_ap_titles(aarq_apdu) 189 _, called_ae_qualifier = _get_ae_qualifiers(aarq_apdu) 190 _, called_ap_invocation_identifier = \ 191 _get_ap_invocation_identifiers(aarq_apdu) 192 _, called_ae_invocation_identifier = \ 193 _get_ae_invocation_identifiers(aarq_apdu) 194 195 aarq_user_data = (syntax_names.get_name(aarq_external.indirect_ref), 196 aarq_external.data) 197 198 user_validate_result = await aio.call(self._validate_cb, syntax_names, 199 aarq_user_data) 200 201 aare_apdu = _aare_apdu(syntax_names, 202 user_validate_result, 203 called_ap_title, called_ae_qualifier, 204 called_ap_invocation_identifier, 205 called_ae_invocation_identifier) 206 return _acse_syntax_name, _encode(aare_apdu) 207 208 async def _on_connection(self, copp_conn): 209 try: 210 try: 211 aarq_apdu = _decode(copp_conn.conn_req_user_data[1]) 212 aare_apdu = _decode(copp_conn.conn_res_user_data[1]) 213 214 calling_ap_title, called_ap_title = _get_ap_titles(aarq_apdu) 215 calling_ae_qualifier, called_ae_qualifier = _get_ae_qualifiers( 216 aarq_apdu) 217 218 conn = Connection(copp_conn, aarq_apdu, aare_apdu, 219 called_ap_title, calling_ap_title, 220 called_ae_qualifier, calling_ae_qualifier, 221 self._receive_queue_size, 222 self._send_queue_size) 223 224 except Exception: 225 await aio.uncancellable(_close_copp(copp_conn, _abrt_apdu(1))) 226 raise 227 228 try: 229 await aio.call(self._connection_cb, conn) 230 231 except BaseException: 232 await aio.uncancellable(conn.async_close()) 233 raise 234 235 except Exception as e: 236 self._log.error("error creating new incomming connection: %s", 237 e, exc_info=e) 238 return 239 240 if not self._bind_connections: 241 return 242 243 try: 244 await conn.wait_closed() 245 246 except BaseException: 247 await aio.uncancellable(conn.async_close()) 248 raise 249 250 251class Connection(aio.Resource): 252 """ACSE connection 253 254 For creating new connection see `connect` or `listen`. 255 256 """ 257 258 def __init__(self, 259 conn: copp.Connection, 260 aarq_apdu: asn1.Value, 261 aare_apdu: asn1.Value, 262 local_ap_title: asn1.ObjectIdentifier | None, 263 remote_ap_title: asn1.ObjectIdentifier | None, 264 local_ae_qualifier: int | None, 265 remote_ae_qualifier: int | None, 266 receive_queue_size: int, 267 send_queue_size: int): 268 aarq_external = aarq_apdu[1]['user-information'][0] 269 aare_external = aare_apdu[1]['user-information'][0] 270 271 conn_req_user_data = ( 272 conn.syntax_names.get_name(aarq_external.indirect_ref), 273 aarq_external.data) 274 conn_res_user_data = ( 275 conn.syntax_names.get_name(aare_external.indirect_ref), 276 aare_external.data) 277 278 self._conn = conn 279 self._conn_req_user_data = conn_req_user_data 280 self._conn_res_user_data = conn_res_user_data 281 self._loop = asyncio.get_running_loop() 282 self._info = _connection_info_from_copp( 283 info=conn.info, 284 local_ap_title=local_ap_title, 285 local_ae_qualifier=local_ae_qualifier, 286 remote_ap_title=remote_ap_title, 287 remote_ae_qualifier=remote_ae_qualifier,) 288 self._close_apdu = _abrt_apdu(0) 289 self._receive_queue = aio.Queue(receive_queue_size) 290 self._send_queue = aio.Queue(send_queue_size) 291 self._async_group = aio.Group() 292 self._log = _create_connection_logger(self._info) 293 294 self.async_group.spawn(aio.call_on_cancel, self._on_close) 295 self.async_group.spawn(self._receive_loop) 296 self.async_group.spawn(self._send_loop) 297 self.async_group.spawn(aio.call_on_done, conn.wait_closing(), 298 self.close) 299 300 @property 301 def async_group(self) -> aio.Group: 302 """Async group""" 303 return self._async_group 304 305 @property 306 def info(self) -> ConnectionInfo: 307 """Connection info""" 308 return self._info 309 310 @property 311 def conn_req_user_data(self) -> copp.IdentifiedEntity: 312 """Connect request's user data""" 313 return self._conn_req_user_data 314 315 @property 316 def conn_res_user_data(self) -> copp.IdentifiedEntity: 317 """Connect response's user data""" 318 return self._conn_res_user_data 319 320 async def receive(self) -> copp.IdentifiedEntity: 321 """Receive data""" 322 try: 323 return await self._receive_queue.get() 324 325 except aio.QueueClosedError: 326 raise ConnectionError() 327 328 async def send(self, data: copp.IdentifiedEntity): 329 """Send data""" 330 try: 331 await self._send_queue.put((data, None)) 332 333 except aio.QueueClosedError: 334 raise ConnectionError() 335 336 async def drain(self): 337 """Drain output buffer""" 338 try: 339 future = self._loop.create_future() 340 await self._send_queue.put((None, future)) 341 await future 342 343 except aio.QueueClosedError: 344 raise ConnectionError() 345 346 async def _on_close(self): 347 await _close_copp(self._conn, self._close_apdu) 348 349 def _close(self, apdu): 350 if not self.is_open: 351 return 352 353 self._close_apdu = apdu 354 self._async_group.close() 355 356 async def _receive_loop(self): 357 try: 358 while True: 359 syntax_name, entity = await self._conn.receive() 360 361 if syntax_name == _acse_syntax_name: 362 if entity[0] == 'abrt': 363 close_apdu = None 364 365 elif entity[0] == 'rlrq': 366 close_apdu = _rlre_apdu() 367 368 else: 369 close_apdu = _abrt_apdu(1) 370 371 self._close(close_apdu) 372 break 373 374 await self._receive_queue.put((syntax_name, entity)) 375 376 except ConnectionError: 377 pass 378 379 except Exception as e: 380 self._log.error("receive loop error: %s", e, exc_info=e) 381 382 finally: 383 self._close(_abrt_apdu(1)) 384 self._receive_queue.close() 385 386 async def _send_loop(self): 387 future = None 388 try: 389 while True: 390 data, future = await self._send_queue.get() 391 392 if data is None: 393 await self._conn.drain() 394 395 else: 396 await self._conn.send(data) 397 398 if future and not future.done(): 399 future.set_result(None) 400 401 except ConnectionError: 402 pass 403 404 except Exception as e: 405 self._log.error("send loop error: %s", e, exc_info=e) 406 407 finally: 408 self._close(_abrt_apdu(1)) 409 self._send_queue.close() 410 411 while True: 412 if future and not future.done(): 413 future.set_result(None) 414 if self._send_queue.empty(): 415 break 416 _, future = self._send_queue.get_nowait() 417 418 419async def _close_copp(copp_conn, apdu): 420 data = (_acse_syntax_name, _encode(apdu)) if apdu else None 421 await copp_conn.async_close(data) 422 423 424def _get_ap_titles(aarq_apdu): 425 calling = None 426 if 'calling-AP-title' in aarq_apdu[1]: 427 if aarq_apdu[1]['calling-AP-title'][0] == 'ap-title-form2': 428 calling = aarq_apdu[1]['calling-AP-title'][1] 429 430 called = None 431 if 'called-AP-title' in aarq_apdu[1]: 432 if aarq_apdu[1]['called-AP-title'][0] == 'ap-title-form2': 433 called = aarq_apdu[1]['called-AP-title'][1] 434 435 return calling, called 436 437 438def _get_ae_qualifiers(aarq_apdu): 439 calling = None 440 if 'calling-AE-qualifier' in aarq_apdu[1]: 441 if aarq_apdu[1]['calling-AE-qualifier'][0] == 'ae-qualifier-form2': 442 calling = aarq_apdu[1]['calling-AE-qualifier'][1] 443 444 called = None 445 if 'called-AE-qualifier' in aarq_apdu[1]: 446 if aarq_apdu[1]['called-AE-qualifier'][0] == 'ae-qualifier-form2': 447 called = aarq_apdu[1]['called-AE-qualifier'][1] 448 449 return calling, called 450 451 452def _get_ap_invocation_identifiers(aarq_apdu): 453 calling = aarq_apdu[1].get('calling-AP-invocation-identifier') 454 called = aarq_apdu[1].get('called-AP-invocation-identifier') 455 return calling, called 456 457 458def _get_ae_invocation_identifiers(aarq_apdu): 459 calling = aarq_apdu[1].get('calling-AE-invocation-identifier') 460 called = aarq_apdu[1].get('called-AE-invocation-identifier') 461 return calling, called 462 463 464def _aarq_apdu(syntax_names, app_context_name, 465 calling_ap_title, called_ap_title, 466 calling_ae_qualifier, called_ae_qualifier, 467 user_data): 468 aarq_apdu = 'aarq', {'application-context-name': app_context_name} 469 470 if calling_ap_title is not None: 471 aarq_apdu[1]['calling-AP-title'] = 'ap-title-form2', calling_ap_title 472 473 if called_ap_title is not None: 474 aarq_apdu[1]['called-AP-title'] = 'ap-title-form2', called_ap_title 475 476 if calling_ae_qualifier is not None: 477 aarq_apdu[1]['calling-AE-qualifier'] = ('ae-qualifier-form2', 478 calling_ae_qualifier) 479 480 if called_ae_qualifier is not None: 481 aarq_apdu[1]['called-AE-qualifier'] = ('ae-qualifier-form2', 482 called_ae_qualifier) 483 484 if user_data: 485 aarq_apdu[1]['user-information'] = [ 486 asn1.External(direct_ref=_encoder.syntax_name, 487 indirect_ref=syntax_names.get_id(user_data[0]), 488 data=user_data[1])] 489 490 return aarq_apdu 491 492 493def _aare_apdu(syntax_names, user_data, 494 responding_ap_title, responding_ae_qualifier, 495 responding_ap_invocation_identifier, 496 responding_ae_invocation_identifier): 497 aare_apdu = 'aare', { 498 'application-context-name': user_data[0], 499 'result': 0, 500 'result-source-diagnostic': ('acse-service-user', 0), 501 'user-information': [ 502 asn1.External(direct_ref=_encoder.syntax_name, 503 indirect_ref=syntax_names.get_id(user_data[0]), 504 data=user_data[1])]} 505 506 if responding_ap_title is not None: 507 aare_apdu[1]['responding-AP-title'] = ('ap-title-form2', 508 responding_ap_title) 509 510 if responding_ae_qualifier is not None: 511 aare_apdu[1]['responding-AE-qualifier'] = ('ae-qualifier-form2', 512 responding_ae_qualifier) 513 514 if responding_ap_invocation_identifier is not None: 515 aare_apdu[1]['responding-AP-invocation-identifier'] = \ 516 responding_ap_invocation_identifier 517 518 if responding_ae_invocation_identifier is not None: 519 aare_apdu[1]['responding-AE-invocation-identifier'] = \ 520 responding_ae_invocation_identifier 521 522 return aare_apdu 523 524 525def _abrt_apdu(source): 526 return 'abrt', {'abort-source': source} 527 528 529def _rlre_apdu(): 530 return 'rlre', {} 531 532 533def _encode(value): 534 return _encoder.encode_value(asn1.TypeRef('ACSE-1', 'ACSE-apdu'), value) 535 536 537def _decode(entity): 538 return _encoder.decode_value(asn1.TypeRef('ACSE-1', 'ACSE-apdu'), entity) 539 540 541def _connection_info_to_net(info): 542 if isinstance(info, TcpConnectionInfo): 543 return net.TcpConnectionInfo(name=info.name, 544 local_addr=info.local_addr, 545 remote_addr=info.remote_addr) 546 547 if isinstance(info, UnixConnectionInfo): 548 return net.UnixConnectionInfo(name=info.name, 549 addr=info.addr) 550 551 raise TypeError('unsupported info type') 552 553 554def _connection_info_from_copp(info, local_ap_title, local_ae_qualifier, 555 remote_ap_title, remote_ae_qualifier,): 556 if isinstance(info, copp.TcpConnectionInfo): 557 return TcpConnectionInfo(name=info.name, 558 local_addr=info.local_addr, 559 local_tsel=info.local_tsel, 560 local_ssel=info.local_ssel, 561 local_psel=info.local_psel, 562 local_ap_title=local_ap_title, 563 local_ae_qualifier=local_ae_qualifier, 564 remote_addr=info.remote_addr, 565 remote_tsel=info.remote_tsel, 566 remote_ssel=info.remote_ssel, 567 remote_psel=info.remote_psel, 568 remote_ap_title=remote_ap_title, 569 remote_ae_qualifier=remote_ae_qualifier) 570 571 if isinstance(info, copp.UnixConnectionInfo): 572 return UnixConnectionInfo(name=info.name, 573 addr=info.addr, 574 local_tsel=info.local_tsel, 575 local_ssel=info.local_ssel, 576 local_psel=info.local_psel, 577 local_ap_title=local_ap_title, 578 local_ae_qualifier=local_ae_qualifier, 579 remote_tsel=info.remote_tsel, 580 remote_ssel=info.remote_ssel, 581 remote_psel=info.remote_psel, 582 remote_ap_title=remote_ap_title, 583 remote_ae_qualifier=remote_ae_qualifier) 584 585 raise TypeError('unsupported info type') 586 587 588def _create_server_logger(info): 589 extra = {'meta': {'type': 'AcseServer', 590 **net.server_info_to_json(info)}} 591 592 return logging.LoggerAdapter(mlog, extra) 593 594 595def _create_connection_logger(info): 596 net_info = _connection_info_to_net(info) 597 598 extra = {'meta': {'type': 'AcseConnection', 599 **net.connection_info_to_json(net_info)}} 600 601 return logging.LoggerAdapter(mlog, extra)
28class TcpConnectionInfo(typing.NamedTuple): 29 name: str | None 30 local_addr: net.TcpAddress 31 local_tsel: int | None 32 local_ssel: int | None 33 local_psel: int | None 34 local_ap_title: asn1.ObjectIdentifier | None 35 local_ae_qualifier: int | None 36 remote_addr: net.TcpAddress 37 remote_tsel: int | None 38 remote_ssel: int | None 39 remote_psel: int | None 40 remote_ap_title: asn1.ObjectIdentifier | None 41 remote_ae_qualifier: int | None
TcpConnectionInfo(name, local_addr, local_tsel, local_ssel, local_psel, local_ap_title, local_ae_qualifier, remote_addr, remote_tsel, remote_ssel, remote_psel, remote_ap_title, remote_ae_qualifier)
Create new instance of TcpConnectionInfo(name, local_addr, local_tsel, local_ssel, local_psel, local_ap_title, local_ae_qualifier, remote_addr, remote_tsel, remote_ssel, remote_psel, remote_ap_title, remote_ae_qualifier)
44class UnixConnectionInfo(typing.NamedTuple): 45 name: str | None 46 addr: net.UnixAddress 47 local_tsel: int | None 48 local_ssel: int | None 49 local_psel: int | None 50 local_ap_title: asn1.ObjectIdentifier | None 51 local_ae_qualifier: int | None 52 remote_tsel: int | None 53 remote_ssel: int | None 54 remote_psel: int | None 55 remote_ap_title: asn1.ObjectIdentifier | None 56 remote_ae_qualifier: int | None
UnixConnectionInfo(name, addr, local_tsel, local_ssel, local_psel, local_ap_title, local_ae_qualifier, remote_tsel, remote_ssel, remote_psel, remote_ap_title, remote_ae_qualifier)
Create new instance of UnixConnectionInfo(name, addr, local_tsel, local_ssel, local_psel, local_ap_title, local_ae_qualifier, remote_tsel, remote_ssel, remote_psel, remote_ap_title, remote_ae_qualifier)
Validate callback
Connection callback
70async def connect(addr: net.StreamAddress, 71 syntax_name_list: list[asn1.ObjectIdentifier], 72 app_context_name: asn1.ObjectIdentifier, 73 user_data: copp.IdentifiedEntity | None = None, 74 *, 75 local_ap_title: asn1.ObjectIdentifier | None = None, 76 remote_ap_title: asn1.ObjectIdentifier | None = None, 77 local_ae_qualifier: int | None = None, 78 remote_ae_qualifier: int | None = None, 79 acse_receive_queue_size: int = 1024, 80 acse_send_queue_size: int = 1024, 81 **kwargs 82 ) -> 'Connection': 83 """Connect to ACSE server 84 85 Additional arguments are passed directly to `hat.drivers.copp.connect` 86 (`syntax_names` is set by this coroutine). 87 88 """ 89 syntax_names = copp.SyntaxNames([_acse_syntax_name, *syntax_name_list]) 90 aarq_apdu = _aarq_apdu(syntax_names, app_context_name, 91 local_ap_title, remote_ap_title, 92 local_ae_qualifier, remote_ae_qualifier, 93 user_data) 94 copp_user_data = _acse_syntax_name, _encode(aarq_apdu) 95 conn = await copp.connect(addr, syntax_names, copp_user_data, **kwargs) 96 97 try: 98 aare_apdu_syntax_name, aare_apdu_entity = conn.conn_res_user_data 99 if aare_apdu_syntax_name != _acse_syntax_name: 100 raise Exception("invalid syntax name") 101 102 aare_apdu = _decode(aare_apdu_entity) 103 if aare_apdu[0] != 'aare' or aare_apdu[1]['result'] != 0: 104 raise Exception("invalid apdu") 105 106 calling_ap_title, called_ap_title = _get_ap_titles(aarq_apdu) 107 calling_ae_qualifier, called_ae_qualifier = _get_ae_qualifiers( 108 aarq_apdu) 109 return Connection(conn, aarq_apdu, aare_apdu, 110 calling_ap_title, called_ap_title, 111 calling_ae_qualifier, called_ae_qualifier, 112 acse_receive_queue_size, acse_send_queue_size) 113 114 except Exception: 115 await aio.uncancellable(_close_copp(conn, _abrt_apdu(1))) 116 raise
Connect to ACSE server
Additional arguments are passed directly to hat.drivers.copp.connect
(syntax_names is set by this coroutine).
119async def listen(validate_cb: ValidateCb, 120 connection_cb: ConnectionCb, 121 addr: net.StreamAddress = net.TcpAddress('0.0.0.0', 102), 122 *, 123 bind_connections: bool = False, 124 acse_receive_queue_size: int = 1024, 125 acse_send_queue_size: int = 1024, 126 **kwargs 127 ) -> 'Server': 128 """Create ACSE listening server 129 130 Additional arguments are passed directly to `hat.drivers.copp.listen`. 131 132 Args: 133 validate_cb: callback function or coroutine called on new 134 incomming connection request prior to creating connection object 135 connection_cb: new connection callback 136 addr: local listening address 137 138 """ 139 server = Server() 140 server._validate_cb = validate_cb 141 server._connection_cb = connection_cb 142 server._bind_connections = bind_connections 143 server._receive_queue_size = acse_receive_queue_size 144 server._send_queue_size = acse_send_queue_size 145 server._log = mlog 146 147 server._srv = await copp.listen(server._on_validate, 148 server._on_connection, 149 addr, 150 bind_connections=False, 151 **kwargs) 152 153 server._log = _create_server_logger(server._srv.info) 154 155 return server
Create ACSE listening server
Additional arguments are passed directly to hat.drivers.copp.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
158class Server(aio.Resource): 159 """ACSE listening server 160 161 For creating new server see `listen`. 162 163 """ 164 165 @property 166 def async_group(self) -> aio.Group: 167 """Async group""" 168 return self._srv.async_group 169 170 @property 171 def info(self) -> net.ServerInfo: 172 """Server info""" 173 return self._srv.info 174 175 async def _on_validate(self, syntax_names, user_data): 176 aarq_apdu_syntax_name, aarq_apdu_entity = user_data 177 if aarq_apdu_syntax_name != _acse_syntax_name: 178 raise Exception('invalid acse syntax name') 179 180 aarq_apdu = _decode(aarq_apdu_entity) 181 if aarq_apdu[0] != 'aarq': 182 raise Exception('not aarq message') 183 184 aarq_external = aarq_apdu[1]['user-information'][0] 185 if aarq_external.direct_ref is not None: 186 if aarq_external.direct_ref != _encoder.syntax_name: 187 raise Exception('invalid encoder identifier') 188 189 _, called_ap_title = _get_ap_titles(aarq_apdu) 190 _, called_ae_qualifier = _get_ae_qualifiers(aarq_apdu) 191 _, called_ap_invocation_identifier = \ 192 _get_ap_invocation_identifiers(aarq_apdu) 193 _, called_ae_invocation_identifier = \ 194 _get_ae_invocation_identifiers(aarq_apdu) 195 196 aarq_user_data = (syntax_names.get_name(aarq_external.indirect_ref), 197 aarq_external.data) 198 199 user_validate_result = await aio.call(self._validate_cb, syntax_names, 200 aarq_user_data) 201 202 aare_apdu = _aare_apdu(syntax_names, 203 user_validate_result, 204 called_ap_title, called_ae_qualifier, 205 called_ap_invocation_identifier, 206 called_ae_invocation_identifier) 207 return _acse_syntax_name, _encode(aare_apdu) 208 209 async def _on_connection(self, copp_conn): 210 try: 211 try: 212 aarq_apdu = _decode(copp_conn.conn_req_user_data[1]) 213 aare_apdu = _decode(copp_conn.conn_res_user_data[1]) 214 215 calling_ap_title, called_ap_title = _get_ap_titles(aarq_apdu) 216 calling_ae_qualifier, called_ae_qualifier = _get_ae_qualifiers( 217 aarq_apdu) 218 219 conn = Connection(copp_conn, aarq_apdu, aare_apdu, 220 called_ap_title, calling_ap_title, 221 called_ae_qualifier, calling_ae_qualifier, 222 self._receive_queue_size, 223 self._send_queue_size) 224 225 except Exception: 226 await aio.uncancellable(_close_copp(copp_conn, _abrt_apdu(1))) 227 raise 228 229 try: 230 await aio.call(self._connection_cb, conn) 231 232 except BaseException: 233 await aio.uncancellable(conn.async_close()) 234 raise 235 236 except Exception as e: 237 self._log.error("error creating new incomming connection: %s", 238 e, exc_info=e) 239 return 240 241 if not self._bind_connections: 242 return 243 244 try: 245 await conn.wait_closed() 246 247 except BaseException: 248 await aio.uncancellable(conn.async_close()) 249 raise
ACSE listening server
For creating new server see listen.
165 @property 166 def async_group(self) -> aio.Group: 167 """Async group""" 168 return self._srv.async_group
Async group
252class Connection(aio.Resource): 253 """ACSE connection 254 255 For creating new connection see `connect` or `listen`. 256 257 """ 258 259 def __init__(self, 260 conn: copp.Connection, 261 aarq_apdu: asn1.Value, 262 aare_apdu: asn1.Value, 263 local_ap_title: asn1.ObjectIdentifier | None, 264 remote_ap_title: asn1.ObjectIdentifier | None, 265 local_ae_qualifier: int | None, 266 remote_ae_qualifier: int | None, 267 receive_queue_size: int, 268 send_queue_size: int): 269 aarq_external = aarq_apdu[1]['user-information'][0] 270 aare_external = aare_apdu[1]['user-information'][0] 271 272 conn_req_user_data = ( 273 conn.syntax_names.get_name(aarq_external.indirect_ref), 274 aarq_external.data) 275 conn_res_user_data = ( 276 conn.syntax_names.get_name(aare_external.indirect_ref), 277 aare_external.data) 278 279 self._conn = conn 280 self._conn_req_user_data = conn_req_user_data 281 self._conn_res_user_data = conn_res_user_data 282 self._loop = asyncio.get_running_loop() 283 self._info = _connection_info_from_copp( 284 info=conn.info, 285 local_ap_title=local_ap_title, 286 local_ae_qualifier=local_ae_qualifier, 287 remote_ap_title=remote_ap_title, 288 remote_ae_qualifier=remote_ae_qualifier,) 289 self._close_apdu = _abrt_apdu(0) 290 self._receive_queue = aio.Queue(receive_queue_size) 291 self._send_queue = aio.Queue(send_queue_size) 292 self._async_group = aio.Group() 293 self._log = _create_connection_logger(self._info) 294 295 self.async_group.spawn(aio.call_on_cancel, self._on_close) 296 self.async_group.spawn(self._receive_loop) 297 self.async_group.spawn(self._send_loop) 298 self.async_group.spawn(aio.call_on_done, conn.wait_closing(), 299 self.close) 300 301 @property 302 def async_group(self) -> aio.Group: 303 """Async group""" 304 return self._async_group 305 306 @property 307 def info(self) -> ConnectionInfo: 308 """Connection info""" 309 return self._info 310 311 @property 312 def conn_req_user_data(self) -> copp.IdentifiedEntity: 313 """Connect request's user data""" 314 return self._conn_req_user_data 315 316 @property 317 def conn_res_user_data(self) -> copp.IdentifiedEntity: 318 """Connect response's user data""" 319 return self._conn_res_user_data 320 321 async def receive(self) -> copp.IdentifiedEntity: 322 """Receive data""" 323 try: 324 return await self._receive_queue.get() 325 326 except aio.QueueClosedError: 327 raise ConnectionError() 328 329 async def send(self, data: copp.IdentifiedEntity): 330 """Send data""" 331 try: 332 await self._send_queue.put((data, None)) 333 334 except aio.QueueClosedError: 335 raise ConnectionError() 336 337 async def drain(self): 338 """Drain output buffer""" 339 try: 340 future = self._loop.create_future() 341 await self._send_queue.put((None, future)) 342 await future 343 344 except aio.QueueClosedError: 345 raise ConnectionError() 346 347 async def _on_close(self): 348 await _close_copp(self._conn, self._close_apdu) 349 350 def _close(self, apdu): 351 if not self.is_open: 352 return 353 354 self._close_apdu = apdu 355 self._async_group.close() 356 357 async def _receive_loop(self): 358 try: 359 while True: 360 syntax_name, entity = await self._conn.receive() 361 362 if syntax_name == _acse_syntax_name: 363 if entity[0] == 'abrt': 364 close_apdu = None 365 366 elif entity[0] == 'rlrq': 367 close_apdu = _rlre_apdu() 368 369 else: 370 close_apdu = _abrt_apdu(1) 371 372 self._close(close_apdu) 373 break 374 375 await self._receive_queue.put((syntax_name, entity)) 376 377 except ConnectionError: 378 pass 379 380 except Exception as e: 381 self._log.error("receive loop error: %s", e, exc_info=e) 382 383 finally: 384 self._close(_abrt_apdu(1)) 385 self._receive_queue.close() 386 387 async def _send_loop(self): 388 future = None 389 try: 390 while True: 391 data, future = await self._send_queue.get() 392 393 if data is None: 394 await self._conn.drain() 395 396 else: 397 await self._conn.send(data) 398 399 if future and not future.done(): 400 future.set_result(None) 401 402 except ConnectionError: 403 pass 404 405 except Exception as e: 406 self._log.error("send loop error: %s", e, exc_info=e) 407 408 finally: 409 self._close(_abrt_apdu(1)) 410 self._send_queue.close() 411 412 while True: 413 if future and not future.done(): 414 future.set_result(None) 415 if self._send_queue.empty(): 416 break 417 _, future = self._send_queue.get_nowait()
259 def __init__(self, 260 conn: copp.Connection, 261 aarq_apdu: asn1.Value, 262 aare_apdu: asn1.Value, 263 local_ap_title: asn1.ObjectIdentifier | None, 264 remote_ap_title: asn1.ObjectIdentifier | None, 265 local_ae_qualifier: int | None, 266 remote_ae_qualifier: int | None, 267 receive_queue_size: int, 268 send_queue_size: int): 269 aarq_external = aarq_apdu[1]['user-information'][0] 270 aare_external = aare_apdu[1]['user-information'][0] 271 272 conn_req_user_data = ( 273 conn.syntax_names.get_name(aarq_external.indirect_ref), 274 aarq_external.data) 275 conn_res_user_data = ( 276 conn.syntax_names.get_name(aare_external.indirect_ref), 277 aare_external.data) 278 279 self._conn = conn 280 self._conn_req_user_data = conn_req_user_data 281 self._conn_res_user_data = conn_res_user_data 282 self._loop = asyncio.get_running_loop() 283 self._info = _connection_info_from_copp( 284 info=conn.info, 285 local_ap_title=local_ap_title, 286 local_ae_qualifier=local_ae_qualifier, 287 remote_ap_title=remote_ap_title, 288 remote_ae_qualifier=remote_ae_qualifier,) 289 self._close_apdu = _abrt_apdu(0) 290 self._receive_queue = aio.Queue(receive_queue_size) 291 self._send_queue = aio.Queue(send_queue_size) 292 self._async_group = aio.Group() 293 self._log = _create_connection_logger(self._info) 294 295 self.async_group.spawn(aio.call_on_cancel, self._on_close) 296 self.async_group.spawn(self._receive_loop) 297 self.async_group.spawn(self._send_loop) 298 self.async_group.spawn(aio.call_on_done, conn.wait_closing(), 299 self.close)
301 @property 302 def async_group(self) -> aio.Group: 303 """Async group""" 304 return self._async_group
Async group
311 @property 312 def conn_req_user_data(self) -> copp.IdentifiedEntity: 313 """Connect request's user data""" 314 return self._conn_req_user_data
Connect request's user data
316 @property 317 def conn_res_user_data(self) -> copp.IdentifiedEntity: 318 """Connect response's user data""" 319 return self._conn_res_user_data
Connect response's user data
321 async def receive(self) -> copp.IdentifiedEntity: 322 """Receive data""" 323 try: 324 return await self._receive_queue.get() 325 326 except aio.QueueClosedError: 327 raise ConnectionError()
Receive data