1use crate::device::DeviceMode;
6use crate::device::block::canonicalize_ioctl_request;
7use crate::device::kobject::DeviceMetadata;
8use crate::fs::sysfs::{BlockDeviceInfo, build_block_device_directory};
9use crate::mm::MemoryAccessorExt;
10use crate::task::dynamic_thread_spawner::SpawnRequestBuilder;
11use crate::task::{CurrentTask, Kernel, KernelThreads};
12use crate::vfs::buffers::{InputBuffer, OutputBuffer};
13use crate::vfs::{FileObject, FileOps, FsString, NamespaceNode, SeekTarget, default_seek};
14use anyhow::{Context as _, Error};
15use block_client::{BlockClient, BufferSlice, MutableBufferSlice, RemoteBlockClient};
16use fidl::endpoints::ClientEnd;
17use fidl_fuchsia_storage_block::BlockMarker;
18use futures::channel::oneshot;
19use futures::executor::block_on;
20use starnix_sync::{LockDepMutex, RemoteBlockDeviceRegistryDevicesLock};
21use starnix_syscalls::{SUCCESS, SyscallArg, SyscallResult};
22use starnix_uapi::device_id::{BLOCK_EXTENDED_MAJOR, DeviceId};
23use starnix_uapi::errors::Errno;
24use starnix_uapi::open_flags::OpenFlags;
25use starnix_uapi::user_address::{MultiArchUserRef, UserRef};
26use starnix_uapi::{BLKGETSIZE, BLKGETSIZE64, errno, error, from_status_like_fdio, off_t};
27use std::collections::btree_map::BTreeMap;
28use std::sync::Arc;
29use std::sync::atomic::{AtomicU32, Ordering};
30
31pub struct RemoteBlockDevice {
33 block_client: Arc<SyncBlockClient>,
34}
35
36impl RemoteBlockDevice {
37 pub fn read(&self, offset: u64, buf: &mut [u8]) -> Result<(), Error> {
38 self.block_client
39 .read_at(MutableBufferSlice::Memory(buf), offset as u64)
40 .context("read_at failed")
41 }
42
43 fn new(
44 kernel: &Kernel,
45 minor: u32,
46 name: &str,
47 block: ClientEnd<BlockMarker>,
48 ) -> Result<Arc<Self>, Errno> {
49 let registry = &kernel.device_registry;
50 let device_name = FsString::from(name);
51 let virtual_block_class = registry.objects.virtual_block_class();
52 let block_client = SyncBlockClient::new(&kernel.kthreads, block)?;
53 let device = Arc::new(Self { block_client });
54 let device_weak = Arc::<RemoteBlockDevice>::downgrade(&device);
55 registry.add_device(
56 kernel,
57 device_name.as_ref(),
58 DeviceMetadata::new(
59 device_name.clone(),
60 DeviceId::new(BLOCK_EXTENDED_MAJOR, minor),
61 DeviceMode::Block,
62 )
63 .with_devtype("disk"),
66 virtual_block_class,
67 |device, dir| build_block_device_directory(device, device_weak, dir),
68 )?;
69 Ok(device)
70 }
71
72 pub fn create_file_ops(&self) -> Box<dyn FileOps> {
73 Box::new(RemoteBlockDeviceFile { block_client: self.block_client.clone() })
74 }
75}
76
77impl BlockDeviceInfo for RemoteBlockDevice {
78 fn size(&self) -> Result<usize, Errno> {
79 (self.block_client.block_count() as usize)
80 .checked_mul(self.block_client.block_size() as usize)
81 .ok_or_else(|| errno!(EINVAL))
82 }
83}
84
85pub struct SyncBlockClient {
86 client: Arc<RemoteBlockClient>,
87 terminate_tx: Option<oneshot::Sender<()>>,
88}
89
90impl SyncBlockClient {
91 fn new(kthreads: &KernelThreads, block: ClientEnd<BlockMarker>) -> Result<Arc<Self>, Errno> {
92 let (init_tx, init_rx) = std::sync::mpsc::channel();
93 let (terminate_tx, terminate_rx) = oneshot::channel();
94
95 let closure = move |_: &CurrentTask| async move {
97 let proxy = block.into_proxy();
98 match RemoteBlockClient::new(proxy).await {
99 Ok(client) => {
100 let _ = init_tx.send(Ok(Arc::new(Self {
101 client: Arc::new(client),
102 terminate_tx: Some(terminate_tx),
103 })));
104 }
105 Err(e) => {
106 let _ = init_tx.send(Err(e));
107 return;
108 }
109 }
110 let _ = terminate_rx.await;
114 };
115
116 let req = SpawnRequestBuilder::new()
117 .with_debug_name("remote-block-client")
118 .with_async_closure(closure)
119 .build();
120 kthreads.spawner().spawn_from_request(req);
121
122 match init_rx.recv() {
123 Ok(Ok(client)) => Ok(client),
124 Ok(Err(status)) => Err(from_status_like_fdio!(status)),
125 Err(_) => Err(errno!(EINVAL)),
126 }
127 }
128
129 fn block_size(&self) -> u32 {
130 self.client.block_size()
131 }
132
133 fn block_count(&self) -> u64 {
134 self.client.block_count()
135 }
136
137 fn read_at(
138 &self,
139 buffer_slice: MutableBufferSlice<'_>,
140 device_offset: u64,
141 ) -> Result<(), zx::Status> {
142 block_on(self.client.read_at(buffer_slice, device_offset))
146 }
147
148 fn write_at(
149 &self,
150 buffer_slice: BufferSlice<'_>,
151 device_offset: u64,
152 ) -> Result<(), zx::Status> {
153 block_on(self.client.write_at(buffer_slice, device_offset))
157 }
158
159 fn flush(&self) -> Result<(), zx::Status> {
160 block_on(self.client.flush())
164 }
165}
166
167impl Drop for SyncBlockClient {
168 fn drop(&mut self) {
169 if let Some(tx) = self.terminate_tx.take() {
170 let _ = tx.send(());
171 }
172 }
173}
174
175struct RemoteBlockDeviceFile {
176 block_client: Arc<SyncBlockClient>,
177}
178
179impl RemoteBlockDeviceFile {
180 fn size(&self) -> Result<usize, Errno> {
181 (self.block_client.block_count() as usize)
182 .checked_mul(self.block_client.block_size() as usize)
183 .ok_or_else(|| errno!(EINVAL))
184 }
185}
186
187impl FileOps for RemoteBlockDeviceFile {
188 fn has_persistent_offsets(&self) -> bool {
189 true
190 }
191
192 fn is_seekable(&self) -> bool {
193 true
194 }
195
196 fn seek(
199 &self,
200 _file: &FileObject,
201 _current_task: &CurrentTask,
202 current_offset: off_t,
203 target: SeekTarget,
204 ) -> Result<off_t, Errno> {
205 default_seek(current_offset, target, || self.size()?.try_into().map_err(|_| errno!(EINVAL)))
206 }
207
208 fn read(
209 &self,
210 _file: &FileObject,
211 _current_task: &CurrentTask,
212 mut offset: usize,
213 data: &mut dyn OutputBuffer,
214 ) -> Result<usize, Errno> {
215 let size = self.size()?;
216 let block_size = self.block_client.block_size() as usize;
217 const MAX_CHUNK_SIZE: usize = 32 * 1024; let mut total_read = 0;
220 while data.available() > 0 && offset < size {
221 let chunk_len = std::cmp::min(data.available(), MAX_CHUNK_SIZE);
222 let chunk_len = std::cmp::min(chunk_len, size - offset);
223 if chunk_len == 0 {
224 break;
225 }
226
227 let aligned_offset = offset - offset % block_size;
228 let end_offset = offset + chunk_len;
229 let aligned_end_offset = std::cmp::min(
230 end_offset.checked_next_multiple_of(block_size).ok_or_else(|| errno!(EINVAL))?,
231 size,
232 );
233 let aligned_data_length = aligned_end_offset - aligned_offset;
234
235 let mut read_data = vec![0u8; aligned_data_length];
236 self.block_client
237 .read_at(MutableBufferSlice::Memory(&mut read_data), aligned_offset as u64)
238 .map_err(|status| from_status_like_fdio!(status))?;
239
240 let read_offset = offset - aligned_offset;
241 let read_end = read_offset + chunk_len;
242 let bytes_written = data.write(&read_data[read_offset..read_end])?;
243
244 offset += bytes_written;
245 total_read += bytes_written;
246
247 if bytes_written < chunk_len {
248 break;
249 }
250 }
251 Ok(total_read)
252 }
253
254 fn write(
255 &self,
256 _file: &FileObject,
257 _current_task: &CurrentTask,
258 mut offset: usize,
259 data: &mut dyn InputBuffer,
260 ) -> Result<usize, Errno> {
261 let size = self.size()?;
262 let block_size = self.block_client.block_size() as usize;
263 const MAX_CHUNK_SIZE: usize = 32 * 1024; let mut total_written = 0;
266 while data.available() > 0 && offset < size {
267 let chunk_len = std::cmp::min(data.available(), MAX_CHUNK_SIZE);
268 let chunk_len = std::cmp::min(chunk_len, size - offset);
269 if chunk_len == 0 {
270 break;
271 }
272
273 let aligned_offset = offset - offset % block_size;
274 let end_offset = offset + chunk_len;
275 let aligned_end_offset = std::cmp::min(
276 end_offset.checked_next_multiple_of(block_size).ok_or_else(|| errno!(EINVAL))?,
277 size,
278 );
279 let aligned_data_length = aligned_end_offset - aligned_offset;
280
281 let mut write_buf = vec![0u8; aligned_data_length];
282
283 let head_unaligned = offset > aligned_offset;
286 let tail_unaligned = end_offset < aligned_end_offset;
287
288 if head_unaligned {
289 self.block_client
290 .read_at(
291 MutableBufferSlice::Memory(&mut write_buf[..block_size]),
292 aligned_offset as u64,
293 )
294 .map_err(|status| from_status_like_fdio!(status))?;
295 }
296
297 if tail_unaligned {
298 let last_block_start = aligned_data_length - block_size;
299 if !head_unaligned || last_block_start > 0 {
302 self.block_client
303 .read_at(
304 MutableBufferSlice::Memory(&mut write_buf[last_block_start..]),
305 (aligned_offset + last_block_start) as u64,
306 )
307 .map_err(|status| from_status_like_fdio!(status))?;
308 }
309 }
310
311 let write_offset = offset - aligned_offset;
312 let write_slice = &mut write_buf[write_offset..write_offset + chunk_len];
313 let write_slice_uninit = unsafe {
316 std::slice::from_raw_parts_mut(
317 write_slice.as_mut_ptr() as *mut std::mem::MaybeUninit<u8>,
318 write_slice.len(),
319 )
320 };
321 let bytes_read = data.read(write_slice_uninit)?;
322
323 self.block_client
324 .write_at(BufferSlice::Memory(&write_buf), aligned_offset as u64)
325 .map_err(|status| from_status_like_fdio!(status))?;
326
327 offset += bytes_read;
328 total_written += bytes_read;
329
330 if bytes_read < chunk_len {
331 break;
332 }
333 }
334 Ok(total_written)
335 }
336
337 fn sync(&self, _file: &FileObject, _current_task: &CurrentTask) -> Result<(), Errno> {
338 self.block_client.flush().map_err(|status| from_status_like_fdio!(status))
339 }
340
341 fn ioctl(
342 &self,
343 _file: &FileObject,
344 current_task: &CurrentTask,
345 request: u32,
346 arg: SyscallArg,
347 ) -> Result<SyscallResult, Errno> {
348 match canonicalize_ioctl_request(current_task, request) {
349 BLKGETSIZE => {
350 let user_size = MultiArchUserRef::<u64, u32>::new(current_task, arg);
351 let size = self.block_client.block_count();
352 current_task.write_multi_arch_object(user_size, size)?;
353 Ok(SUCCESS)
354 }
355 BLKGETSIZE64 => {
356 let user_size = UserRef::<u64>::from(arg);
357 let size = self.size()? as u64;
358 current_task.write_object(user_size, &size)?;
359 Ok(SUCCESS)
360 }
361 _ => error!(ENOTTY),
362 }
363 }
364}
365
366fn open_remote_block_device(
367 current_task: &CurrentTask,
368 id: DeviceId,
369 _node: &NamespaceNode,
370 _flags: OpenFlags,
371) -> Result<Box<dyn FileOps>, Errno> {
372 Ok(current_task.kernel().remote_block_device_registry.open(id.minor())?.create_file_ops())
373}
374
375pub fn remote_block_device_init(kernel: &Kernel) {
376 kernel
377 .device_registry
378 .register_major(
379 "remote-block".into(),
380 DeviceMode::Block,
381 BLOCK_EXTENDED_MAJOR,
382 open_remote_block_device,
383 )
384 .expect("remote block device register failed.");
385}
386
387#[derive(Default)]
388pub struct RemoteBlockDeviceRegistry {
389 devices:
390 LockDepMutex<BTreeMap<u32, Arc<RemoteBlockDevice>>, RemoteBlockDeviceRegistryDevicesLock>,
391 next_minor: AtomicU32,
392}
393
394impl RemoteBlockDeviceRegistry {
395 pub fn create_remote_block_device(
396 &self,
397 kernel: &Kernel,
398 name: &str,
399 block: ClientEnd<BlockMarker>,
400 ) -> Result<(), Error> {
401 let mut devices = self.devices.lock();
402 let minor = self.next_minor.fetch_add(1, Ordering::Relaxed);
403 let device = RemoteBlockDevice::new(kernel, minor, name, block)?;
404 devices.insert(minor, device);
405 Ok(())
406 }
407
408 pub fn open(&self, minor: u32) -> Result<Arc<RemoteBlockDevice>, Errno> {
409 self.devices.lock().get(&minor).ok_or_else(|| errno!(ENODEV)).cloned()
410 }
411}
412
413#[cfg(test)]
414mod tests {
415 use super::*;
416 use crate::testing::{anon_test_file, map_object_anywhere, spawn_kernel_and_run};
417 use crate::vfs::{SeekTarget, VecInputBuffer, VecOutputBuffer};
418 use starnix_uapi::open_flags::OpenFlags;
419 use starnix_uapi::{BLKGETSIZE, BLKGETSIZE64};
420 use vmo_backed_block_server::VmoBackedServer;
421 use zerocopy::FromBytes as _;
422
423 #[::fuchsia::test]
424 async fn test_remote_block_device_registry() {
425 spawn_kernel_and_run(async |current_task| {
426 let kernel = current_task.kernel();
427 remote_block_device_init(kernel);
428 let registry = kernel.remote_block_device_registry.clone();
429 let server =
430 VmoBackedServer::new(2, 512, &[]).expect("Failed to create VmoBackedServer");
431 let (client, server_end) = fidl::endpoints::create_endpoints::<BlockMarker>();
432 std::thread::spawn(move || {
433 let mut executor = fuchsia_async::LocalExecutor::default();
434 executor.run_singlethreaded(async move {
435 use fidl::endpoints::RequestStream;
436 server.serve(server_end.into_stream().cast_stream()).await.unwrap();
437 });
438 });
439
440 registry
441 .create_remote_block_device(kernel, "test", client)
442 .expect("create_remote_block_device failed.");
443
444 let device = registry.open(0).expect("open failed.");
445 let file = anon_test_file(¤t_task, device.create_file_ops(), OpenFlags::RDWR);
446
447 let arg_addr = map_object_anywhere(¤t_task, &0u64);
448 let mut arg = [0u8; 8];
449
450 file.ioctl(¤t_task, BLKGETSIZE64, arg_addr.into()).expect("ioctl failed");
451 current_task.read_memory_to_slice(arg_addr, &mut arg).unwrap();
452 assert_eq!(u64::read_from_bytes(&arg).unwrap(), 1024);
453
454 file.ioctl(¤t_task, BLKGETSIZE, arg_addr.into()).expect("ioctl failed");
455 current_task.read_memory_to_slice(arg_addr, &mut arg).unwrap();
456 assert_eq!(u64::read_from_bytes(&arg).unwrap(), 2);
457
458 let mut buf = VecOutputBuffer::new(256);
461 file.read(¤t_task, &mut buf).expect("read failed.");
462 assert_eq!(buf.data(), &[0u8; 256]);
463
464 let mut buf = VecInputBuffer::from(vec![1u8; 256]);
465 file.seek(¤t_task, SeekTarget::Set(0)).expect("seek failed");
466 file.write(¤t_task, &mut buf).expect("write failed.");
467
468 let mut buf = VecOutputBuffer::new(256);
469 file.seek(¤t_task, SeekTarget::Set(0)).expect("seek failed");
470 file.read(¤t_task, &mut buf).expect("read failed.");
471 assert_eq!(buf.data(), &[1u8; 256]);
472 })
473 .await;
474 }
475
476 #[::fuchsia::test]
477 async fn test_read_write_past_eof() {
478 spawn_kernel_and_run(async |current_task| {
479 let kernel = current_task.kernel();
480 remote_block_device_init(kernel);
481 let registry = kernel.remote_block_device_registry.clone();
482 let server =
483 VmoBackedServer::new(2, 512, &[]).expect("Failed to create VmoBackedServer");
484 let (client, server_end) = fidl::endpoints::create_endpoints::<BlockMarker>();
485 std::thread::spawn(move || {
486 let mut executor = fuchsia_async::LocalExecutor::default();
487 executor.run_singlethreaded(async move {
488 use fidl::endpoints::RequestStream;
489 server.serve(server_end.into_stream().cast_stream()).await.unwrap();
490 });
491 });
492
493 registry
494 .create_remote_block_device(kernel, "test", client)
495 .expect("create_remote_block_device failed.");
496
497 let device = registry.open(0).expect("open failed.");
498 let file = anon_test_file(¤t_task, device.create_file_ops(), OpenFlags::RDWR);
499
500 file.seek(¤t_task, SeekTarget::End(0)).expect("seek failed");
501 let mut buf = VecOutputBuffer::new(512);
502 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 0);
503
504 let mut buf = VecInputBuffer::from(vec![1u8; 512]);
505 assert_eq!(file.write(¤t_task, &mut buf).expect("write failed."), 0);
506 })
507 .await;
508 }
509
510 #[::fuchsia::test]
511 async fn test_unaligned_read_write_spanning_blocks() {
512 spawn_kernel_and_run(async |current_task| {
513 let kernel = current_task.kernel();
514 remote_block_device_init(kernel);
515 let registry = kernel.remote_block_device_registry.clone();
516 let server =
518 VmoBackedServer::new(3, 512, &[]).expect("Failed to create VmoBackedServer");
519 let (client, server_end) = fidl::endpoints::create_endpoints::<BlockMarker>();
520 std::thread::spawn(move || {
521 let mut executor = fuchsia_async::LocalExecutor::default();
522 executor.run_singlethreaded(async move {
523 use fidl::endpoints::RequestStream;
524 server.serve(server_end.into_stream().cast_stream()).await.unwrap();
525 });
526 });
527
528 registry
529 .create_remote_block_device(kernel, "test", client)
530 .expect("create_remote_block_device failed.");
531
532 let device = registry.open(0).expect("open failed.");
533 let file = anon_test_file(¤t_task, device.create_file_ops(), OpenFlags::RDWR);
534
535 let mut buf = VecInputBuffer::from(vec![0xAAu8; 516]);
539 file.seek(¤t_task, SeekTarget::Set(510)).expect("seek failed");
540 assert_eq!(file.write(¤t_task, &mut buf).expect("write failed."), 516);
541
542 let mut buf = VecOutputBuffer::new(516);
544 file.seek(¤t_task, SeekTarget::Set(510)).expect("seek failed");
545 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 516);
546 assert_eq!(buf.data(), &[0xAAu8; 516]);
547
548 let mut buf = VecOutputBuffer::new(1);
550 file.seek(¤t_task, SeekTarget::Set(509)).expect("seek failed");
551 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 1);
552 assert_eq!(buf.data(), &[0u8]);
553
554 let mut buf = VecOutputBuffer::new(1);
555 file.seek(¤t_task, SeekTarget::Set(1026)).expect("seek failed");
556 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 1);
557 assert_eq!(buf.data(), &[0u8]);
558 })
559 .await;
560 }
561
562 #[::fuchsia::test]
563 async fn test_exact_eof_boundary() {
564 spawn_kernel_and_run(async |current_task| {
565 let kernel = current_task.kernel();
566 remote_block_device_init(kernel);
567 let registry = kernel.remote_block_device_registry.clone();
568 let server =
570 VmoBackedServer::new(2, 512, &[]).expect("Failed to create VmoBackedServer");
571 let (client, server_end) = fidl::endpoints::create_endpoints::<BlockMarker>();
572 std::thread::spawn(move || {
573 let mut executor = fuchsia_async::LocalExecutor::default();
574 executor.run_singlethreaded(async move {
575 use fidl::endpoints::RequestStream;
576 server.serve(server_end.into_stream().cast_stream()).await.unwrap();
577 });
578 });
579
580 registry
581 .create_remote_block_device(kernel, "test", client)
582 .expect("create_remote_block_device failed.");
583
584 let device = registry.open(0).expect("open failed.");
585 let file = anon_test_file(¤t_task, device.create_file_ops(), OpenFlags::RDWR);
586
587 let mut buf = VecInputBuffer::from(vec![0xBBu8; 4]);
590 file.seek(¤t_task, SeekTarget::Set(1020)).expect("seek failed");
591 assert_eq!(file.write(¤t_task, &mut buf).expect("write failed."), 4);
592
593 let mut buf = VecOutputBuffer::new(4);
595 file.seek(¤t_task, SeekTarget::Set(1020)).expect("seek failed");
596 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 4);
597 assert_eq!(buf.data(), &[0xBBu8; 4]);
598
599 let mut buf = VecOutputBuffer::new(5);
601 file.seek(¤t_task, SeekTarget::Set(1020)).expect("seek failed");
602 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 4);
603 assert_eq!(buf.data()[..4], [0xBBu8; 4]);
604 })
605 .await;
606 }
607
608 #[::fuchsia::test]
609 async fn test_rmw_preserves_data() {
610 spawn_kernel_and_run(async |current_task| {
611 let kernel = current_task.kernel();
612 remote_block_device_init(kernel);
613 let registry = kernel.remote_block_device_registry.clone();
614 let server = VmoBackedServer::new(3, 512, &[0xFFu8; 1536])
617 .expect("Failed to create VmoBackedServer");
618 let (client, server_end) = fidl::endpoints::create_endpoints::<BlockMarker>();
619 std::thread::spawn(move || {
620 let mut executor = fuchsia_async::LocalExecutor::default();
621 executor.run_singlethreaded(async move {
622 use fidl::endpoints::RequestStream;
623 server.serve(server_end.into_stream().cast_stream()).await.unwrap();
624 });
625 });
626
627 registry
628 .create_remote_block_device(kernel, "test", client)
629 .expect("create_remote_block_device failed.");
630
631 let device = registry.open(0).expect("open failed.");
632 let file = anon_test_file(¤t_task, device.create_file_ops(), OpenFlags::RDWR);
633
634 let mut buf = VecInputBuffer::from(vec![0xAAu8; 100]);
638 file.seek(¤t_task, SeekTarget::Set(600)).expect("seek failed");
639 assert_eq!(file.write(¤t_task, &mut buf).expect("write failed."), 100);
640
641 let mut buf = VecOutputBuffer::new(512);
643 file.seek(¤t_task, SeekTarget::Set(512)).expect("seek failed");
644 assert_eq!(file.read(¤t_task, &mut buf).expect("read failed."), 512);
645 let data = buf.data();
646
647 assert_eq!(&data[0..88], &[0xFFu8; 88]);
649 assert_eq!(&data[88..188], &[0xAAu8; 100]);
651 assert_eq!(&data[188..512], &[0xFFu8; 324]);
653 })
654 .await;
655 }
656}