summaryrefslogtreecommitdiff
path: root/src/selection_pipe.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/selection_pipe.zig')
-rw-r--r--src/selection_pipe.zig214
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);
+}