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