Skip to main content

wash_runtime/plugin/
wasi_blobstore.rs

1//! # WASI Blobstore Memory Plugin
2//!
3//! This module implements an in-memory blobstore plugin for the wasmCloud runtime,
4//! providing the `wasi:blobstore@0.2.0-draft` interface for development and testing scenarios.
5
6use std::{
7    collections::{HashMap, HashSet},
8    sync::Arc,
9    time::SystemTime,
10};
11
12const WASI_BLOBSTORE_ID: &str = "wasi-blobstore";
13use tokio::sync::RwLock;
14use wasmtime::component::Resource;
15use wasmtime_wasi::{
16    InputStream, OutputStream,
17    pipe::{MemoryInputPipe, MemoryOutputPipe},
18};
19
20use crate::{
21    engine::ctx::Ctx,
22    engine::workload::{ResolvedWorkload, WorkloadComponent},
23    plugin::HostPlugin,
24    wit::{WitInterface, WitWorld},
25};
26
27mod bindings {
28    wasmtime::component::bindgen!({
29        world: "blobstore",
30        trappable_imports: true,
31        async: true,
32        with: {
33            "wasi:io": ::wasmtime_wasi::bindings::io,
34            "wasi:blobstore/container/container": String,
35            "wasi:blobstore/container/stream-object-names": crate::plugin::wasi_blobstore::StreamObjectNamesHandle,
36            "wasi:blobstore/types/incoming-value": crate::plugin::wasi_blobstore::IncomingValueHandle,
37            "wasi:blobstore/types/outgoing-value": crate::plugin::wasi_blobstore::OutgoingValueHandle,
38        },
39    });
40}
41
42use bindings::wasi::blobstore::{
43    container::Error as ContainerError,
44    types::{
45        ContainerMetadata, ContainerName, Error as BlobstoreError, ObjectId, ObjectMetadata,
46        ObjectName,
47    },
48};
49
50/// Metadata for an object stored in memory
51#[derive(Clone, Debug)]
52pub struct ObjectData {
53    pub name: String,
54    pub container: String,
55    pub data: Vec<u8>,
56    pub created_at: u64,
57}
58
59/// In-memory container representation
60#[derive(Clone, Debug)]
61pub struct ContainerData {
62    pub name: String,
63    pub created_at: u64,
64    pub objects: HashMap<String, ObjectData>,
65}
66
67/// Resource representation for an incoming value (data being read)
68pub type IncomingValueHandle = Vec<u8>;
69
70/// Resource representation for an outgoing value (data being written)
71pub struct OutgoingValueHandle {
72    pub pipe: MemoryOutputPipe,
73    pub container_name: Option<String>,
74    pub object_name: Option<String>,
75}
76
77/// Resource representation for streaming object names
78#[derive(Debug)]
79pub struct StreamObjectNamesHandle {
80    pub container_name: String,
81    pub workload_id: String,
82    pub objects: Vec<String>,
83    pub position: usize,
84}
85
86/// Memory-based blobstore plugin
87#[derive(Clone, Default)]
88pub struct WasiBlobstore {
89    /// Storage for all containers, keyed by store context ID
90    storage: Arc<RwLock<HashMap<String, HashMap<String, ContainerData>>>>,
91    /// The maximum size for objects stored in the blobstore
92    max_object_size: usize,
93}
94
95impl WasiBlobstore {
96    pub fn new(max_object_size: Option<usize>) -> Self {
97        Self {
98            storage: Arc::new(RwLock::new(HashMap::new())),
99            max_object_size: max_object_size.unwrap_or(1_000_000), // 1mb limit by default
100        }
101    }
102
103    fn get_timestamp() -> u64 {
104        SystemTime::now()
105            .duration_since(SystemTime::UNIX_EPOCH)
106            .unwrap_or_default()
107            .as_secs()
108    }
109}
110
111// Implementation for the main blobstore interface
112impl bindings::wasi::blobstore::blobstore::Host for Ctx {
113    async fn create_container(
114        &mut self,
115        name: ContainerName,
116    ) -> anyhow::Result<Result<Resource<String>, BlobstoreError>> {
117        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
118            return Ok(Err("blobstore plugin not available".to_string()));
119        };
120
121        let mut storage = plugin.storage.write().await;
122        let workload_storage = storage.entry(self.id.clone()).or_default();
123
124        if workload_storage.contains_key(&name) {
125            return Ok(Err(format!("container '{name}' already exists")));
126        }
127
128        let container_data = ContainerData {
129            name: name.clone(),
130            created_at: WasiBlobstore::get_timestamp(),
131            objects: HashMap::new(),
132        };
133
134        workload_storage.insert(name.clone(), container_data);
135        let resource = self.table.push(name)?;
136        Ok(Ok(resource))
137    }
138
139    async fn get_container(
140        &mut self,
141        name: ContainerName,
142    ) -> anyhow::Result<Result<Resource<String>, BlobstoreError>> {
143        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
144            return Ok(Err("blobstore plugin not available".to_string()));
145        };
146
147        let storage = plugin.storage.read().await;
148        let empty_map = HashMap::new();
149        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
150
151        if !workload_storage.contains_key(&name) {
152            return Ok(Err(format!("container '{name}' does not exist")));
153        }
154
155        let resource = self.table.push(name)?;
156        Ok(Ok(resource))
157    }
158
159    async fn delete_container(
160        &mut self,
161        name: ContainerName,
162    ) -> anyhow::Result<Result<(), BlobstoreError>> {
163        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
164            return Ok(Err("blobstore plugin not available".to_string()));
165        };
166
167        let mut storage = plugin.storage.write().await;
168        let workload_storage = storage.entry(self.id.clone()).or_default();
169
170        workload_storage.remove(&name);
171        Ok(Ok(()))
172    }
173
174    async fn container_exists(
175        &mut self,
176        name: ContainerName,
177    ) -> anyhow::Result<Result<bool, BlobstoreError>> {
178        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
179            return Ok(Err("blobstore plugin not available".to_string()));
180        };
181
182        let storage = plugin.storage.read().await;
183        let empty_map = HashMap::new();
184        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
185
186        Ok(Ok(workload_storage.contains_key(&name)))
187    }
188
189    async fn copy_object(
190        &mut self,
191        src: ObjectId,
192        dest: ObjectId,
193    ) -> anyhow::Result<Result<(), BlobstoreError>> {
194        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
195            return Ok(Err("blobstore plugin not available".to_string()));
196        };
197
198        let mut storage = plugin.storage.write().await;
199        let workload_storage = storage.entry(self.id.clone()).or_default();
200
201        // Get source object data (clone to avoid borrow conflicts)
202        let src_object_data = {
203            let src_container = match workload_storage.get(&src.container) {
204                Some(container) => container,
205                None => {
206                    return Ok(Err(format!(
207                        "source container '{}' does not exist",
208                        src.container
209                    )));
210                }
211            };
212
213            match src_container.objects.get(&src.object) {
214                Some(object) => object.clone(),
215                None => {
216                    return Ok(Err(format!(
217                        "source object '{}' does not exist",
218                        src.object
219                    )));
220                }
221            }
222        };
223
224        // Ensure destination container exists and copy object
225        let dest_container = match workload_storage.get_mut(&dest.container) {
226            Some(container) => container,
227            None => {
228                return Ok(Err(format!(
229                    "destination container '{}' does not exist",
230                    dest.container
231                )));
232            }
233        };
234
235        let mut copied_object = src_object_data;
236        copied_object.name = dest.object.clone();
237        copied_object.container = dest.container.clone();
238        copied_object.created_at = WasiBlobstore::get_timestamp();
239
240        dest_container.objects.insert(dest.object, copied_object);
241        Ok(Ok(()))
242    }
243
244    async fn move_object(
245        &mut self,
246        src: ObjectId,
247        dest: ObjectId,
248    ) -> anyhow::Result<Result<(), BlobstoreError>> {
249        // First copy the object
250        let _ = self.copy_object(src.clone(), dest).await?;
251
252        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
253            return Ok(Err("blobstore plugin not available".to_string()));
254        };
255
256        // Then delete the source
257        let mut storage = plugin.storage.write().await;
258        let workload_storage = storage.entry(self.id.clone()).or_default();
259
260        if let Some(src_container) = workload_storage.get_mut(&src.container) {
261            src_container.objects.remove(&src.object);
262        }
263
264        Ok(Ok(()))
265    }
266}
267
268// Resource host trait implementations - these handle the lifecycle of each resource type
269impl bindings::wasi::blobstore::container::HostContainer for Ctx {
270    async fn name(
271        &mut self,
272        container: Resource<String>,
273    ) -> anyhow::Result<Result<String, ContainerError>> {
274        let container_name = self.table.get(&container)?;
275        Ok(Ok(container_name.clone()))
276    }
277
278    async fn info(
279        &mut self,
280        container: Resource<String>,
281    ) -> anyhow::Result<Result<ContainerMetadata, ContainerError>> {
282        let container_name = self.table.get(&container)?;
283
284        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
285            return Ok(Err("blobstore plugin not available".to_string()));
286        };
287
288        let storage = plugin.storage.read().await;
289        let empty_map = HashMap::new();
290        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
291
292        match workload_storage.get(container_name) {
293            Some(container_data) => Ok(Ok(ContainerMetadata {
294                name: container_data.name.clone(),
295                created_at: container_data.created_at,
296            })),
297            None => Ok(Err(format!("container '{container_name}' does not exist"))),
298        }
299    }
300
301    async fn get_data(
302        &mut self,
303        container: Resource<String>,
304        name: ObjectName,
305        start: u64,
306        end: u64,
307    ) -> anyhow::Result<Result<Resource<IncomingValueHandle>, ContainerError>> {
308        let container_name = self.table.get(&container)?;
309
310        tracing::debug!(
311            container = container_name,
312            object = name,
313            start = start,
314            end = end,
315            workload_id = self.id,
316            "Getting object data from container"
317        );
318
319        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
320            tracing::error!("blobstore plugin not available for get_data");
321            return Ok(Err("blobstore plugin not available".to_string()));
322        };
323
324        let storage = plugin.storage.read().await;
325        let empty_map = HashMap::new();
326        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
327
328        match workload_storage.get(container_name) {
329            Some(container_data) => match container_data.objects.get(&name) {
330                Some(object_data) => {
331                    let start_idx = start.min(object_data.data.len() as u64) as usize;
332                    let end_idx = end.min(object_data.data.len() as u64) as usize;
333                    let data_slice = object_data.data[start_idx..end_idx].to_vec();
334
335                    tracing::debug!(
336                        container = container_name,
337                        object = name,
338                        original_size = object_data.data.len(),
339                        slice_size = data_slice.len(),
340                        start_idx = start_idx,
341                        end_idx = end_idx,
342                        "Retrieved object data slice"
343                    );
344
345                    let resource = self.table.push(data_slice)?;
346                    Ok(Ok(resource))
347                }
348                None => {
349                    tracing::warn!(
350                        container = container_name,
351                        object = name,
352                        "Object does not exist in container"
353                    );
354                    Ok(Err(format!("object '{name}' does not exist")))
355                }
356            },
357            None => {
358                tracing::warn!(
359                    container = container_name,
360                    workload_id = self.id,
361                    "Container does not exist for workload"
362                );
363                Ok(Err(format!("container '{container_name}' does not exist")))
364            }
365        }
366    }
367
368    async fn write_data(
369        &mut self,
370        container: Resource<String>,
371        name: ObjectName,
372        data: Resource<OutgoingValueHandle>,
373    ) -> anyhow::Result<Result<(), ContainerError>> {
374        let container_name = self.table.get(&container)?.clone();
375
376        tracing::debug!(
377            container = container_name,
378            object = name,
379            workload_id = self.id,
380            "Initiating write_data for object"
381        );
382
383        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
384            tracing::error!("blobstore plugin not available for write_data");
385            return Ok(Err("blobstore plugin not available".to_string()));
386        };
387
388        // Verify the container exists
389        let storage = plugin.storage.read().await;
390        let empty_map = HashMap::new();
391        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
392
393        if !workload_storage.contains_key(&container_name) {
394            tracing::warn!(
395                container = container_name,
396                workload_id = self.id,
397                "Container does not exist for write_data"
398            );
399            return Ok(Err(format!("container '{container_name}' does not exist")));
400        }
401        drop(storage);
402
403        // Store the container and object names - actual writing happens in finish()
404        let outgoing_handle = self.table.get_mut(&data)?;
405        outgoing_handle.container_name = Some(container_name.clone());
406        outgoing_handle.object_name = Some(name.clone());
407
408        tracing::debug!(
409            container = container_name,
410            object = name,
411            "write_data setup complete, actual write will happen in finish()"
412        );
413
414        Ok(Ok(()))
415    }
416
417    async fn list_objects(
418        &mut self,
419        container: Resource<String>,
420    ) -> anyhow::Result<Result<Resource<StreamObjectNamesHandle>, ContainerError>> {
421        let container_name = self.table.get(&container)?;
422
423        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
424            return Ok(Err("blobstore plugin not available".to_string()));
425        };
426
427        let storage = plugin.storage.read().await;
428        let empty_map = HashMap::new();
429        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
430
431        match workload_storage.get(container_name) {
432            Some(container_data) => {
433                let objects: Vec<String> = container_data.objects.keys().cloned().collect();
434                let handle = StreamObjectNamesHandle {
435                    container_name: container_name.clone(),
436                    workload_id: self.id.clone(),
437                    objects,
438                    position: 0,
439                };
440                let resource = self.table.push(handle)?;
441                Ok(Ok(resource))
442            }
443            None => Ok(Err(format!("container '{container_name}' does not exist"))),
444        }
445    }
446
447    async fn delete_object(
448        &mut self,
449        container: Resource<String>,
450        name: ObjectName,
451    ) -> anyhow::Result<Result<(), ContainerError>> {
452        let container_name = self.table.get(&container)?;
453
454        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
455            return Ok(Err("blobstore plugin not available".to_string()));
456        };
457
458        let mut storage = plugin.storage.write().await;
459        let workload_storage = storage.entry(self.id.clone()).or_default();
460
461        match workload_storage.get_mut(container_name) {
462            Some(container_data) => {
463                container_data.objects.remove(&name);
464                Ok(Ok(()))
465            }
466            None => Ok(Err(format!("container '{container_name}' does not exist"))),
467        }
468    }
469
470    async fn delete_objects(
471        &mut self,
472        container: Resource<String>,
473        names: Vec<ObjectName>,
474    ) -> anyhow::Result<Result<(), ContainerError>> {
475        let container_name = self.table.get(&container)?;
476
477        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
478            return Ok(Err("blobstore plugin not available".to_string()));
479        };
480
481        let mut storage = plugin.storage.write().await;
482        let workload_storage = storage.entry(self.id.clone()).or_default();
483
484        match workload_storage.get_mut(container_name) {
485            Some(container_data) => {
486                for name in names {
487                    container_data.objects.remove(&name);
488                }
489                Ok(Ok(()))
490            }
491            None => Ok(Err(format!("container '{container_name}' does not exist"))),
492        }
493    }
494
495    async fn has_object(
496        &mut self,
497        container: Resource<String>,
498        name: ObjectName,
499    ) -> anyhow::Result<Result<bool, ContainerError>> {
500        let container_name = self.table.get(&container)?;
501
502        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
503            return Ok(Err("blobstore plugin not available".to_string()));
504        };
505
506        let storage = plugin.storage.read().await;
507        let empty_map = HashMap::new();
508        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
509
510        match workload_storage.get(container_name) {
511            Some(container_data) => Ok(Ok(container_data.objects.contains_key(&name))),
512            None => Ok(Err(format!("container '{container_name}' does not exist"))),
513        }
514    }
515
516    async fn object_info(
517        &mut self,
518        container: Resource<String>,
519        name: ObjectName,
520    ) -> anyhow::Result<Result<ObjectMetadata, ContainerError>> {
521        let container_name = self.table.get(&container)?;
522
523        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
524            return Ok(Err("blobstore plugin not available".to_string()));
525        };
526
527        let storage = plugin.storage.read().await;
528        let empty_map = HashMap::new();
529        let workload_storage = storage.get(&self.id).unwrap_or(&empty_map);
530
531        match workload_storage.get(container_name) {
532            Some(container_data) => match container_data.objects.get(&name) {
533                Some(object_data) => Ok(Ok(ObjectMetadata {
534                    name: object_data.name.clone(),
535                    container: object_data.container.clone(),
536                    created_at: object_data.created_at,
537                    size: object_data.data.len() as u64,
538                })),
539                None => Ok(Err(format!("object '{name}' does not exist"))),
540            },
541            None => Ok(Err(format!("container '{container_name}' does not exist"))),
542        }
543    }
544
545    async fn clear(
546        &mut self,
547        container: Resource<String>,
548    ) -> anyhow::Result<Result<(), ContainerError>> {
549        let container_name = self.table.get(&container)?;
550
551        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
552            return Ok(Err("blobstore plugin not available".to_string()));
553        };
554
555        let mut storage = plugin.storage.write().await;
556        let workload_storage = storage.entry(self.id.clone()).or_default();
557
558        match workload_storage.get_mut(container_name) {
559            Some(container_data) => {
560                container_data.objects.clear();
561                Ok(Ok(()))
562            }
563            None => Ok(Err(format!("container '{container_name}' does not exist"))),
564        }
565    }
566
567    async fn drop(&mut self, rep: Resource<String>) -> anyhow::Result<()> {
568        // Container resource cleanup - resource table handles this automatically
569        tracing::debug!(
570            workload_id = self.id,
571            resource_id = ?rep,
572            "Dropping container resource"
573        );
574        self.table.delete(rep)?;
575        Ok(())
576    }
577}
578
579impl bindings::wasi::blobstore::container::HostStreamObjectNames for Ctx {
580    async fn read_stream_object_names(
581        &mut self,
582        stream: Resource<StreamObjectNamesHandle>,
583        len: u64,
584    ) -> anyhow::Result<Result<(Vec<ObjectName>, bool), ContainerError>> {
585        let stream_handle = self.table.get_mut(&stream)?;
586
587        let remaining = stream_handle
588            .objects
589            .len()
590            .saturating_sub(stream_handle.position);
591        let to_read = (len as usize).min(remaining);
592
593        let mut objects = Vec::new();
594        for i in 0..to_read {
595            if let Some(obj_name) = stream_handle.objects.get(stream_handle.position + i) {
596                objects.push(obj_name.clone());
597            }
598        }
599
600        stream_handle.position += to_read;
601        let is_end = stream_handle.position >= stream_handle.objects.len();
602
603        Ok(Ok((objects, is_end)))
604    }
605
606    async fn skip_stream_object_names(
607        &mut self,
608        stream: Resource<StreamObjectNamesHandle>,
609        num: u64,
610    ) -> anyhow::Result<Result<(u64, bool), ContainerError>> {
611        let stream_handle = self.table.get_mut(&stream)?;
612
613        let remaining = stream_handle
614            .objects
615            .len()
616            .saturating_sub(stream_handle.position);
617        let to_skip = (num as usize).min(remaining);
618
619        stream_handle.position += to_skip;
620        let is_end = stream_handle.position >= stream_handle.objects.len();
621
622        Ok(Ok((to_skip as u64, is_end)))
623    }
624
625    async fn drop(&mut self, rep: Resource<StreamObjectNamesHandle>) -> anyhow::Result<()> {
626        // StreamObjectNames resource cleanup
627        tracing::debug!(
628            workload_id = self.id,
629            resource_id = ?rep,
630            "Dropping StreamObjectNames resource"
631        );
632        self.table.delete(rep)?;
633        Ok(())
634    }
635}
636
637impl bindings::wasi::blobstore::types::HostOutgoingValue for Ctx {
638    async fn new_outgoing_value(&mut self) -> anyhow::Result<Resource<OutgoingValueHandle>> {
639        tracing::debug!(workload_id = self.id, "Creating new OutgoingValue");
640
641        let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
642            tracing::error!("blobstore plugin not available in new_outgoing_value");
643            return Err(anyhow::anyhow!("blobstore plugin not available"));
644        };
645
646        let handle = OutgoingValueHandle {
647            pipe: MemoryOutputPipe::new(plugin.max_object_size),
648            container_name: None,
649            object_name: None,
650        };
651
652        tracing::debug!(
653            workload_id = self.id,
654            pipe_capacity = plugin.max_object_size,
655            "Created OutgoingValueHandle with MemoryOutputPipe"
656        );
657
658        match self.table.push(handle) {
659            Ok(resource) => {
660                tracing::debug!(
661                    workload_id = self.id,
662                    resource_id = ?resource,
663                    "Successfully pushed OutgoingValueHandle to resource table"
664                );
665                Ok(resource)
666            }
667            Err(e) => {
668                tracing::error!(
669                    workload_id = self.id,
670                    error = ?e,
671                    "Failed to push OutgoingValueHandle to resource table in new_outgoing_value"
672                );
673                Err(e.into())
674            }
675        }
676    }
677
678    async fn outgoing_value_write_body(
679        &mut self,
680        outgoing_value: Resource<OutgoingValueHandle>,
681    ) -> anyhow::Result<Result<Resource<bindings::wasi::io0_2_1::streams::OutputStream>, ()>> {
682        tracing::debug!(workload_id = self.id, "outgoing_value_write_body called");
683
684        let handle = match self.table.get_mut(&outgoing_value) {
685            Ok(h) => {
686                tracing::debug!(
687                    workload_id = self.id,
688                    "Successfully retrieved OutgoingValueHandle from table"
689                );
690                h
691            }
692            Err(e) => {
693                tracing::error!(
694                    workload_id = self.id,
695                    error = ?e,
696                    "Failed to get OutgoingValueHandle from table"
697                );
698                return Err(e.into());
699            }
700        };
701
702        tracing::debug!(
703            workload_id = self.id,
704            "Creating boxed OutputStream from pipe"
705        );
706
707        // Return the pipe as the output stream - this is the same pipe that will be read in finish()
708        let boxed: Box<dyn OutputStream> = Box::new(handle.pipe.clone());
709
710        tracing::debug!(
711            workload_id = self.id,
712            "Attempting to push OutputStream to resource table"
713        );
714
715        match self.table.push(boxed) {
716            Ok(stream) => {
717                tracing::debug!(
718                    workload_id = self.id,
719                    stream_resource_id = ?stream,
720                    "Successfully pushed OutputStream to resource table"
721                );
722                Ok(Ok(stream))
723            }
724            Err(e) => {
725                tracing::error!(
726                    workload_id = self.id,
727                    error = ?e,
728                    error_type = std::any::type_name::<anyhow::Error>(),
729                    "Failed to push OutputStream to resource table - this is likely the TryFromIntError source"
730                );
731                Err(e.into())
732            }
733        }
734    }
735
736    async fn finish(
737        &mut self,
738        outgoing_value: Resource<OutgoingValueHandle>,
739    ) -> anyhow::Result<Result<(), BlobstoreError>> {
740        tracing::debug!(workload_id = self.id, "finish() called for OutgoingValue");
741
742        let handle = self.table.delete(outgoing_value)?;
743
744        tracing::debug!(
745            container_name = ?handle.container_name,
746            object_name = ?handle.object_name,
747            "Retrieved OutgoingValueHandle in finish()"
748        );
749
750        // If we have container and object names, perform the actual write
751        if let (Some(container_name), Some(object_name)) =
752            (&handle.container_name, &handle.object_name)
753        {
754            let Some(plugin) = self.get_plugin::<WasiBlobstore>(WASI_BLOBSTORE_ID) else {
755                tracing::error!("blobstore plugin not available in finish()");
756                return Ok(Err("blobstore plugin not available".to_string()));
757            };
758
759            // Get the data from the pipe
760            let data_bytes = handle.pipe.contents();
761
762            tracing::debug!(
763                container = container_name,
764                object = object_name,
765                pipe_data_size = data_bytes.len(),
766                workload_id = self.id,
767                "Retrieved data from pipe in finish()"
768            );
769
770            let mut storage = plugin.storage.write().await;
771            let workload_storage = storage.entry(self.id.clone()).or_default();
772
773            match workload_storage.get_mut(container_name) {
774                Some(container_data) => {
775                    let object_data = ObjectData {
776                        name: object_name.clone(),
777                        container: container_name.clone(),
778                        data: data_bytes.to_vec(),
779                        created_at: WasiBlobstore::get_timestamp(),
780                    };
781                    container_data
782                        .objects
783                        .insert(object_name.clone(), object_data);
784
785                    tracing::debug!(
786                        container = container_name,
787                        object = object_name,
788                        size = data_bytes.len(),
789                        "Stored object data to container"
790                    );
791                }
792                None => {
793                    tracing::error!(
794                        container = container_name,
795                        workload_id = self.id,
796                        "Container does not exist in finish()"
797                    );
798                    return Ok(Err(format!("container '{container_name}' does not exist")));
799                }
800            }
801        } else {
802            tracing::warn!(
803                workload_id = self.id,
804                "finish() called without container/object names set"
805            );
806        }
807
808        Ok(Ok(()))
809    }
810
811    async fn drop(&mut self, rep: Resource<OutgoingValueHandle>) -> anyhow::Result<()> {
812        tracing::debug!(
813            workload_id = self.id,
814            resource_id = ?rep,
815            "Dropping OutgoingValue resource"
816        );
817        self.table.delete(rep)?;
818        Ok(())
819    }
820}
821
822impl bindings::wasi::blobstore::types::HostIncomingValue for Ctx {
823    async fn incoming_value_consume_sync(
824        &mut self,
825        incoming_value: Resource<IncomingValueHandle>,
826    ) -> anyhow::Result<Result<Vec<u8>, BlobstoreError>> {
827        let data = self.table.delete(incoming_value)?;
828
829        tracing::debug!(
830            workload_id = self.id,
831            data_size = data.len(),
832            "incoming_value_consume_sync returning data"
833        );
834
835        Ok(Ok(data))
836    }
837
838    async fn incoming_value_consume_async(
839        &mut self,
840        incoming_value: Resource<IncomingValueHandle>,
841    ) -> anyhow::Result<
842        Result<Resource<bindings::wasi::blobstore::types::IncomingValueAsyncBody>, BlobstoreError>,
843    > {
844        let data = self.table.get(&incoming_value)?;
845
846        tracing::debug!(
847            workload_id = self.id,
848            data_size = data.len(),
849            "incoming_value_consume_async creating MemoryInputPipe with data"
850        );
851
852        let stream: Box<dyn InputStream> = Box::new(MemoryInputPipe::new(data.clone()));
853        let stream = self.table.push(stream)?;
854
855        tracing::debug!(
856            workload_id = self.id,
857            "incoming_value_consume_async created stream resource"
858        );
859
860        Ok(Ok(stream))
861    }
862
863    async fn size(&mut self, incoming_value: Resource<IncomingValueHandle>) -> anyhow::Result<u64> {
864        let data = self.table.get(&incoming_value)?;
865        Ok(data.len() as u64)
866    }
867
868    async fn drop(&mut self, rep: Resource<IncomingValueHandle>) -> anyhow::Result<()> {
869        tracing::debug!(
870            workload_id = self.id,
871            resource_id = ?rep,
872            "Dropping IncomingValue resource"
873        );
874        self.table.delete(rep)?;
875        Ok(())
876    }
877}
878
879// Note: wasi:io interface implementations are handled automatically by wasmtime-wasi
880// when setting up the Ctx during runtime initialization. The bindgen-generated
881// traits are sealed and can only be implemented on &mut _T types.
882
883// Implement the main types Host trait that combines all resource types
884impl bindings::wasi::blobstore::types::Host for Ctx {}
885
886// Implement the main container Host trait that combines all resource types
887impl bindings::wasi::blobstore::container::Host for Ctx {}
888
889#[async_trait::async_trait]
890impl HostPlugin for WasiBlobstore {
891    fn id(&self) -> &'static str {
892        WASI_BLOBSTORE_ID
893    }
894    fn world(&self) -> WitWorld {
895        WitWorld {
896            imports: HashSet::from([WitInterface::from(
897                "wasi:blobstore/blobstore,container,types@0.2.0-draft",
898            )]),
899            ..Default::default()
900        }
901    }
902
903    async fn on_component_bind(
904        &self,
905        workload_handle: &mut WorkloadComponent,
906        interfaces: std::collections::HashSet<crate::wit::WitInterface>,
907    ) -> anyhow::Result<()> {
908        // Check if any of the interfaces are wasi:blobstore related
909        let has_blobstore = interfaces
910            .iter()
911            .any(|i| i.namespace == "wasi" && i.package == "blobstore");
912
913        if !has_blobstore {
914            tracing::warn!(
915                "WasiBlobstore plugin requested for non-wasi:blobstore interface(s): {:?}",
916                interfaces
917            );
918            return Ok(());
919        }
920
921        // Add blobstore interfaces to the workload's linker
922        // Note: wasi:io interfaces are already added by wasmtime_wasi::add_to_linker_async()
923        // in the engine initialization, so we only need to add the blobstore-specific interfaces
924        tracing::debug!(
925            workload_id = workload_handle.id(),
926            "Adding blobstore interfaces to linker for workload"
927        );
928        let linker = workload_handle.linker();
929
930        bindings::wasi::blobstore::blobstore::add_to_linker(linker, |ctx| ctx)?;
931        bindings::wasi::blobstore::container::add_to_linker(linker, |ctx| ctx)?;
932        bindings::wasi::blobstore::types::add_to_linker(linker, |ctx| ctx)?;
933
934        let id = workload_handle.id();
935
936        tracing::debug!(
937            workload_id = id,
938            "Successfully added blobstore interfaces to linker for workload"
939        );
940
941        // Initialize storage for this component (note: actual storage is per-store-context, this is just a placeholder)
942        let mut storage = self.storage.write().await;
943        storage.insert(id.to_string(), HashMap::new());
944
945        tracing::debug!("WasiBlobstore plugin bound to workload '{id}'");
946
947        Ok(())
948    }
949
950    async fn on_workload_unbind(
951        &self,
952        workload_handle: &ResolvedWorkload,
953        _interfaces: std::collections::HashSet<crate::wit::WitInterface>,
954    ) -> anyhow::Result<()> {
955        let id = workload_handle.id();
956        // Clean up storage for this workload
957        let mut storage = self.storage.write().await;
958        storage.remove(id);
959
960        tracing::debug!("WasiBlobstore plugin unbound from workload '{id}'");
961
962        Ok(())
963    }
964}
965
966#[cfg(test)]
967mod tests {
968    use super::*;
969
970    #[test]
971    fn test_wasi_blobstore_creation() {
972        let blobstore = WasiBlobstore::new(None);
973        assert!(blobstore.storage.try_read().is_ok());
974    }
975
976    #[test]
977    fn test_get_timestamp() {
978        let timestamp = WasiBlobstore::get_timestamp();
979        assert!(timestamp > 0);
980    }
981
982    #[test]
983    fn test_object_data_creation() {
984        let data = ObjectData {
985            name: "test.txt".to_string(),
986            container: "test-container".to_string(),
987            data: b"hello world".to_vec(),
988            created_at: WasiBlobstore::get_timestamp(),
989        };
990
991        assert_eq!(data.name, "test.txt");
992        assert_eq!(data.container, "test-container");
993        assert_eq!(data.data, b"hello world");
994        assert!(data.created_at > 0);
995    }
996
997    #[test]
998    fn test_container_data_creation() {
999        let container = ContainerData {
1000            name: "test-container".to_string(),
1001            created_at: WasiBlobstore::get_timestamp(),
1002            objects: HashMap::new(),
1003        };
1004
1005        assert_eq!(container.name, "test-container");
1006        assert!(container.created_at > 0);
1007        assert!(container.objects.is_empty());
1008    }
1009
1010    #[tokio::test]
1011    async fn test_storage_operations() {
1012        let blobstore = WasiBlobstore::new(None);
1013
1014        // Test write access
1015        {
1016            let mut storage = blobstore.storage.write().await;
1017            storage.insert("workload1".to_string(), HashMap::new());
1018        }
1019
1020        // Test read access
1021        {
1022            let storage = blobstore.storage.read().await;
1023            assert!(storage.contains_key("workload1"));
1024        }
1025    }
1026}