@@ -74,6 +74,7 @@ def stub_message(id: int):
7474 written_at = datetime .datetime (2023 , 3 , 18 , 14 , 15 ),
7575 producer_id = "" ,
7676 data = bytes (),
77+ metadata_items = {},
7778 _partition_session = stub_partition_session (),
7879 _commit_start_offset = 0 ,
7980 _commit_end_offset = 1 ,
@@ -207,6 +208,7 @@ def create_message(
207208 written_at = datetime .datetime (2023 , 2 , 3 , 14 , 16 ),
208209 producer_id = "test-producer-id" ,
209210 data = bytes (),
211+ metadata_items = {},
210212 _partition_session = partition_session ,
211213 _commit_start_offset = partition_session ._next_message_start_commit_offset + offset_delta - 1 ,
212214 _commit_end_offset = partition_session ._next_message_start_commit_offset + offset_delta ,
@@ -250,6 +252,7 @@ def batch_size():
250252 seq_no = message .seqno ,
251253 created_at = message .created_at ,
252254 data = message .data ,
255+ metadata_items = {},
253256 uncompresed_size = len (message .data ),
254257 message_group_id = message .message_group_id ,
255258 )
@@ -445,6 +448,7 @@ async def test_commit_ranges_for_received_messages(
445448 written_at = datetime .datetime (2023 , 3 , 14 , 15 , 42 ),
446449 producer_id = "asd" ,
447450 data = rb"123" ,
451+ metadata_items = {},
448452 _partition_session = None ,
449453 _commit_start_offset = 5 ,
450454 _commit_end_offset = 15 ,
@@ -468,6 +472,7 @@ async def test_commit_ranges_for_received_messages(
468472 written_at = datetime .datetime (2023 , 3 , 14 , 15 , 42 ),
469473 producer_id = "asd" ,
470474 data = gzip .compress (rb"123" ),
475+ metadata_items = {},
471476 _partition_session = None ,
472477 _commit_start_offset = 5 ,
473478 _commit_end_offset = 15 ,
@@ -490,6 +495,7 @@ async def test_commit_ranges_for_received_messages(
490495 offset = 1 ,
491496 written_at = datetime .datetime (2023 , 3 , 14 , 15 , 42 ),
492497 producer_id = "asd" ,
498+ metadata_items = {},
493499 data = rb"123" ,
494500 _partition_session = None ,
495501 _commit_start_offset = 5 ,
@@ -504,6 +510,7 @@ async def test_commit_ranges_for_received_messages(
504510 written_at = datetime .datetime (2023 , 3 , 14 , 15 , 42 ),
505511 producer_id = "asd" ,
506512 data = rb"456" ,
513+ metadata_items = {},
507514 _partition_session = None ,
508515 _commit_start_offset = 5 ,
509516 _commit_end_offset = 15 ,
@@ -527,6 +534,7 @@ async def test_commit_ranges_for_received_messages(
527534 written_at = datetime .datetime (2023 , 3 , 14 , 15 , 42 ),
528535 producer_id = "asd" ,
529536 data = gzip .compress (rb"123" ),
537+ metadata_items = {},
530538 _partition_session = None ,
531539 _commit_start_offset = 5 ,
532540 _commit_end_offset = 15 ,
@@ -540,6 +548,7 @@ async def test_commit_ranges_for_received_messages(
540548 written_at = datetime .datetime (2023 , 3 , 14 , 15 , 42 ),
541549 producer_id = "asd" ,
542550 data = gzip .compress (rb"456" ),
551+ metadata_items = {},
543552 _partition_session = None ,
544553 _commit_start_offset = 5 ,
545554 _commit_end_offset = 15 ,
@@ -766,6 +775,7 @@ async def test_free_buffer_after_partition_stop(self, stream, stream_reader, par
766775 seq_no = 123 ,
767776 created_at = t ,
768777 data = bytes (),
778+ metadata_items = {},
769779 uncompresed_size = message_size ,
770780 message_group_id = "test-message-group" ,
771781 )
@@ -846,6 +856,7 @@ def reader_batch_count():
846856 created_at = created_at ,
847857 data = data ,
848858 uncompresed_size = len (data ),
859+ metadata_items = {},
849860 message_group_id = message_group_id ,
850861 )
851862 ],
@@ -877,6 +888,7 @@ def reader_batch_count():
877888 written_at = written_at ,
878889 producer_id = producer_id ,
879890 data = data ,
891+ metadata_items = {},
880892 _partition_session = partition_session ,
881893 _commit_start_offset = expected_message_offset ,
882894 _commit_end_offset = expected_message_offset + 1 ,
@@ -923,6 +935,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
923935 seq_no = 3 ,
924936 created_at = created_at ,
925937 data = data ,
938+ metadata_items = {},
926939 uncompresed_size = len (data ),
927940 message_group_id = message_group_id ,
928941 )
@@ -944,6 +957,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
944957 seq_no = 2 ,
945958 created_at = created_at2 ,
946959 data = data ,
960+ metadata_items = {},
947961 uncompresed_size = len (data ),
948962 message_group_id = message_group_id ,
949963 )
@@ -960,6 +974,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
960974 seq_no = 3 ,
961975 created_at = created_at3 ,
962976 data = data2 ,
977+ metadata_items = {},
963978 uncompresed_size = len (data2 ),
964979 message_group_id = message_group_id ,
965980 ),
@@ -968,6 +983,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
968983 seq_no = 5 ,
969984 created_at = created_at4 ,
970985 data = data ,
986+ metadata_items = {},
971987 uncompresed_size = len (data ),
972988 message_group_id = message_group_id2 ,
973989 ),
@@ -998,6 +1014,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
9981014 written_at = written_at ,
9991015 producer_id = producer_id ,
10001016 data = data ,
1017+ metadata_items = {},
10011018 _partition_session = partition_session ,
10021019 _commit_start_offset = partition1_mess1_expected_offset ,
10031020 _commit_end_offset = partition1_mess1_expected_offset + 1 ,
@@ -1018,6 +1035,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
10181035 written_at = written_at2 ,
10191036 producer_id = producer_id ,
10201037 data = data ,
1038+ metadata_items = {},
10211039 _partition_session = second_partition_session ,
10221040 _commit_start_offset = partition2_mess1_expected_offset ,
10231041 _commit_end_offset = partition2_mess1_expected_offset + 1 ,
@@ -1038,6 +1056,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
10381056 written_at = written_at2 ,
10391057 producer_id = producer_id2 ,
10401058 data = data2 ,
1059+ metadata_items = {},
10411060 _partition_session = second_partition_session ,
10421061 _commit_start_offset = partition2_mess2_expected_offset ,
10431062 _commit_end_offset = partition2_mess2_expected_offset + 1 ,
@@ -1051,6 +1070,7 @@ async def test_read_batches(self, stream_reader, partition_session, second_parti
10511070 written_at = written_at2 ,
10521071 producer_id = producer_id ,
10531072 data = data ,
1073+ metadata_items = {},
10541074 _partition_session = second_partition_session ,
10551075 _commit_start_offset = partition2_mess3_expected_offset ,
10561076 _commit_end_offset = partition2_mess3_expected_offset + 1 ,
0 commit comments