Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 64 additions & 1 deletion nisshi-storage/src/proxy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ use nisshi_sans_io::{
txn_offset_commit_response::TxnOffsetCommitResponseTopic,
};
use opentelemetry::{KeyValue, metrics::Histogram};
use tokio::sync::Semaphore;
use tokio::{sync::Semaphore, time::Instant};
use tracing::instrument;
use url::Url;
use uuid::Uuid;
Expand Down Expand Up @@ -187,12 +187,19 @@ where
max_wait: Duration,
) -> Result<Vec<deflated::Batch>> {
let start = SystemTime::now();
let started_at = Instant::now();
let _permit = self.semaphore.acquire().await.inspect(|_| {
SEMAPHORE_ACQUIRE_DURATION.record(
elapsed_millis(start),
&[KeyValue::new("operation", "fetch")],
)
})?;

// The wait for the permit is charged to the budget, so that the
// read ends when its caller's budget does, rather than that long
// after the wait.
let max_wait = max_wait.saturating_sub(started_at.elapsed());
Comment thread
solace-aross marked this conversation as resolved.

self.storage
.fetch(topition, offset, min_bytes, max_bytes, isolation, max_wait)
.await
Expand Down Expand Up @@ -554,3 +561,59 @@ where
self.storage.ping().await
}
}

#[cfg(test)]
mod tests {
use nisshi_sans_io::IsolationLevel;
use tokio::time::{Duration, advance};

use super::SemaphoreProxy;
use crate::{
Result, Storage, Topition,
service::{Request, RequestChannelService, Response, bounded_channel},
};

/// A read that waits for a permit reaches storage with the budget left
/// after the wait, so that the read ends when its caller's budget does.
#[tokio::test(start_paused = true)]
async fn fetch_charges_the_permit_wait_to_the_budget() -> Result<()> {
let (tx, mut rx) = bounded_channel(1);
let proxy = SemaphoreProxy::new(RequestChannelService::new(tx));

let permit = proxy.semaphore.clone().acquire_owned().await?;
let tp = Topition::new("abc", 0);

let read = proxy.fetch(
&tp,
0,
1,
1024,
IsolationLevel::ReadUncommitted,
Duration::from_millis(500),
);

let release = async {
advance(Duration::from_millis(200)).await;
drop(permit);
};

let answer = async {
let (request, reply) = rx.recv().await.expect("a fetch request");

let Request::Fetch { max_wait, .. } = request else {
panic!("{request:?}");
};

reply.send(Response::Fetch(Ok(vec![]))).expect("the reader");

max_wait
};

let (fetched, (), max_wait) = tokio::join!(read, release, answer);

assert!(fetched?.is_empty());
assert_eq!(Duration::from_millis(300), max_wait);

Ok(())
}
}
Loading
Loading