1use 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#[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#[derive(Clone, Debug)]
61pub struct ContainerData {
62 pub name: String,
63 pub created_at: u64,
64 pub objects: HashMap<String, ObjectData>,
65}
66
67pub type IncomingValueHandle = Vec<u8>;
69
70pub struct OutgoingValueHandle {
72 pub pipe: MemoryOutputPipe,
73 pub container_name: Option<String>,
74 pub object_name: Option<String>,
75}
76
77#[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#[derive(Clone, Default)]
88pub struct WasiBlobstore {
89 storage: Arc<RwLock<HashMap<String, HashMap<String, ContainerData>>>>,
91 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), }
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
111impl 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 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 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 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 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
268impl 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 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 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 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 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 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 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 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
879impl bindings::wasi::blobstore::types::Host for Ctx {}
885
886impl 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 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 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 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 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 {
1016 let mut storage = blobstore.storage.write().await;
1017 storage.insert("workload1".to_string(), HashMap::new());
1018 }
1019
1020 {
1022 let storage = blobstore.storage.read().await;
1023 assert!(storage.contains_key("workload1"));
1024 }
1025 }
1026}