@@ -423,27 +423,32 @@ use_defaults(Config, [{K, V} | Rest]) ->
423423 false -> use_defaults (Config #{K => V }, Rest )
424424 end .
425425
426- remake_requests (#{q_items := Items } = Sent , St , ? reconnect ) ->
427- # kpro_req { ref = Ref } = Req = make_request (Items , St ),
428- {[ Req ], [ Sent #{req_ref := Ref }]} ;
429- remake_requests (#{q_items := Items } = Sent , St , ? message_too_large ) ->
430- Reqs = lists :map (fun (Item ) -> make_request ([Item ], St ) end , Items ),
426+ resend_requests (#{q_items := Items } = Sent , St , ? reconnect ) ->
427+ { ok , Ref } = send_request (Items , St ),
428+ [ Sent #{req_ref := Ref }];
429+ resend_requests (#{q_items := Items } = Sent , St , ? message_too_large ) ->
430+ Refs = lists :map (fun (Item ) -> { ok , Ref } = send_request ([Item ], St ), Ref end , Items ),
431431 #{attempts := Attempts , q_ack_ref := Q_AckRef } = Sent ,
432- { Reqs , make_sent_items_list (Items , Reqs , Attempts , Q_AckRef )} .
432+ make_sent_items_list (Items , Refs , Attempts , Q_AckRef ).
433433
434434% % only ack replayq when the last message is accepted by Kafka
435- make_sent_items_list ([LastItem ], [# kpro_req { ref = Ref } ], Attempts , Q_AckRef ) ->
435+ make_sent_items_list ([LastItem ], [Ref ], Attempts , Q_AckRef ) ->
436436 [#{req_ref => Ref ,
437437 q_items => [LastItem ],
438438 q_ack_ref => Q_AckRef ,
439439 attempts => Attempts
440440 }];
441- make_sent_items_list ([Item | Items ], [# kpro_req { ref = Ref } | Reqs ], Attempts , Q_AckRef ) ->
441+ make_sent_items_list ([Item | Items ], [Ref | Refs ], Attempts , Q_AckRef ) ->
442442 [#{req_ref => Ref ,
443443 q_items => [Item ],
444444 q_ack_ref => ? no_queue_ack ,
445445 attempts => Attempts
446- } | make_sent_items_list (Items , Reqs , Attempts , Q_AckRef )].
446+ } | make_sent_items_list (Items , Refs , Attempts , Q_AckRef )].
447+
448+ send_request (QueueItems , #{conn := Conn } = St ) ->
449+ # kpro_req {ref = Ref } = Req = make_request (QueueItems , St ),
450+ ok = request_async (Conn , Req ),
451+ {ok , Ref }.
447452
448453make_request (QueueItems ,
449454 #{config := #{required_acks := RequiredAcks ,
@@ -461,13 +466,11 @@ make_request(QueueItems,
461466 }).
462467
463468resend_sent_reqs (#{sent_reqs := SentReqs ,
464- conn := Conn ,
465469 config := Config
466470 } = St , Reason ) ->
467471 F = fun (Sent , Acc ) ->
468472 wolff_metrics :retried_inc (Config , 1 ),
469- {Reqs , NewSentList } = remake_requests (Sent , St , Reason ),
470- lists :foreach (fun (Req ) -> ok = request_async (Conn , Req ) end , Reqs ),
473+ NewSentList = resend_requests (Sent , St , Reason ),
471474 lists :reverse (NewSentList , Acc )
472475 end ,
473476 NewSentReqs = lists :foldl (F , [], queue :to_list (SentReqs )),
0 commit comments