use crate::concurrent::atomic_buffer::AtomicBuffer;
use crate::concurrent::logbuffer::data_frame_header::{self, DataFrameHeaderDefn};
use crate::concurrent::logbuffer::{frame_descriptor, log_buffer_descriptor};
use crate::utils::bit_utils::{align, number_of_trailing_zeroes};
use crate::utils::types::Index;
#[derive(Clone)]
pub struct Header {
buffer: Option<AtomicBuffer>,
offset: Index,
initial_term_id: i32,
position_bits_to_shift: i32,
}
impl Header {
pub fn new(initial_term_id: i32, capacity: Index) -> Self {
Self {
initial_term_id,
offset: 0,
position_bits_to_shift: number_of_trailing_zeroes(capacity),
buffer: None,
}
}
pub fn initial_term_id(&self) -> i32 {
self.initial_term_id
}
pub fn set_initial_term_id(&mut self, initial_term_id: i32) {
self.initial_term_id = initial_term_id;
}
pub fn offset(&self) -> Index {
self.offset
}
pub fn set_offset(&mut self, offset: Index) {
self.offset = offset;
}
pub fn buffer(&self) -> AtomicBuffer {
self.buffer.expect("Buffer not set")
}
pub fn set_buffer(&mut self, buffer: AtomicBuffer) {
self.buffer = Some(buffer);
}
pub fn frame_length(&self) -> Index {
self.buffer.expect("Buffer not set").get::<i32>(self.offset) as Index
}
pub fn session_id(&self) -> i32 {
self.buffer
.expect("Buffer not set")
.get::<i32>(self.offset + *data_frame_header::SESSION_ID_FIELD_OFFSET)
}
pub fn stream_id(&self) -> i32 {
self.buffer
.expect("Buffer not set")
.get::<i32>(self.offset + *data_frame_header::STREAM_ID_FIELD_OFFSET)
}
pub fn term_id(&self) -> i32 {
self.buffer
.expect("Buffer not set")
.get::<i32>(self.offset + *data_frame_header::TERM_ID_FIELD_OFFSET)
}
pub fn term_offset(&self) -> Index {
self.offset
}
pub fn frame_type(&self) -> u16 {
self.buffer
.expect("Buffer not set")
.get::<u16>(self.offset + *data_frame_header::TYPE_FIELD_OFFSET)
}
pub fn flags(&self) -> u8 {
self.buffer
.expect("Buffer not set")
.get::<u8>(self.offset + *data_frame_header::FLAGS_FIELD_OFFSET)
}
pub fn position(&self) -> i64 {
let resulting_offset = align(self.term_offset() + self.frame_length(), frame_descriptor::FRAME_ALIGNMENT);
log_buffer_descriptor::compute_position(
self.term_id(),
resulting_offset,
self.position_bits_to_shift,
self.initial_term_id,
)
}
pub fn reserved_value(&self) -> i64 {
self.buffer
.expect("Buffer not set")
.get::<i64>(self.offset + *data_frame_header::RESERVED_VALUE_FIELD_OFFSET)
}
}
pub struct HeaderWriter {
session_id: i32,
stream_id: i32,
}
impl HeaderWriter {
pub fn new(default_hdr: AtomicBuffer) -> Self {
Self {
session_id: default_hdr.get::<i32>(*data_frame_header::SESSION_ID_FIELD_OFFSET),
stream_id: default_hdr.get::<i32>(*data_frame_header::STREAM_ID_FIELD_OFFSET),
}
}
pub fn write(&self, term_buffer: &AtomicBuffer, offset: Index, length: Index, term_id: i32) {
term_buffer.put_ordered::<i32>(offset, -(length));
unsafe {
let hdr = term_buffer.overlay_struct::<DataFrameHeaderDefn>(offset);
(*hdr).version = data_frame_header::CURRENT_VERSION;
(*hdr).flags = frame_descriptor::BEGIN_FRAG | frame_descriptor::END_FRAG;
(*hdr).frame_type = data_frame_header::HDR_TYPE_DATA;
(*hdr).term_offset = offset;
(*hdr).session_id = self.session_id;
(*hdr).stream_id = self.stream_id;
(*hdr).term_id = term_id;
}
}
}