@@ -681,9 +681,40 @@ try_read(
681681 ),
682682 case RemoteRange of
683683 {RemoteStartOffset , _EndOffset } when RemoteStartOffset >= NextFragment ->
684- % % re-set the state at this offset. Also take the min between this
685- % % and the first local offset.
686- erlang :error (unimplemented );
684+ % % Retention has evicted the next fragment and advanced past it.
685+ % % Jump to the oldest fragment still available. Cancel any
686+ % % in-flight requests and track their refs as cancelled so that
687+ % % late frames are dropped silently.
688+ #? MODULE {
689+ fragment = CurrentFragment ,
690+ requests = Requests ,
691+ cancelled_requests = Cancelled0
692+ } = State0 ,
693+ cancel_requests (Requests ),
694+ Cancelled = maps :merge (
695+ Cancelled0 ,
696+ #{Req => ok || Req := _ <- Requests }
697+ ),
698+ Key = rabbitmq_stream_s3 :fragment_key (StreamId , RemoteStartOffset ),
699+ ? LOG_WARNING (
700+ " Fragment ~20..0B finished but next fragment ~20..0B was evicted. "
701+ " Jumping to oldest available fragment ~20..0B " ,
702+ [CurrentFragment , NextFragment , RemoteStartOffset ]
703+ ),
704+ State = State0 #? MODULE {
705+ fragment = RemoteStartOffset ,
706+ key = Key ,
707+ info = undefined ,
708+ buffer = <<>>,
709+ start_pos = ? FRAGMENT_HEADER_B ,
710+ current_pos = ? FRAGMENT_HEADER_B ,
711+ end_pos = ? FRAGMENT_HEADER_B ,
712+ next = undefined ,
713+ current_not_found = false ,
714+ requests = #{},
715+ cancelled_requests = Cancelled
716+ },
717+ {next_fragment , State , RemoteStartOffset };
687718 {_StartOffset , EndOffset } when EndOffset < NextFragment ->
688719 State = goto_next_fragment (State0 ),
689720 {become_local , State };
0 commit comments