feat: check fetch cancellation before each envelope
This commit is contained in:
parent
aa3950bc80
commit
21931c38d4
2 changed files with 24 additions and 6 deletions
|
|
@ -535,8 +535,8 @@ pub struct Fetched {
|
|||
/// cursor (the default).
|
||||
pub acknowledged: bool,
|
||||
/// `cancel` was raised mid-fetch; the summary covers what completed.
|
||||
/// Every finished page is accounted for locally — cursor moved or
|
||||
/// messages deleted — so fetching again continues rather than repeats.
|
||||
/// Every envelope the summary counted is accounted for locally, so
|
||||
/// fetching again continues rather than repeats.
|
||||
pub cancelled: bool,
|
||||
}
|
||||
|
||||
|
|
@ -551,8 +551,10 @@ pub struct FetchOptions<'a> {
|
|||
pub reset: bool,
|
||||
/// Network timeout in seconds.
|
||||
pub timeout: u64,
|
||||
/// Checked between pages; when raised, `fetch_with` stops and returns
|
||||
/// the partial summary with `cancelled` set.
|
||||
/// Checked between pages and before each envelope; when raised,
|
||||
/// `fetch_with` stops and returns the partial summary with `cancelled`
|
||||
/// set. What was processed is accounted for locally — cursor moved or
|
||||
/// messages deleted — so fetching again continues rather than repeats.
|
||||
pub cancel: Option<&'a std::sync::atomic::AtomicBool>,
|
||||
}
|
||||
|
||||
|
|
@ -632,7 +634,16 @@ pub fn fetch_with(
|
|||
break;
|
||||
}
|
||||
let mut acked: Vec<[u8; ID_LEN]> = Vec::new();
|
||||
let mut page_cancelled = false;
|
||||
for _ in 0..count {
|
||||
// Checked before each envelope as well as between pages, so a
|
||||
// cancel lands within one envelope; the cursor and acks below
|
||||
// then cover exactly the envelopes that were processed.
|
||||
if cancelled() {
|
||||
page_cancelled = true;
|
||||
summary.cancelled = true;
|
||||
break;
|
||||
}
|
||||
let mid: [u8; ID_LEN] = r.take(ID_LEN)?.try_into().unwrap();
|
||||
let received_at = r.i64()?;
|
||||
let flags = r.u8()?;
|
||||
|
|
@ -662,12 +673,19 @@ pub fn fetch_with(
|
|||
}
|
||||
acked.push(mid);
|
||||
}
|
||||
r.done()?;
|
||||
// A cancelled page is not fully parsed, so trailing bytes are
|
||||
// expected and not an error; everything processed is accounted for.
|
||||
if !page_cancelled {
|
||||
r.done()?;
|
||||
}
|
||||
if opts.keep {
|
||||
store.set_cursor(after_time, &after_id)?;
|
||||
} else if !acked.is_empty() {
|
||||
delete_ids(transport, &acked)?;
|
||||
}
|
||||
if page_cancelled {
|
||||
break;
|
||||
}
|
||||
}
|
||||
session.close();
|
||||
if !opts.keep {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue