1use 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(¤t_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 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 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
426pub 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
450pub 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
491pub 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 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 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 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 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 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 assert!(wait.await.expect("await_answer").is_ok());
704
705 {
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 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 assert_matches!(cmp_requeue_response2.count, Some(2));
777 fasync::Timer::new(
780 fasync::MonotonicDuration::from_millis(100).after_now(),
781 )
782 .await;
783 let mut still_pending_count: u32 = 0;
784 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 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 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}