Skip to main content

command_stream/commands/
text.rs

1//! Portable UTF-8 text commands. Input is read a line at a time; head stops
2//! early, tail keeps a ring of the last N lines, and uniq keeps one group.
3//! Sort necessarily retains the complete input. Output channels are progressive.
4use super::{CommandContext, StreamChunk};
5use crate::{trace, CommandResult, VirtualUtils};
6use std::collections::{HashSet, VecDeque};
7use std::io::Cursor;
8use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader};
9
10type TextResult<T> = std::result::Result<T, CommandResult>;
11
12struct Options {
13    count: usize,
14    flags: HashSet<char>,
15    files: Vec<String>,
16}
17
18fn failure(command: &str, message: impl std::fmt::Display) -> CommandResult {
19    CommandResult::error(format!("{command}: {message}\n"))
20}
21
22fn cancelled(command: &str, ctx: &CommandContext) -> TextResult<()> {
23    if ctx.is_cancelled() {
24        Err(CommandResult::error_with_code(
25            format!("{command}: cancelled\n"),
26            130,
27        ))
28    } else {
29        Ok(())
30    }
31}
32
33fn parse_count(command: &str, value: Option<&str>) -> TextResult<usize> {
34    let value = value.unwrap_or("");
35    if value.is_empty() || !value.bytes().all(|byte| byte.is_ascii_digit()) {
36        return Err(failure(
37            command,
38            format!("invalid number of lines: '{value}'"),
39        ));
40    }
41    value
42        .parse::<u64>()
43        .ok()
44        .filter(|number| *number <= 9_007_199_254_740_991)
45        .and_then(|number| usize::try_from(number).ok())
46        .ok_or_else(|| failure(command, format!("invalid number of lines: '{value}'")))
47}
48
49fn flags(command: &str, value: &str, output: &mut HashSet<char>) -> TextResult<()> {
50    let value = match value {
51        "--reverse" => "r",
52        "--numeric-sort" => "n",
53        "--unique" => "u",
54        "--count" => "c",
55        "--repeated" => "d",
56        "--ignore-case" => "i",
57        _ => &value[1..],
58    };
59    let allowed = if command == "sort" { "rnu" } else { "cdui" };
60    for flag in value.chars() {
61        if !allowed.contains(flag) {
62            return Err(failure(command, format!("invalid option '{value}'")));
63        }
64        output.insert(flag);
65    }
66    Ok(())
67}
68
69fn parse(command: &str, args: &[String]) -> TextResult<Options> {
70    let mut result = Options {
71        count: 10,
72        flags: HashSet::new(),
73        files: Vec::new(),
74    };
75    let mut index = 0;
76    let mut options = true;
77    while index < args.len() {
78        let arg = &args[index];
79        if options && arg == "--" {
80            options = false;
81        } else if options && arg.starts_with('-') && arg != "-" {
82            if command == "head" || command == "tail" {
83                let value = if arg == "-n" || arg == "--lines" {
84                    index += 1;
85                    args.get(index).map(String::as_str)
86                } else if let Some(value) = arg.strip_prefix("--lines=") {
87                    Some(value)
88                } else {
89                    Some(arg.strip_prefix("-n").unwrap_or(&arg[1..]))
90                };
91                result.count = parse_count(command, value)?;
92            } else {
93                flags(command, arg, &mut result.flags)?;
94            }
95        } else {
96            result.files.push(arg.clone());
97        }
98        index += 1;
99    }
100    if command == "uniq" && result.files.len() > 2 {
101        return Err(failure(command, "extra operand"));
102    }
103    if result.flags.contains(&'d') && result.flags.contains(&'u') {
104        return Err(failure(command, "cannot combine -d and -u"));
105    }
106    Ok(result)
107}
108
109async fn input(
110    command: &str,
111    file: &str,
112    ctx: &CommandContext,
113) -> TextResult<Box<dyn AsyncBufRead + Send + Unpin>> {
114    cancelled(command, ctx)?;
115    trace("VirtualCommand", &format!("{command}: reading {file}"));
116    if file == "-" {
117        return Ok(Box::new(Cursor::new(
118            ctx.stdin.clone().unwrap_or_default().into_bytes(),
119        )));
120    }
121    let path = VirtualUtils::resolve_path(file, Some(&ctx.get_cwd()));
122    let handle = tokio::fs::File::open(&path)
123        .await
124        .map_err(|error| failure(command, format!("{file}: {error}")))?;
125    if handle
126        .metadata()
127        .await
128        .map_err(|error| failure(command, error))?
129        .is_dir()
130    {
131        return Err(failure(command, format!("{file}: Is a directory")));
132    }
133    Ok(Box::new(BufReader::new(handle)))
134}
135
136async fn next_line(
137    command: &str,
138    reader: &mut (dyn AsyncBufRead + Send + Unpin),
139    ctx: &CommandContext,
140) -> TextResult<Option<String>> {
141    cancelled(command, ctx)?;
142    let mut line = String::new();
143    let count = reader
144        .read_line(&mut line)
145        .await
146        .map_err(|error| failure(command, error))?;
147    Ok((count > 0).then_some(line))
148}
149
150async fn emit(
151    command: &str,
152    text: &str,
153    ctx: &CommandContext,
154    output: &mut String,
155) -> TextResult<()> {
156    cancelled(command, ctx)?;
157    if text.is_empty() {
158        return Ok(());
159    }
160    if let Some(sender) = &ctx.output_tx {
161        sender
162            .send(StreamChunk::Stdout(text.to_string()))
163            .await
164            .map_err(|_| CommandResult::error_with_code("output closed", 130))?;
165    }
166    output.push_str(text);
167    Ok(())
168}
169
170async fn select(command: &str, ctx: &CommandContext, options: &Options) -> TextResult<String> {
171    let files = if options.files.is_empty() {
172        vec!["-".to_string()]
173    } else {
174        options.files.clone()
175    };
176    let mut output = String::new();
177    for (index, file) in files.iter().enumerate() {
178        let mut reader = input(command, file, ctx).await?;
179        if files.len() > 1 {
180            let name = if file == "-" { "standard input" } else { file };
181            emit(
182                command,
183                &format!("{}==> {name} <==\n", if index > 0 { "\n" } else { "" }),
184                ctx,
185                &mut output,
186            )
187            .await?;
188        }
189        if command == "head" {
190            for _ in 0..options.count {
191                let Some(line) = next_line(command, reader.as_mut(), ctx).await? else {
192                    break;
193                };
194                emit(command, &line, ctx, &mut output).await?;
195            }
196        } else {
197            let mut ring = VecDeque::new();
198            while let Some(line) = next_line(command, reader.as_mut(), ctx).await? {
199                if options.count > 0 {
200                    if ring.len() == options.count {
201                        ring.pop_front();
202                    }
203                    ring.push_back(line);
204                }
205            }
206            for line in ring {
207                emit(command, &line, ctx, &mut output).await?;
208            }
209        }
210    }
211    Ok(output)
212}
213
214fn numeric(line: &str) -> f64 {
215    static NUMBER: once_cell::sync::Lazy<regex::Regex> = once_cell::sync::Lazy::new(|| {
216        regex::Regex::new(r"^\s*([+-]?(?:\d+(?:\.\d*)?|\.\d+))").unwrap()
217    });
218    let value = NUMBER
219        .captures(line)
220        .and_then(|matches| matches[1].parse().ok())
221        .unwrap_or(0.0);
222    // Match JavaScript numeric comparison: signed zero shares the same key.
223    if value == 0.0 {
224        0.0
225    } else {
226        value
227    }
228}
229
230async fn sorted(ctx: &CommandContext, options: &Options) -> TextResult<String> {
231    let files = if options.files.is_empty() {
232        vec!["-".to_string()]
233    } else {
234        options.files.clone()
235    };
236    let mut lines = Vec::new();
237    for file in files {
238        let mut reader = input("sort", &file, ctx).await?;
239        while let Some(line) = next_line("sort", reader.as_mut(), ctx).await? {
240            lines.push(line.strip_suffix('\n').unwrap_or(&line).to_string());
241        }
242    }
243    lines.sort_by(|left, right| {
244        let numeric_order = if options.flags.contains(&'n') {
245            numeric(left).total_cmp(&numeric(right))
246        } else {
247            std::cmp::Ordering::Equal
248        };
249        numeric_order.then_with(|| left.cmp(right))
250    });
251    if options.flags.contains(&'r') {
252        lines.reverse();
253    }
254    let mut output = String::new();
255    let mut previous: Option<String> = None;
256    for line in lines {
257        let duplicate = previous.as_ref().is_some_and(|previous| {
258            if options.flags.contains(&'n') {
259                numeric(previous) == numeric(&line)
260            } else {
261                previous == &line
262            }
263        });
264        if !options.flags.contains(&'u') || !duplicate {
265            emit("sort", &format!("{line}\n"), ctx, &mut output).await?;
266        }
267        previous = Some(line);
268    }
269    Ok(output)
270}
271
272fn group(line: &str, count: usize, options: &Options) -> String {
273    if (options.flags.contains(&'d') && count == 1) || (options.flags.contains(&'u') && count > 1) {
274        return String::new();
275    }
276    if options.flags.contains(&'c') {
277        format!("{count:>7} {line}")
278    } else {
279        line.to_string()
280    }
281}
282
283async fn write_group(
284    text: &str,
285    file: &mut Option<tokio::fs::File>,
286    ctx: &CommandContext,
287    output: &mut String,
288) -> TextResult<()> {
289    cancelled("uniq", ctx)?;
290    if let Some(file) = file {
291        file.write_all(text.as_bytes())
292            .await
293            .map_err(|error| failure("uniq", error))
294    } else {
295        emit("uniq", text, ctx, output).await
296    }
297}
298
299async fn unique(ctx: &CommandContext, options: &Options) -> TextResult<String> {
300    let mut reader = input(
301        "uniq",
302        options.files.first().map(String::as_str).unwrap_or("-"),
303        ctx,
304    )
305    .await?;
306    let mut file = match options.files.get(1).filter(|path| path.as_str() != "-") {
307        Some(path) => Some(
308            tokio::fs::File::create(VirtualUtils::resolve_path(path, Some(&ctx.get_cwd())))
309                .await
310                .map_err(|error| failure("uniq", error))?,
311        ),
312        None => None,
313    };
314    let mut output = String::new();
315    let mut previous = String::new();
316    let mut previous_key = String::new();
317    let mut count = 0;
318    while let Some(line) = next_line("uniq", reader.as_mut(), ctx).await? {
319        let text = line.strip_suffix('\n').unwrap_or(&line);
320        let key = if options.flags.contains(&'i') {
321            text.to_lowercase()
322        } else {
323            text.to_string()
324        };
325        if count > 0 && key != previous_key {
326            write_group(
327                &group(&previous, count, options),
328                &mut file,
329                ctx,
330                &mut output,
331            )
332            .await?;
333            count = 0;
334        }
335        if count == 0 {
336            previous = line;
337            previous_key = key;
338        }
339        count += 1;
340    }
341    if count > 0 {
342        write_group(
343            &group(&previous, count, options),
344            &mut file,
345            ctx,
346            &mut output,
347        )
348        .await?;
349    }
350    if let Some(mut file) = file {
351        file.flush().await.map_err(|error| failure("uniq", error))?;
352    }
353    Ok(output)
354}
355
356async fn execute(command: &str, ctx: CommandContext) -> CommandResult {
357    let result = async {
358        let options = parse(command, &ctx.args)?;
359        cancelled(command, &ctx)?;
360        match command {
361            "head" | "tail" => select(command, &ctx, &options).await,
362            "sort" => sorted(&ctx, &options).await,
363            _ => unique(&ctx, &options).await,
364        }
365    }
366    .await;
367    match result {
368        Ok(output) => CommandResult::success(output),
369        Err(error) => error,
370    }
371}
372
373/// Show the first N lines (default 10), preserving original line endings.
374pub async fn head(ctx: CommandContext) -> CommandResult {
375    execute("head", ctx).await
376}
377/// Show the last N lines (default 10), retaining at most N input lines.
378pub async fn tail(ctx: CommandContext) -> CommandResult {
379    execute("tail", ctx).await
380}
381/// Sort complete input using locale-independent UTF-8 byte ordering.
382pub async fn sort(ctx: CommandContext) -> CommandResult {
383    execute("sort", ctx).await
384}
385/// Filter consecutive duplicate groups, optionally counting or selecting them.
386pub async fn uniq(ctx: CommandContext) -> CommandResult {
387    execute("uniq", ctx).await
388}