use std::collections::HashMap;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum ChartType {
Bar,
Line,
Scatter,
Pie,
Area,
Histogram,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DataPoint {
pub x: f64,
pub y: f64,
pub label: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DataSeries {
pub name: String,
pub chart_type: ChartType,
pub points: Vec<DataPoint>,
}
impl DataSeries {
pub fn add_point(&mut self, point: DataPoint) {
self.points.push(point);
}
pub fn min(&self) -> Option<f64> {
self.points
.iter()
.map(|p| p.y)
.filter(|value| value.is_finite())
.min_by(f64::total_cmp)
}
pub fn max(&self) -> Option<f64> {
self.points
.iter()
.map(|p| p.y)
.filter(|value| value.is_finite())
.max_by(f64::total_cmp)
}
pub fn mean(&self) -> Option<f64> {
let mut count = 0u64;
let mut mean = 0.0f64;
for value in self
.points
.iter()
.map(|point| point.y)
.filter(|y| y.is_finite())
{
count = count.saturating_add(1);
let count_f64 = count as f64;
mean = mean * ((count_f64 - 1.0) / count_f64) + value / count_f64;
}
if count == 0 {
return None;
}
Some(mean)
}
pub fn sum(&self) -> f64 {
self.points
.iter()
.map(|p| p.y)
.filter(|value| value.is_finite())
.fold(0.0, |sum, value| {
let next = sum + value;
if next.is_finite() {
next
} else if next.is_sign_negative() {
f64::MIN
} else {
f64::MAX
}
})
}
pub fn sort_by_x(&mut self) {
self.points
.sort_by(|a, b| match (a.x.is_finite(), b.x.is_finite()) {
(true, true) => a.x.total_cmp(&b.x),
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
(false, false) => a.x.total_cmp(&b.x),
});
}
pub fn len(&self) -> usize {
self.points.len()
}
pub fn is_empty(&self) -> bool {
self.points.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum FilterOp {
Eq,
Neq,
Gt,
Gte,
Lt,
Lte,
Contains,
StartsWith,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DataFilter {
pub column: String,
pub op: FilterOp,
pub value: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GroupBy {
pub column: String,
pub aggregation: Aggregation,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum Aggregation {
Count,
Sum,
Avg,
Min,
Max,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum QueryStatus {
Queued,
Running,
Completed,
Failed(String),
Cancelled,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QueryJob {
pub id: String,
pub query: String,
pub filters: Vec<DataFilter>,
pub group_by: Option<GroupBy>,
pub status: QueryStatus,
pub created_at: u64,
pub completed_at: Option<u64>,
}
#[derive(Debug)]
pub struct QueryScheduler {
jobs: Vec<QueryJob>,
indices: HashMap<String, usize>,
next_id: u64,
max_jobs: usize,
}
impl Default for QueryScheduler {
fn default() -> Self {
Self::with_max_jobs(Self::DEFAULT_MAX_JOBS)
}
}
impl QueryScheduler {
pub const DEFAULT_MAX_JOBS: usize = 100_000;
pub fn new() -> Self {
Self::default()
}
pub fn with_max_jobs(max_jobs: usize) -> Self {
Self {
jobs: Vec::new(),
indices: HashMap::new(),
next_id: 0,
max_jobs,
}
}
pub fn submit(
&mut self,
query: String,
filters: Vec<DataFilter>,
group_by: Option<GroupBy>,
created_at: u64,
) -> anyhow::Result<String> {
if self.jobs.len() >= self.max_jobs {
anyhow::bail!("query scheduler capacity of {} jobs reached", self.max_jobs);
}
self.jobs
.try_reserve(1)
.map_err(|error| anyhow::anyhow!("cannot reserve query job storage: {error}"))?;
self.indices
.try_reserve(1)
.map_err(|error| anyhow::anyhow!("cannot reserve query index storage: {error}"))?;
let id = self.allocate_id()?;
let index = self.jobs.len();
self.jobs.push(QueryJob {
id: id.clone(),
query,
filters,
group_by,
status: QueryStatus::Queued,
created_at,
completed_at: None,
});
let replaced = self.indices.insert(id.clone(), index);
debug_assert!(replaced.is_none());
Ok(id)
}
fn allocate_id(&mut self) -> anyhow::Result<String> {
let start = self.next_id;
let mut candidate = start;
loop {
let id = format!("job-{candidate}");
if !self.indices.contains_key(&id) {
self.next_id = candidate.wrapping_add(1);
return Ok(id);
}
candidate = candidate.wrapping_add(1);
if candidate == start {
anyhow::bail!("query job id space exhausted");
}
}
}
pub fn cancel(&mut self, id: &str, completed_at: u64) -> anyhow::Result<()> {
let job = self.job_mut(id)?;
match job.status {
QueryStatus::Queued | QueryStatus::Running => {
job.status = QueryStatus::Cancelled;
job.completed_at = Some(completed_at);
Ok(())
}
_ => anyhow::bail!("job '{id}' is already in terminal state"),
}
}
pub fn start(&mut self, id: &str) -> anyhow::Result<()> {
let job = self.job_mut(id)?;
if job.status != QueryStatus::Queued {
anyhow::bail!("job '{id}' is not queued");
}
job.status = QueryStatus::Running;
Ok(())
}
pub fn complete(&mut self, id: &str, completed_at: u64) -> anyhow::Result<()> {
let job = self.job_mut(id)?;
if job.status != QueryStatus::Running {
anyhow::bail!("job '{id}' is not running");
}
job.status = QueryStatus::Completed;
job.completed_at = Some(completed_at);
Ok(())
}
pub fn fail(
&mut self,
id: &str,
message: impl Into<String>,
completed_at: u64,
) -> anyhow::Result<()> {
let job = self.job_mut(id)?;
if !matches!(job.status, QueryStatus::Queued | QueryStatus::Running) {
anyhow::bail!("job '{id}' is already in terminal state");
}
job.status = QueryStatus::Failed(message.into());
job.completed_at = Some(completed_at);
Ok(())
}
fn job_mut(&mut self, id: &str) -> anyhow::Result<&mut QueryJob> {
let index = self
.indices
.get(id)
.copied()
.ok_or_else(|| anyhow::anyhow!("job '{id}' not found"))?;
self.jobs
.get_mut(index)
.ok_or_else(|| anyhow::anyhow!("job index for '{id}' is invalid"))
}
pub fn get(&self, id: &str) -> Option<&QueryJob> {
self.indices.get(id).and_then(|&index| self.jobs.get(index))
}
pub fn remove(&mut self, id: &str) -> Option<QueryJob> {
let index = self.indices.remove(id)?;
let removed = self.jobs.remove(index);
for (offset, job) in self.jobs[index..].iter().enumerate() {
if let Some(stored_index) = self.indices.get_mut(&job.id) {
*stored_index = index + offset;
}
}
Some(removed)
}
pub fn list(&self) -> &[QueryJob] {
&self.jobs
}
pub const fn max_jobs(&self) -> usize {
self.max_jobs
}
pub fn len(&self) -> usize {
self.jobs.len()
}
pub fn is_empty(&self) -> bool {
self.jobs.is_empty()
}
pub fn completed_jobs(&self) -> Vec<&QueryJob> {
self.jobs
.iter()
.filter(|j| j.status == QueryStatus::Completed)
.collect()
}
pub fn pending_count(&self) -> usize {
self.jobs
.iter()
.filter(|j| matches!(j.status, QueryStatus::Queued | QueryStatus::Running))
.count()
}
}
#[derive(Debug, Default)]
pub struct CsvImporter {
headers: Vec<String>,
}
impl CsvImporter {
pub fn new() -> Self {
Self::default()
}
pub fn parse_header(&mut self, line: &str) -> anyhow::Result<Vec<String>> {
let headers = parse_csv_line(line)?;
if headers.is_empty() || headers.iter().any(String::is_empty) {
anyhow::bail!("CSV headers must be non-empty");
}
let mut unique = std::collections::HashSet::new();
if headers.iter().any(|header| !unique.insert(header.as_str())) {
anyhow::bail!("CSV headers must be unique");
}
self.headers = headers.clone();
Ok(headers)
}
pub fn parse_row(&self, line: &str) -> anyhow::Result<Vec<(String, String)>> {
if self.headers.is_empty() {
anyhow::bail!("headers must be parsed before rows");
}
let values = parse_csv_line(line)?;
if values.len() != self.headers.len() {
anyhow::bail!(
"column count mismatch: expected {}, got {}",
self.headers.len(),
values.len()
);
}
Ok(self
.headers
.iter()
.zip(values.iter())
.map(|(h, v)| (h.clone(), v.clone()))
.collect())
}
pub fn validate_row(&self, line: &str) -> bool {
if self.headers.is_empty() {
return false;
}
parse_csv_line(line).is_ok_and(|values| values.len() == self.headers.len())
}
}
fn parse_csv_line(line: &str) -> anyhow::Result<Vec<String>> {
const MAX_RECORD_BYTES: usize = 16 * 1024 * 1024;
const MAX_FIELDS: usize = 65_536;
if line.len() > MAX_RECORD_BYTES {
anyhow::bail!("CSV record exceeds {MAX_RECORD_BYTES} bytes");
}
#[derive(Clone, Copy)]
enum State {
FieldStart,
Unquoted,
Quoted,
AfterQuote,
}
let mut values = Vec::new();
let mut value = String::new();
let mut state = State::FieldStart;
for character in line.chars() {
match (state, character) {
(State::FieldStart, '"') => state = State::Quoted,
(State::FieldStart, ',') => push_csv_value(&mut values, String::new(), MAX_FIELDS)?,
(State::FieldStart, character) if character.is_ascii_whitespace() => {}
(State::FieldStart, _) => {
value.push(character);
state = State::Unquoted;
}
(State::Unquoted, ',') => {
push_csv_value(&mut values, value.trim().to_string(), MAX_FIELDS)?;
value = String::new();
state = State::FieldStart;
}
(State::Unquoted, '"') => anyhow::bail!("quote inside unquoted CSV field"),
(State::Unquoted, _) => value.push(character),
(State::Quoted, '"') => state = State::AfterQuote,
(State::Quoted, _) => value.push(character),
(State::AfterQuote, '"') => {
value.push('"');
state = State::Quoted;
}
(State::AfterQuote, ',') => {
push_csv_value(&mut values, std::mem::take(&mut value), MAX_FIELDS)?;
state = State::FieldStart;
}
(State::AfterQuote, character) if character.is_ascii_whitespace() => {}
(State::AfterQuote, _) => {
anyhow::bail!("unexpected character after closing CSV quote")
}
}
}
match state {
State::Quoted => anyhow::bail!("unterminated quoted CSV field"),
State::Unquoted => {
push_csv_value(&mut values, value.trim().to_string(), MAX_FIELDS)?;
}
State::FieldStart => push_csv_value(&mut values, String::new(), MAX_FIELDS)?,
State::AfterQuote => push_csv_value(&mut values, value, MAX_FIELDS)?,
}
Ok(values)
}
fn push_csv_value(
values: &mut Vec<String>,
value: String,
max_fields: usize,
) -> anyhow::Result<()> {
if values.len() >= max_fields {
anyhow::bail!("CSV record exceeds {max_fields} fields");
}
values.push(value);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_series() -> DataSeries {
DataSeries {
name: "test".into(),
chart_type: ChartType::Line,
points: vec![
DataPoint {
x: 1.0,
y: 10.0,
label: None,
},
DataPoint {
x: 3.0,
y: 30.0,
label: None,
},
DataPoint {
x: 2.0,
y: 20.0,
label: Some("mid".into()),
},
],
}
}
#[test]
fn data_series_add_point() {
let mut ds = DataSeries {
name: "s".into(),
chart_type: ChartType::Bar,
points: vec![],
};
ds.add_point(DataPoint {
x: 0.0,
y: 5.0,
label: None,
});
assert_eq!(ds.len(), 1);
assert!(!ds.is_empty());
}
#[test]
fn data_series_min_max() {
let ds = sample_series();
assert!((ds.min().unwrap() - 10.0).abs() < f64::EPSILON);
assert!((ds.max().unwrap() - 30.0).abs() < f64::EPSILON);
}
#[test]
fn data_series_mean() {
let ds = sample_series();
assert!((ds.mean().unwrap() - 20.0).abs() < f64::EPSILON);
}
#[test]
fn data_series_sum() {
let ds = sample_series();
assert!((ds.sum() - 60.0).abs() < f64::EPSILON);
}
#[test]
fn data_series_sort_by_x() {
let mut ds = sample_series();
ds.sort_by_x();
let xs: Vec<f64> = ds.points.iter().map(|p| p.x).collect();
assert_eq!(xs, vec![1.0, 2.0, 3.0]);
}
#[test]
fn data_series_empty_stats() {
let ds = DataSeries {
name: "empty".into(),
chart_type: ChartType::Pie,
points: vec![],
};
assert!(ds.min().is_none());
assert!(ds.max().is_none());
assert!(ds.mean().is_none());
assert!(ds.is_empty());
}
#[test]
fn data_series_stats_ignore_non_finite_values_and_sort_them_last() {
let mut series = sample_series();
series.points.push(DataPoint {
x: f64::NAN,
y: f64::NAN,
label: None,
});
assert_eq!(series.min(), Some(10.0));
assert_eq!(series.max(), Some(30.0));
assert_eq!(series.mean(), Some(20.0));
assert_eq!(series.sum(), 60.0);
series.sort_by_x();
assert!(series.points.last().unwrap().x.is_nan());
}
#[test]
fn query_scheduler_submit_and_get() {
let mut sched = QueryScheduler::new();
let id = sched.submit("SELECT *".into(), vec![], None, 10).unwrap();
assert!(sched.get(&id).is_some());
assert_eq!(sched.get(&id).unwrap().status, QueryStatus::Queued);
assert_eq!(sched.pending_count(), 1);
}
#[test]
fn query_scheduler_cancel() {
let mut sched = QueryScheduler::new();
let id = sched.submit("SELECT 1".into(), vec![], None, 10).unwrap();
assert!(sched.cancel(&id, 20).is_ok());
assert_eq!(sched.get(&id).unwrap().status, QueryStatus::Cancelled);
assert_eq!(sched.get(&id).unwrap().completed_at, Some(20));
assert!(sched.cancel(&id, 21).is_err());
}
#[test]
fn query_scheduler_cancel_unknown() {
let mut sched = QueryScheduler::new();
assert!(sched.cancel("nope", 0).is_err());
}
#[test]
fn query_scheduler_list_and_completed() {
let mut sched = QueryScheduler::new();
sched.submit("q1".into(), vec![], None, 1).unwrap();
sched.submit("q2".into(), vec![], None, 2).unwrap();
assert_eq!(sched.list().len(), 2);
assert!(sched.completed_jobs().is_empty());
}
#[test]
fn query_scheduler_runs_jobs_to_terminal_states_and_wraps_ids() {
let mut scheduler = QueryScheduler::new();
scheduler.next_id = u64::MAX;
let max = scheduler.submit("max".into(), vec![], None, 1).unwrap();
let zero = scheduler.submit("zero".into(), vec![], None, 2).unwrap();
assert_eq!(max, format!("job-{}", u64::MAX));
assert_eq!(zero, "job-0");
scheduler.start(&max).unwrap();
scheduler.complete(&max, 123).unwrap();
assert_eq!(scheduler.get(&max).unwrap().completed_at, Some(123));
assert_eq!(scheduler.completed_jobs().len(), 1);
scheduler.fail(&zero, "bad query", 124).unwrap();
assert!(matches!(
scheduler.get(&zero).unwrap().status,
QueryStatus::Failed(_)
));
}
#[test]
fn query_scheduler_enforces_capacity_and_can_release_jobs() {
let mut scheduler = QueryScheduler::with_max_jobs(2);
let first = scheduler.submit("first".into(), vec![], None, 1).unwrap();
let second = scheduler.submit("second".into(), vec![], None, 2).unwrap();
assert_eq!(scheduler.max_jobs(), 2);
assert_eq!(scheduler.len(), 2);
assert!(
scheduler
.submit("overflow".into(), vec![], None, 3)
.is_err()
);
assert_eq!(scheduler.remove(&first).unwrap().query, "first");
let third = scheduler.submit("third".into(), vec![], None, 4).unwrap();
assert!(scheduler.get(&first).is_none());
assert_eq!(scheduler.get(&second).unwrap().query, "second");
assert_eq!(scheduler.get(&third).unwrap().query, "third");
assert_eq!(scheduler.list()[0].id, second);
assert_eq!(scheduler.list()[1].id, third);
}
#[test]
fn csv_importer_parse_header_and_row() {
let mut csv = CsvImporter::new();
let headers = csv.parse_header("name, age, city").unwrap();
assert_eq!(headers, vec!["name", "age", "city"]);
let row = csv.parse_row("Alice, 30, NYC").unwrap();
assert_eq!(row.len(), 3);
assert_eq!(row[0], ("name".into(), "Alice".into()));
assert_eq!(row[1], ("age".into(), "30".into()));
}
#[test]
fn csv_importer_handles_quotes_and_rejects_malformed_headers() {
let mut csv = CsvImporter::new();
assert_eq!(
csv.parse_header("name,notes").unwrap(),
vec!["name", "notes"]
);
let row = csv.parse_row("Alice,\"hello, \"\"world\"\"\"").unwrap();
assert_eq!(row[1].1, "hello, \"world\"");
let spaced = csv.parse_row(" Alice , \" keep me \" ").unwrap();
assert_eq!(spaced[0].1, "Alice");
assert_eq!(spaced[1].1, " keep me ");
assert!(!csv.validate_row("Alice,\"unterminated"));
assert!(csv.parse_header("name,name").is_err());
assert!(csv.parse_header("name,").is_err());
assert!(csv.parse_row("Alice,un\"quoted").is_err());
assert!(csv.parse_row("Alice,\"quoted\"suffix").is_err());
}
#[test]
fn csv_importer_row_without_header() {
let csv = CsvImporter::new();
assert!(csv.parse_row("a,b,c").is_err());
}
#[test]
fn csv_importer_column_mismatch() {
let mut csv = CsvImporter::new();
csv.parse_header("a,b").unwrap();
assert!(csv.parse_row("1,2,3").is_err());
}
#[test]
fn csv_importer_validate_row() {
let mut csv = CsvImporter::new();
assert!(!csv.validate_row("x,y"));
csv.parse_header("a,b").unwrap();
assert!(csv.validate_row("1,2"));
assert!(!csv.validate_row("1,2,3"));
}
#[test]
fn filter_op_serialization() {
let filter = DataFilter {
column: "price".into(),
op: FilterOp::Gte,
value: "100".into(),
};
let json = serde_json::to_string(&filter).unwrap();
let deser: DataFilter = serde_json::from_str(&json).unwrap();
assert_eq!(deser.op, FilterOp::Gte);
assert_eq!(deser.column, "price");
}
#[test]
fn query_status_variants() {
let statuses = vec![
QueryStatus::Queued,
QueryStatus::Running,
QueryStatus::Completed,
QueryStatus::Failed("err".into()),
QueryStatus::Cancelled,
];
for s in &statuses {
let json = serde_json::to_string(s).unwrap();
let deser: QueryStatus = serde_json::from_str(&json).unwrap();
assert_eq!(&deser, s);
}
}
}