diff options
Diffstat (limited to 'web/main.zig')
| -rw-r--r-- | web/main.zig | 411 |
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, + } + } } |
