use std::collections::{BTreeMap, HashSet};
use std::io::{self, Read};
use std::path::Path;
use crate::edge::Edge;
use crate::lifecycle::BaseChange;
use crate::message::Message;
use crate::mutate::{self, Op};
use crate::reads::legacy;
use crate::task::Task;
use crate::taskfile::{exists, write_task};
use crate::verb::Verb;
use crate::wire::Command;
#[path = "import_stream.rs"]
mod stream;
pub fn run(edge: &Edge, input: &mut dyn Read, args: &[String]) -> io::Result<()> {
let flags = parse(args, &edge.default_actor)?;
if let Some(spec) = &flags.legacy {
return button(edge, &flags, spec);
}
let mut text = String::new();
input.read_to_string(&mut text)?;
ingest(edge, &flags, stream::records(&text)?)
}
#[derive(Debug)]
struct Flags {
actor: String,
message: Option<String>,
remote: Option<String>,
legacy: Option<String>,
}
fn parse(args: &[String], default_actor: &str) -> io::Result<Flags> {
let mut f = Flags { actor: default_actor.to_string(), message: None, remote: None, legacy: None };
let mut args = args.iter();
while let Some(arg) = args.next() {
match arg.as_str() {
"--as" => f.actor = value(&mut args, "--as")?,
"-m" | "--message" => f.message = Some(value(&mut args, "-m")?),
"--remote" => f.remote = Some(value(&mut args, "--remote")?),
"--center" => {
let center = value(&mut args, "--center")?;
f.remote.get_or_insert(center);
}
a if legacy::flag(a).is_some() => f.legacy = legacy::flag(a),
other => {
return Err(io::Error::other(format!(
"import: unexpected argument '{other}' — records ride stdin"
)))
}
}
}
Ok(f)
}
fn value(args: &mut std::slice::Iter<'_, String>, flag: &str) -> io::Result<String> {
args.next().cloned().ok_or_else(|| io::Error::other(format!("import: {flag} needs a value")))
}
fn ingest(edge: &Edge, flags: &Flags, balls: Vec<(String, Task)>) -> io::Result<()> {
let n = balls.len();
let base = Ingest { balls, actor: flags.actor.clone(), message: flags.message.clone() };
let op = Op {
actor: flags.actor.clone(),
remote: flags.remote.clone(),
command: Command { op: Verb::Import.token().to_string(), body_change: None },
};
mutate::seal_op(edge, Verb::Import, &op, &base, None)?;
eprintln!("import {n} ball{}", if n == 1 { "" } else { "s" });
Ok(())
}
struct Ingest {
balls: Vec<(String, Task)>,
actor: String,
message: Option<String>,
}
impl BaseChange for Ingest {
fn stage(&self, dir: &Path) -> io::Result<()> {
let mut seen = HashSet::new();
let held: Vec<&str> = self
.balls
.iter()
.map(|(id, _)| id.as_str())
.filter(|id| !seen.insert(*id) || exists(dir, id))
.collect();
if !held.is_empty() {
return Err(io::Error::new(
io::ErrorKind::AlreadyExists,
format!(
"import: id(s) already held: {} — refusing the whole stream, nothing imported (retire the holder or strip the record, then re-run)",
held.join(", ")
),
));
}
for (id, task) in &self.balls {
write_task(dir, id, task)?;
}
Ok(())
}
fn finalize(&self, _dir: &Path) -> io::Result<String> {
let n = self.balls.len();
Message {
verb: Verb::Import,
actor: self.actor.clone(),
id: None,
subject: format!("import {n} ball{}", if n == 1 { "" } else { "s" }),
body: self.message.clone(),
}
.render()
}
fn narrated(&self) -> bool {
self.message.is_some()
}
}
fn button(edge: &Edge, flags: &Flags, spec: &str) -> io::Result<()> {
let balls = legacy::balls(&edge.invocation_path, spec)?;
let mut edges: BTreeMap<String, Vec<String>> = BTreeMap::new();
for (id, task) in &balls {
if let Some(parent) = &task.parent {
edges.entry(parent.clone()).or_default().push(id.clone());
}
}
ingest(edge, flags, balls)?;
for (parent, children) in edges {
let mut args = vec![parent];
for child in children {
args.push("--needs".into());
args.push(child);
}
args.extend(["--as".into(), flags.actor.clone()]);
if let Some(remote) = &flags.remote {
args.extend(["--remote".into(), remote.clone()]);
}
mutate::run(edge, Verb::Update, &args)?;
}
Ok(())
}
#[cfg(test)]
#[path = "import_tests.rs"]
mod tests;