use std::collections::VecDeque;
use std::io;
use std::sync::{Arc, Mutex};
use thiserror::Error;
use crate::exceptions::TemplateProcessingException;
use crate::expression::{TemplateObject, TemplateValue};
use crate::util::Utf16String;
use super::{DataDrivenTemplateSignal, IThrottledTemplateWriterControl};
const SSE_HEAD_EVENT_NAME: &[u16] = &[104, 101, 97, 100];
const SSE_MESSAGE_EVENT_NAME: &[u16] = &[109, 101, 115, 115, 97, 103, 101];
const SSE_TAIL_EVENT_NAME: &[u16] = &[116, 97, 105, 108];
#[derive(Clone, Copy, Debug, Eq, Error, PartialEq)]
pub enum DataDrivenTemplateIteratorError {
#[error("java.util.NoSuchElementException")]
NoSuchElement,
#[error("remove() is not supported in Throttled Iterator")]
RemoveUnsupported,
}
pub struct DataDrivenTemplateIterator<T> {
values: VecDeque<T>,
writer_control: Option<Box<dyn IThrottledTemplateWriterControl + Send>>,
sse_events_prefix: Option<Vec<u16>>,
sse_events_id: i64,
in_step: bool,
feeding_complete: bool,
queried: bool,
signal: DataDrivenTemplateSignal,
}
impl<T> DataDrivenTemplateIterator<T> {
#[must_use]
pub fn new() -> Self {
Self {
values: VecDeque::with_capacity(10),
writer_control: None,
sse_events_prefix: None,
sse_events_id: 0,
in_step: false,
feeding_complete: false,
queried: false,
signal: DataDrivenTemplateSignal::new(),
}
}
pub fn set_writer_control(
&mut self,
writer_control: Box<dyn IThrottledTemplateWriterControl + Send>,
) {
self.writer_control = Some(writer_control);
}
pub fn set_sse_events_prefix(&mut self, sse_events_prefix: Option<&Utf16String>) {
self.sse_events_prefix = sse_events_prefix
.filter(|prefix| !prefix.is_empty())
.map(|prefix| prefix.as_utf16().to_vec());
}
pub const fn set_sse_events_first_id(&mut self, sse_events_first_id: i64) {
self.sse_events_id = sse_events_first_id;
}
pub const fn take_back_last_event_id(&mut self) {
if self.sse_events_id > 0 {
self.sse_events_id -= 1;
}
}
pub fn has_next(&mut self) -> bool {
self.queried = true;
!self.values.is_empty()
}
pub fn next_java(&mut self) -> Result<T, DataDrivenTemplateIteratorError> {
self.queried = true;
self.values
.pop_front()
.ok_or(DataDrivenTemplateIteratorError::NoSuchElement)
}
pub const fn remove(&self) -> Result<(), DataDrivenTemplateIteratorError> {
Err(DataDrivenTemplateIteratorError::RemoveUnsupported)
}
pub fn start_iteration(&mut self) {
self.start_step(SSE_MESSAGE_EVENT_NAME);
}
pub fn finish_iteration(&mut self) -> Result<(), TemplateProcessingException> {
self.finish_step()
}
#[must_use]
pub const fn has_been_queried(&self) -> bool {
self.queried
}
pub(crate) fn is_paused(&mut self) -> bool {
self.queried = true;
self.values.is_empty() && !self.feeding_complete
}
#[must_use]
pub fn continue_buffer_execution(&self) -> bool {
!self.values.is_empty()
}
pub fn feed_buffer<I>(&mut self, new_elements: I)
where
I: IntoIterator<Item = T>,
{
self.values.extend(new_elements);
self.signal.notify();
}
pub fn start_head(&mut self) {
self.start_step(SSE_HEAD_EVENT_NAME);
}
pub fn feeding_complete(&mut self) {
self.feeding_complete = true;
self.signal.notify();
}
#[must_use]
pub fn get_signal(&self) -> DataDrivenTemplateSignal {
self.signal.clone()
}
pub fn start_tail(&mut self) {
self.start_step(SSE_TAIL_EVENT_NAME);
}
pub fn finish_step(&mut self) -> Result<(), TemplateProcessingException> {
if !self.in_step {
return Ok(());
}
self.in_step = false;
if let Some(sse_control) = self
.writer_control
.as_deref_mut()
.and_then(IThrottledTemplateWriterControl::as_sse_control)
{
sse_control.end_event().map_err(Self::processing_error)?;
}
Ok(())
}
pub fn is_step_output_finished(&mut self) -> Result<bool, TemplateProcessingException> {
if self.in_step {
return Ok(false);
}
match self.writer_control.as_deref_mut() {
Some(control) => control
.is_overflown()
.map(|overflown| !overflown)
.map_err(Self::processing_error),
None => Ok(true),
}
}
fn start_step(&mut self, event: &[u16]) {
self.in_step = true;
let id_token: Vec<u16> = self.sse_events_id.to_string().encode_utf16().collect();
let id = self.compose_token(&id_token);
let event = self.compose_token(event);
if let Some(sse_control) = self
.writer_control
.as_deref_mut()
.and_then(IThrottledTemplateWriterControl::as_sse_control)
{
sse_control.start_event(Some(&id), Some(&event));
self.sse_events_id = self.sse_events_id.wrapping_add(1);
}
}
fn compose_token(&self, token: &[u16]) -> Vec<u16> {
let Some(prefix) = self.sse_events_prefix.as_deref() else {
return token.to_vec();
};
let mut result = Vec::with_capacity(prefix.len() + 1 + token.len());
result.extend_from_slice(prefix);
result.push(45);
result.extend_from_slice(token);
result
}
fn processing_error(cause: io::Error) -> TemplateProcessingException {
TemplateProcessingException::with_cause(
Some("Cannot signal end of SSE event".to_owned()),
cause,
)
}
}
impl<T> Default for DataDrivenTemplateIterator<T> {
fn default() -> Self {
Self::new()
}
}
impl DataDrivenTemplateIterator<Arc<TemplateValue>> {
#[must_use]
pub fn shared_template_value() -> (Arc<Mutex<Self>>, Arc<TemplateValue>) {
let iterator = Arc::new(Mutex::new(Self::new()));
let object: Arc<dyn TemplateObject> = iterator.clone();
let value = Arc::new(TemplateValue::Object(object));
(iterator, value)
}
#[must_use]
pub fn to_template_value(iterator: &Arc<Mutex<Self>>) -> Arc<TemplateValue> {
let object: Arc<dyn TemplateObject> = iterator.clone();
Arc::new(TemplateValue::Object(object))
}
}
impl<T> Iterator for DataDrivenTemplateIterator<T> {
type Item = T;
fn next(&mut self) -> Option<Self::Item> {
self.queried = true;
self.values.pop_front()
}
}
impl TemplateObject for Mutex<DataDrivenTemplateIterator<Arc<TemplateValue>>> {
fn class_name(&self) -> &str {
"org.thymeleaf.engine.DataDrivenTemplateIterator"
}
fn to_utf16_string(&self) -> Utf16String {
Utf16String::from_rust_str("org.thymeleaf.engine.DataDrivenTemplateIterator")
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}