summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--.gitignore4
-rw-r--r--README.md51
-rw-r--r--build.zig74
-rw-r--r--build.zig.zon7
-rw-r--r--docs/design.md73
-rw-r--r--docs/results/batch.json60
-rw-r--r--docs/results/fragmented.json60
-rw-r--r--docs/spec.md33
-rw-r--r--docs/validation.md73
-rw-r--r--src/Server.zig194
-rw-r--r--src/client.zig392
-rw-r--r--src/quic.zig302
-rw-r--r--src/root.zig53
-rw-r--r--src/session_test.zig222
-rw-r--r--src/transport.zig207
-rw-r--r--src/wire.zig1153
-rw-r--r--test/differential/go.mod14
-rw-r--r--test/differential/go.sum58
-rw-r--r--test/differential/main.go172
-rw-r--r--test/differential/probe.zig59
-rw-r--r--test/differential/session.go45
-rw-r--r--test/differential/session.zig82
-rw-r--r--test/fuzz.zig70
-rw-r--r--test/quic.zig268
-rw-r--r--test/transport.zig82
25 files changed, 3808 insertions, 0 deletions
diff --git a/.gitignore b/.gitignore
new file mode 100644
index 0000000..af7bb5c
--- /dev/null
+++ b/.gitignore
@@ -0,0 +1,4 @@
+.zig-cache/
+zig-out/
+test/differential/bin/
+test/differential/results/
diff --git a/README.md b/README.md
new file mode 100644
index 0000000..226ec78
--- /dev/null
+++ b/README.md
@@ -0,0 +1,51 @@
+# cloud9
+
+A Zig 0.16 library for base **9P2000** clients, servers, and transports.
+
+The protocol core uses caller-owned buffers and bounded request tables. It has
+no heap allocation, OS calls, threads, or filesystem policy. TCP and Unix stream
+adapters support `std.Io`; nonblocking POSIX adapters support existing event loops.
+Optional QUIC transport takes caller-provided OpenSSL bindings and an ALPN name.
+
+Mounting, namespace discovery, exported trees, permissions, and application
+lifecycle stay with the application. There are no 9P2000.u or 9P2000.L messages.
+
+```sh
+zig build test
+zig build transport-test
+zig build quic-test -Dquic=true # system OpenSSL 3.6+
+zig build fuzz -Doptimize=ReleaseSafe -- 4200 1000000
+zig build differential -Doptimize=ReleaseSafe -- --seed 4200 --rounds 1000
+zig build differential -Doptimize=ReleaseSafe -- --seed 4201 --rounds 100 --chunk 1
+```
+
+The differential harness requires Go, downloads pinned test-only dependencies,
+and saves a corpus, configuration, and JSON report under
+`test/differential/results`. It compares valid messages with 9fans and go9p and
+runs the cloud9 client against go9p's filesystem server over local pipes.
+Malformed-input probes run only against cloud9. No remote targets are contacted.
+
+To use a sibling checkout, add `.cloud9 = .{ .path = "../cloud9" }` to the package
+dependencies, then import its module:
+
+```zig
+const cloud9 = b.dependency("cloud9", .{
+ .target = target,
+ .optimize = optimize,
+}).module("cloud9");
+app.root_module.addImport("cloud9", cloud9);
+```
+
+Start a client with `Client.init(.{ .in = input_buffer, .out = output_buffer })`.
+Submit `.version` first. Drain `output()` through your transport and report the
+number sent with `wrote()`. Feed received bytes using `push()` and collect tagged
+results with `take()`. All base requests are available through `Client.Request`.
+
+For a server, call `Server.receive()` after `push()`. A request borrows the input
+until `release()`. After releasing backend fids and canceling outstanding work,
+answer Tversion with `negotiate()`. Answer other requests with `reply()`. The
+backend implements its filesystem and fid lifecycle; cloud9 checks reply types,
+tags, counts, negotiated frame sizes, and flush completion.
+
+See [design and ownership contracts](docs/design.md),
+[specification references](docs/spec.md), and [validation](docs/validation.md).
diff --git a/build.zig b/build.zig
new file mode 100644
index 0000000..160efdf
--- /dev/null
+++ b/build.zig
@@ -0,0 +1,74 @@
+const std = @import("std");
+pub fn build(b: *std.Build) void {
+ const target = b.standardTargetOptions(.{});
+ const optimize = b.standardOptimizeOption(.{});
+ const module = b.addModule("cloud9", .{
+ .root_source_file = b.path("src/root.zig"),
+ .target = target,
+ .optimize = optimize,
+ });
+ const tests = b.addTest(.{ .root_module = module });
+ b.step("test", "Run protocol and session tests").dependOn(&b.addRunArtifact(tests).step);
+ const probe = b.addExecutable(.{
+ .name = "cloud9-probe",
+ .root_module = b.createModule(.{
+ .root_source_file = b.path("test/differential/probe.zig"),
+ .target = target,
+ .optimize = optimize,
+ .link_libc = true,
+ .imports = &.{.{ .name = "cloud9", .module = module }},
+ }),
+ });
+ b.installArtifact(probe);
+ const session_probe = b.addExecutable(.{
+ .name = "cloud9-session-probe",
+ .root_module = b.createModule(.{
+ .root_source_file = b.path("test/differential/session.zig"),
+ .target = target,
+ .optimize = optimize,
+ .link_libc = true,
+ .imports = &.{.{ .name = "cloud9", .module = module }},
+ }),
+ });
+
+ const diff = b.addSystemCommand(&.{ "go", "run", "." });
+ diff.setCwd(b.path("test/differential"));
+ diff.addArg("--probe");
+ diff.addArtifactArg(probe);
+ diff.addArg("--session-probe");
+ diff.addArtifactArg(session_probe);
+ if (b.args) |args| diff.addArgs(args);
+ b.step("differential", "Compare valid traffic with pinned 9fans and go9p").dependOn(&diff.step);
+ const transport_tests = b.addTest(.{ .root_module = b.createModule(.{
+ .root_source_file = b.path("test/transport.zig"),
+ .target = target,
+ .optimize = optimize,
+ .link_libc = true,
+ .imports = &.{.{ .name = "cloud9", .module = module }},
+ }) });
+ b.step("transport-test", "Test TCP, Unix and std.Io adapters").dependOn(&b.addRunArtifact(transport_tests).step);
+ if (b.option(bool, "quic", "Test QUIC with system OpenSSL 3.6+") orelse false) {
+ const header = b.addWriteFiles().add("openssl.h", "#include <openssl/ssl.h>\n#include <openssl/quic.h>\n#include <openssl/err.h>\n");
+ const translated = b.addTranslateC(.{ .root_source_file = header, .target = target, .optimize = optimize, .link_libc = true });
+ const ssl = translated.createModule();
+ const tests_quic = b.addTest(.{ .root_module = b.createModule(.{
+ .root_source_file = b.path("test/quic.zig"),
+ .target = target,
+ .optimize = optimize,
+ .link_libc = true,
+ .imports = &.{ .{ .name = "cloud9", .module = module }, .{ .name = "openssl", .module = ssl } },
+ }) });
+ tests_quic.root_module.linkSystemLibrary("ssl", .{});
+ tests_quic.root_module.linkSystemLibrary("crypto", .{});
+ b.step("quic-test", "Test optional QUIC transport").dependOn(&b.addRunArtifact(tests_quic).step);
+ }
+ const fuzz = b.addExecutable(.{ .name = "cloud9-fuzz", .root_module = b.createModule(.{
+ .root_source_file = b.path("test/fuzz.zig"),
+ .target = target,
+ .optimize = optimize,
+ .imports = &.{.{ .name = "cloud9", .module = module }},
+ }) });
+ const run_fuzz = b.addRunArtifact(fuzz);
+ if (b.args) |args| run_fuzz.addArgs(args);
+ b.step("fuzz", "Run deterministic local decoder and server probes (seed, iterations)").dependOn(&run_fuzz.step);
+}
diff --git a/build.zig.zon b/build.zig.zon
new file mode 100644
index 0000000..0dbc3cf
--- /dev/null
+++ b/build.zig.zon
@@ -0,0 +1,7 @@
+.{
+ .name = .cloud9,
+ .fingerprint = 0xc8b2d5eaa3adfca,
+ .version = "0.1.0",
+ .minimum_zig_version = "0.16.0",
+ .paths = .{ "build.zig", "build.zig.zon", "src", "test", "docs", "README.md" },
+}
diff --git a/docs/design.md b/docs/design.md
new file mode 100644
index 0000000..ed9a66e
--- /dev/null
+++ b/docs/design.md
@@ -0,0 +1,73 @@
+# API and ownership
+
+The public entry point is `src/root.zig`. `Msg` covers every base request and
+response, excluding the reserved, nonexistent Terror. Decode returns borrowed
+strings/data. Encode checks complete output size before writing; input payloads
+may already occupy their exact final position in the output frame, but arbitrary
+input/output overlap is not supported. Stat uses the protocol's inner length;
+Rstat and Twstat add and validate the outer length separately.
+
+`Client` supports 16 ordinary outstanding requests and one additional flush.
+Replies can arrive out of order. Tags remain reserved while flushes are pending,
+even when the original request has already completed. Multiple flushes and
+flushing a flush are accounted for. A completed result's slices remain valid
+until the next `take()` or `hangup()`. `push()` only appends and never compacts.
+Submit copies request payloads into the output buffer. An iounit returned by an
+open/create is per fid; the caller must apply it when selecting atomic I/O sizes.
+`maxRead()` and `maxWrite()` describe frame limits, not that per-fid guarantee.
+Renegotiation requires an empty client request window.
+
+`Server` supports 64 ordinary outstanding requests and one additional flush.
+Receive errors terminate the connection. A full request table produces NoTags;
+output backpressure instead returns null without consuming the next request.
+`receive()` yields at most one frame. `release()` ends its borrow and advances
+input. A backend retaining any request data must copy it before release.
+A successful reply copies output and frees the request tag. An output-capacity
+error leaves that tag pending so the backend can drain and retry.
+
+The server validates transaction structure, not filesystem policy. Backends
+own fid maps, authentication state, file handles, access modes, complete directory
+records, atomic wstat, and cancellations. They must retire canceled work before
+reusing its tag; a tag alone cannot distinguish a stale backend completion from a
+new request. A version event requires backend cancellation and fid cleanup before
+`negotiate()`. A partially sent response drains before a new version is accepted.
+An Rflush promises no later reply to oldtag. Replying Rerror to a flush is refused.
+
+Caller-provided input/output buffers must be disjoint and remain stable for the
+session. No protocol API takes an allocator. Server and client request tables are
+fixed-size. Copies and work per frame are bounded by buffer size and fixed table
+capacities; transports and application backends can have their own allocations.
+
+# Transports
+
+`transport.readFrame` and `writeFrame` adapt `std.Io.Reader` and `std.Io.Writer`.
+The writer remains buffered until its owner flushes it. Reader errors leave the
+stream unsuitable for further frame processing; close the connection.
+
+`transport.connect` and `listen` use `std.Io.net` for TCP and Unix streams. The
+POSIX `connectFd`, `listenFd`, `acceptFd`, `read`, `write`, and `wait` APIs fit
+caller-owned poll loops. They configure nonblocking descriptors and close-on-exec;
+read/write return null for retry, and read returns zero for EOF. `wait` takes an
+absolute monotonic deadline. The application closes descriptors, limits its
+connections, selects addresses and deadlines, and manages Unix path permissions
+and removal. No mounting or namespace policy is performed.
+
+`Quic(OpenSSL, alpn)` is an optional OpenSSL 3.6+ adapter with a single ordered
+bidirectional stream per connection. It preserves the previous Pardes transport:
+ephemeral self-signed certificates, no peer authentication. Applications needing
+identity verification must provide that policy before using it across a trust
+boundary. OpenSSL allocations and handshake costs are outside the allocation-free
+protocol core. Accepted connections must close before their shared listener.
+
+# Pardes integration
+
+Pardes consumes the sibling package through build.zig.zon. Its `src/9p.zig` now
+adapts filesystem requests and retains editor/GPIO capacities and permissions.
+`src/9p_io.zig` retains mounting, discovery, the editor's event-loop scheduling,
+connection limits, and error presentation. Protocol bytes, session validation,
+and TCP/Unix/QUIC transport implementation come from cloud9. Pardes selects the
+existing `pardes-9p` ALPN for compatibility.
+
+The separate `05-zig-p4` firmware build adds cloud9 to the GPIO application's
+module map. UART hardware access remains in the board firmware. No hardware was
+flashed by this change.
diff --git a/docs/results/batch.json b/docs/results/batch.json
new file mode 100644
index 0000000..b2d4d76
--- /dev/null
+++ b/docs/results/batch.json
@@ -0,0 +1,60 @@
+{
+ "cloud9": {
+ "allocation_evidence": "no allocator or allocation calls in codec",
+ "bytes_in": 3803125,
+ "bytes_out": 3803125,
+ "codec_allocations": 0,
+ "codec_ns": 11364700,
+ "implementation": "cloud9",
+ "operations": 540000,
+ "read_calls": 54001,
+ "write_calls": 27000
+ },
+ "configuration": {
+ "chunk": 65536,
+ "frames": 27000,
+ "go": "go1.26.5-X:nodwarf5",
+ "platform": "linux/amd64",
+ "repeats": 20,
+ "rounds": 1000,
+ "seed": 4200
+ },
+ "correctness": "all generated base-9P2000 frames agree byte-for-byte",
+ "measurement_notes": [
+ "Cloud9 codec timer excludes pipe I/O; Go timers include in-memory stream parsing and byte comparison. Do not treat these as equivalent throughput benchmarks.",
+ "Go allocations use runtime.MemStats deltas on a single-process run; cloud9 codec has no allocation API.",
+ "Cloud9 I/O counters are actual stdin/stdout syscalls; Go read counters are io.Reader calls, not syscalls.",
+ "These are wire-conformance probes, not differential filesystem-server semantics."
+ ],
+ "references": [
+ {
+ "implementation": "9fans v0.0.7",
+ "operations": 540000,
+ "codec_ns": 116434888,
+ "codec_allocations": 2509437,
+ "allocated_bytes": 486133240,
+ "read_calls": 1080000,
+ "bytes": 76062500
+ },
+ {
+ "implementation": "go9p v1.18.0",
+ "operations": 540000,
+ "codec_ns": 94922755,
+ "codec_allocations": 3130709,
+ "allocated_bytes": 289069480,
+ "read_calls": 1080000,
+ "bytes": 76062500
+ }
+ ],
+ "session": {
+ "bytes_read": 3156839,
+ "bytes_written": 646229,
+ "checks": "walk/open/write/read/stat/EOF/clunk, content oracle, fragmented requests and replies",
+ "read_calls": 454124,
+ "requests": 7003,
+ "rounds": 1000,
+ "seed": 4200,
+ "server": "go9p v1.18.0 StaticFile",
+ "write_calls": 217407
+ }
+}
diff --git a/docs/results/fragmented.json b/docs/results/fragmented.json
new file mode 100644
index 0000000..a49b23a
--- /dev/null
+++ b/docs/results/fragmented.json
@@ -0,0 +1,60 @@
+{
+ "cloud9": {
+ "allocation_evidence": "no allocator or allocation calls in codec",
+ "bytes_in": 369117,
+ "bytes_out": 369117,
+ "codec_allocations": 0,
+ "codec_ns": 1327036,
+ "implementation": "cloud9",
+ "operations": 54000,
+ "read_calls": 369118,
+ "write_calls": 369117
+ },
+ "configuration": {
+ "chunk": 1,
+ "frames": 2700,
+ "go": "go1.26.5-X:nodwarf5",
+ "platform": "linux/amd64",
+ "repeats": 20,
+ "rounds": 100,
+ "seed": 4201
+ },
+ "correctness": "all generated base-9P2000 frames agree byte-for-byte",
+ "measurement_notes": [
+ "Cloud9 codec timer excludes pipe I/O; Go timers include in-memory stream parsing and byte comparison. Do not treat these as equivalent throughput benchmarks.",
+ "Go allocations use runtime.MemStats deltas on a single-process run; cloud9 codec has no allocation API.",
+ "Cloud9 I/O counters are actual stdin/stdout syscalls; Go read counters are io.Reader calls, not syscalls.",
+ "These are wire-conformance probes, not differential filesystem-server semantics."
+ ],
+ "references": [
+ {
+ "implementation": "9fans v0.0.7",
+ "operations": 54000,
+ "codec_ns": 48042988,
+ "codec_allocations": 250519,
+ "allocated_bytes": 48088256,
+ "read_calls": 7382340,
+ "bytes": 7382340
+ },
+ {
+ "implementation": "go9p v1.18.0",
+ "operations": 54000,
+ "codec_ns": 43224680,
+ "codec_allocations": 312801,
+ "allocated_bytes": 28165920,
+ "read_calls": 7382340,
+ "bytes": 7382340
+ }
+ ],
+ "session": {
+ "bytes_read": 282585,
+ "bytes_written": 63853,
+ "checks": "walk/open/write/read/stat/EOF/clunk, content oracle, fragmented requests and replies",
+ "read_calls": 40697,
+ "requests": 703,
+ "rounds": 100,
+ "seed": 4201,
+ "server": "go9p v1.18.0 StaticFile",
+ "write_calls": 21488
+ }
+}
diff --git a/docs/spec.md b/docs/spec.md
new file mode 100644
index 0000000..4537620
--- /dev/null
+++ b/docs/spec.md
@@ -0,0 +1,33 @@
+# Specification baseline
+
+The baseline is Plan 9's 9P2000 manual, read before defining the package boundary.
+Reference implementations are comparison partners, not the specification.
+
+| Requirement | Primary source | Implementation |
+|---|---|---|
+| Little-endian framing, counted strings, tags, qids | [intro(9P)](https://9fans.github.io/plan9port/man/man9/intro.html) | `wire.zig` |
+| First-message negotiation, NOTAG, msize, session reset | [version(9P)](https://9fans.github.io/plan9port/man/man9/version.html) | `Client`, `Server.negotiate`; backend releases fids |
+| Authentication and attach | [attach(9P)](https://9fans.github.io/plan9port/man/man9/attach.html) | All wire/client messages; backend authentication |
+| Tag reservation and cancellation completion | [flush(9P)](https://9fans.github.io/plan9port/man/man9/flush.html) | Client reservations; server Rflush retires oldtag |
+| Walk bounds, cloning and partial walks | [walk(9P)](https://9fans.github.io/plan9port/man/man9/walk.html) | Codec and count checks; backend commits only full walks |
+| Opening, creation, permissions and iounit | [open(9P)](https://9fans.github.io/plan9port/man/man9/open.html) | All wire/client messages; backend filesystem behavior |
+| Read/write counts and directory records | [read(9P)](https://9fans.github.io/plan9port/man/man9/read.html) | Reply count checks; backend directory offsets and records |
+| Fid release, including failed removal | [clunk(9P)](https://9fans.github.io/plan9port/man/man9/clunk.html), [remove(9P)](https://9fans.github.io/plan9port/man/man9/remove.html) | Client messages; backend releases handles |
+| Double stat lengths and atomic metadata changes | [stat(9P)](https://9fans.github.io/plan9port/man/man9/stat.html) | `Stat`, wire length checks; backend atomic updates |
+| Error replies | [error(9P)](https://9fans.github.io/plan9port/man/man9/error.html) | Counted errors, frame-size truncation, no Rerror for Tflush |
+
+The local 24-byte session floor and fixed request capacities are resource choices,
+not additional wire requirements. TCP, Unix, and QUIC adapters carry unchanged
+9P2000 frames. QUIC ALPN and authentication policy are not defined by 9P2000.
+
+[Tiger Style](https://github.com/tigerbeetle/tigerbeetle/blob/main/docs/TIGER_STYLE.md)
+informs bounded memory, explicit ownership, assertions for internal invariants,
+recoverable errors for external input, and deterministic tests. Zig standard
+library conventions take precedence for names: functions use camelCase, values
+use snake_case, and types use PascalCase. The code is formatted with `zig fmt`.
+
+The wire codec and initial client were extracted from Pardes and checked against
+these rules. Cloud9 adds the reusable server connection, missing client
+operations, cancellation bookkeeping, bounds checks, shared transports, and
+independent conformance harness. Pardes regression tests remain with its adapter;
+wire regression tests moved with the codec.
diff --git a/docs/validation.md b/docs/validation.md
new file mode 100644
index 0000000..acc35c3
--- /dev/null
+++ b/docs/validation.md
@@ -0,0 +1,73 @@
+# Validation and reproduction
+
+Run the commands in README.md. The Go modules in `test/differential/go.mod` and
+`go.sum` pin the reference implementations and their test-only dependencies.
+They are not library dependencies. The valid-traffic corpus exercises all 27
+message types, zero and nonzero payloads, 0–16 walk elements, UTF-8 names, stat
+records, wstat sentinel values, and varying integer fields. Use different seeds
+and `--chunk 1` for byte-at-a-time stream probes. Generated cases, their seed and
+configuration, and reports are saved before a mismatch can terminate the run.
+
+The live session probe connects cloud9's client to go9p's StaticFile backend over
+local pipes. A seeded in-memory content oracle checks writes at varying offsets,
+reads, stat lengths, EOF, open/walk/clunk, and fid reuse. It does not claim coverage
+of every reference-server filesystem policy. Pardes's separate Python 9P client
+exercises the integrated cloud9 server over Unix, TCP, and optional QUIC.
+
+The standalone `zig build fuzz -- <seed> <iterations>` runner mutates frames only
+for cloud9. It checks decode/re-encode identity and server pre-negotiation handling.
+The native `zig build test --fuzz=10000` route is also present through std.testing,
+but the installed Zig 0.16.0 build currently fails compiling its own fuzz test
+runner due to incompatible StackTrace types. The standalone runner avoids that
+toolchain failure; it is deterministic mutation testing, not coverage-guided.
+
+# Interpreting measurements
+
+Reports retain raw counts. Cloud9 codec timings exclude pipe I/O. Reference Go
+measurements include in-memory io.Reader parsing and byte comparison. These are
+not directly comparable end-to-end throughput numbers. Cloud9 syscall counts
+cover one traversal of the corpus, while Go reader calls cover all repetitions.
+Counts must be normalized before comparison, and an io.Reader call is not a
+syscall. Byte totals distinguish corpus transport from repeated codec work.
+
+Go allocations are runtime.MemStats deltas across the measurement. Cloud9's zero
+codec allocation count is structural evidence: no allocator or allocation calls
+exist in that path. It is not a claim about process startup, std.Io, OpenSSL,
+filesystem backends, or the test harness. There are no portable performance gates
+or claims about kernel/disk performance in these reports.
+
+If this machine's /tmp quota is exhausted, set TMPDIR to an owned directory on a
+filesystem with space. Unix socket tests need a short absolute path because of
+sockaddr_un's path limit. No unrelated temporary files need to be removed.
+
+# Recorded run — 2026-09-14
+
+Linux x86_64, Zig 0.16.0, OpenSSL 3.6.3, Go 1.26.5. These results describe this
+checkout and machine, not a production-readiness or complete-conformance claim.
+
+| Check | Result |
+|---|---|
+| Cloud9 Debug tests | 24 protocol/session + 3 transport + 6 QUIC passed |
+| Cloud9 ReleaseSafe tests | The same 33 tests passed |
+| Deterministic mutation probes | 1,000,000 iterations, seed 4200; 442,898 accepted, 557,102 rejected |
+| Valid wire differential, seed 4200 | 27,000 frames; 540,000 codec operations per implementation |
+| Valid wire differential, seed 4201, chunk 1 | 2,700 frames; 54,000 codec operations per implementation |
+| Live go9p server, seed 4200 | 7,003 requests passed |
+| Live go9p server, seed 4201 | 703 requests passed |
+| Pardes protocol/filesystem adapter tests | 31 passed |
+| Pardes native transport/client tests (`9p-io-test`) | 8 passed both with and without QUIC |
+| Pardes terminal build with QUIC | Built with Zig grammar enabled |
+| Pardes dedicated Python Unix/TCP/QUIC integration | Passed IPv4/IPv6 and headless/TTY variants |
+| ESP32-P4 GPIO image | Built successfully; no flash performed |
+
+Raw measurement reports: [batch](results/batch.json),
+[byte-at-a-time](results/fragmented.json).
+
+The full Pardes editor integration script reached its syntax-style assertion and
+failed because `fn` and `if` were not bold in the current theme. Syntax coloring
+was enabled on the second run. This was left unchanged; the dedicated network
+integration function was run independently and passed. A complete core-test run
+passed earlier, but later full reruns encountered hard-coded /tmp shell-fixture
+creation failures from the machine's quota. The final protocol and transport
+checks were run separately. The `--fuzz` compiler-runner limitation is described
+above. Firmware hardware behavior and non-Linux native transports were not run.
diff --git a/src/Server.zig b/src/Server.zig
new file mode 100644
index 0000000..a652f99
--- /dev/null
+++ b/src/Server.zig
@@ -0,0 +1,194 @@
+//! A bounded, caller-driven 9P2000 server connection.
+//! The backend owns fids, authentication, permissions, and filesystem operations.
+//! receive() borrows one input frame until release(). Pending operations may outlive
+//! that frame only if the backend copies their strings/data. reply() copies output.
+const std = @import("std");
+const assert = std.debug.assert;
+const wire = @import("wire.zig");
+const Server = @This();
+
+in: []u8,
+out: []u8,
+in_len: usize = 0,
+frame: u32 = 0,
+out_len: usize = 0,
+out_off: usize = 0,
+msize: u32 = 0,
+dead: bool = false,
+pending: [65]Pending = @splat(.{}),
+versioning: bool = false,
+
+const Pending = struct {
+ request: ?wire.Type = null,
+ tag: u16 = 0,
+ count: u32 = 0,
+ oldtag: u16 = wire.notag,
+};
+
+pub const Error = wire.Error || error{ Protocol, NoTags, UnknownTag, WrongReply, TooLarge };
+pub const Options = struct { in: []u8, out: []u8 };
+/// Local resource floor, not a minimum imposed by the 9P specification.
+pub const msize_min: u32 = 24;
+
+pub fn init(options: Options) Server {
+ assert(options.in.len >= msize_min);
+ assert(options.out.len >= msize_min);
+ return .{ .in = options.in, .out = options.out };
+}
+
+pub fn push(s: *Server, bytes: []const u8) usize {
+ if (s.dead) return 0;
+ const n = @min(bytes.len, s.in.len - s.in_len);
+ @memcpy(s.in[s.in_len..][0..n], bytes[0..n]);
+ s.in_len += n;
+ return n;
+}
+
+pub fn output(s: *const Server) []const u8 {
+ return s.out[s.out_off..s.out_len];
+}
+
+pub fn wrote(s: *Server, n: usize) void {
+ assert(n <= s.output().len);
+ s.out_off += n;
+ if (s.out_off == s.out_len) {
+ s.out_off = 0;
+ s.out_len = 0;
+ }
+}
+
+pub fn hasRoom(s: *Server) bool {
+ if (s.out_off != 0) {
+ const n = s.output().len;
+ std.mem.copyForwards(u8, s.out[0..n], s.output());
+ s.out_off = 0;
+ s.out_len = n;
+ }
+ return s.out.len - s.out_len >= @max(s.msize, msize_min);
+}
+
+/// A null result means more input or output drainage is needed.
+/// Any receive error is terminal: the stream cannot be safely resynchronized.
+pub fn receive(s: *Server) Error!?wire.Decoded {
+ assert(s.frame == 0);
+ if (s.dead) return null;
+ errdefer s.dead = true;
+ const len = wire.frameLen(s.in[0..s.in_len]) orelse return null;
+ if (len < wire.header_len or len > s.in.len) return error.Protocol;
+ if (len > s.in_len) return null;
+ const got = try wire.decode(s.in[0..len]);
+ const kind = got.msg.msgType();
+ if (!wire.isT(kind)) return error.Protocol;
+ if (!s.hasRoom()) return null;
+ if (kind == .tversion) {
+ if (got.tag != wire.notag) return error.Protocol;
+ // Finish any partially transmitted response before starting a new session.
+ if (s.output().len != 0) return null;
+ s.pending = @splat(.{});
+ s.versioning = true;
+ s.msize = 0;
+ } else {
+ if (s.msize == 0 or s.versioning) return error.Protocol;
+ if (len > s.msize or got.tag == wire.notag) return error.Protocol;
+ if (s.find(got.tag) != null) return error.Protocol;
+ const slot = s.free(kind) orelse return error.NoTags;
+ slot.* = .{ .request = kind, .tag = got.tag, .count = switch (got.msg) {
+ .tread => |m| m.count,
+ .twrite => |m| @intCast(m.data.len),
+ .twalk => |m| m.nwname,
+ else => 0,
+ }, .oldtag = if (got.msg == .tflush) got.msg.tflush.oldtag else wire.notag };
+ }
+ s.frame = len;
+ return got;
+}
+
+pub fn release(s: *Server) void {
+ assert(s.frame != 0 and s.frame <= s.in_len);
+ const n = s.in_len - s.frame;
+ std.mem.copyForwards(u8, s.in[0..n], s.in[s.frame..s.in_len]);
+ s.in_len = n;
+ s.frame = 0;
+}
+
+/// Call only for a received Tversion, after aborting backend work and releasing fids.
+/// Suffixes may fall back to base 9P2000; arbitrary strings beginning with 9P may not.
+pub fn negotiate(s: *Server, want: u32, version: []const u8) Error!void {
+ assert(s.versioning);
+ const size: u32 = @intCast(@min(want, s.in.len, s.out.len, std.math.maxInt(u32)));
+ if (size < msize_min) return error.TooLarge;
+ const base = std.mem.sliceTo(version, '.');
+ const known = std.mem.eql(u8, base, "9P2000");
+ try s.append(wire.notag, .{ .rversion = .{
+ .msize = size,
+ .version = if (known) "9P2000" else "unknown",
+ } }, size);
+ s.msize = if (known) size else 0;
+ s.versioning = false;
+}
+
+/// An Rflush is a backend promise: no further response for oldtag will be sent.
+/// A canceled backend operation must be retired before its tag can be reused.
+pub fn reply(s: *Server, tag: u16, msg: wire.Msg) Error!void {
+ if (s.dead) return error.Protocol;
+ const slot = s.find(tag) orelse return error.UnknownTag;
+ const request = slot.request.?;
+ const kind = msg.msgType();
+ if (request == .tflush and kind != .rflush) return error.WrongReply;
+ if (kind != .rerror and @intFromEnum(kind) != @intFromEnum(request) + 1)
+ return error.WrongReply;
+ switch (msg) {
+ .rread => |m| if (m.data.len > slot.count) return error.WrongReply,
+ .rwrite => |m| if (m.count > slot.count) return error.WrongReply,
+ .rwalk => |m| {
+ if (m.nwqid > slot.count) return error.WrongReply;
+ if (m.nwqid == 0 and slot.count != 0) return error.WrongReply;
+ },
+ else => {},
+ }
+ var response = msg;
+ if (response == .rerror) {
+ const cap = @min(s.msize - wire.header_len - 2, std.math.maxInt(u16));
+ response.rerror.ename = response.rerror.ename[0..@min(response.rerror.ename.len, cap)];
+ }
+ try s.append(tag, response, s.msize);
+ const oldtag = slot.oldtag;
+ slot.* = .{};
+ if (kind == .rflush) {
+ if (s.find(oldtag)) |old| old.* = .{};
+ }
+}
+
+fn append(s: *Server, tag: u16, msg: wire.Msg, limit: u32) Error!void {
+ if (try wire.encodedLen(msg) > limit) return error.TooLarge;
+ _ = s.hasRoom();
+ const bytes = try wire.encode(msg, tag, s.out[s.out_len..]);
+ s.out_len += bytes.len;
+}
+
+fn find(s: *Server, tag: u16) ?*Pending {
+ for (&s.pending) |*slot| {
+ if (slot.request != null and slot.tag == tag) return slot;
+ }
+ return null;
+}
+
+fn free(s: *Server, request: wire.Type) ?*Pending {
+ for (&s.pending, 0..) |*slot, i| {
+ // A full ordinary request window must still admit cancellation.
+ if (i == s.pending.len - 1 and request != .tflush) continue;
+ if (slot.request == null) return slot;
+ }
+ return null;
+}
+
+pub fn hangup(s: *Server) void {
+ s.dead = true;
+ s.pending = @splat(.{});
+ s.in_len = 0;
+ s.frame = 0;
+ s.out_len = 0;
+ s.out_off = 0;
+ s.msize = 0;
+ s.versioning = false;
+}
diff --git a/src/client.zig b/src/client.zig
new file mode 100644
index 0000000..5b428dd
--- /dev/null
+++ b/src/client.zig
@@ -0,0 +1,392 @@
+//! Caller-driven 9P2000 client. No allocator, sockets, threads, or filesystem policy.
+//! Result slices remain valid until the next take() or hangup().
+const std = @import("std");
+const assert = std.debug.assert;
+const wire = @import("wire.zig");
+const Msg = wire.Msg;
+const Stat = wire.Stat;
+const Qid = wire.Qid;
+const notag = wire.notag;
+const nofid = wire.nofid;
+const max_welem = wire.max_welem;
+const header_len = wire.header_len;
+const isT = wire.isT;
+const encode = wire.encode;
+const decode = wire.decode;
+const frameLen = wire.frameLen;
+const Decoded = wire.Decoded;
+const totalLen = wire.encodedLen;
+pub const msize_min: u32 = 24;
+pub const max_tags: usize = 16;
+
+const twrite_header: usize = header_len + 4 + 8 + 4;
+
+const rread_header: usize = header_len + 4;
+
+comptime {
+ assert(twrite_header == 23);
+ assert(rread_header == 11);
+ assert(max_tags <= notag);
+ assert(msize_min > rread_header);
+ assert(msize_min > twrite_header);
+}
+
+pub const ClientError = error{
+ NoTags,
+ NoSpace,
+ TooLarge,
+ Handshake,
+ Dead,
+ BadRequest,
+};
+
+pub const Client = struct {
+ in: []u8,
+ out: []u8,
+
+ in_len: usize = 0,
+ frame: u32 = 0,
+ out_len: usize = 0,
+ out_off: usize = 0,
+
+ msize: u32 = 0,
+ asked: u32 = 0,
+ versioning: bool = false,
+ dead: bool = false,
+
+ tags: [max_tags + 1]Slot = @splat(.{}),
+
+ const Slot = struct {
+ op: ?Op = null,
+ count: u32 = 0,
+ oldtag: u16 = notag,
+ completed: bool = false,
+ };
+
+ pub const Op = enum { version, auth, attach, flush, walk, open, create, read, write, clunk, remove, stat, wstat };
+
+ pub const Request = union(Op) {
+ version: struct { msize: u32 = 0 },
+ auth: struct { afid: u32, uname: []const u8, aname: []const u8 = "" },
+ attach: struct { fid: u32, afid: u32 = nofid, uname: []const u8, aname: []const u8 = "" },
+ flush: struct { oldtag: u16 },
+ walk: struct { fid: u32, newfid: u32, names: []const []const u8 },
+ open: struct { fid: u32, mode: u8 },
+ create: struct { fid: u32, name: []const u8, perm: u32, mode: u8 },
+ read: struct { fid: u32, offset: u64, count: u32 },
+ write: struct { fid: u32, offset: u64, data: []const u8 },
+ clunk: struct { fid: u32 },
+ remove: struct { fid: u32 },
+ stat: struct { fid: u32 },
+ wstat: struct { fid: u32, stat: Stat },
+ };
+
+ pub const Result = union(enum) {
+ fail: []const u8,
+ version: struct { msize: u32, version: []const u8 },
+ auth: Qid,
+ flush: void,
+ create: struct { qid: Qid, iounit: u32 },
+ remove: void,
+ wstat: void,
+ attach: Qid,
+ walk: struct { nwqid: u16, wqid: [max_welem]Qid },
+ open: struct { qid: Qid, iounit: u32 },
+ read: []const u8,
+ write: u32,
+ clunk: void,
+ stat: Stat,
+ };
+
+ pub const Done = struct {
+ tag: u16,
+ op: Op,
+ result: Result,
+ };
+
+ pub const Options = struct {
+ in: []u8,
+ out: []u8,
+ };
+
+ pub fn init(opts: Options) Client {
+ assert(opts.in.len >= msize_min);
+ assert(opts.out.len >= msize_min);
+ return .{ .in = opts.in, .out = opts.out };
+ }
+
+ pub fn hangup(c: *Client) void {
+ c.dead = true;
+ c.tags = @splat(.{});
+ c.versioning = false;
+ c.msize = 0;
+ c.asked = 0;
+ c.in_len = 0;
+ c.frame = 0;
+ c.out_len = 0;
+ c.out_off = 0;
+ }
+
+ pub fn push(c: *Client, bytes: []const u8) usize {
+ if (c.dead) return 0;
+ const n = @min(bytes.len, c.in.len - c.in_len);
+ @memcpy(c.in[c.in_len..][0..n], bytes[0..n]);
+ c.in_len += n;
+ return n;
+ }
+
+ pub fn output(c: *const Client) []const u8 {
+ return c.out[c.out_off..c.out_len];
+ }
+
+ pub fn wrote(c: *Client, n: usize) void {
+ assert(n <= c.out_len - c.out_off);
+ c.out_off += n;
+ if (c.out_off == c.out_len) {
+ c.out_off = 0;
+ c.out_len = 0;
+ }
+ }
+
+ fn compact(c: *Client) void {
+ assert(c.out_off <= c.out_len);
+ const n = c.out_len - c.out_off;
+ std.mem.copyForwards(u8, c.out[0..n], c.out[c.out_off..c.out_len]);
+ c.out_off = 0;
+ c.out_len = n;
+ }
+
+ fn dropFrame(c: *Client) void {
+ assert(c.frame != 0);
+ assert(c.frame <= c.in_len);
+ const n = c.frame;
+ std.mem.copyForwards(u8, c.in[0 .. c.in_len - n], c.in[n..c.in_len]);
+ c.in_len -= n;
+ c.frame = 0;
+ }
+
+ pub fn maxRead(c: *const Client) u32 {
+ if (c.msize == 0) return 0;
+ return c.msize - @as(u32, @intCast(rread_header));
+ }
+
+ pub fn maxWrite(c: *const Client) u32 {
+ if (c.msize == 0) return 0;
+ return c.msize - @as(u32, @intCast(twrite_header));
+ }
+
+ pub fn pending(c: *const Client) usize {
+ var n: usize = @intFromBool(c.versioning);
+ for (c.tags) |t| n += @intFromBool(t.op != null);
+ return n;
+ }
+
+ pub fn submit(c: *Client, req: Request) ClientError!u16 {
+ if (c.dead) return error.Dead;
+ if (req == .version) return c.beginVersion(req.version.msize);
+ if (c.msize == 0 or c.versioning) return error.Handshake;
+
+ const msg: Msg = switch (req) {
+ .version => unreachable, // handled above
+ .auth => |m| blk: {
+ if (m.afid == nofid) return error.BadRequest;
+ break :blk .{ .tauth = .{ .afid = m.afid, .uname = m.uname, .aname = m.aname } };
+ },
+ .flush => |m| blk: {
+ if (m.oldtag == notag) return error.BadRequest;
+ break :blk .{ .tflush = .{ .oldtag = m.oldtag } };
+ },
+ .create => |m| blk: {
+ if (m.fid == nofid or m.name.len == 0) return error.BadRequest;
+ if (std.mem.indexOfAny(u8, m.name, "/\x00") != null) return error.BadRequest;
+ if (std.mem.eql(u8, m.name, ".") or std.mem.eql(u8, m.name, "..")) return error.BadRequest;
+ break :blk .{ .tcreate = .{ .fid = m.fid, .name = m.name, .perm = m.perm, .mode = m.mode } };
+ },
+ .remove => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .tremove = .{ .fid = m.fid } };
+ },
+ .wstat => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .twstat = .{ .fid = m.fid, .stat = m.stat } };
+ },
+ .attach => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .tattach = .{
+ .fid = m.fid,
+ .afid = m.afid,
+ .uname = m.uname,
+ .aname = m.aname,
+ } };
+ },
+ .walk => |m| blk: {
+ if (m.fid == nofid or m.newfid == nofid) return error.BadRequest;
+ if (m.names.len > max_welem) return error.BadRequest;
+ var w: [max_welem][]const u8 = @splat("");
+ for (m.names, 0..) |n, i| {
+ if (n.len == 0) return error.BadRequest;
+ if (std.mem.indexOfAny(u8, n, "/\x00") != null) return error.BadRequest;
+ w[i] = n;
+ }
+ break :blk .{ .twalk = .{
+ .fid = m.fid,
+ .newfid = m.newfid,
+ .nwname = @intCast(m.names.len),
+ .wname = w,
+ } };
+ },
+ .open => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .topen = .{ .fid = m.fid, .mode = m.mode } };
+ },
+ .read => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ if (m.count > c.maxRead()) return error.TooLarge;
+ break :blk .{ .tread = .{ .fid = m.fid, .offset = m.offset, .count = m.count } };
+ },
+ .write => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .twrite = .{ .fid = m.fid, .offset = m.offset, .data = m.data } };
+ },
+ .clunk => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .tclunk = .{ .fid = m.fid } };
+ },
+ .stat => |m| blk: {
+ if (m.fid == nofid) return error.BadRequest;
+ break :blk .{ .tstat = .{ .fid = m.fid } };
+ },
+ };
+ const need = totalLen(msg) catch return error.TooLarge;
+ if (need > c.msize) return error.TooLarge;
+
+ const op = std.meta.activeTag(req);
+ const tag = c.claim(op, if (req == .flush) req.flush.oldtag else notag) orelse return error.NoTags;
+ errdefer c.tags[tag] = .{};
+ try c.emit(tag, msg);
+ switch (req) {
+ .read => |m| c.tags[tag].count = m.count,
+ .write => |m| c.tags[tag].count = @intCast(m.data.len),
+ .walk => |m| c.tags[tag].count = @intCast(m.names.len),
+ .flush => |m| {
+ c.tags[tag].oldtag = m.oldtag;
+ },
+ else => {},
+ }
+ return tag;
+ }
+
+ fn beginVersion(c: *Client, want: u32) ClientError!u16 {
+ if (c.pending() != 0) return error.Handshake;
+ const cap: u32 = @intCast(@min(c.in.len, c.out.len, std.math.maxInt(u32)));
+ const m = @min(if (want == 0) cap else want, cap);
+ if (m < msize_min) return error.BadRequest;
+ try c.emit(notag, .{ .tversion = .{ .msize = m, .version = "9P2000" } });
+ c.msize = 0;
+ c.asked = m;
+ c.versioning = true;
+ return notag;
+ }
+
+ fn reserved(c: *const Client, tag: u16) bool {
+ for (c.tags) |slot| {
+ if (slot.op == .flush and !slot.completed and slot.oldtag == tag) return true;
+ }
+ return false;
+ }
+
+ fn claim(c: *Client, op: Op, avoid: u16) ?u16 {
+ for (&c.tags, 0..) |*t, i| {
+ // Keep one tag available to cancel a fully occupied request window.
+ if (i >= max_tags and op != .flush) continue;
+ if (t.op != null or i == avoid or c.reserved(@intCast(i))) continue;
+ t.* = .{ .op = op };
+ return @intCast(i);
+ }
+ return null;
+ }
+
+ fn emit(c: *Client, tag: u16, msg: Msg) ClientError!void {
+ if (c.out_off != 0) c.compact();
+ const bytes = encode(msg, tag, c.out[c.out_len..]) catch return error.NoSpace;
+ c.out_len += bytes.len;
+ }
+
+ pub fn take(c: *Client) ?Done {
+ if (c.frame != 0) c.dropFrame();
+ if (c.dead) return null;
+ const len = frameLen(c.in[0..c.in_len]) orelse return null;
+ if (len < header_len or len > c.in.len) return c.die();
+ if (c.msize != 0 and len > c.msize) return c.die();
+ if (len > c.in_len) return null;
+ c.frame = len;
+ const got = decode(c.in[0..len]) catch return c.die();
+ return c.consume(got);
+ }
+
+ fn die(c: *Client) ?Done {
+ c.dead = true;
+ return null;
+ }
+
+ fn consume(c: *Client, got: Decoded) ?Done {
+ if (isT(got.msg.msgType())) return c.die();
+ if (got.msg == .rversion) return c.version(got);
+ if (c.versioning or c.msize == 0) return c.die();
+ if (got.tag >= c.tags.len) return c.die();
+ const slot = &c.tags[got.tag];
+ const op = slot.op orelse return c.die();
+ if (slot.completed) return c.die();
+ const result: Result = switch (got.msg) {
+ .rerror => |m| if (op == .flush) return c.die() else .{ .fail = m.ename },
+ .rauth => |m| if (op != .auth) return c.die() else .{ .auth = m.aqid },
+ .rcreate => |m| if (op != .create) return c.die() else .{ .create = .{ .qid = m.qid, .iounit = m.iounit } },
+ .rremove => if (op != .remove) return c.die() else .remove,
+ .rwstat => if (op != .wstat) return c.die() else .wstat,
+ .rflush => blk: {
+ if (op != .flush) return c.die();
+ if (slot.oldtag < c.tags.len) c.tags[slot.oldtag] = .{};
+ break :blk .flush;
+ },
+ .rattach => |m| if (op != .attach) return c.die() else .{ .attach = m.qid },
+ .rwalk => |m| if (op != .walk or m.nwqid > slot.count or (m.nwqid == 0 and slot.count != 0)) return c.die() else .{
+ .walk = .{ .nwqid = m.nwqid, .wqid = m.wqid },
+ },
+ .ropen => |m| if (op != .open) return c.die() else .{
+ .open = .{ .qid = m.qid, .iounit = m.iounit },
+ },
+ .rread => |m| blk: {
+ if (op != .read) return c.die();
+ if (m.data.len > slot.count) return c.die();
+ break :blk .{ .read = m.data };
+ },
+ .rwrite => |m| if (op != .write or m.count > slot.count) return c.die() else .{ .write = m.count },
+ .rclunk => if (op != .clunk) return c.die() else .clunk,
+ .rstat => |m| if (op != .stat) return c.die() else .{ .stat = m.stat },
+ else => return c.die(),
+ };
+ slot.completed = true;
+ // Multiple flushes can reserve the same tag; flushing a flush can release
+ // a reservation on a completed original operation.
+ for (&c.tags, 0..) |*candidate, tag| {
+ if (candidate.completed and !c.reserved(@intCast(tag))) candidate.* = .{};
+ }
+ return .{ .tag = got.tag, .op = op, .result = result };
+ }
+
+ fn version(c: *Client, got: Decoded) ?Done {
+ if (!c.versioning) return c.die();
+ if (got.tag != notag) return c.die();
+ const m = got.msg.rversion;
+ if (m.msize > c.asked or m.msize < msize_min) return c.die();
+ c.versioning = false;
+ if (std.mem.eql(u8, m.version, "9P2000")) {
+ c.msize = m.msize;
+ } else if (!std.mem.eql(u8, m.version, "unknown")) {
+ return c.die();
+ }
+ return .{ .tag = notag, .op = .version, .result = .{
+ .version = .{ .msize = m.msize, .version = m.version },
+ } };
+ }
+};
diff --git a/src/quic.zig b/src/quic.zig
new file mode 100644
index 0000000..c78a79d
--- /dev/null
+++ b/src/quic.zig
@@ -0,0 +1,302 @@
+const std = @import("std");
+const libc = std.c;
+
+/// QUIC transport with caller-selected ALPN. OpenSSL bindings are injected so
+/// applications control linkage. Certificates are ephemeral and peers unauthenticated.
+pub fn Quic(comptime ssl: type, comptime protocol: []const u8) type {
+ if (protocol.len == 0 or protocol.len > 255) @compileError("invalid ALPN length");
+ return struct {
+ comptime {
+ if (ssl.OPENSSL_VERSION_NUMBER < 0x30600000)
+ @compileError("9P over QUIC requires OpenSSL 3.6 or newer");
+ }
+
+ pub const alpn = protocol;
+ pub const Error = error{ Tls, Socket, SocketFlags, SocketOption, Bind, Address, Closed, InvalidWrite };
+
+ pub const Listener = struct {
+ fd: c_int,
+ handle: *ssl.SSL,
+ address: std.Io.net.IpAddress,
+
+ pub fn init(address: std.Io.net.IpAddress) Error!Listener {
+ ssl.ERR_clear_error();
+ const ctx = ssl.SSL_CTX_new(ssl.OSSL_QUIC_server_method()) orelse return error.Tls;
+ defer ssl.SSL_CTX_free(ctx);
+ const key = ssl.EVP_PKEY_Q_keygen(null, null, "EC", @as([*:0]const u8, "prime256v1")) orelse return error.Tls;
+ defer ssl.EVP_PKEY_free(key);
+ const cert = ssl.X509_new() orelse return error.Tls;
+ defer ssl.X509_free(cert);
+ if (ssl.X509_set_version(cert, 2) != 1 or
+ ssl.ASN1_INTEGER_set(ssl.X509_get_serialNumber(cert), 1) != 1 or
+ ssl.X509_gmtime_adj(ssl.X509_getm_notBefore(cert), -60) == null or
+ ssl.X509_gmtime_adj(ssl.X509_getm_notAfter(cert), 365 * 24 * 60 * 60) == null or
+ ssl.X509_set_pubkey(cert, key) != 1) return error.Tls;
+ const name = ssl.X509_get_subject_name(cert) orelse return error.Tls;
+ if (ssl.X509_NAME_add_entry_by_txt(name, "CN", ssl.MBSTRING_ASC, protocol.ptr, @intCast(protocol.len), -1, 0) != 1 or
+ ssl.X509_set_issuer_name(cert, name) != 1 or
+ ssl.X509_sign(cert, key, ssl.EVP_sha256()) <= 0 or
+ ssl.SSL_CTX_use_certificate(ctx, cert) != 1 or
+ ssl.SSL_CTX_use_PrivateKey(ctx, key) != 1) return error.Tls;
+ ssl.SSL_CTX_set_verify(ctx, ssl.SSL_VERIFY_NONE, null);
+ ssl.SSL_CTX_set_alpn_select_cb(ctx, selectAlpn, null);
+ var addr: libc.sockaddr.storage = undefined;
+ const addr_len = sockaddr(address, &addr);
+ const fd = try udp(addr.family);
+ errdefer _ = libc.close(fd);
+ if (libc.bind(fd, @ptrCast(&addr), addr_len) != 0) return error.Bind;
+ var actual_len: libc.socklen_t = @sizeOf(@TypeOf(addr));
+ if (libc.getsockname(fd, @ptrCast(&addr), &actual_len) != 0) return error.Address;
+ var actual = address;
+ actual.setPort(switch (address) {
+ .ip4 => std.mem.bigToNative(u16, @as(*const libc.sockaddr.in, @ptrCast(&addr)).port),
+ .ip6 => std.mem.bigToNative(u16, @as(*const libc.sockaddr.in6, @ptrCast(&addr)).port),
+ });
+ const handle = ssl.SSL_new_listener(ctx, 0) orelse return error.Tls;
+ errdefer ssl.SSL_free(handle);
+ if (ssl.SSL_set_fd(handle, fd) != 1 or ssl.SSL_set_blocking_mode(handle, 0) != 1 or
+ ssl.SSL_listen(handle) != 1) return error.Tls;
+ return .{ .fd = fd, .handle = handle, .address = actual };
+ }
+
+ pub fn accept(l: *Listener) Error!?Connection {
+ ssl.ERR_clear_error();
+ const handle = ssl.SSL_accept_connection(l.handle, ssl.SSL_ACCEPT_CONNECTION_NO_BLOCK) orelse {
+ if (ssl.ERR_peek_error() != 0) return error.Tls;
+ return null;
+ };
+ errdefer ssl.SSL_free(handle);
+ if (ssl.SSL_set_default_stream_mode(handle, ssl.SSL_DEFAULT_STREAM_MODE_NONE) != 1 or
+ ssl.SSL_set_blocking_mode(handle, 0) != 1) return error.Tls;
+ return .{ .handle = handle };
+ }
+
+ pub fn events(l: *Listener) Error!void {
+ ssl.ERR_clear_error();
+ if (ssl.SSL_handle_events(l.handle) != 1) return error.Tls;
+ }
+
+ pub fn poll(l: *const Listener) libc.pollfd {
+ return pollFd(l.handle, l.fd);
+ }
+
+ pub fn nextDue(l: *const Listener) ?i32 {
+ return due(l.handle);
+ }
+
+ // Accepted connections must be released before the shared UDP socket.
+ pub fn deinit(l: *Listener) void {
+ ssl.SSL_free(l.handle);
+ _ = libc.close(l.fd);
+ l.* = undefined;
+ }
+ };
+
+ pub const Connection = struct {
+ handle: *ssl.SSL,
+ stream: ?*ssl.SSL = null,
+ fd: c_int = -1,
+ pending_write_len: usize = 0,
+
+ pub fn dial(address: std.Io.net.IpAddress) Error!Connection {
+ ssl.ERR_clear_error();
+ const ctx = ssl.SSL_CTX_new(ssl.OSSL_QUIC_client_method()) orelse return error.Tls;
+ defer ssl.SSL_CTX_free(ctx);
+ ssl.SSL_CTX_set_verify(ctx, ssl.SSL_VERIFY_NONE, null);
+ const fd = try udp(if (address == .ip4) libc.AF.INET else libc.AF.INET6);
+ errdefer _ = libc.close(fd);
+ const handle = ssl.SSL_new(ctx) orelse return error.Tls;
+ errdefer ssl.SSL_free(handle);
+ if (ssl.SSL_set_fd(handle, fd) != 1 or ssl.SSL_set_blocking_mode(handle, 0) != 1 or
+ ssl.SSL_set_default_stream_mode(handle, ssl.SSL_DEFAULT_STREAM_MODE_NONE) != 1) return error.Tls;
+ const protocols = [_]u8{alpn.len} ++ protocol[0..protocol.len].*;
+ if (ssl.SSL_set_alpn_protos(handle, &protocols, protocols.len) != 0) return error.Tls;
+ const peer = ssl.BIO_ADDR_new() orelse return error.Tls;
+ defer ssl.BIO_ADDR_free(peer);
+ const made = switch (address) {
+ .ip4 => |ip| ssl.BIO_ADDR_rawmake(peer, libc.AF.INET, &ip.bytes, ip.bytes.len, std.mem.nativeToBig(u16, ip.port)),
+ .ip6 => |ip| ssl.BIO_ADDR_rawmake(peer, libc.AF.INET6, &ip.bytes, ip.bytes.len, std.mem.nativeToBig(u16, ip.port)),
+ };
+ if (made != 1 or ssl.SSL_set1_initial_peer_addr(handle, peer) != 1) return error.Tls;
+ return .{ .handle = handle, .fd = fd };
+ }
+
+ pub fn handshake(c: *Connection) Error!bool {
+ ssl.ERR_clear_error();
+ var close_info: ssl.SSL_CONN_CLOSE_INFO = undefined;
+ if (ssl.SSL_get_conn_close_info(c.handle, &close_info, @sizeOf(@TypeOf(close_info))) == 1)
+ return error.Closed;
+ if (ssl.SSL_is_init_finished(c.handle) == 1) return true;
+ const rc = if (c.fd >= 0) ssl.SSL_connect(c.handle) else ssl.SSL_accept(c.handle);
+ if (rc == 1) return true;
+ try retry(c.handle, rc);
+ return false;
+ }
+
+ fn ready(c: *Connection) Error!bool {
+ if (!try c.handshake()) return false;
+ if (c.stream != null) return true;
+ ssl.ERR_clear_error();
+ const stream = if (c.fd >= 0)
+ ssl.SSL_new_stream(c.handle, ssl.SSL_STREAM_FLAG_NO_BLOCK)
+ else
+ ssl.SSL_accept_stream(c.handle, ssl.SSL_ACCEPT_STREAM_NO_BLOCK);
+ if (stream == null) {
+ if (ssl.ERR_peek_error() != 0) return error.Tls;
+ return false;
+ }
+ errdefer ssl.SSL_free(stream);
+ if (ssl.SSL_set_blocking_mode(stream, 0) != 1 or
+ ssl.SSL_get_stream_id(stream) != 0 or
+ ssl.SSL_set_incoming_stream_policy(c.handle, ssl.SSL_INCOMING_STREAM_POLICY_REJECT, 0) != 1)
+ return error.Tls;
+ _ = ssl.SSL_set_mode(stream, ssl.SSL_MODE_ACCEPT_MOVING_WRITE_BUFFER);
+ c.stream = stream;
+ return true;
+ }
+
+ pub fn read(c: *Connection, bytes: []u8) Error!?usize {
+ if (!try c.ready()) return null;
+ ssl.ERR_clear_error();
+ var len: usize = 0;
+ const rc = ssl.SSL_read_ex(c.stream, bytes.ptr, bytes.len, &len);
+ if (rc == 1) return len;
+ if (ssl.SSL_get_error(c.stream, rc) == ssl.SSL_ERROR_ZERO_RETURN) return 0;
+ try retry(c.stream.?, rc);
+ return null;
+ }
+
+ pub fn pending(c: *const Connection) bool {
+ var close_info: ssl.SSL_CONN_CLOSE_INFO = undefined;
+ if (ssl.SSL_get_conn_close_info(c.handle, &close_info, @sizeOf(@TypeOf(close_info))) == 1) return true;
+ if (c.stream) |stream| {
+ var item: ssl.SSL_POLL_ITEM = .{ .desc = ssl.SSL_as_poll_descriptor(stream), .events = ssl.SSL_POLL_EVENT_RE, .revents = 0 };
+ const timeout: ssl.struct_timeval = .{ .tv_sec = 0, .tv_usec = 0 };
+ if (ssl.SSL_poll(&item, 1, @sizeOf(@TypeOf(item)), &timeout, ssl.SSL_POLL_FLAG_NO_HANDLE_EVENTS, null) != 1) return true;
+ return item.revents != 0;
+ }
+ return ssl.SSL_get_accept_stream_queue_len(c.handle) != 0;
+ }
+
+ pub fn write(c: *Connection, bytes: []const u8) Error!usize {
+ if (!try c.ready()) return 0;
+ if (bytes.len < c.pending_write_len) return error.InvalidWrite;
+ const requested = if (c.pending_write_len != 0) c.pending_write_len else bytes.len;
+ if (requested == 0) return 0;
+ ssl.ERR_clear_error();
+ var len: usize = 0;
+ const rc = ssl.SSL_write_ex(c.stream, bytes.ptr, requested, &len);
+ if (rc == 1) {
+ c.pending_write_len = 0;
+ return len;
+ }
+ try retry(c.stream.?, rc);
+ c.pending_write_len = requested;
+ return 0;
+ }
+
+ pub fn conclude(c: *Connection) Error!void {
+ if (c.pending_write_len != 0) return error.InvalidWrite;
+ if (!try c.ready()) return error.Closed;
+ ssl.ERR_clear_error();
+ if (ssl.SSL_stream_conclude(c.stream, 0) != 1) return error.Tls;
+ }
+
+ pub fn events(c: *Connection) Error!void {
+ if (c.fd < 0) return;
+ ssl.ERR_clear_error();
+ if (ssl.SSL_handle_events(c.handle) != 1) return error.Tls;
+ }
+
+ pub fn poll(c: *const Connection) ?libc.pollfd {
+ return if (c.fd >= 0) pollFd(c.handle, c.fd) else null;
+ }
+
+ pub fn nextDue(c: *const Connection) ?i32 {
+ return if (c.fd >= 0) due(c.handle) else null;
+ }
+
+ pub fn deinit(c: *Connection) void {
+ ssl.ERR_clear_error();
+ _ = ssl.SSL_shutdown_ex(c.handle, ssl.SSL_SHUTDOWN_FLAG_RAPID | ssl.SSL_SHUTDOWN_FLAG_NO_STREAM_FLUSH | ssl.SSL_SHUTDOWN_FLAG_NO_BLOCK, null, 0);
+ ssl.SSL_free(c.stream);
+ ssl.SSL_free(c.handle);
+ if (c.fd >= 0) _ = libc.close(c.fd);
+ c.* = undefined;
+ }
+ };
+
+ fn udp(family: u16) Error!c_int {
+ const fd = libc.socket(family, libc.SOCK.DGRAM, 0);
+ if (fd < 0) return error.Socket;
+ errdefer _ = libc.close(fd);
+ const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0));
+ if (flags < 0) return error.SocketFlags;
+ var options: libc.O = @bitCast(@as(u32, @bitCast(flags)));
+ options.NONBLOCK = true;
+ if (libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(options))))) != 0 or
+ libc.fcntl(fd, libc.F.SETFD, @as(c_int, libc.FD_CLOEXEC)) != 0) return error.SocketFlags;
+ if (family == libc.AF.INET6) {
+ const enabled: c_int = 1;
+ const v6only = if (@import("builtin").os.tag.isDarwin()) 27 else libc.IPV6.V6ONLY;
+ if (libc.setsockopt(fd, libc.IPPROTO.IPV6, v6only, &enabled, @sizeOf(c_int)) != 0) return error.SocketOption;
+ }
+ return fd;
+ }
+
+ fn sockaddr(address: std.Io.net.IpAddress, out: *libc.sockaddr.storage) libc.socklen_t {
+ switch (address) {
+ .ip4 => |ip| {
+ const addr: *libc.sockaddr.in = @ptrCast(out);
+ addr.* = .{ .port = std.mem.nativeToBig(u16, ip.port), .addr = @bitCast(ip.bytes) };
+ return @sizeOf(libc.sockaddr.in);
+ },
+ .ip6 => |ip| {
+ const addr: *libc.sockaddr.in6 = @ptrCast(out);
+ addr.* = .{ .port = std.mem.nativeToBig(u16, ip.port), .addr = ip.bytes, .flowinfo = 0, .scope_id = 0 };
+ return @sizeOf(libc.sockaddr.in6);
+ },
+ }
+ }
+
+ fn pollFd(handle: *ssl.SSL, fd: c_int) libc.pollfd {
+ var result: libc.pollfd = .{ .fd = fd, .events = 0, .revents = 0 };
+ if (ssl.SSL_net_read_desired(handle) == 1) result.events |= libc.POLL.IN;
+ if (ssl.SSL_net_write_desired(handle) == 1) result.events |= libc.POLL.OUT;
+ return result;
+ }
+
+ fn due(handle: *ssl.SSL) ?i32 {
+ var tv: ssl.struct_timeval = undefined;
+ var infinite: c_int = undefined;
+ if (ssl.SSL_get_event_timeout(handle, &tv, &infinite) != 1) return 0;
+ if (infinite != 0) return null;
+ const ms = @as(i128, tv.tv_sec) * 1000 + @divFloor(@as(i128, tv.tv_usec) + 999, 1000);
+ return @intCast(std.math.clamp(ms, 0, std.math.maxInt(i32)));
+ }
+
+ fn retry(handle: *ssl.SSL, rc: c_int) Error!void {
+ switch (ssl.SSL_get_error(handle, rc)) {
+ ssl.SSL_ERROR_WANT_READ, ssl.SSL_ERROR_WANT_WRITE => {},
+ ssl.SSL_ERROR_ZERO_RETURN => return error.Closed,
+ else => return error.Tls,
+ }
+ }
+
+ fn selectAlpn(_: ?*ssl.SSL, out: [*c][*c]const u8, outlen: [*c]u8, input: [*c]const u8, len: c_uint, _: ?*anyopaque) callconv(.c) c_int {
+ var offset: usize = 0;
+ while (offset < len) {
+ const size = input[offset];
+ offset += 1;
+ if (size > len - offset) return ssl.SSL_TLSEXT_ERR_ALERT_FATAL;
+ if (std.mem.eql(u8, input[offset..][0..size], alpn)) {
+ out.* = input + offset;
+ outlen.* = size;
+ return ssl.SSL_TLSEXT_ERR_OK;
+ }
+ offset += size;
+ }
+ return ssl.SSL_TLSEXT_ERR_ALERT_FATAL;
+ }
+ };
+}
diff --git a/src/root.zig b/src/root.zig
new file mode 100644
index 0000000..47303be
--- /dev/null
+++ b/src/root.zig
@@ -0,0 +1,53 @@
+//! Allocation-free 9P2000 protocol library. See docs/design.md for ownership contracts.
+pub const wire = @import("wire.zig");
+pub const Error = wire.Error;
+pub const Type = wire.Type;
+pub const isT = wire.isT;
+pub const header_len = wire.header_len;
+pub const qid_len = wire.qid_len;
+pub const stat_fixed = wire.stat_fixed;
+pub const notag = wire.notag;
+pub const nofid = wire.nofid;
+pub const max_welem = wire.max_welem;
+pub const iohdrsz = wire.iohdrsz;
+pub const Qid = wire.Qid;
+pub const Stat = wire.Stat;
+pub const Msg = wire.Msg;
+pub const Decoded = wire.Decoded;
+pub const frameLen = wire.frameLen;
+pub const encode = wire.encode;
+pub const decode = wire.decode;
+pub const encodedLen = wire.encodedLen;
+pub const qtdir = wire.qtdir;
+pub const qtappend = wire.qtappend;
+pub const qtexcl = wire.qtexcl;
+pub const qtmount = wire.qtmount;
+pub const qtauth = wire.qtauth;
+pub const qttmp = wire.qttmp;
+pub const qtfile = wire.qtfile;
+pub const dmdir = wire.dmdir;
+pub const dmappend = wire.dmappend;
+pub const dmexcl = wire.dmexcl;
+pub const dmmount = wire.dmmount;
+pub const dmauth = wire.dmauth;
+pub const dmtmp = wire.dmtmp;
+pub const dmperm = wire.dmperm;
+pub const Client = @import("client.zig").Client;
+pub const ClientError = @import("client.zig").ClientError;
+pub const max_tags = @import("client.zig").max_tags;
+pub const Server = @import("Server.zig");
+pub const oread: u8 = 0;
+pub const owrite: u8 = 1;
+pub const ordwr: u8 = 2;
+pub const oexec: u8 = 3;
+pub const otrunc: u8 = 16;
+pub const ocexec: u8 = 32;
+pub const orclose: u8 = 64;
+test {
+ @import("std").testing.refAllDecls(@This());
+}
+pub const transport = @import("transport.zig");
+pub const Quic = @import("quic.zig").Quic;
+test {
+ _ = @import("session_test.zig");
+}
diff --git a/src/session_test.zig b/src/session_test.zig
new file mode 100644
index 0000000..8600d03
--- /dev/null
+++ b/src/session_test.zig
@@ -0,0 +1,222 @@
+const std = @import("std");
+const testing = std.testing;
+const c9 = @import("root.zig");
+const Pair = struct {
+ client_in: [4096]u8 = undefined,
+ client_out: [4096]u8 = undefined,
+ server_in: [4096]u8 = undefined,
+ server_out: [8192]u8 = undefined,
+ client: c9.Client = undefined,
+ server: c9.Server = undefined,
+
+ fn init(p: *Pair) !void {
+ p.client = .init(.{ .in = &p.client_in, .out = &p.client_out });
+ p.server = .init(.{ .in = &p.server_in, .out = &p.server_out });
+ _ = try p.client.submit(.{ .version = .{} });
+ const request = try p.nextRequest();
+ try p.server.negotiate(request.msg.tversion.msize, request.msg.tversion.version);
+ p.server.release();
+ try testing.expectEqual(c9.Client.Op.version, (try p.nextResult()).op);
+ }
+
+ fn nextRequest(p: *Pair) !c9.Decoded {
+ // Exercise every frame boundary through single-byte delivery.
+ while (p.client.output().len != 0) {
+ try testing.expectEqual(@as(usize, 1), p.server.push(p.client.output()[0..1]));
+ p.client.wrote(1);
+ if (try p.server.receive()) |request| return request;
+ }
+ return (try p.server.receive()) orelse error.NoRequest;
+ }
+
+ fn nextResult(p: *Pair) !c9.Client.Done {
+ while (p.server.output().len != 0) {
+ try testing.expectEqual(@as(usize, 1), p.client.push(p.server.output()[0..1]));
+ p.server.wrote(1);
+ if (p.client.take()) |result| return result;
+ }
+ return p.client.take() orelse error.NoResult;
+ }
+
+ fn exchange(p: *Pair, req: c9.Client.Request, response: c9.Msg) !void {
+ const tag = try p.client.submit(req);
+ const request = try p.nextRequest();
+ try testing.expectEqual(tag, request.tag);
+ try p.server.reply(tag, response);
+ p.server.release();
+ const result = try p.nextResult();
+ try testing.expectEqual(std.meta.activeTag(req), result.op);
+ try testing.expectEqualStrings(@tagName(req), @tagName(result.result));
+ }
+};
+
+const qid: c9.Qid = .{ .type = 0, .version = 1, .path = 2 };
+const stat: c9.Stat = .{
+ .type = 0,
+ .dev = 0,
+ .qid = qid,
+ .mode = 0o600,
+ .atime = 0,
+ .mtime = 0,
+ .length = 0,
+ .name = "file",
+ .uid = "user",
+ .gid = "group",
+ .muid = "user",
+};
+
+test "all base 9P2000 client operations through fragmented server transport" {
+ var p: Pair = .{};
+ try p.init();
+ try p.exchange(.{ .auth = .{ .afid = 1, .uname = "user" } }, .{ .rauth = .{ .aqid = qid } });
+ try p.exchange(.{ .attach = .{ .fid = 2, .afid = 1, .uname = "user" } }, .{ .rattach = .{ .qid = qid } });
+ try p.exchange(.{ .walk = .{ .fid = 2, .newfid = 3, .names = &.{"file"} } }, .{ .rwalk = .{ .nwqid = 1, .wqid = @splat(qid) } });
+ try p.exchange(.{ .open = .{ .fid = 3, .mode = 2 } }, .{ .ropen = .{ .qid = qid, .iounit = 0 } });
+ try p.exchange(.{ .create = .{ .fid = 2, .name = "new", .perm = 0o600, .mode = 2 } }, .{ .rcreate = .{ .qid = qid, .iounit = 0 } });
+ try p.exchange(.{ .read = .{ .fid = 3, .offset = 0, .count = 3 } }, .{ .rread = .{ .data = "abc" } });
+ try p.exchange(.{ .write = .{ .fid = 3, .offset = 0, .data = "abc" } }, .{ .rwrite = .{ .count = 2 } });
+ try p.exchange(.{ .stat = .{ .fid = 3 } }, .{ .rstat = .{ .stat = stat } });
+ try p.exchange(.{ .wstat = .{ .fid = 3, .stat = stat } }, .rwstat);
+ try p.exchange(.{ .clunk = .{ .fid = 3 } }, .rclunk);
+ try p.exchange(.{ .remove = .{ .fid = 2 } }, .rremove);
+ try p.exchange(.{ .flush = .{ .oldtag = 42 } }, .rflush);
+}
+
+test "flush holds oldtag until Rflush, including a completed original request" {
+ for ([_]bool{ false, true }) |complete| {
+ var p: Pair = .{};
+ try p.init();
+ const oldtag = try p.client.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 1 } });
+ _ = try p.nextRequest();
+ p.server.release();
+ const flush = try p.client.submit(.{ .flush = .{ .oldtag = oldtag } });
+ _ = try p.nextRequest();
+ p.server.release();
+ if (complete) {
+ try p.server.reply(oldtag, .{ .rread = .{ .data = "x" } });
+ try testing.expectEqual(oldtag, (try p.nextResult()).tag);
+ try testing.expectEqual(@as(usize, 2), p.client.pending());
+ }
+ const other = try p.client.submit(.{ .stat = .{ .fid = 1 } });
+ try testing.expect(other != oldtag and other != flush);
+ try p.server.reply(flush, .rflush);
+ try testing.expectEqual(flush, (try p.nextResult()).tag);
+ try testing.expectEqual(@as(usize, 1), p.client.pending());
+ try testing.expectError(error.UnknownTag, p.server.reply(oldtag, .{ .rread = .{ .data = "x" } }));
+ const reused = try p.client.submit(.{ .stat = .{ .fid = 1 } });
+ try testing.expectEqual(oldtag, reused);
+ }
+}
+
+test "response bounds and reply types are checked before output is changed" {
+ var p: Pair = .{};
+ try p.init();
+ const tag = try p.client.submit(.{ .write = .{ .fid = 1, .offset = 0, .data = "a" } });
+ _ = try p.nextRequest();
+ p.server.release();
+ try testing.expectError(error.WrongReply, p.server.reply(tag, .{ .rwrite = .{ .count = 2 } }));
+ try testing.expectError(error.WrongReply, p.server.reply(tag, .rclunk));
+ try testing.expectEqual(@as(usize, 0), p.server.output().len);
+ try p.server.reply(tag, .{ .rwrite = .{ .count = 1 } });
+ _ = try p.nextResult();
+}
+
+test "stat outer length overflow is rejected without writing output" {
+ var name: [65535 - 48]u8 = @splat('a');
+ var entry = stat;
+ entry.name = &name;
+ entry.uid = "";
+ entry.gid = "";
+ entry.muid = "";
+ var buffer: [65550]u8 = @splat(0xaa);
+ try testing.expectError(error.Overlong, c9.encode(.{ .rstat = .{ .stat = entry } }, 1, &buffer));
+ try testing.expect(std.mem.allEqual(u8, &buffer, 0xaa));
+}
+
+test "arbitrary input decoder round trip" {
+ try testing.fuzz({}, fuzzDecode, .{});
+}
+
+fn fuzzDecode(_: void, smith: *testing.Smith) !void {
+ var input_buffer: [65536]u8 = undefined;
+ const input = input_buffer[0..smith.slice(&input_buffer)];
+ const decoded = c9.decode(input) catch return;
+ var buffer: [65536]u8 = undefined;
+ const encoded = try c9.encode(decoded.msg, decoded.tag, &buffer);
+ try testing.expectEqualSlices(u8, input, encoded);
+}
+
+test "flush can cancel a full request window and reserves unused oldtags" {
+ var p: Pair = .{};
+ try p.init();
+ _ = try p.client.submit(.{ .flush = .{ .oldtag = 2 } });
+ const next = try p.client.submit(.{ .stat = .{ .fid = 0 } });
+ try testing.expect(next != 2);
+ const third = try p.client.submit(.{ .stat = .{ .fid = 0 } });
+ try testing.expect(third != 2);
+ p.client.hangup();
+ try p.init();
+ for (0..c9.max_tags) |_| _ = try p.client.submit(.{ .stat = .{ .fid = 0 } });
+ try testing.expectError(error.NoTags, p.client.submit(.{ .stat = .{ .fid = 0 } }));
+ const flush = try p.client.submit(.{ .flush = .{ .oldtag = 0 } });
+ try testing.expectEqual(@as(u16, c9.max_tags), flush);
+}
+
+test "multiple flushes retain reservations until each response and can themselves be flushed" {
+ var p: Pair = .{};
+ try p.init();
+ const oldtag = try p.client.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 1 } });
+ _ = try p.nextRequest();
+ p.server.release();
+ const first = try p.client.submit(.{ .flush = .{ .oldtag = oldtag } });
+ _ = try p.nextRequest();
+ p.server.release();
+ const second = try p.client.submit(.{ .flush = .{ .oldtag = oldtag } });
+ _ = try p.nextRequest();
+ p.server.release();
+ try p.server.reply(first, .rflush);
+ _ = try p.nextResult();
+ const held = try p.client.submit(.{ .stat = .{ .fid = 1 } });
+ try testing.expect(held != oldtag);
+ try p.server.reply(second, .rflush);
+ _ = try p.nextResult();
+ try testing.expectEqual(oldtag, try p.client.submit(.{ .stat = .{ .fid = 1 } }));
+
+ try p.init();
+ const read_tag = try p.client.submit(.{ .read = .{ .fid = 1, .offset = 0, .count = 1 } });
+ _ = try p.nextRequest();
+ p.server.release();
+ const flush_tag = try p.client.submit(.{ .flush = .{ .oldtag = read_tag } });
+ _ = try p.nextRequest();
+ p.server.release();
+ const cancel_flush = try p.client.submit(.{ .flush = .{ .oldtag = flush_tag } });
+ _ = try p.nextRequest();
+ p.server.release();
+ try p.server.reply(read_tag, .{ .rread = .{ .data = "x" } });
+ _ = try p.nextResult();
+ try p.server.reply(cancel_flush, .rflush);
+ _ = try p.nextResult();
+ try testing.expectEqual(@as(usize, 0), p.client.pending());
+}
+
+test "server reserves cancellation capacity and rejects duplicate pending tags" {
+ var input: [4096]u8 = undefined;
+ var output: [4096]u8 = undefined;
+ var server: c9.Server = .init(.{ .in = &input, .out = &output });
+ server.msize = 4096;
+ var frame: [64]u8 = undefined;
+ for (0..64) |i| {
+ const bytes = try c9.encode(.{ .tstat = .{ .fid = 1 } }, @intCast(i), &frame);
+ _ = server.push(bytes);
+ _ = (try server.receive()).?;
+ server.release();
+ }
+ _ = server.push(try c9.encode(.{ .tflush = .{ .oldtag = 0 } }, 100, &frame));
+ _ = (try server.receive()).?;
+ server.release();
+ try server.reply(100, .rflush);
+ server.wrote(server.output().len);
+ _ = server.push(try c9.encode(.{ .tstat = .{ .fid = 1 } }, 1, &frame));
+ try testing.expectError(error.Protocol, server.receive());
+ try testing.expect(server.dead);
+}
diff --git a/src/transport.zig b/src/transport.zig
new file mode 100644
index 0000000..d7619be
--- /dev/null
+++ b/src/transport.zig
@@ -0,0 +1,207 @@
+//! Stream transport adapters. Namespace paths, mounting, connection limits, and
+//! event-loop scheduling belong to the application. POSIX descriptors are owned
+//! by the caller; std.Io streams retain the standard library ownership contract.
+const std = @import("std");
+const libc = std.c;
+const wire = @import("wire.zig");
+pub const darwin = @import("builtin").os.tag.isDarwin();
+pub const sun_path_len = @typeInfo(@FieldType(libc.sockaddr.un, "path")).array.len;
+pub const Address = union(enum) { unix: [:0]const u8, tcp: std.Io.net.IpAddress };
+pub const Error = error{ Socket, Flags, SocketOption, BadAddress, Bind, Listen, Connect, Closed, Timeout, Io };
+
+/// Read one frame from any std.Io reader, including TCP, Unix sockets, or files.
+pub fn readFrame(reader: *std.Io.Reader, buffer: []u8, msize: u32) ![]const u8 {
+ if (buffer.len < wire.header_len) return error.NoSpace;
+ try reader.readSliceAll(buffer[0..4]);
+ const len = wire.frameLen(buffer[0..4]).?;
+ if (len < wire.header_len) return error.BadValue;
+ if (len > msize or len > buffer.len) return error.Overlong;
+ try reader.readSliceAll(buffer[4..len]);
+ return buffer[0..len];
+}
+
+/// Batching and flushing are explicit: this does not flush the writer.
+pub fn writeFrame(writer: *std.Io.Writer, frame: []const u8, msize: u32) !void {
+ if (frame.len > msize) return error.Overlong;
+ _ = try wire.decode(frame);
+ try writer.writeAll(frame);
+}
+
+pub fn connect(io: std.Io, address: Address) !std.Io.net.Stream {
+ return switch (address) {
+ .tcp => |ip| ip.connect(io, .{ .mode = .stream }),
+ .unix => |path| (try std.Io.net.UnixAddress.init(path)).connect(io),
+ };
+}
+
+pub fn listen(io: std.Io, address: Address, backlog: u31) !std.Io.net.Server {
+ return switch (address) {
+ .tcp => |ip| ip.listen(io, .{ .kernel_backlog = backlog }),
+ .unix => |path| (try std.Io.net.UnixAddress.init(path)).listen(io, .{ .kernel_backlog = backlog }),
+ };
+}
+
+pub fn nowMs() i64 {
+ var ts: libc.timespec = undefined;
+ if (libc.clock_gettime(.MONOTONIC, &ts) != 0) return std.math.maxInt(i64);
+ return @as(i64, ts.sec) * std.time.ms_per_s + @divTrunc(ts.nsec, std.time.ns_per_ms);
+}
+
+pub fn configure(fd: c_int, tcp: bool) Error!void {
+ const flags = libc.fcntl(fd, libc.F.GETFL, @as(c_int, 0));
+ if (flags < 0) return error.Flags;
+ var options: libc.O = @bitCast(@as(u32, @bitCast(flags)));
+ options.NONBLOCK = true;
+ if (libc.fcntl(fd, libc.F.SETFL, @as(c_int, @bitCast(@as(u32, @bitCast(options))))) != 0 or
+ libc.fcntl(fd, libc.F.SETFD, @as(c_int, libc.FD_CLOEXEC)) != 0) return error.Flags;
+ const on: c_int = 1;
+ if (tcp and libc.setsockopt(fd, libc.IPPROTO.TCP, libc.TCP.NODELAY, &on, @sizeOf(c_int)) != 0)
+ return error.SocketOption;
+ if (comptime darwin) {
+ if (libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.NOSIGPIPE, &on, @sizeOf(c_int)) != 0)
+ return error.SocketOption;
+ }
+}
+
+pub fn ipSockaddr(ip: std.Io.net.IpAddress, out: *libc.sockaddr.storage) libc.socklen_t {
+ return switch (ip) {
+ .ip4 => |a| blk: {
+ const addr: *libc.sockaddr.in = @ptrCast(out);
+ addr.* = .{ .port = std.mem.nativeToBig(u16, a.port), .addr = @bitCast(a.bytes) };
+ break :blk @sizeOf(libc.sockaddr.in);
+ },
+ .ip6 => |a| blk: {
+ const addr: *libc.sockaddr.in6 = @ptrCast(out);
+ addr.* = .{ .port = std.mem.nativeToBig(u16, a.port), .addr = a.bytes, .flowinfo = 0, .scope_id = a.interface.index };
+ break :blk @sizeOf(libc.sockaddr.in6);
+ },
+ };
+}
+
+fn sockaddr(address: Address, out: *libc.sockaddr.storage) Error!libc.socklen_t {
+ return switch (address) {
+ .tcp => |ip| ipSockaddr(ip, out),
+ .unix => |path| blk: {
+ if (path.len >= sun_path_len or std.mem.indexOfScalar(u8, path, 0) != null)
+ return error.BadAddress;
+ const addr: *libc.sockaddr.un = @ptrCast(out);
+ addr.* = .{ .path = @splat(0) };
+ @memcpy(addr.path[0..path.len], path);
+ break :blk @sizeOf(libc.sockaddr.un);
+ },
+ };
+}
+
+pub fn wait(fd: c_int, events: i16, deadline_ms: i64) Error!void {
+ while (true) {
+ const left = deadline_ms -| nowMs();
+ if (left <= 0) return error.Timeout;
+ var fds = [1]libc.pollfd{.{ .fd = fd, .events = events, .revents = 0 }};
+ const rc = libc.poll(&fds, 1, @intCast(@min(left, std.math.maxInt(c_int))));
+ if (rc < 0) {
+ if (libc.errno(rc) == .INTR) continue;
+ return error.Io;
+ }
+ if (rc == 0) continue;
+ if (nowMs() >= deadline_ms) return error.Timeout;
+ if (fds[0].revents & events != 0) return;
+ return error.Closed;
+ }
+}
+
+/// Connect a nonblocking socket with an absolute monotonic deadline.
+pub fn connectFd(address: Address, deadline_ms: i64) Error!c_int {
+ var addr: libc.sockaddr.storage = undefined;
+ const len = try sockaddr(address, &addr);
+ const fd = libc.socket(addr.family, libc.SOCK.STREAM, 0);
+ if (fd < 0) return error.Socket;
+ errdefer close(fd);
+ try configure(fd, address == .tcp);
+ const rc = libc.connect(fd, @ptrCast(&addr), len);
+ if (rc != 0) {
+ switch (libc.errno(rc)) {
+ .INPROGRESS, .ALREADY, .INTR => {},
+ else => return error.Connect,
+ }
+ try wait(fd, @intCast(libc.POLL.OUT), deadline_ms);
+ var status: c_int = 0;
+ var size: libc.socklen_t = @sizeOf(c_int);
+ if (libc.getsockopt(fd, libc.SOL.SOCKET, libc.SO.ERROR, @ptrCast(&status), &size) != 0)
+ return error.Connect;
+ if (status != 0) return error.Connect;
+ }
+ return fd;
+}
+
+/// Does not unlink Unix paths or alter their permissions.
+pub fn listenFd(address: Address, backlog: u31) Error!c_int {
+ var addr: libc.sockaddr.storage = undefined;
+ const len = try sockaddr(address, &addr);
+ const fd = libc.socket(addr.family, libc.SOCK.STREAM, 0);
+ if (fd < 0) return error.Socket;
+ errdefer close(fd);
+ try configure(fd, false);
+ if (address == .tcp) {
+ const on: c_int = 1;
+ if (libc.setsockopt(fd, libc.SOL.SOCKET, libc.SO.REUSEADDR, &on, @sizeOf(c_int)) != 0)
+ return error.SocketOption;
+ if (address.tcp == .ip6) {
+ const v6only = if (darwin) 27 else std.os.linux.IPV6.V6ONLY;
+ if (libc.setsockopt(fd, libc.IPPROTO.IPV6, v6only, &on, @sizeOf(c_int)) != 0)
+ return error.SocketOption;
+ }
+ }
+ if (libc.bind(fd, @ptrCast(&addr), len) != 0) return error.Bind;
+ if (libc.listen(fd, backlog) != 0) return error.Listen;
+ return fd;
+}
+
+pub fn acceptFd(listener: c_int, tcp: bool) Error!?c_int {
+ const fd = libc.accept(listener, null, null);
+ if (fd < 0) return switch (libc.errno(fd)) {
+ .AGAIN, .INTR, .CONNABORTED => null,
+ else => error.Socket,
+ };
+ errdefer close(fd);
+ try configure(fd, tcp);
+ return fd;
+}
+
+/// null means retry after readiness; zero means EOF. Empty buffers are forbidden.
+pub fn read(fd: c_int, buffer: []u8) Error!?usize {
+ std.debug.assert(buffer.len > 0);
+ const n = libc.read(fd, buffer.ptr, buffer.len);
+ if (n < 0) return switch (libc.errno(n)) {
+ .INTR, .AGAIN => null,
+ else => error.Io,
+ };
+ return @intCast(n);
+}
+
+pub fn write(fd: c_int, bytes: []const u8) Error!?usize {
+ std.debug.assert(bytes.len > 0);
+ const n = libc.send(fd, bytes.ptr, bytes.len, if (darwin) 0 else libc.MSG.NOSIGNAL);
+ if (n < 0) return switch (libc.errno(n)) {
+ .INTR, .AGAIN => null,
+ else => error.Io,
+ };
+ if (n == 0) return error.Closed;
+ return @intCast(n);
+}
+
+pub fn close(fd: c_int) void {
+ _ = libc.close(fd);
+}
+
+/// Probe a Unix listener without changing namespace entries. Uncertainty is live.
+pub fn isListening(path: [:0]const u8) bool {
+ var addr: libc.sockaddr.un = .{ .path = @splat(0) };
+ if (path.len + 1 > sun_path_len) return true; // cannot ask; assume occupied
+ @memcpy(addr.path[0 .. path.len + 1], path[0 .. path.len + 1]);
+ const fd = libc.socket(libc.AF.UNIX, libc.SOCK.STREAM, 0);
+ if (fd < 0) return true;
+ defer _ = libc.close(fd);
+ configure(fd, false) catch return true;
+ if (libc.connect(fd, @ptrCast(&addr), @sizeOf(@TypeOf(addr))) == 0) return true;
+ return libc.errno(-1) != .CONNREFUSED;
+}
diff --git a/src/wire.zig b/src/wire.zig
new file mode 100644
index 0000000..e093450
--- /dev/null
+++ b/src/wire.zig
@@ -0,0 +1,1153 @@
+//! Base 9P2000 wire format. Decoded strings and data borrow the input frame.
+//! Encoding performs a complete size check before writing the caller-owned buffer.
+const std = @import("std");
+const assert = std.debug.assert;
+
+pub const Error = error{
+ Truncated,
+ Overlong,
+ BadTag,
+ BadValue,
+ Trailing,
+ NoSpace,
+};
+
+pub const Type = enum(u8) {
+ tversion = 100,
+ rversion = 101,
+ tauth = 102,
+ rauth = 103,
+ tattach = 104,
+ rattach = 105,
+ terror = 106,
+ rerror = 107,
+ tflush = 108,
+ rflush = 109,
+ twalk = 110,
+ rwalk = 111,
+ topen = 112,
+ ropen = 113,
+ tcreate = 114,
+ rcreate = 115,
+ tread = 116,
+ rread = 117,
+ twrite = 118,
+ rwrite = 119,
+ tclunk = 120,
+ rclunk = 121,
+ tremove = 122,
+ rremove = 123,
+ tstat = 124,
+ rstat = 125,
+ twstat = 126,
+ rwstat = 127,
+ _,
+};
+
+pub fn isT(t: Type) bool {
+ return @intFromEnum(t) % 2 == 0;
+}
+
+pub const header_len: usize = 4 + 1 + 2;
+
+pub const qid_len: usize = 1 + 4 + 8;
+
+pub const stat_fixed: usize = 2 + qid_len + 5 * 2 + 4 * 4 + 8;
+
+pub const notag: u16 = 0xFFFF;
+
+pub const nofid: u32 = 0xFFFF_FFFF;
+
+pub const max_welem: usize = 16;
+
+const test_msize: u32 = 4096;
+
+pub const iohdrsz: u32 = 24;
+
+pub const qtdir: u8 = 0x80;
+pub const qtappend: u8 = 0x40;
+pub const qtexcl: u8 = 0x20;
+pub const qtmount: u8 = 0x10;
+pub const qtauth: u8 = 0x08;
+pub const qttmp: u8 = 0x04;
+pub const qtfile: u8 = 0x00;
+
+pub const dmdir: u32 = 0x8000_0000;
+pub const dmappend: u32 = 0x4000_0000;
+pub const dmexcl: u32 = 0x2000_0000;
+pub const dmmount: u32 = 0x1000_0000;
+pub const dmauth: u32 = 0x0800_0000;
+pub const dmtmp: u32 = 0x0400_0000;
+pub const dmperm: u32 = 0o777;
+
+comptime {
+ assert(header_len == 7);
+ assert(qid_len == 13);
+ assert(stat_fixed == 49);
+ for (std.enums.values(Type)) |t| {
+ const even = @intFromEnum(t) % 2 == 0;
+ assert(isT(t) == even);
+ assert(std.mem.startsWith(u8, @tagName(t), if (even) "t" else "r"));
+ }
+}
+
+pub const Qid = struct {
+ type: u8,
+ version: u32,
+ path: u64,
+
+ pub fn encode(self: Qid, buf: []u8) Error![]u8 {
+ if (buf.len < qid_len) return error.NoSpace;
+ buf[0] = self.type;
+ std.mem.writeInt(u32, buf[1..5], self.version, .little);
+ std.mem.writeInt(u64, buf[5..13], self.path, .little);
+ return buf[0..qid_len];
+ }
+
+ pub fn decode(bytes: []const u8) Error!Qid {
+ if (bytes.len < qid_len) return error.Truncated;
+ return .{
+ .type = bytes[0],
+ .version = std.mem.readInt(u32, bytes[1..5], .little),
+ .path = std.mem.readInt(u64, bytes[5..13], .little),
+ };
+ }
+};
+
+pub const Stat = struct {
+ type: u16,
+ dev: u32,
+ qid: Qid,
+ mode: u32,
+ atime: u32,
+ mtime: u32,
+ length: u64,
+ name: []const u8,
+ uid: []const u8,
+ gid: []const u8,
+ muid: []const u8,
+
+ pub fn size(self: Stat) Error!u16 {
+ const n =
+ 2 + // type
+ 4 + // dev
+ qid_len + // qid: type[1] version[4] path[8]
+ 4 + // mode
+ 4 + // atime
+ 4 + // mtime
+ 8 + // length
+ try stringLen(self.name) +
+ try stringLen(self.uid) +
+ try stringLen(self.gid) +
+ try stringLen(self.muid);
+ assert(n >= stat_fixed - 2);
+ if (n > std.math.maxInt(u16)) return error.Overlong;
+ return @intCast(n);
+ }
+
+ pub fn encode(self: Stat, buf: []u8) Error![]u8 {
+ const n = try self.size();
+ const total = @as(usize, n) + 2;
+ if (buf.len < total) return error.NoSpace;
+ var w: Writer = .init(buf[0..total]);
+ try w.putU16(n);
+ try w.putU16(self.type);
+ try w.putU32(self.dev);
+ try w.putQid(self.qid);
+ try w.putU32(self.mode);
+ try w.putU32(self.atime);
+ try w.putU32(self.mtime);
+ try w.putU64(self.length);
+ try w.putString(self.name);
+ try w.putString(self.uid);
+ try w.putString(self.gid);
+ try w.putString(self.muid);
+ assert(w.n == total);
+ return buf[0..total];
+ }
+
+ pub fn decode(bytes: []const u8) Error!Stat {
+ var r: Reader = .init(bytes);
+ const n = try r.getU16();
+ const body = bytes.len - 2;
+ if (n > body) return error.Truncated;
+ if (n < body) return error.Trailing;
+ const self: Stat = .{
+ .type = try r.getU16(),
+ .dev = try r.getU32(),
+ .qid = try r.getQid(),
+ .mode = try r.getU32(),
+ .atime = try r.getU32(),
+ .mtime = try r.getU32(),
+ .length = try r.getU64(),
+ .name = try r.getString(),
+ .uid = try r.getString(),
+ .gid = try r.getString(),
+ .muid = try r.getString(),
+ };
+ try r.end();
+ return self;
+ }
+};
+
+pub const Msg = union(enum) {
+ tversion: struct { msize: u32, version: []const u8 },
+ rversion: struct { msize: u32, version: []const u8 },
+
+ tauth: struct { afid: u32, uname: []const u8, aname: []const u8 },
+ rauth: struct { aqid: Qid },
+
+ tattach: struct { fid: u32, afid: u32, uname: []const u8, aname: []const u8 },
+ rattach: struct { qid: Qid },
+
+ rerror: struct { ename: []const u8 },
+
+ tflush: struct { oldtag: u16 },
+ rflush: void,
+
+ twalk: struct {
+ fid: u32,
+ newfid: u32,
+ nwname: u16,
+ wname: [max_welem][]const u8 = @splat(""),
+ },
+ rwalk: struct {
+ nwqid: u16,
+ wqid: [max_welem]Qid = @splat(.{ .type = 0, .version = 0, .path = 0 }),
+ },
+
+ topen: struct { fid: u32, mode: u8 },
+ ropen: struct { qid: Qid, iounit: u32 },
+
+ tcreate: struct { fid: u32, name: []const u8, perm: u32, mode: u8 },
+ rcreate: struct { qid: Qid, iounit: u32 },
+
+ tread: struct { fid: u32, offset: u64, count: u32 },
+ rread: struct { data: []const u8 },
+
+ twrite: struct { fid: u32, offset: u64, data: []const u8 },
+ rwrite: struct { count: u32 },
+
+ tclunk: struct { fid: u32 },
+ rclunk: void,
+ tremove: struct { fid: u32 },
+ rremove: void,
+
+ tstat: struct { fid: u32 },
+ rstat: struct { stat: Stat },
+ twstat: struct { fid: u32, stat: Stat },
+ rwstat: void,
+
+ pub fn msgType(msg: Msg) Type {
+ return switch (msg) {
+ inline else => |_, t| @field(Type, @tagName(t)),
+ };
+ }
+};
+
+pub const Decoded = struct {
+ tag: u16,
+ msg: Msg,
+};
+
+pub fn frameLen(prefix: []const u8) ?u32 {
+ if (prefix.len < 4) return null;
+ return std.mem.readInt(u32, prefix[0..4], .little);
+}
+
+pub fn encodedLen(msg: Msg) Error!usize {
+ const body: u64 = switch (msg) {
+ .tversion => |m| 4 + try stringLen(m.version),
+ .rversion => |m| 4 + try stringLen(m.version),
+ .tauth => |m| 4 + try stringLen(m.uname) + try stringLen(m.aname),
+ .rauth => qid_len,
+ .tattach => |m| 4 + 4 + try stringLen(m.uname) + try stringLen(m.aname),
+ .rattach => qid_len,
+ .rerror => |m| try stringLen(m.ename),
+ .tflush => 2,
+ .rflush => 0,
+ .twalk => |m| blk: {
+ if (m.nwname > max_welem) return error.Overlong;
+ var n: usize = 4 + 4 + 2;
+ for (m.wname[0..m.nwname]) |name| n += try stringLen(name);
+ break :blk n;
+ },
+ .rwalk => |m| blk: {
+ if (m.nwqid > max_welem) return error.Overlong;
+ break :blk 2 + @as(usize, m.nwqid) * qid_len;
+ },
+ .topen => 4 + 1,
+ .ropen => qid_len + 4,
+ .tcreate => |m| 4 + try stringLen(m.name) + 4 + 1,
+ .rcreate => qid_len + 4,
+ .tread => 4 + 8 + 4,
+ .rread => |m| try dataLen(m.data),
+ .twrite => |m| 4 + 8 + try dataLen(m.data),
+ .rwrite => 4,
+ .tclunk => 4,
+ .rclunk => 0,
+ .tremove => 4,
+ .rremove => 0,
+ .tstat => 4,
+ .rstat => |m| 2 + try statLen(m.stat),
+ .twstat => |m| 4 + 2 + try statLen(m.stat),
+ .rwstat => 0,
+ };
+ const total = header_len + body;
+ if (total > std.math.maxInt(u32)) return error.Overlong;
+ return @intCast(total);
+}
+
+fn statLen(stat: Stat) Error!usize {
+ const n = @as(usize, try stat.size()) + 2;
+ if (n > std.math.maxInt(u16)) return error.Overlong;
+ return n;
+}
+
+fn stringLen(s: []const u8) Error!usize {
+ if (s.len > std.math.maxInt(u16)) return error.Overlong;
+ return 2 + s.len;
+}
+
+fn dataLen(d: []const u8) Error!u64 {
+ if (d.len > std.math.maxInt(u32)) return error.Overlong;
+ return 4 + @as(u64, d.len);
+}
+
+pub fn encode(msg: Msg, tag: u16, buf: []u8) Error![]u8 {
+ const total = try encodedLen(msg);
+ if (total > buf.len) return error.NoSpace;
+
+ var w: Writer = .init(buf[0..total]);
+ try w.putU32(@intCast(total));
+ try w.putByte(@intFromEnum(msg.msgType()));
+ try w.putU16(tag);
+
+ switch (msg) {
+ .tversion => |m| {
+ try w.putU32(m.msize);
+ try w.putString(m.version);
+ },
+ .rversion => |m| {
+ try w.putU32(m.msize);
+ try w.putString(m.version);
+ },
+ .tauth => |m| {
+ try w.putU32(m.afid);
+ try w.putString(m.uname);
+ try w.putString(m.aname);
+ },
+ .rauth => |m| try w.putQid(m.aqid),
+ .tattach => |m| {
+ try w.putU32(m.fid);
+ try w.putU32(m.afid);
+ try w.putString(m.uname);
+ try w.putString(m.aname);
+ },
+ .rattach => |m| try w.putQid(m.qid),
+ .rerror => |m| try w.putString(m.ename),
+ .tflush => |m| try w.putU16(m.oldtag),
+ .rflush => {},
+ .twalk => |m| {
+ try w.putU32(m.fid);
+ try w.putU32(m.newfid);
+ try w.putU16(m.nwname);
+ for (m.wname[0..m.nwname]) |name| try w.putString(name);
+ },
+ .rwalk => |m| {
+ try w.putU16(m.nwqid);
+ for (m.wqid[0..m.nwqid]) |qid| try w.putQid(qid);
+ },
+ .topen => |m| {
+ try w.putU32(m.fid);
+ try w.putByte(m.mode);
+ },
+ .ropen => |m| {
+ try w.putQid(m.qid);
+ try w.putU32(m.iounit);
+ },
+ .tcreate => |m| {
+ try w.putU32(m.fid);
+ try w.putString(m.name);
+ try w.putU32(m.perm);
+ try w.putByte(m.mode);
+ },
+ .rcreate => |m| {
+ try w.putQid(m.qid);
+ try w.putU32(m.iounit);
+ },
+ .tread => |m| {
+ try w.putU32(m.fid);
+ try w.putU64(m.offset);
+ try w.putU32(m.count);
+ },
+ .rread => |m| {
+ try w.putU32(@intCast(m.data.len));
+ try w.putBytes(m.data);
+ },
+ .twrite => |m| {
+ try w.putU32(m.fid);
+ try w.putU64(m.offset);
+ try w.putU32(@intCast(m.data.len));
+ try w.putBytes(m.data);
+ },
+ .rwrite => |m| try w.putU32(m.count),
+ .tclunk => |m| try w.putU32(m.fid),
+ .rclunk => {},
+ .tremove => |m| try w.putU32(m.fid),
+ .rremove => {},
+ .tstat => |m| try w.putU32(m.fid),
+ .rstat => |m| {
+ try w.putU16(try m.stat.size() + 2);
+ try w.putStat(m.stat);
+ },
+ .twstat => |m| {
+ try w.putU32(m.fid);
+ try w.putU16(try m.stat.size() + 2);
+ try w.putStat(m.stat);
+ },
+ .rwstat => {},
+ }
+
+ assert(w.n == total);
+ return buf[0..total];
+}
+
+pub fn decode(bytes: []const u8) Error!Decoded {
+ if (bytes.len < header_len) return error.Truncated;
+ const size = std.mem.readInt(u32, bytes[0..4], .little);
+ if (size < header_len) return error.BadValue;
+ if (size > bytes.len) return error.Truncated;
+ if (size < bytes.len) return error.Trailing;
+
+ const t: Type = @enumFromInt(bytes[4]);
+ const tag = std.mem.readInt(u16, bytes[5..7], .little);
+
+ var r: Reader = .init(bytes[header_len..size]);
+
+ const msg: Msg = switch (t) {
+ .tversion => .{ .tversion = .{ .msize = try r.getU32(), .version = try r.getString() } },
+ .rversion => .{ .rversion = .{ .msize = try r.getU32(), .version = try r.getString() } },
+ .tauth => .{ .tauth = .{
+ .afid = try r.getU32(),
+ .uname = try r.getString(),
+ .aname = try r.getString(),
+ } },
+ .rauth => .{ .rauth = .{ .aqid = try r.getQid() } },
+ .tattach => .{ .tattach = .{
+ .fid = try r.getU32(),
+ .afid = try r.getU32(),
+ .uname = try r.getString(),
+ .aname = try r.getString(),
+ } },
+ .rattach => .{ .rattach = .{ .qid = try r.getQid() } },
+ .rerror => .{ .rerror = .{ .ename = try r.getString() } },
+ .tflush => .{ .tflush = .{ .oldtag = try r.getU16() } },
+ .rflush => .rflush,
+ .twalk => blk: {
+ var m: Msg = .{ .twalk = .{
+ .fid = try r.getU32(),
+ .newfid = try r.getU32(),
+ .nwname = try r.getU16(),
+ } };
+ if (m.twalk.nwname > max_welem) return error.Overlong;
+ for (m.twalk.wname[0..m.twalk.nwname]) |*name| name.* = try r.getString();
+ break :blk m;
+ },
+ .rwalk => blk: {
+ var m: Msg = .{ .rwalk = .{ .nwqid = try r.getU16() } };
+ if (m.rwalk.nwqid > max_welem) return error.Overlong;
+ for (m.rwalk.wqid[0..m.rwalk.nwqid]) |*qid| qid.* = try r.getQid();
+ break :blk m;
+ },
+ .topen => .{ .topen = .{ .fid = try r.getU32(), .mode = try r.getByte() } },
+ .ropen => .{ .ropen = .{ .qid = try r.getQid(), .iounit = try r.getU32() } },
+ .tcreate => .{ .tcreate = .{
+ .fid = try r.getU32(),
+ .name = try r.getString(),
+ .perm = try r.getU32(),
+ .mode = try r.getByte(),
+ } },
+ .rcreate => .{ .rcreate = .{ .qid = try r.getQid(), .iounit = try r.getU32() } },
+ .tread => .{ .tread = .{
+ .fid = try r.getU32(),
+ .offset = try r.getU64(),
+ .count = try r.getU32(),
+ } },
+ .rread => .{ .rread = .{ .data = try r.getData() } },
+ .twrite => .{ .twrite = .{
+ .fid = try r.getU32(),
+ .offset = try r.getU64(),
+ .data = try r.getData(),
+ } },
+ .rwrite => .{ .rwrite = .{ .count = try r.getU32() } },
+ .tclunk => .{ .tclunk = .{ .fid = try r.getU32() } },
+ .rclunk => .rclunk,
+ .tremove => .{ .tremove = .{ .fid = try r.getU32() } },
+ .rremove => .rremove,
+ .tstat => .{ .tstat = .{ .fid = try r.getU32() } },
+ .rstat => .{ .rstat = .{ .stat = try Stat.decode(try r.getBlob16()) } },
+ .twstat => .{ .twstat = .{
+ .fid = try r.getU32(),
+ .stat = try Stat.decode(try r.getBlob16()),
+ } },
+ .rwstat => .rwstat,
+ .terror, _ => return error.BadTag,
+ };
+
+ try r.end();
+ return .{ .tag = tag, .msg = msg };
+}
+
+const Writer = struct {
+ buf: []u8,
+ n: usize = 0,
+
+ fn init(buf: []u8) Writer {
+ return .{ .buf = buf };
+ }
+
+ fn room(w: *Writer, k: usize) Error![]u8 {
+ if (w.buf.len - w.n < k) return error.NoSpace;
+ defer w.n += k;
+ return w.buf[w.n..][0..k];
+ }
+
+ fn putByte(w: *Writer, v: u8) Error!void {
+ (try w.room(1))[0] = v;
+ }
+
+ fn putU16(w: *Writer, v: u16) Error!void {
+ std.mem.writeInt(u16, (try w.room(2))[0..2], v, .little);
+ }
+
+ fn putU32(w: *Writer, v: u32) Error!void {
+ std.mem.writeInt(u32, (try w.room(4))[0..4], v, .little);
+ }
+
+ fn putU64(w: *Writer, v: u64) Error!void {
+ std.mem.writeInt(u64, (try w.room(8))[0..8], v, .little);
+ }
+
+ fn putBytes(w: *Writer, v: []const u8) Error!void {
+ const target = try w.room(v.len);
+ // Permit payloads staged at their final position in the output frame.
+ if (target.ptr != v.ptr) @memcpy(target, v);
+ }
+
+ fn putString(w: *Writer, v: []const u8) Error!void {
+ assert(v.len <= std.math.maxInt(u16));
+ try w.putU16(@intCast(v.len));
+ try w.putBytes(v);
+ }
+
+ fn putQid(w: *Writer, v: Qid) Error!void {
+ _ = try v.encode(try w.room(qid_len));
+ }
+
+ fn putStat(w: *Writer, v: Stat) Error!void {
+ const total = @as(usize, try v.size()) + 2;
+ _ = try v.encode(try w.room(total));
+ }
+};
+
+const Reader = struct {
+ bytes: []const u8,
+ i: usize = 0,
+
+ fn init(bytes: []const u8) Reader {
+ return .{ .bytes = bytes };
+ }
+
+ fn take(r: *Reader, n: usize) Error![]const u8 {
+ if (r.bytes.len - r.i < n) return error.Truncated;
+ defer r.i += n;
+ return r.bytes[r.i..][0..n];
+ }
+
+ fn getByte(r: *Reader) Error!u8 {
+ return (try r.take(1))[0];
+ }
+
+ fn getU16(r: *Reader) Error!u16 {
+ return std.mem.readInt(u16, (try r.take(2))[0..2], .little);
+ }
+
+ fn getU32(r: *Reader) Error!u32 {
+ return std.mem.readInt(u32, (try r.take(4))[0..4], .little);
+ }
+
+ fn getU64(r: *Reader) Error!u64 {
+ return std.mem.readInt(u64, (try r.take(8))[0..8], .little);
+ }
+
+ fn getString(r: *Reader) Error![]const u8 {
+ return r.take(try r.getU16());
+ }
+
+ fn getData(r: *Reader) Error![]const u8 {
+ return r.take(try r.getU32());
+ }
+
+ fn getBlob16(r: *Reader) Error![]const u8 {
+ return r.take(try r.getU16());
+ }
+
+ fn getQid(r: *Reader) Error!Qid {
+ return Qid.decode(try r.take(qid_len));
+ }
+
+ fn end(r: *Reader) Error!void {
+ if (r.i != r.bytes.len) return error.Trailing;
+ }
+};
+
+const testing = std.testing;
+
+fn roundTrip(buf: []u8, tag: u16, msg: Msg) !Msg {
+ const bytes = try encode(msg, tag, buf);
+ try testing.expectEqual(bytes.len, frameLen(bytes).?);
+ const got = try decode(bytes);
+ try testing.expectEqual(tag, got.tag);
+ try testing.expectEqual(msg.msgType(), got.msg.msgType());
+ try expectMsgEqual(msg, got.msg);
+ return got.msg;
+}
+
+fn expectStatEqual(want: Stat, have: Stat) !void {
+ try testing.expectEqual(want.type, have.type);
+ try testing.expectEqual(want.dev, have.dev);
+ try testing.expectEqual(want.qid, have.qid);
+ try testing.expectEqual(want.mode, have.mode);
+ try testing.expectEqual(want.atime, have.atime);
+ try testing.expectEqual(want.mtime, have.mtime);
+ try testing.expectEqual(want.length, have.length);
+ try testing.expectEqualStrings(want.name, have.name);
+ try testing.expectEqualStrings(want.uid, have.uid);
+ try testing.expectEqualStrings(want.gid, have.gid);
+ try testing.expectEqualStrings(want.muid, have.muid);
+}
+
+fn expectMsgEqual(want: Msg, have: Msg) !void {
+ switch (want) {
+ .tversion => |w| {
+ try testing.expectEqual(w.msize, have.tversion.msize);
+ try testing.expectEqualStrings(w.version, have.tversion.version);
+ },
+ .rversion => |w| {
+ try testing.expectEqual(w.msize, have.rversion.msize);
+ try testing.expectEqualStrings(w.version, have.rversion.version);
+ },
+ .tauth => |w| {
+ try testing.expectEqual(w.afid, have.tauth.afid);
+ try testing.expectEqualStrings(w.uname, have.tauth.uname);
+ try testing.expectEqualStrings(w.aname, have.tauth.aname);
+ },
+ .rauth => |w| try testing.expectEqual(w.aqid, have.rauth.aqid),
+ .tattach => |w| {
+ try testing.expectEqual(w.fid, have.tattach.fid);
+ try testing.expectEqual(w.afid, have.tattach.afid);
+ try testing.expectEqualStrings(w.uname, have.tattach.uname);
+ try testing.expectEqualStrings(w.aname, have.tattach.aname);
+ },
+ .rattach => |w| try testing.expectEqual(w.qid, have.rattach.qid),
+ .rerror => |w| try testing.expectEqualStrings(w.ename, have.rerror.ename),
+ .tflush => |w| try testing.expectEqual(w.oldtag, have.tflush.oldtag),
+ .rflush, .rclunk, .rremove, .rwstat => {},
+ .twalk => |w| {
+ try testing.expectEqual(w.fid, have.twalk.fid);
+ try testing.expectEqual(w.newfid, have.twalk.newfid);
+ try testing.expectEqual(w.nwname, have.twalk.nwname);
+ for (w.wname[0..w.nwname], have.twalk.wname[0..w.nwname]) |a, b|
+ try testing.expectEqualStrings(a, b);
+ },
+ .rwalk => |w| {
+ try testing.expectEqual(w.nwqid, have.rwalk.nwqid);
+ for (w.wqid[0..w.nwqid], have.rwalk.wqid[0..w.nwqid]) |a, b|
+ try testing.expectEqual(a, b);
+ },
+ .topen => |w| {
+ try testing.expectEqual(w.fid, have.topen.fid);
+ try testing.expectEqual(w.mode, have.topen.mode);
+ },
+ .ropen => |w| {
+ try testing.expectEqual(w.qid, have.ropen.qid);
+ try testing.expectEqual(w.iounit, have.ropen.iounit);
+ },
+ .tcreate => |w| {
+ try testing.expectEqual(w.fid, have.tcreate.fid);
+ try testing.expectEqualStrings(w.name, have.tcreate.name);
+ try testing.expectEqual(w.perm, have.tcreate.perm);
+ try testing.expectEqual(w.mode, have.tcreate.mode);
+ },
+ .rcreate => |w| {
+ try testing.expectEqual(w.qid, have.rcreate.qid);
+ try testing.expectEqual(w.iounit, have.rcreate.iounit);
+ },
+ .tread => |w| {
+ try testing.expectEqual(w.fid, have.tread.fid);
+ try testing.expectEqual(w.offset, have.tread.offset);
+ try testing.expectEqual(w.count, have.tread.count);
+ },
+ .rread => |w| try testing.expectEqualStrings(w.data, have.rread.data),
+ .twrite => |w| {
+ try testing.expectEqual(w.fid, have.twrite.fid);
+ try testing.expectEqual(w.offset, have.twrite.offset);
+ try testing.expectEqualStrings(w.data, have.twrite.data);
+ },
+ .rwrite => |w| try testing.expectEqual(w.count, have.rwrite.count),
+ .tclunk => |w| try testing.expectEqual(w.fid, have.tclunk.fid),
+ .tremove => |w| try testing.expectEqual(w.fid, have.tremove.fid),
+ .tstat => |w| try testing.expectEqual(w.fid, have.tstat.fid),
+ .rstat => |w| try expectStatEqual(w.stat, have.rstat.stat),
+ .twstat => |w| {
+ try testing.expectEqual(w.fid, have.twstat.fid);
+ try expectStatEqual(w.stat, have.twstat.stat);
+ },
+ }
+}
+
+const sample_qid: Qid = .{ .type = qtdir, .version = 3, .path = 0x0102_0304_0506_0708 };
+
+const sample_stat: Stat = .{
+ .type = 0,
+ .dev = 0,
+ .qid = sample_qid,
+ .mode = dmdir | 0o755,
+ .atime = 1,
+ .mtime = 2,
+ .length = 0,
+ .name = "body",
+ .uid = "goblin",
+ .gid = "goblin",
+ .muid = "goblin",
+};
+
+test "9p: the type numbers and their parity are the protocol's own" {
+ try testing.expectEqual(@as(u8, 100), @intFromEnum(Type.tversion));
+ try testing.expectEqual(@as(u8, 106), @intFromEnum(Type.terror));
+ try testing.expectEqual(@as(u8, 107), @intFromEnum(Type.rerror));
+ try testing.expectEqual(@as(u8, 126), @intFromEnum(Type.twstat));
+ try testing.expectEqual(@as(u8, 127), @intFromEnum(Type.rwstat));
+
+ try testing.expectEqual(@as(usize, 28), std.enums.values(Type).len);
+ for (std.enums.values(Type), 100..) |t, want| try testing.expectEqual(@as(u8, @intCast(want)), @intFromEnum(t));
+
+ try testing.expect(isT(.tversion));
+ try testing.expect(!isT(.rversion));
+ try testing.expect(isT(.twstat));
+ try testing.expect(!isT(.rwstat));
+
+ try testing.expectEqual(@as(u16, 0xFFFF), notag);
+ try testing.expectEqual(@as(u32, 0xFFFF_FFFF), nofid);
+ try testing.expectEqual(@as(usize, 16), max_welem);
+}
+
+test "9p: a qid is thirteen bytes" {
+ var buf: [32]u8 = undefined;
+ const bytes = try sample_qid.encode(&buf);
+ try testing.expectEqual(qid_len, bytes.len);
+ try testing.expectEqual(@as(usize, 13), bytes.len);
+ try testing.expectEqual(sample_qid, try Qid.decode(bytes));
+ try testing.expectError(error.Truncated, Qid.decode(bytes[0..12]));
+ try testing.expectError(error.NoSpace, sample_qid.encode(buf[0..12]));
+}
+
+test "9p: an encoded stat is size() + 2 bytes" {
+ var buf: [256]u8 = undefined;
+ const bytes = try sample_stat.encode(&buf);
+ const n = try sample_stat.size();
+ try testing.expectEqual(@as(usize, n) + 2, bytes.len);
+ try testing.expectEqual(@as(u16, 69), n);
+ try testing.expectEqual(stat_fixed - 2 + 22, n);
+ try testing.expectEqual(n, std.mem.readInt(u16, bytes[0..2], .little));
+ try expectStatEqual(sample_stat, try Stat.decode(bytes));
+
+ const bare: Stat = .{
+ .type = 0,
+ .dev = 0,
+ .qid = .{ .type = qtfile, .version = 0, .path = 0 },
+ .mode = 0,
+ .atime = 0,
+ .mtime = 0,
+ .length = 0,
+ .name = "",
+ .uid = "",
+ .gid = "",
+ .muid = "",
+ };
+ try testing.expectEqual(@as(u16, 47), try bare.size());
+ try testing.expectEqual(@as(usize, 49), (try bare.encode(&buf)).len);
+}
+
+test "9p: every message round-trips" {
+ var buf: [512]u8 = undefined;
+
+ _ = try roundTrip(&buf, notag, .{ .tversion = .{ .msize = 8192, .version = "9P2000" } });
+ _ = try roundTrip(&buf, notag, .{ .rversion = .{ .msize = 8192, .version = "9P2000" } });
+ _ = try roundTrip(&buf, notag, .{ .rversion = .{ .msize = test_msize, .version = "unknown" } });
+ _ = try roundTrip(&buf, 1, .{ .tauth = .{ .afid = 1, .uname = "goblin", .aname = "" } });
+ _ = try roundTrip(&buf, 1, .{ .rauth = .{ .aqid = .{ .type = qtauth, .version = 0, .path = 9 } } });
+ _ = try roundTrip(&buf, 2, .{ .tattach = .{ .fid = 0, .afid = nofid, .uname = "goblin", .aname = "" } });
+ _ = try roundTrip(&buf, 2, .{ .rattach = .{ .qid = sample_qid } });
+ _ = try roundTrip(&buf, 3, .{ .rerror = .{ .ename = "no such file" } });
+ _ = try roundTrip(&buf, 4, .{ .tflush = .{ .oldtag = 3 } });
+ _ = try roundTrip(&buf, 4, .rflush);
+ _ = try roundTrip(&buf, 5, .{ .twalk = .{ .fid = 0, .newfid = 1, .nwname = 2, .wname = .{ "7", "body" } ++ @as([max_welem - 2][]const u8, @splat("")) } });
+ _ = try roundTrip(&buf, 5, .{ .rwalk = .{ .nwqid = 2, .wqid = .{ sample_qid, sample_qid } ++ @as([max_welem - 2]Qid, @splat(sample_qid)) } });
+ _ = try roundTrip(&buf, 6, .{ .topen = .{ .fid = 1, .mode = 0 } });
+ _ = try roundTrip(&buf, 6, .{ .ropen = .{ .qid = sample_qid, .iounit = 8192 - iohdrsz } });
+ _ = try roundTrip(&buf, 7, .{ .tcreate = .{ .fid = 1, .name = "new", .perm = dmdir | 0o777, .mode = 2 } });
+ _ = try roundTrip(&buf, 7, .{ .rcreate = .{ .qid = sample_qid, .iounit = 0 } });
+ _ = try roundTrip(&buf, 8, .{ .tread = .{ .fid = 1, .offset = 0xdead_beef_cafe, .count = 4096 } });
+ _ = try roundTrip(&buf, 8, .{ .rread = .{ .data = "hello" } });
+ _ = try roundTrip(&buf, 8, .{ .rread = .{ .data = "" } });
+ _ = try roundTrip(&buf, 9, .{ .twrite = .{ .fid = 1, .offset = 0, .data = "Edit ,d" } });
+ _ = try roundTrip(&buf, 9, .{ .twrite = .{ .fid = 1, .offset = 0, .data = "" } });
+ _ = try roundTrip(&buf, 9, .{ .rwrite = .{ .count = 7 } });
+ _ = try roundTrip(&buf, 10, .{ .tclunk = .{ .fid = 1 } });
+ _ = try roundTrip(&buf, 10, .rclunk);
+ _ = try roundTrip(&buf, 11, .{ .tremove = .{ .fid = 1 } });
+ _ = try roundTrip(&buf, 11, .rremove);
+ _ = try roundTrip(&buf, 12, .{ .tstat = .{ .fid = 1 } });
+ _ = try roundTrip(&buf, 12, .{ .rstat = .{ .stat = sample_stat } });
+ _ = try roundTrip(&buf, 13, .{ .twstat = .{ .fid = 1, .stat = sample_stat } });
+ _ = try roundTrip(&buf, 13, .rwstat);
+
+ try testing.expectEqual(@as(usize, 27), @typeInfo(Msg).@"union".fields.len);
+ try testing.expectEqual(std.enums.values(Type).len - 1, @typeInfo(Msg).@"union".fields.len);
+}
+
+test "9p: empty and maximum-length strings survive the trip" {
+ var buf: [70_000]u8 = undefined;
+
+ const empty = try roundTrip(&buf, 1, .{ .tattach = .{ .fid = 0, .afid = nofid, .uname = "", .aname = "" } });
+ try testing.expectEqual(@as(usize, 0), empty.tattach.uname.len);
+ try testing.expectEqual(@as(usize, header_len + 4 + 4 + 2 + 2), (try encode(empty, 1, &buf)).len);
+
+ var big: [65_536]u8 = undefined;
+ @memset(&big, 'x');
+ const max = big[0..std.math.maxInt(u16)];
+ const got = try roundTrip(&buf, 1, .{ .rerror = .{ .ename = max } });
+ try testing.expectEqual(@as(usize, 65_535), got.rerror.ename.len);
+ try testing.expectError(error.Overlong, encode(.{ .rerror = .{ .ename = &big } }, 1, &buf));
+
+ var wide = sample_stat;
+ wide.name = max;
+ try testing.expectError(error.Overlong, wide.size());
+ try testing.expectError(error.Overlong, encode(.{ .rstat = .{ .stat = wide } }, 1, &buf));
+}
+
+test "9p: Twalk carries 0, 1 and 16 elements and refuses 17" {
+ var buf: [512]u8 = undefined;
+
+ const zero = try roundTrip(&buf, 1, .{ .twalk = .{ .fid = 0, .newfid = 1, .nwname = 0 } });
+ try testing.expectEqual(@as(u16, 0), zero.twalk.nwname);
+ try testing.expectEqual(@as(usize, header_len + 4 + 4 + 2), (try encode(zero, 1, &buf)).len);
+
+ _ = try roundTrip(&buf, 1, .{ .twalk = .{
+ .fid = 0,
+ .newfid = 1,
+ .nwname = 1,
+ .wname = .{"body"} ++ @as([max_welem - 1][]const u8, @splat("")),
+ } });
+
+ const names: [max_welem][]const u8 = .{ "a", "b", "c", "d", "e", "f", "g", "h", "i", "j", "k", "l", "m", "n", "o", "p" };
+ const full = try roundTrip(&buf, 1, .{ .twalk = .{ .fid = 0, .newfid = 1, .nwname = max_welem, .wname = names } });
+ try testing.expectEqual(@as(u16, 16), full.twalk.nwname);
+ for (names, full.twalk.wname[0..max_welem]) |a, b| try testing.expectEqualStrings(a, b);
+ _ = try roundTrip(&buf, 1, .{ .rwalk = .{ .nwqid = max_welem, .wqid = @splat(sample_qid) } });
+
+ var raw: [256]u8 = undefined;
+ const bad = blk: {
+ var w: Writer = .init(&raw);
+ try w.putU32(0); // patched below
+ try w.putByte(@intFromEnum(Type.twalk));
+ try w.putU16(1);
+ try w.putU32(0);
+ try w.putU32(1);
+ try w.putU16(17);
+ for (0..17) |i| try w.putString(&[_]u8{@intCast('a' + i)});
+ std.mem.writeInt(u32, raw[0..4], @intCast(w.n), .little);
+ break :blk raw[0..w.n];
+ };
+ try testing.expectEqual(@as(usize, header_len + 4 + 4 + 2 + 17 * 3), bad.len);
+ try testing.expectError(error.Overlong, decode(bad));
+
+ const bad_r = blk: {
+ var w: Writer = .init(&raw);
+ try w.putU32(0);
+ try w.putByte(@intFromEnum(Type.rwalk));
+ try w.putU16(1);
+ try w.putU16(17);
+ for (0..17) |_| try w.putQid(sample_qid);
+ std.mem.writeInt(u32, raw[0..4], @intCast(w.n), .little);
+ break :blk raw[0..w.n];
+ };
+ try testing.expectError(error.Overlong, decode(bad_r));
+}
+
+test "9p: the stat double length" {
+ var buf: [512]u8 = undefined;
+
+ var good: [512]u8 = undefined;
+ const n = blk: {
+ const bytes = try encode(.{ .rstat = .{ .stat = sample_stat } }, 1, &buf);
+ @memcpy(good[0..bytes.len], bytes);
+ break :blk bytes.len;
+ };
+ const inner = try sample_stat.size();
+ try testing.expectEqual(inner + 2, std.mem.readInt(u16, good[header_len..][0..2], .little));
+ try testing.expectEqual(inner, std.mem.readInt(u16, good[header_len + 2 ..][0..2], .little));
+ try testing.expectEqual(header_len + 2 + @as(usize, inner) + 2, n);
+
+ const w_bytes = try encode(.{ .twstat = .{ .fid = 7, .stat = sample_stat } }, 1, &buf);
+ try testing.expectEqual(inner + 2, std.mem.readInt(u16, w_bytes[header_len + 4 ..][0..2], .little));
+ try testing.expectEqual(inner, std.mem.readInt(u16, w_bytes[header_len + 6 ..][0..2], .little));
+
+ var off: [512]u8 = undefined;
+
+ @memcpy(off[0..n], good[0..n]);
+ std.mem.writeInt(u16, off[header_len..][0..2], inner, .little);
+ try testing.expectError(error.Truncated, decode(off[0..n]));
+
+ @memcpy(off[0..n], good[0..n]);
+ std.mem.writeInt(u16, off[header_len..][0..2], inner + 4, .little);
+ try testing.expectError(error.Truncated, decode(off[0..n]));
+
+ @memcpy(off[0..n], good[0..n]);
+ std.mem.writeInt(u16, off[header_len + 2 ..][0..2], inner + 2, .little);
+ try testing.expectError(error.Truncated, decode(off[0..n]));
+
+ @memcpy(off[0..n], good[0..n]);
+ std.mem.writeInt(u16, off[header_len + 2 ..][0..2], inner - 1, .little);
+ try testing.expectError(error.Trailing, decode(off[0..n]));
+}
+
+fn expectTruncatedAtEveryBoundary(full: []const u8) !void {
+ var scratch: [1024]u8 = undefined;
+ var n: usize = 0;
+ while (n < full.len) : (n += 1) {
+ try testing.expectError(error.Truncated, decode(full[0..n]));
+ if (n < header_len) continue;
+ @memcpy(scratch[0..n], full[0..n]);
+ std.mem.writeInt(u32, scratch[0..4], @intCast(n), .little);
+ try testing.expectError(error.Truncated, decode(scratch[0..n]));
+ }
+ _ = try decode(full);
+}
+
+test "9p: truncation at every field boundary is refused" {
+ var buf: [512]u8 = undefined;
+
+ try expectTruncatedAtEveryBoundary(try encode(
+ .{ .tversion = .{ .msize = 8192, .version = "9P2000" } },
+ notag,
+ &buf,
+ ));
+ try expectTruncatedAtEveryBoundary(try encode(.{ .twalk = .{
+ .fid = 1,
+ .newfid = 2,
+ .nwname = 3,
+ .wname = .{ "usr", "", "bin" } ++ @as([max_welem - 3][]const u8, @splat("")),
+ } }, 1, &buf));
+ try expectTruncatedAtEveryBoundary(try encode(
+ .{ .tread = .{ .fid = 1, .offset = 0x0102_0304_0506_0708, .count = 8168 } },
+ 1,
+ &buf,
+ ));
+ try expectTruncatedAtEveryBoundary(try encode(.{ .rstat = .{ .stat = sample_stat } }, 1, &buf));
+ try expectTruncatedAtEveryBoundary(try encode(.{ .rread = .{ .data = "12345678" } }, 1, &buf));
+ try expectTruncatedAtEveryBoundary(try encode(
+ .{ .rwalk = .{ .nwqid = 3, .wqid = @splat(sample_qid) } },
+ 1,
+ &buf,
+ ));
+ try expectTruncatedAtEveryBoundary(try encode(.{ .twstat = .{ .fid = 1, .stat = sample_stat } }, 1, &buf));
+}
+
+test "9p: a size field that disagrees with the buffer is refused" {
+ var buf: [512]u8 = undefined;
+ const bytes = try encode(.{ .tclunk = .{ .fid = 1 } }, 1, &buf);
+ try testing.expectEqual(@as(usize, 11), bytes.len);
+
+ var raw: [64]u8 = undefined;
+ @memcpy(raw[0..bytes.len], bytes);
+
+ for ([_]u32{ 12, 13, 64, 1 << 20, std.math.maxInt(u32) }) |claim| {
+ std.mem.writeInt(u32, raw[0..4], claim, .little);
+ try testing.expectError(error.Truncated, decode(raw[0..bytes.len]));
+ }
+
+ std.mem.writeInt(u32, raw[0..4], 10, .little);
+ try testing.expectError(error.Trailing, decode(raw[0..bytes.len]));
+
+ for ([_]u32{ 0, 1, 6 }) |claim| {
+ std.mem.writeInt(u32, raw[0..4], claim, .little);
+ try testing.expectError(error.BadValue, decode(raw[0..bytes.len]));
+ try testing.expectError(error.BadValue, decode(raw[0..header_len]));
+ }
+}
+
+test "9p: an unknown or illegal type byte is refused" {
+ var buf: [512]u8 = undefined;
+ const bytes = try encode(.{ .tclunk = .{ .fid = 1 } }, 1, &buf);
+ var raw: [64]u8 = undefined;
+ @memcpy(raw[0..bytes.len], bytes);
+
+ for ([_]u8{ 0, 1, 8, 12, 99, 106, 128, 255 }) |t| {
+ raw[4] = t;
+ try testing.expectError(error.BadTag, decode(raw[0..bytes.len]));
+ }
+
+ var t: u16 = 0;
+ while (t <= 255) : (t += 1) {
+ raw[4] = @intCast(t);
+ const defined = t >= 100 and t <= 127 and t != @intFromEnum(Type.terror);
+ if (decode(raw[0..bytes.len])) |got| {
+ try testing.expectEqual(@as(u8, @intCast(t)), @intFromEnum(got.msg.msgType()));
+ try testing.expect(t == @intFromEnum(Type.tclunk) or
+ t == @intFromEnum(Type.tremove) or
+ t == @intFromEnum(Type.tstat) or
+ t == @intFromEnum(Type.rwrite));
+ } else |err| {
+ if (!defined) try testing.expectEqual(Error.BadTag, err);
+ }
+ }
+}
+
+test "9p: trailing bytes inside the size are refused" {
+ var raw: [64]u8 = undefined;
+
+ var w: Writer = .init(&raw);
+ try w.putU32(12);
+ try w.putByte(@intFromEnum(Type.tclunk));
+ try w.putU16(1);
+ try w.putU32(7);
+ try w.putByte(0xAA);
+ try testing.expectEqual(@as(usize, 12), w.n);
+ try testing.expectError(error.Trailing, decode(raw[0..12]));
+
+ w = .init(&raw);
+ try w.putU32(header_len + 2 + 2);
+ try w.putByte(@intFromEnum(Type.tflush));
+ try w.putU16(1);
+ try w.putU16(3);
+ try w.putU16(3);
+ try testing.expectError(error.Trailing, decode(raw[0..w.n]));
+}
+
+test "9p: frameLen needs four bytes" {
+ var buf: [512]u8 = undefined;
+ const bytes = try encode(.{ .tread = .{ .fid = 1, .offset = 0, .count = 8168 } }, 1, &buf);
+ try testing.expectEqual(@as(usize, 23), bytes.len);
+
+ for (0..4) |n| try testing.expectEqual(@as(?u32, null), frameLen(bytes[0..n]));
+ try testing.expectEqual(@as(?u32, 23), frameLen(bytes[0..4]));
+ try testing.expectEqual(@as(?u32, 23), frameLen(bytes));
+
+ var raw: [4]u8 = .{ 0xFF, 0xFF, 0xFF, 0xFF };
+ try testing.expectEqual(@as(?u32, std.math.maxInt(u32)), frameLen(&raw));
+ raw = .{ 0, 0, 0, 0 };
+ try testing.expectEqual(@as(?u32, 0), frameLen(&raw));
+}
+
+test "9p: encode refuses a short buffer and writes nothing" {
+ var buf: [512]u8 = undefined;
+ const want = (try encode(.{ .rstat = .{ .stat = sample_stat } }, 1, &buf)).len;
+
+ var n: usize = 0;
+ while (n < want) : (n += 1) {
+ var scratch: [512]u8 = @splat(0xAA);
+ try testing.expectError(error.NoSpace, encode(.{ .rstat = .{ .stat = sample_stat } }, 1, scratch[0..n]));
+ for (scratch) |b| try testing.expectEqual(@as(u8, 0xAA), b);
+ }
+
+ var exact: [512]u8 = @splat(0xAA);
+ try testing.expectEqual(want, (try encode(.{ .rstat = .{ .stat = sample_stat } }, 1, exact[0..want])).len);
+ try testing.expectEqual(@as(u8, 0xAA), exact[want]);
+}
+
+test "9p: byte for byte against u9fs convS2M" {
+ var buf: [512]u8 = undefined;
+
+ try testing.expectEqualSlices(u8, &.{
+ 0x13, 0x00, 0x00, 0x00, // size = 19
+ 0x64, // Tversion = 100
+ 0xff, 0xff, // NOTAG
+ 0x00, 0x20, 0x00, 0x00, // msize = 8192
+ 0x06, 0x00, // n = 6
+ '9', 'P',
+ '2', '0',
+ '0', '0',
+ }, try encode(.{ .tversion = .{ .msize = 8192, .version = "9P2000" } }, notag, &buf));
+
+ try testing.expectEqualSlices(u8, &.{
+ 0x1b, 0x00, 0x00, 0x00, // size = 27
+ 0x6e, // Twalk = 110
+ 0x01, 0x00, // tag = 1
+ 0x01, 0x00, 0x00, 0x00, // fid = 1
+ 0x02, 0x00, 0x00, 0x00, // newfid = 2
+ 0x02, 0x00, // nwname = 2
+ 0x03, 0x00,
+ 'u', 's',
+ 'r', 0x03,
+ 0x00, 'b',
+ 'i', 'n',
+ }, try encode(.{ .twalk = .{
+ .fid = 1,
+ .newfid = 2,
+ .nwname = 2,
+ .wname = .{ "usr", "bin" } ++ @as([max_welem - 2][]const u8, @splat("")),
+ } }, 1, &buf));
+
+ try testing.expectEqualSlices(u8, &.{
+ 0x0e, 0x00, 0x00, 0x00, // size = 14
+ 0x75, // Rread = 117
+ 0x09, 0x00, // tag = 9
+ 0x03, 0x00, 0x00, 0x00, // count = 3
+ 'a', 'b', 'c',
+ }, try encode(.{ .rread = .{ .data = "abc" } }, 9, &buf));
+
+ const one: Stat = .{
+ .type = 0,
+ .dev = 0,
+ .qid = .{ .type = qtdir, .version = 1, .path = 2 },
+ .mode = dmdir | 0o755,
+ .atime = 3,
+ .mtime = 4,
+ .length = 0,
+ .name = "a",
+ .uid = "u",
+ .gid = "g",
+ .muid = "m",
+ };
+ try testing.expectEqual(@as(u16, 51), try one.size());
+ try testing.expectEqualSlices(u8, &.{
+ 0x3e, 0x00, 0x00, 0x00, // size = 62
+ 0x7d, // Rstat = 125
+ 0x07, 0x00, // tag = 7
+ 0x35, 0x00, // OUTER count = 53 = 51 + 2
+ 0x33, 0x00, // stat size = 51, excluding these two
+ 0x00, 0x00, // type
+ 0x00, 0x00, 0x00, 0x00, // dev
+ 0x80, // qid.type = QTDIR
+ 0x01, 0x00, 0x00, 0x00, // qid.version = 1
+ 0x02, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // qid.path = 2
+ 0xed, 0x01, 0x00, 0x80, // mode = DMDIR | 0755
+ 0x03, 0x00, 0x00, 0x00, // atime
+ 0x04, 0x00, 0x00, 0x00, // mtime
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, // length
+ 0x01, 0x00, 'a', // name
+ 0x01, 0x00, 'u', // uid
+ 0x01, 0x00, 'g', // gid
+ 0x01, 0x00, 'm', // muid
+ }, try encode(.{ .rstat = .{ .stat = one } }, 7, &buf));
+
+ try testing.expectEqual(qtdir, @as(u8, @intCast(dmdir >> 24)));
+ try testing.expectEqual(qtappend, @as(u8, @intCast(dmappend >> 24)));
+ try testing.expectEqual(qtexcl, @as(u8, @intCast(dmexcl >> 24)));
+ try testing.expectEqual(qtauth, @as(u8, @intCast(dmauth >> 24)));
+ try testing.expectEqual(qttmp, @as(u8, @intCast(dmtmp >> 24)));
+ try testing.expectEqual(@as(u32, 0o777), dmperm);
+}
diff --git a/test/differential/go.mod b/test/differential/go.mod
new file mode 100644
index 0000000..43f6ef2
--- /dev/null
+++ b/test/differential/go.mod
@@ -0,0 +1,14 @@
+module cloud9-differential
+
+go 1.26.5
+
+require (
+ 9fans.net/go v0.0.7
+ github.com/knusbaum/go9p v1.18.0
+)
+
+require (
+ github.com/Plan9-Archive/libauth v0.0.0-20180917063427-d1ca9e94969d // indirect
+ github.com/emersion/go-sasl v0.0.0-20200509203442-7bfe0ed36a21 // indirect
+ github.com/fhs/mux9p v0.3.1 // indirect
+)
diff --git a/test/differential/go.sum b/test/differential/go.sum
new file mode 100644
index 0000000..20c7aea
--- /dev/null
+++ b/test/differential/go.sum
@@ -0,0 +1,58 @@
+9fans.net/go v0.0.2/go.mod h1:lfPdxjq9v8pVQXUMBCx5EO5oLXWQFlKRQgs1kEkjoIM=
+9fans.net/go v0.0.7 h1:H5CsYJTf99C8EYAQr+uSoEJnLP/iZU8RmDuhyk30iSM=
+9fans.net/go v0.0.7/go.mod h1:Rxvbbc1e+1TyGMjAvLthGTyO97t+6JMQ6ly+Lcs9Uf0=
+dmitri.shuralyov.com/gpu/mtl v0.0.0-20201218220906-28db891af037/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU=
+github.com/BurntSushi/xgb v0.0.0-20160522181843-27f122750802/go.mod h1:IVnqGOEym/WlBOVXweHU+Q+/VP0lqqI8lqeDx9IjBqo=
+github.com/Plan9-Archive/libauth v0.0.0-20180917063427-d1ca9e94969d h1:xH/U6K+HYxh1480TkQYRqRO8F2RJsg+R6wFiVJzdldg=
+github.com/Plan9-Archive/libauth v0.0.0-20180917063427-d1ca9e94969d/go.mod h1:UKp8dv9aeaZoQFWin7eQXtz89iHly1YAFZNn3MCutmQ=
+github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8=
+github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
+github.com/emersion/go-sasl v0.0.0-20200509203442-7bfe0ed36a21 h1:OJyUGMJTzHTd1XQp98QTaHernxMYzRaOasRir9hUlFQ=
+github.com/emersion/go-sasl v0.0.0-20200509203442-7bfe0ed36a21/go.mod h1:iL2twTeMvZnrg54ZoPDNfJaJaqy0xIQFuBdrLsmspwQ=
+github.com/fhs/mux9p v0.3.1 h1:x1UswUWZoA9vrA02jfisndCq3xQm+wrQUxUt5N99E08=
+github.com/fhs/mux9p v0.3.1/go.mod h1:F4hwdenmit0WDoNVT2VMWlLJrBVCp/8UhzJa7scfjEQ=
+github.com/go-gl/glfw/v3.3/glfw v0.0.0-20200222043503-6f7a984d4dc4/go.mod h1:tQ2UAYgL5IevRw8kRxooKSPJfGvJ9fJQFa0TUsXzTg8=
+github.com/hanwen/go-fuse v1.0.0/go.mod h1:unqXarDXqzAk0rt98O2tVndEPIpUgLD9+rwFisZH3Ok=
+github.com/hanwen/go-fuse/v2 v2.0.3/go.mod h1:0EQM6aH2ctVpvZ6a+onrQ/vaykxh2GH7hy3e13vzTUY=
+github.com/knusbaum/go9p v1.18.0 h1:/Y67RNvNKX1ZV1IOdnO1lIetiF0X+CumOyvEc0011GI=
+github.com/knusbaum/go9p v1.18.0/go.mod h1:HtMoJKqZUe1Oqag5uJqG5RKQ9gWPSP+wolsnLLv44r8=
+github.com/kylelemons/godebug v0.0.0-20170820004349-d65d576e9348/go.mod h1:B69LEHPfb2qLo0BaaOLcbitczOKLWTsrBG9LczfCD4k=
+github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
+github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
+github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
+github.com/stretchr/testify v1.4.0 h1:2E4SXV/wtOkTonXsotYi4li6zVWxYlZuYNCXe9XRJyk=
+github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
+golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
+golang.org/x/crypto v0.0.0-20190510104115-cbcb75029529/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
+golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
+golang.org/x/exp v0.0.0-20190731235908-ec7cb31e5a56/go.mod h1:JhuoJpWY28nO4Vef9tZUw9qufEGTyX1+7lmHxV5q5G4=
+golang.org/x/exp v0.0.0-20210405174845-4513512abef3/go.mod h1:I6l2HNBLBZEcrOoCpyKLdY2lHoRZ8lI4x60KMCQDft4=
+golang.org/x/image v0.0.0-20190227222117-0694c2d4d067/go.mod h1:kZ7UVZpmo3dzQBMxlp+ypCbDeSB+sBbTgSJuh5dn5js=
+golang.org/x/image v0.0.0-20190802002840-cff245a6509b/go.mod h1:FeLwcggjj3mMvU+oOTbSwawSJRM1uh48EjtB4UJZlP0=
+golang.org/x/mobile v0.0.0-20190312151609-d3739f865fa6/go.mod h1:z+o9i4GpDbdi3rU15maQ/Ox0txvL9dWGYEHz965HBQE=
+golang.org/x/mobile v0.0.0-20201217150744-e6ae53a27f4f/go.mod h1:skQtrUTUwhdJvXM/2KKJzY8pDgNr9I/FOMqDVRPBUS4=
+golang.org/x/mobile v0.0.0-20210220033013-bdb1ca9a1e08/go.mod h1:skQtrUTUwhdJvXM/2KKJzY8pDgNr9I/FOMqDVRPBUS4=
+golang.org/x/mod v0.1.0/go.mod h1:0QHyrYULN0/3qlju5TqG8bIK38QM8yzMo5ekMj3DlcY=
+golang.org/x/mod v0.1.1-0.20191105210325-c90efee705ee/go.mod h1:QqPTAvyqsEbceGzBzNggFXnrqF1CaUcvgkdR5Ot7KZg=
+golang.org/x/mod v0.1.1-0.20191209134235-331c550502dd/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
+golang.org/x/mod v0.3.1-0.20200828183125-ce943fd02449/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
+golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
+golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
+golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
+golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
+golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
+golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
+golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/sys v0.0.0-20191001151750-bb3f8db39f24/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/sys v0.0.0-20201020230747-6e5568b54d1a/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/sys v0.0.0-20210415045647-66c3f260301c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
+golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
+golang.org/x/tools v0.0.0-20190312151545-0bb0c0a6e846/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs=
+golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
+golang.org/x/tools v0.0.0-20200117012304-6edc0a871e69/go.mod h1:TB2adYChydJhpapKDTa4BR/hXlZSLoq2Wpct/0txZ28=
+golang.org/x/tools v0.0.0-20200207183749-b753a1ba74fa/go.mod h1:TB2adYChydJhpapKDTa4BR/hXlZSLoq2Wpct/0txZ28=
+golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
+golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
+gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
+gopkg.in/yaml.v2 v2.2.2 h1:ZCJp+EgiOT7lHqUV2J862kp8Qj64Jo6az82+3Td9dZw=
+gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
diff --git a/test/differential/main.go b/test/differential/main.go
new file mode 100644
index 0000000..7a2538d
--- /dev/null
+++ b/test/differential/main.go
@@ -0,0 +1,172 @@
+// Valid-traffic differential conformance and measurement runner. No network targets.
+package main
+
+import (
+ p9 "9fans.net/go/plan9"
+ "bytes"
+ "encoding/binary"
+ "encoding/json"
+ "flag"
+ "fmt"
+ other "github.com/knusbaum/go9p/proto"
+ "io"
+ "math/rand"
+ "os"
+ "os/exec"
+ "path/filepath"
+ "runtime"
+ "time"
+)
+
+type reader struct {
+ data []byte
+ chunk, calls int
+}
+
+func (r *reader) Read(p []byte) (int, error) {
+ r.calls++
+ if len(r.data) == 0 {
+ return 0, io.EOF
+ }
+ n := min(len(p), len(r.data), r.chunk)
+ copy(p, r.data[:n])
+ r.data = r.data[n:]
+ return n, nil
+}
+func check(err error) {
+ if err != nil {
+ panic(err)
+ }
+}
+func marshal(v any) []byte { b, e := json.MarshalIndent(v, "", " "); check(e); return append(b, '\n') }
+func corpus(seed int64, rounds int) [][]byte {
+ r := rand.New(rand.NewSource(seed))
+ frames := [][]byte{}
+ names := []string{"a", "file", "directory", "日本語", "ação"}
+ for i := 0; i < rounds; i++ {
+ q := p9.Qid{Type: 0, Vers: r.Uint32(), Path: r.Uint64()}
+ d := p9.Dir{Qid: q, Mode: 0600, Atime: r.Uint32(), Mtime: r.Uint32(), Length: r.Uint64(), Name: names[i%len(names)], Uid: "user", Gid: "group", Muid: "user"}
+ stat, e := d.Bytes()
+ check(e)
+ if i%2 == 1 {
+ d.Null()
+ stat, e = d.Bytes()
+ check(e)
+ }
+ data := make([]byte, []int{0, 1, 7, 127, 1024, 8192}[i%6])
+ _, e = r.Read(data)
+ check(e)
+ wn := make([]string, i%17)
+ wq := make([]p9.Qid, i%17)
+ for j := range wn {
+ wn[j] = names[r.Intn(len(names))]
+ wq[j] = q
+ }
+ for typ := uint8(100); typ <= 127; typ++ {
+ if typ == 106 {
+ continue
+ }
+ f := p9.Fcall{Type: typ, Tag: uint16(r.Intn(65535)), Fid: r.Uint32(), Newfid: r.Uint32(), Afid: p9.NOFID, Uname: "user", Aname: "", Version: "9P2000", Msize: 65536, Oldtag: uint16(r.Intn(65535)), Ename: "permission denied", Qid: q, Aqid: p9.Qid{Type: p9.QTAUTH, Vers: q.Vers, Path: q.Path}, Iounit: 0, Name: d.Name, Perm: 0600, Mode: uint8(i % 4), Offset: r.Uint64(), Count: uint32(len(data)), Data: data, Wname: wn, Wqid: wq, Stat: stat}
+ if typ == 100 || typ == 101 {
+ f.Tag = p9.NOTAG
+ }
+ if typ == 114 {
+ f.Name = names[i%len(names)]
+ }
+ b, e := f.Bytes()
+ check(e)
+ frames = append(frames, b)
+ }
+ }
+ return frames
+}
+
+type metric struct {
+ Implementation string `json:"implementation"`
+ Operations uint64 `json:"operations"`
+ CodecNS int64 `json:"codec_ns"`
+ Allocations uint64 `json:"codec_allocations"`
+ AllocatedBytes uint64 `json:"allocated_bytes"`
+ ReadCalls uint64 `json:"read_calls"`
+ Bytes uint64 `json:"bytes"`
+}
+
+func measure(name string, frames [][]byte, repeats, chunk int) metric {
+ m := metric{Implementation: name}
+ var before, after runtime.MemStats
+ runtime.GC()
+ runtime.ReadMemStats(&before)
+ for index, frame := range frames {
+ start := time.Now()
+ for j := 0; j < repeats; j++ {
+ rd := reader{data: frame, chunk: chunk}
+ var out []byte
+ if name == "9fans v0.0.7" {
+ f, e := p9.ReadFcall(&rd)
+ check(e)
+ out, e = f.Bytes()
+ check(e)
+ } else {
+ f, e := other.ParseCall(&rd)
+ check(e)
+ out = f.Compose()
+ }
+ if !bytes.Equal(frame, out) {
+ panic(fmt.Sprintf("%s frame %d type %d differs", name, index, frame[4]))
+ }
+ m.ReadCalls += uint64(rd.calls)
+ m.Bytes += uint64(len(frame))
+ m.Operations++
+ }
+ m.CodecNS += time.Since(start).Nanoseconds()
+ }
+ runtime.ReadMemStats(&after)
+ m.Allocations = after.Mallocs - before.Mallocs
+ m.AllocatedBytes = after.TotalAlloc - before.TotalAlloc
+ return m
+}
+func main() {
+ seed := flag.Int64("seed", 4200, "deterministic valid-traffic seed")
+ rounds := flag.Int("rounds", 100, "27 messages per round")
+ repeats := flag.Int("repeats", 20, "codec repetitions")
+ chunk := flag.Int("chunk", 65536, "maximum read/write chunk; use 1 to fragment every byte")
+ session := flag.String("session-probe", "", "run live sessions against go9p")
+ probe := flag.String("probe", "../../zig-out/bin/cloud9-probe", "cloud9 executable")
+ output := flag.String("output", "results", "artifact directory")
+ flag.Parse()
+ if *rounds < 1 || *rounds > 10000 || *repeats < 1 || *repeats > 10000 || *chunk < 1 {
+ panic("invalid bounds")
+ }
+ frames := corpus(*seed, *rounds)
+ all := bytes.Join(frames, nil)
+ check(os.MkdirAll(*output, 0755))
+ check(os.WriteFile(filepath.Join(*output, "corpus.9p"), all, 0644))
+ config := map[string]any{"seed": *seed, "rounds": *rounds, "repeats": *repeats, "chunk": *chunk, "frames": len(frames), "go": runtime.Version(), "platform": runtime.GOOS + "/" + runtime.GOARCH}
+ check(os.WriteFile(filepath.Join(*output, "config.json"), marshal(config), 0644))
+ metrics := []metric{measure("9fans v0.0.7", frames, *repeats, *chunk), measure("go9p v1.18.0", frames, *repeats, *chunk)}
+ var input bytes.Buffer
+ for _, v := range []uint32{uint32(len(frames)), uint32(*repeats), uint32(*chunk)} {
+ check(binary.Write(&input, binary.LittleEndian, v))
+ }
+ input.Write(all)
+ command := exec.Command(*probe)
+ command.Stdin = &input
+ var stderr bytes.Buffer
+ command.Stderr = &stderr
+ out, e := command.Output()
+ if e != nil {
+ panic(fmt.Sprintf("cloud9: %v: %s", e, stderr.String()))
+ }
+ if !bytes.Equal(all, out) {
+ panic("cloud9 round-trip differs; replay corpus.9p with saved configuration")
+ }
+ var cloud map[string]any
+ check(json.Unmarshal(stderr.Bytes(), &cloud))
+ report := map[string]any{"configuration": config, "correctness": "all generated base-9P2000 frames agree byte-for-byte", "cloud9": cloud, "references": metrics, "measurement_notes": []string{"Cloud9 codec timer excludes pipe I/O; Go timers include in-memory stream parsing and byte comparison. Do not treat these as equivalent throughput benchmarks.", "Go allocations use runtime.MemStats deltas on a single-process run; cloud9 codec has no allocation API.", "Cloud9 I/O counters are actual stdin/stdout syscalls; Go read counters are io.Reader calls, not syscalls.", "These are wire-conformance probes, not differential filesystem-server semantics."}}
+ if *session != "" {
+ report["session"] = sessionProbe(*session, *seed, *rounds)
+ }
+ b := marshal(report)
+ check(os.WriteFile(filepath.Join(*output, "report.json"), b, 0644))
+ os.Stdout.Write(b)
+}
diff --git a/test/differential/probe.zig b/test/differential/probe.zig
new file mode 100644
index 0000000..ec7c88e
--- /dev/null
+++ b/test/differential/probe.zig
@@ -0,0 +1,59 @@
+const std = @import("std");
+const c9 = @import("cloud9");
+const libc = std.c;
+var reads: u64 = 0;
+var writes: u64 = 0;
+var bytes_in: u64 = 0;
+var bytes_out: u64 = 0;
+
+fn readAll(bytes: []u8, chunk: usize) !void {
+ var offset: usize = 0;
+ while (offset < bytes.len) {
+ reads += 1;
+ const n = libc.read(0, bytes[offset..].ptr, @min(chunk, bytes.len - offset));
+ if (n <= 0) return error.Input;
+ offset += @intCast(n);
+ bytes_in += @intCast(n);
+ }
+}
+fn writeAll(bytes: []const u8, chunk: usize) !void {
+ var offset: usize = 0;
+ while (offset < bytes.len) {
+ writes += 1;
+ const n = libc.write(1, bytes[offset..].ptr, @min(chunk, bytes.len - offset));
+ if (n <= 0) return error.Output;
+ offset += @intCast(n);
+ bytes_out += @intCast(n);
+ }
+}
+fn nowNs() u64 {
+ var ts: libc.timespec = undefined;
+ std.debug.assert(libc.clock_gettime(.MONOTONIC, &ts) == 0);
+ return @as(u64, @intCast(ts.sec)) * std.time.ns_per_s + @as(u64, @intCast(ts.nsec));
+}
+pub fn main() !void {
+ var header: [12]u8 = undefined;
+ try readAll(&header, header.len);
+ const count = std.mem.readInt(u32, header[0..4], .little);
+ const repeats = std.mem.readInt(u32, header[4..8], .little);
+ const chunk = std.mem.readInt(u32, header[8..12], .little);
+ if (count > 1_000_000 or repeats == 0 or repeats > 10000 or chunk == 0) return error.Options;
+ var input: [65536]u8 = undefined;
+ var output: [65536]u8 = undefined;
+ var elapsed_ns: u64 = 0;
+ for (0..count) |_| {
+ try readAll(input[0..4], chunk);
+ const len = c9.frameLen(input[0..4]).?;
+ if (len < c9.header_len or len > input.len) return error.Frame;
+ try readAll(input[4..len], chunk);
+ const start = nowNs();
+ for (0..repeats) |_| {
+ const decoded = try c9.decode(input[0..len]);
+ const encoded = try c9.encode(decoded.msg, decoded.tag, &output);
+ std.mem.doNotOptimizeAway(encoded);
+ }
+ elapsed_ns += nowNs() - start;
+ try writeAll(output[0..len], chunk);
+ }
+ std.debug.print("{{\"implementation\":\"cloud9\",\"codec_ns\":{d},\"operations\":{d},\"read_calls\":{d},\"write_calls\":{d},\"bytes_in\":{d},\"bytes_out\":{d},\"codec_allocations\":0,\"allocation_evidence\":\"no allocator or allocation calls in codec\"}}\n", .{ elapsed_ns, @as(u64, count) * repeats, reads, writes, bytes_in - header.len, bytes_out });
+}
diff --git a/test/differential/session.go b/test/differential/session.go
new file mode 100644
index 0000000..3fda214
--- /dev/null
+++ b/test/differential/session.go
@@ -0,0 +1,45 @@
+package main
+
+import (
+ "bytes"
+ "context"
+ "encoding/json"
+ "fmt"
+ "github.com/knusbaum/go9p"
+ "github.com/knusbaum/go9p/fs"
+ "os/exec"
+ "strconv"
+ "time"
+)
+
+func sessionProbe(binary string, seed int64, rounds int) map[string]any {
+ tree, root := fs.NewFS("user", "user", 0755)
+ check(root.AddChild(fs.NewStaticFile(tree.NewStat("file", "user", "user", 0600), nil)))
+ ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
+ defer cancel()
+ command := exec.CommandContext(ctx, binary, strconv.FormatInt(seed, 10), strconv.Itoa(rounds))
+ input, e := command.StdinPipe()
+ check(e)
+ output, e := command.StdoutPipe()
+ check(e)
+ var stderr bytes.Buffer
+ command.Stderr = &stderr
+ check(command.Start())
+ done := make(chan error, 1)
+ go func() { done <- go9p.ServeReadWriter(output, input, tree.Server()) }()
+ e = command.Wait()
+ input.Close()
+ if e != nil {
+ panic(fmt.Sprintf("session probe: %v: %s", e, stderr.String()))
+ }
+ select {
+ case <-done:
+ case <-time.After(time.Second):
+ panic("reference server did not stop")
+ }
+ var result map[string]any
+ check(json.Unmarshal(stderr.Bytes(), &result))
+ result["server"] = "go9p v1.18.0 StaticFile"
+ result["checks"] = "walk/open/write/read/stat/EOF/clunk, content oracle, fragmented requests and replies"
+ return result
+}
diff --git a/test/differential/session.zig b/test/differential/session.zig
new file mode 100644
index 0000000..3c5d2de
--- /dev/null
+++ b/test/differential/session.zig
@@ -0,0 +1,82 @@
+//! Local interoperability client. The parent connects stdio to a reference server.
+const std = @import("std");
+const c9 = @import("cloud9");
+const libc = std.c;
+const Session = struct {
+ client: c9.Client,
+ requests: usize = 0,
+ bytes_read: usize = 0,
+ bytes_written: usize = 0,
+ read_calls: usize = 0,
+ write_calls: usize = 0,
+
+ fn ask(s: *Session, request: c9.Client.Request) !c9.Client.Result {
+ const tag = try s.client.submit(request);
+ s.requests += 1;
+ while (s.client.output().len != 0) {
+ const output = s.client.output();
+ s.write_calls += 1;
+ // Fragment requests across their header and body fields.
+ const n = libc.write(1, output.ptr, @min(output.len, 3));
+ if (n <= 0) return error.Write;
+ s.bytes_written += @intCast(n);
+ s.client.wrote(@intCast(n));
+ }
+ while (true) {
+ if (s.client.take()) |done| {
+ if (done.tag != tag) return error.Tag;
+ if (done.result == .fail) {
+ std.debug.print("remote error: {s}\n", .{done.result.fail});
+ return error.Remote;
+ }
+ return done.result;
+ }
+ if (s.client.dead) return error.Protocol;
+ var buffer: [7]u8 = undefined;
+ s.read_calls += 1;
+ const n = libc.read(0, &buffer, buffer.len);
+ if (n <= 0) return error.Read;
+ s.bytes_read += @intCast(n);
+ if (s.client.push(buffer[0..@intCast(n)]) != n) return error.InputFull;
+ }
+ }
+};
+
+pub fn main(init: std.process.Init) !void {
+ const args = try init.minimal.args.toSlice(init.arena.allocator());
+ if (args.len != 3) return error.Arguments;
+ const seed = try std.fmt.parseInt(u64, args[1], 10);
+ const rounds = try std.fmt.parseInt(usize, args[2], 10);
+ if (rounds == 0 or rounds > 10000) return error.Arguments;
+ var random: std.Random.DefaultPrng = .init(seed);
+ var input: [8192]u8 = undefined;
+ var output: [8192]u8 = undefined;
+ var session: Session = .{ .client = .init(.{ .in = &input, .out = &output }) };
+ const version = (try session.ask(.{ .version = .{} })).version;
+ if (version.msize > input.len or !std.mem.eql(u8, version.version, "9P2000")) return error.Version;
+ _ = try session.ask(.{ .attach = .{ .fid = 0, .uname = "user" } });
+ var expected: [4096]u8 = @splat(0);
+ var length: usize = 0;
+ for (0..rounds) |_| {
+ const walk = (try session.ask(.{ .walk = .{ .fid = 0, .newfid = 1, .names = &.{"file"} } })).walk;
+ if (walk.nwqid != 1 or walk.wqid[0].type & c9.qtdir != 0) return error.Walk;
+ _ = try session.ask(.{ .open = .{ .fid = 1, .mode = c9.ordwr } });
+ const offset = random.random().uintLessThan(usize, 2048);
+ const count = random.random().uintLessThan(usize, 1024) + 1;
+ var data: [1024]u8 = undefined;
+ random.random().bytes(data[0..count]);
+ const written = (try session.ask(.{ .write = .{ .fid = 1, .offset = offset, .data = data[0..count] } })).write;
+ if (written != count) return error.WriteCount;
+ @memcpy(expected[offset..][0..count], data[0..count]);
+ length = @max(length, offset + count);
+ const result = (try session.ask(.{ .read = .{ .fid = 1, .offset = 0, .count = expected.len } })).read;
+ if (!std.mem.eql(u8, result, expected[0..length])) return error.Contents;
+ const stat = (try session.ask(.{ .stat = .{ .fid = 1 } })).stat;
+ if (stat.length != length or !std.mem.eql(u8, stat.name, "file")) return error.Stat;
+ const eof = (try session.ask(.{ .read = .{ .fid = 1, .offset = length, .count = 1 } })).read;
+ if (eof.len != 0) return error.Eof;
+ _ = try session.ask(.{ .clunk = .{ .fid = 1 } });
+ }
+ _ = try session.ask(.{ .clunk = .{ .fid = 0 } });
+ std.debug.print("{{\"requests\":{d},\"bytes_read\":{d},\"bytes_written\":{d},\"read_calls\":{d},\"write_calls\":{d},\"seed\":{d},\"rounds\":{d}}}\n", .{ session.requests, session.bytes_read, session.bytes_written, session.read_calls, session.write_calls, seed, rounds });
+}
diff --git a/test/fuzz.zig b/test/fuzz.zig
new file mode 100644
index 0000000..b7d9f21
--- /dev/null
+++ b/test/fuzz.zig
@@ -0,0 +1,70 @@
+//! Deterministic local robustness probes; reference implementations see valid traffic only.
+const std = @import("std");
+const c9 = @import("cloud9");
+pub fn main(init: std.process.Init) !void {
+ const args = try init.minimal.args.toSlice(init.arena.allocator());
+ const seed = if (args.len > 1) try std.fmt.parseInt(u64, args[1], 10) else 4200;
+ const iterations = if (args.len > 2) try std.fmt.parseInt(u32, args[2], 10) else 100000;
+ if (args.len > 3 or iterations > 10000000) return error.Arguments;
+ var prng: std.Random.DefaultPrng = .init(seed);
+ const random = prng.random();
+ var accepted: usize = 0;
+ var rejected: usize = 0;
+ var buffer: [8192]u8 = undefined;
+ var encoded: [8192]u8 = undefined;
+ var payload: [4096]u8 = undefined;
+ var server_in: [8192]u8 = undefined;
+ var server_out: [8192]u8 = undefined;
+ for (0..iterations) |_| {
+ random.bytes(&payload);
+ const data = payload[0..random.uintLessThan(usize, payload.len)];
+ const qid: c9.Qid = .{ .type = random.int(u8), .version = random.int(u32), .path = random.int(u64) };
+ const stat: c9.Stat = .{
+ .type = 0,
+ .dev = 0,
+ .qid = qid,
+ .mode = random.int(u32),
+ .atime = 0,
+ .mtime = 0,
+ .length = random.int(u64),
+ .name = "file",
+ .uid = "user",
+ .gid = "group",
+ .muid = "",
+ };
+ const messages = [_]c9.Msg{
+ .{ .tversion = .{ .msize = 8192, .version = "9P2000" } },
+ .{ .twalk = .{ .fid = 0, .newfid = 1, .nwname = 16, .wname = @splat("dir") } },
+ .{ .rwalk = .{ .nwqid = 16, .wqid = @splat(qid) } },
+ .{ .twrite = .{ .fid = random.int(u32), .offset = random.int(u64), .data = data } },
+ .{ .rread = .{ .data = data } },
+ .{ .rstat = .{ .stat = stat } },
+ .{ .twstat = .{ .fid = 1, .stat = stat } },
+ .{ .tflush = .{ .oldtag = random.int(u16) } },
+ .{ .tread = .{ .fid = 1, .offset = random.int(u64), .count = random.int(u32) } },
+ };
+ const msg = messages[random.uintLessThan(usize, messages.len)];
+ var frame: []u8 = try c9.encode(msg, random.int(u16), &buffer);
+ // Leave some valid seeds unchanged; probe truncation and altered fields locally.
+ switch (random.uintLessThan(u8, 4)) {
+ 0 => {},
+ 1 => frame = frame[0..random.uintLessThan(usize, frame.len)],
+ 2 => frame[random.uintLessThan(usize, frame.len)] ^= random.int(u8),
+ 3 => std.mem.writeInt(u32, frame[0..4], random.int(u32), .little),
+ else => unreachable,
+ }
+ if (c9.decode(frame)) |decoded| {
+ const output = try c9.encode(decoded.msg, decoded.tag, &encoded);
+ if (!std.mem.eql(u8, frame, output)) return error.RoundTrip;
+ accepted += 1;
+ } else |_| rejected += 1;
+ var server: c9.Server = .init(.{ .in = &server_in, .out = &server_out });
+ _ = server.push(frame);
+ if (server.receive() catch null) |request| {
+ if (request.msg != .tversion) return error.BeforeVersion;
+ server.negotiate(request.msg.tversion.msize, request.msg.tversion.version) catch {};
+ server.release();
+ }
+ }
+ std.debug.print("seed={d} iterations={d} accepted={d} rejected={d}\n", .{ seed, iterations, accepted, rejected });
+}
diff --git a/test/quic.zig b/test/quic.zig
new file mode 100644
index 0000000..48ff839
--- /dev/null
+++ b/test/quic.zig
@@ -0,0 +1,268 @@
+const std = @import("std");
+const libc = std.c;
+const ssl = @import("openssl");
+const Quic = @import("cloud9").Quic(ssl, "cloud9-test");
+const Listener = Quic.Listener;
+const Connection = Quic.Connection;
+const TestPair = struct {
+ listener: *Listener,
+ client: Connection,
+ server: ?Connection = null,
+
+ fn init(listener: *Listener) !TestPair {
+ var p: TestPair = .{ .listener = listener, .client = try .dial(listener.address) };
+ errdefer p.deinit();
+ const deadline = testNow() + 3000;
+ while (true) {
+ try listener.events();
+ try p.client.events();
+ if (p.server == null) p.server = try listener.accept();
+ const connected = try p.client.handshake();
+ if (p.server) |*server| if (connected and try server.handshake()) return p;
+ try p.wait(deadline);
+ }
+ }
+
+ fn wait(p: *TestPair, deadline: i64) !void {
+ const remaining = deadline - testNow();
+ if (remaining <= 0) return error.Deadline;
+ var timeout: i32 = @intCast(@min(remaining, std.math.maxInt(i32)));
+ if (p.listener.nextDue()) |ms| timeout = @min(timeout, ms);
+ if (p.client.nextDue()) |ms| timeout = @min(timeout, ms);
+ var fds = [_]libc.pollfd{ p.listener.poll(), p.client.poll().? };
+ const rc = libc.poll(&fds, fds.len, timeout);
+ if (rc < 0 and libc.errno(rc) != .INTR) return error.Poll;
+ if (testNow() >= deadline) return error.Deadline;
+ try p.listener.events();
+ try p.client.events();
+ }
+
+ fn transfer(p: *TestPair, from_client: bool, bytes: []const u8, fragment: usize) !void {
+ const writer = if (from_client) &p.client else &p.server.?;
+ const reader = if (from_client) &p.server.? else &p.client;
+ var sent: usize = 0;
+ var received: usize = 0;
+ var buffer: [8192]u8 = undefined;
+ const deadline = testNow() + 3000;
+ while (received < bytes.len) {
+ const written = if (sent != bytes.len) try writer.write(bytes[sent..][0..@min(fragment, bytes.len - sent)]) else 0;
+ sent += written;
+ const count = try reader.read(buffer[0..@min(fragment, buffer.len)]);
+ if (count) |n| {
+ try std.testing.expect(n > 0 and n <= bytes.len - received);
+ try std.testing.expectEqualSlices(u8, bytes[received..][0..n], buffer[0..n]);
+ received += n;
+ }
+ if (testNow() >= deadline) return error.Deadline;
+ if (received != bytes.len and written == 0 and count == null) {
+ try p.wait(deadline);
+ } else {
+ try p.listener.events();
+ try p.client.events();
+ }
+ }
+ try std.testing.expectEqual(bytes.len, sent);
+ }
+
+ fn finish(p: *TestPair) !void {
+ try p.client.conclude();
+ try p.server.?.conclude();
+ var buffer: [16]u8 = undefined;
+ var a = false;
+ var b = false;
+ const deadline = testNow() + 3000;
+ while (!a or !b) {
+ if (!a) if (try p.client.read(&buffer)) |n| {
+ try std.testing.expectEqual(@as(usize, 0), n);
+ a = true;
+ };
+ if (!b) if (try p.server.?.read(&buffer)) |n| {
+ try std.testing.expectEqual(@as(usize, 0), n);
+ b = true;
+ };
+ if (!a or !b) try p.wait(deadline);
+ }
+ }
+
+ fn deinit(p: *TestPair) void {
+ if (p.server) |*server| server.deinit();
+ p.client.deinit();
+ }
+};
+
+fn testNow() i64 {
+ var ts: libc.timespec = undefined;
+ std.debug.assert(libc.clock_gettime(.MONOTONIC, &ts) == 0);
+ return @as(i64, @intCast(ts.sec)) * 1000 + @divFloor(@as(i64, @intCast(ts.nsec)), 1_000_000);
+}
+
+test "QUIC fragmented 9P frames reconnect and stream EOF over IPv4 and IPv6" {
+ const request = "\x13\x00\x00\x00\x64\xff\xff\x00\x20\x00\x00\x06\x00" ++ "9P2000";
+ const response = "\x13\x00\x00\x00\x65\xff\xff\x00\x20\x00\x00\x06\x00" ++ "9P2000";
+ const read_request = "\x17\x00\x00\x00\x74\x01\x00\x02\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x0b\x00\x00\x00";
+ const read_response = "\x16\x00\x00\x00\x75\x01\x00\x0b\x00\x00\x00" ++ "hello ninep";
+ for ([_][]const u8{ "127.0.0.1", "::1" }) |host| {
+ var listener = try Listener.init(try .parse(host, 0));
+ defer listener.deinit();
+ try std.testing.expect(listener.address.getPort() != 0);
+ if (listener.address == .ip6) {
+ var enabled: c_int = 0;
+ var len: libc.socklen_t = @sizeOf(c_int);
+ const v6only = if (@import("builtin").os.tag.isDarwin()) 27 else libc.IPV6.V6ONLY;
+ try std.testing.expectEqual(@as(c_int, 0), libc.getsockopt(listener.fd, libc.IPPROTO.IPV6, v6only, @ptrCast(&enabled), &len));
+ try std.testing.expectEqual(@as(c_int, 1), enabled);
+ }
+ for (0..3) |_| {
+ var pair = try TestPair.init(&listener);
+ defer pair.deinit();
+ try std.testing.expect(pair.server.?.poll() == null);
+ try std.testing.expect(pair.server.?.nextDue() == null);
+ for ([_]usize{ 1, 2, 7, 64 }) |fragment| {
+ try pair.transfer(true, request, fragment);
+ try pair.transfer(false, response, fragment);
+ try pair.transfer(true, read_request, fragment);
+ try pair.transfer(false, read_response, fragment);
+ }
+ try pair.finish();
+ }
+ }
+}
+
+test "QUIC backpressure retries a moved prefix while new replies are appended" {
+ var listener = try Listener.init(try .parse("127.0.0.1", 0));
+ defer listener.deinit();
+ var pair = try TestPair.init(&listener);
+ defer pair.deinit();
+ try pair.transfer(true, "hello", 1);
+ const bytes: [8192]u8 = @splat(0x5a);
+ var total: usize = 0;
+ const deadline = testNow() + 3000;
+ while (total < 64 * 1024 * 1024) {
+ const written = try pair.client.write(&bytes);
+ total += written;
+ if (written == 0) break;
+ try pair.listener.events();
+ try pair.client.events();
+ if (testNow() >= deadline) return error.Deadline;
+ }
+ try std.testing.expect(total > 0 and total < 64 * 1024 * 1024);
+ try std.testing.expectEqual(bytes.len, pair.client.pending_write_len);
+ try std.testing.expectError(error.InvalidWrite, pair.client.write(bytes[0..1]));
+ var buffer: [8192]u8 = undefined;
+ var received: usize = 0;
+ // A pending SSL write may already have delivered a prefix before it
+ // reports completion. Leave that prefix for the retry check below.
+ while (received < total) {
+ if (try pair.server.?.read(buffer[0..@min(buffer.len, total - received)])) |n| {
+ try std.testing.expect(n > 0);
+ try std.testing.expect(std.mem.allEqual(u8, buffer[0..n], 0x5a));
+ received += n;
+ } else try pair.wait(deadline);
+ }
+ try std.testing.expectEqual(total, received);
+ var moved: [8192 + 7]u8 = undefined;
+ @memcpy(moved[0..bytes.len], &bytes);
+ @memset(moved[bytes.len..], 0x6b);
+ while (true) {
+ const n = try pair.client.write(&moved);
+ if (n != 0) {
+ try std.testing.expectEqual(bytes.len, n);
+ break;
+ }
+ try pair.wait(deadline);
+ }
+ received = 0;
+ while (received < bytes.len) {
+ if (try pair.server.?.read(&buffer)) |n| {
+ try std.testing.expect(n > 0);
+ try std.testing.expect(std.mem.allEqual(u8, buffer[0..n], 0x5a));
+ received += n;
+ } else try pair.wait(deadline);
+ }
+ try std.testing.expectEqual(bytes.len, received);
+ try pair.transfer(true, moved[bytes.len..], 7);
+ try pair.finish();
+}
+
+test "QUIC owner can enforce handshake and read deadlines without busy polling" {
+ var listener = try Listener.init(try .parse("127.0.0.1", 0));
+ defer listener.deinit();
+ var client = try Connection.dial(listener.address);
+ defer client.deinit();
+ const handshake_deadline = testNow() + 100;
+ var turns: usize = 0;
+ while (testNow() < handshake_deadline) {
+ try std.testing.expect(!try client.handshake());
+ try client.events();
+ var timeout: i32 = @intCast(@max(0, handshake_deadline - testNow()));
+ if (client.nextDue()) |ms| timeout = @min(timeout, ms);
+ var fds = [_]libc.pollfd{client.poll().?};
+ const rc = libc.poll(&fds, fds.len, timeout);
+ if (rc < 0 and libc.errno(rc) != .INTR) return error.Poll;
+ turns += 1;
+ }
+ try std.testing.expect(turns < 100);
+ var serving = try Listener.init(try .parse("127.0.0.1", 0));
+ defer serving.deinit();
+ var pair = try TestPair.init(&serving);
+ defer pair.deinit();
+ try pair.transfer(true, "request", 2);
+ var buffer: [32]u8 = undefined;
+ const read_deadline = testNow() + 100;
+ turns = 0;
+ while (true) {
+ try std.testing.expect(try pair.client.read(&buffer) == null);
+ pair.wait(read_deadline) catch |err| {
+ try std.testing.expectEqual(error.Deadline, err);
+ break;
+ };
+ turns += 1;
+ }
+ try std.testing.expect(turns < 100);
+}
+
+test "QUIC bind failure leaves the existing listener usable" {
+ var listener = try Listener.init(try .parse("127.0.0.1", 0));
+ defer listener.deinit();
+ for (0..8) |_| try std.testing.expectError(error.Bind, Listener.init(listener.address));
+ var pair = try TestPair.init(&listener);
+ defer pair.deinit();
+ try pair.transfer(true, "still listening", 3);
+ try pair.transfer(false, "still serving", 2);
+ try pair.finish();
+}
+
+test "QUIC pending reports buffered bytes and EOF without UDP readiness" {
+ var listener = try Listener.init(try .parse("127.0.0.1", 0));
+ defer listener.deinit();
+ var pair = try TestPair.init(&listener);
+ defer pair.deinit();
+ try std.testing.expect(!pair.server.?.pending());
+ try std.testing.expectEqual(@as(usize, 3), try pair.client.write("abc"));
+ const deadline = testNow() + 3000;
+ while (!pair.server.?.pending()) try pair.wait(deadline);
+ var buffer: [3]u8 = undefined;
+ try std.testing.expectEqual(@as(?usize, 1), try pair.server.?.read(buffer[0..1]));
+ try std.testing.expectEqual(@as(u8, 'a'), buffer[0]);
+ try listener.events();
+ try pair.client.events();
+ try std.testing.expect(pair.server.?.pending());
+ try std.testing.expectEqual(@as(?usize, 2), try pair.server.?.read(buffer[1..]));
+ try std.testing.expectEqualStrings("abc", &buffer);
+ try std.testing.expect(!pair.server.?.pending());
+ try pair.client.conclude();
+ while (!pair.server.?.pending()) try pair.wait(deadline);
+ try std.testing.expectEqual(@as(?usize, 0), try pair.server.?.read(&buffer));
+}
+
+test "QUIC moved listener retains its in-memory identity" {
+ var original = try Listener.init(try .parse("127.0.0.1", 0));
+ var listener = original;
+ original = undefined;
+ defer listener.deinit();
+ var pair = try TestPair.init(&listener);
+ defer pair.deinit();
+ try pair.transfer(true, "moved listener", 2);
+ try pair.transfer(false, "same identity", 3);
+ try pair.finish();
+}
diff --git a/test/transport.zig b/test/transport.zig
new file mode 100644
index 0000000..3d43ddb
--- /dev/null
+++ b/test/transport.zig
@@ -0,0 +1,82 @@
+const std = @import("std");
+const c9 = @import("cloud9");
+const t = c9.transport;
+const testing = std.testing;
+
+test "standard readers and writers preserve frame boundaries" {
+ var bytes: [128]u8 = undefined;
+ const frame = try c9.encode(.{ .tversion = .{ .msize = 4096, .version = "9P2000" } }, c9.notag, &bytes);
+ var source: testing.Reader = .init(&.{}, &.{.{ .buffer = frame }});
+ source.artificial_limit = .limited(1);
+ var target: [128]u8 = undefined;
+ const got = try t.readFrame(&source.interface, &target, 4096);
+ try testing.expectEqualSlices(u8, frame, got);
+ var output: [128]u8 = undefined;
+ var writer: std.Io.Writer = .fixed(&output);
+ try t.writeFrame(&writer, got, 4096);
+ try testing.expectEqualSlices(u8, frame, writer.buffered());
+}
+
+test "POSIX TCP and Unix streams handle retry and EOF" {
+ var dir = testing.tmpDir(.{});
+ defer dir.cleanup();
+ const io = testing.io;
+ var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
+ const parent_len = try dir.dir.realPath(io, &path_buffer);
+ const parent = path_buffer[0..parent_len];
+ var unix_buffer: [t.sun_path_len]u8 = undefined;
+ const unix = try std.fmt.bufPrintSentinel(&unix_buffer, "{s}/9p", .{parent}, 0);
+ for ([_]t.Address{ .{ .tcp = .{ .ip4 = .loopback(0) } }, .{ .unix = unix } }) |address| {
+ const listener = try t.listenFd(address, 4);
+ defer t.close(listener);
+ var destination = address;
+ if (address == .tcp) {
+ var actual: std.c.sockaddr.in = undefined;
+ var len: std.c.socklen_t = @sizeOf(@TypeOf(actual));
+ try testing.expectEqual(@as(c_int, 0), std.c.getsockname(listener, @ptrCast(&actual), &len));
+ destination.tcp.setPort(std.mem.bigToNative(u16, actual.port));
+ }
+ const client = try t.connectFd(destination, t.nowMs() + 1000);
+ defer t.close(client);
+ const server = (try t.acceptFd(listener, address == .tcp)) orelse return error.NoConnection;
+ var buffer: [32]u8 = undefined;
+ try testing.expectEqual(@as(?usize, null), try t.read(server, &buffer));
+ try testing.expectEqual(@as(?usize, 3), try t.write(client, "abc"));
+ try t.wait(server, @intCast(std.c.POLL.IN), t.nowMs() + 1000);
+ try testing.expectEqual(@as(?usize, 3), try t.read(server, &buffer));
+ try testing.expectEqualStrings("abc", buffer[0..3]);
+ t.close(server);
+ try t.wait(client, @intCast(std.c.POLL.IN), t.nowMs() + 1000);
+ try testing.expectEqual(@as(?usize, 0), try t.read(client, &buffer));
+ }
+}
+
+test "std.Io TCP and Unix listener and client adapters" {
+ var dir = testing.tmpDir(.{});
+ defer dir.cleanup();
+ const io = testing.io;
+ var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
+ const parent_len = try dir.dir.realPath(io, &path_buffer);
+ const parent = path_buffer[0..parent_len];
+ var unix_buffer: [t.sun_path_len]u8 = undefined;
+ const unix = try std.fmt.bufPrintSentinel(&unix_buffer, "{s}/9p", .{parent}, 0);
+ for ([_]t.Address{ .{ .tcp = .{ .ip4 = .loopback(0) } }, .{ .unix = unix } }) |address| {
+ var listener = try t.listen(io, address, 4);
+ defer listener.deinit(io);
+ const destination: t.Address = if (address == .tcp) .{ .tcp = listener.socket.address } else address;
+ const client = try t.connect(io, destination);
+ defer client.close(io);
+ const server = try listener.accept(io);
+ defer server.close(io);
+ var output: [128]u8 = undefined;
+ var writer = client.writer(io, &output);
+ var bytes: [128]u8 = undefined;
+ const frame = try c9.encode(.rflush, 1, &bytes);
+ try t.writeFrame(&writer.interface, frame, 4096);
+ try writer.interface.flush();
+ var input: [128]u8 = undefined;
+ var reader = server.reader(io, &input);
+ var target: [128]u8 = undefined;
+ try testing.expectEqualSlices(u8, frame, try t.readFrame(&reader.interface, &target, 4096));
+ }
+}