pub mod bucket;
use std::{
collections::{self, HashSet},
io,
task::Poll,
};
use crate::{HeaderValue, StatusCode};
use bytes::Bytes;
use futures::{StreamExt, TryStreamExt};
use lazy_static::lazy_static;
use regex::Regex;
use time::OffsetDateTime;
use crate::{header, jetstream::response, Error, Message};
use self::bucket::Status;
use super::{
consumer::DeliverPolicy,
stream::{RawMessage, Republish, StorageType, Stream},
};
fn kv_operation_from_maybe_headers(maybe_headers: Option<&String>) -> Operation {
if let Some(headers) = maybe_headers {
return match headers.as_str() {
KV_OPERATION_DELETE => Operation::Delete,
KV_OPERATION_PURGE => Operation::Purge,
_ => Operation::Put,
};
}
Operation::Put
}
fn kv_operation_from_stream_message(message: &RawMessage) -> Operation {
kv_operation_from_maybe_headers(message.headers.as_ref())
}
lazy_static! {
static ref VALID_BUCKET_RE: Regex = Regex::new(r#"\A[a-zA-Z0-9_-]+\z"#).unwrap();
static ref VALID_KEY_RE: Regex = Regex::new(r#"\A[-/_=\.a-zA-Z0-9]+\z"#).unwrap();
}
pub(crate) const MAX_HISTORY: i64 = 64;
const ALL_KEYS: &str = ">";
const KV_OPERATION: &str = "KV-Operation";
const KV_OPERATION_DELETE: &str = "DEL";
const KV_OPERATION_PURGE: &str = "PURGE";
const KV_OPERATION_PUT: &str = "PUT";
const NATS_ROLLUP: &str = "Nats-Rollup";
const ROLLUP_SUBJECT: &str = "sub";
pub(crate) fn is_valid_bucket_name(bucket_name: &str) -> bool {
VALID_BUCKET_RE.is_match(bucket_name)
}
pub(crate) fn is_valid_key(key: &str) -> bool {
if key.is_empty() || key.starts_with('.') || key.ends_with('.') {
return false;
}
VALID_KEY_RE.is_match(key)
}
#[derive(Debug, Default)]
pub struct Config {
pub bucket: String,
pub description: String,
pub max_value_size: i32,
pub history: i64,
pub max_age: std::time::Duration,
pub max_bytes: i64,
pub storage: StorageType,
pub num_replicas: usize,
pub republish: Option<Republish>,
}
#[derive(Debug, Clone, Copy, Eq, PartialEq)]
pub enum Operation {
Put,
Delete,
Purge,
}
#[derive(Debug, Clone)]
pub struct Store {
pub name: String,
pub stream_name: String,
pub prefix: String,
pub stream: Stream,
}
impl Store {
pub async fn status(&self) -> Result<Status, Error> {
let info = self.stream.info.clone();
Ok(Status {
info,
bucket: self.name.to_string(),
})
}
pub async fn put<T: AsRef<str>>(&self, key: T, value: bytes::Bytes) -> Result<u64, Error> {
if !is_valid_key(key.as_ref()) {
return Err(Box::new(io::Error::new(
io::ErrorKind::InvalidInput,
"invalid key",
)));
}
let subject = format!("{}{}", self.prefix.as_str(), key.as_ref());
let publish_ack = self.stream.context.publish(subject, value).await?;
Ok(publish_ack.sequence)
}
pub async fn entry<T: Into<String>>(&self, key: T) -> Result<Option<Entry>, Error> {
let key: String = key.into();
if !is_valid_key(key.as_ref()) {
return Err(Box::new(io::Error::new(
io::ErrorKind::InvalidInput,
"invalid key",
)));
}
let subject = format!("{}{}", self.prefix.as_str(), &key);
match self
.stream
.get_last_raw_message_by_subject(subject.as_str())
.await
{
Ok(message) => {
let operation = kv_operation_from_stream_message(&message);
let nats_message = Message::try_from(message.clone())?;
if nats_message.status == Some(StatusCode::NO_RESPONDERS) {
return Ok(None);
}
let entry = Entry {
bucket: self.name.clone(),
key,
value: nats_message.payload.to_vec(),
revision: message.sequence,
created: message.time,
operation,
delta: 0,
};
Ok(Some(entry))
}
Err(err) => {
let e: std::io::Error = *err.downcast().unwrap();
let d = e.get_ref().unwrap();
let de = d.downcast_ref::<response::Error>().unwrap();
if de.code == 10037 {
return Ok(None);
}
Err(Box::new(e))
}
}
}
pub async fn watch<T: AsRef<str>>(&self, key: T) -> Result<Watch<'_>, Error> {
let subject = format!("{}{}", self.prefix.as_str(), key.as_ref());
let consumer = self
.stream
.create_consumer(super::consumer::push::OrderedConfig {
deliver_subject: self.stream.context.client.new_inbox(),
description: Some("kv watch consumer".to_string()),
filter_subject: subject,
replay_policy: super::consumer::ReplayPolicy::Instant,
deliver_policy: DeliverPolicy::New,
..Default::default()
})
.await?;
Ok(Watch {
subscription: consumer.messages().await?,
prefix: self.prefix.clone(),
bucket: self.name.clone(),
})
}
pub async fn watch_all(&self) -> Result<Watch<'_>, Error> {
self.watch(ALL_KEYS).await
}
pub async fn get<T: Into<String>>(&self, key: T) -> Result<Option<Vec<u8>>, Error> {
match self.entry(key).await {
Ok(Some(entry)) => match entry.operation {
Operation::Put => Ok(Some(entry.value)),
_ => Ok(None),
},
Ok(None) => Ok(None),
Err(err) => Err(err),
}
}
pub async fn update<T: AsRef<str>>(
&self,
key: T,
value: Bytes,
revision: u64,
) -> Result<u64, Error> {
if !is_valid_key(key.as_ref()) {
return Err(Box::new(io::Error::new(
io::ErrorKind::InvalidInput,
"invalid key",
)));
}
let subject = format!("{}{}", self.prefix.as_str(), key.as_ref());
let mut headers = crate::HeaderMap::default();
headers.insert(
header::NATS_EXPECTED_LAST_SUBJECT_SEQUENCE,
HeaderValue::from(revision),
);
self.stream
.context
.publish_with_headers(subject, headers, value)
.await
.map(|publish_ack| publish_ack.sequence)
}
pub async fn delete<T: AsRef<str>>(&self, key: T) -> Result<(), Error> {
if !is_valid_key(key.as_ref()) {
return Err(Box::new(io::Error::new(
io::ErrorKind::InvalidInput,
"invalid key",
)));
}
let subject = format!("{}{}", self.prefix.as_str(), key.as_ref());
let mut headers = crate::HeaderMap::default();
headers.insert(KV_OPERATION, KV_OPERATION_DELETE.parse::<HeaderValue>()?);
self.stream
.context
.publish_with_headers(subject, headers, "".into())
.await?;
Ok(())
}
pub async fn purge<T: AsRef<str>>(&self, key: T) -> Result<(), Error> {
if !is_valid_key(key.as_ref()) {
return Err(Box::new(io::Error::new(
io::ErrorKind::InvalidInput,
"invalid key",
)));
}
let subject = format!("{}{}", self.prefix.as_str(), key.as_ref());
let mut headers = crate::HeaderMap::default();
headers.insert(KV_OPERATION, HeaderValue::from(KV_OPERATION_PURGE));
headers.insert(NATS_ROLLUP, HeaderValue::from(ROLLUP_SUBJECT));
self.stream
.context
.publish_with_headers(subject, headers, "".into())
.await?;
Ok(())
}
pub async fn history<T: AsRef<str>>(&self, key: T) -> Result<History<'_>, Error> {
if !is_valid_key(key.as_ref()) {
return Err(Box::new(io::Error::new(
io::ErrorKind::InvalidInput,
"invalid key",
)));
}
let subject = format!("{}{}", self.prefix.as_str(), key.as_ref());
let consumer = self
.stream
.create_consumer(super::consumer::push::OrderedConfig {
deliver_subject: self.stream.context.client.new_inbox(),
description: Some("kv history consumer".to_string()),
filter_subject: subject,
replay_policy: super::consumer::ReplayPolicy::Instant,
..Default::default()
})
.await?;
Ok(History {
subscription: consumer.messages().await?,
done: false,
prefix: self.prefix.clone(),
bucket: self.name.clone(),
})
}
pub async fn keys(&self) -> Result<collections::hash_set::IntoIter<String>, Error> {
let subject = format!("{}>", self.prefix.as_str());
let consumer = self
.stream
.create_consumer(super::consumer::push::OrderedConfig {
deliver_subject: self.stream.context.client.new_inbox(),
description: Some("kv history consumer".to_string()),
filter_subject: subject,
headers_only: true,
replay_policy: super::consumer::ReplayPolicy::Instant,
..Default::default()
})
.await?;
let mut entries = History {
done: consumer.info.num_pending == 0,
subscription: consumer.messages().await?,
prefix: self.prefix.clone(),
bucket: self.name.clone(),
};
let mut keys = HashSet::new();
while let Some(entry) = entries.try_next().await? {
keys.insert(entry.key);
}
Ok(keys.into_iter())
}
}
pub struct Watch<'a> {
subscription: super::consumer::push::Ordered<'a>,
prefix: String,
bucket: String,
}
impl<'a> futures::Stream for Watch<'a> {
type Item = Result<Entry, Error>;
fn poll_next(
mut self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
match self.subscription.poll_next_unpin(cx) {
Poll::Ready(message) => match message {
None => Poll::Ready(None),
Some(message) => {
let message = message?;
let info = message.info()?;
let operation = match message
.headers
.as_ref()
.and_then(|headers| headers.get(KV_OPERATION))
.unwrap_or(&HeaderValue::from(KV_OPERATION_PUT))
.iter()
.next()
.unwrap()
.as_str()
{
KV_OPERATION_DELETE => Operation::Delete,
KV_OPERATION_PURGE => Operation::Purge,
_ => Operation::Put,
};
let key = message
.subject
.strip_prefix(&self.prefix)
.map(|s| s.to_string())
.unwrap();
Poll::Ready(Some(Ok(Entry {
bucket: self.bucket.clone(),
key,
value: message.payload.to_vec(),
revision: info.stream_sequence,
created: info.published,
delta: info.pending,
operation,
})))
}
},
std::task::Poll::Pending => Poll::Pending,
}
}
fn size_hint(&self) -> (usize, Option<usize>) {
(0, None)
}
}
pub struct History<'a> {
subscription: super::consumer::push::Ordered<'a>,
done: bool,
prefix: String,
bucket: String,
}
impl<'a> futures::Stream for History<'a> {
type Item = Result<Entry, Error>;
fn poll_next(
mut self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
if self.done {
return Poll::Ready(None);
}
match self.subscription.poll_next_unpin(cx) {
Poll::Ready(message) => match message {
None => Poll::Ready(None),
Some(message) => {
let message = message?;
let info = message.info()?;
if info.pending == 0 {
self.done = true;
}
let operation = match message
.headers
.as_ref()
.and_then(|headers| headers.get(KV_OPERATION))
.unwrap_or(&HeaderValue::from(KV_OPERATION_PUT))
.iter()
.next()
.unwrap()
.as_str()
{
KV_OPERATION_DELETE => Operation::Delete,
KV_OPERATION_PURGE => Operation::Purge,
_ => Operation::Put,
};
let key = message
.subject
.strip_prefix(&self.prefix)
.map(|s| s.to_string())
.unwrap();
Poll::Ready(Some(Ok(Entry {
bucket: self.bucket.clone(),
key,
value: message.payload.to_vec(),
revision: info.stream_sequence,
created: info.published,
delta: info.pending,
operation,
})))
}
},
std::task::Poll::Pending => Poll::Pending,
}
}
fn size_hint(&self) -> (usize, Option<usize>) {
(0, None)
}
}
#[derive(Debug, Clone)]
pub struct Entry {
pub bucket: String,
pub key: String,
pub value: Vec<u8>,
pub revision: u64,
pub delta: u64,
pub created: OffsetDateTime,
pub operation: Operation,
}