summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/9p_io.zig')
-rw-r--r--src/9p_io.zig86
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 {