Skip to main content

conmonrs/
rpc.rs

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
63/// capnp_text_list takes text_list as an input and outputs list of text.
64macro_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    /// Retrieve version information from the server.
95    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    /// Create a new container for the provided parameters.
123    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                        // Attach the stderr output to the error message
187                        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                // register grandchild with server
200                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    /// Execute a command in sync inside of a container.
224    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                        // register grandchild with server
284                        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    /// Attach to a running container.
323    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    /// Rotate all log drivers for a running container.
362    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    /// Adjust the window size of a container running inside of a terminal.
385    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    /// Create a new set of namespaces.
410    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        // Validate that the container actually exists
503        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}