Skip to main content

starnix_kernel_runner/
serve_protocols.rs

1// Copyright 2022 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use crate::Container;
6use anyhow::{Context as _, Error};
7use fidl::endpoints::{RequestStream, ServerEnd};
8use fidl_fuchsia_component_runner as frunner;
9use fidl_fuchsia_element as felement;
10use fidl_fuchsia_io as fio;
11use fidl_fuchsia_memory_attribution as fattribution;
12use fidl_fuchsia_posix as fposix;
13use fidl_fuchsia_starnix_binder as fbinder;
14use fidl_fuchsia_starnix_container as fstarcontainer;
15use fuchsia_async::{
16    DurationExt, {self as fasync},
17};
18use fuchsia_rcu::RcuReadScope;
19use futures::channel::oneshot;
20use futures::{
21    AsyncReadExt, AsyncWriteExt, Future, FutureExt, StreamExt, TryFutureExt, TryStreamExt, pin_mut,
22    select,
23};
24use starnix_core::execution::{create_init_child_process, execute_task_with_prerun_result};
25use starnix_core::fs::fuchsia::create_fuchsia_pipe;
26use starnix_core::task::dynamic_thread_spawner::SpawnRequestBuilder;
27use starnix_core::task::{CurrentTask, ExitStatus, Kernel, ProcessEntryRef};
28use starnix_core::vfs::buffers::{VecInputBuffer, VecOutputBuffer};
29use starnix_core::vfs::file_server::serve_file_at;
30use starnix_core::vfs::socket::VsockSocket;
31use starnix_core::vfs::{FdFlags, FileHandle};
32use starnix_logging::{log_error, log_warn};
33use starnix_modules_devpts::create_main_and_replica;
34use starnix_modules_framebuffer::Framebuffer;
35use starnix_task_command::TaskCommand;
36use starnix_uapi::auth::Credentials;
37use starnix_uapi::errors::Errno;
38use starnix_uapi::open_flags::OpenFlags;
39use starnix_uapi::signals::UncheckedSignal;
40use starnix_uapi::{errno, error, uapi};
41use std::ffi::CString;
42
43use super::start_component;
44
45pub fn expose_root(
46    system_task: &CurrentTask,
47    server_end: ServerEnd<fio::DirectoryMarker>,
48) -> Result<(), Error> {
49    let root_file = system_task.open_file("/".into(), OpenFlags::RDONLY)?;
50    serve_file_at(server_end.into_channel().into(), system_task, &root_file, Credentials::root())?;
51    Ok(())
52}
53
54pub async fn serve_component_runner(
55    request_stream: frunner::ComponentRunnerRequestStream,
56    system_task: &CurrentTask,
57) -> Result<(), Error> {
58    request_stream
59        .try_for_each_concurrent(None, |event| async {
60            match event {
61                frunner::ComponentRunnerRequest::Start { start_info, controller, .. } => {
62                    if let Err(e) = start_component(start_info, controller, system_task).await {
63                        log_error!("failed to start component: {:?}", e);
64                    }
65                }
66                frunner::ComponentRunnerRequest::_UnknownMethod { ordinal, .. } => {
67                    log_warn!("Unknown ComponentRunner request: {ordinal}");
68                }
69            }
70            Ok(())
71        })
72        .await
73        .map_err(Error::from)
74}
75
76fn to_winsize(window_size: Option<fstarcontainer::ConsoleWindowSize>) -> uapi::winsize {
77    window_size
78        .map(|window_size| uapi::winsize {
79            ws_row: window_size.rows,
80            ws_col: window_size.cols,
81            ws_xpixel: window_size.x_pixels,
82            ws_ypixel: window_size.y_pixels,
83        })
84        .unwrap_or(uapi::winsize::default())
85}
86
87async fn spawn_console(
88    kernel: &Kernel,
89    payload: fstarcontainer::ControllerSpawnConsoleRequest,
90) -> Result<Result<u8, fstarcontainer::SpawnConsoleError>, Error> {
91    if let (Some(console_in), Some(console_out), Some(binary_path)) =
92        (payload.console_in, payload.console_out, payload.binary_path)
93    {
94        let binary_path = CString::new(binary_path)?;
95        let argv = payload
96            .argv
97            .unwrap_or(vec![])
98            .into_iter()
99            .map(CString::new)
100            .collect::<Result<Vec<_>, _>>()?;
101        let environ = payload
102            .environ
103            .unwrap_or(vec![])
104            .into_iter()
105            .map(CString::new)
106            .collect::<Result<Vec<_>, _>>()?;
107        let window_size = to_winsize(payload.window_size);
108        let current_task = create_init_child_process(
109            &kernel.weak_self.upgrade().expect("Kernel must still be alive"),
110            TaskCommand::new(binary_path.as_bytes()),
111            Credentials::with_ids(0, 0),
112            None,
113        )?;
114        let (sender, receiver) = oneshot::channel();
115        let pty = execute_task_with_prerun_result(
116            current_task,
117            move |current_task| {
118                let executable = current_task.open_file_for_exec(
119                    starnix_core::vfs::FdNumber::AT_FDCWD,
120                    binary_path.as_bytes().into(),
121                    OpenFlags::empty(),
122                )?;
123                current_task.exec(executable, binary_path, argv, environ)?;
124                let (pty, pts) = create_main_and_replica(&current_task, window_size)?;
125                let fd_flags = FdFlags::empty();
126                assert_eq!(0, current_task.add_file(pts.clone(), fd_flags)?.raw());
127                assert_eq!(1, current_task.add_file(pts.clone(), fd_flags)?.raw());
128                assert_eq!(2, current_task.add_file(pts, fd_flags)?.raw());
129                Ok(pty)
130            },
131            move |result| {
132                let _ = match result {
133                    Ok(ExitStatus::Exit(exit_code)) => sender.send(Ok(exit_code)),
134                    _ => sender.send(Err(fstarcontainer::SpawnConsoleError::Canceled)),
135                };
136            },
137            None,
138        )?;
139        let _ = forward_to_pty(kernel, console_in, console_out, pty).map_err(|e| {
140            log_error!("failed to forward to terminal {:?}", e);
141        });
142
143        Ok(receiver.await?)
144    } else {
145        Ok(Err(fstarcontainer::SpawnConsoleError::InvalidArgs))
146    }
147}
148
149pub async fn serve_container_controller(
150    request_stream: fstarcontainer::ControllerRequestStream,
151    system_task: &CurrentTask,
152) -> Result<(), Error> {
153    request_stream
154        .map_err(Error::from)
155        .try_for_each_concurrent(None, |event| async {
156            match event {
157                fstarcontainer::ControllerRequest::VsockConnect {
158                    payload:
159                        fstarcontainer::ControllerVsockConnectRequest { port, bridge_socket, .. },
160                    ..
161                } => {
162                    let Some(port) = port else {
163                        log_warn!("vsock connection missing port");
164                        return Ok(());
165                    };
166                    let Some(bridge_socket) = bridge_socket else {
167                        log_warn!("vsock connection missing bridge_socket");
168                        return Ok(());
169                    };
170                    connect_to_vsock(port, bridge_socket, system_task).await.unwrap_or_else(|e| {
171                        log_error!("failed to connect to vsock {:?}", e);
172                    });
173                }
174                fstarcontainer::ControllerRequest::SpawnConsole { payload, responder } => {
175                    responder.send(spawn_console(system_task.kernel(), payload).await?)?;
176                }
177                fstarcontainer::ControllerRequest::GetVmoReferences { payload, responder } => {
178                    if let Some(koid) = payload.koid {
179                        let mut results = vec![];
180                        for thread_group in
181                            system_task.kernel().pids.get_thread_groups(&RcuReadScope::new())
182                        {
183                            if let Ok(leader) = thread_group.leader.get_task() {
184                                if let Ok(files) = leader.files() {
185                                    let fds = files.get_all_fds();
186                                    for fd in fds {
187                                        if let Ok(file) = files.get(fd) {
188                                            if let Ok(memory) = file.get_memory(
189                                                system_task,
190                                                None,
191                                                starnix_core::mm::ProtectionFlags::READ,
192                                            ) {
193                                                let memory_koid = memory
194                                                    .info()
195                                                    .expect("Failed to get memory info")
196                                                    .koid;
197                                                if memory_koid.raw_koid() == koid {
198                                                    let process_name = thread_group
199                                                        .process
200                                                        .get_name()
201                                                        .unwrap_or_default();
202                                                    results.push(fstarcontainer::VmoReference {
203                                                        process_name: Some(
204                                                            process_name.to_string(),
205                                                        ),
206                                                        pid: Some(leader.get_pid() as u64),
207                                                        fd: Some(fd.raw()),
208                                                        koid: Some(koid),
209                                                        ..Default::default()
210                                                    });
211                                                }
212                                            }
213                                        }
214                                    }
215                                }
216                            }
217                        }
218                        let _ =
219                            responder.send(&fstarcontainer::ControllerGetVmoReferencesResponse {
220                                references: Some(results),
221                                ..Default::default()
222                            });
223                    }
224                }
225                fstarcontainer::ControllerRequest::GetJobHandle { responder } => {
226                    let _result = responder.send(fstarcontainer::ControllerGetJobHandleResponse {
227                        job: Some(
228                            fuchsia_runtime::job_default()
229                                .duplicate_handle(zx::Rights::SAME_RIGHTS)
230                                .expect("Failed to dup handle"),
231                        ),
232                        ..Default::default()
233                    });
234                }
235                fstarcontainer::ControllerRequest::SendSignal {
236                    payload:
237                        fstarcontainer::ControllerSendSignalRequest {
238                            pid: Some(pid),
239                            signal: Some(signal),
240                            ..
241                        },
242                    responder,
243                } => {
244                    let pids = &system_task.kernel().pids;
245                    if let Some(ProcessEntryRef::Process(target_thread_group)) =
246                        pids.get(pid).ok().and_then(|p| p.get_process())
247                    {
248                        #[allow(
249                            clippy::undocumented_unsafe_blocks,
250                            reason = "Force documented unsafe blocks in Starnix"
251                        )]
252                        match unsafe {
253                            target_thread_group.send_signal_unchecked_debug(
254                                system_task,
255                                UncheckedSignal::new(signal),
256                            )
257                        } {
258                            Ok(_) => {
259                                let _result = responder.send(Ok(()));
260                            }
261                            Err(_) => {
262                                let _result =
263                                    responder.send(Err(fstarcontainer::SignalError::InvalidSignal));
264                            }
265                        };
266                    } else {
267                        let _result =
268                            responder.send(Err(fstarcontainer::SignalError::InvalidTarget));
269                    };
270                }
271                // The request did not contain both a signal and a target pid.
272                fstarcontainer::ControllerRequest::SendSignal { responder, .. } => {
273                    log_error!("malformed SendSignal request");
274                    let _result = responder.send(Err(fstarcontainer::SignalError::InvalidTarget));
275                }
276                fstarcontainer::ControllerRequest::SetSyscallLogFilter { payload, responder } => {
277                    if let Some(process_name) = payload.process_name {
278                        system_task.kernel().add_syscall_log_filter(&process_name);
279                        let _ = responder.send(Ok(()));
280                    } else {
281                        let _ = responder.send(Err(
282                            fstarcontainer::SetSyscallLogFilterError::MissingProcessName,
283                        ));
284                    }
285                }
286                fstarcontainer::ControllerRequest::ClearSyscallLogFilters { responder } => {
287                    system_task.kernel().clear_syscall_log_filters();
288                    let _ = responder.send();
289                }
290                fstarcontainer::ControllerRequest::_UnknownMethod { .. } => (),
291            }
292            Ok(())
293        })
294        .await
295}
296
297async fn connect_to_vsock(
298    port: u32,
299    bridge_socket: fidl::Socket,
300    system_task: &CurrentTask,
301) -> Result<(), Error> {
302    let socket = loop {
303        if let Ok(socket) = system_task.kernel().default_abstract_vsock_namespace.lookup(&port) {
304            break socket;
305        };
306        fasync::Timer::new(fasync::MonotonicDuration::from_millis(100).after_now()).await;
307    };
308
309    let pipe =
310        create_fuchsia_pipe(system_task, bridge_socket, OpenFlags::RDWR | OpenFlags::NONBLOCK)?;
311    socket.downcast_socket::<VsockSocket>().unwrap().remote_connection(
312        &socket,
313        system_task,
314        pipe,
315    )?;
316
317    Ok(())
318}
319
320fn forward_to_pty(
321    kernel: &Kernel,
322    console_in: fidl::Socket,
323    console_out: fidl::Socket,
324    pty: FileHandle,
325) -> Result<(), Error> {
326    // Matches fuchsia.io.Transfer capacity, somewhat arbitrarily.
327    const BUFFER_CAPACITY: usize = 8192;
328
329    let mut rx = fuchsia_async::Socket::from_socket(console_in);
330    let mut tx = fuchsia_async::Socket::from_socket(console_out);
331    let pty_sink = pty.clone();
332    let closure = async move |current_task: &CurrentTask| {
333        let _result: Result<(), Error> = (async || {
334            let mut buffer = vec![0u8; BUFFER_CAPACITY];
335            loop {
336                let bytes = rx.read(&mut buffer[..]).await?;
337                if bytes == 0 {
338                    return Ok(());
339                }
340                pty_sink.write(current_task, &mut VecInputBuffer::new(&buffer[..bytes]))?;
341            }
342        })()
343        .await;
344    };
345    let req = SpawnRequestBuilder::new()
346        .with_debug_name("forward-to-pty-in")
347        .with_async_closure(closure)
348        .build();
349    kernel.kthreads.spawner().spawn_from_request(req);
350
351    let pty_source = pty;
352    let closure = move |current_task: &CurrentTask| {
353        let _result: Result<(), Error> =
354            fasync::LocalExecutor::default().run_singlethreaded(async {
355                let mut buffer = VecOutputBuffer::new(BUFFER_CAPACITY);
356                loop {
357                    buffer.reset();
358                    let bytes = pty_source.read(current_task, &mut buffer)?;
359                    if bytes == 0 {
360                        return Ok(());
361                    }
362                    tx.write_all(buffer.data()).await?;
363                }
364            });
365    };
366    let req = SpawnRequestBuilder::new()
367        .with_debug_name("forward-to-pty-out")
368        .with_sync_closure(closure)
369        .build();
370    kernel.kthreads.spawner().spawn_from_request(req);
371
372    Ok(())
373}
374
375pub async fn serve_container_info(
376    request_stream: fstarcontainer::InfoRequestStream,
377) -> Result<(), Error> {
378    request_stream
379        .map_err(Error::from)
380        .try_for_each_concurrent(None, |event| async {
381            match event {
382                fstarcontainer::InfoRequest::GetJobKoid { responder } => {
383                    let koid = fuchsia_runtime::job_default()
384                        .koid()
385                        .unwrap_or(zx::Koid::from_raw(zx::sys::ZX_KOID_INVALID));
386                    responder.send(&fstarcontainer::InfoGetJobKoidResponse {
387                        koid: Some(koid.raw_koid()),
388                        ..Default::default()
389                    })?;
390                }
391                fstarcontainer::InfoRequest::_UnknownMethod { ordinal, .. } => {
392                    starnix_logging::log_warn!("Unknown InfoRequest method: {}", ordinal);
393                }
394            }
395            Ok(())
396        })
397        .await
398}
399
400pub async fn serve_graphical_presenter(
401    mut request_stream: felement::GraphicalPresenterRequestStream,
402    kernel: &Kernel,
403) -> Result<(), Error> {
404    while let Some(request) = request_stream.next().await {
405        match request.context("reading graphical presenter request")? {
406            felement::GraphicalPresenterRequest::PresentView {
407                view_spec,
408                annotation_controller: _,
409                view_controller_request: _,
410                responder,
411            } => match view_spec.viewport_creation_token {
412                Some(token) => {
413                    let fb = Framebuffer::get(kernel).context("getting framebuffer from kernel")?;
414                    fb.present_view(token);
415                    let _ = responder.send(Ok(()));
416                }
417                None => {
418                    let _ = responder.send(Err(felement::PresentViewError::InvalidArgs));
419                }
420            },
421        }
422    }
423    Ok(())
424}
425
426/// Serves the memory attribution provider for the Kernel ELF component.
427pub fn serve_memory_attribution_provider_elfkernel(
428    mut request_stream: fattribution::ProviderRequestStream,
429    container: &Container,
430) -> impl Future<Output = Result<(), Error>> {
431    let observer = container.new_memory_attribution_observer(request_stream.control_handle());
432    async move {
433        while let Some(event) = request_stream.try_next().await? {
434            match event {
435                fattribution::ProviderRequest::Get { responder } => {
436                    observer.next(responder);
437                }
438                fattribution::ProviderRequest::_UnknownMethod {
439                    ordinal, control_handle, ..
440                } => {
441                    log_error!("Invalid request to AttributionProvider: {ordinal}");
442                    control_handle.shutdown_with_epitaph(zx::Status::INVALID_ARGS);
443                }
444            }
445        }
446        Ok(())
447    }
448}
449
450/// Serves the memory attribution provider for the Container component.
451pub fn serve_memory_attribution_provider_container(
452    mut request_stream: fattribution::ProviderRequestStream,
453    kernel: &Kernel,
454) -> impl Future<Output = ()> + use<> {
455    let observer = kernel.new_memory_attribution_observer(request_stream.control_handle());
456    async move {
457        while let Some(event) = request_stream
458            .try_next()
459            .await
460            .inspect_err(|err| {
461                log_warn!("Error while serving container memory attribution: {:?}", err)
462            })
463            .ok()
464            .flatten()
465        {
466            match event {
467                fattribution::ProviderRequest::Get { responder } => {
468                    observer.next(responder);
469                }
470                fattribution::ProviderRequest::_UnknownMethod {
471                    ordinal, control_handle, ..
472                } => {
473                    log_error!("Invalid request to AttributionProvider: {ordinal}");
474                    control_handle.shutdown_with_epitaph(zx::Status::INVALID_ARGS);
475                }
476            }
477        }
478    }
479}
480
481async fn select_first<O>(f1: impl Future<Output = O>, f2: impl Future<Output = O>) -> O {
482    let f1 = f1.fuse();
483    let f2 = f2.fuse();
484    pin_mut!(f1, f2);
485    select! {
486        f1 = f1 => f1,
487        f2 = f2 => f2,
488    }
489}
490
491/// Serve the LutexController protocol.
492pub async fn serve_lutex_controller(
493    request_stream: fbinder::LutexControllerRequestStream,
494    current_task: &CurrentTask,
495) -> Result<(), Error> {
496    let kernel = current_task.kernel();
497    request_stream
498        .map_err(Error::from)
499        .try_for_each_concurrent(None, |event| async move {
500            match event {
501                fbinder::LutexControllerRequest::WaitBitset { payload, responder } => {
502                    let deadline_and_receiver = (|| {
503                        let vmo = payload.vmo.ok_or_else(|| errno!(EINVAL))?;
504                        let offset = payload.offset.ok_or_else(|| errno!(EINVAL))?;
505                        let value = payload.value.ok_or_else(|| errno!(EINVAL))?;
506                        let mask = payload.mask.unwrap_or(u32::MAX);
507                        let deadline = payload.deadline.map(zx::MonotonicInstant::from_nanos);
508                        kernel
509                            .shared_futexes
510                            .external_wait(vmo.into(), offset, value, mask)
511                            .map(|(event, receiver)| (deadline, event, receiver))
512                    })();
513                    let result = match deadline_and_receiver {
514                        Ok((deadline, event, receiver)) => {
515                            // We construct a specific `wait_fut` to explicitly bind the lifecycle
516                            // of the `event` to the duration of this wait operation. If the FIDL
517                            // client disconnects (or the future is otherwise dropped/cancelled),
518                            // the `event` is immediately dropped. This zeros the strong reference
519                            // count on the `InterruptibleEvent`, signaling to the `FutexTable`
520                            // that this external waiter is now stale and should be garbage
521                            // collected on its next cleanup pass.
522                            let wait_fut = async move {
523                                let _event = event;
524                                let receiver = receiver.map_err(|_| errno!(EINTR));
525                                if let Some(deadline) = deadline {
526                                    let timer =
527                                        fasync::Timer::new(deadline).map(|_| error!(ETIMEDOUT));
528                                    select_first(timer, receiver).await
529                                } else {
530                                    receiver.await
531                                }
532                            };
533                            wait_fut.await
534                        }
535                        Err(e) => Err(e),
536                    };
537                    let result = result.map_err(|e: Errno| {
538                        fposix::Errno::from_primitive(e.code.error_code() as i32)
539                            .unwrap_or(fposix::Errno::Einval)
540                    });
541                    responder
542                        .send(result)
543                        .context("Unable to send LutexControllerRequest::WaitBitset response")?;
544                }
545                fbinder::LutexControllerRequest::WakeBitset { payload, responder } => {
546                    let result = (|| {
547                        let vmo = payload.vmo.ok_or_else(|| errno!(EINVAL))?;
548                        let offset = payload.offset.ok_or_else(|| errno!(EINVAL))?;
549                        let count = payload.count.ok_or_else(|| errno!(EINVAL))?;
550                        let mask = payload.mask.unwrap_or(u32::MAX);
551                        kernel.shared_futexes.external_wake(
552                            vmo.into(),
553                            offset,
554                            count as usize,
555                            mask,
556                        )
557                    })();
558                    let result = result
559                        .map(|count| fbinder::WakeResponse {
560                            count: Some(count as u64),
561                            ..fbinder::WakeResponse::default()
562                        })
563                        .map_err(|e: Errno| {
564                            fposix::Errno::from_primitive(e.code.error_code() as i32)
565                                .unwrap_or(fposix::Errno::Einval)
566                        });
567                    responder
568                        .send(result)
569                        .context("Unable to send LutexControllerRequest::WakeBitset response")?;
570                }
571                fbinder::LutexControllerRequest::CmpRequeue { payload, responder } => {
572                    let result = (|| {
573                        let first_vmo = payload.first_vmo.ok_or_else(|| errno!(EINVAL))?;
574                        let first_offset = payload.first_offset.ok_or_else(|| errno!(EINVAL))?;
575                        let second_vmo = payload.second_vmo;
576                        let second_offset = payload.second_offset.ok_or_else(|| errno!(EINVAL))?;
577                        let wake_count = payload.wake_count.ok_or_else(|| errno!(EINVAL))?;
578                        let requeue_count = payload.requeue_count.ok_or_else(|| errno!(EINVAL))?;
579                        let cmp_val = payload.cmp_val;
580                        kernel.shared_futexes.external_requeue(
581                            first_vmo.into(),
582                            first_offset,
583                            second_vmo.map(Into::into),
584                            second_offset,
585                            wake_count as usize,
586                            requeue_count as usize,
587                            cmp_val,
588                        )
589                    })();
590                    let result = result
591                        .map(|count| fbinder::CmpRequeueResponse {
592                            count: Some(count as u64),
593                            ..fbinder::CmpRequeueResponse::default()
594                        })
595                        .map_err(|e: Errno| {
596                            fposix::Errno::from_primitive(e.code.error_code() as i32)
597                                .unwrap_or(fposix::Errno::Einval)
598                        });
599                    responder
600                        .send(result)
601                        .context("Unable to send LutexControllerRequest::Requeue response")?;
602                }
603                fbinder::LutexControllerRequest::_UnknownMethod { ordinal, .. } => {
604                    log_warn!("Unknown LutexController ordinal: {}", ordinal);
605                }
606            }
607            Ok(())
608        })
609        .await
610        .context("failed fbinder::LutexController request")
611}
612
613#[cfg(test)]
614mod tests {
615    use super::*;
616    use assert_matches::assert_matches;
617    use fidl_fuchsia_posix as fposix;
618    use starnix_core::testing::*;
619    use zx;
620
621    #[fuchsia::test]
622    async fn lutex_controller_test() {
623        spawn_kernel_and_run(async |current_task| {
624            let (sender, receiver) = oneshot::channel::<()>();
625            current_task.kernel.kthreads.spawn_future(
626                {
627                    let kernel = current_task.kernel.clone();
628                    move || async move {
629                        let (lutex_controller, stream) = fidl::endpoints::create_proxy_and_stream::<
630                            fbinder::LutexControllerMarker,
631                        >();
632
633                        // Spawn the server
634                        let server_fut =
635                            serve_lutex_controller(stream, kernel.kthreads.system_task());
636
637                        let client_fut = async move {
638                            const VMO_SIZE: usize = 4 * 1024;
639                            let vmo = zx::Vmo::create(VMO_SIZE as u64).expect("Vmo::create");
640                            // Wait on an incorrect value.
641                            let wait = lutex_controller
642                                .wait_bitset(fbinder::WaitBitsetRequest {
643                                    vmo: Some(
644                                        vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
645                                            .expect("duplicate vmo"),
646                                    ),
647                                    offset: Some(0),
648                                    value: Some(1),
649                                    ..Default::default()
650                                })
651                                .await
652                                .expect("got_answer");
653                            assert_matches!(wait, Err(fposix::Errno::Eagain));
654
655                            // Wait with a timeout
656                            let wait = lutex_controller
657                                .wait_bitset(fbinder::WaitBitsetRequest {
658                                    vmo: Some(
659                                        vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
660                                            .expect("duplicate vmo"),
661                                    ),
662                                    offset: Some(0),
663                                    value: Some(0),
664                                    deadline: Some(0),
665                                    ..Default::default()
666                                })
667                                .await
668                                .expect("got_answer");
669                            assert_matches!(wait, Err(fposix::Errno::Etimedout));
670
671                            let mut wait = Box::pin(
672                                lutex_controller.wait_bitset(fbinder::WaitBitsetRequest {
673                                    vmo: Some(
674                                        vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
675                                            .expect("duplicate vmo"),
676                                    ),
677                                    offset: Some(0),
678                                    value: Some(0),
679                                    deadline: None,
680                                    ..Default::default()
681                                }),
682                            );
683                            // The wait is correct, the future should stay pending until a wake.
684                            assert!(futures::poll!(&mut wait).is_pending());
685
686                            let waken_up: fbinder::WakeResponse = lutex_controller
687                                .wake_bitset(fbinder::WakeBitsetRequest {
688                                    vmo: Some(
689                                        vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
690                                            .expect("duplicate vmo"),
691                                    ),
692                                    offset: Some(0),
693                                    count: Some(1),
694                                    mask: None,
695                                    ..Default::default()
696                                })
697                                .await
698                                .expect("wake_answer")
699                                .expect("wake_response");
700                            assert_eq!(waken_up.count, Some(1));
701
702                            // The wait should now return.
703                            assert!(wait.await.expect("await_answer").is_ok());
704
705                            // test cmp_requeue
706                            {
707                                let vmo2 = zx::Vmo::create(VMO_SIZE as u64).expect("Vmo::create");
708                                let mut wait1 = Box::pin(
709                                    lutex_controller.wait_bitset(fbinder::WaitBitsetRequest {
710                                        vmo: Some(
711                                            vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
712                                                .expect("duplicate vmo"),
713                                        ),
714                                        offset: Some(4),
715                                        value: Some(0),
716                                        deadline: None,
717                                        ..Default::default()
718                                    }),
719                                );
720                                let mut wait2 = Box::pin(
721                                    lutex_controller.wait_bitset(fbinder::WaitBitsetRequest {
722                                        vmo: Some(
723                                            vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
724                                                .expect("duplicate vmo"),
725                                        ),
726                                        offset: Some(4),
727                                        value: Some(0),
728                                        deadline: None,
729                                        ..Default::default()
730                                    }),
731                                );
732                                let cmp_requeue_response = lutex_controller
733                                    .cmp_requeue(fbinder::CmpRequeueRequest {
734                                        first_vmo: Some(
735                                            vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
736                                                .expect("duplicate vmo"),
737                                        ),
738                                        first_offset: Some(4),
739                                        second_vmo: Some(
740                                            vmo2.duplicate_handle(zx::Rights::SAME_RIGHTS)
741                                                .expect("duplicate vmo"),
742                                        ),
743                                        second_offset: Some(8),
744                                        wake_count: Some(1),
745                                        requeue_count: Some(1),
746                                        cmp_val: Some(0xBAD),
747                                        ..Default::default()
748                                    })
749                                    .await
750                                    .expect("got_answer")
751                                    .unwrap_err();
752                                // 0xBAD is not 0x0, so EAGAIN expected
753                                assert_matches!(cmp_requeue_response, fposix::Errno::Eagain);
754                                let cmp_requeue_response2: fbinder::CmpRequeueResponse =
755                                    lutex_controller
756                                        .cmp_requeue(fbinder::CmpRequeueRequest {
757                                            first_vmo: Some(
758                                                vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)
759                                                    .expect("duplicate vmo"),
760                                            ),
761                                            first_offset: Some(4),
762                                            second_vmo: Some(
763                                                vmo2.duplicate_handle(zx::Rights::SAME_RIGHTS)
764                                                    .expect("duplicate vmo"),
765                                            ),
766                                            second_offset: Some(8),
767                                            wake_count: Some(1),
768                                            requeue_count: Some(1),
769                                            cmp_val: Some(0x0),
770                                            ..Default::default()
771                                        })
772                                        .await
773                                        .expect("got_answer")
774                                        .expect("got_response");
775                                // 1 woken, 1 re-queued
776                                assert_matches!(cmp_requeue_response2.count, Some(2));
777                                // give wait1 and wait2 a chance to complete; we don't rely on this
778                                // being long enough however
779                                fasync::Timer::new(
780                                    fasync::MonotonicDuration::from_millis(100).after_now(),
781                                )
782                                .await;
783                                let mut still_pending_count: u32 = 0;
784                                // Up to one of wait1 and wait2 can be completed by this point, so at
785                                // least one must still be pending.
786                                let wait1_poll = futures::poll!(&mut wait1);
787                                if wait1_poll.is_pending() {
788                                    still_pending_count += 1;
789                                }
790                                let wait2_poll = futures::poll!(&mut wait2);
791                                if wait2_poll.is_pending() {
792                                    still_pending_count += 1;
793                                }
794                                assert!(still_pending_count >= 1);
795                                // at this point exactly one of wait1 or wait2 will remain pending until
796                                // we wake it via vmo2 - we don't know which of wait1 or wait2 we're
797                                // waking here (which is fine)
798                                let vmo2_wake_response: fbinder::WakeResponse = lutex_controller
799                                    .wake_bitset(fbinder::WakeBitsetRequest {
800                                        vmo: Some(
801                                            vmo2.duplicate_handle(zx::Rights::SAME_RIGHTS)
802                                                .expect("duplicate vmo"),
803                                        ),
804                                        offset: Some(8),
805                                        count: Some(1),
806                                        mask: None,
807                                        ..Default::default()
808                                    })
809                                    .await
810                                    .expect("wake_answer")
811                                    .expect("wake_response");
812                                assert_eq!(vmo2_wake_response.count, Some(1));
813                                // Now we know that both wait1 and wait2 are unblocked, so we can
814                                // wait on both of them.
815                                if wait1_poll.is_pending() {
816                                    assert!(wait1.await.expect("await_answer").is_ok());
817                                }
818                                if wait2_poll.is_pending() {
819                                    assert!(wait2.await.expect("await_answer").is_ok());
820                                }
821                            }
822                        };
823
824                        let (server_res, _) = futures::join!(server_fut, client_fut);
825                        server_res.expect("server failed");
826                        let _ = sender.send(());
827                    }
828                },
829                "lutex_controller_test",
830            );
831            receiver.await.expect("test failed");
832        })
833        .await;
834    }
835
836    #[fuchsia::test]
837    async fn container_info_test() {
838        let (info_proxy, stream) =
839            fidl::endpoints::create_proxy_and_stream::<fstarcontainer::InfoMarker>();
840
841        let server_fut = serve_container_info(stream);
842        let client_fut = async move {
843            let response = info_proxy.get_job_koid().await.expect("get_job_koid failed");
844            let expected_koid = fuchsia_runtime::job_default()
845                .koid()
846                .unwrap_or(zx::Koid::from_raw(zx::sys::ZX_KOID_INVALID));
847            assert_eq!(response.koid, Some(expected_koid.raw_koid()));
848            assert_ne!(response.koid, Some(zx::sys::ZX_KOID_INVALID));
849            drop(info_proxy);
850        };
851
852        let (server_res, ()) = futures::join!(server_fut, client_fut);
853        server_res.expect("server failed");
854    }
855}