1#[cfg(feature = "tracing")]
2use crate::telemetry::Telemetry;
3use crate::{
4 capnp_util,
5 child::Child,
6 container_io::{ContainerIO, SharedContainerIO},
7 container_log::ContainerLog,
8 pause::Pause,
9 server::{GenerateRuntimeArgs, Server},
10 version::Version,
11};
12use anyhow::{Context, format_err};
13use capnp::{Error, capability::Promise};
14use capnp_rpc::pry;
15use conmon_common::conmon_capnp::conmon;
16use std::{
17 path::{Path, PathBuf},
18 process,
19 rc::Rc,
20 str,
21 time::Duration,
22};
23use tokio::time::Instant;
24use tracing::{Instrument, debug, debug_span, error};
25use uuid::Uuid;
26
27macro_rules! pry_err {
28 ($x:expr_2021) => {
29 pry!(capnp_err!($x))
30 };
31}
32
33macro_rules! capnp_err {
34 ($x:expr_2021) => {
35 $x.map_err(|e| Error::failed(format!("{:#}", e)))
36 };
37}
38
39#[cfg(feature = "tracing")]
40fn set_parent_context(
41 reader: capnp::struct_list::Reader<'_, conmon::text_text_map_entry::Owned>,
42) -> Result<(), Error> {
43 Telemetry::set_parent_context(reader).map_err(|e| Error::failed(format!("{:#}", e)))
44}
45
46#[cfg(not(feature = "tracing"))]
47fn set_parent_context(
48 _reader: capnp::struct_list::Reader<'_, conmon::text_text_map_entry::Owned>,
49) -> Result<(), Error> {
50 Ok(())
51}
52
53macro_rules! new_root_span {
54 ($name:expr_2021, $container_id:expr_2021) => {
55 debug_span!(
56 $name,
57 container_id = $container_id,
58 uuid = %Uuid::new_v4()
59 )
60 };
61}
62
63macro_rules! capnp_text_list {
65 ($x:expr_2021) => {
66 pry!(pry!($x).iter().collect::<Result<Vec<_>, _>>())
67 };
68}
69
70macro_rules! capnp_vec_str {
71 ($x:expr_2021) => {
72 pry!(
73 capnp_text_list!($x)
74 .iter()
75 .map(|s| s.to_string())
76 .collect::<Result<Vec<_>, _>>()
77 )
78 };
79}
80
81macro_rules! capnp_vec_path {
82 ($x:expr_2021) => {
83 pry!(
84 capnp_text_list!($x)
85 .iter()
86 .map(|s| s.to_str().map(|x| PathBuf::from(x)))
87 .collect::<Result<Vec<_>, _>>()
88 )
89 };
90}
91
92#[allow(refining_impl_trait_reachable)]
93impl conmon::Server for Server {
94 fn version(
96 self: Rc<Server>,
97 params: conmon::VersionParams,
98 mut results: conmon::VersionResults,
99 ) -> Promise<(), capnp::Error> {
100 debug!("Got a version request");
101 let req = pry!(pry!(params.get()).get_request());
102
103 let span = debug_span!("version", uuid = %Uuid::new_v4());
104 let _enter = span.enter();
105 pry!(set_parent_context(pry!(req.get_metadata())));
106
107 let version = Version::new(req.get_verbose());
108 let mut response = results.get().init_response();
109 response.set_process_id(process::id());
110 response.set_version(version.version());
111 response.set_tag(version.tag());
112 response.set_commit(version.commit());
113 response.set_build_date(version.build_date());
114 response.set_target(version.target());
115 response.set_rust_version(version.rust_version());
116 response.set_cargo_version(version.cargo_version());
117 response.set_cargo_tree(version.cargo_tree());
118
119 Promise::ok(())
120 }
121
122 fn create_container(
124 self: Rc<Server>,
125 params: conmon::CreateContainerParams,
126 mut results: conmon::CreateContainerResults,
127 ) -> Promise<(), capnp::Error> {
128 let req = pry!(pry!(params.get()).get_request());
129 let id = pry!(pry!(req.get_id()).to_string());
130
131 let span = new_root_span!("create_container", id.as_str());
132 let _enter = span.enter();
133 pry!(set_parent_context(pry!(req.get_metadata())));
134
135 let cleanup_cmd: Vec<String> = capnp_vec_str!(req.get_cleanup_cmd());
136
137 debug!("Got a create container request");
138
139 let log_drivers = pry!(req.get_log_drivers());
140 let container_log = pry_err!(ContainerLog::from(log_drivers));
141 let mut container_io =
142 pry_err!(ContainerIO::new(req.get_terminal(), container_log.clone()));
143
144 let bundle_path = Path::new(pry!(pry!(req.get_bundle_path()).to_str()));
145 let pidfile = bundle_path.join("pidfile");
146 debug!("PID file is {}", pidfile.display());
147
148 let child_reaper = self.reaper().clone();
149 let global_args = pry!(req.get_global_args());
150 let command_args = pry!(req.get_command_args());
151 let cgroup_manager = pry!(req.get_cgroup_manager());
152 let args = GenerateRuntimeArgs {
153 config: self.config(),
154 id: &id,
155 container_io: &container_io,
156 pidfile: &pidfile,
157 cgroup_manager,
158 };
159 let args = pry_err!(args.create_args(bundle_path, global_args, command_args));
160 let stdin = req.get_stdin();
161 let runtime = self.config().runtime().clone();
162 let exit_paths = capnp_vec_path!(req.get_exit_paths());
163 let oom_exit_paths = capnp_vec_path!(req.get_oom_exit_paths());
164 let env_vars = pry!(req.get_env_vars().and_then(capnp_util::into_map));
165
166 let additional_fds = pry_err!(self.fd_socket().take_all(pry!(req.get_additional_fds())));
167 let leak_fds = pry_err!(self.fd_socket().take_all(pry!(req.get_leak_fds())));
168
169 Promise::from_future(
170 async move {
171 capnp_err!(container_log.write().await.init().await)?;
172
173 let (grandchild_pid, token) = capnp_err!(match child_reaper
174 .create_child(
175 runtime,
176 args,
177 stdin,
178 &mut container_io,
179 &pidfile,
180 env_vars,
181 additional_fds,
182 )
183 .await
184 {
185 Err(e) => {
186 let (_, stderr, _) =
188 capnp_err!(container_io.read_all_with_timeout(None).await)?;
189 if !stderr.is_empty() {
190 let stderr_str = str::from_utf8(&stderr)?;
191 Err(format_err!("{:#}: {}", e, stderr_str))
192 } else {
193 Err(e)
194 }
195 }
196 res => res,
197 })?;
198
199 let io = SharedContainerIO::new(container_io);
201 let child = Child::new(
202 id,
203 grandchild_pid,
204 exit_paths,
205 oom_exit_paths,
206 None,
207 io,
208 cleanup_cmd,
209 token,
210 );
211 capnp_err!(child_reaper.watch_grandchild(child, leak_fds))?;
212
213 results
214 .get()
215 .init_response()
216 .set_container_pid(grandchild_pid);
217 Ok(())
218 }
219 .instrument(debug_span!("promise")),
220 )
221 }
222
223 fn exec_sync_container(
225 self: Rc<Server>,
226 params: conmon::ExecSyncContainerParams,
227 mut results: conmon::ExecSyncContainerResults,
228 ) -> Promise<(), capnp::Error> {
229 let req = pry!(pry!(params.get()).get_request());
230 let id = pry!(pry!(req.get_id()).to_string());
231
232 let span = new_root_span!("exec_sync_container", id.as_str());
233 let _enter = span.enter();
234 pry!(set_parent_context(pry!(req.get_metadata())));
235
236 let timeout = req.get_timeout_sec();
237
238 let pidfile =
239 ContainerIO::temp_file_name(Some(self.config().runtime_dir()), "exec_sync", "pid");
240
241 debug!("Got exec sync container request with timeout {}", timeout);
242
243 let runtime = self.config().runtime().clone();
244 let child_reaper = self.reaper().clone();
245
246 let logger = ContainerLog::new();
247 let mut container_io = pry_err!(ContainerIO::new(req.get_terminal(), logger));
248
249 let command = pry!(req.get_command());
250 let env_vars = pry!(req.get_env_vars().and_then(capnp_util::into_map));
251 let cgroup_manager = pry!(req.get_cgroup_manager());
252
253 let args = GenerateRuntimeArgs {
254 config: self.config(),
255 id: &id,
256 container_io: &container_io,
257 pidfile: &pidfile,
258 cgroup_manager,
259 };
260 let args = pry_err!(args.exec_sync_args(command));
261
262 Promise::from_future(
263 async move {
264 match child_reaper
265 .create_child(
266 &runtime,
267 &args,
268 false,
269 &mut container_io,
270 &pidfile,
271 env_vars,
272 vec![],
273 )
274 .await
275 {
276 Ok((grandchild_pid, token)) => {
277 let time_to_timeout = if timeout > 0 {
278 Some(Instant::now() + Duration::from_secs(timeout))
279 } else {
280 None
281 };
282 let mut resp = results.get().init_response();
283 let io = SharedContainerIO::new(container_io);
285 let io_clone = io.clone();
286 let child = Child::new(
287 id,
288 grandchild_pid,
289 vec![],
290 vec![],
291 time_to_timeout,
292 io_clone,
293 vec![],
294 token.clone(),
295 );
296
297 let mut exit_rx = capnp_err!(child_reaper.watch_grandchild(child, vec![]))?;
298
299 let (stdout, stderr, timed_out) =
300 capnp_err!(io.read_all_with_timeout(time_to_timeout).await)?;
301
302 let exit_data = capnp_err!(exit_rx.recv().await)?;
303 resp.set_stdout(&stdout);
304 resp.set_stderr(&stderr);
305 resp.set_exit_code(exit_data.exit_code);
306 if timed_out || exit_data.timed_out {
307 resp.set_timed_out(true);
308 }
309 }
310 Err(e) => {
311 error!("Unable to create child: {:#}", e);
312 let mut resp = results.get().init_response();
313 resp.set_exit_code(-2);
314 }
315 }
316 Ok(())
317 }
318 .instrument(debug_span!("promise")),
319 )
320 }
321
322 fn attach_container(
324 self: Rc<Server>,
325 params: conmon::AttachContainerParams,
326 _: conmon::AttachContainerResults,
327 ) -> Promise<(), capnp::Error> {
328 let req = pry!(pry!(params.get()).get_request());
329 let id = pry_err!(pry_err!(req.get_id()).to_str());
330
331 let span = new_root_span!("attach_container", id);
332 let _enter = span.enter();
333 pry!(set_parent_context(pry!(req.get_metadata())));
334
335 debug!("Got a attach container request",);
336
337 let exec_session_id = pry_err!(pry_err!(req.get_exec_session_id()).to_str());
338 if !exec_session_id.is_empty() {
339 debug!("Using exec session id {}", exec_session_id);
340 }
341
342 let socket_path = pry!(pry!(req.get_socket_path()).to_string());
343 let child = pry_err!(self.reaper().get(id));
344 let stop_after_stdin_eof = req.get_stop_after_stdin_eof();
345
346 Promise::from_future(
347 async move {
348 capnp_err!(
349 child
350 .io()
351 .attach()
352 .await
353 .add(&socket_path, child.token().clone(), stop_after_stdin_eof)
354 .await
355 )
356 }
357 .instrument(debug_span!("promise")),
358 )
359 }
360
361 fn reopen_log_container(
363 self: Rc<Server>,
364 params: conmon::ReopenLogContainerParams,
365 _: conmon::ReopenLogContainerResults,
366 ) -> Promise<(), capnp::Error> {
367 let req = pry!(pry!(params.get()).get_request());
368 let id = pry_err!(pry_err!(req.get_id()).to_str());
369
370 let span = new_root_span!("reopen_log_container", id);
371 let _enter = span.enter();
372 pry!(set_parent_context(pry!(req.get_metadata())));
373
374 debug!("Got a reopen container log request");
375
376 let child = pry_err!(self.reaper().get(id));
377
378 Promise::from_future(
379 async move { capnp_err!(child.io().logger().await.write().await.reopen().await) }
380 .instrument(debug_span!("promise")),
381 )
382 }
383
384 fn set_window_size_container(
386 self: Rc<Server>,
387 params: conmon::SetWindowSizeContainerParams,
388 _: conmon::SetWindowSizeContainerResults,
389 ) -> Promise<(), capnp::Error> {
390 let req = pry!(pry!(params.get()).get_request());
391 let id = pry_err!(pry_err!(req.get_id()).to_str());
392
393 let span = new_root_span!("set_window_size_container", id);
394 let _enter = span.enter();
395 pry!(set_parent_context(pry!(req.get_metadata())));
396
397 debug!("Got a set window size container request");
398
399 let child = pry_err!(self.reaper().get(id));
400 let width = req.get_width();
401 let height = req.get_height();
402
403 Promise::from_future(
404 async move { capnp_err!(child.io().resize(width, height).await) }
405 .instrument(debug_span!("promise")),
406 )
407 }
408
409 fn create_namespaces(
411 self: Rc<Server>,
412 params: conmon::CreateNamespacesParams,
413 mut results: conmon::CreateNamespacesResults,
414 ) -> Promise<(), capnp::Error> {
415 debug!("Got a create namespaces request");
416 let req = pry!(pry!(params.get()).get_request());
417 let pod_id = pry_err!(pry_err!(req.get_pod_id()).to_str());
418
419 if pod_id.is_empty() {
420 return Promise::err(Error::failed("no pod ID provided".into()));
421 }
422
423 let span = new_root_span!("create_namespaces", pod_id);
424 let _enter = span.enter();
425 pry!(set_parent_context(pry!(req.get_metadata())));
426
427 let pause = pry_err!(Pause::init_shared(
428 pry!(pry!(req.get_base_path()).to_str()),
429 pod_id,
430 pry!(req.get_namespaces()),
431 capnp_vec_str!(req.get_uid_mappings()),
432 capnp_vec_str!(req.get_gid_mappings()),
433 ));
434
435 let response = results.get().init_response();
436 let mut namespaces =
437 response.init_namespaces(pry_err!(pause.namespaces().len().try_into()));
438
439 for (idx, namespace) in pause.namespaces().iter().enumerate() {
440 let mut ns = namespaces.reborrow().get(pry_err!(idx.try_into()));
441 ns.set_path(
442 namespace
443 .path(pause.base_path(), pod_id)
444 .display()
445 .to_string(),
446 );
447 ns.set_type(namespace.to_capnp_namespace());
448 }
449
450 Promise::ok(())
451 }
452
453 fn start_fd_socket(
454 self: Rc<Server>,
455 params: conmon::StartFdSocketParams,
456 mut results: conmon::StartFdSocketResults,
457 ) -> Promise<(), capnp::Error> {
458 let req = pry!(pry!(params.get()).get_request());
459
460 let span = debug_span!(
461 "start_fd_socket",
462 uuid = %Uuid::new_v4()
463 );
464 let _enter = span.enter();
465 pry!(set_parent_context(pry!(req.get_metadata())));
466
467 debug!("Got a start fd socket request");
468
469 let path = self.config().fd_socket();
470 let fd_socket = self.fd_socket().clone();
471
472 Promise::from_future(
473 async move {
474 let path = capnp_err!(fd_socket.start(path).await)?;
475
476 let mut resp = results.get().init_response();
477 resp.set_path(capnp_err!(path.to_str().context("fd_socket path to str"))?);
478
479 Ok(())
480 }
481 .instrument(debug_span!("promise")),
482 )
483 }
484
485 fn serve_exec_container(
486 self: Rc<Server>,
487 params: conmon::ServeExecContainerParams,
488 mut results: conmon::ServeExecContainerResults,
489 ) -> Promise<(), capnp::Error> {
490 debug!("Got a serve exec container request");
491 let req = pry!(pry!(params.get()).get_request());
492
493 let span = debug_span!(
494 "serve_exec_container",
495 uuid = %Uuid::new_v4()
496 );
497 let _enter = span.enter();
498 pry!(set_parent_context(pry!(req.get_metadata())));
499
500 let id = pry_err!(pry_err!(req.get_id()).to_string());
501
502 pry_err!(self.reaper().get(&id));
504
505 let command = capnp_vec_str!(req.get_command());
506 let (tty, stdin, stdout, stderr) = (
507 req.get_tty(),
508 req.get_stdin(),
509 req.get_stdout(),
510 req.get_stderr(),
511 );
512
513 let streaming_server = self.streaming_server().clone();
514 let child_reaper = self.reaper().clone();
515 let container_io = pry_err!(ContainerIO::new(tty, ContainerLog::new()));
516 let config = self.config().clone();
517 let cgroup_manager = pry!(req.get_cgroup_manager());
518
519 Promise::from_future(
520 async move {
521 capnp_err!(
522 streaming_server
523 .write()
524 .await
525 .start_if_required()
526 .await
527 .context("start streaming server if required")
528 )?;
529
530 let url = streaming_server
531 .read()
532 .await
533 .exec_url(
534 child_reaper,
535 container_io,
536 config,
537 cgroup_manager,
538 id.into_boxed_str(),
539 command,
540 stdin,
541 stdout,
542 stderr,
543 )
544 .await;
545
546 results.get().init_response().set_url(&url);
547 Ok(())
548 }
549 .instrument(debug_span!("promise")),
550 )
551 }
552
553 fn serve_attach_container(
554 self: Rc<Server>,
555 params: conmon::ServeAttachContainerParams,
556 mut results: conmon::ServeAttachContainerResults,
557 ) -> Promise<(), capnp::Error> {
558 debug!("Got a serve attach container request");
559 let req = pry!(pry!(params.get()).get_request());
560
561 let span = debug_span!(
562 "serve_attach_container",
563 uuid = %Uuid::new_v4()
564 );
565 let _enter = span.enter();
566 pry!(set_parent_context(pry!(req.get_metadata())));
567
568 let id = pry_err!(pry_err!(req.get_id()).to_str());
569 let (stdin, stdout, stderr) = (req.get_stdin(), req.get_stdout(), req.get_stderr());
570
571 let streaming_server = self.streaming_server().clone();
572 let child = pry_err!(self.reaper().get(id));
573
574 Promise::from_future(
575 async move {
576 capnp_err!(
577 streaming_server
578 .write()
579 .await
580 .start_if_required()
581 .await
582 .context("start streaming server")
583 )?;
584
585 let url = streaming_server
586 .read()
587 .await
588 .attach_url(child, stdin, stdout, stderr)
589 .await;
590
591 results.get().init_response().set_url(&url);
592 Ok(())
593 }
594 .instrument(debug_span!("promise")),
595 )
596 }
597
598 fn serve_port_forward_container(
599 self: Rc<Server>,
600 params: conmon::ServePortForwardContainerParams,
601 mut results: conmon::ServePortForwardContainerResults,
602 ) -> Promise<(), capnp::Error> {
603 debug!("Got a serve port forward container request");
604 let req = pry!(pry!(params.get()).get_request());
605
606 let span = debug_span!(
607 "serve_port_forward_container",
608 uuid = %Uuid::new_v4()
609 );
610 let _enter = span.enter();
611 pry!(set_parent_context(pry!(req.get_metadata())));
612
613 let net_ns_path = pry_err!(pry_err!(req.get_net_ns_path()).to_string());
614 let streaming_server = self.streaming_server().clone();
615
616 Promise::from_future(
617 async move {
618 capnp_err!(
619 streaming_server
620 .write()
621 .await
622 .start_if_required()
623 .await
624 .context("start streaming server if required")
625 )?;
626
627 let url = streaming_server
628 .read()
629 .await
630 .port_forward_url(net_ns_path)
631 .await;
632
633 results.get().init_response().set_url(&url);
634 Ok(())
635 }
636 .instrument(debug_span!("promise")),
637 )
638 }
639}