//! Native runner and boundary-owned values for `|`: pipe every editor //! selection through one shell command. The core imports only the plain value //! types below; native shells call `Job.copy` before starting a worker, then //! `runOne` on that worker. No subprocess or borrowed core memory reaches the //! editor/event-loop thread. const std = @import("std"); /// A deliberately finite answer. One selection cannot retain more than 1 MiB /// and a multi-selection command cannot retain more than 4 MiB in total. /// stderr is diagnostic-only and has a smaller independent ceiling. Crossing /// any ceiling fails the whole atomic request. pub const max_stdout_bytes: usize = 1024 * 1024; pub const max_total_stdout_bytes: usize = 4 * 1024 * 1024; pub const max_stderr_bytes: usize = 64 * 1024; pub const command_timeout_seconds: u64 = 10; pub const Input = struct { bytes: []const u8 }; /// Borrowed view exposed by the core while an effect is being drained. pub const Request = struct { id: u32, command: []const u8, cwd: []const u8, inputs: []const Input, }; /// Worker-owned snapshot. `copy` is intentionally called synchronously while /// draining the effect: the worker can start after arbitrary later edits and /// still owns exactly the command, directory, and selection bytes submitted. pub const Job = struct { id: u32, command: []u8, cwd: []u8, inputs: [][]u8, pub fn copy(gpa: std.mem.Allocator, request: Request) !*Job { const job = try gpa.create(Job); errdefer gpa.destroy(job); job.* = .{ .id = request.id, .command = try gpa.dupe(u8, request.command), .cwd = &.{}, .inputs = &.{}, }; errdefer gpa.free(job.command); job.cwd = try gpa.dupe(u8, request.cwd); errdefer gpa.free(job.cwd); job.inputs = try gpa.alloc([]u8, request.inputs.len); errdefer gpa.free(job.inputs); var made: usize = 0; errdefer for (job.inputs[0..made]) |input| gpa.free(input); for (request.inputs, 0..) |input, i| { job.inputs[i] = try gpa.dupe(u8, input.bytes); made += 1; } return job; } pub fn deinit(job: *Job, gpa: std.mem.Allocator) void { gpa.free(job.command); gpa.free(job.cwd); for (job.inputs) |input| gpa.free(input); gpa.free(job.inputs); gpa.destroy(job); } }; /// WHY a filter produced nothing, carried home so somebody can be told. /// /// Until this existed the runner read the command's stderr into memory and /// then FREED IT UNREAD — the one artifact that explains a failure, discarded /// two lines after it arrived — and every caller answered a failed filter with /// a bare `return`. `| trr a-z A-Z` (a typo), `| grep nomatch` (exit 1), /// `| jq .` on bad JSON: all of them did nothing, said nothing, and left the /// text alone with no way to find out why. pub const Failure = struct { pub const Kind = enum { /// The command never started: no `/bin/sh`, a cwd that is gone, a NUL /// in the command, a fork that failed. spawn, /// Still running at `command_timeout_seconds`. timeout, /// Past `max_stdout_bytes` / `max_stderr_bytes` / the job total. too_large, /// Ran, and exited nonzero. `code` says which. exit, /// Killed by a signal. signal, /// A read, a write or an allocation failed under us. io, }; kind: Kind = .io, /// Which selection this was, so a report over several cursors can say. index: u32 = 0, /// The exit status, when `kind` is `.exit`. code: u8 = 0, /// stderr exactly as the command wrote it, OWNED by the response. Empty /// when the command said nothing, which is why `kind` and `code` exist. stderr: []u8 = &.{}, }; /// One worker answer. `outputs` owns each slice; a failure owns an empty list /// and, usually, the command's own account of itself in `failure`. pub const Response = struct { id: u32, success: bool, outputs: [][]u8, failure: ?Failure = null, pub fn deinit(response: *Response, gpa: std.mem.Allocator) void { for (response.outputs) |output| gpa.free(output); if (response.outputs.len > 0) gpa.free(response.outputs); if (response.failure) |f| if (f.stderr.len > 0) gpa.free(f.stderr); response.* = undefined; } }; /// What one invocation came to: the bytes, or the reason there are none. pub const Outcome = union(enum) { ok: []u8, failed: Failure }; /// The in-flight set a host keeps while pipes run off its loop. /// /// tty.zig and gui.zig each had this verbatim — same `finish` walk, same /// `cancelAll`, same 16 — differing only in whether `add` asserted or returned /// a bool. It lives here beside the Job it tracks so a third host (the AppKit /// shell, which had no pipe support at all) does not have to grow a fourth. /// /// Bounded on purpose: a filter is a user gesture, and sixteen concurrent ones /// is already more than anybody means. `add` returning false is the host's cue /// to answer the request as failed rather than to queue it. pub const Tasks = struct { pub const capacity = 16; pub const Task = struct { id: u32, future: std.Io.Future(anyerror!void), }; items: [capacity]Task = undefined, len: usize = 0, pub fn full(tasks: *const Tasks) bool { return tasks.len == tasks.items.len; } pub fn add(tasks: *Tasks, task: Task) bool { if (tasks.full()) return false; tasks.items[tasks.len] = task; tasks.len += 1; return true; } /// Join the one that answered and drop it, preserving order so `cancelAll` /// stays deterministic. pub fn finish(tasks: *Tasks, io: std.Io, id: u32) void { for (tasks.items[0..tasks.len], 0..) |*task, i| if (task.id == id) { task.future.await(io) catch {}; tasks.len -= 1; std.mem.copyForwards(Task, tasks.items[i..tasks.len], tasks.items[i + 1 .. tasks.len + 1]); return; }; } pub fn cancelAll(tasks: *Tasks, io: std.Io) void { for (tasks.items[0..tasks.len]) |*task| task.future.cancel(io) catch {}; tasks.len = 0; } }; const WriterContext = struct { io: std.Io, file: std.Io.File, input: []const u8, ok: bool = false, }; fn writeInput(context: *WriterContext) void { defer context.file.close(context.io); context.file.writeStreamingAll(context.io, context.input) catch return; context.ok = true; } /// Run one POSIX shell command with exact stdin, concurrently draining stdout /// and stderr so a producer cannot deadlock against a full pipe. The caller is /// already a worker. Nonzero exit, signal, timeout, IO failure, or either /// output ceiling is reported as `null`; only an exited-zero command transfers /// ownership of stdout to the caller. pub fn runOne( gpa: std.mem.Allocator, io: std.Io, command: []const u8, cwd: []const u8, input: []const u8, ) Outcome { const fail = struct { fn k(kind: Failure.Kind) Outcome { return .{ .failed = .{ .kind = kind } }; } }; if (std.mem.indexOfScalar(u8, command, 0) != null or std.mem.indexOfScalar(u8, cwd, 0) != null) return fail.k(.spawn); var child = std.process.spawn(io, .{ .argv = &.{ "/bin/sh", "-c", command }, .cwd = if (cwd.len == 0) .inherit else .{ .path = cwd }, .stdin = .pipe, .stdout = .pipe, .stderr = .pipe, }) catch return fail.k(.spawn); var writer_context: WriterContext = .{ .io = io, .file = child.stdin.?, .input = input, }; child.stdin = null; // writer_context owns and closes this endpoint var writer: ?std.Thread = std.Thread.spawn(.{}, writeInput, .{&writer_context}) catch { writer_context.file.close(io); child.kill(io); return fail.k(.spawn); }; // On every early return kill first, unblocking a command which never read // stdin, then join the short-lived writer before its borrowed input dies. defer if (writer) |thread| thread.join(); defer child.kill(io); var multi_reader_buffer: std.Io.File.MultiReader.Buffer(2) = undefined; var multi_reader: std.Io.File.MultiReader = undefined; multi_reader.init(gpa, io, multi_reader_buffer.toStreams(), &.{ child.stdout.?, child.stderr.? }); defer multi_reader.deinit(); const stdout_reader = multi_reader.reader(0); const stderr_reader = multi_reader.reader(1); const deadline: std.Io.Timeout = (std.Io.Timeout{ .duration = .{ .clock = .awake, .raw = .fromSeconds(command_timeout_seconds), } }).toDeadline(io); while (multi_reader.fill(64, deadline)) |_| { if (stdout_reader.buffered().len > max_stdout_bytes or stderr_reader.buffered().len > max_stderr_bytes) return fail.k(.too_large); } else |err| switch (err) { error.EndOfStream => {}, // The deadline is the only one of these a human is likely to cause, // and it is the one they are least able to guess at: ten seconds of // nothing used to be followed by nothing. error.Timeout => return fail.k(.timeout), else => return fail.k(.io), } if (stdout_reader.buffered().len > max_stdout_bytes or stderr_reader.buffered().len > max_stderr_bytes) return fail.k(.too_large); multi_reader.checkAnyError() catch return fail.k(.io); const term = child.wait(io) catch return fail.k(.io); writer.?.join(); writer = null; const stdout = multi_reader.toOwnedSlice(0) catch return fail.k(.io); const stderr = multi_reader.toOwnedSlice(1) catch { gpa.free(stdout); return fail.k(.io); }; // stderr is NOT freed here any more. It is the command's own account of // what went wrong, and it goes home with the failure. errdefer gpa.free(stderr); // A STDIN WRITE THAT ENDED EARLY IS NOT A FAILURE. `writer_context.ok` was // part of this condition, so `| head -1` over a selection bigger than the // pipe buffer failed — the command closed stdin after the line it wanted, // the write got EPIPE, and a filter that had done exactly its job reported // nothing. helix joins its input task and ignores the result for this // reason; the exit status is the whole verdict. switch (term) { .exited => |code| if (code != 0) { gpa.free(stdout); return .{ .failed = .{ .kind = .exit, .code = code, .stderr = stderr } }; }, else => { gpa.free(stdout); return .{ .failed = .{ .kind = .signal, .stderr = stderr } }; }, } gpa.free(stderr); return .{ .ok = stdout }; } /// Invoke the command independently for every selection. The response is all /// or nothing: one failure frees every earlier stdout and returns a failed, /// empty answer for the core to ignore. pub fn runJob(gpa: std.mem.Allocator, io: std.Io, job: *const Job) Response { var response: Response = .{ .id = job.id, .success = false, .outputs = &.{} }; const outputs = gpa.alloc([]u8, job.inputs.len) catch return response; var made: usize = 0; var total: usize = 0; for (job.inputs, 0..) |input, i| { const output = switch (runOne(gpa, io, job.command, job.cwd, input)) { .ok => |bytes| bytes, .failed => |f| { // WHICH selection, because with several cursors "it failed" is // not enough to go looking with. var owned = f; owned.index = @intCast(i); response.failure = owned; break; }, }; if (std.math.add(usize, total, output.len) catch null) |next_total| { if (next_total <= max_total_stdout_bytes) { outputs[i] = output; total = next_total; made += 1; continue; } } gpa.free(output); response.failure = .{ .kind = .too_large, .index = @intCast(i) }; break; } if (made != job.inputs.len) { for (outputs[0..made]) |output| gpa.free(output); gpa.free(outputs); return response; } response.success = true; response.outputs = outputs; return response; } test "native pipe runner preserves stdin/stdout bytes and reports how it failed" { const gpa = std.testing.allocator; const io = std.testing.io; const output = switch (runOne(gpa, io, "printf 'prefix:'; cat; printf '\\n'", "/tmp", "a\x00b\n")) { .ok => |bytes| bytes, .failed => return error.PipeCommandFailed, }; defer gpa.free(output); try std.testing.expectEqualSlices(u8, "prefix:a\x00b\n\n", output); // A NONZERO EXIT COMES HOME WITH ITS REASON. The status and the command's // own stderr are the whole of what a human needs to fix a typo'd filter, // and both used to be freed on the floor. switch (runOne(gpa, io, "echo trouble >&2; exit 7", "/tmp", "")) { .ok => |bytes| { gpa.free(bytes); return error.PipeShouldHaveFailed; }, .failed => |f| { defer gpa.free(f.stderr); try std.testing.expectEqual(Failure.Kind.exit, f.kind); try std.testing.expectEqual(@as(u8, 7), f.code); try std.testing.expectEqualSlices(u8, "trouble\n", f.stderr); }, } // ...and a command that stops reading its stdin SUCCEEDS. `| head -1` over // a selection bigger than the pipe buffer takes the line it wanted and // closes the pipe; the write gets EPIPE, and that used to fail the filter // even though it had done exactly its job. const big = try gpa.alloc(u8, 512 * 1024); defer gpa.free(big); @memset(big, 'x'); big[0] = 'a'; big[1] = '\n'; switch (runOne(gpa, io, "head -1", "/tmp", big)) { .ok => |bytes| { defer gpa.free(bytes); try std.testing.expectEqualSlices(u8, "a\n", bytes); }, .failed => return error.EarlyStdinCloseShouldNotFail, } }