//! 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); } }; /// One worker answer. `outputs` owns each slice; a failure owns an empty list. pub const Response = struct { id: u32, success: bool, outputs: [][]u8, 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); response.* = undefined; } }; /// 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, ) ?[]u8 { if (std.mem.indexOfScalar(u8, command, 0) != null or std.mem.indexOfScalar(u8, cwd, 0) != null) return null; 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 null; 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 null; }; // 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 null; } else |err| switch (err) { error.EndOfStream => {}, else => return null, } if (stdout_reader.buffered().len > max_stdout_bytes or stderr_reader.buffered().len > max_stderr_bytes) return null; multi_reader.checkAnyError() catch return null; const term = child.wait(io) catch return null; writer.?.join(); writer = null; const stdout = multi_reader.toOwnedSlice(0) catch return null; const stderr = multi_reader.toOwnedSlice(1) catch { gpa.free(stdout); return null; }; defer gpa.free(stderr); const exited_zero = switch (term) { .exited => |code| code == 0, else => false, }; if (!writer_context.ok or !exited_zero) { gpa.free(stdout); return null; } return 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 = runOne(gpa, io, job.command, job.cwd, input) orelse 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); 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 rejects failure" { const gpa = std.testing.allocator; const io = std.testing.io; const output = runOne(gpa, io, "printf 'prefix:'; cat; printf '\\n'", "/tmp", "a\x00b\n") orelse return error.PipeCommandFailed; defer gpa.free(output); try std.testing.expectEqualSlices(u8, "prefix:a\x00b\n\n", output); try std.testing.expect(runOne(gpa, io, "printf ignored; exit 7", "/tmp", "") == null); }