@@ -22,15 +22,11 @@ def __init__(self):
2222 self .counter = 0
2323 self .partition_idx = 0
2424
25- async def read_handler (
26- self , datum : sourcer .ReadRequest
27- ) -> AsyncIterator [sourcer .Message ]:
25+ async def read_handler (self , datum : sourcer .ReadRequest ) -> AsyncIterator [sourcer .Message ]:
2826 """
2927 The simple source generates messages with incrementing numbers.
3028 """
31- _LOGGER .info (
32- f"Read request: num_records={ datum .num_records } , timeout_ms={ datum .timeout_ms } "
33- )
29+ _LOGGER .info (f"Read request: num_records={ datum .num_records } , timeout_ms={ datum .timeout_ms } " )
3430
3531 # Generate the requested number of messages
3632 for _ in range (datum .num_records ):
@@ -66,19 +62,15 @@ async def ack_handler(self, request: sourcer.AckRequest) -> None:
6662 """
6763 _LOGGER .info (f"Acknowledging { len (request .offsets )} offsets" )
6864 for offset in request .offsets :
69- _LOGGER .debug (
70- f"Acked offset: { offset .offset .decode ('utf-8' )} , partition: { offset .partition_id } "
71- )
65+ _LOGGER .debug (f"Acked offset: { offset .offset .decode ('utf-8' )} , partition: { offset .partition_id } " )
7266
7367 async def nack_handler (self , request : sourcer .NackRequest ) -> None :
7468 """
7569 The simple source negatively acknowledges the offsets.
7670 """
7771 _LOGGER .info (f"Negatively acknowledging { len (request .offsets )} offsets" )
7872 for offset in request .offsets :
79- _LOGGER .warning (
80- f"Nacked offset: { offset .offset .decode ('utf-8' )} , partition: { offset .partition_id } "
81- )
73+ _LOGGER .warning (f"Nacked offset: { offset .offset .decode ('utf-8' )} , partition: { offset .partition_id } " )
8274
8375 async def pending_handler (self ) -> sourcer .PendingResponse :
8476 """
@@ -92,6 +84,7 @@ async def partitions_handler(self) -> sourcer.PartitionsResponse:
9284 """
9385 return sourcer .PartitionsResponse (partitions = [self .partition_idx ])
9486
87+
9588async def start ():
9689 server = sourcer .SourceAsyncServer ()
9790
0 commit comments