Programming Language

Nocter

A self-contained systems language built around simplicity, encapsulation, and foolproof design.

/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
}