Merge nucleic/brisk-yarn-dingo-50se into dev
This commit is contained in:
@@ -368,6 +368,7 @@ impl<'a, SE: extensions::ShellExtensions> SimpleCommand<'a, SE> {
|
|||||||
argv,
|
argv,
|
||||||
cwd: self.shell.working_dir().to_path_buf(),
|
cwd: self.shell.working_dir().to_path_buf(),
|
||||||
redirects,
|
redirects,
|
||||||
|
pipeline: self.params.nash_pipeline,
|
||||||
};
|
};
|
||||||
let token = match gate.on_exec(start_event) {
|
let token = match gate.on_exec(start_event) {
|
||||||
crate::gate::Verdict::Allow(token) => token,
|
crate::gate::Verdict::Allow(token) => token,
|
||||||
@@ -835,6 +836,11 @@ pub(crate) async fn invoke_command_in_subshell_and_get_output(
|
|||||||
// Get our own set of parameters we can customize and use.
|
// Get our own set of parameters we can customize and use.
|
||||||
let mut params = params.clone();
|
let mut params = params.clone();
|
||||||
params.process_group_policy = ProcessGroupPolicy::SameProcessGroup;
|
params.process_group_policy = ProcessGroupPolicy::SameProcessGroup;
|
||||||
|
// nash: a command substitution is an *argument* to the stage, not the stage itself, and it
|
||||||
|
// runs during expansion — i.e. before the stage's own command. Letting it inherit the slot
|
||||||
|
// made `grep "$(printf x)" f | wc -l` report `printf x` as the producing stage
|
||||||
|
// (docs/NASH.md §5.1).
|
||||||
|
params.nash_pipeline = None;
|
||||||
|
|
||||||
// Set up pipe so we can read the output.
|
// Set up pipe so we can read the output.
|
||||||
let (reader, writer) = std::io::pipe()?;
|
let (reader, writer) = std::io::pipe()?;
|
||||||
|
|||||||
@@ -22,6 +22,24 @@ pub struct ExecStart {
|
|||||||
/// nash: the redirections applied to this command (docs/NASH.md §5.2). The
|
/// nash: the redirections applied to this command (docs/NASH.md §5.2). The
|
||||||
/// observer reads back the affected byte range after the command completes.
|
/// observer reads back the affected byte range after the command completes.
|
||||||
pub redirects: Vec<RedirectRecord>,
|
pub redirects: Vec<RedirectRecord>,
|
||||||
|
/// nash: the pipeline stage this command occupies, when it runs as part of
|
||||||
|
/// `a | b | c` (docs/NASH.md §5.1). `None` for a standalone command.
|
||||||
|
pub pipeline: Option<PipelineSlot>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Where a command sits in the pipeline it belongs to (docs/NASH.md §5.1).
|
||||||
|
///
|
||||||
|
/// Shares its `id` with the [`PipeEvent`]s of the same pipeline, which is what
|
||||||
|
/// lets an observer name the commands either side of a link (`curl … | sh`)
|
||||||
|
/// instead of reporting bare stage indices.
|
||||||
|
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||||
|
pub struct PipelineSlot {
|
||||||
|
/// Correlates this stage with the pipeline's links; see [`PipeEvent::pipeline_id`].
|
||||||
|
pub id: u64,
|
||||||
|
/// Zero-based position of this stage within the pipeline.
|
||||||
|
pub index: usize,
|
||||||
|
/// Total number of stages in the pipeline.
|
||||||
|
pub len: usize,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// How a redirection's data is captured after the command runs (docs/NASH.md §5.2).
|
/// How a redirection's data is captured after the command runs (docs/NASH.md §5.2).
|
||||||
|
|||||||
@@ -41,6 +41,12 @@ pub struct ExecutionParameters {
|
|||||||
// `gate::ExecStart` for post-run read-back (docs/NASH.md §5.2). Empty unless
|
// `gate::ExecStart` for post-run read-back (docs/NASH.md §5.2). Empty unless
|
||||||
// observation is active and the command has redirects.
|
// observation is active and the command has redirects.
|
||||||
pub(crate) nash_redirects: Vec<crate::gate::RedirectRecord>,
|
pub(crate) nash_redirects: Vec<crate::gate::RedirectRecord>,
|
||||||
|
// nash: the pipeline stage the command being built belongs to, stamped onto
|
||||||
|
// its `gate::ExecStart` so the observer can name the stages either side of a
|
||||||
|
// pipe link (docs/NASH.md §5.1). Inherited by nested commands — a command
|
||||||
|
// inside a stage's subshell or brace group belongs to that same stage — and
|
||||||
|
// overwritten only when a nested pipeline starts its own.
|
||||||
|
pub(crate) nash_pipeline: Option<crate::gate::PipelineSlot>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl ExecutionParameters {
|
impl ExecutionParameters {
|
||||||
@@ -488,19 +494,23 @@ async fn spawn_pipeline_processes(
|
|||||||
let mut spawn_results = VecDeque::new();
|
let mut spawn_results = VecDeque::new();
|
||||||
let mut process_group_id: Option<i32> = None;
|
let mut process_group_id: Option<i32> = None;
|
||||||
|
|
||||||
|
// nash: when observation is active, one id per multi-stage pipeline correlates
|
||||||
|
// both its tapped links and the exec events of the stages either side of them
|
||||||
|
// (docs/NASH.md §5.3). `None` on the unobserved fast path and for a lone
|
||||||
|
// command, where there is no link to tap and nothing to correlate.
|
||||||
|
let nash_pipeline_id = if pipeline_len > 1 && crate::gate::is_active() {
|
||||||
|
Some(crate::gate::next_pipeline_id())
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
};
|
||||||
|
|
||||||
// Create pipes to use between commands, but only bother doing so if there's more than one
|
// Create pipes to use between commands, but only bother doing so if there's more than one
|
||||||
// command.
|
// command.
|
||||||
if pipeline_len > 1 {
|
if pipeline_len > 1 {
|
||||||
// nash: when observation is active, interpose a tee on each link so the
|
|
||||||
// bytes flowing `a | b` are captured (docs/NASH.md §5.3). Off otherwise —
|
|
||||||
// no extra pipe/thread on the unobserved fast path.
|
|
||||||
let nash_pipeline_id = if crate::gate::is_active() {
|
|
||||||
Some(crate::gate::next_pipeline_id())
|
|
||||||
} else {
|
|
||||||
None
|
|
||||||
};
|
|
||||||
for link in 0..(pipeline_len - 1) {
|
for link in 0..(pipeline_len - 1) {
|
||||||
if let Some(pipeline_id) = nash_pipeline_id {
|
if let Some(pipeline_id) = nash_pipeline_id {
|
||||||
|
// nash: interpose a tee on the link so the bytes flowing `a | b` are
|
||||||
|
// captured (docs/NASH.md §5.3).
|
||||||
// Producer writes `a_write`; a copier moves bytes to `b_write`;
|
// Producer writes `a_write`; a copier moves bytes to `b_write`;
|
||||||
// consumer reads `b_read`. Slots stay identical to the untapped
|
// consumer reads `b_read`. Slots stay identical to the untapped
|
||||||
// case, so the pop-ordering below is unchanged.
|
// case, so the pop-ordering below is unchanged.
|
||||||
@@ -539,6 +549,17 @@ async fn spawn_pipeline_processes(
|
|||||||
// Set up parameters appropriate for this command.
|
// Set up parameters appropriate for this command.
|
||||||
let mut cmd_params = params.clone();
|
let mut cmd_params = params.clone();
|
||||||
|
|
||||||
|
// nash: stamp this stage's slot so its exec event can be joined to the links
|
||||||
|
// either side of it. Only when this pipeline is itself tapped — a lone command
|
||||||
|
// must keep any slot inherited from the stage that contains it.
|
||||||
|
if let Some(id) = nash_pipeline_id {
|
||||||
|
cmd_params.nash_pipeline = Some(crate::gate::PipelineSlot {
|
||||||
|
id,
|
||||||
|
index: current_pipeline_index,
|
||||||
|
len: pipeline_len,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
// Install pipes.
|
// Install pipes.
|
||||||
if let Some(Some(reader)) = pipe_readers.pop() {
|
if let Some(Some(reader)) = pipe_readers.pop() {
|
||||||
cmd_params.open_files.set_fd(OpenFiles::STDIN_FD, reader);
|
cmd_params.open_files.set_fd(OpenFiles::STDIN_FD, reader);
|
||||||
|
|||||||
Reference in New Issue
Block a user