1use 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(¤t_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 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 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
396pub 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
420pub 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
461pub 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 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 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 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 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 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 assert!(wait.await.expect("await_answer").is_ok());
674
675 {
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 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 assert_matches!(cmp_requeue_response2.count, Some(2));
747 fasync::Timer::new(
750 fasync::MonotonicDuration::from_millis(100).after_now(),
751 )
752 .await;
753 let mut still_pending_count: u32 = 0;
754 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 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 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}