refactor: move execute_in_batches function to db module for better reusability
This commit is contained in:
parent
e965276754
commit
6f842cc2ce
3 changed files with 31 additions and 54 deletions
23
src/db.rs
23
src/db.rs
|
|
@ -1,7 +1,28 @@
|
||||||
use crate::error::AppError;
|
use crate::error::AppError;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use worker::{D1Database, Env};
|
use worker::{D1Database, D1PreparedStatement, Env};
|
||||||
|
|
||||||
pub fn get_db(env: &Arc<Env>) -> Result<D1Database, AppError> {
|
pub fn get_db(env: &Arc<Env>) -> Result<D1Database, AppError> {
|
||||||
env.d1("vault1").map_err(AppError::Worker)
|
env.d1("vault1").map_err(AppError::Worker)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Execute D1 statements in batches, allowing batch_size 0 to run everything at once.
|
||||||
|
pub async fn execute_in_batches(
|
||||||
|
db: &D1Database,
|
||||||
|
statements: Vec<D1PreparedStatement>,
|
||||||
|
batch_size: usize,
|
||||||
|
) -> Result<(), AppError> {
|
||||||
|
if statements.is_empty() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
if batch_size == 0 {
|
||||||
|
db.batch(statements).await?;
|
||||||
|
} else {
|
||||||
|
for chunk in statements.chunks(batch_size) {
|
||||||
|
db.batch(chunk.to_vec()).await?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -12,27 +12,6 @@ use crate::error::AppError;
|
||||||
use crate::models::cipher::{Cipher, CipherData, CipherRequestData, CreateCipherRequest};
|
use crate::models::cipher::{Cipher, CipherData, CipherRequestData, CreateCipherRequest};
|
||||||
use axum::extract::Path;
|
use axum::extract::Path;
|
||||||
|
|
||||||
/// Execute D1 statements in batches, allowing batch_size 0 to run everything at once.
|
|
||||||
async fn execute_in_batches(
|
|
||||||
db: &worker::D1Database,
|
|
||||||
statements: Vec<D1PreparedStatement>,
|
|
||||||
batch_size: usize,
|
|
||||||
) -> Result<(), AppError> {
|
|
||||||
if statements.is_empty() {
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
|
|
||||||
if batch_size == 0 {
|
|
||||||
db.batch(statements).await?;
|
|
||||||
} else {
|
|
||||||
for chunk in statements.chunks(batch_size) {
|
|
||||||
db.batch(chunk.to_vec()).await?;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
#[worker::send]
|
#[worker::send]
|
||||||
pub async fn create_cipher(
|
pub async fn create_cipher(
|
||||||
claims: Claims,
|
claims: Claims,
|
||||||
|
|
@ -221,7 +200,7 @@ pub async fn soft_delete_ciphers_bulk(
|
||||||
let db = db::get_db(&env)?;
|
let db = db::get_db(&env)?;
|
||||||
let now = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ").to_string();
|
let now = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ").to_string();
|
||||||
let batch_size = get_batch_size(&env);
|
let batch_size = get_batch_size(&env);
|
||||||
let mut statements = Vec::with_capacity(payload.ids.len());
|
let mut statements: Vec<D1PreparedStatement> = Vec::with_capacity(payload.ids.len());
|
||||||
|
|
||||||
for id in payload.ids {
|
for id in payload.ids {
|
||||||
let stmt = query!(
|
let stmt = query!(
|
||||||
|
|
@ -236,7 +215,7 @@ pub async fn soft_delete_ciphers_bulk(
|
||||||
statements.push(stmt);
|
statements.push(stmt);
|
||||||
}
|
}
|
||||||
|
|
||||||
execute_in_batches(&db, statements, batch_size).await?;
|
db::execute_in_batches(&db, statements, batch_size).await?;
|
||||||
|
|
||||||
Ok(Json(()))
|
Ok(Json(()))
|
||||||
}
|
}
|
||||||
|
|
@ -273,7 +252,7 @@ pub async fn hard_delete_ciphers_bulk(
|
||||||
) -> Result<Json<()>, AppError> {
|
) -> Result<Json<()>, AppError> {
|
||||||
let db = db::get_db(&env)?;
|
let db = db::get_db(&env)?;
|
||||||
let batch_size = get_batch_size(&env);
|
let batch_size = get_batch_size(&env);
|
||||||
let mut statements = Vec::with_capacity(payload.ids.len());
|
let mut statements: Vec<D1PreparedStatement> = Vec::with_capacity(payload.ids.len());
|
||||||
|
|
||||||
for id in payload.ids {
|
for id in payload.ids {
|
||||||
let stmt = query!(
|
let stmt = query!(
|
||||||
|
|
@ -287,7 +266,7 @@ pub async fn hard_delete_ciphers_bulk(
|
||||||
statements.push(stmt);
|
statements.push(stmt);
|
||||||
}
|
}
|
||||||
|
|
||||||
execute_in_batches(&db, statements, batch_size).await?;
|
db::execute_in_batches(&db, statements, batch_size).await?;
|
||||||
|
|
||||||
Ok(Json(()))
|
Ok(Json(()))
|
||||||
}
|
}
|
||||||
|
|
@ -350,7 +329,7 @@ pub async fn restore_ciphers_bulk(
|
||||||
let now = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ").to_string();
|
let now = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ").to_string();
|
||||||
let batch_size = get_batch_size(&env);
|
let batch_size = get_batch_size(&env);
|
||||||
let ids = payload.ids;
|
let ids = payload.ids;
|
||||||
let mut update_statements = Vec::with_capacity(ids.len());
|
let mut update_statements: Vec<D1PreparedStatement> = Vec::with_capacity(ids.len());
|
||||||
|
|
||||||
for id in ids.iter() {
|
for id in ids.iter() {
|
||||||
let stmt = query!(
|
let stmt = query!(
|
||||||
|
|
@ -365,7 +344,7 @@ pub async fn restore_ciphers_bulk(
|
||||||
update_statements.push(stmt);
|
update_statements.push(stmt);
|
||||||
}
|
}
|
||||||
|
|
||||||
execute_in_batches(&db, update_statements, batch_size).await?;
|
db::execute_in_batches(&db, update_statements, batch_size).await?;
|
||||||
|
|
||||||
let mut restored_ciphers = Vec::with_capacity(ids.len());
|
let mut restored_ciphers = Vec::with_capacity(ids.len());
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -11,30 +11,7 @@ use crate::models::cipher::{Cipher, CipherData};
|
||||||
use crate::models::folder::Folder;
|
use crate::models::folder::Folder;
|
||||||
use crate::models::import::ImportRequest;
|
use crate::models::import::ImportRequest;
|
||||||
|
|
||||||
use super::{get_batch_size};
|
use super::get_batch_size;
|
||||||
|
|
||||||
/// Execute statements in batches. If batch_size is 0, execute all in one batch.
|
|
||||||
async fn execute_in_batches(
|
|
||||||
db: &worker::D1Database,
|
|
||||||
statements: Vec<D1PreparedStatement>,
|
|
||||||
batch_size: usize,
|
|
||||||
) -> Result<(), AppError> {
|
|
||||||
if statements.is_empty() {
|
|
||||||
return Ok(());
|
|
||||||
}
|
|
||||||
|
|
||||||
if batch_size == 0 {
|
|
||||||
// Execute all statements in one batch
|
|
||||||
db.batch(statements).await?;
|
|
||||||
} else {
|
|
||||||
// Execute in chunks of batch_size
|
|
||||||
for chunk in statements.chunks(batch_size) {
|
|
||||||
db.batch(chunk.to_vec()).await?;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
#[worker::send]
|
#[worker::send]
|
||||||
pub async fn import_data(
|
pub async fn import_data(
|
||||||
|
|
@ -74,7 +51,7 @@ pub async fn import_data(
|
||||||
}
|
}
|
||||||
|
|
||||||
// Execute folder inserts in batches
|
// Execute folder inserts in batches
|
||||||
execute_in_batches(&db, folder_statements, batch_size).await?;
|
db::execute_in_batches(&db, folder_statements, batch_size).await?;
|
||||||
|
|
||||||
// Process folder relationships
|
// Process folder relationships
|
||||||
for relationship in payload.folder_relationships {
|
for relationship in payload.folder_relationships {
|
||||||
|
|
@ -146,7 +123,7 @@ pub async fn import_data(
|
||||||
}
|
}
|
||||||
|
|
||||||
// Execute cipher inserts in batches
|
// Execute cipher inserts in batches
|
||||||
execute_in_batches(&db, cipher_statements, batch_size).await?;
|
db::execute_in_batches(&db, cipher_statements, batch_size).await?;
|
||||||
|
|
||||||
Ok(Json(()))
|
Ok(Json(()))
|
||||||
}
|
}
|
||||||
Loading…
Reference in a new issue