perf(pi-shell): docker batch — frequency-dedup logs, kubectl arms, helm glog strip, compose-legacy dispatch

- filter_docker_logs/filter_logs now collapse repeated lines to order-
  preserving '(×N)' counts via shared compact_log_lines helper.
- kubectl apply|delete|rollout|scale|create|wait|label|annotate now
  route through head_tail_dedup instead of raw passthrough.
- filter_helm strips client-go glog 'W####' warnings and kube-config
  permission warnings before table compaction.
- normalize_program maps standalone 'docker-compose' (v1) to 'docker'
  so it routes through the docker filter with proper subcommand detection.

Op: compress
This commit is contained in:
metaphorics
2026-06-11 10:40:13 +09:00
parent b3ceaa44ee
commit 1f8bf85213
3 changed files with 527 additions and 51 deletions
+121 -1
View File
@@ -22,7 +22,42 @@ pub fn detect_tokens(tokens: &[String]) -> Option<CommandIdentity> {
let tokens = strip_launch_prefix(tokens)?;
let (program, rest) = tokens.split_first()?;
let normalized = normalize_program(program)?;
let subcommand = detect_subcommand(&normalized, rest);
let is_docker_compose = program
.rsplit('/')
.next()
.is_some_and(|n| n.eq_ignore_ascii_case("docker-compose"));
let subcommand = if is_docker_compose {
// docker-compose v1 flags that consume a value; skip them so the
// real action (up/down/ps/logs/etc.) is found, matching docker compose
// routing in docker.rs.
first_non_global_arg(
rest,
&[
"-f",
"--file",
"--profile",
"-p",
"--project-name",
"--env-file",
"--parallel",
"--progress",
"--project-directory",
"--workdir",
"-w",
"--ansi",
"--log-level",
"-H",
"--host",
"--tlscacert",
"--tlscert",
"--tlskey",
],
&["--compatibility", "--dry-run", "--verbose", "-v", "--no-ansi"],
&[],
)
} else {
detect_subcommand(&normalized, rest)
};
Some(CommandIdentity { program: normalized, subcommand })
}
@@ -95,6 +130,7 @@ fn normalize_program(program: &str) -> Option<String> {
Some(match lowered.as_str() {
"gradlew.bat" => "gradlew".to_string(),
"mvnw.cmd" => "mvnw".to_string(),
"docker-compose" => "docker".to_string(),
_ => lowered,
})
}
@@ -759,3 +795,87 @@ fn npx_workspace_value_is_skipped_in_subcommand_detection() {
assert_eq!(command.program, "npx");
assert_eq!(command.subcommand.as_deref(), Some("vitest"));
}
#[test]
fn normalizes_docker_compose_to_docker() {
let command = detect("docker-compose up").expect("docker-compose command is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("up"));
let command = detect("/usr/local/bin/docker-compose logs")
.expect("path-prefixed docker-compose is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("logs"));
}
#[test]
fn docker_composer_is_not_normalized() {
let command = detect("docker-composer up").expect("docker-composer command is detected");
assert_eq!(command.program, "docker-composer");
assert_eq!(command.subcommand.as_deref(), Some("up"));
}
#[test]
fn docker_compose_with_flags_finds_real_subcommand() {
// Value-taking compose flags must be skipped so the real action is found.
let command =
detect("docker-compose -p myproj up").expect("docker-compose with -p flag is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("up"));
let command = detect("docker-compose --profile logs up")
.expect("docker-compose with --profile flag is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("up"));
let command = detect("docker-compose -f docker-compose.yml ps")
.expect("docker-compose with -f flag is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("ps"));
}
#[test]
fn docker_compose_global_value_flags_skip_correctly() {
// --log-level consumes its value; the next token is the subcommand.
let command = detect("docker-compose --log-level debug up")
.expect("docker-compose --log-level is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("up"));
// -H consumes its value.
let command =
detect("docker-compose -H tcp://host:2376 ps").expect("docker-compose -H is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("ps"));
// --host=inline also works.
let command = detect("docker-compose --host=tcp://host:2376 logs")
.expect("docker-compose --host= is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("logs"));
// --tlscacert consumes its value.
let command = detect("docker-compose --tlscacert /path/ca.pem up")
.expect("docker-compose --tlscacert is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("up"));
// --tlscert consumes its value.
let command = detect("docker-compose --tlscert /path/cert.pem ps")
.expect("docker-compose --tlscert is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("ps"));
// --tlskey consumes its value.
let command = detect("docker-compose --tlskey /path/key.pem logs")
.expect("docker-compose --tlskey is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("logs"));
}
#[test]
fn docker_compose_without_flags_still_works() {
let command = detect("docker-compose up").expect("docker-compose command is detected");
assert_eq!(command.program, "docker");
assert_eq!(command.subcommand.as_deref(), Some("up"));
}
+370 -49
View File
@@ -21,7 +21,17 @@ pub fn supports(subcommand: Option<&str>) -> bool {
| "install"
| "upgrade"
| "template"
| "lint"
| "lint" | "apply"
| "delete"
| "rollout"
| "scale"
| "create"
| "wait" | "label"
| "annotate"
| "up" | "down"
| "start"
| "stop" | "restart"
| "rm"
)
)
}
@@ -56,6 +66,12 @@ fn filter_docker(ctx: &MinimizerCtx<'_>, input: &str, exit_code: i32) -> String
input.to_string()
};
}
// start/stop/restart/rm output is container names or confirmation messages;
// compact_build_or_progress would strip legitimate lines that happen to
// contain progress substrings (e.g. a container named "my-Downloading-app").
if is_docker_lifecycle_command(ctx) {
return head_tail_dedup(input);
}
compact_build_or_progress(input)
}
@@ -87,6 +103,9 @@ fn filter_kubectl(ctx: &MinimizerCtx<'_>, input: &str, exit_code: i32) -> String
Some("describe") => {
primitives::head_tail_lines(&primitives::dedup_consecutive_lines(input), 120, 80)
},
Some("apply" | "delete" | "rollout" | "scale" | "create" | "wait" | "label" | "annotate") => {
head_tail_dedup(input)
},
_ => compact_build_or_progress(input),
}
}
@@ -398,14 +417,65 @@ fn filter_helm(ctx: &MinimizerCtx<'_>, input: &str, exit_code: i32) -> String {
if exit_code != 0 {
return input.to_string();
}
let cleaned = strip_helm_noise(input);
match ctx.subcommand {
Some("list" | "ls" | "status") => compact_table(input, 20),
Some("install" | "upgrade" | "lint") => compact_build_or_progress(input),
Some("list" | "ls" | "status") => compact_table(&cleaned, 20),
Some("install" | "upgrade" | "lint") => compact_build_or_progress(&cleaned),
Some("template") => input.to_string(),
_ => head_tail_dedup(input),
_ => head_tail_dedup(&cleaned),
}
}
fn strip_helm_noise(input: &str) -> String {
let mut out = String::new();
for line in input.lines() {
let trimmed = line.trim_start();
if is_glog_prefix(trimmed) {
continue;
}
if trimmed.starts_with("WARNING: Kubernetes configuration file is") {
continue;
}
out.push_str(line);
out.push('\n');
}
out
}
fn is_glog_prefix(line: &str) -> bool {
let bytes = line.as_bytes();
if bytes.len() < 6
|| bytes[0] != b'W'
|| bytes[5] != b' '
|| !bytes[1..5].iter().all(|b| b.is_ascii_digit())
{
return false;
}
// Validate month (01-12) and day (01-31) to avoid stripping legitimate
// output that happens to start with W + 4 digits + space.
let month = (bytes[1] - b'0') * 10 + (bytes[2] - b'0');
let day = (bytes[3] - b'0') * 10 + (bytes[4] - b'0');
if month < 1 || month > 12 || day < 1 || day > 31 {
return false;
}
// Require a time stamp (hh:mm...) after the space to distinguish real
// glog lines from space-separated table output where a release name
// happens to match W + MMDD + space.
if bytes.len() < 11 {
return false;
}
if !bytes[6].is_ascii_digit() || !bytes[7].is_ascii_digit() {
return false;
}
if bytes[8] != b':' {
return false;
}
if !bytes[9].is_ascii_digit() || !bytes[10].is_ascii_digit() {
return false;
}
true
}
/// Returns `true` when `tok` is a known docker-compose option that consumes
/// the next token as its value (i.e. is space-separated, not `--flag=value`).
fn compose_option_consumes_next(tok: &str) -> bool {
@@ -419,7 +489,7 @@ fn compose_option_consumes_next(tok: &str) -> bool {
| "--progress"
| "--project-directory"
| "--project-name"
| "--workdir"
| "-p" | "--workdir"
| "-w"
)
}
@@ -501,6 +571,32 @@ fn is_compose_listing_action(command: &str) -> bool {
}
}
fn is_docker_lifecycle_command(ctx: &MinimizerCtx<'_>) -> bool {
matches!(ctx.subcommand, Some("start" | "stop" | "restart" | "rm"))
|| ctx.subcommand == Some("compose") && is_compose_lifecycle_action(ctx.command)
}
fn is_compose_lifecycle_action(command: &str) -> bool {
let mut tokens = command
.split_whitespace()
.skip_while(|token| *token != "compose");
if tokens.next() != Some("compose") {
return false;
}
loop {
match tokens.next() {
None => return false,
Some(tok)
if tok.starts_with('-') && !tok.contains('=') && compose_option_consumes_next(tok) =>
{
tokens.next(); // skip value
},
Some(tok) if tok.starts_with('-') => {}, // skip boolean flag
Some(tok) => return matches!(tok, "start" | "stop" | "restart" | "rm"),
}
}
}
fn docker_listing_requests_table(command: &str) -> bool {
let mut tokens = command.split_whitespace();
while let Some(token) = tokens.next() {
@@ -524,58 +620,29 @@ fn docker_format_requests_table(format: &str) -> bool {
fn filter_logs(input: &str) -> String {
let without_empty_runs = drop_repeated_blank_lines(input);
let deduped = primitives::dedup_consecutive_lines(&without_empty_runs);
primitives::head_tail_lines(&deduped, 120, 80)
crate::minimizer::filters::system::compact_log_lines(
&without_empty_runs,
120,
80,
crate::minimizer::filters::system::normalize_log_line,
)
}
fn filter_docker_logs(input: &str) -> String {
let without_empty_runs = drop_repeated_blank_lines(input);
let deduped = dedup_consecutive_log_lines(&without_empty_runs);
primitives::head_tail_lines(&deduped, 120, 80)
crate::minimizer::filters::system::compact_log_lines(&without_empty_runs, 120, 80, log_dedup_key)
}
fn dedup_consecutive_log_lines(input: &str) -> String {
let mut out = String::new();
let mut previous: Option<&str> = None;
let mut previous_key: Option<&str> = None;
let mut count = 0usize;
for line in input.lines() {
let key = log_dedup_key(line);
if previous_key == Some(key) {
count += 1;
continue;
}
flush_repeated_log_line(&mut out, previous, count);
previous = Some(line);
previous_key = Some(key);
count = 1;
}
flush_repeated_log_line(&mut out, previous, count);
out
}
fn flush_repeated_log_line(out: &mut String, line: Option<&str>, count: usize) {
let Some(line) = line else {
return;
};
out.push_str(line);
if count > 1 {
out.push_str(" (×");
out.push_str(&count.to_string());
out.push(')');
}
out.push('\n');
}
fn log_dedup_key(line: &str) -> &str {
fn log_dedup_key(line: &str) -> String {
if let Some((service, message)) = line.split_once('|') {
let service = service.trim();
if is_compose_log_service(service) {
return message.trim_start();
let message = message.trim_start();
let normalized = crate::minimizer::filters::system::normalize_log_line(message);
return format!("{}|{}", service, normalized);
}
}
line
crate::minimizer::filters::system::normalize_log_line(line)
}
fn is_compose_log_service(value: &str) -> bool {
@@ -691,8 +758,11 @@ mod tests {
fn dedups_compose_service_prefixed_log_messages() {
let input = "api-1 | ready\napi-2 | ready\napi | ready\nworker | busy\n";
let out = filter_docker_logs(input);
assert!(out.contains("api-1 | ready (×3)"));
assert!(out.contains("worker | busy"));
// Different services with the same message must NOT be collapsed together.
assert!(out.contains("api-1 | ready"), "distinct service must be preserved: {out}");
assert!(out.contains("api-2 | ready"), "distinct service must be preserved: {out}");
assert!(out.contains("api | ready"), "distinct service must be preserved: {out}");
assert!(out.contains("worker | busy"), "distinct service must be preserved: {out}");
}
#[test]
@@ -706,7 +776,10 @@ mod tests {
};
let input = "api-1 | ready\napi-2 | ready\napi | ready\n";
let out = filter(&compose_ctx, input, 0).text;
assert!(out.contains("api-1 | ready (×3)"));
// Per-service dedup: same message from different services stays distinct.
assert!(out.contains("api-1 | ready"), "distinct service line must be preserved: {out}");
assert!(out.contains("api-2 | ready"), "distinct service line must be preserved: {out}");
assert!(out.contains("api | ready"), "distinct service line must be preserved: {out}");
}
#[test]
@@ -1204,4 +1277,252 @@ mod tests {
let out = filter(&ctx, input.as_str(), 0).text;
assert!(out.contains("rows"), "-owide is a table format and must be compacted, got: {out}");
}
// ── Global normalized log dedup tests ───────────────────────────────
#[test]
fn dedups_non_consecutive_log_lines_globally() {
let input = "ready\nstarting\nready\nstarting\nready\ndone\n";
let out = filter_docker_logs(input);
assert!(out.contains("ready (×3)"), "global dedup should count all occurrences: {out}");
assert!(out.contains("starting (×2)"), "global dedup should count all occurrences: {out}");
assert!(out.contains("done"), "global dedup should preserve unique lines: {out}");
}
#[test]
fn dedups_compose_service_prefixed_lines_non_consecutive() {
let input = "api-1 | ready\nworker | busy\napi-2 | ready\nworker | busy\n";
let out = filter_docker_logs(input);
assert!(
out.contains("api-1 | ready"),
"compose service-prefixed lines should dedup per service: {out}"
);
assert!(
out.contains("api-2 | ready"),
"compose service-prefixed lines should preserve distinct services: {out}"
);
assert!(
out.contains("worker | busy (×2)"),
"compose same-service lines should dedup globally: {out}"
);
}
// ── kubectl uncovered subcommands tests ───────────────────────────────
#[test]
fn kubectl_apply_output_condensed() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let kubectl_ctx = ctx("kubectl", Some("apply"), &cfg);
let mut input = String::new();
for i in 0..250 {
input.push_str(&format!("deployment.apps/app-{i} unchanged\n"));
}
let out = filter(&kubectl_ctx, &input, 0).text;
assert!(
out.contains("lines omitted") || out.contains("lines truncated"),
"kubectl apply output should be head/tail condensed: {out}"
);
}
// ── helm glog strip tests ───────────────────────────────────────────
#[test]
fn helm_list_strips_glog_warning_lines() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let helm_ctx = ctx("helm", Some("list"), &cfg);
let input = "W0115 10:30:00 warning message from internal\nNAME: my-release\nLAST DEPLOYED: \
Mon Jan 15 10:30:00 2024\nNAMESPACE: default\nSTATUS: deployed\nREVISION: 3\n";
let out = filter(&helm_ctx, input, 0).text;
assert!(!out.contains("W0115"), "glog W-prefixed lines should be stripped: {out}");
assert!(out.contains("NAME: my-release"), "release info should be preserved: {out}");
assert!(out.contains("STATUS: deployed"), "release info should be preserved: {out}");
}
#[test]
fn helm_list_strips_kubernetes_config_warning() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let helm_ctx = ctx("helm", Some("list"), &cfg);
let input = "WARNING: Kubernetes configuration file is group-readable\nNAME: \
my-release\nSTATUS: deployed\n";
let out = filter(&helm_ctx, input, 0).text;
assert!(
!out.contains("WARNING: Kubernetes configuration file"),
"Kubernetes config warning should be stripped: {out}"
);
assert!(out.contains("NAME: my-release"), "release info should be preserved: {out}");
}
// ── Blocking-issue regression tests ─────────────────────────────────
#[test]
fn docker_logs_preserves_blank_lines() {
let input = "a\n\nb\n";
let out = filter_docker_logs(input);
assert_eq!(out, "a\n\nb\n", "blank lines must not be silently dropped: {out}");
}
#[test]
fn docker_logs_preserves_non_consecutive_blank_lines() {
let input = "a\n\nb\n\nc\n";
let out = filter_docker_logs(input);
assert_eq!(
out, "a\n\nb\n\nc\n",
"non-consecutive blank lines must not be globally deduplicated: {out}"
);
}
#[test]
fn docker_logs_dedups_empty_compose_messages_per_service() {
let input = "api | \napi | \napi | \nworker | \nworker | \n";
let out = filter_docker_logs(input);
assert!(out.contains("api | (×3)"), "same-service empty messages should dedup: {out}");
assert!(out.contains("worker | (×2)"), "same-service empty messages should dedup: {out}");
}
#[test]
fn docker_start_preserves_container_names_with_progress_keywords() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
for sub in &["start", "stop", "restart", "rm"] {
let ctx = MinimizerCtx {
program: "docker",
subcommand: Some(sub),
command: &format!("docker {sub} my-Downloading-app"),
config: &cfg,
};
let input = "my-Downloading-app\n";
let out = filter(&ctx, input, 0).text;
assert_eq!(
out, input,
"docker {sub} must not strip container names matching progress heuristics"
);
}
}
#[test]
fn docker_compose_start_preserves_container_names_with_progress_keywords() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
for action in &["start", "stop", "restart", "rm"] {
let ctx = MinimizerCtx {
program: "docker",
subcommand: Some("compose"),
command: &format!("docker compose {action} my-Downloading-app"),
config: &cfg,
};
let input = "my-Downloading-app\n";
let out = filter(&ctx, input, 0).text;
assert_eq!(
out, input,
"docker compose {action} must not strip container names matching progress heuristics"
);
}
}
#[test]
fn helm_list_preserves_release_names_starting_with_w_digit() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let helm_ctx = ctx("helm", Some("list"), &cfg);
let input = "NAME\tNAMESPACE\nW2024-app\tdefault\nmy-app\tdefault\n";
let out = filter(&helm_ctx, input, 0).text;
assert!(
out.contains("W2024-app"),
"release names starting with W+digits must not be stripped: {out}"
);
assert!(out.contains("my-app"), "other releases must be preserved: {out}");
}
#[test]
fn glog_prefix_requires_trailing_space() {
assert!(is_glog_prefix("W0115 10:30:00 warning"));
assert!(!is_glog_prefix("W2024-app"));
assert!(!is_glog_prefix("W123"));
assert!(!is_glog_prefix("W12345"));
}
#[test]
fn glog_prefix_rejects_invalid_month_day() {
// W2024 has month=20 (invalid) — must not be treated as glog.
assert!(!is_glog_prefix("W2024 default"));
assert!(!is_glog_prefix("W0015 10:30:00 warning"));
assert!(!is_glog_prefix("W1301 10:30:00 warning"));
assert!(!is_glog_prefix("W1232 10:30:00 warning"));
}
#[test]
fn glog_prefix_rejects_space_separated_without_time() {
// W1231 looks like a valid glog date, but "default" is not a time stamp.
assert!(!is_glog_prefix("W1231 default"));
// Valid glog with time stamp must still pass.
assert!(is_glog_prefix("W1231 10:30:00 warning"));
}
#[test]
fn helm_list_preserves_release_names_starting_with_w_digit_space() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let helm_ctx = ctx("helm", Some("list"), &cfg);
// W2024 default looks like a glog prefix but month=20 is invalid.
let input = "NAME\tNAMESPACE\nW2024 default\tdefault\nmy-app\tdefault\n";
let out = filter(&helm_ctx, input, 0).text;
assert!(
out.contains("W2024 default"),
"release names like 'W2024 default' must not be stripped: {out}"
);
assert!(out.contains("my-app"), "other releases must be preserved: {out}");
}
#[test]
fn helm_list_preserves_space_separated_valid_w_digit_release_names() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let helm_ctx = ctx("helm", Some("list"), &cfg);
// W1231 default has a valid month/day but lacks a time stamp, so it must
// NOT be stripped as a glog line.
let input = "NAME NAMESPACE\nW1231 default\nmy-app default\n";
let out = filter(&helm_ctx, input, 0).text;
assert!(
out.contains("W1231 default"),
"space-separated release name W1231 must not be stripped: {out}"
);
assert!(out.contains("my-app default"), "other releases must be preserved: {out}");
}
#[test]
fn docker_compose_p_project_name_routes_logs() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let ctx = MinimizerCtx {
program: "docker",
subcommand: Some("compose"),
command: "docker compose -p myproj logs api",
config: &cfg,
};
assert!(is_log_command(&ctx), "docker compose -p myproj logs must be a log command");
}
#[test]
fn docker_compose_p_project_name_routes_ps() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let ctx = MinimizerCtx {
program: "docker",
subcommand: Some("compose"),
command: "docker compose -p myproj ps",
config: &cfg,
};
assert!(
is_table_command(&ctx),
"docker compose -p myproj ps must be classified as a table command"
);
}
#[test]
fn docker_compose_p_project_name_routes_start() {
let cfg = MinimizerConfig { enabled: true, ..Default::default() };
let ctx = MinimizerCtx {
program: "docker",
subcommand: Some("compose"),
command: "docker compose -p myproj start api",
config: &cfg,
};
assert!(
is_docker_lifecycle_command(&ctx),
"docker compose -p myproj start must be classified as a lifecycle command"
);
}
}
@@ -223,7 +223,7 @@ struct LogLine {
count: usize,
}
fn normalize_log_line(line: &str) -> String {
pub(super) fn normalize_log_line(line: &str) -> String {
let without_timestamp = strip_leading_timestamp(line.trim());
let mut out = String::new();
for token in without_timestamp.split_whitespace() {
@@ -310,6 +310,41 @@ fn flush_digits(out: &mut String, digits: &mut String) {
digits.clear();
}
pub(super) fn compact_log_lines(
input: &str,
head: usize,
tail: usize,
key_fn: impl Fn(&str) -> String,
) -> String {
let mut unique: Vec<LogLine> = Vec::new();
let mut by_key: HashMap<String, usize> = HashMap::new();
for (idx, line) in input.lines().enumerate() {
let key = if line.trim().is_empty() {
// Blank lines are section separators; do not globally deduplicate them.
// drop_repeated_blank_lines already collapsed consecutive blanks.
format!("<blank-{idx}>")
} else {
let key = key_fn(line);
if key.is_empty() {
line.to_string()
} else {
key
}
};
if let Some(index) = by_key.get(&key).copied() {
if let Some(entry) = unique.get_mut(index) {
entry.count += 1;
}
} else {
by_key.insert(key, unique.len());
unique.push(LogLine { original: line.to_string(), count: 1 });
}
}
render_counted_lines(&unique, head, tail)
}
fn render_counted_lines(lines: &[LogLine], head: usize, tail: usize) -> String {
let mut out = String::new();
if lines.len() <= head + tail {