Skip to main content

isb_server/server/
mod.rs

1//! `isb serve`: the MCP server layer, with the tools supplied by the embedder.
2//!
3//! Security model:
4//! - TCP listeners bind loopback only. Remote clients arrive through a
5//!   cloudflared tunnel, behind a Cloudflare Access application with Managed
6//!   OAuth; Cloudflare runs the OAuth flow and isb stays the resource origin.
7//! - With Access configured, every `/mcp` request must carry a valid
8//!   `Cf-Access-Jwt-Assertion` for the application's audience, so a request
9//!   reaching the port by another route is still refused.
10//! - A TCP listener without Access is refused unless the embedder opts in, and
11//!   then only accepts browser requests whose `Origin` is localhost, which
12//!   blocks DNS rebinding.
13//! - The unix socket (0600, in a 0700 directory) is the trusted local path:
14//!   filesystem permissions are the gate, and its callers are
15//!   [`Caller::Local`], which [`Caller::is_trusted`] reports.
16//! - `/healthz` never requires auth and reveals only what the embedder puts in it.
17//! - A listener can carry extra [`Routes`] (`isb serve` mounts the identity
18//!   endpoints, `/api/v1/auth/*`, this way). They authenticate their own
19//!   callers; with Access configured they sit behind it, as `/mcp` does.
20//! - [`Listener::public_routes`] are served ahead of Access, for requests
21//!   that carry their own credential (app webhooks, signed by the sender).
22//! - [`Listener::preview`] sees every request first, and takes the ones
23//!   addressed to a preview host (a workspace port's own origin), which it
24//!   authenticates itself; nothing of isb's (UI, API, headers) is served
25//!   on those hosts.
26
27pub mod access;
28pub mod aliases;
29pub use isb_core::serve_client as client;
30#[cfg(test)]
31mod client_tests;
32pub mod http;
33pub mod mcp;
34pub mod openapi;
35pub mod service;
36pub mod ssh;
37pub mod ssh_config;
38pub mod system_service;
39pub mod tailnet;
40pub mod terminal;
41
42use std::path::PathBuf;
43use std::sync::Arc;
44use std::time::Duration;
45
46use serde_json::Value;
47
48pub use access::{AccessValidator, Identity};
49pub use http::Shutdown;
50pub use mcp::{Authenticated, Caller, Hooks, Registry, Tool, ToolHandler, ToolPolicy};
51
52use crate::error::{Error, Result};
53use http::{Handler, HttpListener, HttpServer, Limits};
54use mcp::Endpoint;
55
56/// Health for `GET /healthz`: `(ok, body)`. 200 when ok, else 503. Called on
57/// every probe, so it must be cheap.
58pub type Healthz = Arc<dyn Fn() -> (bool, Value) + Send + Sync>;
59
60/// Extra routes on a listener: `Some` answers the request, `None` leaves it
61/// to the server (a 404).
62pub type Routes = Arc<dyn Fn(&http::Request) -> Option<http::Response> + Send + Sync>;
63
64#[derive(Clone)]
65pub enum ListenerKind {
66    /// `host:port`, which must resolve to loopback only.
67    Tcp(String),
68    Unix(PathBuf),
69}
70
71impl std::fmt::Debug for ListenerKind {
72    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
73        match self {
74            ListenerKind::Tcp(a) => write!(f, "Tcp({a:?})"),
75            ListenerKind::Unix(p) => write!(f, "Unix({p:?})"),
76        }
77    }
78}
79
80/// One address the server answers on, with its own gate and tool policy.
81#[derive(Clone)]
82pub struct Listener {
83    pub kind: ListenerKind,
84    /// Required on TCP unless `allow_unauthenticated`; refused on unix.
85    pub access: Option<Arc<AccessValidator>>,
86    pub policy: ToolPolicy,
87    /// Serve TCP with no Access validation, trusting whatever reaches the port.
88    pub allow_unauthenticated: bool,
89    /// A TCP listener on a tailnet address rather than loopback.
90    pub tailnet: bool,
91    /// Paths other than `/healthz` and `/mcp`.
92    pub routes: Option<Routes>,
93    /// Routes that authenticate every request themselves and are served
94    /// even with Access configured (webhooks, signed by their sender).
95    pub public_routes: Option<Routes>,
96    /// Requests for preview hosts, by `Host`, ahead of everything else.
97    pub preview: Option<Routes>,
98    /// Authentication and authorization the embedder supplies.
99    pub hooks: mcp::Hooks,
100}
101
102impl std::fmt::Debug for Listener {
103    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
104        f.debug_struct("Listener")
105            .field("kind", &self.kind)
106            .field("access", &self.access)
107            .field("policy", &self.policy)
108            .field("allow_unauthenticated", &self.allow_unauthenticated)
109            .field("routes", &self.routes.is_some())
110            .finish()
111    }
112}
113
114impl Listener {
115    pub fn tcp(addr: impl Into<String>) -> Self {
116        Self::new(ListenerKind::Tcp(addr.into()))
117    }
118
119    pub fn unix(path: impl Into<PathBuf>) -> Self {
120        Self::new(ListenerKind::Unix(path.into()))
121    }
122
123    fn new(kind: ListenerKind) -> Self {
124        Listener {
125            kind,
126            access: None,
127            policy: ToolPolicy::default(),
128            allow_unauthenticated: false,
129            tailnet: false,
130            routes: None,
131            public_routes: None,
132            preview: None,
133            hooks: mcp::Hooks::default(),
134        }
135    }
136
137    /// Serve `routes` on this listener too.
138    pub fn hooks(mut self, h: mcp::Hooks) -> Self {
139        self.hooks = h;
140        self
141    }
142
143    pub fn routes(mut self, r: Routes) -> Self {
144        self.routes = Some(r);
145        self
146    }
147
148    /// Serve `r` ahead of Access: only for routes whose every request
149    /// carries its own credential.
150    pub fn public_routes(mut self, r: Routes) -> Self {
151        self.public_routes = Some(r);
152        self
153    }
154
155    /// Let `r` take requests by their `Host` before anything else: it
156    /// answers only for hosts that are its own.
157    pub fn preview(mut self, r: Routes) -> Self {
158        self.preview = Some(r);
159        self
160    }
161
162    pub fn access(mut self, v: AccessValidator) -> Self {
163        self.access = Some(Arc::new(v));
164        self
165    }
166
167    /// [`Self::access`], sharing a validator (and its key cache).
168    pub fn access_shared(mut self, v: Arc<AccessValidator>) -> Self {
169        self.access = Some(v);
170        self
171    }
172
173    /// Bind a tailnet address instead of loopback (no Access: only tailnet
174    /// peers reach it).
175    pub fn tailnet(mut self, yes: bool) -> Self {
176        self.tailnet = yes;
177        self
178    }
179
180    pub fn policy(mut self, p: ToolPolicy) -> Self {
181        self.policy = p;
182        self
183    }
184
185    pub fn allow_unauthenticated(mut self, yes: bool) -> Self {
186        self.allow_unauthenticated = yes;
187        self
188    }
189
190    /// The unix socket is trusted as itself; a TCP caller is trusted only
191    /// as a superadmin, per request.
192    pub fn is_trusted(&self) -> bool {
193        matches!(self.kind, ListenerKind::Unix(_))
194    }
195
196    fn check(&self) -> Result<()> {
197        match (&self.kind, &self.access) {
198            (ListenerKind::Unix(p), Some(_)) => Err(Error::invalid(format!(
199                "unix socket {}: Cloudflare Access applies to TCP listeners only",
200                p.display()
201            ))),
202            (ListenerKind::Tcp(a), Some(_)) if self.tailnet => Err(Error::invalid(format!(
203                "{a}: Cloudflare Access applies to loopback listeners (behind the tunnel), not a tailnet one"
204            ))),
205            (ListenerKind::Tcp(a), None) if !self.allow_unauthenticated => {
206                Err(Error::invalid(format!(
207                    "TCP listener {a} needs Cloudflare Access (team domain and audience), \
208                     or an explicit opt-in to serve it unauthenticated"
209                )))
210            }
211            _ => Ok(()),
212        }
213    }
214
215    fn describe(&self, tools: usize) -> String {
216        match (&self.kind, &self.access) {
217            (ListenerKind::Unix(p), _) => {
218                format!("unix:{} (trusted local, {tools} tools)", p.display())
219            }
220            (ListenerKind::Tcp(a), Some(v)) => format!(
221                "http://{a}/mcp (Cloudflare Access: {}, {tools} tools)",
222                v.issuer()
223            ),
224            (ListenerKind::Tcp(a), None) if self.tailnet => format!(
225                "http://{a}/mcp on the tailnet ({tools} tools): callers sign in with isb API tokens \
226                 or sessions, or are superadmins by tailnet identity"
227            ),
228            (ListenerKind::Tcp(a), None) if self.hooks.authorize.is_some() => format!(
229                "http://{a}/mcp ({tools} tools) without Cloudflare Access: callers sign in \
230                 with isb API tokens or sessions"
231            ),
232            (ListenerKind::Tcp(a), None) => format!(
233                "http://{a}/mcp ({tools} tools) WITHOUT Cloudflare Access: anything that \
234                 reaches this port can call these tools"
235            ),
236        }
237    }
238}
239
240/// Where the CLI and the server meet: `$ISB_SERVE_SOCKET`, else
241/// `$XDG_RUNTIME_DIR/isb/serve.sock`, else a per-uid directory under /tmp.
242/// On macOS, where the daemon runs inside the `isb machine`, it is that
243/// machine's forwarded socket, `~/.isb/machine/isb/serve.sock`.
244pub fn default_socket_path() -> PathBuf {
245    if let Some(s) = std::env::var_os("ISB_SERVE_SOCKET").filter(|s| !s.is_empty()) {
246        return PathBuf::from(s);
247    }
248    #[cfg(target_os = "macos")]
249    if let Ok(s) = crate::machine::serve_socket(crate::machine::DEFAULT_NAME) {
250        return s;
251    }
252    if let Some(d) = std::env::var_os("XDG_RUNTIME_DIR").filter(|s| !s.is_empty()) {
253        return PathBuf::from(d).join("isb/serve.sock");
254    }
255    std::env::temp_dir()
256        .join(format!("isb-{}", rustix::process::getuid().as_raw()))
257        .join("serve.sock")
258}
259
260/// Serve until SIGINT or SIGTERM.
261pub fn serve(listeners: Vec<Listener>, registry: Registry, healthz: Healthz) -> Result<()> {
262    serve_until(listeners, registry, healthz, Shutdown::on_signals()?)
263}
264
265/// [`serve`], with a registry the embedder also keeps (to serve it on
266/// listeners it adds later, [`spawn_private`]).
267pub fn serve_shared(
268    listeners: Vec<Listener>,
269    registry: Arc<Registry>,
270    healthz: Healthz,
271) -> Result<()> {
272    serve_until_shared(listeners, registry, healthz, Shutdown::on_signals()?)
273}
274
275/// What a listener's requests go to: `/healthz`, `/mcp`, the REST surface
276/// and its routes, through its hooks and policy.
277pub fn handler(l: &Listener, registry: Arc<Registry>, healthz: Healthz) -> Handler {
278    let ep = Endpoint {
279        registry,
280        policy: l.policy.clone(),
281        access: l.access.clone(),
282        healthz,
283        routes: l.routes.clone(),
284        public_routes: l.public_routes.clone(),
285        hooks: l.hooks.clone(),
286    };
287    let preview = l.preview.clone();
288    Arc::new(move |r: &http::Request| {
289        if let Some(resp) = preview.as_ref().and_then(|p| p(r)) {
290            return resp;
291        }
292        ep.handle(r)
293    })
294}
295
296/// Serve `handler` on a private (RFC 1918) address that is not loopback,
297/// such as an org bridge's, until `stop` (or a bind failure, returned at
298/// once). The caller's `handler` decides who gets in: nothing about such
299/// an address keeps anyone out.
300pub fn spawn_private(
301    addr: std::net::SocketAddr,
302    handler: Handler,
303    stop: Shutdown,
304) -> Result<std::thread::JoinHandle<()>> {
305    let sock = HttpListener::bind_tcp_private(addr)?;
306    let server = HttpServer::new(Limits::default(), stop);
307    std::thread::Builder::new()
308        .name(format!("isb-listen-{addr}"))
309        .spawn(move || {
310            if let Err(e) = server.run(sock, handler) {
311                eprintln!("isb serve: listener {addr}: {e}");
312            }
313            server.drain(Duration::from_secs(5));
314        })
315        .map_err(|e| Error::Protocol(format!("cannot start a listener thread: {e}")))
316}
317
318/// Serve until `shutdown` is triggered, then give in-flight requests up to 10s.
319/// Every listener is bound before any is served, so a bad one fails startup.
320pub fn serve_until(
321    listeners: Vec<Listener>,
322    registry: Registry,
323    healthz: Healthz,
324    shutdown: Shutdown,
325) -> Result<()> {
326    serve_until_shared(listeners, Arc::new(registry), healthz, shutdown)
327}
328
329/// [`serve_until`] with a shared registry.
330pub fn serve_until_shared(
331    listeners: Vec<Listener>,
332    registry: Arc<Registry>,
333    healthz: Healthz,
334    shutdown: Shutdown,
335) -> Result<()> {
336    if listeners.is_empty() {
337        return Err(Error::invalid("isb serve needs at least one listener"));
338    }
339    let mut bound: Vec<(HttpListener, Handler)> = Vec::new();
340    for l in &listeners {
341        l.check()?;
342        let sock = match &l.kind {
343            ListenerKind::Tcp(a) if l.tailnet => HttpListener::bind_tcp_tailnet(a)?,
344            ListenerKind::Tcp(a) => HttpListener::bind_tcp(a)?,
345            ListenerKind::Unix(p) => HttpListener::bind_unix(p)?,
346        };
347        let tools = registry
348            .tools()
349            .iter()
350            .filter(|t| l.policy.allows(&t.name))
351            .count();
352        let line = l.describe(tools);
353        if l.access.is_none() && !l.is_trusted() && l.hooks.authorize.is_none() {
354            eprintln!("isb serve: WARNING: {line}");
355        } else {
356            eprintln!("isb serve: listening on {line}");
357        }
358        bound.push((sock, handler(l, registry.clone(), healthz.clone())));
359    }
360    let server = HttpServer::new(Limits::default(), shutdown.clone());
361    let threads: Vec<_> = bound
362        .into_iter()
363        .map(|(sock, h)| {
364            let (srv, stop) = (server.clone(), shutdown.clone());
365            std::thread::spawn(move || {
366                let r = srv.run(sock, h);
367                // One listener failing takes the others down with it.
368                stop.trigger();
369                r
370            })
371        })
372        .collect();
373    let mut first = Ok(());
374    for t in threads {
375        let r = t
376            .join()
377            .unwrap_or_else(|_| Err(Error::Protocol("listener thread panicked".into())));
378        if first.is_ok() {
379            first = r;
380        }
381    }
382    server.drain(Duration::from_secs(10));
383    if server.active() > 0 {
384        eprintln!(
385            "isb serve: exiting with {} request(s) still running",
386            server.active()
387        );
388    }
389    first
390}
391
392#[cfg(test)]
393mod tests {
394    use super::*;
395    use serde_json::json;
396
397    #[test]
398    fn listener_config_is_checked() {
399        assert!(Listener::tcp("127.0.0.1:0").check().is_err());
400        assert!(
401            Listener::tcp("127.0.0.1:0")
402                .allow_unauthenticated(true)
403                .check()
404                .is_ok()
405        );
406        let v = || AccessValidator::new("team.cloudflareaccess.com", "aud").unwrap();
407        assert!(Listener::unix("/x").access(v()).check().is_err());
408        assert!(Listener::tcp("127.0.0.1:0").access(v()).check().is_ok());
409        assert!(Listener::unix("/x").is_trusted());
410        let r = serve_until(
411            vec![Listener::tcp("0.0.0.0:0").allow_unauthenticated(true)],
412            Registry::new(),
413            Arc::new(|| (true, json!({}))),
414            Shutdown::new(),
415        );
416        assert!(r.unwrap_err().to_string().contains("not loopback"));
417    }
418}