backbone-bucket 0.2.0

Bucket Bounded Context: File Storage Module for Backbone Framework
Documentation
//! Bulk Operations for FileShare aggregate
//!
//! Generated by backbone-schema. Do not edit manually.
//!
//! High-performance bulk operations with batching and error recovery.

use std::sync::Arc;
use anyhow::Result;
use tokio::time::{sleep, Duration};

use crate::domain::entity::FileShare;
use crate::domain::repositories::FileShareRepository;
use super::{BulkOperationConfig, BulkOperationProgress, BulkOperationResult, BulkOperationError};

/// Bulk operations handler for FileShare
pub struct FileShareBulkOperations<R: FileShareRepository> {
    repository: Arc<R>,
    config: BulkOperationConfig,
}

impl<R: FileShareRepository> FileShareBulkOperations<R> {
    pub fn new(repository: Arc<R>) -> Self {
        Self {
            repository,
            config: BulkOperationConfig::default(),
        }
    }

    pub fn with_config(repository: Arc<R>, config: BulkOperationConfig) -> Self {
        Self { repository, config }
    }

    // =========================================================================
    // Bulk Create
    // =========================================================================

    /// Bulk create file_share entities with batching and error recovery
    pub async fn bulk_create(&self, entities: Vec<FileShare>) -> Result<BulkOperationResult<FileShare>> {
        let total = entities.len();
        let mut progress = BulkOperationProgress::new(total, self.config.batch_size);
        let mut result = BulkOperationResult::new(progress.clone());

        // Process in batches
        for (batch_idx, batch) in entities.chunks(self.config.batch_size).enumerate() {
            progress.current_batch = batch_idx + 1;

            for (idx, entity) in batch.iter().enumerate() {
                let global_idx = batch_idx * self.config.batch_size + idx;
                let mut retries = 0;

                loop {
                    match self.repository.save(entity).await {
                        Ok(saved) => {
                            result.successful.push(saved);
                            progress.successful_items += 1;
                            break;
                        }
                        Err(e) => {
                            retries += 1;
                            if retries >= self.config.max_retries {
                                result.failed.push(BulkOperationError {
                                    index: global_idx,
                                    item_id: Some(entity.id.to_string()),
                                    error: e.to_string(),
                                    retries,
                                });
                                progress.failed_items += 1;
                                if !self.config.continue_on_error {
                                    result.progress = progress;
                                    return Ok(result);
                                }
                                break;
                            }
                            sleep(Duration::from_millis(self.config.retry_delay_ms)).await;
                        }
                    }
                }
                progress.processed_items += 1;
            }
        }

        result.progress = progress;
        Ok(result)
    }

    // =========================================================================
    // Bulk Update
    // =========================================================================

    /// Bulk update file_share entities with batching and error recovery
    pub async fn bulk_update(&self, entities: Vec<FileShare>) -> Result<BulkOperationResult<FileShare>> {
        let total = entities.len();
        let mut progress = BulkOperationProgress::new(total, self.config.batch_size);
        let mut result = BulkOperationResult::new(progress.clone());

        for (batch_idx, batch) in entities.chunks(self.config.batch_size).enumerate() {
            progress.current_batch = batch_idx + 1;

            for (idx, entity) in batch.iter().enumerate() {
                let global_idx = batch_idx * self.config.batch_size + idx;
                let id = entity.id.to_string();
                let mut retries = 0;

                loop {
                    match self.repository.update(&id, entity).await {
                        Ok(Some(updated)) => {
                            result.successful.push(updated);
                            progress.successful_items += 1;
                            break;
                        }
                        Ok(None) => {
                            result.failed.push(BulkOperationError {
                                index: global_idx,
                                item_id: Some(id.clone()),
                                error: "Entity not found".to_string(),
                                retries: 0,
                            });
                            progress.failed_items += 1;
                            break;
                        }
                        Err(e) => {
                            retries += 1;
                            if retries >= self.config.max_retries {
                                result.failed.push(BulkOperationError {
                                    index: global_idx,
                                    item_id: Some(id.clone()),
                                    error: e.to_string(),
                                    retries,
                                });
                                progress.failed_items += 1;
                                if !self.config.continue_on_error {
                                    result.progress = progress;
                                    return Ok(result);
                                }
                                break;
                            }
                            sleep(Duration::from_millis(self.config.retry_delay_ms)).await;
                        }
                    }
                }
                progress.processed_items += 1;
            }
        }

        result.progress = progress;
        Ok(result)
    }

    // =========================================================================
    // Bulk Delete
    // =========================================================================

