diff options
Diffstat (limited to 'src/9p_io.zig')
| -rw-r--r-- | src/9p_io.zig | 86 |
1 files changed, 68 insertions, 18 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 6f05e3ab..78ff5d59 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -1947,14 +1947,8 @@ pub const Client = struct { } }; - fn transact( - s: *Session, - names: []const []const u8, - out: *std.Io.Writer.Allocating, - remote: *RemoteError, - write_bytes: ?[]const u8, - read_limit: usize, - ) !void { + /// Version, attach and a walk to `names`: the fid there and its qid. + fn walkTo(s: *Session, names: []const []const u8, remote: *RemoteError) !struct { fid: u32, qid: ninep.Qid } { _ = try s.ask(.{ .version = .{} }, remote); if (s.cl.msize == 0) return Error.Botch; @@ -1978,6 +1972,20 @@ pub const Client = struct { next = if (next == 1) 2 else 1; i += n; } + return .{ .fid = cur, .qid = here }; + } + + fn transact( + s: *Session, + names: []const []const u8, + out: *std.Io.Writer.Allocating, + remote: *RemoteError, + write_bytes: ?[]const u8, + read_limit: usize, + ) !void { + const at = try walkTo(s, names, remote); + const cur = at.fid; + const here = at.qid; defer s.dropNoWait(cur); const directory = here.type & ninep.qtdir != 0; @@ -2036,15 +2044,20 @@ pub const Client = struct { read_limit: usize, ) ![]u8 { if (comptime !supported) return Error.Dial; + const s = try startSession(gpa, sock, display_path); + defer endSession(gpa, s); + var out: std.Io.Writer.Allocating = .init(gpa); + errdefer out.deinit(); + try transact(s, names, &out, remote, write_bytes, read_limit); + return out.toOwnedSlice(); + } + + /// A connected session whose requests have the usual deadline. + fn startSession(gpa: std.mem.Allocator, sock: Dial, display_path: []const u8) !*Session { const deadline = nowMs() +| budget_ms; const s = try gpa.create(Session); s.* = .{ .fd = -1, .deadline = deadline, .display_path = display_path }; - defer { - if (quic_enabled and s.quic != null) { - s.quic.?.deinit(); - } else if (s.fd >= 0) _ = libc.close(s.fd); - gpa.destroy(s); - } + errdefer endSession(gpa, s); if (sock == .quic) { if (comptime quic_enabled) { s.quic = quic.Connection.dial(sock.quic) catch return Error.Dial; @@ -2052,11 +2065,48 @@ pub const Client = struct { } else return Error.QuicUnavailable; } else s.fd = try connect(sock, deadline); s.cl = .init(.{ .in = &s.in, .out = &s.out }); + return s; + } - var out: std.Io.Writer.Allocating = .init(gpa); - errdefer out.deinit(); - try transact(s, names, &out, remote, write_bytes, read_limit); - return out.toOwnedSlice(); + fn endSession(gpa: std.mem.Allocator, s: *Session) void { + if (quic_enabled and s.quic != null) { + s.quic.?.deinit(); + } else if (s.fd >= 0) _ = libc.close(s.fd); + gpa.destroy(s); + } + + /// Opens `path` on one connection and writes `first` to it (/log's + /// `follow new`), all with the usual deadline; then, once `ready(ctx)` + /// has said it is not done already, reads it with no deadline, a read at + /// a time, until `record(ctx, bytes)` says done. An end of file or a + /// dropped connection is Hangup. What `pardes --wait` blocks on. + pub fn follow( + gpa: std.mem.Allocator, + dial: []const u8, + path: []const u8, + first: []const u8, + ctx: anytype, + comptime ready: fn (@TypeOf(ctx)) bool, + comptime record: fn (@TypeOf(ctx), []const u8) bool, + ) !void { + if (comptime !supported) return error.Unsupported; + var names: [max_depth][]const u8 = undefined; + const n = try elements(path, &names); + var sock_buf: [sun_path_len]u8 = undefined; + const sock = try resolve(&sock_buf, dial); + var remote: RemoteError = .{}; + const s = try startSession(gpa, sock, path); + defer endSession(gpa, s); + const fid = (try walkTo(s, names[0..n], &remote)).fid; + _ = try s.ask(.{ .open = .{ .fid = fid, .mode = ninep.ordwr } }, &remote); + _ = try s.ask(.{ .write = .{ .fid = fid, .offset = 0, .data = first } }, &remote); + if (ready(ctx)) return; + s.deadline = std.math.maxInt(i64); + while (true) { + const data = (try s.ask(.{ .read = .{ .fid = fid, .offset = 0, .count = s.cl.maxRead() } }, &remote)).read; + if (data.len == 0) return Error.Hangup; + if (record(ctx, data)) return; + } } fn connect(sock: Dial, deadline: i64) Error!c_int { |
