use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime};
use tokio::sync::broadcast;
use tokio::task::JoinHandle;
use tokio::time::interval;
use vantage_diorama::ChangeEvent;
use vantage_vista::Vista;
use crate::live_folder::listing_shell::FolderListingShell;
use crate::live_folder::size_shell::FolderSizeShell;
use crate::live_folder::tree::simulate_second;
mod listing_shell;
mod size_shell;
pub mod tree;
pub use tree::{Entry, EntryKind, Tree, format_ts};
const EVENT_CAPACITY: usize = 1024;
const LOOP_TICK: Duration = Duration::from_secs(1);
pub const EVENT_TYPES: &[(&str, u32)] = &[
("user_signup", 1),
("payment_succeeded", 2),
("payment_failed", 1),
("page_view", 10),
("api_call", 8),
("search_query", 6),
("file_upload", 3),
("email_sent", 4),
("notification", 5),
("error_reported", 1),
];
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum PushMode {
Poll,
Notify,
}
#[derive(Clone, Debug)]
pub struct LiveFolderConfig {
pub requests_per_sec: u64,
pub bytes_per_request: (u64, u64),
pub chunk_threshold: u64,
pub error_pct_per_sec: f64,
pub error_size: (u64, u64),
pub event_bump: (u64, u64),
pub backfill: Duration,
}
impl Default for LiveFolderConfig {
fn default() -> Self {
Self {
requests_per_sec: 100,
bytes_per_request: (60, 100),
chunk_threshold: 288_000,
error_pct_per_sec: 1.0,
error_size: (500, 3_000),
event_bump: (2_000, 4_000),
backfill: Duration::ZERO,
}
}
}
#[derive(Clone)]
pub struct LiveFolderSim {
inner: Arc<Inner>,
events: broadcast::Sender<ChangeEvent>,
_task: Arc<AbortOnDrop>,
}
pub(crate) struct SimState {
pub(crate) tree: Tree,
}
pub(crate) struct Inner {
pub(crate) cfg: LiveFolderConfig,
pub(crate) state: Mutex<SimState>,
}
impl LiveFolderSim {
pub fn new(cfg: LiveFolderConfig) -> Self {
let now = SystemTime::now();
let inner = Arc::new(Inner {
cfg: cfg.clone(),
state: Mutex::new(SimState {
tree: Tree::new(now),
}),
});
if cfg.backfill > Duration::ZERO {
let start = now
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
.saturating_sub(cfg.backfill.as_secs());
let end = now
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let mut state = inner.state.lock().unwrap();
for s in start..end {
let when = SystemTime::UNIX_EPOCH + Duration::from_secs(s);
simulate_second(&mut state.tree, &cfg, when);
}
}
let (events, _) = broadcast::channel(EVENT_CAPACITY);
let task = spawn_loop(inner.clone(), events.clone());
Self {
inner,
events,
_task: Arc::new(AbortOnDrop(task)),
}
}
pub fn listing_vista(
&self,
name: impl Into<String>,
path: impl Into<String>,
) -> (Vista, broadcast::Sender<ChangeEvent>) {
let shell = FolderListingShell::new(self.inner.clone(), path.into());
(Vista::new(name, Box::new(shell)), self.events.clone())
}
pub fn size_vista(&self, name: impl Into<String>) -> Vista {
Vista::new(name, Box::new(FolderSizeShell::new(self.inner.clone())))
}
pub fn size_augment(&self) -> vantage_diorama::Augmentation {
vantage_diorama::Augmentation {
detail: vantage_diorama::Detail::Fixed(std::sync::Arc::new(
self.size_vista("folder_size"),
)),
source: vantage_diorama::Source::Column {
from: "path".to_string(),
to: None,
},
fetch: vantage_diorama::Fetch::PerRow,
merge: vantage_diorama::MergeRule {
columns: vec!["size".to_string(), "file_count".to_string()],
},
}
}
pub fn snapshot(&self) -> Vec<(String, Entry)> {
let state = self.inner.state.lock().unwrap();
let mut out = Vec::new();
fn walk(prefix: &str, entry: &Entry, out: &mut Vec<(String, Entry)>) {
for (name, child) in &entry.children {
let path = if prefix.is_empty() {
name.clone()
} else {
format!("{prefix}/{name}")
};
out.push((path.clone(), child.clone()));
if child.kind == EntryKind::Folder {
walk(&path, child, out);
}
}
}
walk("", &state.tree.root, &mut out);
out
}
pub fn config(&self) -> LiveFolderConfig {
self.inner.cfg.clone()
}
}
fn spawn_loop(inner: Arc<Inner>, events: broadcast::Sender<ChangeEvent>) -> JoinHandle<()> {
tokio::spawn(async move {
let mut ticker = interval(LOOP_TICK);
loop {
ticker.tick().await;
let now = SystemTime::now();
{
let mut state = inner.state.lock().unwrap();
simulate_second(&mut state.tree, &inner.cfg, now);
}
let _ = events.send(ChangeEvent::Invalidated);
}
})
}
struct AbortOnDrop(JoinHandle<()>);
impl Drop for AbortOnDrop {
fn drop(&mut self) {
self.0.abort();
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use vantage_dataset::prelude::ReadableValueSet;
fn cfg_no_backfill() -> LiveFolderConfig {
LiveFolderConfig {
backfill: Duration::ZERO,
..LiveFolderConfig::default()
}
}
#[tokio::test]
async fn listing_vista_reads_from_the_live_tree() {
let cfg = LiveFolderConfig {
backfill: Duration::from_secs(3600),
error_pct_per_sec: 100.0, ..LiveFolderConfig::default()
};
let sim = LiveFolderSim::new(cfg);
let (root_vista, _tx) = sim.listing_vista("root", "");
let rows = root_vista.list_values().await.unwrap();
assert!(!rows.is_empty(), "root listing should have day folders");
}
#[tokio::test]
async fn size_vista_returns_none_for_unknown_path() {
let sim = LiveFolderSim::new(cfg_no_backfill());
let vista = sim.size_vista("sizes");
let got = vista.get_value("does/not/exist".to_string()).await.unwrap();
assert!(got.is_none());
}
#[tokio::test]
async fn size_vista_returns_size_for_a_real_folder() {
let sim = LiveFolderSim::new(LiveFolderConfig {
backfill: Duration::from_secs(120),
..LiveFolderConfig::default()
});
let vista = sim.size_vista("sizes");
let snap = sim.snapshot();
let any_folder = snap
.iter()
.find(|(_, e)| e.kind == EntryKind::Folder && !e.children.is_empty())
.map(|(p, _)| p.clone())
.expect("backfill produced at least one folder");
let rec = vista
.get_value(&any_folder)
.await
.unwrap()
.expect("folder resolves");
assert!(rec.get("size").is_some());
assert!(rec.get("file_count").is_some());
}
#[tokio::test]
async fn listing_vista_supports_subdir_traversal() {
let sim = LiveFolderSim::new(LiveFolderConfig {
backfill: Duration::from_secs(3600),
..LiveFolderConfig::default()
});
let (ymd_vista, _tx) = sim.listing_vista("ymd", "");
let ymd_rows = ymd_vista.list_values().await.unwrap();
let date_row = ymd_rows
.iter()
.find(|(_, r)| r.get("kind").and_then(|v| v.as_text()) == Some("folder"))
.map(|(_, r)| r.clone())
.expect("at least one date folder");
let sub = ymd_vista.get_ref("subdir", &date_row).expect("subdir ref");
let sub_rows = sub.list_values().await.unwrap();
assert!(!sub_rows.is_empty(), "date folder should list its children");
let parent = date_row
.get("path")
.and_then(|v| v.as_text())
.map(|s| s.to_string())
.unwrap_or_default();
for (_, rec) in &sub_rows {
let path = rec
.get("path")
.and_then(|v| v.as_text())
.map(|s| s.to_string())
.unwrap_or_default();
assert!(
path.starts_with(&parent),
"child path {path:?} should start with parent {parent:?}"
);
}
}
#[tokio::test]
async fn listing_leaves_folder_size_unfilled_for_the_augment() {
let sim = LiveFolderSim::new(LiveFolderConfig {
backfill: Duration::from_secs(3600),
error_pct_per_sec: 100.0, ..LiveFolderConfig::default()
});
let (root, _tx) = sim.listing_vista("root", "");
let rows = root.list_values().await.unwrap();
assert!(!rows.is_empty());
for (_, rec) in &rows {
assert_eq!(rec.get("kind").and_then(|v| v.as_text()), Some("folder"));
assert!(rec.get("size").is_none(), "folder rows leave size unfilled");
}
let file_path = sim
.snapshot()
.into_iter()
.find(|(_, e)| e.kind == EntryKind::File)
.map(|(p, _)| p)
.expect("backfill produced files");
let parent = file_path.rsplit_once('/').map(|(p, _)| p).unwrap_or("");
let (listing, _tx) = sim.listing_vista("parent", parent);
let rows = listing.list_values().await.unwrap();
let file_row = rows
.values()
.find(|r| r.get("kind").and_then(|v| v.as_text()) == Some("file"))
.expect("parent folder lists the file");
assert!(
file_row.get("size").is_some(),
"file rows carry their own size"
);
}
#[tokio::test]
async fn size_augment_hydrates_listing_rows_over_one_dio() {
use std::sync::Arc;
let sim = LiveFolderSim::new(LiveFolderConfig {
backfill: Duration::from_secs(1800),
..LiveFolderConfig::default()
});
let (listing, _tx) = sim.listing_vista("root", "");
let lens = Arc::new(
vantage_diorama::Lens::new()
.cache_in_memory()
.viewport_debounce(Duration::from_millis(1))
.build()
.expect("lens builds"),
);
let dio = lens.make_dio(listing).await.expect("make_dio").augment(
Arc::new(vantage_vista_factory::VistaCatalog::new()),
vec![sim.size_augment()],
);
let scenery = dio.table_scenery().open().await.expect("scenery opens");
scenery.set_viewport(0..10);
let mut augmented = false;
for _ in 0..100 {
let n = scenery.row_count();
augmented = (0..n).filter_map(|i| scenery.row(i)).any(|row| {
row.record.get("file_count").is_some()
&& matches!(
row.record.get("size"),
Some(ciborium::Value::Integer(i)) if i128::from(*i) > 0
)
});
if augmented {
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
assert!(
augmented,
"a folder row gained a positive size + file_count from the augment"
);
}
}