diff options
Diffstat (limited to 'src/selection_pipe.zig')
| -rw-r--r-- | src/selection_pipe.zig | 120 |
1 files changed, 118 insertions, 2 deletions
diff --git a/src/selection_pipe.zig b/src/selection_pipe.zig index 4fb024bb..fb6587f9 100644 --- a/src/selection_pipe.zig +++ b/src/selection_pipe.zig @@ -30,6 +30,8 @@ pub const Request = struct { commands: []const []const u8 = &.{}, cwds: []const []const u8 = &.{}, shell: []const u8 = "", + /// Nonzero: `stop(token)` from another thread kills its commands. + token: u32 = 0, }; /// Worker-owned snapshot. `copy` is intentionally called synchronously while @@ -44,12 +46,14 @@ pub const Job = struct { commands: [][]u8 = &.{}, cwds: [][]u8 = &.{}, shell: []u8 = &.{}, + token: u32 = 0, 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 = &.{}, @@ -228,9 +232,71 @@ pub fn runOne( cwd: []const u8, input: []const u8, ) Outcome { - return runIn(gpa, io, "", command, cwd, input); + 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, @@ -239,6 +305,7 @@ pub fn runIn( command: []const u8, cwd: []const u8, input: []const u8, + token: u32, ) Outcome { const fail = struct { fn k(kind: Failure.Kind) Outcome { @@ -263,7 +330,10 @@ pub fn runIn( .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 = .{ .io = io, @@ -280,6 +350,10 @@ pub fn runIn( // stdin, then join the short-lived writer before its borrowed input dies. defer if (writer) |thread| thread.join(); defer child.kill(io); + // Before the shell is reaped, while its group id cannot be another's. + defer if (child.id != null) std.posix.kill(-pid, .KILL) catch {}; + 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; @@ -306,6 +380,7 @@ pub fn runIn( 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; @@ -339,6 +414,12 @@ pub fn runIn( 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. @@ -350,7 +431,11 @@ pub fn runJob(gpa: std.mem.Allocator, io: std.Io, job: *const Job) Response { for (job.inputs, 0..) |input, i| { 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; - const output = switch (runIn(gpa, io, job.shell, command, cwd, input)) { + 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)) { .ok => |bytes| bytes, .failed => |f| { // WHICH selection, because with several cursors "it failed" is @@ -468,3 +553,34 @@ test "a request may give each input its own command and directory, run by the sh 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); +} |
