1212import logging
1313import struct
1414from dataclasses import dataclass
15- from functools import partial
1615from pathlib import Path
1716
1817from pycyphal2 import Arrival , DeliveryError , NackError , Node , SendError
3231
3332@dataclass (frozen = True )
3433class FileReadRequest :
35- read_offset : int
36- file_path : str
34+ _read_offset : int
35+ _file_path : str
36+
37+ @property
38+ def read_offset (self ) -> int :
39+ return self ._read_offset
40+
41+ @property
42+ def file_path (self ) -> str :
43+ return self ._file_path
44+
45+ @staticmethod
46+ def deserialize (payload : bytes ) -> FileReadRequest | None :
47+ if len (payload ) < REQUEST_HEADER_SIZE :
48+ return None
49+ read_offset , path_len = struct .unpack_from (REQUEST_HEADER_FORMAT , payload )
50+ if path_len == 0 or path_len > PATH_MAX_LEN :
51+ return None
52+ path_end = REQUEST_HEADER_SIZE + path_len
53+ if len (payload ) != path_end :
54+ return None
55+ try :
56+ file_path = payload [REQUEST_HEADER_SIZE :path_end ].decode ("utf8" )
57+ except UnicodeDecodeError :
58+ return None
59+ return FileReadRequest (read_offset , file_path )
3760
3861
3962@dataclass (frozen = True )
4063class FileReadResponse :
41- error : int
42- data : bytes
43-
44-
45- def _decode_request (payload : bytes ) -> FileReadRequest | None :
46- if len (payload ) < REQUEST_HEADER_SIZE :
47- return None
48- read_offset , path_len = struct .unpack_from (REQUEST_HEADER_FORMAT , payload )
49- if path_len == 0 or path_len > PATH_MAX_LEN :
50- return None
51- path_end = REQUEST_HEADER_SIZE + path_len
52- if len (payload ) != path_end :
53- return None
54- try :
55- file_path = payload [REQUEST_HEADER_SIZE :path_end ].decode ("utf8" )
56- except UnicodeDecodeError :
57- return None
58- return FileReadRequest (read_offset = read_offset , file_path = file_path )
64+ _error : int
65+ _data : bytes
66+
67+ @property
68+ def error (self ) -> int :
69+ return self ._error
5970
71+ @property
72+ def data (self ) -> bytes :
73+ return self ._data
6074
61- def _encode_response ( response : FileReadResponse ) -> bytes :
62- if len (response . data ) > DATA_MAX :
63- raise ValueError (f"Response data is too large: { len (response . data )} " )
64- return struct .pack (RESPONSE_HEADER_FORMAT , response . error , len (response . data )) + response . data
75+ def serialize ( self ) -> bytes :
76+ if len (self . _data ) > DATA_MAX :
77+ raise ValueError (f"Response data is too large: { len (self . _data )} " )
78+ return struct .pack (RESPONSE_HEADER_FORMAT , self . _error , len (self . _data )) + self . _data
6579
6680
6781def _errno_from_exception (ex : BaseException ) -> int :
@@ -78,13 +92,13 @@ def _read_chunk(file_path: str, offset: int) -> FileReadResponse:
7892 file .seek (offset )
7993 data = file .read (DATA_MAX )
8094 except (OSError , ValueError , OverflowError ) as ex :
81- return FileReadResponse (error = _errno_from_exception (ex ), data = b"" )
82- return FileReadResponse (error = 0 , data = data )
95+ return FileReadResponse (_errno_from_exception (ex ), b"" )
96+ return FileReadResponse (0 , data )
8397
8498
8599async def _serve_request (arrival : Arrival , request : FileReadRequest ) -> None :
86100 response = _read_chunk (request .file_path , request .read_offset )
87- payload = _encode_response ( response )
101+ payload = response . serialize ( )
88102 _logger .info (
89103 "responding: file=%r offset=%d size=%d error=%d" ,
90104 request .file_path ,
@@ -104,36 +118,20 @@ async def _serve_request(arrival: Arrival, request: FileReadRequest) -> None:
104118 _logger .warning ("response send failed: remote=%016x error=%s" , arrival .breadcrumb .remote_id , ex )
105119
106120
107- def _on_task_done (tasks : set [asyncio .Task [None ]], task : asyncio .Task [None ]) -> None :
108- tasks .discard (task )
109- if task .cancelled ():
110- return
111- exc = task .exception ()
112- if exc is not None :
113- _logger .error ("file request task failed: %s" , exc )
114-
115-
116121async def run () -> None :
117122 transport = UDPTransport .new ()
118123 node = Node .new (transport , NAME )
119124 sub = node .subscribe (TOPIC )
120- tasks : set [asyncio .Task [None ]] = set ()
121125 _logger .info ("file server ready on %r via %s" , TOPIC , transport )
122126 try :
123127 async for arrival in sub :
124- request = _decode_request (arrival .message )
128+ request = FileReadRequest . deserialize (arrival .message )
125129 if request is None :
126130 _logger .debug ("dropping malformed request of size %d" , len (arrival .message ))
127131 continue
128- task = asyncio .create_task (_serve_request (arrival , request ), name = f"file:{ arrival .breadcrumb .tag } " )
129- tasks .add (task )
130- task .add_done_callback (partial (_on_task_done , tasks ))
132+ await _serve_request (arrival , request )
131133 finally :
132134 sub .close ()
133- for task in list (tasks ):
134- task .cancel ()
135- if tasks :
136- await asyncio .gather (* tasks , return_exceptions = True )
137135 node .close ()
138136 transport .close ()
139137
0 commit comments