Skip to content

Commit c798708

Browse files
authored
feat(cache-client): expose cache-only gateway operations (#517)
1 parent c25f10b commit c798708

1 file changed

Lines changed: 165 additions & 1 deletion

File tree

crates/talon-cache-client/src/block_reader.rs

Lines changed: 165 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ use crate::membership_cache::{MembershipCache, MembershipSnapshot};
3030
use crate::metrics::ReadStats;
3131
use crate::placement_cache::{Cached, PlacementCache, RefreshReason};
3232
use crate::pool::ConnectionPool;
33+
use crate::range_stream::CacheReadError;
3334
use crate::read_plan::plan_read;
3435
use crate::worker_client::{WorkerClient, WorkerError};
3536

@@ -152,6 +153,79 @@ impl BlockReader {
152153
&self.stats
153154
}
154155

156+
/// Read a versioned block slice only when it is already resident in Talon.
157+
///
158+
/// Unlike [`read_block`](Self::read_block), this operation cannot invoke a
159+
/// worker backend. It is therefore suitable for a gateway that must use a
160+
/// request-scoped client capability for every origin miss.
161+
pub async fn read_cached_block(
162+
&self,
163+
block: &BlockId,
164+
offset_in_block: u32,
165+
len: u32,
166+
now_ms: u64,
167+
) -> Result<Vec<u8>, CacheReadError> {
168+
let placement = match self.cache.get(block, now_ms) {
169+
Some(cached) => cached,
170+
None => self
171+
.resolve_and_cache(block, now_ms)
172+
.await
173+
.map_err(cache_block_error)?,
174+
};
175+
let offset = block.offset + u64::from(offset_in_block);
176+
let mut last_error = None;
177+
for address in &placement.replicas {
178+
let worker = WorkerClient::with_pool(address.clone(), Arc::clone(&self.worker_pool));
179+
match worker
180+
.fetch_cached_range(&block.object, &block.version, offset, u64::from(len))
181+
.await
182+
{
183+
Ok(bytes) if bytes.len() == len as usize => return Ok(bytes),
184+
Ok(bytes) => {
185+
last_error = Some(CacheReadError::Protocol(
186+
WorkerError::RangeLengthMismatch {
187+
expected: u64::from(len),
188+
actual: bytes.len() as u64,
189+
}
190+
.to_string(),
191+
));
192+
}
193+
Err(error) => last_error = Some(error.into()),
194+
}
195+
}
196+
Err(last_error.unwrap_or_else(|| {
197+
CacheReadError::Unavailable("placement contained no worker addresses".into())
198+
}))
199+
}
200+
201+
/// Admit one complete versioned block to its current primary owner.
202+
///
203+
/// The worker validates alignment, object length, and exact body length and
204+
/// commits atomically without invoking its configured backend.
205+
pub async fn admit_block(
206+
&self,
207+
block: &BlockId,
208+
object_len: u64,
209+
body: &[u8],
210+
now_ms: u64,
211+
) -> Result<(), CacheReadError> {
212+
let placement = match self.cache.get(block, now_ms) {
213+
Some(cached) => cached,
214+
None => self
215+
.resolve_and_cache(block, now_ms)
216+
.await
217+
.map_err(cache_block_error)?,
218+
};
219+
let address = placement
220+
.replicas
221+
.first()
222+
.ok_or_else(|| CacheReadError::Unavailable("block has no primary owner".into()))?;
223+
WorkerClient::with_pool(address.clone(), Arc::clone(&self.worker_pool))
224+
.admit_cached_block(block, object_len, body)
225+
.await
226+
.map_err(Into::into)
227+
}
228+
155229
/// Read `len` bytes at `offset_in_block` within `block`.
156230
///
157231
/// Resolves placement (cache hit, else local Maglev lookup),
@@ -474,6 +548,14 @@ impl BlockReader {
474548
}
475549
}
476550

