use std::io::Write;
use std::sync::Arc;
use tabnas_render::{TextOut, WriteOut};
use tabnas_transduce::{AbortFlag, Code, Duplicates, Fail, Limits, Metrics, Selector, Sink};
use crate::ast::Sources;
use crate::interp::{Runtime, MAX_PLAN_STEPS};
use crate::lower::{EventSink, Lowering, Out};
use crate::resolve::{resolve, Resolved};
use crate::value::{Plan, Protocol, Seq, Val};
use crate::{desugar, parse_file, stdlib};
pub use crate::lower::Renderer;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Output {
Text,
TableRows,
JsonEvents,
}
impl Output {
pub fn as_str(self) -> &'static str {
match self {
Output::Text => "Text",
Output::TableRows => "TableRows/1",
Output::JsonEvents => "JsonEvents/1",
}
}
}
#[derive(Clone)]
pub struct Program {
resolved: Arc<Resolved>,
sources: Sources,
output: Output,
native: bool,
duplicates: Duplicates,
abort: AbortFlag,
result: Val,
}
impl std::fmt::Debug for Program {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Program")
.field("file", &self.resolved.file)
.field("output", &self.output())
.field("native", &self.native)
.finish()
}
}
pub fn compile(src: &str, file: &str) -> Result<Program, Fail> {
compile_sources(&[Source { file, text: src }])
}
#[derive(Clone, Copy, Debug)]
pub struct Source<'a> {
pub file: &'a str,
pub text: &'a str,
}
pub fn compile_sources(sources: &[Source<'_>]) -> Result<Program, Fail> {
on_stack(|| {
let mut named: Vec<(Arc<str>, Arc<str>)> = Vec::with_capacity(sources.len());
for source in sources {
if named.iter().any(|(file, _)| &**file == source.file) {
return Err(Fail::new(
Code::DslTypeError,
format!(
"duplicate_file: {} is given twice; each source has its own name",
source.file
),
));
}
named.push((Arc::from(source.file), Arc::from(source.text)));
}
if named.is_empty() {
return Err(crate::check::no_export());
}
let linked = Sources::several(named);
let mut forms = Vec::new();
for source in sources {
let in_file = |fail: Fail| {
if linked.names_files() {
fail.in_file(source.file)
} else {
fail
}
};
let parsed = parse_file(source.text, source.file).map_err(in_file)?;
forms.extend(desugar::program(parsed, source.text).map_err(in_file)?);
}
let resolved = Arc::new(resolve(forms, &linked, &stdlib::outer)?);
let checked = crate::check::program(&resolved, &linked)?;
Program::build_here(
resolved,
linked,
checked.output,
true,
Duplicates::Reject,
AbortFlag::new(),
)
})
}
fn on_stack<T: Send>(work: impl FnOnce() -> Result<T, Fail> + Send) -> Result<T, Fail> {
std::thread::scope(|scope| {
let thread = std::thread::Builder::new()
.name("alchemy-compile".into())
.stack_size(crate::STACK_BYTES)
.spawn_scoped(scope, work)
.map_err(|error| {
Fail::new(
Code::ResourceLimitExceeded,
format!("no thread could be started to compile the program: {error}"),
)
})?;
thread
.join()
.unwrap_or_else(|panic| std::panic::resume_unwind(panic))
})
}
impl Program {
fn build(
resolved: Arc<Resolved>,
sources: Sources,
output: Output,
native: bool,
duplicates: Duplicates,
abort: AbortFlag,
) -> Result<Program, Fail> {
on_stack(|| Program::build_here(resolved, sources, output, native, duplicates, abort))
}
fn build_here(
resolved: Arc<Resolved>,
sources: Sources,
output: Output,
native: bool,
duplicates: Duplicates,
abort: AbortFlag,
) -> Result<Program, Fail> {
let result = Runtime::new(resolved.clone(), sources.clone())
.with_native(native)
.with_duplicates(duplicates)
.with_fuel(Some(MAX_PLAN_STEPS))
.with_abort(abort.clone())
.export()?;
Ok(Program {
resolved,
sources,
output,
native,
duplicates,
abort,
result,
})
}
pub fn with_native(&self, native: bool) -> Result<Program, Fail> {
Program::build(
self.resolved.clone(),
self.sources.clone(),
self.output,
native,
self.duplicates,
self.abort.clone(),
)
}
pub fn with_duplicates(&self, duplicates: Duplicates) -> Result<Program, Fail> {
Program::build(
self.resolved.clone(),
self.sources.clone(),
self.output,
self.native,
duplicates,
self.abort.clone(),
)
}
pub fn with_abort(&self, abort: AbortFlag) -> Program {
Program {
abort,
..self.clone()
}
}
pub fn native(&self) -> bool {
self.native
}
pub fn duplicates(&self) -> Duplicates {
self.duplicates
}
pub fn file(&self) -> &str {
&self.resolved.file
}
pub fn resolved(&self) -> &Arc<Resolved> {
&self.resolved
}
pub fn result(&self) -> &Val {
&self.result
}
pub fn output(&self) -> Output {
self.output
}
pub fn plan_output(&self) -> Output {
match &self.result {
Val::Stream(plan) => match plan.protocol() {
Protocol::JsonEvents => Output::JsonEvents,
_ => Output::TableRows,
},
_ => Output::Text,
}
}
pub fn row_selector(&self) -> Option<&Selector> {
match root_stage(self.plan()?) {
Plan::TableFromJson {
binding: Val::Record(fields),
..
} => match fields.get("rows") {
Some(Val::Selector(s)) => Some(&**s),
_ => None,
},
Plan::Select { selector, .. } => Some(selector),
Plan::Route { specs, .. } => {
let multi: Vec<&Selector> = specs
.iter()
.map(|s| &s.selector)
.filter(|s| s.is_multi())
.collect();
match (multi.as_slice(), specs.len()) {
([one], _) => Some(one),
([], 1) => Some(&specs[0].selector),
_ => None,
}
}
_ => None,
}
}
pub fn explain(&self) -> String {
crate::effects::explain(self)
}
pub fn explain_json(&self) -> serde_json::Value {
crate::effects::explain_json(self)
}
fn plan(&self) -> Option<&Arc<Plan>> {
match &self.result {
Val::Stream(p) | Val::Text(p) => Some(p),
_ => None,
}
}
pub fn sink(
&self,
out: Box<dyn Write + Send>,
render: Option<Renderer>,
limits: &Limits,
metrics: Arc<Metrics>,
) -> Result<Box<dyn Sink + Send>, Fail> {
let out = WriteOut::new(out)
.with_limits(limits)
.with_metrics(metrics.clone());
self.sink_out(Box::new(out), render, limits, metrics)
}
pub fn sink_out(
&self,
out: Box<dyn TextOut + Send>,
render: Option<Renderer>,
limits: &Limits,
metrics: Arc<Metrics>,
) -> Result<EventSink, Fail> {
let rt = Arc::new(
Runtime::new(self.resolved.clone(), self.sources.clone())
.with_native(self.native)
.with_duplicates(self.duplicates)
.with_limits(limits)
.with_abort(self.abort.clone()),
);
let out: Out = out;
Lowering::new(rt, limits, metrics).sink(&self.result, out, render)
}
}
fn root_stage(plan: &Plan) -> &Plan {
let mut here = plan;
loop {
let next: &Plan = match here {
Plan::Input => return here,
Plan::Route { source, .. }
| Plan::Select { source, .. }
| Plan::Events { source }
| Plan::ScanEmit { source, .. }
| Plan::Map { source, .. }
| Plan::Filter { source, .. }
| Plan::TableFromJson { source, .. }
| Plan::Records { source }
| Plan::CsvTable { source, .. }
| Plan::Csv { source, .. }
| Plan::Json { source } => {
if matches!(**source, Plan::Input) {
return here;
}
source
}
Plan::ConcatMap {
items: Seq::Stream(source),
..
}
| Plan::Join {
items: Seq::Stream(source),
..
} => source,
Plan::Concat { items, live } => match live.and_then(|i| items.get(i)) {
Some(Val::Text(inner)) | Some(Val::Stream(inner)) => inner,
_ => return here,
},
Plan::Replace {
source: Val::Text(inner),
..
} => inner,
_ => return here,
};
here = next;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::lower::tests::{PROGRAM, RECORDS};
#[test]
fn a_program_knows_its_output_and_row_selector() {
let p = compile(PROGRAM, "t.alc").unwrap();
assert_eq!(p.output(), Output::Text);
assert_eq!(
p.row_selector().unwrap().to_string(),
".response.payload.deep.records[*]"
);
assert!(p.native());
let slow = p.with_native(false).unwrap();
assert!(!slow.native());
assert_eq!(
slow.row_selector().map(ToString::to_string).as_deref(),
Some(".response.payload.deep.records[*]")
);
let table = compile(&PROGRAM.replace(" csv csv-options\n", ""), "t.alc").unwrap();
assert_eq!(table.output(), Output::TableRows);
let echo = compile("def export [input] input", "t.alc").unwrap();
assert_eq!(echo.output(), Output::JsonEvents);
assert_eq!(echo.row_selector(), None);
let select = compile(
"def export [input] (join \",\" (select (path \"a\" each-index) input))",
"t.alc",
)
.unwrap();
assert_eq!(select.output(), Output::Text);
assert_eq!(select.row_selector().unwrap().to_string(), ".a[*]");
}
#[test]
fn the_sink_writes_through_a_write() {
let p = compile(PROGRAM, "t.alc").unwrap();
let buffer = crate::lower::tests::Shared::default();
let limits = Limits::default();
let mut sink = p
.sink(Box::new(buffer.clone()), None, &limits, Metrics::new())
.unwrap();
let datum = tabnas_transduce::Datum::from_json(&serde_json::from_str(RECORDS).unwrap());
tabnas_transduce::walk_datum(&datum, &mut sink).unwrap();
sink.event(tabnas_transduce::JsonEvent::End).unwrap();
let bytes = buffer.0.lock().unwrap().clone();
assert_eq!(
String::from_utf8(bytes).unwrap(),
crate::lower::tests::EXPECTED_CSV
);
}
#[test]
fn compile_reports_reader_resolver_and_export_failures() {
assert_eq!(
compile("(a b", "t.alc").unwrap_err().code,
tabnas_transduce::Code::DslParseError
);
let f = compile("def export [input] (nope input)", "t.alc").unwrap_err();
assert!(f.message.starts_with("unknown_name: "), "{f}");
let f = compile("def x 1", "t.alc").unwrap_err();
assert!(f.message.starts_with("no_export: "), "{f}");
}
}