Skip to main content

starnix_core/device/
remote_block_device.rs

1// Copyright 2024 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::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
31/// A block device, backed by a partition hosted by Fuchsia.
32pub 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            // It is not generally true that all remote block devices are disks, but it is true for
64            // the ones we currently support. Future work should allow us to set this dynamically.
65            .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        // Spawn a thread to run the executor.
96        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            // RemoteBlockClient::new() spawns a future on this closure's executor to handle the
111            // block fifo. This keeps the executor alive and handling the fifo until the client is
112            // dropped.
113            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        // TODO(https://fxbug.dev/475530917): block_on is uninterruptible. Once there is an
143        // interruptible block_on, switch to that. For now this is okay because we expect the block
144        // to be brief in most cases.
145        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        // TODO(https://fxbug.dev/475530917): block_on is uninterruptible. Once there is an
154        // interruptible block_on, switch to that. For now this is okay because we expect the block
155        // to be brief in most cases.
156        block_on(self.client.write_at(buffer_slice, device_offset))
157    }
158
159    fn flush(&self) -> Result<(), zx::Status> {
160        // TODO(https://fxbug.dev/475530917): block_on is uninterruptible. Once there is an
161        // interruptible block_on, switch to that. For now this is okay because we expect the block
162        // to be brief in most cases.
163        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    // Manually implement seek, because default_eof_offset uses st_size (which is not used for block
197    // devices).
198    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; // 32KB
218
219        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; // 32KB
264
265        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            // Read-Modify-Write: If the write is not block-aligned at the start or end,
284            // we need to read the existing data for the first and/or last block to preserve it.
285            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 we already read the first block and it's the same as the last block, don't
300                // read again.
301                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            // SAFETY: We are writing to a buffer of u8, which is always initialized.
314            // We can safely cast &mut [u8] to &mut [MaybeUninit<u8>].
315            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(&current_task, device.create_file_ops(), OpenFlags::RDWR);
446
447            let arg_addr = map_object_anywhere(&current_task, &0u64);
448            let mut arg = [0u8; 8];
449
450            file.ioctl(&current_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(&current_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            // Deliberately read with a non-block-aligned buffer size. These reads come from
459            // uncontrolled sources so we need to be able to handle the alignment ourselves.
460            let mut buf = VecOutputBuffer::new(256);
461            file.read(&current_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(&current_task, SeekTarget::Set(0)).expect("seek failed");
466            file.write(&current_task, &mut buf).expect("write failed.");
467
468            let mut buf = VecOutputBuffer::new(256);
469            file.seek(&current_task, SeekTarget::Set(0)).expect("seek failed");
470            file.read(&current_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(&current_task, device.create_file_ops(), OpenFlags::RDWR);
499
500            file.seek(&current_task, SeekTarget::End(0)).expect("seek failed");
501            let mut buf = VecOutputBuffer::new(512);
502            assert_eq!(file.read(&current_task, &mut buf).expect("read failed."), 0);
503
504            let mut buf = VecInputBuffer::from(vec![1u8; 512]);
505            assert_eq!(file.write(&current_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            // 3 blocks of 512 bytes = 1536 bytes
517            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(&current_task, device.create_file_ops(), OpenFlags::RDWR);
534
535            // Write spanning across block boundaries (e.g., from 510 to 1026)
536            // Start at 510 (2 bytes before end of 1st block)
537            // Write 516 bytes (2 bytes in 1st block, 512 bytes in 2nd block, 2 bytes in 3rd block)
538            let mut buf = VecInputBuffer::from(vec![0xAAu8; 516]);
539            file.seek(&current_task, SeekTarget::Set(510)).expect("seek failed");
540            assert_eq!(file.write(&current_task, &mut buf).expect("write failed."), 516);
541
542            // Read back the data
543            let mut buf = VecOutputBuffer::new(516);
544            file.seek(&current_task, SeekTarget::Set(510)).expect("seek failed");
545            assert_eq!(file.read(&current_task, &mut buf).expect("read failed."), 516);
546            assert_eq!(buf.data(), &[0xAAu8; 516]);
547
548            // Verify surrounding data is still 0
549            let mut buf = VecOutputBuffer::new(1);
550            file.seek(&current_task, SeekTarget::Set(509)).expect("seek failed");
551            assert_eq!(file.read(&current_task, &mut buf).expect("read failed."), 1);
552            assert_eq!(buf.data(), &[0u8]);
553
554            let mut buf = VecOutputBuffer::new(1);
555            file.seek(&current_task, SeekTarget::Set(1026)).expect("seek failed");
556            assert_eq!(file.read(&current_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            // 2 blocks of 512 bytes = 1024 bytes
569            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(&current_task, device.create_file_ops(), OpenFlags::RDWR);
586
587            // Write ending exactly at EOF (1024)
588            // Start at 1020, write 4 bytes
589            let mut buf = VecInputBuffer::from(vec![0xBBu8; 4]);
590            file.seek(&current_task, SeekTarget::Set(1020)).expect("seek failed");
591            assert_eq!(file.write(&current_task, &mut buf).expect("write failed."), 4);
592
593            // Read back
594            let mut buf = VecOutputBuffer::new(4);
595            file.seek(&current_task, SeekTarget::Set(1020)).expect("seek failed");
596            assert_eq!(file.read(&current_task, &mut buf).expect("read failed."), 4);
597            assert_eq!(buf.data(), &[0xBBu8; 4]);
598
599            // Try to read past EOF from 1020 (request 5 bytes)
600            let mut buf = VecOutputBuffer::new(5);
601            file.seek(&current_task, SeekTarget::Set(1020)).expect("seek failed");
602            assert_eq!(file.read(&current_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            // 3 blocks of 512 bytes = 1536 bytes
615            // Initialize with a known pattern (0xFF)
616            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(&current_task, device.create_file_ops(), OpenFlags::RDWR);
633
634            // Write a small chunk in the middle of the second block (offset 600, length 100)
635            // Block 1 is 512-1024. 600 is inside.
636            // This should trigger RMW for the second block.
637            let mut buf = VecInputBuffer::from(vec![0xAAu8; 100]);
638            file.seek(&current_task, SeekTarget::Set(600)).expect("seek failed");
639            assert_eq!(file.write(&current_task, &mut buf).expect("write failed."), 100);
640
641            // Read back the entire second block to verify
642            let mut buf = VecOutputBuffer::new(512);
643            file.seek(&current_task, SeekTarget::Set(512)).expect("seek failed");
644            assert_eq!(file.read(&current_task, &mut buf).expect("read failed."), 512);
645            let data = buf.data();
646
647            // 512 to 600 (88 bytes) should be 0xFF
648            assert_eq!(&data[0..88], &[0xFFu8; 88]);
649            // 600 to 700 (100 bytes) should be 0xAA
650            assert_eq!(&data[88..188], &[0xAAu8; 100]);
651            // 700 to 1024 (324 bytes) should be 0xFF
652            assert_eq!(&data[188..512], &[0xFFu8; 324]);
653        })
654        .await;
655    }
656}