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.zig291
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);
+}