use std::io::{Read, Write};
use std::os::unix::net::UnixStream;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use gnitz_wire::{ColumnDef, TypeCode, WireStatus};
use crate::protocol::error::ProtocolError;
use crate::protocol::transport::{poll_fd, ClientTransport};
use crate::{BatchAppender, RelName, Schema, Session, ZSetBatch};
pub(crate) fn rel(schema: &str, name: &str) -> RelName {
RelName::new(schema, name).unwrap()
}
pub(crate) fn kv_schema(v: TypeCode) -> Schema {
Schema {
columns: vec![
ColumnDef::new("pk", TypeCode::U64, false),
ColumnDef::new("v", v, false),
],
pk_cols: vec![0],
}
}
pub(crate) fn kv_rows(rows: &[(u64, i64, i64)]) -> ZSetBatch {
let schema = kv_schema(TypeCode::I64);
let mut b = ZSetBatch::new(&schema);
let mut app = BatchAppender::new(&mut b);
for &(pk, v, w) in rows {
app.add_row(pk as u128, w).i64_val(v);
}
b
}
pub(crate) struct Peer(pub(crate) UnixStream);
impl Peer {
pub(crate) fn send(&self, payload: &[u8]) {
write_frame(&self.0, payload);
}
pub(crate) fn send_bytes(&self, bytes: &[u8]) {
(&self.0).write_all(bytes).unwrap();
}
pub(crate) fn recv(&self) -> Vec<u8> {
read_frame(&self.0)
}
}
pub(crate) fn transport_pair() -> (ClientTransport, Peer) {
let (a, b) = UnixStream::pair().unwrap();
(ClientTransport::unix(a).unwrap(), Peer(b))
}
pub(crate) fn pass(t: &mut ClientTransport) -> (Vec<Vec<u8>>, Result<(), ProtocolError>) {
let mut frames = Vec::new();
loop {
let read = t.read(|f| {
frames.push(f.into_owned());
Ok(())
});
match read {
Ok(true) => {}
Ok(false) => return (frames, Ok(())),
Err(e) => return (frames, Err(e)),
}
}
}
pub(crate) fn recv_frames(t: &mut ClientTransport, n: usize) -> Vec<Vec<u8>> {
let mut frames = Vec::new();
while frames.len() < n {
poll_fd(t.as_raw_fd(), libc::POLLIN, None).unwrap();
let (got, end) = pass(t);
end.unwrap();
frames.extend(got);
}
frames
}
pub(crate) fn flush_all(t: &mut ClientTransport) {
loop {
t.flush().unwrap();
if !t.wants_write() {
return;
}
poll_fd(t.as_raw_fd(), libc::POLLOUT, None).unwrap();
}
}
pub(crate) fn session_pair() -> (Session, Peer) {
let (t, peer) = transport_pair();
(Session::over(t), peer)
}
pub(crate) fn reply_ctrl(tid: u64, lsn: u64) -> Vec<u8> {
let hdr = gnitz_wire::control::ControlHeader {
target_id: tid,
arg0: lsn,
..Default::default()
};
crate::encode_frame(hdr, &[], None, None)
}
pub(crate) fn reply_status(tid: u64, status: WireStatus, text: &str) -> Vec<u8> {
let hdr = gnitz_wire::control::ControlHeader {
target_id: tid,
status,
..Default::default()
};
crate::encode_frame(hdr, text.as_bytes(), None, None)
}
pub(crate) fn interrupt_self_until(stop: Arc<AtomicBool>) -> std::thread::JoinHandle<()> {
extern "C" fn noop(_: libc::c_int) {}
unsafe {
let mut sa: libc::sigaction = std::mem::zeroed();
sa.sa_sigaction = noop as extern "C" fn(libc::c_int) as usize;
libc::sigemptyset(&mut sa.sa_mask);
libc::sigaction(libc::SIGUSR1, &sa, std::ptr::null_mut());
}
let me = unsafe { libc::pthread_self() };
std::thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
std::thread::sleep(std::time::Duration::from_millis(1));
unsafe { libc::pthread_kill(me, libc::SIGUSR1) };
}
})
}
pub(crate) fn framed(payload: &[u8]) -> Vec<u8> {
let mut v = gnitz_wire::frame_len_prefix(payload.len()).to_vec();
v.extend_from_slice(payload);
v
}
pub(crate) fn io_kind<T>(r: &Result<T, ProtocolError>) -> Option<std::io::ErrorKind> {
match r {
Err(ProtocolError::IoError(e)) => Some(e.kind()),
_ => None,
}
}
pub(crate) fn write_frame(mut w: impl Write, payload: &[u8]) {
w.write_all(&framed(payload)).unwrap();
w.flush().unwrap();
}
pub(crate) fn read_frame(mut r: impl Read) -> Vec<u8> {
let mut hdr = [0u8; gnitz_wire::FRAME_LEN_PREFIX_BYTES];
r.read_exact(&mut hdr).unwrap();
let mut payload = vec![0u8; u32::from_le_bytes(hdr) as usize];
r.read_exact(&mut payload).unwrap();
payload
}
pub(crate) fn encode_wal_block(batch: &ZSetBatch) -> Vec<u8> {
let mut out = Vec::new();
gnitz_wire::wal::append_block(&batch.wire_regions(), 0, &mut out);
out
}
pub(crate) fn decode_wal_block(data: &[u8], schema: &Schema) -> Result<ZSetBatch, ProtocolError> {
let mut sink = ZSetBatch::new(schema);
crate::protocol::wal_block::decode_wal_block_into(&mut sink, data, schema)?;
Ok(sink)
}