    /// Bulk delete file_share entities by IDs with batching and error recovery
    pub async fn bulk_delete(&self, ids: Vec<String>) -> Result<BulkOperationResult<String>> {
        let total = ids.len();
        let mut progress = BulkOperationProgress::new(total, self.config.batch_size);
        let mut result = BulkOperationResult::new(progress.clone());

        for (batch_idx, batch) in ids.chunks(self.config.batch_size).enumerate() {
            progress.current_batch = batch_idx + 1;

            for (idx, id) in batch.iter().enumerate() {
                let global_idx = batch_idx * self.config.batch_size + idx;
                let mut retries = 0;

                loop {
                    match self.repository.delete(id).await {
                        Ok(true) => {
                            result.successful.push(id.clone());
                            progress.successful_items += 1;
                            break;
                        }
                        Ok(false) => {
                            result.failed.push(BulkOperationError {
                                index: global_idx,
                                item_id: Some(id.clone()),
                                error: "Entity not found".to_string(),
                                retries: 0,
                            });
                            progress.failed_items += 1;
                            break;
                        }
                        Err(e) => {
                            retries += 1;
                            if retries >= self.config.max_retries {
                                result.failed.push(BulkOperationError {
                                    index: global_idx,
                                    item_id: Some(id.clone()),
                                    error: e.to_string(),
                                    retries,
                                });
                                progress.failed_items += 1;
                                if !self.config.continue_on_error {
                                    result.progress = progress;
                                    return Ok(result);
                                }
                                break;
                            }
                            sleep(Duration::from_millis(self.config.retry_delay_ms)).await;
                        }
                    }
                }
                progress.processed_items += 1;
            }
        }

        result.progress = progress;
        Ok(result)
    }

    // =========================================================================
    // Bulk Soft Delete
    // =========================================================================

    /// Bulk soft delete file_share entities by IDs
    pub async fn bulk_soft_delete(&self, ids: Vec<String>) -> Result<BulkOperationResult<String>> {
        let total = ids.len();
        let mut progress = BulkOperationProgress::new(total, self.config.batch_size);
        let mut result = BulkOperationResult::new(progress.clone());

        for (batch_idx, batch) in ids.chunks(self.config.batch_size).enumerate() {
            progress.current_batch = batch_idx + 1;

            for (idx, id) in batch.iter().enumerate() {
                let global_idx = batch_idx * self.config.batch_size + idx;

                match self.repository.soft_delete(id).await {
                    Ok(true) => {
                        result.successful.push(id.clone());
                        progress.successful_items += 1;
                    }
                    Ok(false) => {
                        result.failed.push(BulkOperationError::new(
                            global_idx,
                            Some(id.clone()),
                            "Entity not found".to_string(),
                        ));
                        progress.failed_items += 1;
                    }
                    Err(e) => {
                        result.failed.push(BulkOperationError::new(
                            global_idx,
                            Some(id.clone()),
                            e.to_string(),
                        ));
                        progress.failed_items += 1;
                        if !self.config.continue_on_error {
                            result.progress = progress;
                            return Ok(result);
                        }
                    }
                }
                progress.processed_items += 1;
            }
        }

        result.progress = progress;
        Ok(result)
    }

    /// Bulk restore soft-deleted file_share entities by IDs
    pub async fn bulk_restore(&self, ids: Vec<String>) -> Result<BulkOperationResult<FileShare>> {
        let total = ids.len();
        let mut progress = BulkOperationProgress::new(total, self.config.batch_size);
        let mut result = BulkOperationResult::new(progress.clone());

        for (batch_idx, batch) in ids.chunks(self.config.batch_size).enumerate() {
            progress.current_batch = batch_idx + 1;

            for (idx, id) in batch.iter().enumerate() {
                let global_idx = batch_idx * self.config.batch_size + idx;

                match self.repository.restore(id).await {
                    Ok(Some(restored)) => {
                        result.successful.push(restored);
                        progress.successful_items += 1;
                    }
                    Ok(None) => {
                        result.failed.push(BulkOperationError::new(
                            global_idx,
                            Some(id.clone()),
                            "Entity not found in trash".to_string(),
                        ));
                        progress.failed_items += 1;
                    }
                    Err(e) => {
                        result.failed.push(BulkOperationError::new(
                            global_idx,
                            Some(id.clone()),
                            e.to_string(),
                        ));
                        progress.failed_items += 1;
                        if !self.config.continue_on_error {
                            result.progress = progress;
                            return Ok(result);
                        }
                    }
                }
                progress.processed_items += 1;
            }
        }

        result.progress = progress;
        Ok(result)
    }

}