summaryrefslogtreecommitdiff
path: root/web/probe.zig
blob: a6f5ceb22be46e2a53d714c96f5e6faf3dc48b89 (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
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;
            }
        }
    };
}