summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--README.md14
-rw-r--r--build.zig32
-rw-r--r--docs/design.md47
-rw-r--r--docs/http.md99
-rw-r--r--test/web/http_fs.mjs131
-rw-r--r--test/web/mux.zig310
-rw-r--r--web/httpfs.zig583
-rw-r--r--web/main.zig411
-rw-r--r--web/mux.zig812
-rw-r--r--web/probe.zig83
-rw-r--r--web/static/app.mjs183
-rw-r--r--web/static/index.html70
-rw-r--r--web/static/style.css22
13 files changed, 2632 insertions, 165 deletions
diff --git a/README.md b/README.md
index 04c68fa..c06331f 100644
--- a/README.md
+++ b/README.md
@@ -100,9 +100,17 @@ own sources, tests, README and a `build.zig` fragment that the root
`build.zig` imports and enables with a `-D<name>` toggle (`zig build --help`
lists the steps). New related programs follow the same layout.
-* [`web/`](docs/http.md) — the HTTP/WebSocket gateway `9web` and its
- WebAssembly browser client (`web/main.zig`, `web/client.zig`, assets under
- `web/static/`). Always built; steps `serve`, `web`, `http-test`, `e2e`.
+* [`web/`](docs/http.md) — the HTTP/WebSocket gateway `9web`, a 9P
+ multiplexer and its WebAssembly browser client (`web/main.zig`,
+ `web/mux.zig`, `web/httpfs.zig`, `web/probe.zig`, `web/client.zig`, assets
+ under `web/static/`). It holds one upstream 9P connection and fans it out to
+ many downstreams: the browser over a WebSocket, plain 9P clients over an
+ optional `--serve` TCP/Unix listener (a mux like 9pserve), and an HTTP view
+ of the tree under `/fs/` (`GET`/`PUT`/`DELETE`, directory JSON, `Range`, and
+ `?follow=1` server-sent events). `--probe` embeds 9proc for
+ self-introspection. The page adds a tree view, create/rename/delete, upload,
+ a stat panel and live follow. Always built; steps `serve`, `web`,
+ `http-test`, `mux-test`, `http-fs-test`, `e2e`.
* [`9proc/`](9proc/README.md) — a 9P debug/introspection server as
a library (module `9proc`, exported next to `cloud9`; freestanding
core that is a backend of `fs.Server`, Linux debug layer) and its demo
diff --git a/build.zig b/build.zig
index 13e54ec..334da9f 100644
--- a/build.zig
+++ b/build.zig
@@ -36,7 +36,6 @@ pub fn build(b: *std.Build) void {
.root_source_file = b.path("web/main.zig"),
.target = target,
.optimize = optimize,
- .link_libc = true,
.imports = &.{.{ .name = "cloud9", .module = module }},
}) });
bridge.root_module.addAnonymousImport("client.wasm", .{ .root_source_file = wasm.getEmittedBin() });
@@ -58,6 +57,19 @@ pub fn build(b: *std.Build) void {
.imports = &.{.{ .name = "cloud9", .module = module }},
}) });
b.step("http-test", "Test HTTP WebSocket transport framing").dependOn(&b.addRunArtifact(http_tests).step);
+ const mux_mod = b.createModule(.{
+ .root_source_file = b.path("web/mux.zig"),
+ .target = target,
+ .optimize = optimize,
+ .imports = &.{.{ .name = "cloud9", .module = module }},
+ });
+ const mux_tests = b.addTest(.{ .root_module = b.createModule(.{
+ .root_source_file = b.path("test/web/mux.zig"),
+ .target = target,
+ .optimize = optimize,
+ .imports = &.{ .{ .name = "cloud9", .module = module }, .{ .name = "mux", .module = mux_mod } },
+ }) });
+ b.step("mux-test", "Test the 9web multiplexer (shared upstream, fid isolation, reconnect)").dependOn(&b.addRunArtifact(mux_tests).step);
const e2e = b.addSystemCommand(&.{ "node", "test/web/e2e.mjs" });
e2e.addArtifactArg(bridge);
e2e.addArtifactArg(native_http);
@@ -165,6 +177,24 @@ pub fn build(b: *std.Build) void {
if (i.debug_test_step) |s| programs_test.dependOn(s);
}
+ // 9web can embed 9proc for self-introspection (--probe). The module is
+ // linked into 9web only when 9proc is enabled (Linux); otherwise the probe
+ // code compiles to a stub selected by the have_probe build option.
+ {
+ const web_opts = b.addOptions();
+ web_opts.addOption(bool, "have_probe", proc != null);
+ bridge.root_module.addImport("build_options", web_opts.createModule());
+ if (proc) |i| bridge.root_module.addImport("9proc", i.module);
+ }
+
+ // The HTTP /fs mapping suite drives 9web against a 9proc-demo upstream.
+ if (proc) |i| if (i.demo) |demo| {
+ const hfs = b.addSystemCommand(&.{ "node", "test/web/http_fs.mjs" });
+ hfs.addArtifactArg(demo);
+ hfs.addArtifactArg(bridge);
+ b.step("http-fs-test", "Test the HTTP /fs view against a 9proc-demo upstream").dependOn(&hfs.step);
+ };
+
const ns: ?ns_build.Artifacts = if (want_9ns) ns_build.add(b, .{
.target = target,
.optimize = optimize,
diff --git a/docs/design.md b/docs/design.md
index 3ad9338..f34a3ea 100644
--- a/docs/design.md
+++ b/docs/design.md
@@ -173,6 +173,53 @@ that waits on another thread must stop waiting once `stop()` has begun.
Unix socket paths are the application's: the runner neither unlinks,
chmods nor removes them.
+# Multiplexer (9web)
+
+`9web` (`web/`) is a gateway and a 9P multiplexer. It keeps **one** upstream
+connection — a single `cloud9.Client`, all sixteen tags — and fans it out to
+any number of downstream sessions: the browser WASM client over a WebSocket at
+`/_cloud9/9p`, plain 9P clients over an optional `--serve` TCP/Unix listener,
+and an HTTP view of the tree under `/fs/`. `web/mux.zig` owns the shared
+upstream; `web/httpfs.zig` is the HTTP view; `web/probe.zig` embeds 9proc under
+`--probe`.
+
+The mux is a **frame-level remux**, not an `fs.Server`/`serve.Runner` backend.
+Each downstream request is forwarded upstream on a remapped tag with its fids
+remapped into the shared upstream fid space, and each upstream reply is routed
+back by tag. `Tversion` is answered locally (the upstream is negotiated once at
+connect); `Tattach` is forwarded, so each downstream gets its own upstream tree
+root; `Tflush` is forwarded as `Tflush`. Downstream fid spaces are isolated by
+construction — two downstreams never share an upstream fid. The choice is
+deliberate: the file-server engine re-decomposes each request into filesystem
+operations and re-encodes directories with synthetic stats and an entry-index
+cursor, which would lose the upstream's real directory stats, iounit and qids
+and turn the readdir byte offset into an index mapping. Forwarding frames
+unchanged is exactly the "forward each request upstream" contract, the shape
+plan9port's 9pserve has, and it keeps the upstream's bytes intact. `serve.Runner`
+remains the right tool for a *server* whose backend is a real filesystem; a
+transparent *proxy* is not that.
+
+The shared upstream runs one reader task and any number of downstream forwarder
+tasks under one mutex; the tag, fid and pending tables are comptime-sized and
+nothing allocates per request. Sockets are only written under the mutex with the
+client's own output buffer (disjoint from the stream writer's), and replies are
+encoded into per-tag buffers and sent to downstreams outside the lock so a slow
+downstream cannot stall the client. The WebSocket and `--serve` downstreams
+drive the mux the same way — each reads whole frames (WebSocket binary messages
+or `transport.readFrame`) and calls `forward`; the runner's socket-only listener
+is not used for them because WebSocket framing and the HTTP `/fs` view do not fit
+it. When all sixteen upstream tags are in flight, a further downstream request
+waits for one to free (Tflush keeps its reserved seventeenth tag). On an upstream
+I/O failure the reader fails every outstanding request, reconnects, bumps a
+generation, and answers any downstream request that still names a pre-reconnect
+fid with EIO until it re-attaches.
+
+The HTTP `/fs` view uses the same shared upstream at the fid level through a
+blocking `rpc`: it allocates upstream fids from the mux pool and walks, stats,
+reads and writes directly, so `GET`/`PUT`/`DELETE`, directory JSON, `Range` and
+`?follow=1` server-sent events all ride the one upstream connection. A parked
+`follow` read is cancelled with an upstream `Tflush` when its client goes away.
+
# Related programs
Programs built on the library ship from this repository as `cloud9/<name>/`,
diff --git a/docs/http.md b/docs/http.md
index 4e92e5e..5b51e2f 100644
--- a/docs/http.md
+++ b/docs/http.md
@@ -4,14 +4,84 @@ cloud9 provides a reusable HTTP transport and a standalone browser gateway.
The transport carries unchanged base 9P2000 messages. Filesystem permissions,
mounting, UART configuration, and board drivers remain with their applications.
+`9web` is a **multiplexer**: it holds one upstream 9P connection (a single
+`cloud9.Client`, all sixteen tags) and fans it out to any number of downstream
+sessions — the browser WASM client over a WebSocket, plain 9P clients over an
+optional TCP/Unix listener (a mux like plan9port's 9pserve), and an HTTP view
+of the tree under `/fs/`. Each downstream keeps its own fid space; two
+downstreams never share an upstream fid. On an upstream failure the gateway
+reconnects and invalidates downstream fids (their next request answers EIO
+until they re-attach).
+
```mermaid
flowchart LR
- Browser[Browser: cloud9 WASM client] -->|WebSocket / HTTPS| Gateway[HTTP gateway]
+ Browser[Browser: cloud9 WASM client] -->|WebSocket| Gateway[9web mux]
Native[Native cloud9 client] -->|WebSocket / HTTPS| Gateway
- Gateway -->|TCP or Unix stream| Server[9P server]
- Gateway -->|Serial byte stream| Board[9P firmware / ESP32]
+ Curl[curl / fetch / SSE] -->|HTTP /fs| Gateway
+ P9[9ns / plan9port 9p] -->|TCP / Unix 9P| Gateway
+ Gateway -->|one shared TCP / Unix / serial connection| Server[9P server]
+```
+
+## The multiplexer
+
+Downstream requests are forwarded upstream on remapped tags with their fids
+remapped into the shared upstream fid space; upstream replies are routed back to
+the originating downstream by tag. `Tversion` is answered locally (the upstream
+session is negotiated once at connect); `Tattach` is forwarded so every
+downstream gets its own upstream tree root; `Tflush` is forwarded as `Tflush`.
+Frames are forwarded unchanged, so directory reads, qids, iounits and stat
+fields keep the upstream's exact values — the mux does not re-derive them.
+
+This is deliberately a frame-level remux rather than an `fs.Server`/`serve.Runner`
+backend: the file-server engine re-decomposes each request into filesystem
+operations and re-encodes directories with synthetic stats and an entry-index
+cursor, which would lose the upstream's real directory stats and complicate the
+readdir byte offset. Forwarding requests unchanged is exactly the "one upstream,
+many downstreams" contract. See [design.md](design.md#multiplexer-9web).
+
+The tables are comptime-sized and nothing allocates per request: the shared
+upstream has one reader task and any number of downstream forwarder tasks
+serialized by one mutex, sixteen ordinary tags plus one flush tag, a fixed
+upstream fid pool, and a bounded downstream-connection table. When all sixteen
+upstream tags are in flight, further downstream requests wait for one to free.
+
+## The HTTP view of the tree
+
+`/fs/<path>` maps HTTP onto the 9P tree — file operations, not verbs:
+
+| Request | 9P | Result |
+| --- | --- | --- |
+| `GET /fs/<file>` | walk, open, read | bytes; `Content-Type` sniffed (`text/plain` or `application/octet-stream`); `Range` honored (206) |
+| `GET /fs/<file>?follow=1` | walk, open, blocking reads | `text/event-stream`; each read that returns data is one SSE `data:` event; also triggered by `Accept: text/event-stream` |
+| `GET /fs/<dir>` | walk, open, read | JSON array of `{name, dir, length, mode, mtime, qid}`; an HTML listing with `Accept: text/html` |
+| `HEAD /fs/<path>` | walk, stat | headers: `content-length`, `x-9p-qid`, `x-9p-mode`, `x-9p-mtime` |
+| `PUT /fs/<file>` | walk or create, open, write | write; creates if missing, truncates unless `?append=1` |
+| `PUT /fs/<dir>/` | walk, create | create a directory (trailing slash) |
+| `DELETE /fs/<path>` | walk, remove | remove |
+
+A 9P `Rerror` maps to a status with the ename in the body: *no such
+file* → 404, *permission* → 403, *exists* / *not empty* / *is a directory* →
+409, otherwise 500. Path components are percent-decoded once; encoded slashes,
+NULs and malformed escapes are rejected. `..` is refused.
+
+Verify against a running `9proc-demo` upstream:
+
+```sh
+9proc-demo --unix /tmp/up.sock &
+9web --listen 127.0.0.1:8080 --upstream unix:/tmp/up.sock
+curl -s localhost:8080/fs/ # directory JSON
+curl -s localhost:8080/fs/runtime/fn/now # a live file
+curl -s -X PUT --data-binary hi localhost:8080/fs/scratch/x # write
+curl -N localhost:8080/fs/self/log # follow a blocking file (Pardes)
```
+## Self-introspection
+
+`--probe unix:PATH` (off by default) embeds [9proc](../9proc/README.md) and
+serves the gateway's own live counters as files — `/runtime/fn/{connections,
+downstreams,upstream,reconnects,negotiated}` — plus `/threads`, on a private 9P
+listener. Read-only; serve it on a Unix socket or loopback.
+
## Run
```sh
@@ -71,21 +141,24 @@ Options:
| Option | Default | Meaning |
| --- | --- | --- |
| `--listen IP:PORT` | `127.0.0.1:8080` | HTTP listener; IPv4 or bracketed IPv6 |
-| `--upstream tcp:IP:PORT` | `tcp:127.0.0.1:564` | A new 9P connection per WebSocket |
-| `--upstream unix:PATH` | — | A new Unix connection per WebSocket |
-| `--upstream file:PATH` | — | An already configured duplex device, one session at a time |
+| `--upstream tcp:IP:PORT` | `tcp:127.0.0.1:564` | The one shared upstream 9P connection (TCP) |
+| `--upstream unix:PATH` | — | The one shared upstream 9P connection (Unix) |
+| `--upstream file:PATH` | — | An already configured duplex device (serial), shared by all downstreams |
+| `--serve tcp:IP:PORT\|unix:PATH` | — | Also serve plain 9P downstreams (for 9ns / plan9port clients) |
+| `--probe unix:PATH\|tcp:IP:PORT` | — | Embed 9proc for self-introspection (off by default) |
| `--origin SCHEME://HOST[:PORT]` | Listener origin | Public origin when a reverse proxy serves the gateway |
| `--user NAME` | `user` | Browser's 9P attach user (`uname`), at most 256 bytes |
| `--tree NAME` | empty | Browser's named export (`aname`), at most 256 bytes |
-| `--timeout-ms N` | `300000` | Maximum connection lifetime, including HTTP headers; 0 disables it |
+| `--timeout-ms N` | `300000` | Maximum HTTP/WebSocket connection lifetime; 0 disables it |
-The gateway bounds concurrent HTTP/WebSocket connections at 32, HTTP headers at
+The gateway bounds concurrent downstream connections at 64, HTTP headers at
8 KiB, and 9P frames at 64 KiB. Excess connections close. A connection ending or
-expiring cancels both relay directions before releasing buffers and descriptors.
-Each network connection has a separate upstream fid/tag space. A device lease
-prevents multiple clients from mixing transactions on one physical UART stream.
-A new serial session starts with Tversion; the device owner is responsible for
-link reset/recovery if an interrupted physical link leaves stale bytes in transit.
+expiring drops the session and reclaims its upstream fids. Downstream fid spaces
+are isolated by the mux; the one shared upstream connection carries them all,
+which also means one serial link is multiplexed rather than leased — the mux
+serializes transactions over it. The upstream session starts with a single
+Tversion at connect; the device owner is responsible for link reset/recovery if
+an interrupted physical link leaves stale bytes in transit.
For a configured serial device:
diff --git a/test/web/http_fs.mjs b/test/web/http_fs.mjs
new file mode 100644
index 0000000..8ac4961
--- /dev/null
+++ b/test/web/http_fs.mjs
@@ -0,0 +1,131 @@
+// HTTP /fs mapping tests: spawn 9proc-demo as the upstream and 9web as the
+// gateway, then exercise GET (file/dir/Range/SSE), HEAD, PUT (create/truncate/
+// append), DELETE, mkdir, 404 mapping, and concurrent fan-out over the one
+// shared upstream. No npm packages; Node 22+ builtins only.
+import assert from 'node:assert/strict';
+import { spawn } from 'node:child_process';
+import { once } from 'node:events';
+import path from 'node:path';
+import os from 'node:os';
+import fs from 'node:fs/promises';
+
+const [demoBin, webBin] = process.argv.slice(2).map(p => path.resolve(p));
+if (!webBin) throw new Error('usage: node test/web/http_fs.mjs 9proc-demo 9web');
+
+const children = [];
+const delay = ms => new Promise(r => setTimeout(r, ms));
+function start(cmd, args) {
+ const c = spawn(cmd, args, { stdio: ['ignore', 'pipe', 'pipe'] });
+ c.out = ''; c.err = '';
+ c.stdout.on('data', b => { c.out += b; });
+ c.stderr.on('data', b => { c.err += b; });
+ c.on('error', e => { c.err += e.message; });
+ children.push(c);
+ return c;
+}
+async function wait(check, msg, timeout = 8000) {
+ const end = Date.now() + timeout;
+ while (Date.now() < end) { try { const r = await check(); if (r) return r; } catch {} await delay(40); }
+ throw new Error('timeout: ' + msg);
+}
+
+let sock, url;
+try {
+ sock = path.join(await fs.mkdtemp(path.join(os.tmpdir(), '9web-')), 'up');
+ const demo = start(demoBin, ['--unix', sock]);
+ await wait(() => demo.err.includes('listening') || demo.out.includes('listening'), '9proc-demo ready');
+ const web = start(webBin, ['--listen', '127.0.0.1:0', '--upstream', `unix:${sock}`, '--timeout-ms', '0']);
+ const m = await wait(() => web.err.match(/9web (http:\/\/127\.0\.0\.1:\d+)/), '9web ready');
+ url = m[1];
+
+ // GET directory -> JSON with real stats.
+ const root = await (await fetch(`${url}/fs/`)).json();
+ const names = root.map(e => e.name);
+ assert.ok(names.includes('README') && names.includes('runtime') && names.includes('scratch'), 'root listing');
+ const readme = root.find(e => e.name === 'README');
+ assert.equal(readme.dir, false);
+ assert.ok(readme.qid && typeof readme.qid.path === 'number' || typeof readme.qid.path === 'bigint');
+
+ // GET directory as HTML.
+ const htmlRes = await fetch(`${url}/fs/`, { headers: { accept: 'text/html' } });
+ assert.match(htmlRes.headers.get('content-type'), /text\/html/);
+ assert.match(await htmlRes.text(), /<li>README<\/li>/);
+
+ // GET file.
+ const rd = await fetch(`${url}/fs/README`);
+ assert.equal(rd.status, 200);
+ assert.match(rd.headers.get('content-type'), /text\/plain/);
+ assert.match(await rd.text(), /9P2000/);
+
+ // GET a dynamic (length-0) file streams its content.
+ const now = (await (await fetch(`${url}/fs/runtime/fn/now`)).text()).trim();
+ assert.match(now, /^\d+$/, 'runtime/fn/now');
+
+ // HEAD.
+ const hd = await fetch(`${url}/fs/README`, { method: 'HEAD' });
+ assert.equal(hd.status, 200);
+ assert.ok(Number(hd.headers.get('content-length')) > 0);
+ assert.ok(hd.headers.get('x-9p-qid'));
+
+ // Range.
+ const rg = await fetch(`${url}/fs/README`, { headers: { range: 'bytes=0-4' } });
+ assert.equal(rg.status, 206);
+ assert.equal((await rg.text()).length, 5);
+
+ // 404 mapping.
+ assert.equal((await fetch(`${url}/fs/nope/missing`)).status, 404);
+
+ // PUT create + GET + truncate + append + DELETE round trip on /scratch.
+ let put = await fetch(`${url}/fs/scratch/rt.txt`, { method: 'PUT', body: 'hello mux' });
+ assert.equal(put.status, 201);
+ assert.equal(await (await fetch(`${url}/fs/scratch/rt.txt`)).text(), 'hello mux');
+ put = await fetch(`${url}/fs/scratch/rt.txt`, { method: 'PUT', body: 'AA' });
+ assert.equal(put.status, 200); // existed -> truncate
+ assert.equal(await (await fetch(`${url}/fs/scratch/rt.txt`)).text(), 'AA');
+ await fetch(`${url}/fs/scratch/rt.txt?append=1`, { method: 'PUT', body: 'BB' });
+ assert.equal(await (await fetch(`${url}/fs/scratch/rt.txt`)).text(), 'AABB');
+ assert.equal((await fetch(`${url}/fs/scratch/rt.txt`, { method: 'DELETE' })).status, 204);
+ assert.equal((await fetch(`${url}/fs/scratch/rt.txt`)).status, 404);
+
+ // mkdir via trailing slash.
+ assert.equal((await fetch(`${url}/fs/scratch/sub/`, { method: 'PUT' })).status, 201);
+ assert.equal((await fetch(`${url}/fs/scratch/sub`)).status, 200);
+
+ // Larger body: multi-frame streamed PUT then GET byte-exact.
+ const big = Buffer.alloc(150000);
+ for (let i = 0; i < big.length; i++) big[i] = (i * 17) % 251;
+ assert.equal((await fetch(`${url}/fs/scratch/big.bin`, { method: 'PUT', body: big })).status, 201);
+ const back = Buffer.from(await (await fetch(`${url}/fs/scratch/big.bin`)).arrayBuffer());
+ assert.ok(back.equals(big), 'byte-exact large round trip');
+
+ // SSE: follow a file streams its content as data: events.
+ const ac = new AbortController();
+ const sse = await fetch(`${url}/fs/README?follow=1`, { signal: ac.signal });
+ assert.match(sse.headers.get('content-type'), /text\/event-stream/);
+ const reader = sse.body.getReader();
+ let text = '';
+ while (text.length < 40) {
+ const { value, done } = await reader.read();
+ if (done) break;
+ text += Buffer.from(value).toString();
+ }
+ ac.abort();
+ assert.match(text, /data: /, 'SSE data event');
+
+ // Fan-out: many concurrent requests share the one upstream (16 tags).
+ const results = await Promise.all(Array.from({ length: 16 }, () =>
+ fetch(`${url}/fs/runtime/fn/now`).then(r => r.status).catch(e => 'ERR:' + e.message)));
+ const bad = results.filter(s => s !== 200);
+ assert.ok(bad.length === 0, 'concurrent fan-out, non-200: ' + JSON.stringify(bad));
+
+ console.log('http_fs: all checks passed');
+} catch (e) {
+ console.error(e.stack || e);
+ for (const c of children) if (c.err) console.error(c.spawnargs.join(' '), '\n', c.err.slice(-2000));
+ process.exitCode = 1;
+} finally {
+ for (const c of children) if (c.exitCode === null) c.kill('SIGTERM');
+ await delay(200);
+ for (const c of children) if (c.exitCode === null) c.kill('SIGKILL');
+ if (sock) await fs.rm(path.dirname(sock), { recursive: true, force: true }).catch(() => {});
+}
diff --git a/test/web/mux.zig b/test/web/mux.zig
new file mode 100644
index 0000000..ab95281
--- /dev/null
+++ b/test/web/mux.zig
@@ -0,0 +1,310 @@
+//! Unit tests for the 9web multiplexer (web/mux.zig): one shared upstream 9P
+//! connection fanned out to several downstream sessions. The upstream is an
+//! in-process `serve.Runner` over an `fs.Server` backend with a blocking
+//! "event" file (its reads park until a write to "data" wakes them). The tests
+//! drive the mux with hand-built 9P frames and check fid isolation, a parked
+//! read released by a 9P-side write, Tflush forwarding, and reconnection with
+//! downstream fids invalidated.
+const std = @import("std");
+const c9 = @import("cloud9");
+const mux = @import("mux");
+const wire = c9.wire;
+const fs = c9.fs;
+const serve = c9.serve;
+const transport = c9.transport;
+const Io = std.Io;
+const testing = std.testing;
+
+const msize = 8192;
+const M = mux.Mux(.{ .msize = msize, .upstream_fids = 64, .fids_per_conn = 32, .downstreams = 8 });
+
+// -- the upstream backend: root/{event,data} ---------------------------------
+
+const root_node = 1;
+const event_node = 2;
+const data_node = 3;
+
+const EventFs = struct {
+ pub const Req = fs.Req;
+ pub const Reply = fs.Reply;
+
+ mutex: Io.Mutex = .init,
+ posted: ?[]const u8 = null,
+ store: [256]u8 = undefined,
+
+ fn attrOf(node: u64) fs.Attr {
+ return switch (node) {
+ root_node => .{ .name = "/", .node = root_node, .dir = true, .mode = 0o500 },
+ event_node => .{ .name = "event", .node = event_node, .mode = 0o400 },
+ data_node => .{ .name = "data", .node = data_node, .mode = 0o600 },
+ else => unreachable,
+ };
+ }
+};
+
+const Runner = serve.Runner(EventFs, .{ .fid_capacity = 32, .slot_capacity = 8 }, .{ .msize = msize, .connections = 2, .listeners = 1 });
+
+fn serveFn(ctx: ?*anyopaque, conn: *Runner.Conn, req: fs.Req) void {
+ const st: *EventFs = @ptrCast(@alignCast(ctx.?));
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ const fail: fs.Reply = .fail(req.tag, fs.E.NOENT);
+ switch (req.op) {
+ .lookup => {
+ if (std.mem.eql(u8, req.data, "..")) return conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(root_node) }, "");
+ if (std.mem.eql(u8, req.data, "event")) return conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(event_node) }, "");
+ if (std.mem.eql(u8, req.data, "data")) return conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(data_node) }, "");
+ return conn.reply(&fail, "");
+ },
+ .getattr, .setattr => conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(req.node) }, ""),
+ .open => conn.reply(&.{ .tag = req.tag, .handle = 7 }, ""),
+ .release => conn.reply(&.{ .tag = req.tag }, ""),
+ .readdir => conn.reply(&.{ .tag = req.tag }, ""),
+ .read => {
+ if (req.node == event_node) {
+ if (st.posted) |bytes| {
+ st.posted = null;
+ return conn.reply(&.{ .tag = req.tag }, bytes);
+ }
+ return conn.reply(&.{ .tag = req.tag, .status = .again }, "");
+ }
+ conn.reply(&.{ .tag = req.tag }, ""); // data reads EOF
+ },
+ .write => {
+ const n = @min(req.data.len, st.store.len);
+ @memcpy(st.store[0..n], req.data[0..n]);
+ st.posted = st.store[0..n];
+ conn.reply(&.{ .tag = req.tag, .written = @intCast(n) }, "");
+ conn.wake(); // retry the parked event read on this connection
+ },
+ }
+}
+
+// -- a downstream session that captures replies -------------------------------
+
+const Down = struct {
+ m: *M,
+ conn: M.Conn = undefined,
+ mutex: Io.Mutex = .init,
+ cond: Io.Condition = .init,
+ ring: [8][msize]u8 = undefined,
+ len: [8]usize = @splat(0),
+ head: usize = 0,
+ tail: usize = 0,
+ count: usize = 0,
+
+ fn init(d: *Down, m: *M) void {
+ d.* = .{ .m = m };
+ d.conn = M.Conn.init(m, .{ .ctx = d, .send = sink });
+ }
+ fn deinit(d: *Down) void {
+ d.conn.deinit();
+ }
+ fn sink(ctx: *anyopaque, frame: []const u8) anyerror!void {
+ const d: *Down = @ptrCast(@alignCast(ctx));
+ d.mutex.lockUncancelable(testing.io);
+ defer d.mutex.unlock(testing.io);
+ const i = d.tail;
+ const n = @min(frame.len, msize);
+ @memcpy(d.ring[i][0..n], frame[0..n]);
+ d.len[i] = n;
+ d.tail = (i + 1) % d.ring.len;
+ d.count += 1;
+ d.cond.signal(testing.io);
+ }
+ fn send(d: *Down, frame: []const u8) void {
+ d.m.forward(&d.conn, frame);
+ }
+ fn recv(d: *Down) !wire.Decoded {
+ d.mutex.lockUncancelable(testing.io);
+ defer d.mutex.unlock(testing.io);
+ var tries: usize = 0;
+ while (d.count == 0) : (tries += 1) {
+ if (tries > 4000) return error.Timeout;
+ d.cond.wait(testing.io, &d.mutex) catch return error.Canceled;
+ }
+ const i = d.head;
+ d.head = (i + 1) % d.ring.len;
+ d.count -= 1;
+ return wire.decode(d.ring[i][0..d.len[i]]);
+ }
+};
+
+// -- frame builders -----------------------------------------------------------
+
+var enc_buf: [msize]u8 = undefined;
+fn enc(msg: wire.Msg, tag: u16) []const u8 {
+ return wire.encode(msg, tag, &enc_buf) catch unreachable;
+}
+
+fn handshake(d: *Down) !void {
+ d.send(enc(.{ .tversion = .{ .msize = msize, .version = "9P2000" } }, wire.notag));
+ const v = try d.recv();
+ try testing.expectEqual(wire.Type.rversion, v.msg.msgType());
+ d.send(enc(.{ .tattach = .{ .fid = 0, .afid = wire.nofid, .uname = "u", .aname = "" } }, 1));
+ const a = try d.recv();
+ try testing.expectEqual(wire.Type.rattach, a.msg.msgType());
+}
+
+fn walk1(d: *Down, newfid: u32, name: []const u8) !void {
+ var wn: [wire.max_welem][]const u8 = @splat("");
+ wn[0] = name;
+ d.send(enc(.{ .twalk = .{ .fid = 0, .newfid = newfid, .nwname = 1, .wname = wn } }, 2));
+ const r = try d.recv();
+ try testing.expectEqual(wire.Type.rwalk, r.msg.msgType());
+ try testing.expectEqual(@as(u16, 1), r.msg.rwalk.nwqid);
+}
+
+// -- the rig ------------------------------------------------------------------
+
+const Rig = struct {
+ dir: testing.TmpDir,
+ path_buf: [std.fs.max_path_bytes]u8 = undefined,
+ sock_buf: [transport.sun_path_len]u8 = undefined,
+ sock: [:0]const u8 = undefined,
+ fsx: EventFs = .{},
+ runner: *Runner = undefined,
+ m: *M = undefined,
+ reader: Io.Future(void) = undefined,
+
+ fn start(r: *Rig) !void {
+ const io = testing.io;
+ r.dir = testing.tmpDir(.{});
+ const plen = try r.dir.dir.realPath(io, &r.path_buf);
+ r.sock = try std.fmt.bufPrintSentinel(&r.sock_buf, "{s}/up", .{r.path_buf[0..plen]}, 0);
+ r.runner = try testing.allocator.create(Runner);
+ r.runner.init(.{ .io = io, .root = root_node, .handler = .{ .serve = serveFn, .ctx = &r.fsx } });
+ _ = try r.runner.listen(.{ .unix = r.sock }, 4);
+ r.m = try testing.allocator.create(M);
+ r.m.init(io, .{ .network = .{ .unix = r.sock } }, "u", "");
+ try r.m.connect();
+ r.reader = try io.concurrent(M.readerLoop, .{r.m});
+ }
+
+ fn stopUpstream(r: *Rig) void {
+ r.runner.stop();
+ Io.Dir.cwd().deleteFile(testing.io, r.sock) catch {};
+ }
+
+ fn restartUpstream(r: *Rig) !void {
+ r.runner.init(.{ .io = testing.io, .root = root_node, .handler = .{ .serve = serveFn, .ctx = &r.fsx } });
+ _ = try r.runner.listen(.{ .unix = r.sock }, 4);
+ }
+
+ fn end(r: *Rig) void {
+ r.m.stop();
+ r.reader.cancel(testing.io);
+ r.runner.stop();
+ testing.allocator.destroy(r.runner);
+ testing.allocator.destroy(r.m);
+ r.dir.cleanup();
+ }
+};
+
+test "mux: two downstreams share one upstream with isolated fids" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ var b: Down = undefined;
+ b.init(rig.m);
+ defer b.deinit();
+ try handshake(&a);
+ try handshake(&b);
+ // Both use downstream fid 1; the mux maps them to distinct upstream fids.
+ try walk1(&a, 1, "event");
+ try walk1(&b, 1, "data");
+ // A opens its event; B opens its data. Independent handles.
+ a.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.oread } }, 3));
+ try testing.expectEqual(wire.Type.ropen, (try a.recv()).msg.msgType());
+ b.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.owrite } }, 3));
+ try testing.expectEqual(wire.Type.ropen, (try b.recv()).msg.msgType());
+}
+
+test "mux: a parked read is released by another downstream's write" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ var b: Down = undefined;
+ b.init(rig.m);
+ defer b.deinit();
+ try handshake(&a);
+ try handshake(&b);
+ try walk1(&a, 1, "event");
+ try walk1(&b, 1, "data");
+ a.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.oread } }, 3));
+ _ = try a.recv();
+ b.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.owrite } }, 3));
+ _ = try b.recv();
+ // A's read parks upstream (no reply yet).
+ a.send(enc(.{ .tread = .{ .fid = 1, .offset = 0, .count = 128 } }, 4));
+ try testing.io.sleep(.fromMilliseconds(30), .awake);
+ // B writes "hello", which posts the event and wakes A's parked read.
+ b.send(enc(.{ .twrite = .{ .fid = 1, .offset = 0, .data = "hello" } }, 4));
+ try testing.expectEqual(wire.Type.rwrite, (try b.recv()).msg.msgType());
+ const rr = try a.recv();
+ try testing.expectEqual(wire.Type.rread, rr.msg.msgType());
+ try testing.expectEqualStrings("hello", rr.msg.rread.data);
+}
+
+test "mux: Tflush is forwarded and cancels a parked read" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ try handshake(&a);
+ try walk1(&a, 1, "event");
+ a.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.oread } }, 3));
+ _ = try a.recv();
+ a.send(enc(.{ .tread = .{ .fid = 1, .offset = 0, .count = 128 } }, 5));
+ try testing.io.sleep(.fromMilliseconds(30), .awake);
+ a.send(enc(.{ .tflush = .{ .oldtag = 5 } }, 6));
+ // Expect the interrupted read (Rerror) and the Rflush, in either order.
+ var saw_err = false;
+ var saw_flush = false;
+ for (0..2) |_| {
+ const got = try a.recv();
+ switch (got.msg.msgType()) {
+ .rerror => saw_err = true,
+ .rflush => saw_flush = true,
+ else => return error.Unexpected,
+ }
+ }
+ try testing.expect(saw_err and saw_flush);
+}
+
+test "mux: an upstream drop errors downstream, and it reconnects" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ try handshake(&a);
+ try walk1(&a, 1, "event");
+ // Drop the upstream. An in-flight/next request must error, not hang.
+ rig.stopUpstream();
+ try testing.io.sleep(.fromMilliseconds(50), .awake);
+ a.send(enc(.{ .tstat = .{ .fid = 1 } }, 7));
+ const err = try a.recv();
+ try testing.expectEqual(wire.Type.rerror, err.msg.msgType());
+ // Bring the upstream back; the mux reconnects and a fresh attach works.
+ try rig.restartUpstream();
+ var tries: usize = 0;
+ while (tries < 200) : (tries += 1) {
+ try testing.io.sleep(.fromMilliseconds(20), .awake);
+ if (std.mem.eql(u8, rig.m.upstreamState(), "connected")) break;
+ }
+ var b: Down = undefined;
+ b.init(rig.m);
+ defer b.deinit();
+ try handshake(&b); // a new session on the reconnected upstream
+ try testing.expect(rig.m.reconnectCount() >= 1);
+}
diff --git a/web/httpfs.zig b/web/httpfs.zig
new file mode 100644
index 0000000..fa783e5
--- /dev/null
+++ b/web/httpfs.zig
@@ -0,0 +1,583 @@
+//! An HTTP view of the upstream 9P tree under `/fs/<path>`, served through the
+//! shared `Mux` at the fid level. It is file operations, not verbs: HTTP
+//! methods map onto the tree.
+//!
+//! GET /fs/<file> bytes (Content-Type sniffed; Range honored)
+//! GET /fs/<file>?follow=1 server-sent events, one per blocking read
+//! GET /fs/<dir> JSON array of entries (HTML with Accept: text/html)
+//! HEAD /fs/<path> stat as headers
+//! PUT /fs/<file> write (create if missing; truncate unless ?append=1)
+//! PUT /fs/<dir>/ create a directory
+//! DELETE /fs/<path> remove
+//!
+//! A 9P Rerror becomes 404/403/409/500 with the ename in the body. Two
+//! caller-owned scratch buffers are used: `big` (>= msize) for read/dir/body
+//! data, `small` (>= a few KiB) for control replies and response framing.
+const std = @import("std");
+const c9 = @import("cloud9");
+const wire = c9.wire;
+const Server = std.http.Server;
+const Status = std.http.Status;
+const Request = c9.Client.Request;
+
+const max_components = 64;
+const name_bytes = 8192;
+
+pub fn handle(m: anytype, request: *Server.Request, path: []const u8, big: []u8, small: []u8) void {
+ const M = @TypeOf(m.*);
+ var fs: FS(M) = .{ .m = m, .request = request, .big = big, .small = small };
+ fs.serve(path) catch |err| switch (err) {
+ error.Handled => {},
+ error.NotFound => fs.fail(.not_found, "No such file or directory\n"),
+ error.BadPath => fs.fail(.bad_request, "Invalid path\n"),
+ error.Upstream, error.Disconnected => fs.fail(.bad_gateway, "9P upstream unavailable\n"),
+ else => fs.fail(.internal_server_error, "Internal error\n"),
+ };
+ fs.cleanup();
+}
+
+fn FS(comptime M: type) type {
+ return struct {
+ const Self = @This();
+
+ m: *M,
+ request: *Server.Request,
+ big: []u8,
+ small: []u8,
+ root_fid: ?u32 = null,
+ leaf_fid: ?u32 = null,
+ responded: bool = false,
+
+ fn rpcSmall(self: *Self, req: Request) !wire.Decoded {
+ var r: M.Rpc = .{ .buf = self.small };
+ return self.m.rpc(req, &r);
+ }
+
+ fn fail(self: *Self, status: Status, text: []const u8) void {
+ if (self.responded) return;
+ self.responded = true;
+ self.request.respond(text, .{ .status = status, .keep_alive = false, .extra_headers = &.{
+ .{ .name = "content-type", .value = "text/plain; charset=utf-8" },
+ .{ .name = "cache-control", .value = "no-store" },
+ .{ .name = "x-content-type-options", .value = "nosniff" },
+ } }) catch {};
+ }
+
+ /// A 9P Rerror: map to a status and answer, then stop the request.
+ fn rerr(self: *Self, ename: []const u8) anyerror {
+ var buf: [160]u8 = undefined;
+ const body = std.fmt.bufPrint(&buf, "{s}\n", .{ename[0..@min(ename.len, 128)]}) catch "error\n";
+ self.fail(mapError(ename), body);
+ return error.Handled;
+ }
+
+ fn cleanup(self: *Self) void {
+ if (self.leaf_fid) |f| {
+ _ = self.rpcSmall(.{ .clunk = .{ .fid = f } }) catch {};
+ self.leaf_fid = null;
+ }
+ if (self.root_fid) |f| {
+ _ = self.rpcSmall(.{ .clunk = .{ .fid = f } }) catch {};
+ self.root_fid = null;
+ }
+ }
+
+ // -- session helpers ----------------------------------------------
+
+ fn attach(self: *Self) !void {
+ if (self.root_fid != null) return;
+ const fid = self.m.takeFid() orelse return error.Upstream;
+ const got = try self.rpcSmall(.{ .attach = .{ .fid = fid, .uname = self.m.user, .aname = self.m.tree } });
+ if (got.msg == .rerror) {
+ self.m.dropFid(fid);
+ return self.rerr(got.msg.rerror.ename);
+ }
+ if (got.msg != .rattach) {
+ self.m.dropFid(fid);
+ return error.Upstream;
+ }
+ self.root_fid = fid;
+ }
+
+ /// Walks from the root to `names`, leaving `leaf_fid` on the target and
+ /// returning its stat. A missing component is `error.NotFound`.
+ fn walkTo(self: *Self, names: [][]const u8) !wire.Stat {
+ try self.attach();
+ if (self.leaf_fid) |old| {
+ _ = self.rpcSmall(.{ .clunk = .{ .fid = old } }) catch {};
+ self.leaf_fid = null;
+ }
+ const leaf = self.m.takeFid() orelse return error.Upstream;
+ var start: usize = 0;
+ var src = self.root_fid.?;
+ var first = true;
+ while (first or start < names.len) {
+ const batch = names[start..@min(start + 16, names.len)];
+ const got = try self.rpcSmall(.{ .walk = .{ .fid = src, .newfid = leaf, .names = batch } });
+ if (got.msg == .rerror) {
+ self.m.dropFid(leaf);
+ return self.rerr(got.msg.rerror.ename);
+ }
+ if (got.msg != .rwalk or got.msg.rwalk.nwqid != batch.len) {
+ self.m.dropFid(leaf); // short walk: newfid was not bound
+ return error.NotFound;
+ }
+ self.leaf_fid = leaf;
+ src = leaf;
+ start += batch.len;
+ first = false;
+ if (batch.len == 0) break;
+ }
+ const st = try self.rpcSmall(.{ .stat = .{ .fid = leaf } });
+ if (st.msg == .rerror) return self.rerr(st.msg.rerror.ename);
+ if (st.msg != .rstat) return error.Upstream;
+ return st.msg.rstat.stat;
+ }
+
+ // -- dispatch ------------------------------------------------------
+
+ fn serve(self: *Self, path: []const u8) !void {
+ const rel = path[3..]; // after "/fs"
+ const trailing = rel.len != 0 and rel[rel.len - 1] == '/';
+ var storage: [name_bytes]u8 = undefined;
+ var names: [max_components][]const u8 = undefined;
+ const count = try decodePath(rel, &storage, &names);
+ const comps = names[0..count];
+
+ switch (self.request.head.method) {
+ .GET => try self.get(comps),
+ .HEAD => try self.head(comps),
+ .PUT => if (trailing) try self.mkdir(comps) else try self.put(comps),
+ .DELETE => try self.delete(comps),
+ else => self.fail(.method_not_allowed, "Method not allowed\n"),
+ }
+ }
+
+ fn head(self: *Self, comps: [][]const u8) !void {
+ const st = try self.walkTo(comps);
+ const dir = st.qid.type & c9.qtdir != 0;
+ var len_buf: [24]u8 = undefined;
+ const clen = std.fmt.bufPrint(&len_buf, "{d}", .{st.length}) catch "0";
+ var qid_buf: [48]u8 = undefined;
+ const qid = std.fmt.bufPrint(&qid_buf, "{d}.{d}.{d}", .{ st.qid.type, st.qid.version, st.qid.path }) catch "";
+ var mode_buf: [16]u8 = undefined;
+ const mode = std.fmt.bufPrint(&mode_buf, "{o}", .{st.mode}) catch "";
+ var mtime_buf: [16]u8 = undefined;
+ const mtime = std.fmt.bufPrint(&mtime_buf, "{d}", .{st.mtime}) catch "";
+ self.responded = true;
+ self.request.respond("", .{ .status = .ok, .keep_alive = false, .transfer_encoding = .none, .extra_headers = &.{
+ .{ .name = "content-type", .value = if (dir) "application/json" else "application/octet-stream" },
+ .{ .name = "content-length", .value = clen },
+ .{ .name = "x-9p-qid", .value = qid },
+ .{ .name = "x-9p-mode", .value = mode },
+ .{ .name = "x-9p-mtime", .value = mtime },
+ .{ .name = "accept-ranges", .value = "bytes" },
+ .{ .name = "cache-control", .value = "no-store" },
+ } }) catch {};
+ }
+
+ fn get(self: *Self, comps: [][]const u8) !void {
+ const follow = self.wantsFollow();
+ const st = try self.walkTo(comps);
+ if (st.qid.type & c9.qtdir != 0) return self.listDir(st);
+ if (follow) return self.followFile();
+ return self.getFile(st);
+ }
+
+ // -- files ---------------------------------------------------------
+
+ fn getFile(self: *Self, st: wire.Stat) !void {
+ const leaf = self.leaf_fid.?;
+ const opened = try self.rpcSmall(.{ .open = .{ .fid = leaf, .mode = c9.oread } });
+ if (opened.msg == .rerror) return self.rerr(opened.msg.rerror.ename);
+ if (opened.msg != .ropen) return error.Upstream;
+
+ const range = self.parseRange(st.length);
+ const status: Status = if (range.partial) .partial_content else .ok;
+
+ var offset: u64 = range.start;
+ const first = try self.readAt(leaf, offset, self.chunk());
+ if (first.msg == .rerror) return self.rerr(first.msg.rerror.ename);
+ if (first.msg != .rread) return error.Upstream;
+ const ctype = sniff(first.msg.rread.data);
+
+ var cr_buf: [64]u8 = undefined;
+ var headers: [4]std.http.Header = undefined;
+ var nh: usize = 0;
+ headers[nh] = .{ .name = "content-type", .value = ctype };
+ nh += 1;
+ headers[nh] = .{ .name = "accept-ranges", .value = "bytes" };
+ nh += 1;
+ headers[nh] = .{ .name = "cache-control", .value = "no-store" };
+ nh += 1;
+ if (range.partial) {
+ headers[nh] = .{ .name = "content-range", .value = std.fmt.bufPrint(&cr_buf, "bytes {d}-{d}/{d}", .{ range.start, range.end - 1, st.length }) catch return error.Upstream };
+ nh += 1;
+ }
+
+ self.responded = true;
+ var bw = try self.request.respondStreaming(self.small, .{
+ .respond_options = .{ .status = status, .keep_alive = false, .extra_headers = headers[0..nh] },
+ });
+ // A synthetic/control file reports length 0 but streams content, so
+ // clamp to the length only for an explicit Range.
+ const limit: u64 = if (range.partial) range.end - range.start else std.math.maxInt(u64);
+ var sent: u64 = 0;
+ var data = first.msg.rread.data;
+ while (true) {
+ const take = @min(@as(u64, data.len), limit - sent);
+ try bw.writer.writeAll(data[0..@intCast(take)]);
+ sent += take;
+ if (sent >= limit or data.len == 0) break;
+ offset += data.len;
+ const r = try self.readAt(leaf, offset, self.chunk());
+ if (r.msg != .rread) break;
+ data = r.msg.rread.data;
+ if (data.len == 0) break;
+ }
+ try bw.end();
+ }
+
+ fn followFile(self: *Self) !void {
+ const leaf = self.leaf_fid.?;
+ const opened = try self.rpcSmall(.{ .open = .{ .fid = leaf, .mode = c9.oread } });
+ if (opened.msg == .rerror) return self.rerr(opened.msg.rerror.ename);
+ if (opened.msg != .ropen) return error.Upstream;
+
+ self.responded = true;
+ var bw = try self.request.respondStreaming(self.small, .{
+ .respond_options = .{ .status = .ok, .keep_alive = false, .extra_headers = &.{
+ .{ .name = "content-type", .value = "text/event-stream; charset=utf-8" },
+ .{ .name = "cache-control", .value = "no-store" },
+ .{ .name = "x-content-type-options", .value = "nosniff" },
+ } },
+ });
+ bw.writer.writeAll(": follow\n\n") catch return;
+ flushBody(&bw) catch return;
+ var offset: u64 = 0;
+ while (true) {
+ var r: M.Rpc = .{ .buf = self.big };
+ // On client disconnect / deadline this rpc is cancelled; the
+ // mux Tflushes the parked upstream read before returning.
+ const got = self.m.rpc(.{ .read = .{ .fid = leaf, .offset = offset, .count = self.chunk() } }, &r) catch return;
+ if (got.msg != .rread) return;
+ const data = got.msg.rread.data;
+ if (data.len == 0) {
+ bw.writer.writeAll("event: eof\ndata:\n\n") catch return;
+ flushBody(&bw) catch return;
+ return;
+ }
+ offset += data.len;
+ self.writeSse(&bw, data) catch return;
+ flushBody(&bw) catch return;
+ }
+ }
+
+ /// Drains the body writer's buffer as a chunk, then flushes to the peer.
+ fn flushBody(bw: anytype) !void {
+ try bw.writer.flush();
+ try bw.flush();
+ }
+
+ fn writeSse(self: *Self, bw: anytype, data: []const u8) !void {
+ _ = self;
+ var it = std.mem.splitScalar(u8, std.mem.trimEnd(u8, data, "\n"), '\n');
+ while (it.next()) |line| {
+ try bw.writer.writeAll("data: ");
+ try bw.writer.writeAll(line);
+ try bw.writer.writeByte('\n');
+ }
+ try bw.writer.writeByte('\n');
+ }
+
+ // -- directories ---------------------------------------------------
+
+ fn listDir(self: *Self, st: wire.Stat) !void {
+ const html = self.wantsHtml();
+ const leaf = self.leaf_fid.?;
+ const opened = try self.rpcSmall(.{ .open = .{ .fid = leaf, .mode = c9.oread } });
+ if (opened.msg == .rerror) return self.rerr(opened.msg.rerror.ename);
+ if (opened.msg != .ropen) return error.Upstream;
+
+ // JSON/HTML is built into `big`; directory data is read into `small`.
+ var w: std.Io.Writer = .fixed(self.big);
+ if (html) {
+ w.writeAll("<!doctype html><meta charset=utf-8><title>") catch return error.Upstream;
+ writeHtml(&w, st.name);
+ w.writeAll("</title><ul>") catch return error.Upstream;
+ } else {
+ w.writeByte('[') catch return error.Upstream;
+ }
+ var offset: u64 = 0;
+ var first = true;
+ while (true) {
+ var r: M.Rpc = .{ .buf = self.small };
+ const got = try self.m.rpc(.{ .read = .{ .fid = leaf, .offset = offset, .count = self.chunk() } }, &r);
+ if (got.msg != .rread) break;
+ const data = got.msg.rread.data;
+ if (data.len == 0) break;
+ offset += data.len;
+ var rest: []const u8 = data;
+ while (rest.len >= 2) {
+ const size: usize = @as(usize, std.mem.readInt(u16, rest[0..2], .little)) + 2;
+ if (size > rest.len) break;
+ const entry = c9.Stat.decode(rest[0..size]) catch break;
+ rest = rest[size..];
+ if (html) {
+ w.writeAll("<li>") catch return error.Upstream;
+ writeHtml(&w, entry.name);
+ if (entry.qid.type & c9.qtdir != 0) w.writeByte('/') catch {};
+ w.writeAll("</li>") catch return error.Upstream;
+ } else {
+ if (!first) w.writeByte(',') catch return error.Upstream;
+ first = false;
+ writeEntryJson(&w, entry) catch return error.Upstream;
+ }
+ }
+ }
+ if (html) w.writeAll("</ul>") catch return error.Upstream else w.writeByte(']') catch return error.Upstream;
+ self.responded = true;
+ self.request.respond(w.buffered(), .{ .status = .ok, .keep_alive = false, .extra_headers = &.{
+ .{ .name = "content-type", .value = if (html) "text/html; charset=utf-8" else "application/json" },
+ .{ .name = "cache-control", .value = "no-store" },
+ .{ .name = "x-content-type-options", .value = "nosniff" },
+ } }) catch {};
+ }
+
+ // -- writes --------------------------------------------------------
+
+ fn put(self: *Self, comps: [][]const u8) !void {
+ if (comps.len == 0) return error.BadPath;
+ const append = self.wantsAppend();
+ const existed = self.tryWalk(comps) catch |err| switch (err) {
+ error.NotFound => false,
+ else => return err,
+ };
+ var write_off: u64 = 0;
+ if (existed) {
+ const leaf = self.leaf_fid.?;
+ if (append) {
+ const st = try self.rpcSmall(.{ .stat = .{ .fid = leaf } });
+ if (st.msg == .rstat) write_off = st.msg.rstat.stat.length;
+ const o = try self.rpcSmall(.{ .open = .{ .fid = leaf, .mode = c9.owrite } });
+ if (o.msg == .rerror) return self.rerr(o.msg.rerror.ename);
+ if (o.msg != .ropen) return error.Upstream;
+ } else {
+ const o = try self.rpcSmall(.{ .open = .{ .fid = leaf, .mode = c9.owrite | c9.otrunc } });
+ if (o.msg == .rerror) return self.rerr(o.msg.rerror.ename);
+ if (o.msg != .ropen) return error.Upstream;
+ }
+ } else {
+ _ = try self.walkTo(comps[0 .. comps.len - 1]);
+ const parent = self.leaf_fid.?;
+ const cr = try self.rpcSmall(.{ .create = .{ .fid = parent, .name = comps[comps.len - 1], .perm = 0o644, .mode = c9.owrite } });
+ if (cr.msg == .rerror) return self.rerr(cr.msg.rerror.ename);
+ if (cr.msg != .rcreate) return error.Upstream;
+ }
+ const written = try self.streamBody(self.leaf_fid.?, write_off);
+ var buf: [48]u8 = undefined;
+ const body = std.fmt.bufPrint(&buf, "{{\"written\":{d}}}\n", .{written}) catch "{}\n";
+ self.responded = true;
+ self.request.respond(body, .{ .status = if (existed) .ok else .created, .keep_alive = false, .extra_headers = &.{
+ .{ .name = "content-type", .value = "application/json" },
+ .{ .name = "cache-control", .value = "no-store" },
+ } }) catch {};
+ }
+
+ fn mkdir(self: *Self, comps: [][]const u8) !void {
+ if (comps.len == 0) return error.BadPath;
+ _ = try self.walkTo(comps[0 .. comps.len - 1]);
+ const parent = self.leaf_fid.?;
+ const cr = try self.rpcSmall(.{ .create = .{ .fid = parent, .name = comps[comps.len - 1], .perm = c9.dmdir | 0o755, .mode = c9.oread } });
+ if (cr.msg == .rerror) return self.rerr(cr.msg.rerror.ename);
+ if (cr.msg != .rcreate) return error.Upstream;
+ self.fail(.created, "created\n");
+ }
+
+ fn delete(self: *Self, comps: [][]const u8) !void {
+ if (comps.len == 0) return error.BadPath;
+ _ = try self.walkTo(comps);
+ const leaf = self.leaf_fid.?;
+ const rm = try self.rpcSmall(.{ .remove = .{ .fid = leaf } });
+ self.m.dropFid(leaf); // Tremove clunks the fid upstream regardless
+ self.leaf_fid = null;
+ if (rm.msg == .rerror) return self.rerr(rm.msg.rerror.ename);
+ self.fail(.no_content, "");
+ }
+
+ /// walkTo that reports missing as `false` rather than an HTTP error.
+ fn tryWalk(self: *Self, comps: [][]const u8) !bool {
+ _ = self.walkTo(comps) catch |err| switch (err) {
+ error.NotFound => return false,
+ else => return err,
+ };
+ return true;
+ }
+
+ fn streamBody(self: *Self, fid: u32, start: u64) !u64 {
+ const body_buf = self.big[0..16384];
+ const cbuf = self.big[16384..];
+ const br = self.request.readerExpectContinue(body_buf) catch return error.Upstream;
+ const maxw = @min(self.m.negotiatedMsize() -| 23, cbuf.len);
+ var offset: u64 = start;
+ var total: u64 = 0;
+ while (true) {
+ const n = br.readSliceShort(cbuf[0..maxw]) catch return error.Upstream;
+ if (n == 0) break;
+ var done: usize = 0;
+ while (done < n) {
+ const w = try self.rpcSmall(.{ .write = .{ .fid = fid, .offset = offset + done, .data = cbuf[done..n] } });
+ if (w.msg == .rerror) return self.rerr(w.msg.rerror.ename);
+ if (w.msg != .rwrite or w.msg.rwrite.count == 0) return error.Upstream;
+ done += w.msg.rwrite.count;
+ }
+ offset += n;
+ total += n;
+ }
+ return total;
+ }
+
+ // -- small helpers -------------------------------------------------
+
+ fn readAt(self: *Self, fid: u32, offset: u64, count: u32) !wire.Decoded {
+ var r: M.Rpc = .{ .buf = self.big };
+ return self.m.rpc(.{ .read = .{ .fid = fid, .offset = offset, .count = count } }, &r);
+ }
+
+ fn chunk(self: *Self) u32 {
+ const ms = self.m.negotiatedMsize();
+ return if (ms > 11) ms - 11 else 0;
+ }
+
+ const Range = struct { start: u64, end: u64, partial: bool };
+ fn parseRange(self: *Self, length: u64) Range {
+ var it = self.request.iterateHeaders();
+ while (it.next()) |h| {
+ if (!std.ascii.eqlIgnoreCase(h.name, "range")) continue;
+ if (!std.mem.startsWith(u8, h.value, "bytes=")) break;
+ const spec = h.value[6..];
+ const dash = std.mem.indexOfScalar(u8, spec, '-') orelse break;
+ const start = std.fmt.parseInt(u64, spec[0..dash], 10) catch break;
+ var end: u64 = length;
+ if (dash + 1 < spec.len) {
+ const e = std.fmt.parseInt(u64, spec[dash + 1 ..], 10) catch break;
+ end = @min(e + 1, length);
+ }
+ if (start >= length or start >= end) break;
+ return .{ .start = start, .end = end, .partial = true };
+ }
+ return .{ .start = 0, .end = length, .partial = false };
+ }
+
+ fn wantsFollow(self: *Self) bool {
+ if (std.mem.indexOf(u8, self.request.head.target, "follow=1") != null) return true;
+ return self.acceptHas("text/event-stream");
+ }
+ fn wantsHtml(self: *Self) bool {
+ return self.acceptHas("text/html");
+ }
+ fn wantsAppend(self: *Self) bool {
+ return std.mem.indexOf(u8, self.request.head.target, "append=1") != null;
+ }
+ fn acceptHas(self: *Self, what: []const u8) bool {
+ var it = self.request.iterateHeaders();
+ while (it.next()) |h| {
+ if (std.ascii.eqlIgnoreCase(h.name, "accept") and std.mem.indexOf(u8, h.value, what) != null) return true;
+ }
+ return false;
+ }
+ };
+}
+
+fn mapError(ename: []const u8) Status {
+ if (contains(ename, "No such file") or contains(ename, "does not exist") or contains(ename, "not found") or contains(ename, "unknown")) return .not_found;
+ if (contains(ename, "permission") or contains(ename, "not permitted") or contains(ename, "denied")) return .forbidden;
+ if (contains(ename, "exists")) return .conflict;
+ if (contains(ename, "not empty") or contains(ename, "in use") or contains(ename, "Is a directory") or contains(ename, "Not a directory")) return .conflict;
+ return .internal_server_error;
+}
+
+fn contains(haystack: []const u8, needle: []const u8) bool {
+ return std.mem.indexOf(u8, haystack, needle) != null;
+}
+
+fn sniff(data: []const u8) []const u8 {
+ const n = @min(data.len, 1024);
+ for (data[0..n]) |b| if (b == 0) return "application/octet-stream";
+ return "text/plain; charset=utf-8";
+}
+
+/// Decodes a `/fs`-relative path into non-empty, percent-decoded components.
+fn decodePath(rel: []const u8, storage: []u8, names: *[max_components][]const u8) !usize {
+ var out: usize = 0;
+ var count: usize = 0;
+ var it = std.mem.splitScalar(u8, rel, '/');
+ while (it.next()) |raw| {
+ if (raw.len == 0) continue;
+ var i: usize = 0;
+ const begin = out;
+ while (i < raw.len) {
+ const ch = raw[i];
+ if (ch == '%') {
+ if (i + 2 >= raw.len) return error.BadPath;
+ const hi = hex(raw[i + 1]) orelse return error.BadPath;
+ const lo = hex(raw[i + 2]) orelse return error.BadPath;
+ const byte = (hi << 4) | lo;
+ if (byte == 0 or byte == '/') return error.BadPath;
+ if (out >= storage.len) return error.BadPath;
+ storage[out] = byte;
+ out += 1;
+ i += 3;
+ } else {
+ if (out >= storage.len) return error.BadPath;
+ storage[out] = ch;
+ out += 1;
+ i += 1;
+ }
+ }
+ const comp = storage[begin..out];
+ if (std.mem.eql(u8, comp, ".")) {
+ out = begin;
+ continue;
+ }
+ if (std.mem.eql(u8, comp, "..")) return error.BadPath;
+ if (count >= max_components) return error.BadPath;
+ names[count] = comp;
+ count += 1;
+ }
+ return count;
+}
+
+fn hex(c: u8) ?u8 {
+ return switch (c) {
+ '0'...'9' => c - '0',
+ 'a'...'f' => c - 'a' + 10,
+ 'A'...'F' => c - 'A' + 10,
+ else => null,
+ };
+}
+
+fn writeHtml(w: *std.Io.Writer, text: []const u8) void {
+ for (text) |c| switch (c) {
+ '<' => w.writeAll("&lt;") catch {},
+ '>' => w.writeAll("&gt;") catch {},
+ '&' => w.writeAll("&amp;") catch {},
+ '"' => w.writeAll("&quot;") catch {},
+ else => w.writeByte(c) catch {},
+ };
+}
+
+fn writeEntryJson(w: *std.Io.Writer, e: wire.Stat) !void {
+ const dir = e.qid.type & c9.qtdir != 0;
+ try w.writeAll("{\"name\":");
+ try std.json.Stringify.value(e.name, .{}, w);
+ try w.print(",\"dir\":{s},\"length\":{d},\"mode\":{d},\"mtime\":{d},\"qid\":{{\"type\":{d},\"version\":{d},\"path\":{d}}}}}", .{
+ if (dir) "true" else "false",
+ e.length,
+ e.mode,
+ e.mtime,
+ e.qid.type,
+ e.qid.version,
+ e.qid.path,
+ });
+}
diff --git a/web/main.zig b/web/main.zig
index e6aa421..7222d91 100644
--- a/web/main.zig
+++ b/web/main.zig
@@ -1,15 +1,29 @@
-//! HTTP application; no web or namespace policy is added to the cloud9 library.
+//! 9web: an HTTP/WebSocket gateway and 9P multiplexer.
+//!
+//! One shared upstream 9P connection (`mux.Mux`, all 16 tags) is fanned out to
+//! any number of downstream sessions: the browser WASM client over a WebSocket
+//! at `/_cloud9/9p`, plain 9P clients over an optional `--serve` TCP/Unix
+//! listener (a 9pserve-style mux for 9ns and plan9port), and an HTTP view of
+//! the tree under `/fs/<path>`. No web or namespace policy is added to the
+//! cloud9 library; this file owns listeners, connection bounds and routing.
const std = @import("std");
const c9 = @import("cloud9");
+const mux = @import("mux.zig");
+const httpfs = @import("httpfs.zig");
+const probe = @import("probe.zig");
const Io = std.Io;
+
const max_frame = 65536;
-const max_connections = 32;
-var serial_busy: std.atomic.Value(bool) = .init(false);
-const Upstream = union(enum) { network: c9.transport.Address, file: []const u8 };
-var connections: std.atomic.Value(u32) = .init(0);
+const max_connections = 64;
+
+pub const MuxT = mux.Mux(.{
+ .msize = max_frame,
+ .upstream_fids = 2048,
+ .fids_per_conn = 256,
+ .downstreams = max_connections,
+});
const Config = struct {
- upstream: Upstream,
host: []const u8,
origin: []const u8,
public_host: []const u8,
@@ -17,11 +31,91 @@ const Config = struct {
browser_config: []const u8,
};
+/// One downstream slot: buffers, the HTTP/WebSocket streams and the mux
+/// session, all in stable memory for the process lifetime so the mux reader
+/// task may reference a slot's sink after the slot's socket has closed.
+const Session = struct {
+ used: std.atomic.Value(bool) = .init(false),
+ app: *App = undefined,
+ stream: Io.net.Stream = undefined,
+ reader: Io.net.Stream.Reader = undefined,
+ writer: Io.net.Stream.Writer = undefined,
+ in: [16384]u8 = undefined,
+ out: [16384]u8 = undefined,
+ frame: [max_frame]u8 = undefined,
+ body: [16384]u8 = undefined,
+ ws: c9.http.WebSocket = undefined,
+ conn: MuxT.Conn = undefined,
+ wmutex: Io.Mutex = .init,
+ is_ws: bool = false,
+
+ fn sink(s: *Session) mux.Sink {
+ return .{ .ctx = s, .send = sendReply };
+ }
+
+ fn sendReply(ctx: *anyopaque, frame: []const u8) anyerror!void {
+ const s: *Session = @ptrCast(@alignCast(ctx));
+ s.wmutex.lockUncancelable(s.app.io);
+ defer s.wmutex.unlock(s.app.io);
+ if (s.is_ws) {
+ try s.ws.send(frame, .binary, null);
+ } else {
+ try c9.transport.writeFrame(&s.writer.interface, frame, max_frame);
+ try s.writer.interface.flush();
+ }
+ }
+};
+
+pub const App = struct {
+ io: Io,
+ m: *MuxT,
+ config: Config,
+ sessions: [max_connections]Session = @splat(.{}),
+ active: std.atomic.Value(u32) = .init(0),
+
+ fn take(app: *App) ?*Session {
+ for (&app.sessions) |*s| {
+ if (s.used.cmpxchgStrong(false, true, .acq_rel, .monotonic) == null) {
+ _ = app.active.fetchAdd(1, .monotonic);
+ return s;
+ }
+ }
+ return null;
+ }
+ fn giveBack(app: *App, s: *Session) void {
+ s.used.store(false, .release);
+ _ = app.active.fetchSub(1, .monotonic);
+ }
+ pub fn connectionCount(app: *App) u32 {
+ return app.active.load(.acquire);
+ }
+};
+
+fn parseUpstream(allocator: std.mem.Allocator, text: []const u8) !mux.Address {
+ if (std.mem.startsWith(u8, text, "tcp:"))
+ return .{ .network = .{ .tcp = try Io.net.IpAddress.parseLiteral(text[4..]) } };
+ if (std.mem.startsWith(u8, text, "unix:"))
+ return .{ .network = .{ .unix = try allocator.dupeZ(u8, text[5..]) } };
+ if (std.mem.startsWith(u8, text, "file:"))
+ return .{ .file = text[5..] };
+ return error.InvalidUpstream;
+}
+
+fn parseServe(allocator: std.mem.Allocator, text: []const u8) !c9.transport.Address {
+ if (std.mem.startsWith(u8, text, "tcp:"))
+ return .{ .tcp = try Io.net.IpAddress.parseLiteral(text[4..]) };
+ if (std.mem.startsWith(u8, text, "unix:"))
+ return .{ .unix = try allocator.dupeZ(u8, text[5..]) };
+ return error.InvalidServe;
+}
+
pub fn main(init: std.process.Init) !void {
const allocator = init.arena.allocator();
const args = try init.minimal.args.toSlice(allocator);
var listen_text: []const u8 = "127.0.0.1:8080";
var upstream_text: []const u8 = "tcp:127.0.0.1:564";
+ var serve_text: ?[]const u8 = null;
+ var probe_text: ?[]const u8 = null;
var timeout_ms: u32 = 300000;
var user: []const u8 = "user";
var tree: []const u8 = "";
@@ -29,115 +123,192 @@ pub fn main(init: std.process.Init) !void {
var i: usize = 1;
while (i < args.len) : (i += 1) {
if (std.mem.eql(u8, args[i], "--help")) {
- std.debug.print("Usage: 9web [--listen IP:PORT] [--upstream tcp:IP:PORT|unix:PATH|file:PATH] [--origin https://HOST:PORT] [--timeout-ms 300000] [--user NAME] [--tree NAME]\nDefault: http://127.0.0.1:8080 -> tcp:127.0.0.1:564\n", .{});
+ std.debug.print(
+ \\Usage: 9web [--listen IP:PORT] [--upstream tcp:IP:PORT|unix:PATH|file:PATH]
+ \\ [--serve tcp:IP:PORT|unix:PATH] [--probe unix:PATH|tcp:IP:PORT]
+ \\ [--origin https://HOST:PORT] [--timeout-ms 300000] [--user NAME] [--tree NAME]
+ \\One shared upstream, many downstreams. Default: http://127.0.0.1:8080 -> tcp:127.0.0.1:564
+ \\
+ , .{});
return;
}
if (i + 1 == args.len) return error.MissingArgument;
- if (std.mem.eql(u8, args[i], "--listen")) {
- i += 1;
- listen_text = args[i];
- } else if (std.mem.eql(u8, args[i], "--upstream")) {
- i += 1;
- upstream_text = args[i];
- } else if (std.mem.eql(u8, args[i], "--origin")) {
- i += 1;
- public_origin = args[i];
- } else if (std.mem.eql(u8, args[i], "--timeout-ms")) {
- i += 1;
- timeout_ms = try std.fmt.parseInt(u32, args[i], 10);
- } else if (std.mem.eql(u8, args[i], "--user")) {
- i += 1;
- user = args[i];
- } else if (std.mem.eql(u8, args[i], "--tree")) {
- i += 1;
- tree = args[i];
- } else return error.UnknownArgument;
+ i += 1;
+ const v = args[i];
+ const k = args[i - 1];
+ if (std.mem.eql(u8, k, "--listen")) listen_text = v else if (std.mem.eql(u8, k, "--upstream")) upstream_text = v else if (std.mem.eql(u8, k, "--serve")) serve_text = v else if (std.mem.eql(u8, k, "--probe")) probe_text = v else if (std.mem.eql(u8, k, "--origin")) public_origin = v else if (std.mem.eql(u8, k, "--timeout-ms")) timeout_ms = try std.fmt.parseInt(u32, v, 10) else if (std.mem.eql(u8, k, "--user")) user = v else if (std.mem.eql(u8, k, "--tree")) tree = v else return error.UnknownArgument;
}
- const address = try Io.net.IpAddress.parseLiteral(listen_text);
- const upstream: Upstream = if (std.mem.startsWith(u8, upstream_text, "tcp:"))
- .{ .network = .{ .tcp = try Io.net.IpAddress.parseLiteral(upstream_text[4..]) } }
- else if (std.mem.startsWith(u8, upstream_text, "unix:"))
- .{ .network = .{ .unix = try allocator.dupeZ(u8, upstream_text[5..]) } }
- else if (std.mem.startsWith(u8, upstream_text, "file:"))
- .{ .file = upstream_text[5..] }
- else
- return error.InvalidUpstream;
+ if (user.len > 256 or tree.len > 256) return error.AttachNameTooLong;
+
const io = init.io;
+ const address = try Io.net.IpAddress.parseLiteral(listen_text);
+ const upstream = try parseUpstream(allocator, upstream_text);
+
+ var m: MuxT = undefined;
+ m.init(io, upstream, user, tree);
+ m.connect() catch {
+ std.debug.print("9web: upstream {s} unavailable at startup; will retry.\n", .{upstream_text});
+ };
+ defer m.stop();
+
var listener = c9.transport.listen(io, .{ .tcp = address }, max_connections) catch |err| {
if (err == error.AddressInUse) {
std.debug.print("9web: cannot listen on {s}: address already in use.\nUse --listen 127.0.0.1:0 to select a free port; the URL is printed at startup.\n", .{listen_text});
- } else {
- std.debug.print("9web: cannot listen on {s}: {s}\n", .{ listen_text, @errorName(err) });
- }
+ } else std.debug.print("9web: cannot listen on {s}: {s}\n", .{ listen_text, @errorName(err) });
return err;
};
defer listener.deinit(io);
+
const host = try std.fmt.allocPrint(allocator, "{f}", .{listener.socket.address});
const origin_text = public_origin orelse try std.fmt.allocPrint(allocator, "http://{s}", .{host});
const origin_uri = try std.Uri.parse(origin_text);
if ((!std.mem.eql(u8, origin_uri.scheme, "https") and !std.mem.eql(u8, origin_uri.scheme, "http")) or
origin_uri.host == null or origin_uri.user != null or origin_uri.password != null or
origin_uri.path.percent_encoded.len != 0 or origin_uri.query != null or origin_uri.fragment != null) return error.InvalidOrigin;
- if (user.len > 256 or tree.len > 256) return error.AttachNameTooLong;
const browser_config = try std.json.Stringify.valueAlloc(allocator, .{ .user = user, .tree = tree }, .{});
- const config: Config = .{ .upstream = upstream, .host = host, .origin = origin_text, .public_host = origin_text[origin_uri.scheme.len + 3 ..], .timeout_ms = timeout_ms, .browser_config = browser_config };
- std.debug.print("9web {s} -> {s}\n", .{ config.origin, upstream_text });
+
+ var app: App = .{ .io = io, .m = &m, .config = .{
+ .host = host,
+ .origin = origin_text,
+ .public_host = origin_text[origin_uri.scheme.len + 3 ..],
+ .timeout_ms = timeout_ms,
+ .browser_config = browser_config,
+ } };
+
+ std.debug.print("9web {s} -> {s}\n", .{ app.config.origin, upstream_text });
+
var group: Io.Group = .init;
defer group.cancel(io);
+ try group.concurrent(io, MuxT.readerLoop, .{&m});
+
+ // Optional 9P listener for plan9port / 9ns clients (a mux like 9pserve).
+ var serve_listener: ?Io.net.Server = null;
+ if (serve_text) |t| {
+ const addr = try parseServe(allocator, t);
+ serve_listener = c9.transport.listen(io, addr, max_connections) catch |err| {
+ std.debug.print("9web: cannot --serve {s}: {s}\n", .{ t, @errorName(err) });
+ return err;
+ };
+ std.debug.print("9web: serving 9P on {s}\n", .{t});
+ try group.concurrent(io, serveLoop, .{ &app, &serve_listener.?, &group });
+ }
+ defer if (serve_listener) |*l| l.deinit(io);
+
+ // Optional self-introspection: embed 9proc under --probe.
+ const Introspect = probe.Probe(App);
+ if (probe_text) |t| {
+ Introspect.start(&app, allocator, t) catch |err| {
+ std.debug.print("9web: --probe {s} failed: {s}\n", .{ t, @errorName(err) });
+ };
+ }
+ defer Introspect.stop();
+
while (true) {
- const stream = try listener.accept(io);
- if (connections.fetchAdd(1, .monotonic) >= max_connections) {
- _ = connections.fetchSub(1, .monotonic);
+ const stream = listener.accept(io) catch |err| switch (err) {
+ error.Canceled => return,
+ else => {
+ io.sleep(.fromMilliseconds(20), .awake) catch return;
+ continue;
+ },
+ };
+ const s = app.take() orelse {
stream.close(io);
continue;
- }
- group.concurrent(io, handle, .{ io, stream, config }) catch |err| {
- _ = connections.fetchSub(1, .monotonic);
+ };
+ s.app = &app;
+ s.stream = stream;
+ group.concurrent(io, handleHttp, .{s}) catch {
+ app.giveBack(s);
stream.close(io);
- return err;
+ continue;
};
}
}
-fn handle(io: Io, stream: Io.net.Stream, config: Config) void {
- defer _ = connections.fetchSub(1, .monotonic);
- defer stream.close(io);
- var work: Connection = .{ .io = io, .stream = stream, .config = config };
- var group: Io.Group = .init;
- defer group.cancel(io);
- group.concurrent(io, Connection.run, .{&work}) catch return;
- const timeout: Io.Timeout = if (config.timeout_ms == 0) .none else .{ .duration = .{ .raw = .fromMilliseconds(config.timeout_ms), .clock = .awake } };
- work.done.waitTimeout(io, timeout) catch {};
+/// Accepts plain 9P clients on the `--serve` listener; each runs a raw mux
+/// session for its lifetime.
+fn serveLoop(app: *App, listener: *Io.net.Server, group: *Io.Group) void {
+ const io = app.io;
+ while (true) {
+ const stream = listener.accept(io) catch return;
+ const s = app.take() orelse {
+ stream.close(io);
+ continue;
+ };
+ s.app = app;
+ s.stream = stream;
+ group.concurrent(io, rawSession, .{s}) catch {
+ app.giveBack(s);
+ stream.close(io);
+ };
+ }
}
-const Connection = struct {
- io: Io,
- stream: Io.net.Stream,
- config: Config,
- done: Io.Event = .unset,
- fn run(connection: *Connection) void {
- defer connection.done.set(connection.io);
- serve(connection.io, connection.stream, connection.config) catch {};
+/// A downstream 9P session over a raw stream (`--serve`).
+fn rawSession(s: *Session) void {
+ const app = s.app;
+ const io = app.io;
+ defer {
+ s.stream.close(io);
+ app.giveBack(s);
}
-};
+ s.is_ws = false;
+ s.reader = s.stream.reader(io, &s.in);
+ s.writer = s.stream.writer(io, &s.out);
+ s.conn = MuxT.Conn.init(app.m, s.sink());
+ defer s.conn.deinit();
+ while (true) {
+ const frame = c9.transport.readFrame(&s.reader.interface, &s.frame, max_frame) catch return;
+ app.m.forward(&s.conn, frame);
+ }
+}
fn respond(request: *std.http.Server.Request, content: []const u8, mime: []const u8, status: std.http.Status) !void {
try request.respond(content, .{ .status = status, .keep_alive = false, .extra_headers = &.{
.{ .name = "content-type", .value = mime },
.{ .name = "cache-control", .value = "no-store" },
.{ .name = "x-content-type-options", .value = "nosniff" },
- .{ .name = "content-security-policy", .value = "default-src 'self'; script-src 'self' 'wasm-unsafe-eval'; style-src 'self'; connect-src 'self'; frame-ancestors 'none'; base-uri 'none'" },
+ .{ .name = "content-security-policy", .value = "default-src 'self'; script-src 'self' 'wasm-unsafe-eval'; style-src 'self' 'unsafe-inline'; connect-src 'self'; frame-ancestors 'none'; base-uri 'none'" },
} });
}
-fn serve(io: Io, stream: Io.net.Stream, config: Config) !void {
- var input: [max_frame + 14]u8 = undefined;
- var output: [8192]u8 = undefined;
- var reader = stream.reader(io, &input);
- var writer = stream.writer(io, &output);
- var http: std.http.Server = .init(&reader.interface, &writer.interface);
+fn handleHttp(s: *Session) void {
+ const app = s.app;
+ const io = app.io;
+ var timed: TimedSession = .{ .s = s };
+ var group: Io.Group = .init;
+ defer group.cancel(io);
+ group.concurrent(io, TimedSession.run, .{&timed}) catch {
+ s.stream.close(io);
+ app.giveBack(s);
+ return;
+ };
+ const timeout: Io.Timeout = if (app.config.timeout_ms == 0) .none else .{ .duration = .{ .raw = .fromMilliseconds(app.config.timeout_ms), .clock = .awake } };
+ timed.done.waitTimeout(io, timeout) catch {};
+ s.stream.close(io);
+ group.cancel(io);
+ app.giveBack(s);
+}
+
+const TimedSession = struct {
+ s: *Session,
+ done: Io.Event = .unset,
+ fn run(t: *TimedSession) void {
+ defer t.done.set(t.s.app.io);
+ serveHttp(t.s) catch {};
+ }
+};
+
+fn serveHttp(s: *Session) !void {
+ const app = s.app;
+ const io = app.io;
+ const config = app.config;
+ s.reader = s.stream.reader(io, &s.in);
+ s.writer = s.stream.writer(io, &s.out);
+ var http: std.http.Server = .init(&s.reader.interface, &s.writer.interface);
http.reader.max_head_len = 8192;
var request = try http.receiveHead();
+
var host: ?[]const u8 = null;
var origin: ?[]const u8 = null;
var version: ?[]const u8 = null;
@@ -147,58 +318,60 @@ fn serve(io: Io, stream: Io.net.Stream, config: Config) !void {
if (std.ascii.eqlIgnoreCase(header.name, "origin")) origin = header.value;
if (std.ascii.eqlIgnoreCase(header.name, "sec-websocket-version")) version = header.value;
}
- if (!std.mem.eql(u8, host orelse "", config.host) and !std.mem.eql(u8, host orelse "", config.public_host)) return respond(&request, "Unexpected Host\n", "text/plain", .forbidden);
- if (request.head.method != .GET) return respond(&request, "Use GET\n", "text/plain", .method_not_allowed);
+ if (!std.mem.eql(u8, host orelse "", config.host) and !std.mem.eql(u8, host orelse "", config.public_host))
+ return respond(&request, "Unexpected Host\n", "text/plain", .forbidden);
+
const path = std.mem.sliceTo(request.head.target, '?');
+
+ // The WebSocket endpoint.
+ if (std.mem.eql(u8, path, "/_cloud9/9p")) {
+ if (request.head.method != .GET) return respond(&request, "Use GET\n", "text/plain", .method_not_allowed);
+ if (!std.mem.eql(u8, origin orelse "", config.origin)) return respond(&request, "Unexpected Origin\n", "text/plain", .forbidden);
+ if (!std.mem.eql(u8, version orelse "", "13")) return respond(&request, "WebSocket version 13 required\n", "text/plain", .bad_request);
+ return wsSession(s, &request);
+ }
+
+ // The HTTP view of the tree.
+ if (std.mem.eql(u8, path, "/fs") or std.mem.startsWith(u8, path, "/fs/"))
+ return httpfs.handle(app.m, &request, path, &s.frame, &s.body);
+
+ // Static assets and configuration.
+ if (request.head.method != .GET and request.head.method != .HEAD)
+ return respond(&request, "Use GET\n", "text/plain", .method_not_allowed);
if (std.mem.eql(u8, path, "/_cloud9/config.json")) return respond(&request, config.browser_config, "application/json", .ok);
if (std.mem.eql(u8, path, "/_cloud9/app.mjs")) return respond(&request, @embedFile("static/app.mjs"), "text/javascript; charset=utf-8", .ok);
if (std.mem.eql(u8, path, "/_cloud9/client.mjs")) return respond(&request, @embedFile("static/client.mjs"), "text/javascript; charset=utf-8", .ok);
if (std.mem.eql(u8, path, "/_cloud9/style.css")) return respond(&request, @embedFile("static/style.css"), "text/css; charset=utf-8", .ok);
if (std.mem.eql(u8, path, "/_cloud9/cloud9.wasm")) return respond(&request, @embedFile("client.wasm"), "application/wasm", .ok);
- if (!std.mem.eql(u8, path, "/_cloud9/9p")) {
- if (std.mem.eql(u8, path, "/_cloud9") or std.mem.startsWith(u8, path, "/_cloud9/"))
- return respond(&request, "Not found\n", "text/plain", .not_found);
- // File URLs boot the same browser client. The client resolves the path
- // in the upstream 9P namespace, never in this machine's filesystem.
- return respond(&request, @embedFile("static/index.html"), "text/html; charset=utf-8", .ok);
- }
- if (!std.mem.eql(u8, origin orelse "", config.origin)) return respond(&request, "Unexpected Origin\n", "text/plain", .forbidden);
- if (!std.mem.eql(u8, version orelse "", "13")) return respond(&request, "WebSocket version 13 required\n", "text/plain", .bad_request);
- switch (config.upstream) {
- .network => |address| {
- const upstream = c9.transport.connect(io, address) catch return respond(&request, "9P upstream unavailable\n", "text/plain", .bad_gateway);
- defer upstream.close(io);
- var in_buffer: [8192]u8 = undefined;
- var out_buffer: [8192]u8 = undefined;
- var upstream_reader = upstream.reader(io, &in_buffer);
- var upstream_writer = upstream.writer(io, &out_buffer);
- try bridge(io, &request, &upstream_reader.interface, &upstream_writer.interface);
- },
- .file => |path_name| {
- if (serial_busy.swap(true, .acquire)) return respond(&request, "Device already in use\n", "text/plain", .service_unavailable);
- defer serial_busy.store(false, .release);
- const file = Io.Dir.cwd().openFile(io, path_name, .{ .mode = .read_write }) catch return respond(&request, "9P device unavailable\n", "text/plain", .bad_gateway);
- defer file.close(io);
- var in_buffer: [8192]u8 = undefined;
- var out_buffer: [8192]u8 = undefined;
- var upstream_reader = file.readerStreaming(io, &in_buffer);
- var upstream_writer = file.writerStreaming(io, &out_buffer);
- try bridge(io, &request, &upstream_reader.interface, &upstream_writer.interface);
- },
- }
+ if (std.mem.eql(u8, path, "/_cloud9") or std.mem.startsWith(u8, path, "/_cloud9/"))
+ return respond(&request, "Not found\n", "text/plain", .not_found);
+ // Any other path boots the browser client (it resolves the path in 9P).
+ return respond(&request, @embedFile("static/index.html"), "text/html; charset=utf-8", .ok);
}
-fn bridge(io: Io, request: *std.http.Server.Request, reader: *Io.Reader, writer: *Io.Writer) !void {
- var ws = try c9.http.accept(request);
- var requests: [max_frame]u8 = undefined;
- var replies: [max_frame]u8 = undefined;
- var relay: c9.http.Bridge = .{
- .socket = &ws,
- .upstream_reader = reader,
- .upstream_writer = writer,
- .request_buffer = &requests,
- .reply_buffer = &replies,
- .frame_limit = max_frame,
- };
- try relay.run(io, .none);
+fn wsSession(s: *Session, request: *std.http.Server.Request) !void {
+ const app = s.app;
+ s.is_ws = true;
+ s.ws = try c9.http.accept(request);
+ s.conn = MuxT.Conn.init(app.m, s.sink());
+ defer s.conn.deinit();
+ while (true) {
+ const message = s.ws.receive(&s.frame) catch return;
+ switch (message.opcode) {
+ .binary => app.m.forward(&s.conn, message.data),
+ .ping => {
+ s.wmutex.lockUncancelable(app.io);
+ defer s.wmutex.unlock(app.io);
+ s.ws.send(message.data, .pong, null) catch return;
+ },
+ .pong => {},
+ .close => {
+ s.wmutex.lockUncancelable(app.io);
+ defer s.wmutex.unlock(app.io);
+ s.ws.send(message.data, .close, null) catch {};
+ return;
+ },
+ else => return,
+ }
+ }
}
diff --git a/web/mux.zig b/web/mux.zig
new file mode 100644
index 0000000..23f2e6c
--- /dev/null
+++ b/web/mux.zig
@@ -0,0 +1,812 @@
+//! A transparent 9P2000 multiplexer: one shared upstream `cloud9.Client`
+//! connection (all 16 tags) fanned out to any number of downstream sessions.
+//!
+//! Each downstream request is forwarded upstream on a remapped tag with its
+//! fids remapped into the shared upstream fid space; each upstream reply is
+//! routed back to the originating downstream by tag. Tversion is answered
+//! locally (the upstream session is negotiated once at connect); Tattach is
+//! forwarded so every downstream gets its own upstream tree root; Tflush is
+//! forwarded as Tflush. Downstream fid spaces are isolated by construction:
+//! two downstreams never share an upstream fid.
+//!
+//! The upstream is driven by one reader task and any number of downstream
+//! forwarder tasks, all serialized by `lock`. Nothing here allocates per
+//! request: the tag/fid/pending tables are comptime-sized. On an upstream I/O
+//! failure the reader reconnects, bumps `generation`, and every downstream's
+//! fids from an older generation are answered EIO until it re-attaches.
+//!
+//! This is deliberately *not* an `fs.Server`/`serve.Runner` backend: the
+//! engine re-decomposes each request into filesystem operations and re-encodes
+//! directories with synthetic stats and an entry-index cursor, which loses the
+//! upstream's real directory stats, iounit and qids and complicates the
+//! readdir byte offset. A frame-level remux forwards requests unchanged, which
+//! is exactly the "forward each request upstream" contract and matches the way
+//! plan9port's 9pserve multiplexes. The HTTP `/fs` view (see httpfs.zig) uses
+//! the same `Mux` at the fid level directly.
+const std = @import("std");
+const c9 = @import("cloud9");
+const Io = std.Io;
+const wire = c9.wire;
+const transport = c9.transport;
+
+const notag = wire.notag;
+const nofid = wire.nofid;
+
+pub const Limits = struct {
+ /// Largest 9P frame on either side; sizes the shared upstream buffers.
+ msize: u32 = 65536,
+ /// Upstream fids the shared connection may hold across all downstreams.
+ upstream_fids: usize = 4096,
+ /// Fids one downstream session may hold at once.
+ fids_per_conn: usize = 512,
+ /// Downstream sessions routed at once (for introspection only).
+ downstreams: usize = 256,
+};
+
+/// How the mux hands an upstream reply back to a downstream. `send` is called
+/// off the mux lock and must serialize writes on that downstream itself.
+pub const Sink = struct {
+ ctx: *anyopaque,
+ send: *const fn (ctx: *anyopaque, frame: []const u8) anyerror!void,
+};
+
+pub const Address = union(enum) {
+ network: transport.Address,
+ /// An already configured duplex device (serial), one session at a time.
+ file: []const u8,
+};
+
+pub const Error = error{
+ Disconnected,
+ TooManyFids,
+ UnknownFid,
+ BadRequest,
+ Upstream,
+};
+
+pub fn Mux(comptime limits: Limits) type {
+ if (limits.msize < 256) @compileError("mux msize too small");
+ return struct {
+ const Self = @This();
+ pub const msize = limits.msize;
+
+ io: Io,
+ address: Address,
+ user: []const u8,
+ tree: []const u8,
+
+ client: c9.Client = undefined,
+ in: [msize]u8 = undefined,
+ out: [msize]u8 = undefined,
+ rbuf: [msize]u8 = undefined,
+ stage: [msize]u8 = undefined,
+ /// One send buffer per upstream tag (16 ordinary + 1 flush).
+ frames: [17][msize]u8 = undefined,
+ /// A scratch buffer to serialize one submitted frame out of the client.
+ wbuf: [msize]u8 = undefined,
+
+ stream: ?Io.net.Stream = null,
+ file: ?Io.File = null,
+ reader: Io.net.Stream.Reader = undefined,
+ writer: Io.net.Stream.Writer = undefined,
+ freader: Io.File.Reader = undefined,
+ fwriter: Io.File.Writer = undefined,
+ ureader: *Io.Reader = undefined,
+ uwriter: *Io.Writer = undefined,
+
+ mutex: Io.Mutex = .init,
+ wmutex: Io.Mutex = .init,
+ cond: Io.Condition = .init,
+ negotiated: u32 = 0,
+ version: [16]u8 = undefined,
+ version_len: usize = 0,
+ generation: u32 = 1,
+ connected: bool = false,
+ stopping: std.atomic.Value(bool) = .init(0 != 0),
+
+ /// Counts for the /probe introspection.
+ downstreams: std.atomic.Value(u32) = .init(0),
+ upstream_up: std.atomic.Value(u32) = .init(0),
+ reconnects: std.atomic.Value(u32) = .init(0),
+
+ fid_used: [limits.upstream_fids]bool = @splat(false),
+ fid_next: usize = 0,
+
+ pending: [17]Pending = @splat(.{}),
+
+ const Kind = enum { plain, walk, clunk, flush };
+ /// A blocking fid-level RPC used by the HTTP `/fs` view. The reply
+ /// frame is copied into `buf` before the reader advances, so its
+ /// borrowed data stays valid until the waiter consumes it.
+ pub const Rpc = struct {
+ buf: []u8,
+ len: usize = 0,
+ ready: bool = false,
+ failed: bool = false,
+ /// The upstream tag this call holds while in flight (for `flushRpc`).
+ tag: u16 = 0,
+ event: Io.Event = .unset,
+ };
+ const Pending = struct {
+ active: bool = false,
+ kind: Kind = .plain,
+ sink: Sink = undefined,
+ conn: ?*Conn = null,
+ rpc: ?*Rpc = null,
+ down_tag: u16 = 0,
+ /// walk: the downstream/upstream newfid and whether it was in place.
+ newfid_down: u32 = 0,
+ newfid_up: u32 = 0,
+ walk_names: u16 = 0,
+ inplace: bool = false,
+ /// clunk/remove: the downstream fid to drop on completion.
+ clunk_down: u32 = 0,
+ /// flush: the upstream tag it cancels.
+ flush_up: u16 = 0,
+ };
+
+ /// A downstream session: its own fid map and the generation it belongs
+ /// to. `sink` routes replies; `write` on the transport must serialize.
+ pub const Conn = struct {
+ mux: *Self,
+ sink: Sink,
+ generation: u32 = 0,
+ alive: bool = true,
+ fids: [limits.fids_per_conn]FidMap = @splat(.{}),
+
+ const FidMap = struct { used: bool = false, down: u32 = 0, up: u32 = 0 };
+
+ pub fn init(m: *Self, sink: Sink) Conn {
+ _ = m.downstreams.fetchAdd(1, .monotonic);
+ return .{ .mux = m, .sink = sink, .generation = m.generation };
+ }
+
+ /// Drops every upstream fid this session holds and cancels its
+ /// outstanding requests, then unregisters it.
+ pub fn deinit(c: *Conn) void {
+ const m = c.mux;
+ m.lock();
+ c.alive = false;
+ // Orphan any in-flight replies bound for this session.
+ for (&m.pending) |*p| if (p.active and p.conn == c) {
+ p.conn = null;
+ };
+ // Best-effort clunk of every live upstream fid.
+ if (c.generation == m.generation and m.connected) {
+ for (&c.fids) |*e| if (e.used) {
+ m.clunkUpstreamLocked(e.up);
+ e.used = false;
+ };
+ m.flushOutput();
+ }
+ m.unlock();
+ _ = m.downstreams.fetchSub(1, .monotonic);
+ }
+
+ fn mapFind(c: *Conn, down: u32) ?*FidMap {
+ for (&c.fids) |*e| if (e.used and e.down == down) return e;
+ return null;
+ }
+ fn mapAdd(c: *Conn, down: u32, up: u32) void {
+ for (&c.fids) |*e| if (!e.used) {
+ e.* = .{ .used = true, .down = down, .up = up };
+ return;
+ };
+ unreachable; // caller checked capacity via free upstream fid
+ }
+ };
+
+ pub fn init(m: *Self, io: Io, address: Address, user: []const u8, tree: []const u8) void {
+ m.* = .{ .io = io, .address = address, .user = user, .tree = tree };
+ }
+
+ fn lock(m: *Self) void {
+ m.mutex.lockUncancelable(m.io);
+ }
+ fn unlock(m: *Self) void {
+ m.mutex.unlock(m.io);
+ }
+
+ pub fn upstreamState(m: *Self) []const u8 {
+ return if (m.upstream_up.load(.acquire) != 0) "connected" else "disconnected";
+ }
+ pub fn downstreamCount(m: *Self) u32 {
+ return m.downstreams.load(.acquire);
+ }
+ pub fn reconnectCount(m: *Self) u32 {
+ return m.reconnects.load(.acquire);
+ }
+ pub fn negotiatedMsize(m: *Self) u32 {
+ return m.negotiated;
+ }
+
+ // -- upstream connection ------------------------------------------
+
+ fn dial(m: *Self) !void {
+ switch (m.address) {
+ .network => |addr| {
+ const s = try transport.connect(m.io, addr);
+ m.stream = s;
+ m.reader = s.reader(m.io, &m.rbuf);
+ m.writer = s.writer(m.io, &m.wbuf);
+ m.ureader = &m.reader.interface;
+ m.uwriter = &m.writer.interface;
+ },
+ .file => |path| {
+ const f = try Io.Dir.cwd().openFile(m.io, path, .{ .mode = .read_write });
+ m.file = f;
+ m.freader = f.readerStreaming(m.io, &m.rbuf);
+ m.fwriter = f.writerStreaming(m.io, &m.wbuf);
+ m.ureader = &m.freader.interface;
+ m.uwriter = &m.fwriter.interface;
+ },
+ }
+ }
+
+ fn closeUpstream(m: *Self) void {
+ if (m.stream) |s| {
+ s.close(m.io);
+ m.stream = null;
+ }
+ if (m.file) |f| {
+ f.close(m.io);
+ m.file = null;
+ }
+ }
+
+ /// Connects and negotiates the shared session. Call once before the
+ /// reader task runs; returns an error if the upstream is unreachable.
+ pub fn connect(m: *Self) !void {
+ try m.dial();
+ errdefer m.closeUpstream();
+ m.client = .init(.{ .in = &m.in, .out = &m.out });
+ try m.handshake();
+ m.connected = true;
+ m.upstream_up.store(1, .release);
+ }
+
+ fn handshake(m: *Self) !void {
+ _ = m.client.submit(.{ .version = .{ .msize = msize } }) catch return error.Upstream;
+ try m.sendClientOutput();
+ const done = try m.readOne();
+ if (done.op != .version) return error.Upstream;
+ m.negotiated = done.result.version.msize;
+ const v = done.result.version.version;
+ m.version_len = @min(v.len, m.version.len);
+ @memcpy(m.version[0..m.version_len], v[0..m.version_len]);
+ if (!std.mem.eql(u8, m.version[0..m.version_len], "9P2000")) return error.Upstream;
+ }
+
+ /// Writes whatever the client has staged (used only during handshake).
+ fn sendClientOutput(m: *Self) !void {
+ const bytes = m.client.output();
+ m.uwriter.writeAll(bytes) catch return error.Upstream;
+ m.uwriter.flush() catch return error.Upstream;
+ m.client.wrote(bytes.len);
+ }
+
+ fn readOne(m: *Self) !c9.Client.Done {
+ while (true) {
+ if (m.client.take()) |done| return done;
+ if (m.client.dead) return error.Upstream;
+ const frame = transport.readFrame(m.ureader, &m.stage, msize) catch return error.Upstream;
+ if (m.client.push(frame) != frame.len) return error.Upstream;
+ }
+ }
+
+ // -- upstream fid pool --------------------------------------------
+
+ fn allocFid(m: *Self) ?u32 {
+ var i: usize = 0;
+ while (i < limits.upstream_fids) : (i += 1) {
+ const idx = (m.fid_next + i) % limits.upstream_fids;
+ if (!m.fid_used[idx]) {
+ m.fid_used[idx] = true;
+ m.fid_next = (idx + 1) % limits.upstream_fids;
+ return @intCast(idx);
+ }
+ }
+ return null;
+ }
+ fn freeFid(m: *Self, fid: u32) void {
+ if (fid < limits.upstream_fids) m.fid_used[fid] = false;
+ }
+
+ // -- forwarding ----------------------------------------------------
+
+ /// Handles one downstream frame. Tversion is answered locally through
+ /// the sink; every other message is remapped and forwarded upstream,
+ /// its reply delivered later by the reader task. On a synchronous
+ /// failure it answers the downstream with an Rerror.
+ pub fn forward(m: *Self, c: *Conn, frame: []const u8) void {
+ const got = wire.decode(frame) catch {
+ c.sink.send(c.sink.ctx, m.errorFrame(&m.wbuf, 0, "protocol botch")) catch {};
+ return;
+ };
+ if (!wire.isT(got.msg.msgType())) {
+ m.rerror(c, got.tag, "protocol botch");
+ return;
+ }
+ switch (got.msg) {
+ .tversion => |v| {
+ // The shared upstream session is already negotiated. Answer
+ // locally and reset this downstream's fid space.
+ m.lock();
+ for (&c.fids) |*e| if (e.used and c.generation == m.generation and m.connected) {
+ m.clunkUpstreamLocked(e.up);
+ e.used = false;
+ } else {
+ e.used = false;
+ };
+ m.flushOutput();
+ const use: u32 = @min(@min(v.msize, m.negotiated), msize);
+ c.generation = m.generation;
+ m.unlock();
+ const reply = wire.encode(.{ .rversion = .{ .msize = use, .version = "9P2000" } }, got.tag, &m.wbuf) catch return;
+ c.sink.send(c.sink.ctx, reply) catch {};
+ },
+ else => m.forwardRequest(c, got),
+ }
+ }
+
+ fn rerror(m: *Self, c: *Conn, tag: u16, text: []const u8) void {
+ var buf: [wire.header_len + 2 + 128]u8 = undefined;
+ const frame = m.errorFrame(&buf, tag, text);
+ c.sink.send(c.sink.ctx, frame) catch {};
+ }
+
+ fn errorFrame(m: *Self, buf: []u8, tag: u16, text: []const u8) []const u8 {
+ _ = m;
+ const t = text[0..@min(text.len, 128)];
+ return wire.encode(.{ .rerror = .{ .ename = t } }, tag, buf) catch buf[0..0];
+ }
+
+ fn forwardRequest(m: *Self, c: *Conn, got: wire.Decoded) void {
+ m.lock();
+ defer m.unlock();
+
+ if (!m.connected) {
+ m.unlock();
+ m.rerror(c, got.tag, "Transport endpoint is not connected");
+ m.lock();
+ return;
+ }
+ if (c.generation != m.generation) {
+ // Stale session after a reconnect: fids are gone. Reset and,
+ // unless this is a fresh attach, answer EIO.
+ for (&c.fids) |*e| e.used = false;
+ c.generation = m.generation;
+ if (got.msg != .tattach) {
+ m.unlock();
+ m.rerror(c, got.tag, "Input/output error");
+ m.lock();
+ return;
+ }
+ }
+
+ // Translate to a client request with upstream fids.
+ var pend: Pending = .{ .active = true, .sink = c.sink, .conn = c, .down_tag = got.tag };
+ const req: c9.Client.Request = switch (got.msg) {
+ .tattach => |a| blk: {
+ if (c.mapFind(a.fid) != null) return m.syncErr(c, got.tag, "fid already in use");
+ const up = m.allocFid() orelse return m.syncErr(c, got.tag, "Too many open files in system");
+ pend.kind = .walk; // reuse walk bookkeeping to bind on success
+ pend.newfid_down = a.fid;
+ pend.newfid_up = up;
+ pend.walk_names = 0;
+ pend.inplace = false;
+ break :blk .{ .attach = .{ .fid = up, .uname = a.uname, .aname = a.aname } };
+ },
+ .twalk => |w| blk: {
+ const src = c.mapFind(w.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ var up_new: u32 = src.up;
+ const inplace = w.newfid == w.fid;
+ if (!inplace) {
+ if (c.mapFind(w.newfid) != null) return m.syncErr(c, got.tag, "fid already in use");
+ up_new = m.allocFid() orelse return m.syncErr(c, got.tag, "Too many open files in system");
+ }
+ pend.kind = .walk;
+ pend.newfid_down = w.newfid;
+ pend.newfid_up = up_new;
+ pend.walk_names = w.nwname;
+ pend.inplace = inplace;
+ break :blk .{ .walk = .{ .fid = src.up, .newfid = up_new, .names = w.wname[0..w.nwname] } };
+ },
+ .topen => |o| blk: {
+ const e = c.mapFind(o.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .open = .{ .fid = e.up, .mode = o.mode } };
+ },
+ .tcreate => |cr| blk: {
+ const e = c.mapFind(cr.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .create = .{ .fid = e.up, .name = cr.name, .perm = cr.perm, .mode = cr.mode } };
+ },
+ .tread => |r| blk: {
+ const e = c.mapFind(r.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .read = .{ .fid = e.up, .offset = r.offset, .count = r.count } };
+ },
+ .twrite => |w| blk: {
+ const e = c.mapFind(w.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .write = .{ .fid = e.up, .offset = w.offset, .data = w.data } };
+ },
+ .tclunk => |cl| blk: {
+ const e = c.mapFind(cl.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ pend.kind = .clunk;
+ pend.clunk_down = cl.fid;
+ break :blk .{ .clunk = .{ .fid = e.up } };
+ },
+ .tremove => |rm| blk: {
+ const e = c.mapFind(rm.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ pend.kind = .clunk;
+ pend.clunk_down = rm.fid;
+ break :blk .{ .remove = .{ .fid = e.up } };
+ },
+ .tstat => |s| blk: {
+ const e = c.mapFind(s.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .stat = .{ .fid = e.up } };
+ },
+ .twstat => |s| blk: {
+ const e = c.mapFind(s.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .wstat = .{ .fid = e.up, .stat = s.stat } };
+ },
+ .tflush => |f| blk: {
+ const up = m.findUpTag(c, f.oldtag);
+ if (up == null) {
+ m.unlock();
+ var buf: [wire.header_len]u8 = undefined;
+ const rf = wire.encode(.rflush, got.tag, &buf) catch buf[0..0];
+ c.sink.send(c.sink.ctx, rf) catch {};
+ m.lock();
+ return;
+ }
+ pend.kind = .flush;
+ pend.flush_up = up.?;
+ break :blk .{ .flush = .{ .oldtag = up.? } };
+ },
+ .tauth => return m.syncErr(c, got.tag, "authentication not required"),
+ else => return m.syncErr(c, got.tag, "protocol botch"),
+ };
+
+ const up_tag = m.submitLocked(req) catch |err| {
+ if (pend.kind == .walk and !pend.inplace) m.freeFid(pend.newfid_up);
+ m.unlock();
+ m.rerror(c, got.tag, switch (err) {
+ error.Disconnected => "Transport endpoint is not connected",
+ else => "Input/output error",
+ });
+ m.lock();
+ return;
+ };
+ m.pending[up_tag] = pend;
+ }
+
+ /// A synchronous error while holding the lock: drops the lock to send,
+ /// then reacquires so the deferred unlock stays balanced.
+ fn syncErr(m: *Self, c: *Conn, tag: u16, text: []const u8) void {
+ m.unlock();
+ m.rerror(c, tag, text);
+ m.lock();
+ }
+
+ fn findUpTag(m: *Self, c: *Conn, down_tag: u16) ?u16 {
+ for (&m.pending, 0..) |*p, i| {
+ if (p.active and p.conn == c and p.down_tag == down_tag and p.kind != .flush) return @intCast(i);
+ }
+ return null;
+ }
+
+ /// Submits a request, waiting for a free tag; sends it upstream. The
+ /// caller holds `lock`. Returns the upstream tag.
+ fn submitLocked(m: *Self, req: c9.Client.Request) !u16 {
+ while (true) {
+ if (!m.connected) return error.Disconnected;
+ const tag = m.client.submit(req) catch |err| switch (err) {
+ error.NoTags => {
+ m.cond.wait(m.io, &m.mutex) catch return error.Disconnected;
+ continue;
+ },
+ else => return error.BadRequest,
+ };
+ // The client's output buffer is disjoint from the stream
+ // writer's buffer, so write it out directly under the lock.
+ const bytes = m.client.output();
+ m.writeUpstream(bytes) catch return error.Disconnected;
+ m.client.wrote(bytes.len);
+ return tag;
+ }
+ }
+
+ fn writeUpstream(m: *Self, bytes: []const u8) !void {
+ m.wmutex.lockUncancelable(m.io);
+ defer m.wmutex.unlock(m.io);
+ m.uwriter.writeAll(bytes) catch return error.Upstream;
+ m.uwriter.flush() catch return error.Upstream;
+ }
+
+ fn clunkUpstreamLocked(m: *Self, up: u32) void {
+ // Fire-and-forget clunk to reclaim an upstream fid on disconnect.
+ const tag = m.client.submit(.{ .clunk = .{ .fid = up } }) catch {
+ m.freeFid(up);
+ return;
+ };
+ m.pending[tag] = .{ .active = true, .kind = .clunk, .conn = null, .clunk_down = 0 };
+ const bytes = m.client.output();
+ m.writeUpstream(bytes) catch {};
+ m.client.wrote(bytes.len);
+ m.freeFid(up);
+ }
+
+ fn flushOutput(m: *Self) void {
+ _ = m;
+ }
+
+ // -- synchronous fid-level RPC (for the HTTP /fs view) -------------
+
+ /// Allocates an upstream fid from the shared pool. Returns null when
+ /// exhausted or the upstream is down.
+ pub fn takeFid(m: *Self) ?u32 {
+ m.lock();
+ defer m.unlock();
+ if (!m.connected) return null;
+ return m.allocFid();
+ }
+
+ pub fn dropFid(m: *Self, fid: u32) void {
+ m.lock();
+ defer m.unlock();
+ m.freeFid(fid);
+ }
+
+ /// Blocks until the shared upstream answers `request`, copying the
+ /// reply into `r.buf`. Returns the decoded reply. Fids in `request`
+ /// are upstream fids (from `takeFid`); the caller owns their lifetime.
+ pub fn rpc(m: *Self, request: c9.Client.Request, r: *Rpc) !wire.Decoded {
+ r.* = .{ .buf = r.buf };
+ m.lock();
+ if (!m.connected) {
+ m.unlock();
+ return error.Disconnected;
+ }
+ const tag = m.submitLocked(request) catch |err| {
+ m.unlock();
+ return err;
+ };
+ r.tag = tag;
+ m.pending[tag] = .{ .active = true, .kind = .plain, .conn = null, .rpc = r, .down_tag = 0 };
+ m.unlock();
+ r.event.wait(m.io) catch {
+ // The waiting fiber was cancelled (HTTP timeout, follow
+ // teardown). Its `r` is about to be freed, so make sure the
+ // mux stops referencing it before returning: Tflush the
+ // upstream op and wait, uncancelably, until its pending clears.
+ m.cancelAndDrain(r);
+ return error.Canceled;
+ };
+ if (!r.ready or r.len == 0) return error.Upstream;
+ return wire.decode(r.buf[0..r.len]) catch return error.Upstream;
+ }
+
+ /// Guarantees `r` is no longer referenced by any pending slot before
+ /// the caller frees it. Best-effort Tflush of `r`'s upstream tag (the
+ /// flush's own reply clears `r`'s pending in `route`); then an
+ /// uncancelable wait for that to happen. A late reply or a disconnect
+ /// also clears it, so this returns even if the flush cannot be sent.
+ fn cancelAndDrain(m: *Self, r: *Rpc) void {
+ m.lock();
+ var live = false;
+ for (&m.pending) |*p| if (p.active and p.rpc == r) {
+ live = true;
+ };
+ if (!live) {
+ m.unlock();
+ return;
+ }
+ r.event.reset();
+ if (m.connected) {
+ if (m.client.submit(.{ .flush = .{ .oldtag = r.tag } })) |ftag| {
+ m.pending[ftag] = .{ .active = true, .kind = .flush, .conn = null, .flush_up = r.tag };
+ const bytes = m.client.output();
+ m.writeUpstream(bytes) catch {};
+ m.client.wrote(bytes.len);
+ } else |_| {}
+ }
+ m.unlock();
+ while (true) {
+ m.lock();
+ var still = false;
+ for (&m.pending) |*p| if (p.active and p.rpc == r) {
+ still = true;
+ };
+ m.unlock();
+ if (!still) return;
+ r.event.waitUncancelable(m.io);
+ r.event.reset();
+ }
+ }
+
+ // -- reader task ---------------------------------------------------
+
+ /// Reads upstream replies forever, routing each to its downstream.
+ /// Reconnects on failure. Runs on its own task; ended by `stop`.
+ pub fn readerLoop(m: *Self) void {
+ while (!m.stopping.load(.acquire)) {
+ if (!m.connected) {
+ m.reconnect();
+ if (m.stopping.load(.acquire)) return;
+ continue;
+ }
+ const frame = transport.readFrame(m.ureader, &m.stage, msize) catch {
+ m.onDisconnect();
+ if (m.stopping.load(.acquire)) return;
+ m.reconnect();
+ continue;
+ };
+ m.deliver(frame);
+ }
+ }
+
+ const Outgoing = struct { sink: Sink, tag: u8, len: usize };
+
+ fn deliver(m: *Self, frame: []const u8) void {
+ var outs: [17]Outgoing = undefined;
+ var nouts: usize = 0;
+ m.lock();
+ if (m.client.push(frame) != frame.len) {
+ m.unlock();
+ m.onDisconnect();
+ m.reconnect();
+ return;
+ }
+ while (m.client.take()) |done| {
+ const p = &m.pending[done.tag];
+ if (!p.active) continue;
+ if (p.rpc) |r| {
+ const built = m.route(done, p, r.buf);
+ p.active = false;
+ r.len = built orelse 0;
+ r.failed = done.result == .fail;
+ r.ready = true;
+ r.event.set(m.io);
+ continue;
+ }
+ const built = m.route(done, p, &m.frames[done.tag]);
+ const conn = p.conn;
+ p.active = false;
+ if (built) |len| if (conn != null) {
+ outs[nouts] = .{ .sink = p.sink, .tag = @intCast(done.tag), .len = len };
+ nouts += 1;
+ };
+ }
+ m.cond.broadcast(m.io);
+ m.unlock();
+ // Send outside the lock; a slow downstream cannot stall the client.
+ for (outs[0..nouts]) |o| o.sink.send(o.sink.ctx, m.frames[o.tag][0..o.len]) catch {};
+ }
+
+ /// Builds the downstream reply frame into `buf`, updating fid state.
+ /// Returns the encoded length, or null if nothing should be sent.
+ fn route(m: *Self, done: c9.Client.Done, p: *Pending, buf: []u8) ?usize {
+ // Fid bookkeeping first (independent of whether we send).
+ switch (p.kind) {
+ .walk => {
+ const full = done.result == .attach or
+ (done.result == .walk and done.result.walk.nwqid == p.walk_names);
+ if (p.conn) |c| {
+ if (full) {
+ if (!p.inplace) c.mapAdd(p.newfid_down, p.newfid_up);
+ } else if (!p.inplace) {
+ m.freeFid(p.newfid_up);
+ }
+ } else if (!p.inplace and !full) {
+ m.freeFid(p.newfid_up);
+ } else if (p.conn == null and full and !p.inplace) {
+ // Session gone: reclaim the fid the server just bound.
+ m.clunkUpstreamLocked(p.newfid_up);
+ }
+ },
+ .clunk => {
+ if (p.conn) |c| {
+ if (c.mapFind(p.clunk_down)) |e| {
+ m.freeFid(e.up);
+ e.used = false;
+ }
+ }
+ // fire-and-forget clunk (conn==null): the fid was already freed.
+ },
+ .flush => {
+ // The flushed request gets no reply after Rflush: clear its
+ // pending and wake any blocking rpc waiting on it.
+ const fp = &m.pending[p.flush_up];
+ if (fp.active) {
+ fp.active = false;
+ if (fp.rpc) |rr| {
+ rr.failed = true;
+ rr.ready = false;
+ rr.event.set(m.io);
+ }
+ }
+ },
+ .plain => {},
+ }
+ if (p.conn == null and p.rpc == null) return null;
+
+ const reply: wire.Msg = switch (done.result) {
+ .fail => |ename| .{ .rerror = .{ .ename = ename } },
+ .version => return null,
+ .auth => |q| .{ .rauth = .{ .aqid = q } },
+ .attach => |q| .{ .rattach = .{ .qid = q } },
+ .walk => |w| .{ .rwalk = .{ .nwqid = w.nwqid, .wqid = w.wqid } },
+ .open => |o| .{ .ropen = .{ .qid = o.qid, .iounit = o.iounit } },
+ .create => |cr| .{ .rcreate = .{ .qid = cr.qid, .iounit = cr.iounit } },
+ .read => |data| .{ .rread = .{ .data = data } },
+ .write => |n| .{ .rwrite = .{ .count = n } },
+ .clunk => .rclunk,
+ .remove => .rremove,
+ .stat => |s| .{ .rstat = .{ .stat = s } },
+ .wstat => .rwstat,
+ .flush => .rflush,
+ };
+ const encoded = wire.encode(reply, p.down_tag, buf) catch {
+ return (wire.encode(.{ .rerror = .{ .ename = "Invalid argument" } }, p.down_tag, buf) catch return null).len;
+ };
+ return encoded.len;
+ }
+
+ fn onDisconnect(m: *Self) void {
+ var outs: [17]Outgoing = undefined;
+ var nouts: usize = 0;
+ m.lock();
+ m.connected = false;
+ m.upstream_up.store(0, .release);
+ m.client.hangup();
+ // Fail every outstanding request so no downstream hangs.
+ for (&m.pending, 0..) |*p, i| if (p.active) {
+ p.active = false;
+ if (p.rpc) |r| {
+ r.failed = true;
+ r.ready = false;
+ r.event.set(m.io);
+ } else if (p.conn != null) {
+ const frame = m.errorFrame(&m.frames[i], p.down_tag, "Transport endpoint is not connected");
+ outs[nouts] = .{ .sink = p.sink, .tag = @intCast(i), .len = frame.len };
+ nouts += 1;
+ }
+ };
+ for (&m.fid_used) |*u| u.* = false;
+ m.fid_next = 0;
+ m.closeUpstream();
+ m.cond.broadcast(m.io);
+ m.unlock();
+ for (outs[0..nouts]) |o| o.sink.send(o.sink.ctx, m.frames[o.tag][0..o.len]) catch {};
+ }
+
+ fn reconnect(m: *Self) void {
+ while (!m.stopping.load(.acquire)) {
+ m.io.sleep(.fromMilliseconds(200), .awake) catch return;
+ m.dial() catch continue;
+ m.client = .init(.{ .in = &m.in, .out = &m.out });
+ m.handshake() catch {
+ m.closeUpstream();
+ continue;
+ };
+ m.lock();
+ m.generation +%= 1;
+ if (m.generation == 0) m.generation = 1;
+ m.connected = true;
+ m.upstream_up.store(1, .release);
+ _ = m.reconnects.fetchAdd(1, .monotonic);
+ m.cond.broadcast(m.io);
+ m.unlock();
+ return;
+ }
+ }
+
+ pub fn stop(m: *Self) void {
+ m.stopping.store(true, .release);
+ m.lock();
+ m.closeUpstream();
+ m.connected = false;
+ m.cond.broadcast(m.io);
+ m.unlock();
+ }
+ };
+}
diff --git a/web/probe.zig b/web/probe.zig
new file mode 100644
index 0000000..a6f5ceb
--- /dev/null
+++ b/web/probe.zig
@@ -0,0 +1,83 @@
+//! Optional self-introspection: 9web embeds 9proc under `--probe`, off by
+//! default. It exposes the gateway's own live counters (downstream and HTTP
+//! connection counts, the upstream state, reconnect count, negotiated msize)
+//! as files, plus 9proc's /threads view, on a Unix or TCP 9P listener. Small
+//! and read-only; serve it on a private socket.
+const std = @import("std");
+const builtin = @import("builtin");
+const build_options = @import("build_options");
+
+const enabled = build_options.have_probe and builtin.os.tag == .linux;
+
+pub fn Probe(comptime App: type) type {
+ if (!enabled) return struct {
+ pub fn start(_: *App, _: std.mem.Allocator, _: []const u8) !void {
+ return error.Unsupported;
+ }
+ pub fn stop() void {}
+ };
+
+ const proc9 = @import("9proc");
+ const Writer = std.Io.Writer;
+
+ return struct {
+ const Fns = struct {
+ fn app(ctx: *anyopaque) *App {
+ return @ptrCast(@alignCast(ctx));
+ }
+ pub fn connections(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{app(ctx).connectionCount()});
+ }
+ pub fn downstreams(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{app(ctx).m.downstreamCount()});
+ }
+ pub fn upstream(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.writeAll(app(ctx).m.upstreamState());
+ }
+ pub fn reconnects(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{app(ctx).m.reconnectCount()});
+ }
+ pub fn negotiated(ctx: *anyopaque, w: *Writer) anyerror!void {
+ try w.print("{d}", .{app(ctx).m.negotiatedMsize()});
+ }
+ };
+
+ const cfg: proc9.Config = .{
+ .name = "9web",
+ .fns = Fns,
+ .msize = 16 * 1024,
+ .max_fids = 64,
+ .max_providers = 8,
+ };
+ const S = proc9.Server(cfg);
+ const ProbeT = proc9.linux.Probe(S);
+ const max_clients = 4;
+
+ var shared: S.Shared = undefined;
+ var storage: ProbeT.Storage(max_clients) = undefined;
+ var probe: ProbeT = undefined;
+ var running: bool = false;
+
+ pub fn start(app: *App, allocator: std.mem.Allocator, addr_text: []const u8) !void {
+ _ = allocator;
+ const listen: proc9.linux.Listen = if (std.mem.startsWith(u8, addr_text, "unix:"))
+ .{ .unix = addr_text[5..] }
+ else if (std.mem.startsWith(u8, addr_text, "tcp:"))
+ .{ .tcp = addr_text[4..] }
+ else
+ return error.BadProbeAddress;
+ shared = .init(app);
+ try probe.init(&shared, &storage, .{ .io = app.io, .listen = listen, .msize = cfg.msize });
+ try probe.start();
+ running = true;
+ std.debug.print("9web: probe on {s}\n", .{addr_text});
+ }
+
+ pub fn stop() void {
+ if (running) {
+ probe.stop();
+ running = false;
+ }
+ }
+ };
+}
diff --git a/web/static/app.mjs b/web/static/app.mjs
index a6985f0..429bef7 100644
--- a/web/static/app.mjs
+++ b/web/static/app.mjs
@@ -9,9 +9,12 @@ function status(message, error = false) {
$('status').classList.toggle('error', error);
}
function controls() {
- for (const id of ['go', 'reload']) $(id).disabled = busy;
+ for (const id of ['go', 'reload', 'newfile', 'newdir', 'upload']) $(id).disabled = busy;
$('save').disabled = busy || !current || current.stat.directory || $('editor').hidden || !dirty;
- $('download').disabled = busy || !current;
+ $('download').disabled = busy || !current || current.stat.directory;
+ $('follow').disabled = busy || !current || current.stat.directory;
+ $('rename').disabled = busy || currentPath === '/';
+ $('delete').disabled = busy || currentPath === '/';
$('editor').readOnly = busy;
$('edit-status').textContent = dirty ? 'Unsaved changes.' : 'No unsaved changes.';
document.querySelector('main').setAttribute('aria-busy', String(busy));
@@ -102,6 +105,7 @@ function breadcrumbs(path) {
}
async function navigate(path, historyMode = 'push') {
path = normalizePath(path);
+ stopFollow();
status('Loading…');
const result = await readPath(path);
current = result; currentPath = path; dirty = false;
@@ -140,6 +144,8 @@ async function navigate(path, historyMode = 'push') {
$('edit-actions').hidden = !editable;
status('File loaded.');
}
+ updateStat(result);
+ markTree(path);
}
function mayLeave() { return !dirty || window.confirm('Discard unsaved changes?'); }
document.addEventListener('click', event => {
@@ -180,4 +186,175 @@ window.addEventListener('popstate', () => {
window.addEventListener('beforeunload', event => {
if (dirty) { event.preventDefault(); event.returnValue = ''; }
});
-action(() => navigate(locationPath(), 'none'));
+
+// ---- HTTP /fs view: tree, stat panel, create/rename/delete/upload, follow ----
+// These operate over the gateway's HTTP mapping of the same 9P tree, so they
+// stay independent of the WebSocket editing session above.
+function fsURL(path) { return '/fs' + path.split('/').map(encodeURIComponent).join('/'); }
+function parentOf(path) { return path.replace(/\/[^/]+$/, '') || '/'; }
+function baseName(path) { return path.split('/').filter(Boolean).pop() || ''; }
+function dirContext() { return current && current.stat.directory ? currentPath : parentOf(currentPath); }
+function sortEntries(a, b) { return Number(b.dir) - Number(a.dir) || a.name.localeCompare(b.name); }
+
+async function fetchDir(path) {
+ const response = await fetch(fsURL(path) + (path.endsWith('/') ? '' : '/'), { headers: { accept: 'application/json' } });
+ if (!response.ok) throw new Error(`list failed: ${response.status}`);
+ return response.json();
+}
+
+function renderStat(rows) {
+ const dl = $('stat');
+ dl.replaceChildren();
+ for (const [key, value] of rows) {
+ const dt = document.createElement('dt'); dt.textContent = key;
+ const dd = document.createElement('dd'); dd.textContent = String(value);
+ dl.append(dt, dd);
+ }
+}
+function updateStat(result) {
+ const s = result.stat;
+ const rows = [
+ ['name', s.name || baseName(currentPath) || '/'],
+ ['type', s.directory ? 'directory' : 'file'],
+ ['mode', mode(s)],
+ ['length', s.directory ? '—' : (s.length ?? 0).toLocaleString()],
+ ];
+ renderStat(rows);
+ $('stat-panel').hidden = false;
+ // Augment with qid and mtime from a HEAD, best-effort.
+ fetch(fsURL(currentPath), { method: 'HEAD' }).then(response => {
+ if (!response.ok) return;
+ const qid = response.headers.get('x-9p-qid');
+ const mtime = response.headers.get('x-9p-mtime');
+ const extra = [];
+ if (qid) extra.push(['qid', qid]);
+ if (mtime && mtime !== '0') extra.push(['mtime', new Date(Number(mtime) * 1000).toISOString()]);
+ if (extra.length) renderStat(rows.concat(extra));
+ }).catch(() => {});
+}
+
+function treeItem(entry, path) {
+ const li = document.createElement('li');
+ const node = document.createElement('span');
+ node.className = 'node' + (entry.dir ? ' dir' : '');
+ node.dataset.node = path;
+ const twist = document.createElement('span');
+ twist.className = 'twist';
+ twist.textContent = entry.dir ? '▸' : ' ';
+ node.append(twist, document.createTextNode(entry.name));
+ li.append(node);
+ node.onclick = () => {
+ if (entry.dir) toggleTree(li, path, twist);
+ else if (!busy && mayLeave()) action(() => navigate(path));
+ };
+ return li;
+}
+async function toggleTree(li, path, twist) {
+ const open = li.querySelector(':scope > ul');
+ if (open) { open.remove(); twist.textContent = '▸'; return; }
+ twist.textContent = '▾';
+ const ul = document.createElement('ul');
+ li.append(ul);
+ try {
+ for (const entry of (await fetchDir(path)).sort(sortEntries)) ul.append(treeItem(entry, child(path, entry.name)));
+ } catch (error) { ul.textContent = '(unavailable)'; }
+}
+async function loadTree() {
+ const root = $('filetree');
+ try {
+ const entries = (await fetchDir('/')).sort(sortEntries);
+ root.replaceChildren(...entries.map(entry => treeItem(entry, '/' + entry.name)));
+ } catch { /* upstream not ready; leave the last tree in place */ }
+}
+function markTree(path) {
+ for (const node of $('filetree').querySelectorAll('.node.current')) node.classList.remove('current');
+ const match = $('filetree').querySelector(`.node[data-node="${CSS.escape(path)}"]`);
+ if (match) match.classList.add('current');
+}
+
+let stream = null;
+function stopFollow() {
+ if (stream) { stream.close(); stream = null; }
+ $('follow').classList.remove('active');
+ $('stream').hidden = true;
+ $('stream').textContent = '';
+}
+function startFollow() {
+ if (!current || current.stat.directory) return;
+ $('editor').hidden = true; $('binary').hidden = true; $('edit-actions').hidden = true;
+ const view = $('stream');
+ view.hidden = false; view.textContent = '';
+ $('follow').classList.add('active');
+ stream = new EventSource(fsURL(currentPath) + '?follow=1');
+ stream.onmessage = event => { view.textContent += event.data + '\n'; view.scrollTop = view.scrollHeight; };
+ stream.addEventListener('eof', stopFollow);
+ stream.onerror = () => { status('Stream ended.'); stopFollow(); };
+}
+$('follow').onclick = () => {
+ if ($('follow').classList.contains('active')) { stopFollow(); action(() => navigate(currentPath, 'none')); }
+ else startFollow();
+};
+
+async function mutate(fn, done) {
+ if (busy) return;
+ busy = true; controls();
+ try { await fn(); await loadTree(); status(done); }
+ catch (error) { status(error.message, true); }
+ finally { busy = false; controls(); }
+}
+function needOk(response, what) {
+ if (!response.ok) throw new Error(`${what} failed: ${response.status}`);
+}
+$('newfile').onclick = () => {
+ const name = window.prompt('New file name:');
+ if (!name) return;
+ const target = child(dirContext(), name);
+ mutate(async () => {
+ needOk(await fetch(fsURL(target), { method: 'PUT', body: '' }), 'Create');
+ await navigate(target);
+ }, `Created ${name}.`);
+};
+$('newdir').onclick = () => {
+ const name = window.prompt('New folder name:');
+ if (!name) return;
+ const target = child(dirContext(), name);
+ mutate(async () => {
+ needOk(await fetch(fsURL(target) + '/', { method: 'PUT' }), 'Create folder');
+ await navigate(dirContext(), 'none');
+ }, `Created ${name}/.`);
+};
+$('delete').onclick = () => {
+ if (currentPath === '/' || !window.confirm(`Delete ${currentPath}?`)) return;
+ const parent = parentOf(currentPath);
+ mutate(async () => {
+ needOk(await fetch(fsURL(currentPath), { method: 'DELETE' }), 'Delete');
+ await navigate(parent);
+ }, 'Deleted.');
+};
+$('rename').onclick = () => {
+ if (!current || current.stat.directory) { status('Rename is available for files.', true); return; }
+ const name = window.prompt('Rename to:', baseName(currentPath));
+ if (!name) return;
+ const target = child(parentOf(currentPath), name);
+ const from = currentPath;
+ mutate(async () => {
+ // A copy-then-remove rename over HTTP; not atomic.
+ const bytes = new Uint8Array(await (await fetch(fsURL(from))).arrayBuffer());
+ needOk(await fetch(fsURL(target), { method: 'PUT', body: bytes }), 'Rename (write)');
+ needOk(await fetch(fsURL(from), { method: 'DELETE' }), 'Rename (remove)');
+ await navigate(target);
+ }, `Renamed to ${name}.`);
+};
+$('upload').onclick = () => { if (!busy) $('upload-input').click(); };
+$('upload-input').onchange = () => {
+ const file = $('upload-input').files[0];
+ if (!file) return;
+ const target = child(dirContext(), file.name);
+ mutate(async () => {
+ needOk(await fetch(fsURL(target), { method: 'PUT', body: file }), 'Upload');
+ await navigate(target);
+ }, `Uploaded ${file.name}.`);
+ $('upload-input').value = '';
+};
+
+action(() => navigate(locationPath(), 'none')).then(loadTree);
diff --git a/web/static/index.html b/web/static/index.html
index aeca5d3..b012b3c 100644
--- a/web/static/index.html
+++ b/web/static/index.html
@@ -15,39 +15,57 @@
<nav class="toolbar" aria-label="File actions">
<a class="selected" href="/" data-path="/">files</a>
<button id="reload" type="button">reload</button>
+ <button id="newfile" type="button">new file</button>
+ <button id="newdir" type="button">new folder</button>
+ <button id="rename" type="button">rename</button>
+ <button id="delete" type="button">delete</button>
+ <button id="upload" type="button">upload</button>
+ <input id="upload-input" type="file" hidden>
<form id="path-form">
<label for="path">path</label>
<input id="path" name="path" value="/" aria-label="Path" autocomplete="off" spellcheck="false">
<button id="go">go</button>
</form>
</nav>
- <main>
- <nav id="breadcrumbs" aria-label="Breadcrumb"><a href="/" data-path="/">root</a></nav>
- <div class="summary">
- <p id="status" role="status" aria-live="polite">Loading files…</p>
- <a id="up" href="/" data-path="/" hidden>parent directory</a>
- </div>
- <section id="directory" aria-label="Directory">
- <table>
- <thead><tr><th scope="col">Mode</th><th scope="col">Name</th><th scope="col" class="size">Size</th></tr></thead>
- <tbody id="entries"></tbody>
- </table>
- <p id="empty" hidden>This directory is empty.</p>
- </section>
- <section id="file" hidden>
- <div class="file-heading">
- <h1 id="filename"></h1>
- <span id="file-mode"></span><span id="file-size"></span>
- <button id="download" type="button">download</button>
+ <div class="layout">
+ <aside id="sidebar" aria-label="Tree">
+ <div class="tree-head">tree</div>
+ <ul id="filetree" role="tree"></ul>
+ </aside>
+ <main>
+ <nav id="breadcrumbs" aria-label="Breadcrumb"><a href="/" data-path="/">root</a></nav>
+ <div class="summary">
+ <p id="status" role="status" aria-live="polite">Loading files…</p>
+ <a id="up" href="/" data-path="/" hidden>parent directory</a>
</div>
- <textarea id="editor" aria-label="File contents" spellcheck="false"></textarea>
- <p id="binary" hidden>Binary file. Download to view its contents.</p>
- <div id="edit-actions" class="actions">
- <button id="save" type="button">save changes</button>
- <span id="edit-status">No unsaved changes.</span>
- </div>
- </section>
- </main>
+ <section id="directory" aria-label="Directory">
+ <table>
+ <thead><tr><th scope="col">Mode</th><th scope="col">Name</th><th scope="col" class="size">Size</th></tr></thead>
+ <tbody id="entries"></tbody>
+ </table>
+ <p id="empty" hidden>This directory is empty.</p>
+ </section>
+ <section id="file" hidden>
+ <div class="file-heading">
+ <h1 id="filename"></h1>
+ <span id="file-mode"></span><span id="file-size"></span>
+ <button id="follow" type="button">follow</button>
+ <button id="download" type="button">download</button>
+ </div>
+ <textarea id="editor" aria-label="File contents" spellcheck="false"></textarea>
+ <pre id="stream" hidden aria-label="Live stream"></pre>
+ <p id="binary" hidden>Binary file. Download to view its contents.</p>
+ <div id="edit-actions" class="actions">
+ <button id="save" type="button">save changes</button>
+ <span id="edit-status">No unsaved changes.</span>
+ </div>
+ </section>
+ <section id="stat-panel" aria-label="Stat" hidden>
+ <h2>stat</h2>
+ <dl id="stat"></dl>
+ </section>
+ </main>
+ </div>
<footer>served by cloud9</footer>
</body>
</html>
diff --git a/web/static/style.css b/web/static/style.css
index a8605f5..965406b 100644
--- a/web/static/style.css
+++ b/web/static/style.css
@@ -66,6 +66,26 @@ textarea { display: block; width: 100%; min-height: 420px; resize: vertical; bor
#binary { color: #777; padding: 16px 8px; }
.actions { display: flex; align-items: center; gap: 12px; margin-top: 10px; font-size: 12px; }
#edit-status { color: #777; }
+.layout { display: flex; align-items: flex-start; gap: 16px; }
+#sidebar { flex: 0 0 220px; border-right: 1px solid #e2e2e2; padding: 8px 10px 0 0; min-width: 0; }
+.tree-head { color: #999; font-size: 11px; text-transform: uppercase; letter-spacing: .05em; padding: 4px; }
+#filetree, #filetree ul { list-style: none; margin: 0; padding: 0; }
+#filetree ul { margin-left: 12px; }
+#filetree li { line-height: 1.7; white-space: nowrap; }
+#filetree .node { cursor: pointer; font-family: monospace; overflow-wrap: anywhere; }
+#filetree .node.dir { font-weight: bold; }
+#filetree .node.current { background: #ccc; }
+#filetree .twist { display: inline-block; width: 12px; color: #999; }
+.layout main { flex: 1 1 auto; min-width: 0; }
+#stat-panel { margin-top: 18px; border-top: 1px solid #eee; padding-top: 8px; }
+#stat-panel h2 { font-size: 12px; text-transform: uppercase; letter-spacing: .05em; color: #999; margin: 0 0 6px; }
+#stat { display: grid; grid-template-columns: max-content 1fr; gap: 2px 12px; font: 12px monospace; margin: 0; }
+#stat dt { color: #888; }
+#stat dd { margin: 0; overflow-wrap: anywhere; }
+#stream { display: block; width: 100%; min-height: 300px; max-height: 60vh; overflow: auto; border: 1px solid #ddd; border-top: 0; padding: 10px; margin: 0; font: 13px/1.5 monospace; background: #fbfbfb; white-space: pre-wrap; }
+#follow { margin-left: auto; font-size: 12px; }
+#follow.active { background: #ccc; color: #111; }
+#follow + #download { margin-left: 0; }
footer { margin: 30px 8px 0; padding-top: 7px; border-top: 1px solid #ddd; text-align: right; color: #999; font-size: 11px; }
a:focus-visible, button:focus-visible { outline: 2px solid #07539b; outline-offset: 2px; }
[hidden] { display: none !important; }
@@ -84,4 +104,6 @@ a:focus-visible, button:focus-visible { outline: 2px solid #07539b; outline-offs
.summary { align-items: start; }
.file-heading { gap: 7px 12px; }
textarea { min-height: 350px; }
+ .layout { flex-direction: column; }
+ #sidebar { flex: none; width: 100%; border-right: 0; border-bottom: 1px solid #e2e2e2; max-height: 180px; overflow: auto; padding: 0 0 6px; }
}