use crate::header::FSMetaV1;
use crate::jobfile;
use crate::jobqueue;
use crate::seqfile;
use anyhow::bail;
use clap::Args;
use std::ffi::OsString;
use std::io::Write;
use std::os::unix::ffi::OsStringExt;
#[derive(Debug, Args)]
pub struct QueueOpts {
#[arg(short, long, value_name = "DIR")]
pub queuedir: OsString,
}
#[derive(Debug, Args)]
pub struct QueueInitOpts {
#[clap(flatten)]
pub qopts: QueueOpts,
#[arg(long)]
pub append_only: bool,
}
#[derive(Debug, Args)]
pub struct QueueOptsWithDecoder {
#[clap(flatten)]
pub qopts: QueueOpts,
#[arg(long, short, value_name = "DECODECMD")]
pub decoder: Option<OsString>,
}
#[derive(Debug, Args)]
pub struct QueueJobOpts {
#[clap(flatten)]
pub qdopts: QueueOptsWithDecoder,
#[arg(short, long, value_name = "ID")]
pub job: u64,
}
#[derive(Debug, Args)]
pub struct QueueSetSeqOpts {
#[clap(flatten)]
pub qopts: QueueOpts,
#[arg(value_name = "ID")]
pub nextid: u64,
}
pub fn cmd_queueinit(opts: &QueueInitOpts) -> Result<(), anyhow::Error> {
jobqueue::queueinit(&opts.qopts.queuedir, opts.append_only)
}
pub fn cmd_queuewrite(opts: &QueueOpts) -> Result<(), anyhow::Error> {
let mut stdin_fd = std::io::stdin();
jobqueue::queuewrite(&mut stdin_fd, &opts.queuedir)
}
pub fn cmd_queuels(opts: &QueueOptsWithDecoder) -> Result<(), anyhow::Error> {
let mut entries: Vec<(OsString, FSMetaV1)> =
jobqueue::scanqueue(&opts.qopts.queuedir, &opts.decoder)?
.flatten()
.collect();
entries.sort_by(|(_, ma), (_, mb)| ma.seq.cmp(&mb.seq));
let mut out = std::io::stdout();
println!("{:20} {:27} filename", "ID", "creation timestamp");
for (filename, meta) in entries.into_iter() {
print!(
"{:<20} {:27} ",
meta.seq,
meta.get_datestring_rfc3339_local(),
);
out.write_all(&filename.into_vec())?;
out.write_all(b"\n")?;
}
Ok(())
}
pub fn cmd_queueinfo(opts: &QueueJobOpts) -> Result<(), anyhow::Error> {
let qmap = jobqueue::scanqueue_map(&opts.qdopts.qopts.queuedir, &opts.qdopts.decoder)?;
match qmap.get(&opts.job) {
Some((filename, meta)) => {
meta.render(
Some(filename),
Some(&opts.qdopts.qopts.queuedir),
&mut std::io::stdout(),
)?;
}
None => bail!(
"Job {} not found in directory {:?}",
opts.job,
opts.qdopts.qopts.queuedir
),
}
Ok(())
}
pub fn cmd_queuepayload(opts: &QueueJobOpts) -> Result<(), anyhow::Error> {
let qmap = jobqueue::scanqueue_map(&opts.qdopts.qopts.queuedir, &opts.qdopts.decoder)?;
match qmap.get(&opts.job) {
Some((filename, _meta)) => {
let mut inputinfo = jobqueue::queue_openjob(
&opts.qdopts.qopts.queuedir,
filename,
&opts.qdopts.decoder,
)?;
let mut inhandle = inputinfo.reader.take().unwrap().as_read();
let _meta = jobfile::read_jobfile_header(&mut inhandle)?;
std::io::copy(&mut inhandle, &mut std::io::stdout())?;
}
None => bail!(
"Job {} not found in directory {:?}",
opts.job,
opts.qdopts.qopts.queuedir
),
}
Ok(())
}
pub fn cmd_queuegetnext(opts: &QueueOpts) -> Result<(), anyhow::Error> {
let seqfn = jobqueue::get_seqfile(&opts.queuedir);
let mut lock = seqfile::prepare_seqfile_lock(&seqfn, false)?;
let seqf = seqfile::SeqFile::open(&seqfn, &mut lock)?;
println!("{}", seqf.get_next());
Ok(())
}
pub fn cmd_queuesetnext(opts: &QueueSetSeqOpts) -> Result<(), anyhow::Error> {
let seqfn = jobqueue::get_seqfile(&opts.qopts.queuedir);
let mut lock = seqfile::prepare_seqfile_lock(&seqfn, false)?;
let mut seqf = seqfile::SeqFile::open(&seqfn, &mut lock)?;
seqf.set(opts.nextid)?;
Ok(())
}