use std::borrow::Cow;
use std::collections::hash_map::Entry;
use std::collections::{HashMap, VecDeque};
use std::future::Future;
use std::os::fd::{BorrowedFd, RawFd};
use std::pin::Pin;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::task::{Context, Poll, Waker};
use std::time::{Duration, Instant};
use crate::error::ClientError;
use crate::protocol::wal_block::decode_wal_block_into;
use crate::{
encode_ddl_txn, encode_frame, encode_push_txn, sys_schema, ClientTransport, ProtocolError, PushFamily, Schema,
ZSetBatch,
};
use gnitz_wire::control::{peek_control_block, ControlHeader, DecodedControl};
use gnitz_wire::txn_frame::{self, kept_as};
use gnitz_wire::CONNECT_TIMEOUT;
use gnitz_wire::{ClientVerb, WireConflictMode, WireLane};
use gnitz_wire::{RelClass, RelDescriptorBlob, RelIndex};
pub const MAX_IN_FLIGHT: usize = 4096;
pub const MAX_QUEUED_BYTES: usize = 64 << 20;
#[derive(Debug)]
pub struct ScanReply {
pub schema: Arc<Schema>,
pub batch: ZSetBatch,
pub lsn: Option<u64>,
}
#[derive(Debug)]
pub struct RelDescriptor {
pub tid: u64,
pub class: RelClass,
pub pk_repeats: bool,
pub serial: bool,
pub schema: Arc<Schema>,
pub indexes: Vec<RelIndex>,
pub token: u64,
}
pub use gnitz_wire::control::Target;
impl From<&RelDescriptor> for Target {
fn from(rel: &RelDescriptor) -> Self {
Target { tid: rel.tid, token: rel.token }
}
}
#[derive(Debug)]
pub(crate) struct RawBlock {
frame: Vec<u8>,
block: std::ops::Range<usize>,
}
impl RawBlock {
pub(crate) fn block(&self) -> &[u8] {
&self.frame[self.block.clone()]
}
}
pub use gnitz_wire::txn_frame::DeltaCursor;
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub struct Interest {
pub read: bool,
pub write: bool,
}
impl Interest {
pub const NONE: Interest = Interest { read: false, write: false };
pub const READ: Interest = Interest { read: true, write: false };
pub const WRITE: Interest = Interest { read: false, write: true };
pub fn is_empty(self) -> bool {
self == Interest::NONE
}
pub fn poll_events(self) -> libc::c_short {
(if self.read { libc::POLLIN } else { 0 }) | (if self.write { libc::POLLOUT } else { 0 })
}
pub fn from_revents(revents: libc::c_short) -> Interest {
Interest {
read: revents & (libc::POLLIN | libc::POLLHUP | libc::POLLERR) != 0,
write: revents & libc::POLLOUT != 0,
}
}
}
pub enum Request<'a> {
AllocIds(u64),
AllocSerial { table: Target, count: u64 },
DdlTxn(&'a [(u64, ZSetBatch)]),
PushTxn { families: &'a [PushFamily] },
Push {
target: Target,
schema: &'a Schema,
batch: &'a ZSetBatch,
mode: WireConflictMode,
},
}
impl Request<'_> {
fn encode(self) -> Result<(Vec<u8>, u64), ClientError> {
Ok(match self {
Request::AllocIds(n) => {
let hdr = ControlHeader::naming(ClientVerb::AllocIds, Target::from(0), n);
(encode_frame(hdr, &[], None, None), 0)
}
Request::AllocSerial { table, count } => {
let hdr = ControlHeader::naming(ClientVerb::AllocSerialRange, table, count);
(encode_frame(hdr, &[], None, None), table.tid)
}
Request::DdlTxn(families) => {
for (tid, batch) in families {
if gnitz_wire::sys_family_index(*tid).is_none() {
return Err(ClientError::from(format!("DDL family {tid} is not a system table")));
}
batch.validate(sys_schema(*tid))?;
}
(encode_ddl_txn(families), 0)
}
Request::PushTxn { families } => {
for f in families {
f.batch.validate(&f.schema)?;
}
(encode_push_txn(families), 0)
}
Request::Push { target, schema, batch, mode } => {
batch.validate(schema)?;
let mut hdr = ControlHeader::naming(ClientVerb::Push, target, 0);
hdr.flags.conflict_mode = mode;
(
encode_frame(hdr, &[], Some(&schema.to_block()), Some(batch)),
target.tid,
)
}
})
}
}
enum Slot {
Ack { tid: u64, to: Promise<u64> },
Resolve { to: Promise<Option<Arc<RelDescriptor>>> },
Scan {
tid: u64,
reply_schema: Arc<Schema>,
data: Option<ZSetBatch>,
to: Promise<ScanReply>,
},
Delta {
tid: u64,
reply_schema: Arc<Schema>,
data: Option<ZSetBatch>,
keep: Option<u64>,
to: Promise<(ScanReply, DeltaCursor)>,
},
Multi {
rels: Vec<(u64, Arc<Schema>)>,
replies: Vec<ScanReply>,
data: Option<ZSetBatch>,
to: Promise<Vec<ScanReply>>,
},
DeltaPoll {
views: Vec<u64>,
at: usize,
poll: Option<u64>,
first: u64,
blocks: Vec<RawBlock>,
},
Sync { to: Promise<()> },
}
#[derive(Default)]
enum Awaited<T> {
#[default]
Outstanding,
Arrived(Result<T, ClientError>),
Routed(Box<dyn FnOnce(Result<T, ClientError>) + Send>),
Watched(Waker),
}
struct Cell<T>(Mutex<Awaited<T>>);
impl<T> Cell<T> {
fn lock(&self) -> MutexGuard<'_, Awaited<T>> {
self.0.lock().unwrap_or_else(PoisonError::into_inner)
}
}
pub struct Sent<T>(Arc<Cell<T>>);
pub struct Promise<T>(Option<Arc<Cell<T>>>);
pub fn promise<T>() -> (Promise<T>, Sent<T>) {
let cell = Arc::new(Cell(Mutex::new(Awaited::Outstanding)));
(Promise(Some(Arc::clone(&cell))), Sent(cell))
}
impl<T> Promise<T> {
pub fn fulfil(&mut self, reply: Result<T, ClientError>) {
let Some(cell) = self.0.take() else {
return;
};
let mut state = cell.lock();
match std::mem::take(&mut *state) {
Awaited::Routed(route) => {
drop(state);
route(reply)
}
Awaited::Watched(waiter) => {
*state = Awaited::Arrived(reply);
drop(state);
waiter.wake()
}
Awaited::Outstanding | Awaited::Arrived(_) => *state = Awaited::Arrived(reply),
}
}
}
impl<T> Drop for Promise<T> {
fn drop(&mut self) {
self.fulfil(Err(ClientError::Closed));
}
}
impl<T> Sent<T> {
pub fn ready(value: Result<T, ClientError>) -> Self {
Sent(Arc::new(Cell(Mutex::new(Awaited::Arrived(value)))))
}
pub fn try_take(&mut self) -> Option<Result<T, ClientError>> {
let mut state = self.0.lock();
match std::mem::take(&mut *state) {
Awaited::Arrived(reply) => Some(reply),
other => {
*state = other;
None
}
}
}
pub fn then(self, route: impl FnOnce(Result<T, ClientError>) + Send + 'static) {
let mut state = self.0.lock();
match std::mem::take(&mut *state) {
Awaited::Arrived(reply) => {
drop(state);
route(reply)
}
_ => *state = Awaited::Routed(Box::new(route)),
}
}
}
impl<T> Future for Sent<T> {
type Output = Result<T, ClientError>;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let mut state = self.0.lock();
match std::mem::take(&mut *state) {
Awaited::Arrived(reply) => Poll::Ready(reply),
_ => {
*state = Awaited::Watched(cx.waker().clone());
Poll::Pending
}
}
}
}
pub(crate) type PollEnd = Result<DeltaCursor, ClientError>;
pub(crate) enum Polled {
Block(RawBlock),
End(PollEnd),
}
#[derive(Default)]
struct Polls {
live: u64,
queue: VecDeque<Polled>,
subs: HashMap<u64, Parked>,
server_holds: bool,
}
type Parked = Result<(Vec<RawBlock>, DeltaCursor), ClientError>;
enum Taker {
Reader,
Holder,
Nobody,
}
impl Polls {
fn taker(&self, poll: Option<u64>, kept: u64) -> Taker {
if poll == Some(self.live) {
Taker::Reader
} else if self.subs.contains_key(&kept) {
Taker::Holder
} else {
Taker::Nobody
}
}
fn end(&mut self, poll: Option<u64>, kept: u64, blocks: &mut Vec<RawBlock>, end: PollEnd) {
self.server_holds |= end.is_ok();
let train = std::mem::take(blocks);
match self.taker(poll, kept) {
Taker::Reader => {
if let Ok(cursor) = &end {
self.subs.insert(kept, Ok((Vec::new(), *cursor)));
}
self.queue.push_back(Polled::End(end));
}
Taker::Holder => match (self.subs.get_mut(&kept), end) {
(Some(Ok((held, at))), Ok(cursor)) => {
held.extend(train);
*at = cursor;
}
(Some(parked @ Ok(_)), Err(why)) => *parked = Err(why),
_ => {}
},
Taker::Nobody => {}
}
}
fn untaken(&self) -> bool {
let quiet = |parked: &Parked| matches!(parked, Ok((blocks, _)) if blocks.is_empty());
self.subs.values().any(|parked| !quiet(parked))
}
}
pub struct Session {
transport: ClientTransport,
pending: VecDeque<Slot>,
submitted: u64,
ended: Option<ClientError>,
unread: bool,
polls: Polls,
pushed: Option<Slot>,
syncing: VecDeque<Slot>,
}
impl Session {
pub fn connect(target: &str) -> Result<Self, ClientError> {
Ok(Self::over(ClientTransport::connect(
target,
Instant::now() + CONNECT_TIMEOUT,
)?))
}
pub(crate) fn over(transport: ClientTransport) -> Self {
Session {
transport,
pending: VecDeque::new(),
submitted: 0,
ended: None,
unread: false,
polls: Polls::default(),
pushed: None,
syncing: VecDeque::new(),
}
}
pub fn requests_sent(&self) -> u64 {
self.submitted
}
pub fn as_raw_fd(&self) -> RawFd {
self.transport.as_raw_fd()
}
pub fn as_fd(&self) -> BorrowedFd<'_> {
self.transport.as_fd()
}
pub fn submit(&mut self, req: Request<'_>) -> Sent<u64> {
let (frame, tid) = match req.encode() {
Ok(encoded) => encoded,
Err(why) => return Sent::ready(Err(why)),
};
let (to, sent) = promise();
self.enqueue(frame, Slot::Ack { tid, to });
sent
}
pub fn submit_resolve(&mut self, qname: &str) -> Sent<Option<Arc<RelDescriptor>>> {
let hdr = ControlHeader::naming(ClientVerb::Resolve, Target::from(0), 0);
let (to, sent) = promise();
self.enqueue(encode_frame(hdr, qname.as_bytes(), None, None), Slot::Resolve { to });
sent
}
pub fn submit_scan(
&mut self,
target: Target,
spec: &gnitz_wire::ReadSpec,
reply_schema: &Arc<Schema>,
) -> Sent<ScanReply> {
let hdr = ControlHeader::naming(ClientVerb::ScanSpec, target, reply_schema.layout().layout_digest());
let (to, sent) = promise();
let slot = Slot::Scan {
tid: target.tid,
reply_schema: Arc::clone(reply_schema),
data: None,
to,
};
self.enqueue(encode_frame(hdr, &spec.encode(), None, None), slot);
sent
}
pub fn submit_scan_multi(
&mut self,
scans: &[(Target, &gnitz_wire::ReadSpec, &Arc<Schema>)],
) -> Sent<Vec<ScanReply>> {
if scans.is_empty() {
return Sent::ready(Err(ClientError::from("a multi-read names no relation".to_string())));
}
let specs: Vec<Vec<u8>> = scans.iter().map(|(_, spec, _)| spec.encode()).collect();
let items: Vec<txn_frame::ScanItem> = scans
.iter()
.zip(&specs)
.map(|(&(target, _, schema), spec)| txn_frame::ScanItem {
target,
reply_layout: schema.layout().layout_digest(),
spec,
})
.collect();
let (to, sent) = promise();
let slot = Slot::Multi {
replies: Vec::with_capacity(scans.len()),
rels: scans
.iter()
.map(|&(target, _, schema)| (target.tid, Arc::clone(schema)))
.collect(),
data: None,
to,
};
self.enqueue(txn_frame::encode_scan_multi(&items), slot);
sent
}
pub fn submit_delta_read(
&mut self,
view: Target,
from: Option<DeltaCursor>,
reply_schema: &Arc<Schema>,
spec: &[u8],
keep: Option<u64>,
) -> Sent<(ScanReply, DeltaCursor)> {
let (to, sent) = promise();
let slot = Slot::Delta {
tid: view.tid,
reply_schema: Arc::clone(reply_schema),
data: None,
keep,
to,
};
let item = txn_frame::DeltaPollItem {
view,
from,
reply_layout: reply_schema.layout().layout_digest(),
spec,
};
self.enqueue(txn_frame::encode_delta_poll(&[item], keep), slot);
sent
}
pub(crate) fn submit_sync(&mut self, held: &[u64], wait: Duration) -> Sent<()> {
self.polls.subs.retain(|id, _| held.contains(id));
if !self.polls.server_holds {
return Sent::ready(Ok(()));
}
let wait = if self.polls.untaken() || !self.syncing.is_empty() {
Duration::ZERO
} else {
wait
};
let hdr = ControlHeader::naming(ClientVerb::SyncPushed, Target::from(0), wait_ms(wait));
let named = txn_frame::encode_held(self.polls.subs.keys().copied());
let (to, sent) = promise();
if self.enqueue(encode_frame(hdr, &named, None, None), Slot::Sync { to }) {
self.polls.server_holds = !self.polls.subs.is_empty();
}
sent
}
pub(crate) fn take_pushed(&mut self, id: u64) -> Result<(Vec<RawBlock>, DeltaCursor), ClientError> {
let Entry::Occupied(mut parked) = self.polls.subs.entry(id) else {
return Err(ClientError::from(format!(
"subscription {id} is not held on this connection; read from its cursor and subscribe again"
)));
};
match parked.get_mut() {
Ok((blocks, cursor)) => Ok((std::mem::take(blocks), *cursor)),
Err(_) => parked.remove(),
}
}
pub(crate) fn submit_delta_poll(&mut self, first: u64, views: &[txn_frame::DeltaPollItem]) {
if views.is_empty() {
return;
}
let slot = Slot::DeltaPoll {
views: views.iter().map(|v| v.view.tid).collect(),
at: 0,
poll: Some(self.polls.live),
first,
blocks: Vec::new(),
};
self.enqueue(txn_frame::encode_delta_poll(views, Some(first)), slot);
}
fn enqueue(&mut self, frame: Vec<u8>, mut slot: Slot) -> bool {
let (total, limit) = (frame.len(), gnitz_wire::MAX_FRAME_PAYLOAD);
let refusal = if total > limit {
Some(ClientError::from(format!(
"request frame is {total} bytes, exceeding the {limit}-byte server ingress cap; \
split the request"
)))
} else if let Some(why) = &self.ended {
Some(why.clone())
} else {
self.at_capacity().then(|| {
let queued = self.queued_bytes();
ClientError::from(if self.in_flight() >= MAX_IN_FLIGHT {
format!("connection has {MAX_IN_FLIGHT} requests in flight")
} else {
format!("connection has {queued} unwritten bytes queued, at the {MAX_QUEUED_BYTES}-byte cap")
})
})
};
if let Some(why) = refusal {
slot.fail(why, &mut self.polls);
return false;
}
self.transport.enqueue(frame);
self.submitted += 1;
match slot {
Slot::Sync { .. } => self.syncing.push_back(slot),
_ => self.pending.push_back(slot),
}
true
}
fn in_flight(&self) -> usize {
self.pending.len() + self.syncing.len()
}
pub fn step(&mut self, ready: Interest) {
if ready.read {
self.read();
}
if ready.write && self.ended.is_none() {
if let Err(e) = self.transport.flush() {
self.read();
while self.unread {
self.read();
}
self.end(ClientError::ConnectionLost(e));
}
}
}
fn read(&mut self) {
self.unread = false;
if self.ended.is_some() {
return;
}
let Session {
transport,
pending,
polls,
pushed,
syncing,
..
} = self;
match transport.read(|buf| feed(pending, pushed, syncing, buf, polls)) {
Ok(more) => self.unread = more,
Err(e) => self.end(ClientError::ConnectionLost(e)),
}
}
pub(crate) fn next_polled(&mut self) -> Option<Polled> {
self.polls.queue.pop_front()
}
pub(crate) fn polled_out(&self) -> bool {
self.polls.queue.is_empty()
}
pub(crate) fn abandon_poll(&mut self) {
self.polls.live += 1;
self.polls.queue.clear();
}
pub(crate) fn unread(&self) -> bool {
self.unread
}
fn queued_bytes(&self) -> usize {
self.transport.queued_bytes()
}
pub fn at_capacity(&self) -> bool {
self.in_flight() >= MAX_IN_FLIGHT || self.queued_bytes() >= MAX_QUEUED_BYTES
}
pub fn interest(&self) -> Interest {
if self.ended.is_some() {
return Interest::NONE;
}
Interest {
read: self.in_flight() > 0,
write: self.transport.wants_write(),
}
}
pub fn is_closed(&self) -> bool {
self.ended.is_some()
}
pub(crate) fn end(&mut self, why: ClientError) {
if self.ended.is_some() {
return;
}
self.transport.close();
let Session { pending, syncing, polls, .. } = self;
for mut slot in pending.drain(..).chain(syncing.drain(..)) {
slot.fail(why.clone(), polls);
}
self.ended = Some(why);
}
}
fn wait_ms(wait: Duration) -> u64 {
u64::try_from(wait.as_nanos().div_ceil(1_000_000)).unwrap_or(u64::MAX)
}
fn feed(
pending: &mut VecDeque<Slot>,
pushed: &mut Option<Slot>,
syncing: &mut VecDeque<Slot>,
buf: Cow<'_, [u8]>,
polls: &mut Polls,
) -> Result<(), ProtocolError> {
let ctrl = peek_control_block(&buf).map_err(ProtocolError::DecodeError)?;
let in_train = pushed.is_some();
let queue = match ctrl.hdr.flags.lane {
WireLane::Reply => pending,
_ if in_train => {
return Err(ProtocolError::DecodeError(
"a pushed train split by a frame that is not of it".into(),
))
}
WireLane::PushedTrain => {
*pushed = Some(Slot::DeltaPoll {
views: vec![ctrl.hdr.target_id],
at: 0,
poll: None,
first: ctrl.hdr.arg0,
blocks: Vec::new(),
});
return Ok(());
}
WireLane::SyncAnswer => syncing,
};
let Some(slot) = pushed.as_mut().or(queue.front_mut()) else {
return Err(ProtocolError::DecodeError("reply frame with no request pending".into()));
};
if slot.feed(ctrl, buf, polls)? {
match in_train {
true => *pushed = None,
false => drop(queue.pop_front()),
}
}
Ok(())
}
impl Slot {
fn feed(&mut self, ctrl: DecodedControl, mut buf: Cow<'_, [u8]>, polls: &mut Polls) -> Result<bool, ProtocolError> {
let named = ctrl.hdr.target_id;
let want = match &*self {
Slot::Ack { tid, .. } | Slot::Scan { tid, .. } | Slot::Delta { tid, .. } => Some(*tid),
Slot::Sync { .. } => Some(0),
Slot::Multi { rels, replies, .. } => Some(rels[replies.len()].0),
Slot::DeltaPoll { views, at, .. } => Some(views[*at]),
Slot::Resolve { .. } => None,
};
if let Some(fault) = ctrl.fault(&buf) {
let refused = ClientError::Refused(fault);
return match (&mut *self, want) {
(Slot::DeltaPoll { views, at, poll, first, blocks }, Some(view)) if named != 0 => {
if named != view {
return Err(out_of_order(view, named));
}
polls.end(*poll, kept_as(*first, *at), blocks, Err(refused));
*at += 1;
Ok(*at == views.len())
}
_ => {
self.fail(refused, polls);
Ok(true)
}
};
}
if let Some(want) = want.filter(|&want| want != named) {
return Err(out_of_order(want, named));
}
if ctrl.schema.is_some() && !matches!(self, Slot::Resolve { .. }) {
return Err(ProtocolError::DecodeError(
"a schema block on a reply whose request named its schema".into(),
));
}
let decode = |data: &mut Option<ZSetBatch>, schema: &Arc<Schema>, block: &[u8]| {
decode_wal_block_into(data.get_or_insert_with(|| ZSetBatch::new(schema)), block, schema)
};
match (&mut *self, ctrl.data.clone()) {
(_, None) => {}
(Slot::DeltaPoll { poll, first, at, blocks, .. }, Some(block)) => {
let mut raw = || RawBlock {
frame: std::mem::take(&mut buf).into_owned(),
block: block.clone(),
};
match polls.taker(*poll, kept_as(*first, *at)) {
Taker::Reader => polls.queue.push_back(Polled::Block(raw())),
Taker::Holder => blocks.push(raw()),
Taker::Nobody => {}
}
}
(Slot::Scan { reply_schema, data, .. } | Slot::Delta { reply_schema, data, .. }, Some(r)) => {
decode(data, reply_schema, &buf[r])?
}
(Slot::Multi { rels, replies, data, .. }, Some(r)) => decode(data, &rels[replies.len()].1, &buf[r])?,
(Slot::Ack { .. } | Slot::Resolve { .. } | Slot::Sync { .. }, Some(_)) => {
return Err(ProtocolError::DecodeError(
"a data block on a reply that carries no rows".into(),
))
}
}
if ctrl.hdr.flags.continuation {
return Ok(false);
}
let scan_reply = |schema: &Arc<Schema>, data: &mut Option<ZSetBatch>| ScanReply {
batch: data.take().unwrap_or_else(|| ZSetBatch::new(schema)),
schema: Arc::clone(schema),
lsn: Some(ctrl.hdr.arg0),
};
let cursor = || {
DeltaCursor::from_pair(ctrl.hdr.arg1, ctrl.hdr.arg0)
.ok_or_else(|| ProtocolError::DecodeError("a delta-poll terminal at round 0".into()))
};
match self {
Slot::Ack { to, .. } => to.fulfil(Ok(ctrl.hdr.arg0)),
Slot::Resolve { to } => to.fulfil(Ok(resolve_descriptor(&ctrl, &buf)?)),
Slot::Scan { reply_schema, data, to, .. } => to.fulfil(Ok(scan_reply(reply_schema, data))),
Slot::Delta { reply_schema, data, keep, to, .. } => {
let cursor = cursor()?;
if let Some(id) = *keep {
polls.server_holds = true;
polls.subs.insert(id, Ok((Vec::new(), cursor)));
}
let reply = ScanReply {
lsn: None,
..scan_reply(reply_schema, data)
};
to.fulfil(Ok((reply, cursor)))
}
Slot::Multi { rels, replies, data, to } => {
let reply = scan_reply(&rels[replies.len()].1, data);
replies.push(reply);
if replies.len() < rels.len() {
return Ok(false);
}
to.fulfil(Ok(std::mem::take(replies)))
}
Slot::DeltaPoll { views, at, poll, first, blocks } => {
polls.end(*poll, kept_as(*first, *at), blocks, Ok(cursor()?));
*at += 1;
if *at < views.len() {
return Ok(false);
}
}
Slot::Sync { to } => to.fulfil(Ok(())),
}
Ok(true)
}
fn fail(&mut self, why: ClientError, polls: &mut Polls) {
match self {
Slot::Ack { to, .. } => to.fulfil(Err(why)),
Slot::Resolve { to } => to.fulfil(Err(why)),
Slot::Scan { to, .. } => to.fulfil(Err(why)),
Slot::Delta { to, .. } => to.fulfil(Err(why)),
Slot::Multi { to, .. } => to.fulfil(Err(why)),
Slot::DeltaPoll { views, at, poll, first, blocks } => {
for at in *at..views.len() {
polls.end(*poll, kept_as(*first, at), blocks, Err(why.clone()));
}
}
Slot::Sync { to } => to.fulfil(Err(why)),
}
}
}
fn out_of_order(want: u64, got: u64) -> ProtocolError {
ProtocolError::DecodeError(format!("reply out of order: expected target {want}, got {got}"))
}
fn resolve_descriptor(ctrl: &DecodedControl, frame: &[u8]) -> Result<Option<Arc<RelDescriptor>>, ProtocolError> {
if ctrl.hdr.target_id == 0 {
return Ok(None);
}
let schema = ctrl
.schema
.clone()
.ok_or_else(|| ProtocolError::DecodeError("RESOLVE reply carries no schema".into()))?;
let schema = Arc::new(Schema::from_block(&frame[schema]).map_err(ProtocolError::DecodeError)?);
let desc = RelDescriptorBlob::decode(&frame[ctrl.blob.clone()]).map_err(ProtocolError::DecodeError)?;
Ok(Some(Arc::new(RelDescriptor {
tid: ctrl.hdr.target_id,
class: desc.class,
pk_repeats: desc.pk_repeats,
serial: desc.serial,
schema,
indexes: desc.indexes,
token: ctrl.hdr.arg0,
})))
}
#[cfg(test)]
#[path = "tests/connection.rs"]
mod tests;
#[cfg(test)]
#[path = "benches/connection.rs"]
mod bench;