Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion crates/agent/src/jobs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ where
}

#[allow(dead_code)]
async fn read_stdout(mut reader: async_process::ChildStdio) -> Result<Vec<u8>, Error> {
async fn read_stdout(mut reader: async_process::ChildOutput) -> Result<Vec<u8>, Error> {
let mut buffer = Vec::new();
let _ = tokio::io::copy(&mut reader, &mut buffer)
.await
Expand Down
43 changes: 25 additions & 18 deletions crates/async-process/src/lib.rs
Original file line number Diff line number Diff line change
@@ -1,18 +1,16 @@
pub use std::process::{Command, Output, Stdio};
use tokio::io::{AsyncReadExt, AsyncWriteExt};

#[cfg(unix)]
use std::os::fd::OwnedFd as OwnedImpl;
#[cfg(windows)]
use std::os::fd::OwnedHandle as OwnedImpl;
#[cfg(not(unix))]
compile_error!("async-process supports only unix targets");

pub struct Child {
pid: libc::pid_t,
status: Status,

pub stdin: Option<ChildStdio>,
pub stdout: Option<ChildStdio>,
pub stderr: Option<ChildStdio>,
pub stdin: Option<ChildStdin>,
pub stdout: Option<ChildOutput>,
pub stderr: Option<ChildOutput>,
}

/// State machine for reaping the child process.
Expand All @@ -32,13 +30,23 @@ enum Status {
Abandoned,
}

pub type ChildStdio = tokio::fs::File;
/// A child's stdin pipe.
///
/// Child pipes are non-blocking and driven by the IO driver of the runtime
/// that was current at `Child::from`, and must be polled while it's alive.
///
/// `flush` and `shutdown` are no-ops, as writes go directly to the pipe.
/// Dropping a `ChildStdin` is what closes it, sending EOF to the child.
pub type ChildStdin = tokio::net::unix::pipe::Sender;

/// A child's stdout or stderr pipe. See [`ChildStdin`].
pub type ChildOutput = tokio::net::unix::pipe::Receiver;

impl From<std::process::Child> for Child {
fn from(mut inner: std::process::Child) -> Self {
let stdin = map_stdio(inner.stdin.take());
let stdout = map_stdio(inner.stdout.take());
let stderr = map_stdio(inner.stderr.take());
let stdin = map_stdin(inner.stdin.take());
let stdout = map_output(inner.stdout.take());
let stderr = map_output(inner.stderr.take());

let pid = inner.id() as libc::pid_t;
let status = Status::Reaping(tokio::runtime::Handle::current().spawn_blocking(move || {
Expand Down Expand Up @@ -198,13 +206,12 @@ pub async fn input_output(cmd: &mut Command, input: &[u8]) -> std::io::Result<Ou
})
}

fn map_stdio<F>(f: Option<F>) -> Option<ChildStdio>
where
F: Into<OwnedImpl>,
{
let f: Option<OwnedImpl> = f.map(Into::into);
let f: Option<std::fs::File> = f.map(Into::into);
f.map(Into::into)
fn map_stdin(f: Option<std::process::ChildStdin>) -> Option<ChildStdin> {
Some(ChildStdin::from_owned_fd(f?.into()).expect("piped child stdin is a pipe"))
}

fn map_output<F: Into<std::os::fd::OwnedFd>>(f: Option<F>) -> Option<ChildOutput> {
Some(ChildOutput::from_owned_fd(f?.into()).expect("piped child output is a pipe"))
}

#[cfg(test)]
Expand Down
2 changes: 1 addition & 1 deletion crates/connector-init/src/rpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -202,7 +202,7 @@ where
/// the connector is already gone. Logging it at `warn` puts pure downstream noise beside
/// the causal error, so it's logged at `debug`. Any other error is unexpected and remains
/// a `warn`.
async fn write_stdin(stdin: &mut async_process::ChildStdio, buffer: &[u8]) {
async fn write_stdin(stdin: &mut async_process::ChildStdin, buffer: &[u8]) {
let Err(error) = stdin.write_all(buffer).await else {
return;
};
Expand Down
2 changes: 1 addition & 1 deletion crates/control-plane-api/src/jobs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,7 @@ where
}

#[allow(dead_code)]
async fn read_stdout(mut reader: async_process::ChildStdio) -> Result<Vec<u8>, Error> {
async fn read_stdout(mut reader: async_process::ChildOutput) -> Result<Vec<u8>, Error> {
let mut buffer = Vec::new();
let _ = tokio::io::copy(&mut reader, &mut buffer)
.await
Expand Down
33 changes: 33 additions & 0 deletions crates/json/src/validator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,9 @@ where
keywords: &'s [Keyword<A>],
// Bit flags of this Frame.
flags: u8,
// KIND_* classes of `keywords`, which let per-node passes skip Frames
// having no keyword of interest without scanning them.
kinds: u8,
// Counter retains:
// * The current index into Keyword::Properties (when an object).
// * The number of valid Keyword::Contains applications (when an array).
Expand All @@ -93,6 +96,35 @@ const FLAG_VALID_ANY_OF: u8 = 0x08;
// FLAG_VALID_ONE_OF is set if a Keyword::OneOf in-place application validated.
const FLAG_VALID_ONE_OF: u8 = 0x10;

// Keyword classes of a Frame's schema. Validation visits each document node
// with every active Frame, and most Frames (in-place applications especially)
// have no keyword relevant to a given visit: these classes let a visit skip
// them outright. Frame retains these eight as a u8 (keeping Frame within a
// cache line); the u16 classes beyond them only matter while winding the
// Frame, and their type keeps them from being tested against Frame::kinds.
//
// Keywords applied to object properties.
const KIND_PROPERTIES: u8 = 0x01;
// Keyword::PropertyNames.
const KIND_PROPERTY_NAMES: u8 = 0x02;
// Keywords applied to array items.
const KIND_ITEMS: u8 = 0x04;
// Keywords checked against string nodes.
const KIND_STRING: u8 = 0x08;
// Keywords checked against numeric nodes.
const KIND_NUMBER: u8 = 0x10;
// Keywords checked against a whole object or array, after its children.
// Includes Keyword::Properties, which checks for missing required properties.
const KIND_CONTAINER: u8 = 0x20;
// Keywords checked against every node.
const KIND_NODE: u8 = 0x40;
// Keyword::AnyOf or OneOf, which are checked as the Frame unwinds.
const KIND_ANY_OR_ONE_OF: u8 = 0x80;
// Keywords having in-place applications.
const KIND_IN_PLACE: u16 = 0x0100;
// Keyword::UnevaluatedItems or UnevaluatedProperties.
const KIND_UNEVALUATED: u16 = 0x0200;

struct FrameSpeculative<'s, A>
where
A: Annotation,
Expand Down Expand Up @@ -158,6 +190,7 @@ where
parent_keyword: None,
keywords: &[],
flags: 0,
kinds: 0,
outcomes: Vec::new(),
counter: 0,
speculative: None,
Expand Down
Loading
Loading