1use crate::config::{MysqlColumnMapping, MysqlSinkConfig};
4use async_trait::async_trait;
5use faucet_core::{FaucetError, SchemaEvolution, SqlBaseType, json_schema_base_type};
6use serde_json::Value;
7use sqlx::mysql::MySqlPoolOptions;
8use sqlx::{MySqlConnection, MySqlPool, Row};
9
10pub struct MysqlSink {
12 config: MysqlSinkConfig,
13 pool: MySqlPool,
14}
15
16fn quote_ident_mysql(name: &str) -> String {
21 format!("`{}`", name.replace('`', "``"))
22}
23
24fn mysql_keyword(t: SqlBaseType) -> &'static str {
30 match t {
31 SqlBaseType::Integer => "BIGINT",
32 SqlBaseType::Double => "DOUBLE",
33 SqlBaseType::Boolean => "TINYINT(1)",
34 SqlBaseType::Text => "LONGTEXT",
35 SqlBaseType::Json => "JSON",
36 }
37}
38
39fn build_add_column_sql(table: &str, col: &str, t: SqlBaseType) -> String {
46 format!(
47 "ALTER TABLE {table} ADD COLUMN {} {}",
48 quote_ident_mysql(col),
49 mysql_keyword(t)
50 )
51}
52
53fn build_modify_column_sql(table: &str, col: &str, t: SqlBaseType) -> String {
57 format!(
58 "ALTER TABLE {table} MODIFY COLUMN {} {}",
59 quote_ident_mysql(col),
60 mysql_keyword(t)
61 )
62}
63
64fn mysql_data_type_to_json_schema(data_type: &str, nullable: bool) -> Value {
73 let base = match data_type {
74 "bigint" | "int" | "integer" | "smallint" | "mediumint" | "tinyint" => "integer",
75 "double" | "float" | "decimal" | "numeric" => "number",
76 "json" => "object",
77 _ => "string",
78 };
79 if nullable {
80 serde_json::json!({ "type": [base, "null"] })
81 } else {
82 serde_json::json!({ "type": base })
83 }
84}
85
86fn key_matches_unique_index(
113 unique_indexes: &[std::collections::BTreeSet<String>],
114 key: &[String],
115) -> bool {
116 if key.is_empty() {
117 return false;
118 }
119 let key_set: std::collections::BTreeSet<String> = key.iter().cloned().collect();
120 unique_indexes.contains(&key_set)
121}
122
123fn on_duplicate_clause(key: &[String], all_cols: &[String]) -> String {
124 let updates: Vec<String> = all_cols
125 .iter()
126 .filter(|c| !key.iter().any(|k| k == *c))
127 .map(|c| {
128 let q = quote_ident_mysql(c);
129 format!("{q} = VALUES({q})")
130 })
131 .collect();
132 if updates.is_empty() {
133 let q = quote_ident_mysql(&key[0]);
134 format!("ON DUPLICATE KEY UPDATE {q} = {q}")
135 } else {
136 format!("ON DUPLICATE KEY UPDATE {}", updates.join(", "))
137 }
138}
139
140impl MysqlSink {
141 pub async fn new(config: MysqlSinkConfig) -> Result<Self, FaucetError> {
143 config.write.validate()?;
144 if !matches!(config.write.write_mode, faucet_core::WriteMode::Append)
145 && !matches!(config.column_mapping, MysqlColumnMapping::AutoMap)
146 {
147 return Err(FaucetError::Config(
148 "mysql sink: write_mode upsert/delete requires column_mapping: auto_map \
149 (key columns must be real columns, not inside a JSON blob)"
150 .into(),
151 ));
152 }
153
154 let pool = MySqlPoolOptions::new()
155 .max_connections(config.max_connections)
156 .connect(&config.connection_url)
157 .await
158 .map_err(|e| FaucetError::Sink(format!("MySQL connection failed: {e}")))?;
159
160 let sink = Self { config, pool };
161
162 if !matches!(sink.config.write.write_mode, faucet_core::WriteMode::Append) {
167 sink.assert_key_is_unique_index().await?;
168 }
169
170 Ok(sink)
171 }
172
173 async fn read_unique_indexes(
184 &self,
185 ) -> Result<Vec<std::collections::BTreeSet<String>>, FaucetError> {
186 let rows = sqlx::query(
189 "SELECT CAST(INDEX_NAME AS CHAR) AS INDEX_NAME, \
190 CAST(COLUMN_NAME AS CHAR) AS COLUMN_NAME \
191 FROM INFORMATION_SCHEMA.STATISTICS \
192 WHERE TABLE_NAME = ? AND TABLE_SCHEMA = DATABASE() AND NON_UNIQUE = 0 \
193 ORDER BY INDEX_NAME, SEQ_IN_INDEX",
194 )
195 .bind(&self.config.table_name)
196 .fetch_all(&self.pool)
197 .await
198 .map_err(|e| FaucetError::Sink(format!("failed to query table indexes: {e}")))?;
199
200 let mut by_index: std::collections::BTreeMap<String, std::collections::BTreeSet<String>> =
201 std::collections::BTreeMap::new();
202 for row in &rows {
203 let index_name: String = row.get("INDEX_NAME");
204 let column_name: String = row.get("COLUMN_NAME");
205 by_index.entry(index_name).or_default().insert(column_name);
206 }
207 Ok(by_index.into_values().collect())
208 }
209
210 async fn assert_key_is_unique_index(&self) -> Result<(), FaucetError> {
220 let unique_indexes = self.read_unique_indexes().await?;
221 if unique_indexes.is_empty() {
222 tracing::warn!(
223 table = %self.config.table_name,
224 "mysql sink: no PRIMARY/UNIQUE index found on target table (it may not exist \
225 yet); skipping upsert key validation — the first write will surface a \
226 missing-table or missing-constraint error"
227 );
228 return Ok(());
229 }
230 if !key_matches_unique_index(&unique_indexes, &self.config.write.key) {
231 let available: Vec<String> = unique_indexes
232 .iter()
233 .map(|idx| {
234 let mut cols: Vec<&str> = idx.iter().map(String::as_str).collect();
235 cols.sort_unstable();
236 format!("({})", cols.join(", "))
237 })
238 .collect();
239 return Err(FaucetError::Config(format!(
240 "mysql sink: write_mode {} requires `key` {:?} to exactly match a PRIMARY KEY or \
241 UNIQUE index on table '{}' — MySQL's `ON DUPLICATE KEY UPDATE` resolves on the \
242 table's real unique indexes, so an unmatched key would silently upsert on the \
243 wrong index. Existing unique indexes: {}",
244 self.config.write.write_mode.as_str(),
245 self.config.write.key,
246 self.config.table_name,
247 available.join(", "),
248 )));
249 }
250 Ok(())
251 }
252
253 async fn insert_json(
259 &self,
260 conn: &mut MySqlConnection,
261 records: &[Value],
262 column: &str,
263 ) -> Result<usize, FaucetError> {
264 if records.is_empty() {
265 return Ok(0);
266 }
267
268 let placeholders: Vec<&str> = records.iter().map(|_| "(?)").collect();
270 let insert_sql = format!(
271 "INSERT INTO {} ({}) VALUES {}",
272 quote_ident_mysql(&self.config.table_name),
273 quote_ident_mysql(column),
274 placeholders.join(", ")
275 );
276
277 let mut q = sqlx::query(&insert_sql);
278 for record in records {
279 let json_str = serde_json::to_string(record)
280 .map_err(|e| FaucetError::Sink(format!("failed to serialize record: {e}")))?;
281 q = q.bind(json_str);
282 }
283
284 q.execute(&mut *conn)
285 .await
286 .map_err(|e| FaucetError::Sink(format!("MySQL insert failed: {e}")))?;
287
288 Ok(records.len())
289 }
290
291 async fn insert_auto_map_with_conflict(
305 &self,
306 conn: &mut MySqlConnection,
307 records: &[Value],
308 conflict_key: Option<&[String]>,
309 ) -> Result<usize, FaucetError> {
310 if records.is_empty() {
311 return Ok(0);
312 }
313
314 let columns: Vec<String> = sqlx::query(
316 "SELECT COLUMN_NAME FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME = ? AND TABLE_SCHEMA = DATABASE() ORDER BY ORDINAL_POSITION"
317 )
318 .bind(&self.config.table_name)
319 .fetch_all(&mut *conn)
320 .await
321 .map_err(|e| FaucetError::Sink(format!("failed to query table columns: {e}")))?
322 .iter()
323 .map(|row| row.get::<String, _>("COLUMN_NAME"))
324 .collect();
325
326 if columns.is_empty() {
327 return Err(FaucetError::Sink(format!(
328 "table '{}' has no columns or does not exist",
329 self.config.table_name
330 )));
331 }
332
333 let mut matched_rows: Vec<Vec<(&String, &Value)>> = Vec::with_capacity(records.len());
340 let mut used: std::collections::HashSet<&str> = std::collections::HashSet::new();
341
342 for record in records {
343 let obj = record
344 .as_object()
345 .ok_or_else(|| FaucetError::Sink("AutoMap requires JSON object records".into()))?;
346
347 let matching: Vec<(&String, &Value)> = columns
348 .iter()
349 .filter_map(|col| obj.get(col).map(|v| (col, v)))
350 .collect();
351
352 if matching.is_empty() {
353 tracing::warn!(
354 record_keys = ?obj.keys().collect::<Vec<_>>(),
355 table_columns = ?columns,
356 "record has no keys matching table columns, skipping"
357 );
358 continue;
359 }
360
361 for (c, _) in &matching {
362 used.insert(c.as_str());
363 }
364 matched_rows.push(matching);
365 }
366
367 if matched_rows.is_empty() {
368 return Ok(0);
369 }
370
371 let insert_columns: Vec<String> = columns
373 .iter()
374 .filter(|c| used.contains(c.as_str()))
375 .cloned()
376 .collect();
377
378 let num_cols = insert_columns.len();
379 let num_rows = matched_rows.len();
380 let col_names: Vec<String> = insert_columns
381 .iter()
382 .map(|c| quote_ident_mysql(c))
383 .collect();
384
385 const MAX_MYSQL_PARAMS: usize = 65535;
391 let max_rows_per_insert = (MAX_MYSQL_PARAMS / num_cols).max(1);
392
393 for sub in matched_rows.chunks(max_rows_per_insert) {
394 let row_placeholder = format!("({})", vec!["?"; num_cols].join(", "));
396 let value_tuples: Vec<&str> =
397 (0..sub.len()).map(|_| row_placeholder.as_str()).collect();
398 let base_query = format!(
399 "INSERT INTO {} ({}) VALUES {}",
400 quote_ident_mysql(&self.config.table_name),
401 col_names.join(", "),
402 value_tuples.join(", ")
403 );
404 let query = match conflict_key {
405 Some(key) => format!("{base_query} {}", on_duplicate_clause(key, &insert_columns)),
406 None => base_query,
407 };
408
409 let mut q = sqlx::query(&query);
410 for matched in sub {
411 for col in &insert_columns {
412 let val = matched.iter().find(|(c, _)| *c == col).map(|(_, v)| *v);
413 q = match val {
418 None | Some(Value::Null) => q.bind(None::<String>),
419 Some(Value::Bool(b)) => q.bind(*b),
420 Some(Value::Number(n)) => {
421 if let Some(i) = n.as_i64() {
422 q.bind(i)
423 } else if let Some(f) = n.as_f64() {
424 q.bind(f)
425 } else {
426 q.bind(n.to_string())
428 }
429 }
430 Some(Value::String(s)) => q.bind(s.clone()),
431 Some(v) => q.bind(v.to_string()),
434 };
435 }
436 }
437
438 q.execute(&mut *conn)
439 .await
440 .map_err(|e| FaucetError::Sink(format!("MySQL insert failed: {e}")))?;
441 }
442
443 Ok(num_rows)
444 }
445
446 async fn insert_auto_map(
452 &self,
453 conn: &mut MySqlConnection,
454 records: &[Value],
455 ) -> Result<usize, FaucetError> {
456 self.insert_auto_map_with_conflict(conn, records, None)
457 .await
458 }
459
460 async fn delete_by_keys(
464 &self,
465 conn: &mut MySqlConnection,
466 deletes: &[faucet_core::KeyTuple],
467 ) -> Result<usize, FaucetError> {
468 if deletes.is_empty() {
469 return Ok(0);
470 }
471 let key = &self.config.write.key;
472 let table_ref = quote_ident_mysql(&self.config.table_name);
473 let col_list = key
474 .iter()
475 .map(|k| quote_ident_mysql(k))
476 .collect::<Vec<_>>()
477 .join(", ");
478
479 const MAX_MYSQL_PARAMS: usize = 65535;
480 let per = (MAX_MYSQL_PARAMS / key.len().max(1)).max(1);
481 let mut total = 0usize;
482
483 for chunk in deletes.chunks(per) {
484 let tuples: Vec<String> = chunk
485 .iter()
486 .map(|_| format!("({})", vec!["?"; key.len()].join(", ")))
487 .collect();
488 let sql = format!(
489 "DELETE FROM {table_ref} WHERE ({col_list}) IN ({})",
490 tuples.join(", ")
491 );
492 let mut q = sqlx::query(&sql);
493 for kt in chunk {
494 for (_, v) in &kt.0 {
495 q = match v {
497 Value::Null => q.bind(None::<String>),
498 Value::Bool(b) => q.bind(*b),
499 Value::Number(n) => {
500 if let Some(i) = n.as_i64() {
501 q.bind(i)
502 } else if let Some(f) = n.as_f64() {
503 q.bind(f)
504 } else {
505 q.bind(n.to_string())
506 }
507 }
508 Value::String(s) => q.bind(s.clone()),
509 other => q.bind(other.to_string()),
510 };
511 }
512 }
513 let res = q
514 .execute(&mut *conn)
515 .await
516 .map_err(|e| FaucetError::Sink(format!("MySQL delete failed: {e}")))?;
517 total += res.rows_affected() as usize;
518 }
519 Ok(total)
520 }
521
522 async fn apply_plan(&self, plan: &faucet_core::WritePlan) -> Result<usize, FaucetError> {
526 let mut tx = self
527 .pool
528 .begin()
529 .await
530 .map_err(|e| FaucetError::Sink(format!("MySQL transaction begin failed: {e}")))?;
531
532 let mut affected = 0usize;
533 if !plan.upserts.is_empty() {
534 affected += self
535 .insert_auto_map_with_conflict(&mut tx, &plan.upserts, Some(&self.config.write.key))
536 .await?;
537 }
538 if !plan.deletes.is_empty() {
539 affected += self.delete_by_keys(&mut tx, &plan.deletes).await?;
540 }
541
542 tx.commit()
543 .await
544 .map_err(|e| FaucetError::Sink(format!("MySQL transaction commit failed: {e}")))?;
545 Ok(affected)
546 }
547
548 async fn read_columns(&self) -> Result<Vec<(String, String, bool)>, FaucetError> {
557 let rows = sqlx::query(
561 "SELECT CAST(COLUMN_NAME AS CHAR) AS COLUMN_NAME, \
562 CAST(DATA_TYPE AS CHAR) AS DATA_TYPE, \
563 CAST(IS_NULLABLE AS CHAR) AS IS_NULLABLE \
564 FROM INFORMATION_SCHEMA.COLUMNS \
565 WHERE TABLE_NAME = ? AND TABLE_SCHEMA = DATABASE() ORDER BY ORDINAL_POSITION",
566 )
567 .bind(&self.config.table_name)
568 .fetch_all(&self.pool)
569 .await
570 .map_err(|e| FaucetError::Sink(format!("failed to query table columns: {e}")))?;
571
572 Ok(rows
573 .iter()
574 .map(|row| {
575 (
576 row.get::<String, _>("COLUMN_NAME"),
577 row.get::<String, _>("DATA_TYPE").to_ascii_lowercase(),
580 row.get::<String, _>("IS_NULLABLE")
581 .eq_ignore_ascii_case("YES"),
582 )
583 })
584 .collect())
585 }
586
587 async fn ensure_commit_table(&self) -> Result<(), FaucetError> {
596 let sql = format!(
597 "CREATE TABLE IF NOT EXISTS {t} ({s} VARCHAR(255) PRIMARY KEY, {k} TEXT NOT NULL, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)",
598 t = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TABLE),
599 s = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_SCOPE_COL),
600 k = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TOKEN_COL),
601 );
602 sqlx::query(&sql)
603 .execute(&self.pool)
604 .await
605 .map_err(|e| FaucetError::Sink(format!("MySQL commit-table create failed: {e}")))?;
606 Ok(())
607 }
608}
609
610#[async_trait]
611impl faucet_core::Sink for MysqlSink {
612 fn config_schema(&self) -> serde_json::Value {
613 serde_json::to_value(faucet_core::schema_for!(MysqlSinkConfig))
614 .expect("schema serialization")
615 }
616
617 fn dataset_uri(&self) -> String {
618 format!(
619 "{}?table={}",
620 faucet_core::redact_uri_credentials(&self.config.connection_url),
621 self.config.table_name
622 )
623 }
624
625 async fn check(
631 &self,
632 ctx: &faucet_core::check::CheckContext,
633 ) -> Result<faucet_core::check::CheckReport, FaucetError> {
634 use faucet_core::check::{CheckReport, Probe};
635
636 let started = std::time::Instant::now();
637 let probe =
638 match tokio::time::timeout(ctx.timeout, sqlx::query("SELECT 1").execute(&self.pool))
639 .await
640 {
641 Ok(Ok(_)) => Probe::pass("auth", started.elapsed()),
642 Ok(Err(e)) => Probe::fail_hint(
643 "auth",
644 started.elapsed(),
645 e.to_string(),
646 "check connection_url / credentials / that the database is reachable",
647 ),
648 Err(_) => Probe::fail_hint(
649 "auth",
650 started.elapsed(),
651 "timed out",
652 "check connection_url / credentials / that the database is reachable",
653 ),
654 };
655 Ok(CheckReport::single(probe))
656 }
657
658 fn supported_write_modes(&self) -> &'static [faucet_core::WriteMode] {
659 &[
660 faucet_core::WriteMode::Append,
661 faucet_core::WriteMode::Upsert,
662 faucet_core::WriteMode::Delete,
663 ]
664 }
665
666 fn dedups_by_key(&self) -> bool {
667 self.config.write.dedups_by_key()
668 }
669
670 fn supports_schema_evolution(&self) -> bool {
671 true
672 }
673
674 async fn current_schema(&self) -> Result<Option<serde_json::Value>, FaucetError> {
682 let columns = self.read_columns().await?;
683 if columns.is_empty() {
684 return Ok(None); }
686
687 let mut props = serde_json::Map::new();
688 for (name, data_type, nullable) in columns {
689 props.insert(name, mysql_data_type_to_json_schema(&data_type, nullable));
690 }
691 Ok(Some(
692 serde_json::json!({ "type": "object", "properties": props }),
693 ))
694 }
695
696 async fn evolve_schema(&self, evolution: &SchemaEvolution) -> Result<(), FaucetError> {
706 let table_ref = quote_ident_mysql(&self.config.table_name);
707
708 let current = self.read_columns().await?;
712 let existing: std::collections::HashSet<&str> =
713 current.iter().map(|(n, _, _)| n.as_str()).collect();
714
715 let mut conn = self
716 .pool
717 .acquire()
718 .await
719 .map_err(|e| FaucetError::Sink(format!("MySQL evolve acquire failed: {e}")))?;
720
721 for c in &evolution.additions {
722 if existing.contains(c.name.as_str()) {
725 continue;
726 }
727 let t = json_schema_base_type(&c.to).unwrap_or(SqlBaseType::Text);
728 sqlx::query(&build_add_column_sql(&table_ref, &c.name, t))
729 .execute(&mut *conn)
730 .await
731 .map_err(|e| {
732 FaucetError::Sink(format!("MySQL ADD COLUMN {} failed: {e}", c.name))
733 })?;
734 }
735
736 for c in &evolution.widenings {
737 let t = json_schema_base_type(&c.to).unwrap_or(SqlBaseType::Text);
738 sqlx::query(&build_modify_column_sql(&table_ref, &c.name, t))
739 .execute(&mut *conn)
740 .await
741 .map_err(|e| {
742 FaucetError::Sink(format!("MySQL MODIFY COLUMN {} failed: {e}", c.name))
743 })?;
744 }
745
746 for col in &evolution.relax_nullability {
747 let existing_type = current
751 .iter()
752 .find(|(n, _, _)| n == col)
753 .map(|(_, dt, nullable)| {
754 let fragment = mysql_data_type_to_json_schema(dt, *nullable);
755 json_schema_base_type(&fragment).unwrap_or(SqlBaseType::Text)
756 })
757 .unwrap_or(SqlBaseType::Text);
758 let sql = format!(
759 "ALTER TABLE {table_ref} MODIFY COLUMN {} {} NULL",
760 quote_ident_mysql(col),
761 mysql_keyword(existing_type)
762 );
763 sqlx::query(&sql)
764 .execute(&mut *conn)
765 .await
766 .map_err(|e| FaucetError::Sink(format!("MySQL DROP NOT NULL {col} failed: {e}")))?;
767 }
768
769 Ok(())
770 }
771
772 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
785 if records.is_empty() {
786 return Ok(0);
787 }
788
789 if !matches!(self.config.write.write_mode, faucet_core::WriteMode::Append) {
791 let plan = faucet_core::plan_writes(records, &self.config.write);
792 if let Some((idx, msg)) = plan.failed.first() {
793 return Err(FaucetError::Sink(format!(
794 "mysql {}: row {idx}: {msg}",
795 self.config.write.write_mode.as_str()
796 )));
797 }
798 return self.apply_plan(&plan).await;
799 }
800
801 let mut conn = self
802 .pool
803 .acquire()
804 .await
805 .map_err(|e| FaucetError::Sink(format!("MySQL pool acquire failed: {e}")))?;
806
807 let chunks: Vec<&[Value]> = if self.config.batch_size == 0 {
808 vec![records]
812 } else {
813 records.chunks(self.config.batch_size).collect()
814 };
815
816 let mut total = 0;
817 for chunk in chunks {
818 total += match &self.config.column_mapping {
819 MysqlColumnMapping::Json { column } => {
820 self.insert_json(&mut conn, chunk, column).await?
821 }
822 MysqlColumnMapping::AutoMap => self.insert_auto_map(&mut conn, chunk).await?,
823 };
824 }
825
826 tracing::info!(
827 table = %self.config.table_name,
828 rows = total,
829 "MySQL write complete"
830 );
831 Ok(total)
832 }
833
834 async fn write_batch_partial(
843 &self,
844 records: &[Value],
845 ) -> Result<Vec<faucet_core::RowOutcome>, FaucetError> {
846 if matches!(self.config.write.write_mode, faucet_core::WriteMode::Append) {
847 self.write_batch(records).await?;
848 return Ok(records.iter().map(|_| Ok(())).collect());
849 }
850
851 let plan = faucet_core::plan_writes(records, &self.config.write);
852 self.apply_plan(&plan).await?;
853
854 let mut outcomes: Vec<faucet_core::RowOutcome> = records.iter().map(|_| Ok(())).collect();
855 for (idx, msg) in &plan.failed {
856 outcomes[*idx] = Err(FaucetError::Sink(format!(
857 "mysql {}: {msg}",
858 self.config.write.write_mode.as_str()
859 )));
860 }
861 Ok(outcomes)
862 }
863
864 fn supports_idempotent_writes(&self) -> bool {
865 true
866 }
867
868 async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
869 self.ensure_commit_table().await?;
870 let sql = format!(
871 "SELECT {k} FROM {t} WHERE {s} = ?",
872 t = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TABLE),
873 k = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TOKEN_COL),
874 s = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_SCOPE_COL),
875 );
876 let row = sqlx::query(&sql)
877 .bind(scope)
878 .fetch_optional(&self.pool)
879 .await
880 .map_err(|e| FaucetError::Sink(format!("MySQL token read failed: {e}")))?;
881 Ok(row.map(|r| r.get::<String, _>(0)))
882 }
883
884 async fn write_batch_idempotent(
885 &self,
886 records: &[Value],
887 scope: &str,
888 token: &str,
889 ) -> Result<usize, FaucetError> {
890 self.ensure_commit_table().await?;
891
892 let plan = if matches!(self.config.write.write_mode, faucet_core::WriteMode::Append) {
895 None
896 } else {
897 let plan = faucet_core::plan_writes(records, &self.config.write);
898 if let Some((idx, msg)) = plan.failed.first() {
899 return Err(FaucetError::Sink(format!(
900 "mysql {}: row {idx}: {msg}",
901 self.config.write.write_mode.as_str()
902 )));
903 }
904 Some(plan)
905 };
906
907 let mut tx = self
908 .pool
909 .begin()
910 .await
911 .map_err(|e| FaucetError::Sink(format!("MySQL transaction begin failed: {e}")))?;
912
913 let written = match &plan {
918 Some(plan) => {
919 let mut affected = 0usize;
920 if !plan.upserts.is_empty() {
921 affected += self
922 .insert_auto_map_with_conflict(
923 &mut tx,
924 &plan.upserts,
925 Some(&self.config.write.key),
926 )
927 .await?;
928 }
929 if !plan.deletes.is_empty() {
930 affected += self.delete_by_keys(&mut tx, &plan.deletes).await?;
931 }
932 affected
933 }
934 None => match &self.config.column_mapping {
935 MysqlColumnMapping::Json { column } => {
936 self.insert_json(&mut tx, records, column).await?
937 }
938 MysqlColumnMapping::AutoMap => self.insert_auto_map(&mut tx, records).await?,
939 },
940 };
941
942 let upsert = format!(
943 "INSERT INTO {t} ({s}, {k}) VALUES (?, ?) ON DUPLICATE KEY UPDATE {k} = VALUES({k})",
944 t = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TABLE),
945 s = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_SCOPE_COL),
946 k = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TOKEN_COL),
947 );
948 sqlx::query(&upsert)
949 .bind(scope)
950 .bind(token)
951 .execute(&mut *tx)
952 .await
953 .map_err(|e| FaucetError::Sink(format!("MySQL token upsert failed: {e}")))?;
954
955 tx.commit()
956 .await
957 .map_err(|e| FaucetError::Sink(format!("MySQL transaction commit failed: {e}")))?;
958
959 Ok(written)
960 }
961}
962
963#[cfg(test)]
964mod tests {
965 use super::*;
966
967 #[test]
971 fn commit_token_table_is_the_shared_constant() {
972 assert_eq!(
973 faucet_core::idempotency::COMMIT_TOKEN_TABLE,
974 "_faucet_commit_token"
975 );
976 }
977
978 #[test]
979 fn quote_ident_mysql_simple() {
980 assert_eq!(quote_ident_mysql("my_table"), "`my_table`");
981 }
982
983 #[test]
984 fn quote_ident_mysql_with_backtick() {
985 assert_eq!(quote_ident_mysql("has`tick"), "`has``tick`");
986 }
987
988 #[test]
989 fn quote_ident_mysql_empty() {
990 assert_eq!(quote_ident_mysql(""), "``");
991 }
992
993 #[test]
994 fn quote_ident_mysql_special_chars() {
995 assert_eq!(quote_ident_mysql("table; DROP"), "`table; DROP`");
996 }
997
998 #[test]
999 fn mysql_on_duplicate_clause() {
1000 let clause =
1001 on_duplicate_clause(&["id".to_string()], &["id".to_string(), "name".to_string()]);
1002 assert_eq!(clause, "ON DUPLICATE KEY UPDATE `name` = VALUES(`name`)");
1003 }
1004
1005 #[test]
1006 fn mysql_on_duplicate_all_keys_self_assign() {
1007 let clause = on_duplicate_clause(&["id".to_string()], &["id".to_string()]);
1008 assert_eq!(clause, "ON DUPLICATE KEY UPDATE `id` = `id`");
1009 }
1010
1011 #[test]
1012 fn mysql_on_duplicate_composite_key_partial_update() {
1013 let clause = on_duplicate_clause(
1014 &["a".to_string(), "b".to_string()],
1015 &["a".to_string(), "b".to_string(), "v".to_string()],
1016 );
1017 assert_eq!(clause, "ON DUPLICATE KEY UPDATE `v` = VALUES(`v`)");
1018 }
1019
1020 #[test]
1021 fn mysql_add_column_ddl() {
1022 let sql = build_add_column_sql("`t`", "email", SqlBaseType::Text);
1023 assert_eq!(sql, "ALTER TABLE `t` ADD COLUMN `email` LONGTEXT");
1024
1025 let sql = build_add_column_sql("`t`", "ev`il", SqlBaseType::Integer);
1027 assert_eq!(sql, "ALTER TABLE `t` ADD COLUMN `ev``il` BIGINT");
1028 }
1029
1030 #[test]
1031 fn mysql_modify_column_ddl() {
1032 let sql = build_modify_column_sql("`t`", "score", SqlBaseType::Double);
1033 assert_eq!(sql, "ALTER TABLE `t` MODIFY COLUMN `score` DOUBLE");
1034
1035 let sql = build_modify_column_sql("`t`", "flag", SqlBaseType::Boolean);
1036 assert_eq!(sql, "ALTER TABLE `t` MODIFY COLUMN `flag` TINYINT(1)");
1037 }
1038
1039 #[test]
1040 fn mysql_keyword_mapping() {
1041 assert_eq!(mysql_keyword(SqlBaseType::Integer), "BIGINT");
1042 assert_eq!(mysql_keyword(SqlBaseType::Double), "DOUBLE");
1043 assert_eq!(mysql_keyword(SqlBaseType::Boolean), "TINYINT(1)");
1044 assert_eq!(mysql_keyword(SqlBaseType::Text), "LONGTEXT");
1045 assert_eq!(mysql_keyword(SqlBaseType::Json), "JSON");
1046 }
1047
1048 fn idx(cols: &[&str]) -> std::collections::BTreeSet<String> {
1049 cols.iter().map(|s| s.to_string()).collect()
1050 }
1051
1052 fn keyvec(cols: &[&str]) -> Vec<String> {
1053 cols.iter().map(|s| s.to_string()).collect()
1054 }
1055
1056 #[test]
1057 fn key_matches_single_column_index() {
1058 let indexes = vec![idx(&["id"])];
1059 assert!(key_matches_unique_index(&indexes, &keyvec(&["id"])));
1060 }
1061
1062 #[test]
1063 fn key_matches_composite_index_reordered() {
1064 let indexes = vec![idx(&["a", "b"])];
1066 assert!(key_matches_unique_index(&indexes, &keyvec(&["b", "a"])));
1067 }
1068
1069 #[test]
1070 fn key_does_not_match_subset_of_index() {
1071 let indexes = vec![idx(&["a", "b"])];
1073 assert!(!key_matches_unique_index(&indexes, &keyvec(&["a"])));
1074 }
1075
1076 #[test]
1077 fn key_does_not_match_superset_of_index() {
1078 let indexes = vec![idx(&["a"])];
1080 assert!(!key_matches_unique_index(&indexes, &keyvec(&["a", "b"])));
1081 }
1082
1083 #[test]
1084 fn key_does_not_match_disjoint_index() {
1085 let indexes = vec![idx(&["id"])];
1086 assert!(!key_matches_unique_index(&indexes, &keyvec(&["other"])));
1087 }
1088
1089 #[test]
1090 fn key_matches_one_of_multiple_indexes() {
1091 let indexes = vec![idx(&["id"]), idx(&["email"])];
1093 assert!(key_matches_unique_index(&indexes, &keyvec(&["email"])));
1094 assert!(key_matches_unique_index(&indexes, &keyvec(&["id"])));
1095 assert!(!key_matches_unique_index(&indexes, &keyvec(&["name"])));
1097 assert!(!key_matches_unique_index(
1099 &indexes,
1100 &keyvec(&["id", "email"])
1101 ));
1102 }
1103
1104 #[test]
1105 fn key_does_not_match_empty_index_set() {
1106 let indexes: Vec<std::collections::BTreeSet<String>> = vec![];
1108 assert!(!key_matches_unique_index(&indexes, &keyvec(&["id"])));
1109 }
1110
1111 #[test]
1112 fn empty_key_never_matches() {
1113 let indexes = vec![idx(&["id"])];
1114 assert!(!key_matches_unique_index(&indexes, &keyvec(&[])));
1115 let empty_idx = vec![idx(&[])];
1117 assert!(!key_matches_unique_index(&empty_idx, &keyvec(&[])));
1118 }
1119
1120 #[test]
1121 fn key_matches_composite_index_among_several() {
1122 let indexes = vec![idx(&["id"]), idx(&["tenant", "slug"])];
1123 assert!(key_matches_unique_index(
1124 &indexes,
1125 &keyvec(&["slug", "tenant"])
1126 ));
1127 assert!(!key_matches_unique_index(&indexes, &keyvec(&["tenant"])));
1128 }
1129
1130 #[test]
1131 fn mysql_data_type_round_trips_to_json_schema() {
1132 use serde_json::json;
1133 assert_eq!(
1134 mysql_data_type_to_json_schema("bigint", false),
1135 json!({"type":"integer"})
1136 );
1137 assert_eq!(
1138 mysql_data_type_to_json_schema("int", false),
1139 json!({"type":"integer"})
1140 );
1141 assert_eq!(
1144 mysql_data_type_to_json_schema("tinyint", false),
1145 json!({"type":"integer"})
1146 );
1147 assert_eq!(
1148 mysql_data_type_to_json_schema("double", false),
1149 json!({"type":"number"})
1150 );
1151 assert_eq!(
1152 mysql_data_type_to_json_schema("decimal", false),
1153 json!({"type":"number"})
1154 );
1155 assert_eq!(
1156 mysql_data_type_to_json_schema("json", false),
1157 json!({"type":"object"})
1158 );
1159 assert_eq!(
1160 mysql_data_type_to_json_schema("varchar", false),
1161 json!({"type":"string"})
1162 );
1163 assert_eq!(
1165 mysql_data_type_to_json_schema("datetime", true),
1166 json!({"type":["string","null"]})
1167 );
1168 }
1169}