//! 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"); const filesystem = @import("fs.zig"); /// 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. /// An Edit's `<`, `|` and `>` (src/edit_cmd.zig) run a command of their own /// per input, each in its file's directory: `commands` and `cwds`, when /// given, name input i's, and `shell` (the session's Shell) runs them, /// /bin/sh when it is empty. pub const Request = struct { id: u32, command: []const u8, cwd: []const u8, inputs: []const Input, commands: []const []const u8 = &.{}, cwds: []const []const u8 = &.{}, shell: []const u8 = "", /// Nonzero: `stop(token)` from another thread kills its commands. token: u32 = 0, /// The commands' environment, `NAME=value` each, when given; else this /// process's. env: []const []const u8 = &.{}, }; /// 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, /// Input i's own command and directory, when the request gave them. commands: [][]u8 = &.{}, cwds: [][]u8 = &.{}, shell: []u8 = &.{}, token: u32 = 0, env: [][]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, .token = request.token, .command = try gpa.dupe(u8, request.command), .cwd = &.{}, .inputs = &.{}, }; errdefer gpa.free(job.command); job.cwd = try gpa.dupe(u8, filesystem.localPath(request.cwd) orelse request.cwd); errdefer gpa.free(job.cwd); job.inputs = try dupeAll(gpa, request.inputs.len, request.inputs, false); errdefer freeAll(gpa, job.inputs); job.commands = try dupeAll(gpa, request.commands.len, request.commands, false); errdefer freeAll(gpa, job.commands); job.cwds = try dupeAll(gpa, request.cwds.len, request.cwds, true); errdefer freeAll(gpa, job.cwds); job.shell = try gpa.dupe(u8, request.shell); errdefer gpa.free(job.shell); job.env = try dupeAll(gpa, request.env.len, request.env, false); return job; } /// Owned copies of `items` (Inputs or strings); a directory as the /// host's own path, as `cwd` is. fn dupeAll(gpa: std.mem.Allocator, n: usize, items: anytype, dirs: bool) ![][]u8 { const out = try gpa.alloc([]u8, n); var made: usize = 0; errdefer { for (out[0..made]) |item| gpa.free(item); gpa.free(out); } for (items) |item| { const bytes: []const u8 = if (@TypeOf(item) == Input) item.bytes else item; out[made] = try gpa.dupe(u8, if (dirs) filesystem.localPath(bytes) orelse bytes else bytes); made += 1; } return out; } fn freeAll(gpa: std.mem.Allocator, items: [][]u8) void { for (items) |item| gpa.free(item); gpa.free(items); } pub fn deinit(job: *Job, gpa: std.mem.Allocator) void { gpa.free(job.command); gpa.free(job.cwd); freeAll(gpa, job.inputs); freeAll(gpa, job.commands); freeAll(gpa, job.cwds); gpa.free(job.shell); freeAll(gpa, job.env); 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 { return runIn(gpa, io, "", command, cwd, input, 0, &.{}); } /// The commands running for a stoppable request (`Request.token`), so /// another thread can stop them: each is a process group of its own, which /// `stop` kills whole, a `sleep` under the shell too. const Running = struct { var lock: std.atomic.Mutex = .unlocked; var procs: [64]struct { token: u32, pid: std.posix.pid_t } = undefined; var len: usize = 0; /// Tokens stopped lately: a command not started yet is not. var stopped: [16]u32 = @splat(0); var next: usize = 0; fn take() void { while (!lock.tryLock()) std.atomic.spinLoopHint(); } fn isStopped(token: u32) bool { return std.mem.indexOfScalar(u32, &stopped, token) != null; } /// Whether `pid` may run on: false once its token was stopped. fn add(token: u32, pid: std.posix.pid_t) bool { take(); defer lock.unlock(); if (isStopped(token)) return false; if (len < procs.len) { procs[len] = .{ .token = token, .pid = pid }; len += 1; } return true; } fn remove(pid: std.posix.pid_t) void { take(); defer lock.unlock(); for (procs[0..len], 0..) |p, i| if (p.pid == pid) { procs[i] = procs[len - 1]; len -= 1; return; }; } }; /// Stops the request `token` names: its running commands' process groups /// are killed, and no command of it starts after. From any thread. pub fn stop(token: u32) void { if (token == 0) return; Running.take(); defer Running.lock.unlock(); Running.stopped[Running.next] = token; Running.next = (Running.next + 1) % Running.stopped.len; for (Running.procs[0..Running.len]) |p| if (p.token == token) { std.posix.kill(-p.pid, .KILL) catch {}; }; } /// A process-unique token for a stoppable request. pub fn newToken() u32 { const t = next_token.fetchAdd(1, .monotonic); return if (t == 0) next_token.fetchAdd(1, .monotonic) else t; } var next_token: std.atomic.Value(u32) = .init(1); /// `runOne` through `shell -c` (/bin/sh when it is empty). pub fn runIn( gpa: std.mem.Allocator, io: std.Io, shell: []const u8, command: []const u8, cwd: []const u8, input: []const u8, token: u32, env: []const []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 or std.mem.indexOfScalar(u8, shell, 0) != null) return fail.k(.spawn); // std's spawn runs no code in the child: this thread's mask, which the // child inherits, is cleared across the fork (only the tty's SIGWINCH is // ever blocked, and its default is to be ignored). A handler resets at // exec by itself, and nothing here is ignored. var environ: std.process.Environ.Map = .init(gpa); defer environ.deinit(); for (env) |line| { const eq = std.mem.indexOfScalar(u8, line, '=') orelse continue; environ.put(line[0..eq], line[eq + 1 ..]) catch return fail.k(.spawn); } const none = std.posix.sigemptyset(); var kept: std.posix.sigset_t = undefined; std.posix.sigprocmask(std.posix.SIG.SETMASK, &none, &kept); defer std.posix.sigprocmask(std.posix.SIG.SETMASK, &kept, null); var child = std.process.spawn(io, .{ .argv = &.{ if (shell.len == 0) "/bin/sh" else shell, "-c", command }, .cwd = if (cwd.len == 0) .inherit else .{ .path = cwd }, .environ_map = if (env.len > 0) &environ else null, .stdin = .pipe, .stdout = .pipe, .stderr = .pipe, // A group of its own: a stop or a timeout kills what it started too. .pgid = 0, }) catch return fail.k(.spawn); const pid = child.id.?; 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); // Before the shell is reaped, while its group id cannot be another's. defer if (child.id != null) std.posix.kill(-pid, .KILL) catch {}; defer if (token != 0) Running.remove(pid); if (token != 0 and !Running.add(token, pid)) return fail.k(.signal); 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); if (token != 0) Running.remove(pid); 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 }; } fn tokenStopped(token: u32) bool { Running.take(); defer Running.lock.unlock(); return Running.isStopped(token); } /// 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 command = if (job.commands.len > 0) job.commands[i] else job.command; const cwd = if (job.cwds.len > 0) job.cwds[i] else job.cwd; if (job.token != 0 and tokenStopped(job.token)) { response.failure = .{ .kind = .signal, .index = @intCast(i) }; break; } const output = switch (runIn(gpa, io, job.shell, command, cwd, input, job.token, job.env)) { .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 selection filters use the physical directory of an explicit OS mount" { const gpa = std.testing.allocator; var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); var directory_buf: [4096]u8 = undefined; const directory = directory_buf[0..try tmp.dir.realPath(std.testing.io, &directory_buf)]; var declared_buf: [4102]u8 = undefined; const declared = try std.fmt.bufPrint(&declared_buf, "/n/os{s}", .{directory}); const job = try Job.copy(gpa, .{ .id = 1, .command = "pwd", .cwd = declared, .inputs = &.{.{ .bytes = "" }} }); defer job.deinit(gpa); try std.testing.expectEqualStrings(directory, job.cwd); switch (runOne(gpa, std.testing.io, job.command, job.cwd, "")) { .ok => |bytes| { defer gpa.free(bytes); try std.testing.expectEqualStrings(directory, std.mem.trimEnd(u8, bytes, "\n")); }, .failed => |failure| { gpa.free(failure.stderr); return error.FilterFailed; }, } } 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, } } test "a request may give each input its own command and directory, run by the shell it names" { const gpa = std.testing.allocator; const job = try Job.copy(gpa, .{ .id = 3, .command = "", .cwd = "", .inputs = &.{ .{ .bytes = "" }, .{ .bytes = "abc" } }, .commands = &.{ "pwd", "tr a-z A-Z" }, .cwds = &.{ "/tmp", "/" }, .shell = "/bin/sh", }); defer job.deinit(gpa); var response = runJob(gpa, std.testing.io, job); defer response.deinit(gpa); try std.testing.expect(response.success); try std.testing.expectEqualStrings("/tmp\n", response.outputs[0]); try std.testing.expectEqualStrings("ABC", response.outputs[1]); } test "stop kills a running request's commands, the whole group, and starts none after" { const gpa = std.testing.allocator; const token = newToken(); const job = try Job.copy(gpa, .{ .id = 1, .command = "", .cwd = "", .inputs = &.{ .{ .bytes = "" }, .{ .bytes = "" } }, .commands = &.{ "sleep 30; echo late", "echo never" }, .token = token, }); defer job.deinit(gpa); const Run = struct { fn go(j: *const Job, out: *Response) void { out.* = runJob(std.testing.allocator, std.testing.io, j); } }; var response: Response = undefined; const started = std.Io.Clock.awake.now(std.testing.io); const thread = try std.Thread.spawn(.{}, Run.go, .{ job, &response }); std.Io.sleep(std.testing.io, .fromMilliseconds(300), .awake) catch {}; stop(token); thread.join(); defer response.deinit(gpa); try std.testing.expect(!response.success); try std.testing.expectEqual(Failure.Kind.signal, response.failure.?.kind); try std.testing.expectEqual(@as(u32, 0), response.failure.?.index); const took = started.durationTo(std.Io.Clock.awake.now(std.testing.io)); try std.testing.expect(took.toMilliseconds() < 5000); }