1pub 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
56pub type Healthz = Arc<dyn Fn() -> (bool, Value) + Send + Sync>;
59
60pub type Routes = Arc<dyn Fn(&http::Request) -> Option<http::Response> + Send + Sync>;
63
64#[derive(Clone)]
65pub enum ListenerKind {
66 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#[derive(Clone)]
82pub struct Listener {
83 pub kind: ListenerKind,
84 pub access: Option<Arc<AccessValidator>>,
86 pub policy: ToolPolicy,
87 pub allow_unauthenticated: bool,
89 pub tailnet: bool,
91 pub routes: Option<Routes>,
93 pub public_routes: Option<Routes>,
96 pub preview: Option<Routes>,
98 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 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 pub fn public_routes(mut self, r: Routes) -> Self {
151 self.public_routes = Some(r);
152 self
153 }
154
155 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 pub fn access_shared(mut self, v: Arc<AccessValidator>) -> Self {
169 self.access = Some(v);
170 self
171 }
172
173 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 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
240pub 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
260pub fn serve(listeners: Vec<Listener>, registry: Registry, healthz: Healthz) -> Result<()> {
262 serve_until(listeners, registry, healthz, Shutdown::on_signals()?)
263}
264
265pub 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
275pub 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
296pub 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
318pub 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
329pub 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 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}