#[cfg(doctest)]
#[doc = include_str!("../README.md")]
pub struct ReadmeDoctests;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use serde_json::{json, Value as J};
use fv_streams_types::settings::env;
mod batches;
pub mod compute;
pub use fv_streams_runtime::dataflow;
use fv_streams_runtime::{cluster, placement};
pub use fv_streams_types::decode;
mod keys;
use fv_streams_state as epochs;
mod ops;
mod pipeline;
mod rows;
mod spec;
pub mod steps;
use batches::*;
pub use fv_streams_types::{
Binding, BuildSignal, ControlPlane, InlineSource, OutputSink, SinkCtx, SourceCtx, StageDef, Topology,
};
use keys::*;
use spec::*;
fn now_ms() -> i64 {
chrono::Utc::now().timestamp_millis()
}
pub(crate) fn cpu_families() -> std::collections::BTreeMap<String, u64> {
let mut families: std::collections::BTreeMap<String, u64> = std::collections::BTreeMap::new();
for (name, ms) in &dataflow::process_cpu_by_thread() {
let family = if name.starts_with("rdk:") {
name.trim_end_matches(|c: char| c.is_ascii_digit()).to_string()
} else if name.starts_with("fv-task-") {
"fv-task".to_string()
} else {
name.clone()
};
*families.entry(family).or_default() += ms;
}
families
}
fn trace_cpu(build_id: &str) {
if env("STREAM_TRACE_CPU", "0") != "1" {
return;
}
let families = cpu_families();
let total: u64 = families.values().sum();
let line: Vec<String> = families.iter().map(|(k, v)| format!("{k}={v}")).collect();
eprintln!("stream {build_id}: cpu ms total={total} {}", line.join(" "));
}
pub async fn run_stream(
cp: std::sync::Arc<dyn ControlPlane>,
build_id: String,
pipeline_name: String,
mut record: J,
) -> Result<(), String> {
match run_stream_inner(cp.as_ref(), &build_id, &pipeline_name, &mut record).await {
Ok(()) => Ok(()),
Err(e) => {
eprintln!("stream {build_id}: FAILED — {e}");
record["status"] = json!("FAILED");
record["error"] = json!(e);
record["finishedAt"] = json!(chrono::Utc::now().to_rfc3339());
cp.put_record(&build_id, &record).await;
Err(e)
}
}
}
async fn run_stream_inner(
cp: &dyn ControlPlane,
build_id: &str,
pipeline_name: &str,
record: &mut J,
) -> Result<(), String> {
let topo = cp.topology(pipeline_name).await?;
if topo.stages.is_empty() {
return Err("pipeline has no transforms".into());
}
let plans = pipeline::plan(&topo)?;
pipeline::run_pipeline(cp, build_id, pipeline_name, record, &topo, plans).await
}