use iced::Task;
use std::future::Future;
use std::pin::Pin;
use std::time::Duration;
use tokio::sync::broadcast;
use super::editor_widget::{CursorState, EditorAction, EditorBuffer};
pub(crate) const MAX_INPUT_CHARS: usize = 100_000;
#[derive(Debug, Clone)]
pub(crate) struct PaginationState {
pub(crate) page: usize,
pub(crate) page_size: usize,
pub(crate) total: usize,
}
impl PaginationState {
pub(crate) const fn new(page_size: usize) -> Self {
Self {
page: 0,
page_size,
total: 0,
}
}
pub(crate) const fn total_pages(&self) -> usize {
if self.total == 0 {
0
} else {
self.total.div_ceil(self.page_size)
}
}
pub(crate) fn prev_page(&mut self) -> bool {
if self.page > 0 {
self.page -= 1;
true
} else {
false
}
}
pub(crate) fn next_page(&mut self) -> bool {
if self.page + 1 < self.total_pages() {
self.page += 1;
true
} else {
false
}
}
pub(crate) fn reset(&mut self) {
self.page = 0;
}
pub(crate) fn clamp_page(&mut self, total: usize) -> bool {
let total_pages = total.div_ceil(self.page_size);
if total_pages == 0 {
self.page = 0;
false
} else if self.page >= total_pages {
self.page = total_pages - 1;
true
} else {
false
}
}
pub(crate) fn offset(&self) -> usize {
self.page * self.page_size
}
}
#[derive(Debug, Clone)]
pub(crate) struct AsyncLoadState {
loading: bool,
has_loaded: bool,
error: Option<String>,
}
impl AsyncLoadState {
pub(crate) const fn new() -> Self {
Self {
loading: false,
has_loaded: false,
error: None,
}
}
pub(crate) fn loading(&self) -> bool {
self.loading
}
pub(crate) fn has_loaded(&self) -> bool {
self.has_loaded
}
pub(crate) fn error(&self) -> Option<&str> {
self.error.as_deref()
}
pub(crate) fn start_loading(&mut self) {
self.loading = true;
self.error = None;
}
pub(crate) fn finish_loading(&mut self) {
self.loading = false;
self.has_loaded = true;
}
pub(crate) fn fail(&mut self, error: String) {
self.error = Some(error);
self.loading = false;
}
pub(crate) fn clear_error(&mut self) {
self.error = None;
}
pub(crate) fn set_has_loaded(&mut self) {
self.has_loaded = true;
}
}
pub(crate) const POLLED_LIST_REFRESH_INTERVAL: Duration = Duration::from_secs(1);
pub(crate) struct PolledList<T> {
entries: Vec<T>,
in_flight: bool,
loaded: bool,
error: Option<String>,
}
impl<T> PolledList<T> {
pub(crate) const fn new() -> Self {
Self {
entries: Vec::new(),
in_flight: false,
loaded: false,
error: None,
}
}
pub(crate) fn begin(&mut self) -> bool {
if self.in_flight {
return false;
}
self.in_flight = true;
true
}
pub(crate) fn settle(&mut self, result: Result<Vec<T>, String>) {
self.in_flight = false;
self.loaded = true;
match result {
Ok(entries) => {
self.entries = entries;
self.error = None;
}
Err(error) => self.error = Some(error),
}
}
pub(crate) fn entries(&self) -> &[T] {
&self.entries
}
pub(crate) const fn loaded(&self) -> bool {
self.loaded
}
pub(crate) fn error(&self) -> Option<&str> {
self.error.as_deref()
}
}
pub(crate) struct PaginatedTabState<T> {
pub(crate) entries: Vec<T>,
pub(crate) load_state: AsyncLoadState,
pub(crate) pagination: PaginationState,
pub(crate) search: String,
pub(crate) refresh_generation: u64,
}
impl<T> PaginatedTabState<T> {
pub(crate) fn new(page_size: usize) -> Self {
Self {
entries: Vec::new(),
load_state: AsyncLoadState::new(),
pagination: PaginationState::new(page_size),
search: String::new(),
refresh_generation: 0,
}
}
pub(crate) fn begin_refresh(&mut self) -> u64 {
self.load_state.start_loading();
self.refresh_generation = self.refresh_generation.wrapping_add(1);
self.refresh_generation
}
pub(crate) fn handle_refreshed(
&mut self,
generation: u64,
entries: Vec<T>,
total: usize,
) -> bool {
if generation != self.refresh_generation {
return false;
}
if self.pagination.clamp_page(total) {
self.pagination.total = total;
return true;
}
self.entries = entries;
self.pagination.total = total;
self.load_state.finish_loading();
false
}
pub(crate) fn handle_refresh_error(
&mut self,
generation: u64,
e: String,
set_has_loaded_on_error: bool,
) {
if generation != self.refresh_generation {
return;
}
self.load_state.fail(e);
if set_has_loaded_on_error {
self.load_state.set_has_loaded();
}
}
pub(crate) fn clear_entries(&mut self) {
self.entries.clear();
self.pagination.reset();
}
}
#[derive(Debug, Clone)]
pub(crate) struct DebounceState {
generation: u64,
pending: bool,
}
impl DebounceState {
pub(crate) const fn new() -> Self {
Self {
generation: 0,
pending: false,
}
}
pub(crate) fn trigger(&mut self, ms: u64) -> Task<u64> {
self.generation = self.generation.wrapping_add(1);
self.pending = true;
let current = self.generation;
Task::perform(
super::widgets::debounce_sleep(ms, current),
std::convert::identity,
)
}
#[must_use]
pub(crate) fn should_process(&mut self, generation: u64) -> bool {
if generation == self.generation && self.pending {
self.pending = false;
true
} else {
false
}
}
}
pub(crate) trait UndoableText {
fn text(&self) -> String;
fn cursor(&self) -> CursorState;
fn set_text(&mut self, text: &str);
fn move_to(&mut self, line: usize, col: usize);
}
#[derive(Debug, Clone)]
pub(crate) struct UndoStack {
undo: Vec<UndoSnapshot>,
redo: Vec<UndoSnapshot>,
}
#[derive(Debug, Clone)]
pub(crate) struct UndoSnapshot {
pub(crate) text: String,
pub(crate) cursor: CursorState,
}
impl UndoStack {
const MAX_UNDO_DEPTH: usize = 100;
const LARGE_FILE_UNDO_THRESHOLD: usize = 100_000;
pub(crate) const fn new() -> Self {
Self {
undo: Vec::new(),
redo: Vec::new(),
}
}
pub(crate) fn snap_before_edit(&mut self, content: &impl UndoableText) {
let text = content.text();
let max_depth = if text.len() > Self::LARGE_FILE_UNDO_THRESHOLD {
Self::MAX_UNDO_DEPTH / 2
} else {
Self::MAX_UNDO_DEPTH
};
self.redo.clear();
self.undo.push(UndoSnapshot {
text,
cursor: content.cursor(),
});
if self.undo.len() > max_depth {
self.undo.remove(0);
}
}
fn push_and_pop(
dst: &mut Vec<UndoSnapshot>,
src: &mut Vec<UndoSnapshot>,
content: &impl UndoableText,
) -> Option<UndoSnapshot> {
dst.push(UndoSnapshot {
text: content.text(),
cursor: content.cursor(),
});
src.pop()
}
pub(crate) fn undo(&mut self, content: &impl UndoableText) -> Option<UndoSnapshot> {
Self::push_and_pop(&mut self.redo, &mut self.undo, content)
}
pub(crate) fn redo(&mut self, content: &impl UndoableText) -> Option<UndoSnapshot> {
Self::push_and_pop(&mut self.undo, &mut self.redo, content)
}
pub(crate) fn clear(&mut self) {
self.undo.clear();
self.redo.clear();
}
}
pub(crate) fn broadcast_stream_producer<Msg, T, E>(
capacity: usize,
source: &'static std::sync::OnceLock<tokio::sync::broadcast::Sender<T>>,
mut emit: E,
) -> impl futures_util::Stream<Item = Msg>
where
Msg: Send + 'static,
T: Clone + Send + 'static,
E: FnMut(
&mut iced::futures::channel::mpsc::Sender<Msg>,
Option<T>,
) -> Pin<Box<dyn Future<Output = ()> + Send + '_>>
+ Send
+ 'static,
{
iced::stream::channel(
capacity,
move |mut output: iced::futures::channel::mpsc::Sender<Msg>| async move {
let Some(mut rx) = source.get().and_then(|tx| {
if tx.receiver_count() > 100 {
None
} else {
Some(tx.subscribe())
}
}) else {
return;
};
loop {
match rx.recv().await {
Ok(event) => emit(&mut output, Some(event)).await,
Err(broadcast::error::RecvError::Lagged(_n)) => {
emit(&mut output, None).await;
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
},
)
}
pub(crate) fn coalesced_broadcast_producer<Msg, T>(
capacity: usize,
source: &'static std::sync::OnceLock<tokio::sync::broadcast::Sender<T>>,
window: Duration,
make_msg: impl Fn() -> Msg + Send + 'static,
) -> impl futures_util::Stream<Item = Msg>
where
Msg: Send + 'static,
T: Clone + Send + 'static,
{
iced::stream::channel(
capacity,
move |output: iced::futures::channel::mpsc::Sender<Msg>| async move {
let Some(rx) = source.get().and_then(|tx| {
if tx.receiver_count() > 100 {
None
} else {
Some(tx.subscribe())
}
}) else {
return;
};
coalesce_loop(rx, output, window, make_msg).await;
},
)
}
async fn coalesce_loop<Msg, T, F>(
mut rx: broadcast::Receiver<T>,
mut output: iced::futures::channel::mpsc::Sender<Msg>,
window: Duration,
make_msg: F,
) where
Msg: Send,
T: Clone + Send,
F: Fn() -> Msg + Send,
{
loop {
match rx.recv().await {
Ok(_) | Err(broadcast::error::RecvError::Lagged(_)) => {
let deadline = tokio::time::Instant::now() + window;
loop {
tokio::select! {
() = tokio::time::sleep_until(deadline) => {
let _ = futures_util::SinkExt::send(&mut output, make_msg()).await;
break;
}
recv = rx.recv() => {
match recv {
Ok(_) | Err(broadcast::error::RecvError::Lagged(_)) => {}
Err(broadcast::error::RecvError::Closed) => return,
}
}
}
}
}
Err(broadcast::error::RecvError::Closed) => return,
}
}
}
pub(crate) fn apply_editor_action(
content: &mut EditorBuffer,
undo_stack: &mut UndoStack,
action: EditorAction,
) {
match action {
EditorAction::Undo => {
if let Some(s) = undo_stack.undo(content) {
restore_undo_snapshot(content, Some(s));
}
}
EditorAction::Redo => {
if let Some(s) = undo_stack.redo(content) {
restore_undo_snapshot(content, Some(s));
}
}
other => {
if other.is_edit_action() {
undo_stack.snap_before_edit(content);
}
content.perform_action(other);
}
}
}
fn restore_undo_snapshot(content: &mut EditorBuffer, snapshot: Option<UndoSnapshot>) {
if let Some(snapshot) = snapshot {
content.set_text(&snapshot.text);
content.move_to(snapshot.cursor.line, snapshot.cursor.column);
}
}
#[must_use]
pub(crate) fn focus_navigation_task<Message>(action: &EditorAction) -> Option<iced::Task<Message>> {
match action {
EditorAction::FocusNext => Some(iced::widget::operation::focus_next()),
EditorAction::FocusPrevious => Some(iced::widget::operation::focus_previous()),
_ => None,
}
}
pub(crate) struct SingleLineEditorState {
pub(crate) buffer: EditorBuffer,
undo: UndoStack,
}
impl SingleLineEditorState {
pub(crate) fn new(text: &str) -> Self {
let buffer = EditorBuffer::with_text(text, None);
buffer.set_single_line(true);
Self {
buffer,
undo: UndoStack::new(),
}
}
pub(crate) fn set_text(&mut self, text: &str) {
self.buffer.set_text(text);
self.undo.clear();
}
pub(crate) fn text(&self) -> String {
self.buffer.text()
}
pub(crate) fn clear(&mut self) {
self.buffer.clear();
self.undo.clear();
}
pub(crate) fn apply_action(&mut self, action: EditorAction) {
apply_editor_action(&mut self.buffer, &mut self.undo, action);
}
}
pub(crate) fn send_guard<M: 'static>(
text: &str,
sending: bool,
in_flight_first: bool,
over_limit: impl Fn(usize) -> Task<M>,
) -> Result<&str, Task<M>> {
let trimmed = text.trim();
if trimmed.is_empty() {
return Err(Task::none());
}
let over_limit_task = || {
let count = trimmed.chars().count();
if count > MAX_INPUT_CHARS {
Some(over_limit(count))
} else {
None
}
};
if in_flight_first {
if sending {
return Err(Task::none());
}
if let Some(task) = over_limit_task() {
return Err(task);
}
} else {
if let Some(task) = over_limit_task() {
return Err(task);
}
if sending {
return Err(Task::none());
}
}
Ok(trimmed)
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::Cell;
fn stack_with_snapshot(text: &str) -> UndoStack {
let mut stack = UndoStack::new();
let content = EditorBuffer::with_text(text, None);
stack.snap_before_edit(&content);
stack
}
#[test]
fn undo_restores_snapshot() {
let mut stack = stack_with_snapshot("original");
let modified = EditorBuffer::with_text("modified", None);
let snapshot = stack.undo(&modified).unwrap();
assert_eq!(snapshot.text, "original");
}
#[test]
fn redo_restores_undone_state() {
let mut stack = stack_with_snapshot("original");
let modified = EditorBuffer::with_text("modified", None);
let _ = stack.undo(&modified);
let snapshot = stack.redo(&modified).unwrap();
assert_eq!(snapshot.text, "modified");
}
#[test]
fn new_edit_clears_redo() {
let mut stack = stack_with_snapshot("v1");
let v2 = EditorBuffer::with_text("v2", None);
let _ = stack.undo(&v2);
let v3 = EditorBuffer::with_text("v3", None);
stack.snap_before_edit(&v3);
assert!(stack.redo(&v3).is_none());
}
#[test]
fn snapshot_preserves_cursor() {
let content = EditorBuffer::with_text("line1\nline2\nline3", None);
content.move_to(1, 2);
let mut stack = UndoStack::new();
stack.snap_before_edit(&content);
let modified = EditorBuffer::with_text("changed", None);
let snapshot = stack.undo(&modified).unwrap();
assert_eq!(snapshot.cursor.line, 1);
assert_eq!(snapshot.cursor.column, 2);
}
#[test]
fn apply_editor_action_snapshots_edits() {
let mut buffer = EditorBuffer::with_text("hello", None);
buffer.move_to(0, 5);
let mut stack = UndoStack::new();
apply_editor_action(&mut buffer, &mut stack, EditorAction::Insert('!'));
assert_eq!(buffer.text(), "hello!");
assert_eq!(stack.undo(&buffer).unwrap().text, "hello");
}
#[test]
fn restore_undo_snapshot_restores_text_and_cursor() {
let mut buffer = EditorBuffer::with_text("changed", None);
let snapshot = Some(UndoSnapshot {
text: "line1\nline2\nline3".to_string(),
cursor: CursorState {
line: 1,
column: 3,
selection: None,
},
});
restore_undo_snapshot(&mut buffer, snapshot);
assert_eq!(buffer.text(), "line1\nline2\nline3");
let cursor = buffer.cursor();
assert_eq!(cursor.line, 1);
assert_eq!(cursor.column, 3);
}
fn run_guard(
text: &str,
sending: bool,
in_flight_first: bool,
) -> (bool, Result<&str, Task<()>>) {
let fired = Cell::new(false);
let result = send_guard(text, sending, in_flight_first, |_| {
fired.set(true);
Task::none()
});
(fired.get(), result)
}
#[test]
fn send_guard_rejects_empty_and_trims() {
let (fired, result) = run_guard(" \t ", false, true);
assert!(result.is_err() && !fired, "empty input is a silent noop");
let (_, result) = run_guard(" hello ", false, true);
assert_eq!(result.unwrap(), "hello");
}
#[test]
fn send_guard_limit_boundary() {
let at_limit = "a".repeat(MAX_INPUT_CHARS);
let (fired, result) = run_guard(&at_limit, false, true);
assert!(result.is_ok() && !fired, "at-limit text is accepted");
let over_limit = "a".repeat(MAX_INPUT_CHARS + 1);
let (fired, result) = run_guard(&over_limit, false, true);
assert!(result.is_err() && fired, "over-limit text is rejected");
}
#[test]
fn send_guard_combined_in_flight_and_over_limit() {
let text = "a".repeat(MAX_INPUT_CHARS + 1);
let (fired, result) = run_guard(&text, true, true);
assert!(result.is_err() && !fired, "in-flight first: silent noop");
let (fired, result) = run_guard(&text, true, false);
assert!(result.is_err() && fired, "toast fires over-limit first");
}
#[tokio::test]
async fn coalesce_loop_batches_burst() {
use futures_util::StreamExt;
let (tx, rx) = broadcast::channel::<()>(16);
let (out_tx, mut out_rx) = iced::futures::channel::mpsc::channel::<u32>(16);
let window = Duration::from_millis(50);
let handle = tokio::spawn(coalesce_loop(rx, out_tx, window, || 1u32));
for _ in 0..5 {
let _ = tx.send(());
}
let first = tokio::time::timeout(Duration::from_secs(1), out_rx.next())
.await
.expect("burst should flush within the timeout");
assert_eq!(first, Some(1));
let _ = tx.send(());
let second = tokio::time::timeout(Duration::from_secs(1), out_rx.next())
.await
.expect("trailing event should flush within the timeout");
assert_eq!(second, Some(1));
drop(tx);
let _ = tokio::time::timeout(Duration::from_secs(1), handle)
.await
.expect("loop should end when the source closes");
}
}