/examples/subprocess-pipeline/pipeline.nct
pipeline.nct
use std/io
use std/io.Reader
use std/process.{ChildStdin, ChildStdout, Command, ProcessIo, Stdio}
use std/string.String
use std/task
use std/task.Timeout
use std/time.Duration
use std/vec.Vec
const PIPELINE_RECORDS: usize = 8192
const PIPELINE_BYTES: usize = PIPELINE_RECORDS * 9
async func forward(source: ChildStdout, destination: ChildStdin): usize! {
var reader = move source
var writer = move destination
let transferred = await io.copy(&+reader, &+writer) catch failure {
reader.close()
writer.close()
return move failure
}
reader.close()
writer.close()
return transferred
}
async func capture<R>(source: R): Vec<u8>! where R impl Reader {
var reader = move source
return await reader.read_to_end()?
}
async func run(): i32! {
var producer_io = ProcessIo.inherit()
producer_io.stdout(Stdio.pipe)
producer_io.stderr(Stdio.pipe)
let producer_command = Command.new("./producer.sh")?
var producer = await producer_command.spawn(move producer_io)?
let producer_stdout = producer.take_stdout() otherwise {
return error.new("example.pipeline.endpoint", "producer stdout was not piped")
}
let producer_stderr = producer.take_stderr() otherwise {
return error.new("example.pipeline.endpoint", "producer stderr was not piped")
}
let consumer_command = Command.new("./consumer.sh")?
var consumer = await consumer_command.spawn(ProcessIo.piped())?
let consumer_stdin = consumer.take_stdin() otherwise {
return error.new("example.pipeline.endpoint", "consumer stdin was not piped")
}
let consumer_stdout = consumer.take_stdout() otherwise {
return error.new("example.pipeline.endpoint", "consumer stdout was not piped")
}
let consumer_stderr = consumer.take_stderr() otherwise {
return error.new("example.pipeline.endpoint", "consumer stderr was not piped")
}
let producer_work = task.join(
task.join(
forward(move producer_stdout, move consumer_stdin),
capture(move producer_stderr),
),
producer.wait(),
)
let consumer_work = task.join(
task.join(
capture(move consumer_stdout),
capture(move consumer_stderr),
),
consumer.wait(),
)
let bounded = await task.with_timeout(
task.join(move producer_work, move consumer_work),
Duration.from_seconds(5),
)
let completed = match move bounded {
Timeout.completed(value) { move value }
Timeout.elapsed {
return error.new("example.pipeline.timeout", "subprocess pipeline timed out")
}
}
let producer_result = move completed.0
let producer_streams = move producer_result.0
let transfer_result = move producer_streams.0
let producer_stderr_result = move producer_streams.1
let producer_status_result = move producer_result.1
let transferred = move transfer_result?
let producer_diagnostics = move producer_stderr_result?
let producer_status = move producer_status_result?
let consumer_result = move completed.1
let consumer_streams = move consumer_result.0
let consumer_stdout_result = move consumer_streams.0
let consumer_stderr_result = move consumer_streams.1
let consumer_status_result = move consumer_result.1
let consumer_output = move consumer_stdout_result?
let consumer_diagnostics = move consumer_stderr_result?
let consumer_status = move consumer_status_result?
if !producer_status.success() || !consumer_status.success() {
return error.new("example.pipeline.status", "one pipeline child failed")
}
if transferred != PIPELINE_BYTES || consumer_output.len() != PIPELINE_BYTES
|| producer_diagnostics.len() != PIPELINE_BYTES {
return error.new("example.pipeline.bytes", "pipeline byte counts changed")
}
if consumer_output[0] != 112 {
return error.new("example.pipeline.content", "consumer output prefix changed")
}
let consumer_last = consumer_output.len() - 1
if consumer_output[consumer_last] != 10 {
return error.new("example.pipeline.content", "consumer output suffix changed")
}
if producer_diagnostics[0] != 112 {
return error.new("example.pipeline.content", "producer diagnostics prefix changed")
}
let producer_last = producer_diagnostics.len() - 1
if producer_diagnostics[producer_last] != 10 {
return error.new("example.pipeline.content", "producer diagnostics suffix changed")
}
let summary = String.from_utf8(&consumer_diagnostics)?
if (&summary as &str) != "count=8192\n" {
return error.new("example.pipeline.summary", "consumer record count changed")
}
return 0
}