ocre/bulk.rs
1//! Many rows in one D1 statement (Rails' `insert_all`, `upsert_all`, and an
2//! `update_all` with a value per row).
3//!
4//! The rows travel as one JSON array bound to a single parameter, and SQLite
5//! unpacks it with `json_each`, so a thousand rows cost one query instead of
6//! a thousand. It matters on Workers: the free plan allows 50 D1 queries
7//! per invocation, and D1 binds at most 100 parameters per statement.
8//!
9//! ```no_run
10//! use ocre::{Ctx, Result, bulk};
11//! use serde::Serialize;
12//!
13//! #[derive(Serialize)]
14//! struct Track {
15//! title: String,
16//! plays: i64,
17//! }
18//!
19//! async fn import(ctx: &Ctx, tracks: &[Track]) -> Result<usize> {
20//! let insert = bulk::insert("tracks", &["title", "plays"], tracks)?;
21//! ctx.db()?.execute(&insert.sql, insert.params).await
22//! }
23//! ```
24//!
25//! Each builder returns a [`Statement`]: run it with `db.execute`, or put
26//! several in one `db.batch` to apply them together. Rows are any serde
27//! value serializing to an object (a struct, a map, `json!({...})`); the
28//! listed columns are read from each one (a missing key is `NULL`), other
29//! keys are ignored. Values are read with `->>`, so strings, numbers and
30//! `NULL` are stored as SQLite values, booleans as `1`/`0`, and nested
31//! arrays or objects as their JSON text (for `json` columns).
32//!
33//! # Limits
34//!
35//! A D1 statement and its parameters are capped (100 KB of SQL; a bound
36//! string may be larger, but the whole request is limited too): send a few
37//! thousand small rows per statement, `rows.chunks(500)` for larger sets.
38
39use serde::Serialize;
40use serde_json::Value;
41
42use crate::{Error, IntoParam, Result, Statement};
43
44/// `INSERT INTO table (columns) SELECT ... FROM json_each(?1)`: every row in one statement.
45///
46/// # Errors
47///
48/// [`Error::Internal`] when a name is not a plain SQL identifier, `columns`
49/// is empty, or a row does not serialize to a JSON object.
50///
51/// # Examples
52///
53/// ```
54/// use ocre::{bulk, serde_json::json};
55///
56/// let rows = [json!({"title": "A", "plays": 1}), json!({"title": "B", "plays": 2})];
57/// let insert = bulk::insert("tracks", &["title", "plays"], &rows).unwrap();
58/// assert_eq!(
59/// insert.sql,
60/// "INSERT INTO tracks (title, plays) SELECT value ->> '$.title', value ->> '$.plays' FROM json_each(?1)"
61/// );
62/// assert_eq!(insert.params.len(), 1);
63/// ```
64pub fn insert<T: Serialize>(table: &str, columns: &[&str], rows: &[T]) -> Result<Statement> {
65 let (list, values) = columns_and_values(table, columns)?;
66 let sql = format!("INSERT INTO {table} ({list}) SELECT {values} FROM json_each(?1)");
67 Ok(Statement::new(sql, vec![rows_param(serde_json::to_value(rows))?]))
68}
69
70/// `INSERT ... ON CONFLICT (key) DO UPDATE SET ...`: inserts the rows, and
71/// updates the listed columns of those whose `key` already exists (Rails'
72/// `upsert_all`). `key` must have a unique index; it is one of `columns`.
73///
74/// # Errors
75///
76/// As [`insert`], and when `key` is not one of `columns`.
77///
78/// # Examples
79///
80/// ```
81/// use ocre::{bulk, serde_json::json};
82///
83/// let rows = [json!({"isrc": "FR1", "plays": 10})];
84/// let upsert = bulk::upsert("tracks", "isrc", &["isrc", "plays"], &rows).unwrap();
85/// assert!(upsert.sql.ends_with("FROM json_each(?1) WHERE true ON CONFLICT (isrc) DO UPDATE SET plays = excluded.plays"));
86/// ```
87pub fn upsert<T: Serialize>(table: &str, key: &str, columns: &[&str], rows: &[T]) -> Result<Statement> {
88 let (list, values) = columns_and_values(table, columns)?;
89 if !columns.contains(&key) {
90 return Err(Error::internal(format!("bulk::upsert: the key `{key}` must be one of the columns")));
91 }
92 let sets: Vec<String> =
93 columns.iter().filter(|column| **column != key).map(|column| format!("{column} = excluded.{column}")).collect();
94 let action = if sets.is_empty() { "NOTHING".to_owned() } else { format!("UPDATE SET {}", sets.join(", ")) };
95 // `WHERE true` lets SQLite tell the upsert clause from a join condition.
96 let sql = format!(
97 "INSERT INTO {table} ({list}) SELECT {values} FROM json_each(?1) WHERE true ON CONFLICT ({key}) DO {action}"
98 );
99 Ok(Statement::new(sql, vec![rows_param(serde_json::to_value(rows))?]))
100}
101
102/// `UPDATE table SET ... FROM json_each(?1) WHERE table.key = ...`: a value
103/// per row for many rows in one statement. Each row holds `key` (usually
104/// `id`) and the `columns` to set; `updated_at` is set to now when `touch`.
105///
106/// # Errors
107///
108/// As [`insert`].
109///
110/// # Examples
111///
112/// ```
113/// use ocre::{bulk, serde_json::json};
114///
115/// let rows = [json!({"id": 1, "plays": 11}), json!({"id": 2, "plays": 3})];
116/// let update = bulk::update("tracks", "id", &["plays"], &rows, true).unwrap();
117/// assert_eq!(
118/// update.sql,
119/// "UPDATE tracks SET plays = row.value ->> '$.plays', updated_at = datetime('now') \
120/// FROM json_each(?1) AS row WHERE tracks.id = row.value ->> '$.id'"
121/// );
122/// ```
123pub fn update<T: Serialize>(table: &str, key: &str, columns: &[&str], rows: &[T], touch: bool) -> Result<Statement> {
124 columns_and_values(table, columns)?;
125 identifier(key)?;
126 let mut sets: Vec<String> = columns.iter().map(|column| format!("{column} = row.value ->> '$.{column}'")).collect();
127 if touch {
128 sets.push("updated_at = datetime('now')".to_owned());
129 }
130 let sql = format!(
131 "UPDATE {table} SET {} FROM json_each(?1) AS row WHERE {table}.{key} = row.value ->> '$.{key}'",
132 sets.join(", ")
133 );
134 Ok(Statement::new(sql, vec![rows_param(serde_json::to_value(rows))?]))
135}
136
137/// The column list and the `value ->> '$.column'` expressions, names checked.
138fn columns_and_values(table: &str, columns: &[&str]) -> Result<(String, String)> {
139 identifier(table)?;
140 if columns.is_empty() {
141 return Err(Error::internal("bulk: no columns to write"));
142 }
143 for column in columns {
144 identifier(column)?;
145 }
146 let values: Vec<String> = columns.iter().map(|column| format!("value ->> '$.{column}'")).collect();
147 Ok((columns.join(", "), values.join(", ")))
148}
149
150/// Names are written into the SQL: only plain identifiers, never user input.
151fn identifier(name: &str) -> Result<()> {
152 let mut chars = name.chars();
153 let valid = chars.next().is_some_and(|c| c.is_ascii_alphabetic() || c == '_')
154 && chars.all(|c| c.is_ascii_alphanumeric() || c == '_');
155 if valid { Ok(()) } else { Err(Error::internal(format!("bulk: `{name}` is not a plain SQL identifier"))) }
156}
157
158/// The rows (`serde_json::to_value(rows)`) as one JSON array parameter; not
159/// generic, so each row type adds only the serialization.
160fn rows_param(rows: serde_json::Result<Value>) -> Result<crate::Param> {
161 let rows = rows.map_err(|err| Error::internal(format!("bulk: rows do not serialize: {err}")))?;
162 let objects = rows.as_array().is_some_and(|rows| rows.iter().all(Value::is_object));
163 if !objects {
164 return Err(Error::internal("bulk: each row must serialize to a JSON object (a struct or a map)"));
165 }
166 Ok(rows.to_string().into_param())
167}
168
169#[cfg(test)]
170#[path = "../tests/bulk.rs"]
171mod tests;