551+
fn cache_block_error(error: BlockReadError) -> CacheReadError {
552+
match error {
553+
BlockReadError::Coordinator(error) => error.into(),
554+
BlockReadError::Worker(error) => error.into(),
555+
other => CacheReadError::Unavailable(other.to_string()),
556+
}
557+
}
558+
477559
fn replica_retryable(error: &WorkerError) -> bool {
478560
match error {
479561
WorkerError::Remote(error) => !matches!(
@@ -493,7 +575,9 @@ mod tests {
493575
use talon_core::{Backend, NodeId, NodeInfo, NodeRole, ObjectId, Version};
494576
use talon_transport::frame::{FrameHeader, HEADER_LEN};
495577
use talon_transport::{
496-
decode_request, encode_error, response_header_ok, ControlMessage, RangeRequest,
578+
decode_cached_block_put_header, decode_cached_request, decode_request, encode_error,
579+
encode_typed_error, response_header_ok, ControlMessage, DataErrorCode, MsgType,
580+
RangeRequest,
497581
};
498582
use tokio::io::{AsyncReadExt, AsyncWriteExt};
499583
use tokio::net::TcpListener;
@@ -988,4 +1072,84 @@ mod tests {
9881072
assert!(bytes.is_empty());
9891073
assert_eq!(hits.load(std::sync::atomic::Ordering::SeqCst), 0);
9901074
}
1075+
1076+
#[tokio::test]
1077+
async fn cached_block_read_uses_fail_closed_wire_operation_and_typed_miss() {
1078+
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1079+
let worker = listener.local_addr().unwrap().to_string();
1080+
tokio::spawn(async move {
1081+
let (mut socket, _) = listener.accept().await.unwrap();
1082+
let mut header = [0_u8; HEADER_LEN];
1083+
socket.read_exact(&mut header).await.unwrap();
1084+
let frame = FrameHeader::decode(&header).unwrap();
1085+
assert_eq!(frame.msg_type, MsgType::GetCachedRange);
1086+
let mut body = vec![0_u8; frame.length as usize];
1087+
socket.read_exact(&mut body).await.unwrap();
1088+
let mut encoded = header.to_vec();
1089+
encoded.extend_from_slice(&body);
1090+
let (_, request) = decode_cached_request(&encoded).unwrap();
1091+
assert_eq!(request.version, Version::new("v1"));
1092+
assert_eq!((request.offset, request.len), (block().offset + 7, 11));
1093+
socket
1094+
.write_all(&encode_typed_error(
1095+
frame.request_id,
1096+
DataErrorCode::CacheMiss,
1097+
"block is not resident",
1098+
))
1099+
.await
1100+
.unwrap();
1101+
});
1102+
let coordinator = mock_coordinator(worker).await;
1103+
let reader = BlockReader::new(
1104+
CoordinatorClient::new(coordinator),
1105+
Arc::new(PlacementCache::new(10_000)),
1106+
1,
1107+
);
1108+
1109+
let error = reader
1110+
.read_cached_block(&block(), 7, 11, 0)
1111+
.await
1112+
.unwrap_err();
1113+
assert!(matches!(error, CacheReadError::CacheMiss(_)));
1114+
}
1115+
1116+
#[tokio::test]
1117+
async fn block_admission_resolves_primary_and_sends_exact_body() {
1118+
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1119+
let worker = listener.local_addr().unwrap().to_string();
1120+
let expected = block();
1121+
let expected_on_worker = expected.clone();
1122+
tokio::spawn(async move {
1123+
let (mut socket, _) = listener.accept().await.unwrap();
1124+
let mut header = [0_u8; HEADER_LEN];
1125+
socket.read_exact(&mut header).await.unwrap();
1126+
let frame = FrameHeader::decode(&header).unwrap();
1127+
assert_eq!(frame.msg_type, MsgType::AdmitCachedBlock);
1128+
let mut payload = vec![0_u8; frame.length as usize];
1129+
socket.read_exact(&mut payload).await.unwrap();
1130+
let mut encoded = header.to_vec();
1131+
encoded.extend_from_slice(&payload);
1132+
let (_, request) = decode_cached_block_put_header(&encoded).unwrap();
1133+
assert_eq!(request.block, expected_on_worker);
1134+
assert_eq!(request.object_len, expected_on_worker.offset + 5);
1135+
let mut body = vec![0_u8; request.body_len as usize];
1136+
socket.read_exact(&mut body).await.unwrap();
1137+
assert_eq!(body, b"tail!");
1138+
socket
1139+
.write_all(&response_header_ok(frame.request_id, 0))
1140+
.await
1141+
.unwrap();
1142+
});
1143+
let coordinator = mock_coordinator(worker).await;
1144+
let reader = BlockReader::new(
1145+
CoordinatorClient::new(coordinator),
1146+
Arc::new(PlacementCache::new(10_000)),
1147+
1,
1148+
);
1149+
1150+
reader
1151+
.admit_block(&expected, expected.offset + 5, b"tail!", 0)
1152+
.await
1153+
.unwrap();
1154+
}
9911155
}

0 commit comments

Comments
 (0)