use std::io::Write;
use std::process::ExitCode;
use trex::encoding::Incremental;
use trex::files::{Edit, LineIndex, Source};
use trex::follow::{Followed as Change, Follower};
use trex::records::RecordUnit;
use trex::report::Members;
use trex::window::{Asked, Select};
use crate::cli_files::{ScanHow, Scanning, scan_found, whole_line_matches};
use crate::cli_window::{Part, read_part};
pub(crate) struct StreamPrint<'a> {
pub(crate) json: bool,
pub(crate) format: Option<&'a trex::Template>,
pub(crate) painter: &'a trex::paint::Painter,
pub(crate) lists: bool,
pub(crate) values: Option<&'a crate::ValueView<'a>>,
pub(crate) prefixed: bool,
pub(crate) whole_line: bool,
pub(crate) max_count: Option<usize>,
pub(crate) single: bool,
}
pub(crate) struct StreamRun<'a> {
scanning: &'a Scanning,
shapes: &'a trex::ShapeSet,
name: &'a str,
stream: trex::HeldStream,
printed: usize,
seen: Vec<bool>,
}
impl<'a> StreamRun<'a> {
pub(crate) fn new(
scanning: &'a Scanning,
shapes: &'a trex::ShapeSet,
name: &'a str,
origin: usize,
lines_before: Option<usize>,
) -> Self {
let scanner = match scanning {
Scanning::One(p) => trex::StreamScanner::with_shapes(p.clone(), shapes.clone()),
Scanning::Set(set) => trex::StreamScanner::over_set((**set).clone()),
};
let members = scanning.set().map_or(1, trex::PatternSet::len);
StreamRun {
scanning,
shapes,
name,
stream: trex::HeldStream::new(scanner, origin, lines_before),
printed: 0,
seen: vec![false; members],
}
}
pub(crate) fn commits_early(&self) -> bool {
self.stream.commits_early()
}
pub(crate) fn done(&self, p: &StreamPrint<'_>) -> bool {
p.max_count.is_some_and(|n| self.printed >= n) || (p.single && self.seen.iter().all(|&s| s))
}
pub(crate) fn push(&mut self, bytes: &[u8], p: &StreamPrint<'_>, out: &mut dyn FnMut(&str)) {
let committed = self.stream.push(bytes);
let view = Held {
scanning: self.scanning,
shapes: self.shapes,
name: self.name,
held: self.stream.held(),
base: self.stream.base(),
lines_before: self.stream.lines_before(),
};
view.print(&committed, p, &mut self.printed, &mut self.seen, out);
}
pub(crate) fn finish(self, p: &StreamPrint<'_>, out: &mut dyn FnMut(&str)) -> usize {
let StreamRun { scanning, shapes, name, stream, mut printed, mut seen } = self;
let ended = stream.finish();
let view = Held { scanning, shapes, name, held: &ended.held, base: ended.base, lines_before: ended.lines_before };
view.print(&ended.matches, p, &mut printed, &mut seen, out);
printed
}
}
struct Held<'b> {
scanning: &'b Scanning,
shapes: &'b trex::ShapeSet,
name: &'b str,
held: &'b [u8],
base: usize,
lines_before: Option<usize>,
}
impl Held<'_> {
fn print(
&self,
committed: &[(usize, trex::Span)],
p: &StreamPrint<'_>,
printed: &mut usize,
seen: &mut [bool],
out: &mut dyn FnMut(&str),
) {
if committed.is_empty() {
return;
}
let mut members_in: Vec<usize> = committed.iter().map(|&(member, _)| member).collect();
members_in.sort_unstable();
members_in.dedup();
let mut found: Vec<(usize, usize, trex::Match)> = Vec::with_capacity(committed.len());
for member in members_in {
let at: Vec<usize> = (0..committed.len()).filter(|&k| committed[k].0 == member).collect();
let spans: Vec<trex::Span> = at.iter().map(|&k| committed[k].1).collect();
let pattern = self.scanning.pattern(member);
let resolved = if p.lists {
trex::captures_with_shapes_and_lists(pattern, self.held, self.shapes, &spans)
} else {
trex::captures_with_shapes(pattern, self.held, self.shapes, &spans)
};
found.extend(at.into_iter().zip(resolved).map(|(k, m)| (k, member, m)));
}
found.sort_by_key(|&(k, _, _)| k);
let (mut members, mut matches): (Vec<usize>, Vec<trex::Match>) =
found.into_iter().map(|(_, member, m)| (member, m)).unzip();
let index = LineIndex::new(self.held).within(Some(self.base), self.lines_before);
if p.whole_line {
(matches, members) = whole_line_matches(self.held, matches, members, &index);
}
if p.single {
let first: Vec<bool> = members
.iter()
.map(|&m| {
let new = !seen[m];
seen[m] = true;
new
})
.collect();
matches = matches.into_iter().zip(&first).filter(|(_, f)| **f).map(|(m, _)| m).collect();
members = members.into_iter().zip(&first).filter(|(_, f)| **f).map(|(m, _)| m).collect();
}
if let Some(n) = p.max_count {
let room = n.saturating_sub(*printed);
matches.truncate(room);
members.truncate(room);
}
if matches.is_empty() {
return;
}
let set_members = self.scanning.set().map(|set| Members { set, of: &members });
if let Some(t) = p.format {
for (k, m) in matches.iter().enumerate() {
let (line, col) = index.line_col(self.held, m.start);
let name = set_members.map(|ms| ms.name(k));
let place = trex::ReportAt {
path: self.name,
line,
col,
base: index.base(),
pattern: name.as_deref(),
rule: None,
};
out(&t.render_report(m, self.held, &place));
}
} else if p.json {
for (k, m) in matches.iter().enumerate() {
let at = p.prefixed.then(|| {
let (line, col) = index.line_col(self.held, m.start);
(self.name, line, col)
});
let extra = trex::report::json_extras(m, set_members, k, None);
out(&crate::json_match_full(self.held, m, at, index.offset(0), extra.as_deref(), p.values));
}
} else if p.prefixed {
let context = trex::report::Context { before: 0, after: 0, record: None };
trex::report::context_report(
Some(self.name),
self.held,
&matches,
&index,
&context,
p.painter,
set_members,
None,
out,
);
} else {
trex::report::human_report(self.held, &matches, &index, p.painter, set_members, None, out);
}
*printed += matches.len();
}
}
pub(crate) struct StreamReport<'a> {
pub(crate) json: bool,
pub(crate) require_match: bool,
pub(crate) binary: bool,
pub(crate) format: Option<&'a trex::Template>,
pub(crate) painter: &'a trex::paint::Painter,
pub(crate) lists: bool,
pub(crate) style: trex::typed::ValueStyle,
}
pub(crate) fn stream_stdin(scanning: &Scanning, report: &StreamReport<'_>) -> ExitCode {
use std::io::BufRead;
let StreamReport { json, require_match, binary, format, painter, lists, style } = *report;
let capture_kinds = scanning.capture_kinds();
let value_view = crate::ValueView { kinds: &capture_kinds, style, clock: trex::Clock::current() };
let print = StreamPrint {
json,
format,
painter,
lists,
values: Some(&value_view),
prefixed: false,
whole_line: false,
max_count: None,
single: false,
};
let shapes = trex::ShapeSet::new();
let mut run = StreamRun::new(scanning, &shapes, "-", 0, Some(0));
let mut out = |l: &str| crate::out::line(l);
let stdin = std::io::stdin();
let mut reader = stdin.lock();
let mut total = 0usize;
let mut opening: Vec<u8> = Vec::new();
let mut opened = false;
loop {
let chunk = match reader.fill_buf() {
Ok(c) => c.to_vec(),
Err(e) => {
eprintln!("trex: -: {e}");
return ExitCode::FAILURE;
}
};
if chunk.is_empty() {
break;
}
reader.consume(chunk.len());
total += chunk.len();
let piece = if opened {
chunk
} else {
opening.extend_from_slice(&chunk);
if opening.len() < 4 {
continue;
}
opened = true;
if trex::files::has_bom(&opening) {
return scan_stdin_whole(&mut reader, opening, scanning, report);
}
drop_utf8_mark(&mut opening);
std::mem::take(&mut opening)
};
if !binary && piece.contains(&0) {
eprintln!("trex: - holds a NUL byte and is binary; --binary scans it");
return ExitCode::FAILURE;
}
run.push(&piece, &print, &mut out);
}
if !opened {
if trex::files::has_bom(&opening) {
return scan_stdin_whole(&mut reader, opening, scanning, report);
}
drop_utf8_mark(&mut opening);
if !binary && opening.contains(&0) {
eprintln!("trex: - holds a NUL byte and is binary; --binary scans it");
return ExitCode::FAILURE;
}
run.push(&opening, &print, &mut out);
}
let early = run.commits_early();
let printed = run.finish(&print, &mut out);
if !early && total > 0 {
eprintln!(
"trex: {total} bytes retained until the stream ended; a match of this pattern cannot commit before the end, since it reads a content guard, a lookaround reaching past its lines, a field anchor or a whole-stream axis"
);
}
if !json && format.is_none() && printed == 0 {
crate::out::line("no match");
}
if require_match && printed == 0 { ExitCode::FAILURE } else { ExitCode::SUCCESS }
}
fn drop_utf8_mark(opening: &mut Vec<u8>) {
opening.drain(..trex::encoding::Encoding::declared(opening).1);
}
fn scan_stdin_whole(
reader: &mut impl std::io::Read,
mut opening: Vec<u8>,
scanning: &Scanning,
report: &StreamReport<'_>,
) -> ExitCode {
if let Err(e) = reader.read_to_end(&mut opening) {
eprintln!("trex: -: {e}");
return ExitCode::FAILURE;
}
let input = trex::encoding::decode(opening);
let how = ScanHow {
backend: trex::Backend::Auto,
dual_grain: false,
chunk_size: None,
say: false,
take: None,
lists: report.lists,
single: false,
};
let (matches, of_member) = scan_found(scanning, &input, &trex::ShapeSet::new(), how);
let members = scanning.set().map(|set| Members { set, of: &of_member });
let capture_kinds = scanning.capture_kinds();
let value_view = crate::ValueView { kinds: &capture_kinds, style: report.style, clock: trex::Clock::current() };
let index = LineIndex::new(&input);
if let Some(t) = report.format {
crate::cli_files::print_formatted("-", &input, &matches, &index, t, members, None, false);
} else if report.json {
crate::cli_files::print_json_about(&input, &matches, &index, members, None, Some(&value_view));
} else {
crate::cli_files::print_human_about(&input, &matches, &index, report.painter, members, None);
}
if report.require_match && matches.is_empty() { ExitCode::FAILURE } else { ExitCode::SUCCESS }
}
pub(crate) struct Followed<'a> {
pub(crate) sources: &'a [Source],
pub(crate) select: Option<Select>,
pub(crate) unit: &'a RecordUnit,
pub(crate) binary: bool,
pub(crate) failed: bool,
}
fn say_change(name: &str, change: &Change, doing: &str) {
match change {
Change::Appended { .. } => {}
Change::Truncated => eprintln!("trex: {name}: truncated; {doing} from its start"),
Change::Replaced => eprintln!("trex: {name}: replaced by another file; {doing} from its start"),
Change::Gone => eprintln!("trex: {name}: removed; waiting for a file under its name"),
}
}
fn say_uncommittable(command: &str) {
eprintln!(
"trex {command}: no match of this pattern is final before its input ends, since it reads a content guard, a field anchor or a whole-input axis, and a followed file does not end; run it without --follow"
);
}
pub(crate) fn follow_scan(
f: &Followed<'_>,
scanning: &Scanning,
shapes: &trex::ShapeSet,
p: &StreamPrint<'_>,
) -> ExitCode {
if !StreamRun::new(scanning, shapes, "", 0, None).commits_early() {
say_uncommittable("scan");
return ExitCode::FAILURE;
}
let numbers = p.prefixed || p.format.is_some_and(trex::Template::reads_place);
let asked = Asked { numbers, offsets: true, binary: f.binary, units: false };
let names: Vec<String> = f.sources.iter().map(Source::name).collect();
let mut failed = f.failed;
let mut out = |l: &str| crate::out::line(l);
let mut runs: Vec<StreamRun<'_>> = Vec::new();
let mut decoders: Vec<Incremental> = Vec::new();
let mut owners: Vec<usize> = Vec::new();
let mut followees: Vec<(std::path::PathBuf, usize)> = Vec::new();
for (k, src) in f.sources.iter().enumerate() {
let Source::File(path) = src else {
continue;
};
let part = match read_part(src, f.select, f.unit, asked) {
Ok(part) => part,
Err(e) => {
eprintln!("trex: cannot read {}: {e}", names[k]);
failed = true;
continue;
}
};
if part.binary && !f.binary {
eprintln!("trex: {} holds a NUL byte and is binary; --binary scans it", names[k]);
failed = true;
continue;
}
let mut run = StreamRun::new(scanning, shapes, &names[k], part.base(), part.line_base);
run.push(&part.text, p, &mut out);
runs.push(run);
decoders.push(Incremental::after(part.encoding));
owners.push(k);
followees.push((path.clone(), part.input_len));
}
let exit = |failed: bool| if failed { ExitCode::FAILURE } else { ExitCode::SUCCESS };
if runs.is_empty() || runs.iter().all(|r| r.done(p)) {
return exit(failed);
}
let mut follower = match Follower::new(&followees) {
Ok(follower) => follower,
Err(e) => {
eprintln!("trex: cannot follow: {e}");
return ExitCode::FAILURE;
}
};
if let Some(why) = follower.unnotified() {
eprintln!("trex: no change notifications ({why}); looking at the files once a second");
}
loop {
let (i, change) = match follower.wait() {
Ok(next) => next,
Err(e) => {
eprintln!("trex: cannot follow {e}");
return ExitCode::FAILURE;
}
};
let name = &names[owners[i]];
match change {
Change::Appended { bytes, .. } => {
let text = decoders[i].decode(&bytes);
runs[i].push(&text, p, &mut out);
if runs.iter().all(|r| r.done(p)) {
return exit(failed);
}
}
other => {
say_change(name, &other, "scanning it");
let fresh = StreamRun::new(scanning, shapes, name, 0, Some(0));
std::mem::replace(&mut runs[i], fresh).finish(p, &mut out);
decoders[i] = Incremental::from_start();
}
}
}
}
pub(crate) fn follow_rules(
sources: &[Source],
run: &crate::cli_rules::RulesRun<'_>,
report: &crate::cli_rules::Report<'_>,
asked: Asked,
prefixed: bool,
failed: bool,
) -> ExitCode {
let scan = report.scan;
let on_records = scan.record_rules();
if !on_records.is_empty() {
eprintln!(
"trex scan: {} fire on records, and a record is whole only once the input holding it ends, which a followed file does not; follow the rules that fire on matches",
on_records.join(", ")
);
return ExitCode::FAILURE;
}
if !scan.stream(None, 0, None).commits_early() {
say_uncommittable("scan");
return ExitCode::FAILURE;
}
let names: Vec<String> = sources.iter().map(Source::name).collect();
let mut failed = failed;
let mut errors = 0usize;
let mut streams: Vec<trex::rule_scan::RuleStream<'_>> = Vec::new();
let mut printed: Vec<usize> = Vec::new();
let mut decoders: Vec<Incremental> = Vec::new();
let mut owners: Vec<usize> = Vec::new();
let mut followees: Vec<(std::path::PathBuf, usize)> = Vec::new();
let print = |found: Vec<trex::rule_scan::Found>, count: &mut usize, errors: &mut usize| {
for mut one in found {
if let Some(n) = run.max_count {
one.findings.truncate(n.saturating_sub(*count));
}
let (total, errs) =
report.print_findings(std::slice::from_ref(&one), prefixed, &mut |object| crate::out::line(&object));
*count += total;
*errors += errs;
}
};
for (k, src) in sources.iter().enumerate() {
let Source::File(path) = src else {
continue;
};
let part = match read_part(src, run.windowing.select, &run.unit, asked) {
Ok(part) => part,
Err(e) => {
eprintln!("trex: cannot read {}: {e}", names[k]);
failed = true;
continue;
}
};
if part.binary && !run.binary {
eprintln!("trex: {} holds a NUL byte and is binary; --binary scans it", names[k]);
failed = true;
continue;
}
let mut stream = scan.stream(Some(names[k].as_str()), part.base(), part.line_base);
let mut count = 0usize;
print(stream.push(&part.text), &mut count, &mut errors);
streams.push(stream);
printed.push(count);
decoders.push(Incremental::after(part.encoding));
owners.push(k);
followees.push((path.clone(), part.input_len));
}
let all_done = |printed: &[usize]| run.max_count.is_some_and(|n| printed.iter().all(|&c| c >= n));
let exit = |failed: bool, errors: usize| if failed || errors > 0 { ExitCode::FAILURE } else { ExitCode::SUCCESS };
if streams.is_empty() || all_done(&printed) {
return exit(failed, errors);
}
let mut follower = match Follower::new(&followees) {
Ok(follower) => follower,
Err(e) => {
eprintln!("trex: cannot follow: {e}");
return ExitCode::FAILURE;
}
};
if let Some(why) = follower.unnotified() {
eprintln!("trex: no change notifications ({why}); looking at the files once a second");
}
loop {
let (i, change) = match follower.wait() {
Ok(next) => next,
Err(e) => {
eprintln!("trex: cannot follow {e}");
return ExitCode::FAILURE;
}
};
let name = &names[owners[i]];
match change {
Change::Appended { bytes, .. } => {
let text = decoders[i].decode(&bytes);
print(streams[i].push(&text), &mut printed[i], &mut errors);
if all_done(&printed) {
return exit(failed, errors);
}
}
other => {
say_change(name, &other, "scanning it");
let fresh = scan.stream(Some(name.as_str()), 0, Some(0));
print(std::mem::replace(&mut streams[i], fresh).finish(), &mut printed[i], &mut errors);
printed[i] = 0;
decoders[i] = Incremental::from_start();
}
}
}
}
pub(crate) struct StreamEdit {
stream: trex::HeldStream,
written: usize,
}
pub(crate) type EditsAt<'a> = dyn FnMut(&[u8], &[trex::Span]) -> Vec<Edit> + 'a;
impl StreamEdit {
pub(crate) fn new(pattern: &trex::ast::Pattern, shapes: &trex::ShapeSet, origin: usize) -> Self {
let scanner = trex::StreamScanner::with_shapes(pattern.clone(), shapes.clone());
StreamEdit { stream: trex::HeldStream::new(scanner, origin, None), written: origin }
}
pub(crate) fn commits_early(&self) -> bool {
self.stream.commits_early()
}
pub(crate) fn push(&mut self, bytes: &[u8], edits: &mut EditsAt<'_>, out: &mut dyn Write) -> std::io::Result<()> {
let committed: Vec<trex::Span> = self.stream.push(bytes).into_iter().map(|(_, s)| s).collect();
let base = self.stream.base();
let upto = base + self.stream.settled();
write_edits(self.stream.held(), base, &mut self.written, &committed, upto, edits, out)
}
pub(crate) fn finish(self, edits: &mut EditsAt<'_>, out: &mut dyn Write) -> std::io::Result<()> {
let StreamEdit { stream, mut written } = self;
let ended = stream.finish();
let rest: Vec<trex::Span> = ended.matches.into_iter().map(|(_, s)| s).collect();
let upto = ended.base + ended.held.len();
write_edits(&ended.held, ended.base, &mut written, &rest, upto, edits, out)
}
}
fn write_edits(
held: &[u8],
base: usize,
written: &mut usize,
committed: &[trex::Span],
upto: usize,
edits: &mut EditsAt<'_>,
out: &mut dyn Write,
) -> std::io::Result<()> {
for e in edits(held, committed) {
out.write_all(&held[*written - base..e.start])?;
out.write_all(&e.replacement)?;
*written = base + e.end;
}
if upto > *written {
out.write_all(&held[*written - base..upto - base])?;
*written = upto;
}
out.flush()
}
pub(crate) fn follow_edit(
command: &str,
path: &std::path::Path,
part: &Part,
pattern: &trex::ast::Pattern,
shapes: &trex::ShapeSet,
edits: &mut EditsAt<'_>,
) -> ExitCode {
let name = path.display().to_string();
let mut stream = StreamEdit::new(pattern, shapes, part.base());
if !stream.commits_early() {
say_uncommittable(command);
return ExitCode::FAILURE;
}
let stdout = std::io::stdout();
let mut out = stdout.lock();
let written = |r: std::io::Result<()>| match r {
Ok(()) => true,
Err(e) => {
eprintln!("trex: cannot write the standard output: {e}");
false
}
};
if !written(stream.push(&part.text, edits, &mut out)) {
return ExitCode::FAILURE;
}
let mut follower = match Follower::new(&[(path.to_path_buf(), part.input_len)]) {
Ok(follower) => follower,
Err(e) => {
eprintln!("trex: cannot follow {name}: {e}");
return ExitCode::FAILURE;
}
};
if let Some(why) = follower.unnotified() {
eprintln!("trex: no change notifications ({why}); looking at the file once a second");
}
let mut decoder = Incremental::after(part.encoding);
loop {
let change = match follower.wait() {
Ok((_, change)) => change,
Err(e) => {
eprintln!("trex: cannot follow {e}");
return ExitCode::FAILURE;
}
};
match change {
Change::Appended { bytes, .. } => {
let text = decoder.decode(&bytes);
if !written(stream.push(&text, edits, &mut out)) {
return ExitCode::FAILURE;
}
}
other => {
say_change(&name, &other, "reading it");
let old = std::mem::replace(&mut stream, StreamEdit::new(pattern, shapes, 0));
if !written(old.finish(edits, &mut out)) {
return ExitCode::FAILURE;
}
decoder = Incremental::from_start();
}
}
}
}