hat.drivers.snmp
1from hat.drivers.snmp.common import (Version, 2 ErrorType, 3 CauseType, 4 AuthType, 5 PrivType, 6 Error, 7 Cause, 8 IntegerData, 9 UnsignedData, 10 CounterData, 11 BigCounterData, 12 StringData, 13 ObjectIdData, 14 IpAddressData, 15 TimeTicksData, 16 ArbitraryData, 17 EmptyData, 18 UnspecifiedData, 19 NoSuchObjectData, 20 NoSuchInstanceData, 21 EndOfMibViewData, 22 Data, 23 CommunityName, 24 UserName, 25 Password, 26 EngineId, 27 User, 28 Context, 29 Trap, 30 Inform, 31 GetDataReq, 32 GetNextDataReq, 33 GetBulkDataReq, 34 SetDataReq, 35 Request, 36 Response) 37from hat.drivers.snmp.agent import (V1RequestCb, 38 V2CRequestCb, 39 V3RequestCb, 40 create_agent, 41 Agent) 42from hat.drivers.snmp.manager import (Manager, 43 create_v1_manager, 44 create_v2c_manager, 45 create_v3_manager) 46from hat.drivers.snmp.trap import (V1TrapCb, 47 V2CTrapCb, 48 V2CInformCb, 49 V3TrapCb, 50 V3InformCb, 51 create_trap_listener, 52 TrapListener, 53 TrapSender, 54 create_v1_trap_sender, 55 create_v2c_trap_sender, 56 create_v3_trap_sender) 57 58 59__all__ = ['Version', 60 'ErrorType', 61 'CauseType', 62 'AuthType', 63 'PrivType', 64 'Error', 65 'Cause', 66 'IntegerData', 67 'UnsignedData', 68 'CounterData', 69 'BigCounterData', 70 'StringData', 71 'ObjectIdData', 72 'IpAddressData', 73 'TimeTicksData', 74 'ArbitraryData', 75 'EmptyData', 76 'UnspecifiedData', 77 'NoSuchObjectData', 78 'NoSuchInstanceData', 79 'EndOfMibViewData', 80 'Data', 81 'CommunityName', 82 'UserName', 83 'Password', 84 'EngineId', 85 'User', 86 'Context', 87 'Trap', 88 'Inform', 89 'GetDataReq', 90 'GetNextDataReq', 91 'GetBulkDataReq', 92 'SetDataReq', 93 'Request', 94 'Response', 95 'V1RequestCb', 96 'V2CRequestCb', 97 'V3RequestCb', 98 'create_agent', 99 'Agent', 100 'Manager', 101 'create_v1_manager', 102 'create_v2c_manager', 103 'create_v3_manager', 104 'V1TrapCb', 105 'V2CTrapCb', 106 'V2CInformCb', 107 'V3TrapCb', 108 'V3InformCb', 109 'create_trap_listener', 110 'TrapListener', 111 'TrapSender', 112 'create_v1_trap_sender', 113 'create_v2c_trap_sender', 114 'create_v3_trap_sender']
18class ErrorType(enum.Enum): 19 NO_ERROR = 0 # v1, v2c, v3 20 TOO_BIG = 1 # v1, v2c, v3 21 NO_SUCH_NAME = 2 # v1, v2c, v3 22 BAD_VALUE = 3 # v1, v2c, v3 23 READ_ONLY = 4 # v1, v2c, v3 24 GEN_ERR = 5 # v1, v2c, v3 25 NO_ACCESS = 6 # v2c, v3 26 WRONG_TYPE = 7 # v2c, v3 27 WRONG_LENGTH = 8 # v2c, v3 28 WRONG_ENCODING = 9 # v2c, v3 29 WRONG_VALUE = 10 # v2c, v3 30 NO_CREATION = 11 # v2c, v3 31 INCONSISTENT_VALUE = 12 # v2c, v3 32 RESOURCE_UNAVAILABLE = 13 # v2c, v3 33 COMMIT_FAILED = 14 # v2c, v3 34 UNDO_FAILED = 15 # v2c, v3 35 AUTHORIZATION_ERROR = 16 # v2c, v3 36 NOT_WRITABLE = 17 # v2c, v3 37 INCONSISTENT_NAME = 18 # v2c, v3
40class CauseType(enum.Enum): 41 COLD_START = 0 42 WARM_START = 1 43 LINK_DOWN = 2 44 LINK_UP = 3 45 AUTHENICATION_FAILURE = 4 46 EGP_NEIGHBOR_LOSS = 5 47 ENTERPRISE_SPECIFIC = 6
Error(type, index)
Cause(type, value)
IntegerData(name, value)
UnsignedData(name, value)
CounterData(name, value)
BigCounterData(name, value)
StringData(name, value)
100class ObjectIdData(typing.NamedTuple): 101 name: asn1.ObjectIdentifier 102 value: asn1.ObjectIdentifier
ObjectIdData(name, value)
106class IpAddressData(typing.NamedTuple): 107 name: asn1.ObjectIdentifier 108 value: tuple[int, int, int, int]
IpAddressData(name, value)
TimeTicksData(name, value)
ArbitraryData(name, value)
EmptyData(name,)
UnspecifiedData(name,)
NoSuchObjectData(name,)
NoSuchInstanceData(name,)
EndOfMibViewData(name,)
173class User(typing.NamedTuple): 174 name: UserName 175 auth_type: AuthType | None 176 auth_password: Password | None 177 priv_type: PrivType | None 178 priv_password: Password | None
User(name, auth_type, auth_password, priv_type, priv_password)
Context(engine_id, name)
186class Trap(typing.NamedTuple): 187 cause: Cause | None 188 """cause is available in case of v1""" 189 oid: asn1.ObjectIdentifier 190 timestamp: int 191 data: Collection[Data]
Trap(cause, oid, timestamp, data)
Create new instance of Trap(cause, oid, timestamp, data)
Alias for field number 3
Inform(data,)
Create new instance of Inform(data,)
Alias for field number 0
GetDataReq(names,)
GetNextDataReq(names,)
GetBulkDataReq(names,)
SetDataReq(data,)
Create new instance of SetDataReq(data,)
Alias for field number 0
33async def create_agent(local_addr: net.DatagramAddress = net.UdpAddress('0.0.0.0', 161), # NOQA 34 *, 35 v1_request_cb: V1RequestCb | None = None, 36 v2c_request_cb: V2CRequestCb | None = None, 37 v3_request_cb: V3RequestCb | None = None, 38 authoritative_engine_id: common.EngineId | None = None, 39 users: Collection[common.User] = [], 40 **kwargs 41 ) -> 'Agent': 42 """Create agent""" 43 if isinstance(local_addr, net.UdpAddress): 44 datagram_type = net.DatagramType.UDP 45 46 elif isinstance(local_addr, net.UnixAddress): 47 datagram_type = net.DatagramType.UNIX 48 49 else: 50 raise TypeError('unsupported address type') 51 52 endpoint = await net.create_endpoint(datagram_type=datagram_type, 53 local_addr=local_addr, 54 remote_addr=None, 55 **kwargs) 56 57 try: 58 return Agent(endpoint=endpoint, 59 v1_request_cb=v1_request_cb, 60 v2c_request_cb=v2c_request_cb, 61 v3_request_cb=v3_request_cb, 62 authoritative_engine_id=authoritative_engine_id, 63 users=users) 64 65 except BaseException: 66 await aio.uncancellable(endpoint.async_close()) 67 raise
Create agent
70class Agent(aio.Resource): 71 72 def __init__(self, 73 endpoint: net.Endpoint, 74 v1_request_cb: V1RequestCb | None, 75 v2c_request_cb: V2CRequestCb | None, 76 v3_request_cb: V3RequestCb | None, 77 authoritative_engine_id: common.EngineId | None, 78 users: Collection[common.User]): 79 self._endpoint = endpoint 80 self._v1_request_cb = v1_request_cb 81 self._v2c_request_cb = v2c_request_cb 82 self._v3_request_cb = v3_request_cb 83 self._auth_engine_id = authoritative_engine_id 84 self._auth_keys = {} 85 self._priv_keys = {} 86 self._log = logger.create_logger(mlog, 'SnmpAgent', endpoint.info) 87 self._comm_log = logger.CommunicationLogger(mlog, 'SnmpAgent', 88 endpoint.info) 89 90 for user in users: 91 common.validate_user(user) 92 93 if user.auth_type: 94 key_type = key.auth_type_to_key_type(user.auth_type) 95 self._auth_keys[user.name] = key.create_key( 96 key_type=key_type, 97 password=user.auth_password, 98 engine_id=authoritative_engine_id) 99 100 else: 101 self._auth_keys[user.name] = None 102 103 if user.priv_type: 104 key_type = key.priv_type_to_key_type(user.priv_type) 105 self._priv_keys[user.name] = key.create_key( 106 key_type=key_type, 107 password=user.priv_password, 108 engine_id=authoritative_engine_id) 109 110 else: 111 self._priv_keys[user.name] = None 112 113 self.async_group.spawn(self._receive_loop) 114 115 self.async_group.spawn(aio.call_on_cancel, self._comm_log.log, 116 common.CommLogAction.CLOSE) 117 self._comm_log.log(common.CommLogAction.OPEN) 118 119 @property 120 def async_group(self) -> aio.Group: 121 return self._endpoint.async_group 122 123 def _on_auth_key(self, engine_id, username): 124 if engine_id != self._auth_engine_id: 125 raise Exception('invalid authoritative engine id') 126 127 if username not in self._auth_keys: 128 raise Exception('invalid user') 129 return self._auth_keys[username] 130 131 def _on_priv_key(self, engine_id, username): 132 if engine_id != self._auth_engine_id: 133 raise Exception('invalid authoritative engine id') 134 135 if username not in self._priv_keys: 136 raise Exception('invalid user') 137 return self._priv_keys[username] 138 139 async def _receive_loop(self): 140 try: 141 while True: 142 req_msg_bytes, addr = await self._endpoint.receive() 143 144 try: 145 req_msg = encoder.decode(msg_bytes=req_msg_bytes, 146 auth_key_cb=self._on_auth_key, 147 priv_key_cb=self._on_priv_key) 148 149 except Exception as e: 150 self._log.warning("error decoding message from %s: %s", 151 addr, e, exc_info=e) 152 continue 153 154 self._comm_log.log(common.CommLogAction.RECEIVE, req_msg) 155 156 try: 157 if isinstance(req_msg, encoder.v1.Msg): 158 res_msg = await self._process_v1_req_msg( 159 req_msg=req_msg, 160 addr=addr) 161 162 elif isinstance(req_msg, encoder.v2c.Msg): 163 res_msg = await self._process_v2c_req_msg( 164 req_msg=req_msg, 165 addr=addr) 166 167 elif isinstance(req_msg, encoder.v3.Msg): 168 res_msg = await self._process_v3_req_msg( 169 req_msg=req_msg, 170 addr=addr) 171 172 else: 173 raise ValueError('unsupported message type') 174 175 except Exception as e: 176 self._log.warning("error processing message from %s: %s", 177 addr, e, exc_info=e) 178 continue 179 180 if not res_msg: 181 continue 182 183 try: 184 if isinstance(res_msg, encoder.v3.Msg): 185 auth_key = ( 186 self._on_auth_key(res_msg.authorative_engine.id, 187 res_msg.user) 188 if res_msg.auth else None) 189 priv_key = ( 190 self._on_priv_key(res_msg.authorative_engine.id, 191 res_msg.user) 192 if res_msg.priv else None) 193 194 else: 195 auth_key = None 196 priv_key = None 197 198 res_msg_bytes = encoder.encode(msg=res_msg, 199 auth_key=auth_key, 200 priv_key=priv_key) 201 202 except Exception as e: 203 self._log.warning("error encoding message: %s", 204 e, exc_info=e) 205 continue 206 207 self._comm_log.log(common.CommLogAction.SEND, res_msg) 208 209 self._endpoint.send(res_msg_bytes, addr) 210 211 except ConnectionError: 212 pass 213 214 except Exception as e: 215 self._log.error("receive loop error: %s", e, exc_info=e) 216 217 finally: 218 self.close() 219 220 async def _process_v1_req_msg(self, req_msg, addr): 221 if not self._v1_request_cb: 222 raise Exception('not accepting V1') 223 224 if req_msg.type == encoder.v1.MsgType.GET_REQUEST: 225 req = common.GetDataReq(names=[i.name for i in req_msg.pdu.data]) 226 227 elif req_msg.type == encoder.v1.MsgType.GET_NEXT_REQUEST: 228 req = common.GetNextDataReq( 229 names=[i.name for i in req_msg.pdu.data]) 230 231 elif req_msg.type == encoder.v1.MsgType.SET_REQUEST: 232 req = common.SetDataReq(data=req_msg.pdu.data) 233 234 else: 235 raise Exception('invalid request message type') 236 237 try: 238 res = await aio.call(self._v1_request_cb, addr, req_msg.community, 239 req) 240 241 if isinstance(res, common.Error): 242 if res.type.value > common.ErrorType.GEN_ERR.value: 243 raise Exception('invalid error type') 244 245 res_error = res 246 res_data = [] 247 248 else: 249 res_error = common.Error(common.ErrorType.NO_ERROR, 0) 250 res_data = res 251 252 except Exception as e: 253 self._log.warning("error processing request: %s", e, exc_info=e) 254 255 res_error = common.Error(common.ErrorType.GEN_ERR, 0) 256 res_data = [] 257 258 res_pdu = encoder.v1.BasicPdu( 259 request_id=req_msg.pdu.request_id, 260 error=res_error, 261 data=res_data) 262 263 res_msg = encoder.v1.Msg( 264 type=encoder.v1.MsgType.GET_RESPONSE, 265 community=req_msg.community, 266 pdu=res_pdu) 267 268 return res_msg 269 270 async def _process_v2c_req_msg(self, req_msg, addr): 271 if not self._v2c_request_cb: 272 raise Exception('not accepting V2C') 273 274 if req_msg.type == encoder.v2c.MsgType.GET_REQUEST: 275 req = common.GetDataReq(names=[i.name for i in req_msg.pdu.data]) 276 277 elif req_msg.type == encoder.v2c.MsgType.GET_NEXT_REQUEST: 278 req = common.GetNextDataReq( 279 names=[i.name for i in req_msg.pdu.data]) 280 281 elif req_msg.type == encoder.v2c.MsgType.GET_BULK_REQUEST: 282 req = common.GetBulkDataReq( 283 names=[i.name for i in req_msg.pdu.data]) 284 285 elif req_msg.type == encoder.v2c.MsgType.SET_REQUEST: 286 req = common.SetDataReq(data=req_msg.pdu.data) 287 288 else: 289 raise Exception('invalid request message type') 290 291 try: 292 res = await aio.call(self._v2c_request_cb, addr, req_msg.community, 293 req) 294 295 if isinstance(res, common.Error): 296 res_error = res 297 res_data = [] 298 299 else: 300 res_error = common.Error(common.ErrorType.NO_ERROR, 0) 301 res_data = res 302 303 except Exception as e: 304 self._log.warning("error processing request: %s", e, exc_info=e) 305 306 res_error = common.Error(common.ErrorType.GEN_ERR, 0) 307 res_data = [] 308 309 res_pdu = encoder.v2c.BasicPdu( 310 request_id=req_msg.pdu.request_id, 311 error=res_error, 312 data=res_data) 313 314 res_msg = encoder.v2c.Msg( 315 type=encoder.v2c.MsgType.RESPONSE, 316 community=req_msg.community, 317 pdu=res_pdu) 318 319 return res_msg 320 321 async def _process_v3_req_msg(self, req_msg, addr): 322 if not self._v3_request_cb or self._auth_engine_id is None: 323 raise Exception('not accepting V3') 324 325 if req_msg.type == encoder.v3.MsgType.GET_REQUEST: 326 req = common.GetDataReq(names=[i.name for i in req_msg.pdu.data]) 327 328 elif req_msg.type == encoder.v3.MsgType.GET_NEXT_REQUEST: 329 req = common.GetNextDataReq( 330 names=[i.name for i in req_msg.pdu.data]) 331 332 elif req_msg.type == encoder.v3.MsgType.GET_BULK_REQUEST: 333 req = common.GetBulkDataReq( 334 names=[i.name for i in req_msg.pdu.data]) 335 336 elif req_msg.type == encoder.v3.MsgType.SET_REQUEST: 337 req = common.SetDataReq(data=req_msg.pdu.data) 338 339 else: 340 raise Exception('invalid request message type') 341 342 if req_msg.authorative_engine.id != self._auth_engine_id: 343 if req_msg.reportable: 344 345 # TODO report data and conditions for sending reports 346 347 authorative_engine = encoder.v3.AuthorativeEngine( 348 id=self._auth_engine_id, 349 boots=0, 350 time=round(time.monotonic())) 351 352 res_pdu = encoder.v3.BasicPdu( 353 request_id=req_msg.pdu.request_id, 354 error=common.Error(common.ErrorType.NO_ERROR, 0), 355 data=[]) 356 357 res_msg = encoder.v3.Msg( 358 type=encoder.v3.MsgType.REPORT, 359 id=req_msg.id, 360 reportable=False, 361 auth=False, 362 priv=False, 363 authorative_engine=authorative_engine, 364 user='', 365 context=req_msg.context, 366 pdu=res_pdu) 367 368 return res_msg 369 370 raise Exception('invalid authoritative engine id') 371 372 # TODO check authoritative engine boot and time 373 374 if (req_msg.user not in self._auth_keys or 375 req_msg.user not in self._priv_keys): 376 raise Exception('invalid user') 377 378 if self._auth_keys[req_msg.user] is not None and not req_msg.auth: 379 raise Exception('invalid auth flag') 380 381 if self._priv_keys[req_msg.user] is not None and not req_msg.priv: 382 raise Exception('invalid priv flag') 383 384 try: 385 res = await aio.call(self._v3_request_cb, addr, req_msg.user, 386 req_msg.context, req) 387 388 if isinstance(res, common.Error): 389 res_error = res 390 res_data = [] 391 392 else: 393 res_error = common.Error(common.ErrorType.NO_ERROR, 0) 394 res_data = res 395 396 except Exception as e: 397 self._log.warning("error processing request: %s", e, exc_info=e) 398 399 res_error = common.Error(common.ErrorType.GEN_ERR, 0) 400 res_data = [] 401 402 authorative_engine = encoder.v3.AuthorativeEngine( 403 id=req_msg.authorative_engine.id, 404 boots=0, 405 time=round(time.monotonic())) 406 407 res_pdu = encoder.v3.BasicPdu( 408 request_id=req_msg.pdu.request_id, 409 error=res_error, 410 data=res_data) 411 412 # TODO can we reuse request id for res msg id 413 414 res_msg = encoder.v3.Msg( 415 type=encoder.v3.MsgType.RESPONSE, 416 id=req_msg.id, 417 reportable=False, 418 auth=req_msg.auth, 419 priv=req_msg.priv, 420 authorative_engine=authorative_engine, 421 user=req_msg.user, 422 context=req_msg.context, 423 pdu=res_pdu) 424 425 return res_msg
Resource with lifetime control based on Group.
72 def __init__(self, 73 endpoint: net.Endpoint, 74 v1_request_cb: V1RequestCb | None, 75 v2c_request_cb: V2CRequestCb | None, 76 v3_request_cb: V3RequestCb | None, 77 authoritative_engine_id: common.EngineId | None, 78 users: Collection[common.User]): 79 self._endpoint = endpoint 80 self._v1_request_cb = v1_request_cb 81 self._v2c_request_cb = v2c_request_cb 82 self._v3_request_cb = v3_request_cb 83 self._auth_engine_id = authoritative_engine_id 84 self._auth_keys = {} 85 self._priv_keys = {} 86 self._log = logger.create_logger(mlog, 'SnmpAgent', endpoint.info) 87 self._comm_log = logger.CommunicationLogger(mlog, 'SnmpAgent', 88 endpoint.info) 89 90 for user in users: 91 common.validate_user(user) 92 93 if user.auth_type: 94 key_type = key.auth_type_to_key_type(user.auth_type) 95 self._auth_keys[user.name] = key.create_key( 96 key_type=key_type, 97 password=user.auth_password, 98 engine_id=authoritative_engine_id) 99 100 else: 101 self._auth_keys[user.name] = None 102 103 if user.priv_type: 104 key_type = key.priv_type_to_key_type(user.priv_type) 105 self._priv_keys[user.name] = key.create_key( 106 key_type=key_type, 107 password=user.priv_password, 108 engine_id=authoritative_engine_id) 109 110 else: 111 self._priv_keys[user.name] = None 112 113 self.async_group.spawn(self._receive_loop) 114 115 self.async_group.spawn(aio.call_on_cancel, self._comm_log.log, 116 common.CommLogAction.CLOSE) 117 self._comm_log.log(common.CommLogAction.OPEN)
11class Manager(aio.Resource): 12 13 @abc.abstractmethod 14 async def send(self, req: Request) -> Response: 15 """Send request and wait for response"""
Resource with lifetime control based on Group.
13 @abc.abstractmethod 14 async def send(self, req: Request) -> Response: 15 """Send request and wait for response"""
Send request and wait for response
18async def create_v1_manager(remote_addr: net.DatagramAddress, 19 community: common.CommunityName = 'public', 20 **kwargs 21 ) -> common.Manager: 22 """Create v1 manager""" 23 if isinstance(remote_addr, net.UdpAddress): 24 datagram_type = net.DatagramType.UDP 25 26 elif isinstance(remote_addr, net.UnixAddress): 27 datagram_type = net.DatagramType.UNIX 28 29 else: 30 raise TypeError('unsupported address type') 31 32 endpoint = await net.create_endpoint(datagram_type=datagram_type, 33 local_addr=None, 34 remote_addr=remote_addr, 35 **kwargs) 36 37 try: 38 return V1Manager(endpoint=endpoint, 39 community=community) 40 41 except BaseException: 42 await aio.uncancellable(endpoint.async_close()) 43 raise
Create v1 manager
18async def create_v2c_manager(remote_addr: net.DatagramAddress, 19 community: common.CommunityName = 'public', 20 **kwargs 21 ) -> common.Manager: 22 """Create v2c manager""" 23 if isinstance(remote_addr, net.UdpAddress): 24 datagram_type = net.DatagramType.UDP 25 26 elif isinstance(remote_addr, net.UnixAddress): 27 datagram_type = net.DatagramType.UNIX 28 29 else: 30 raise TypeError('unsupported address type') 31 32 endpoint = await net.create_endpoint(datagram_type=datagram_type, 33 local_addr=None, 34 remote_addr=remote_addr, 35 **kwargs) 36 37 try: 38 return V2CManager(endpoint=endpoint, 39 community=community) 40 41 except BaseException: 42 await aio.uncancellable(endpoint.async_close()) 43 raise
Create v2c manager
26async def create_v3_manager(remote_addr: net.DatagramAddress, 27 context: common.Context | None = None, 28 user: common.User = _default_user, 29 **kwargs 30 ) -> common.Manager: 31 """Create v3 manager""" 32 if isinstance(remote_addr, net.UdpAddress): 33 datagram_type = net.DatagramType.UDP 34 35 elif isinstance(remote_addr, net.UnixAddress): 36 datagram_type = net.DatagramType.UNIX 37 38 else: 39 raise TypeError('unsupported address type') 40 41 endpoint = await net.create_endpoint(datagram_type=datagram_type, 42 local_addr=None, 43 remote_addr=remote_addr, 44 **kwargs) 45 46 try: 47 manager = V3Manager(endpoint=endpoint, 48 context=context, 49 user=user) 50 51 except BaseException: 52 await aio.uncancellable(endpoint.async_close()) 53 raise 54 55 try: 56 await manager.sync() 57 58 except BaseException: 59 await aio.uncancellable(manager.async_close()) 60 raise 61 62 return manager
Create v3 manager
44async def create_trap_listener(local_addr: net.DatagramAddress = net.UdpAddress('0.0.0.0', 162), # NOQA 45 *, 46 v1_trap_cb: V1TrapCb | None = None, 47 v2c_trap_cb: V2CTrapCb | None = None, 48 v2c_inform_cb: V2CInformCb | None = None, 49 v3_trap_cb: V3TrapCb | None = None, 50 v3_inform_cb: V3InformCb | None = None, 51 users: Collection[common.User] = [], 52 **kwargs 53 ) -> 'TrapListener': 54 """Create trap listener""" 55 if isinstance(local_addr, net.UdpAddress): 56 datagram_type = net.DatagramType.UDP 57 58 elif isinstance(local_addr, net.UnixAddress): 59 datagram_type = net.DatagramType.UNIX 60 61 else: 62 raise TypeError('unsupported address type') 63 64 endpoint = await net.create_endpoint(datagram_type=datagram_type, 65 local_addr=local_addr, 66 remote_addr=None, 67 **kwargs) 68 69 try: 70 return TrapListener(endpoint=endpoint, 71 v1_trap_cb=v1_trap_cb, 72 v2c_trap_cb=v2c_trap_cb, 73 v2c_inform_cb=v2c_inform_cb, 74 v3_trap_cb=v3_trap_cb, 75 v3_inform_cb=v3_inform_cb, 76 users=users) 77 78 except BaseException: 79 await aio.uncancellable(endpoint.async_close()) 80 raise
Create trap listener
83class TrapListener(aio.Resource): 84 85 def __init__(self, 86 endpoint: net.Endpoint, 87 v1_trap_cb: V1TrapCb | None, 88 v2c_trap_cb: V2CTrapCb | None, 89 v2c_inform_cb: V2CInformCb | None, 90 v3_trap_cb: V3TrapCb | None, 91 v3_inform_cb: V3InformCb | None, 92 users: Collection[common.User]): 93 self._endpoint = endpoint 94 self._v1_trap_cb = v1_trap_cb 95 self._v2c_trap_cb = v2c_trap_cb 96 self._v2c_inform_cb = v2c_inform_cb 97 self._v3_trap_cb = v3_trap_cb 98 self._v3_inform_cb = v3_inform_cb 99 self._users = {} 100 self._auth_keys = {} 101 self._priv_keys = {} 102 self._log = logger.create_logger(mlog, 'SnmpTrapListener', 103 endpoint.info) 104 self._comm_log = logger.CommunicationLogger(mlog, 'SnmpTrapListener', 105 endpoint.info) 106 107 for user in users: 108 common.validate_user(user) 109 self._users[user.name] = user 110 111 self.async_group.spawn(self._receive_loop) 112 113 self.async_group.spawn(aio.call_on_cancel, self._comm_log.log, 114 common.CommLogAction.CLOSE) 115 self._comm_log.log(common.CommLogAction.OPEN) 116 117 @property 118 def async_group(self) -> aio.Group: 119 """Async group""" 120 return self._endpoint.async_group 121 122 def _on_auth_key(self, engine_id, username): 123 user = self._users.get(username) 124 if not user or not user.auth_type: 125 return 126 127 auth_key = self._auth_keys.get((engine_id, username)) 128 if auth_key: 129 return auth_key 130 131 key_type = key.auth_type_to_key_type(user.auth_type) 132 auth_key = key.create_key(key_type=key_type, 133 password=user.auth_password, 134 engine_id=engine_id) 135 136 self._auth_keys[(engine_id, username)] = auth_key 137 return auth_key 138 139 def _on_priv_key(self, engine_id, username): 140 user = self._users.get(username) 141 if not user or not user.priv_type: 142 return 143 144 priv_key = self._priv_keys.get((engine_id, username)) 145 if priv_key: 146 return priv_key 147 148 key_type = key.priv_type_to_key_type(user.priv_type) 149 priv_key = key.create_key(key_type=key_type, 150 password=user.priv_password, 151 engine_id=engine_id) 152 153 self._priv_keys[(engine_id, username)] = priv_key 154 return priv_key 155 156 async def _receive_loop(self): 157 try: 158 while True: 159 req_msg_bytes, addr = await self._endpoint.receive() 160 161 try: 162 req_msg = encoder.decode(msg_bytes=req_msg_bytes, 163 auth_key_cb=self._on_auth_key, 164 priv_key_cb=self._on_priv_key) 165 166 except Exception as e: 167 self._log.warning("error decoding message from %s: %s", 168 addr, e, exc_info=e) 169 continue 170 171 self._comm_log.log(common.CommLogAction.RECEIVE, req_msg) 172 173 try: 174 if isinstance(req_msg, encoder.v1.Msg): 175 res_msg = await _process_v1_req_msg( 176 req_msg=req_msg, 177 addr=addr, 178 trap_cb=self._v1_trap_cb) 179 180 elif isinstance(req_msg, encoder.v2c.Msg): 181 res_msg = await _process_v2c_req_msg( 182 req_msg=req_msg, 183 addr=addr, 184 trap_cb=self._v2c_trap_cb, 185 inform_cb=self._v2c_inform_cb) 186 187 elif isinstance(req_msg, encoder.v3.Msg): 188 res_msg = await _process_v3_req_msg( 189 req_msg=req_msg, 190 addr=addr, 191 trap_cb=self._v3_trap_cb, 192 inform_cb=self._v3_inform_cb) 193 194 else: 195 raise ValueError('unsupported message type') 196 197 except Exception as e: 198 self._log.warning("error processing message from %s: %s", 199 addr, e, exc_info=e) 200 continue 201 202 if not res_msg: 203 continue 204 205 try: 206 if isinstance(res_msg, encoder.v3.Msg): 207 auth_key = ( 208 self._on_auth_key(res_msg.authorative_engine.id, 209 res_msg.user) 210 if res_msg.auth else None) 211 priv_key = ( 212 self._on_priv_key(res_msg.authorative_engine.id, 213 res_msg.user) 214 if res_msg.priv else None) 215 216 else: 217 auth_key = None 218 priv_key = None 219 220 res_msg_bytes = encoder.encode(msg=res_msg, 221 auth_key=auth_key, 222 priv_key=priv_key) 223 224 except Exception as e: 225 self._log.warning("error encoding message: %s", 226 e, exc_info=e) 227 continue 228 229 self._comm_log.log(common.CommLogAction.SEND, res_msg) 230 231 self._endpoint.send(res_msg_bytes, addr) 232 233 except ConnectionError: 234 pass 235 236 except Exception as e: 237 self._log.error("receive loop error: %s", e, exc_info=e) 238 239 finally: 240 self.close()
Resource with lifetime control based on Group.
85 def __init__(self, 86 endpoint: net.Endpoint, 87 v1_trap_cb: V1TrapCb | None, 88 v2c_trap_cb: V2CTrapCb | None, 89 v2c_inform_cb: V2CInformCb | None, 90 v3_trap_cb: V3TrapCb | None, 91 v3_inform_cb: V3InformCb | None, 92 users: Collection[common.User]): 93 self._endpoint = endpoint 94 self._v1_trap_cb = v1_trap_cb 95 self._v2c_trap_cb = v2c_trap_cb 96 self._v2c_inform_cb = v2c_inform_cb 97 self._v3_trap_cb = v3_trap_cb 98 self._v3_inform_cb = v3_inform_cb 99 self._users = {} 100 self._auth_keys = {} 101 self._priv_keys = {} 102 self._log = logger.create_logger(mlog, 'SnmpTrapListener', 103 endpoint.info) 104 self._comm_log = logger.CommunicationLogger(mlog, 'SnmpTrapListener', 105 endpoint.info) 106 107 for user in users: 108 common.validate_user(user) 109 self._users[user.name] = user 110 111 self.async_group.spawn(self._receive_loop) 112 113 self.async_group.spawn(aio.call_on_cancel, self._comm_log.log, 114 common.CommLogAction.CLOSE) 115 self._comm_log.log(common.CommLogAction.OPEN)
11class TrapSender(aio.Resource): 12 13 @abc.abstractmethod 14 def send_trap(self, trap: Trap): 15 """Send trap""" 16 17 @abc.abstractmethod 18 async def send_inform(self, 19 inform: Inform 20 ) -> Error | None: 21 """Send inform"""
Resource with lifetime control based on Group.
16async def create_v1_trap_sender(remote_addr: net.DatagramAddress, 17 community: common.CommunityName = 'public', 18 **kwargs 19 ) -> common.TrapSender: 20 """Create v1 trap sender""" 21 if isinstance(remote_addr, net.UdpAddress): 22 datagram_type = net.DatagramType.UDP 23 24 elif isinstance(remote_addr, net.UnixAddress): 25 datagram_type = net.DatagramType.UNIX 26 27 else: 28 raise TypeError('unsupported address type') 29 30 endpoint = await net.create_endpoint(datagram_type=datagram_type, 31 local_addr=None, 32 remote_addr=remote_addr, 33 **kwargs) 34 35 try: 36 return V1TrapSender(endpoint=endpoint, 37 community=community) 38 39 except BaseException: 40 await aio.uncancellable(endpoint.async_close()) 41 raise
Create v1 trap sender
18async def create_v2c_trap_sender(remote_addr: net.DatagramAddress, 19 community: common.CommunityName = 'public', 20 **kwargs 21 ) -> common.TrapSender: 22 """Create v2c trap sender""" 23 if isinstance(remote_addr, net.UdpAddress): 24 datagram_type = net.DatagramType.UDP 25 26 elif isinstance(remote_addr, net.UnixAddress): 27 datagram_type = net.DatagramType.UNIX 28 29 else: 30 raise TypeError('unsupported address type') 31 32 endpoint = await net.create_endpoint(datagram_type=datagram_type, 33 local_addr=None, 34 remote_addr=remote_addr, 35 **kwargs) 36 37 try: 38 return V2CTrapSender(endpoint=endpoint, 39 community=community) 40 41 except BaseException: 42 await aio.uncancellable(endpoint.async_close()) 43 raise
Create v2c trap sender
26async def create_v3_trap_sender(remote_addr: net.DatagramAddress, 27 authoritative_engine_id: common.EngineId, 28 context: common.Context | None = None, 29 user: common.User = _default_user, 30 **kwargs 31 ) -> common.TrapSender: 32 """Create v3 trap sender""" 33 if isinstance(remote_addr, net.UdpAddress): 34 datagram_type = net.DatagramType.UDP 35 36 elif isinstance(remote_addr, net.UnixAddress): 37 datagram_type = net.DatagramType.UNIX 38 39 else: 40 raise TypeError('unsupported address type') 41 42 endpoint = await net.create_endpoint(datagram_type=datagram_type, 43 local_addr=None, 44 remote_addr=remote_addr, 45 **kwargs) 46 47 try: 48 return V3TrapSender(endpoint=endpoint, 49 authoritative_engine_id=authoritative_engine_id, 50 context=context, 51 user=user) 52 53 except BaseException: 54 await aio.uncancellable(endpoint.async_close()) 55 raise
Create v3 trap sender