use crate::client::ReadyWriteBatch;
use crate::proto::{PbPutKvReqForBucket, PutKvResponse};
use crate::rpc::api_key::ApiKey;
use crate::rpc::frame::ReadError;
use crate::rpc::frame::WriteError;
use crate::rpc::message::{ReadType, RequestBody, WriteType};
use crate::{TableId, impl_read_type, impl_write_type, proto};
use bytes::{Buf, BufMut};
use prost::Message;
#[allow(dead_code)]
pub struct PutKvRequest {
pub(crate) inner_request: proto::PutKvRequest,
}
#[allow(dead_code)]
impl PutKvRequest {
pub fn new(
table_id: TableId,
ack: i16,
max_request_timeout_ms: i32,
target_columns: Vec<i32>,
ready_batches: &mut [ReadyWriteBatch],
) -> crate::error::Result<Self> {
let mut request = proto::PutKvRequest {
table_id,
acks: ack as i32,
timeout_ms: max_request_timeout_ms,
target_columns,
..Default::default()
};
for ready_batch in ready_batches {
request.buckets_req.push(PbPutKvReqForBucket {
partition_id: ready_batch.table_bucket.partition_id(),
bucket_id: ready_batch.table_bucket.bucket_id(),
records: ready_batch.write_batch.build()?,
original_partition_name: None,
routing_bucket_count: None,
})
}
Ok(PutKvRequest {
inner_request: request,
})
}
}
impl RequestBody for PutKvRequest {
type ResponseBody = PutKvResponse;
const API_KEY: ApiKey = ApiKey::PutKv;
}
impl_write_type!(PutKvRequest);
impl_read_type!(PutKvResponse);