Merge nucleic/brisk-yarn-dingo-50se into dev
This commit is contained in:
@@ -91,6 +91,18 @@ pub struct PipeEvent {
|
|||||||
pub from_index: usize,
|
pub from_index: usize,
|
||||||
/// Zero-based index of the consuming stage.
|
/// Zero-based index of the consuming stage.
|
||||||
pub to_index: usize,
|
pub to_index: usize,
|
||||||
|
/// The producing stage as written, rendered from its AST at spawn time.
|
||||||
|
///
|
||||||
|
/// Carried on the link itself so an observer never has to *infer* which command fed a pipe
|
||||||
|
/// from the exec events around it: a stage that is a compound command (`{ a; b; } | c`) emits
|
||||||
|
/// one exec per command inside it, and none of them says which owned the pipe. Rendering the
|
||||||
|
/// stage answers it outright — and does so at spawn, before the producer has even exited.
|
||||||
|
///
|
||||||
|
/// It is also the *pre-expansion* text, so `curl -H "Bearer $TOKEN" | sh` names the stage with
|
||||||
|
/// `$TOKEN` rather than the expanded secret an argv-derived name would carry.
|
||||||
|
pub from_text: String,
|
||||||
|
/// The consuming stage as written. See [`PipeEvent::from_text`].
|
||||||
|
pub to_text: String,
|
||||||
/// Total bytes that crossed the link.
|
/// Total bytes that crossed the link.
|
||||||
pub total_bytes: u64,
|
pub total_bytes: u64,
|
||||||
/// A bounded prefix of the bytes (observer applies its own caps/redaction).
|
/// A bounded prefix of the bytes (observer applies its own caps/redaction).
|
||||||
|
|||||||
@@ -518,7 +518,17 @@ async fn spawn_pipeline_processes(
|
|||||||
let (b_read, b_write) = std::io::pipe()?;
|
let (b_read, b_write) = std::io::pipe()?;
|
||||||
// Link `link` connects producer stage (len-2-link) → consumer (len-1-link).
|
// Link `link` connects producer stage (len-2-link) → consumer (len-1-link).
|
||||||
let from_index = pipeline_len - 2 - link;
|
let from_index = pipeline_len - 2 - link;
|
||||||
nash_spawn_pipe_tee(a_read, b_write, pipeline_id, from_index, from_index + 1);
|
// Name both ends from their AST *here*, where the stage is unambiguous. Deriving
|
||||||
|
// the name later from exec events cannot distinguish which command inside a
|
||||||
|
// compound stage owned the pipe, and cannot name either end until it exits.
|
||||||
|
nash_spawn_pipe_tee(
|
||||||
|
a_read,
|
||||||
|
b_write,
|
||||||
|
pipeline_id,
|
||||||
|
from_index,
|
||||||
|
nash_stage_text(pipeline, from_index),
|
||||||
|
nash_stage_text(pipeline, from_index + 1),
|
||||||
|
);
|
||||||
pipe_readers.push(Some(openfiles::OpenFile::PipeReader(b_read)));
|
pipe_readers.push(Some(openfiles::OpenFile::PipeReader(b_read)));
|
||||||
pipe_writers.push(Some(openfiles::OpenFile::PipeWriter(a_write)));
|
pipe_writers.push(Some(openfiles::OpenFile::PipeWriter(a_write)));
|
||||||
continue;
|
continue;
|
||||||
@@ -1964,17 +1974,45 @@ pub(crate) async fn setup_redirect(
|
|||||||
// nash: bounded prefix captured per pipe link before reporting (docs/NASH.md §5.3).
|
// nash: bounded prefix captured per pipe link before reporting (docs/NASH.md §5.3).
|
||||||
const NASH_PIPE_CAPTURE_CAP: usize = 64 * 1024;
|
const NASH_PIPE_CAPTURE_CAP: usize = 64 * 1024;
|
||||||
|
|
||||||
|
// nash: longest stage label carried on a pipe event. A stage can be an entire `while` loop, and
|
||||||
|
// this rides on every link of every pipeline — one of the highest-volume kinds nash emits.
|
||||||
|
const NASH_STAGE_TEXT_CAP: usize = 200;
|
||||||
|
|
||||||
|
// nash: the pipeline stage at `index`, rendered back to shell syntax to label a pipe link
|
||||||
|
// (docs/NASH.md §5.3).
|
||||||
|
//
|
||||||
|
// Collapsing and capping happen here rather than in the observer because this is a *label*, and
|
||||||
|
// only brush-core has the AST to render one; redaction still happens observer-side with every
|
||||||
|
// other outbound string. `Display` re-renders the parsed command, so the result is the stage as
|
||||||
|
// parsed rather than byte-identical source — which is what a one-line label wants anyway.
|
||||||
|
fn nash_stage_text(pipeline: &ast::Pipeline, index: usize) -> String {
|
||||||
|
let Some(command) = pipeline.seq.get(index) else {
|
||||||
|
return String::new();
|
||||||
|
};
|
||||||
|
// A compound stage renders multi-line; a label has one line.
|
||||||
|
let collapsed = command.to_string();
|
||||||
|
let mut text = collapsed.split_whitespace().collect::<Vec<_>>().join(" ");
|
||||||
|
if text.chars().count() > NASH_STAGE_TEXT_CAP {
|
||||||
|
text = text.chars().take(NASH_STAGE_TEXT_CAP).collect::<String>();
|
||||||
|
text.push('…');
|
||||||
|
}
|
||||||
|
text
|
||||||
|
}
|
||||||
|
|
||||||
// nash: spawn a copier that moves bytes from `src` (producer's output) to `dst`
|
// nash: spawn a copier that moves bytes from `src` (producer's output) to `dst`
|
||||||
// (consumer's input), mirroring a bounded prefix, then reports the link to the
|
// (consumer's input), mirroring a bounded prefix, then reports the link to the
|
||||||
// gate (docs/NASH.md §5.3). The synchronous copy preserves pipe backpressure;
|
// gate (docs/NASH.md §5.3). The synchronous copy preserves pipe backpressure;
|
||||||
// EOF on `src` or EPIPE on `dst` (consumer gone, e.g. `yes | head`) ends it and
|
// EOF on `src` or EPIPE on `dst` (consumer gone, e.g. `yes | head`) ends it and
|
||||||
// closes `dst` so the consumer sees EOF.
|
// closes `dst` so the consumer sees EOF.
|
||||||
|
// `to_index` is always `from_index + 1` — a link joins adjacent stages — so it is derived rather
|
||||||
|
// than passed; `from_text`/`to_text` are the two stages already rendered by `nash_stage_text`.
|
||||||
fn nash_spawn_pipe_tee(
|
fn nash_spawn_pipe_tee(
|
||||||
mut src: std::io::PipeReader,
|
mut src: std::io::PipeReader,
|
||||||
mut dst: std::io::PipeWriter,
|
mut dst: std::io::PipeWriter,
|
||||||
pipeline_id: u64,
|
pipeline_id: u64,
|
||||||
from_index: usize,
|
from_index: usize,
|
||||||
to_index: usize,
|
from_text: String,
|
||||||
|
to_text: String,
|
||||||
) {
|
) {
|
||||||
std::thread::Builder::new()
|
std::thread::Builder::new()
|
||||||
.name("nash-pipe-tee".into())
|
.name("nash-pipe-tee".into())
|
||||||
@@ -2013,7 +2051,9 @@ fn nash_spawn_pipe_tee(
|
|||||||
crate::gate::gate().on_pipe(crate::gate::PipeEvent {
|
crate::gate::gate().on_pipe(crate::gate::PipeEvent {
|
||||||
pipeline_id,
|
pipeline_id,
|
||||||
from_index,
|
from_index,
|
||||||
to_index,
|
to_index: from_index + 1,
|
||||||
|
from_text,
|
||||||
|
to_text,
|
||||||
total_bytes: total,
|
total_bytes: total,
|
||||||
captured,
|
captured,
|
||||||
truncated,
|
truncated,
|
||||||
|
|||||||
Reference in New Issue
Block a user