summaryrefslogtreecommitdiff
path: root/web/main.zig
diff options
context:
space:
mode:
Diffstat (limited to 'web/main.zig')
-rw-r--r--web/main.zig411
1 files changed, 292 insertions, 119 deletions
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,
+ }
+ }
}