diff options
Diffstat (limited to 'src/selection_pipe.zig')
| -rw-r--r-- | src/selection_pipe.zig | 214 |
1 files changed, 214 insertions, 0 deletions
diff --git a/src/selection_pipe.zig b/src/selection_pipe.zig new file mode 100644 index 00000000..65c4c7fd --- /dev/null +++ b/src/selection_pipe.zig @@ -0,0 +1,214 @@ +//! 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; + } +}; + +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); +} |
