//! 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/`. 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 = 64; pub const MuxT = mux.Mux(.{ .msize = max_frame, .upstream_fids = 2048, .fids_per_conn = 256, .downstreams = max_connections, }); const Config = struct { host: []const u8, origin: []const u8, public_host: []const u8, timeout_ms: u32, 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 = ""; var public_origin: ?[]const u8 = null; 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] \\ [--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; 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; } 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) }); 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; const browser_config = try std.json.Stringify.valueAlloc(allocator, .{ .user = user, .tree = tree }, .{}); 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 = 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; }; s.app = &app; s.stream = stream; group.concurrent(io, handleHttp, .{s}) catch { app.giveBack(s); stream.close(io); continue; }; } } /// 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); }; } } /// 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' 'unsafe-inline'; connect-src 'self'; frame-ancestors 'none'; base-uri 'none'" }, } }); } 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; var headers = request.iterateHeaders(); while (headers.next()) |header| { if (std.ascii.eqlIgnoreCase(header.name, "host")) host = header.value; 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); 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") 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 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, } } }