diff options
Diffstat (limited to 'src/selection_pipe.zig')
| -rw-r--r-- | src/selection_pipe.zig | 291 |
1 files changed, 268 insertions, 23 deletions
diff --git a/src/selection_pipe.zig b/src/selection_pipe.zig index 44cc1709..c3e5e3a6 100644 --- a/src/selection_pipe.zig +++ b/src/selection_pipe.zig @@ -18,11 +18,23 @@ 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 @@ -33,12 +45,19 @@ pub const Job = struct { 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 = &.{}, @@ -46,22 +65,48 @@ pub const Job = struct { 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 gpa.alloc([]u8, request.inputs.len); - errdefer gpa.free(job.inputs); + 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 (job.inputs[0..made]) |input| gpa.free(input); - for (request.inputs, 0..) |input, i| { - job.inputs[i] = try gpa.dupe(u8, input.bytes); + 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 job; + 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); - for (job.inputs) |input| gpa.free(input); - gpa.free(job.inputs); + freeAll(gpa, job.inputs); + freeAll(gpa, job.commands); + freeAll(gpa, job.cwds); + gpa.free(job.shell); + freeAll(gpa, job.env); gpa.destroy(job); } }; @@ -169,17 +214,53 @@ pub const Tasks = struct { } }; +/// The stdin writer's own: its copy of the input, freed by the writer when +/// it is done, so a writer stuck on a command that never reads (and cannot +/// be killed, below) can be let go of without its input dying under it. const WriterContext = struct { + gpa: std.mem.Allocator, io: std.Io, file: std.Io.File, - input: []const u8, - ok: bool = false, + input: []u8, }; fn writeInput(context: *WriterContext) void { - defer context.file.close(context.io); + defer { + context.file.close(context.io); + context.gpa.free(context.input); + context.gpa.destroy(context); + } context.file.writeStreamingAll(context.io, context.input) catch return; - context.ok = true; +} + +/// Kills the command's group, then `abandon`s it. +fn abandonNow(io: std.Io, child: *std.process.Child, pid: std.posix.pid_t) void { + // Before the shell is reaped, while its group id cannot be another's. + std.posix.kill(-pid, .KILL) catch {}; + std.posix.kill(pid, .KILL) catch {}; + abandon(io, child.*); + child.id = null; + child.stdout = null; + child.stderr = null; +} + +/// A command given up on, killed and reaped where nothing waits for it: a +/// process in uninterruptible sleep (a write to a file of this very +/// session, through its mount, waiting on the request whose answer is +/// waiting on this command) dies only when that request is answered, and +/// the answer must not wait on its death. +fn abandon(io: std.Io, child: std.process.Child) void { + const Reap = struct { + fn run(i: std.Io, c: std.process.Child) void { + var reaped = c; + reaped.kill(i); + } + }; + const thread = std.Thread.spawn(.{}, Reap.run, .{ io, child }) catch { + var c = child; + return c.kill(io); + }; + thread.detach(); } /// Run one POSIX shell command with exact stdin, concurrently draining stdout @@ -194,45 +275,146 @@ pub fn runOne( 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) return fail.k(.spawn); + 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 = &.{ "/bin/sh", "-c", command }, + .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 = .{ + const writer_context = gpa.create(WriterContext) catch { + abandonNow(io, &child, pid); + return fail.k(.spawn); + }; + writer_context.* = .{ + .gpa = gpa, .io = io, .file = child.stdin.?, - .input = input, + .input = gpa.dupe(u8, input) catch { + gpa.destroy(writer_context); + abandonNow(io, &child, pid); + return fail.k(.spawn); + }, }; child.stdin = null; // writer_context owns and closes this endpoint - var writer: ?std.Thread = std.Thread.spawn(.{}, writeInput, .{&writer_context}) catch { + var writer: ?std.Thread = std.Thread.spawn(.{}, writeInput, .{writer_context}) catch { writer_context.file.close(io); - child.kill(io); + gpa.free(writer_context.input); + gpa.destroy(writer_context); + abandonNow(io, &child, pid); 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); + // On every early return (a timeout, a ceiling) the answer goes at once: + // the group is killed and the shell left to be reaped elsewhere + // (abandon), and the writer, which owns its input, let go. + defer if (writer) |thread| thread.detach(); + defer if (child.id != null) abandonNow(io, &child, pid); + 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; @@ -259,6 +441,7 @@ pub fn runOne( 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; @@ -292,6 +475,12 @@ pub fn runOne( 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. @@ -301,7 +490,13 @@ pub fn runJob(gpa: std.mem.Allocator, io: std.Io, job: *const Job) 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)) { + 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 @@ -400,3 +595,53 @@ test "native pipe runner preserves stdin/stdout bytes and reports how it failed" .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); +} |
