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
10const SCOPE_COL_WIDTH: usize = 255;
15
16fn scope_key(scope: &str) -> String {
18 faucet_core::idempotency::scope_key(scope, SCOPE_COL_WIDTH)
19}
20
21pub struct MysqlSink {
23 config: MysqlSinkConfig,
24 pool: MySqlPool,
25}
26
27fn quote_ident_mysql(name: &str) -> String {
32 format!("`{}`", name.replace('`', "``"))
33}
34
35fn mysql_keyword(t: SqlBaseType) -> &'static str {
41 match t {
42 SqlBaseType::Integer => "BIGINT",
43 SqlBaseType::Double => "DOUBLE",
44 SqlBaseType::Boolean => "TINYINT(1)",
45 SqlBaseType::Text => "LONGTEXT",
46 SqlBaseType::Json => "JSON",
47 }
48}
49
50fn build_add_column_sql(table: &str, col: &str, t: SqlBaseType) -> String {
57 format!(
58 "ALTER TABLE {table} ADD COLUMN {} {}",
59 quote_ident_mysql(col),
60 mysql_keyword(t)
61 )
62}
63
64fn build_modify_column_sql(table: &str, col: &str, t: SqlBaseType) -> String {
68 format!(
69 "ALTER TABLE {table} MODIFY COLUMN {} {}",
70 quote_ident_mysql(col),
71 mysql_keyword(t)
72 )
73}
74
75fn mysql_data_type_to_json_schema(data_type: &str, nullable: bool) -> Value {
84 let base = match data_type {
85 "bigint" | "int" | "integer" | "smallint" | "mediumint" | "tinyint" => "integer",
86 "double" | "float" | "decimal" | "numeric" => "number",
87 "json" => "object",
88 _ => "string",
89 };
90 if nullable {
91 serde_json::json!({ "type": [base, "null"] })
92 } else {
93 serde_json::json!({ "type": base })
94 }
95}
96
97fn key_matches_unique_index(
124 unique_indexes: &[std::collections::BTreeSet<String>],
125 key: &[String],
126) -> bool {
127 if key.is_empty() {
128 return false;
129 }
130 let key_set: std::collections::BTreeSet<String> = key.iter().cloned().collect();
131 unique_indexes.contains(&key_set)
132}
133
134fn on_duplicate_clause(key: &[String], all_cols: &[String]) -> String {
135 let updates: Vec<String> = all_cols
136 .iter()
137 .filter(|c| !key.iter().any(|k| k == *c))
138 .map(|c| {
139 let q = quote_ident_mysql(c);
140 format!("{q} = VALUES({q})")
141 })
142 .collect();
143 if updates.is_empty() {
144 let q = quote_ident_mysql(&key[0]);
145 format!("ON DUPLICATE KEY UPDATE {q} = {q}")
146 } else {
147 format!("ON DUPLICATE KEY UPDATE {}", updates.join(", "))
148 }
149}
150
151impl MysqlSink {
152 pub async fn new(config: MysqlSinkConfig) -> Result<Self, FaucetError> {
154 config.write.validate()?;
155 if !matches!(config.write.write_mode, faucet_core::WriteMode::Append)
156 && !matches!(config.column_mapping, MysqlColumnMapping::AutoMap)
157 {
158 return Err(FaucetError::Config(
159 "mysql sink: write_mode upsert/delete requires column_mapping: auto_map \
160 (key columns must be real columns, not inside a JSON blob)"
161 .into(),
162 ));
163 }
164
165 let pool = MySqlPoolOptions::new()
166 .max_connections(config.max_connections)
167 .connect(&config.connection_url)
168 .await
169 .map_err(|e| FaucetError::Sink(format!("MySQL connection failed: {e}")))?;
170
171 let sink = Self { config, pool };
172
173 if !matches!(sink.config.write.write_mode, faucet_core::WriteMode::Append) {
178 sink.assert_key_is_unique_index().await?;
179 }
180
181 Ok(sink)
182 }
183
184 async fn read_unique_indexes(
195 &self,
196 ) -> Result<Vec<std::collections::BTreeSet<String>>, FaucetError> {
197 let rows = sqlx::query(
200 "SELECT CAST(INDEX_NAME AS CHAR) AS INDEX_NAME, \
201 CAST(COLUMN_NAME AS CHAR) AS COLUMN_NAME \
202 FROM INFORMATION_SCHEMA.STATISTICS \
203 WHERE TABLE_NAME = ? AND TABLE_SCHEMA = DATABASE() AND NON_UNIQUE = 0 \
204 ORDER BY INDEX_NAME, SEQ_IN_INDEX",
205 )
206 .bind(&self.config.table_name)
207 .fetch_all(&self.pool)
208 .await
209 .map_err(|e| FaucetError::Sink(format!("failed to query table indexes: {e}")))?;
210
211 let mut by_index: std::collections::BTreeMap<String, std::collections::BTreeSet<String>> =
212 std::collections::BTreeMap::new();
213 for row in &rows {
214 let index_name: String = row.get("INDEX_NAME");
215 let column_name: String = row.get("COLUMN_NAME");
216 by_index.entry(index_name).or_default().insert(column_name);
217 }
218 Ok(by_index.into_values().collect())
219 }
220
221 async fn assert_key_is_unique_index(&self) -> Result<(), FaucetError> {
231 let unique_indexes = self.read_unique_indexes().await?;
232 if unique_indexes.is_empty() {
233 tracing::warn!(
234 table = %self.config.table_name,
235 "mysql sink: no PRIMARY/UNIQUE index found on target table (it may not exist \
236 yet); skipping upsert key validation — the first write will surface a \
237 missing-table or missing-constraint error"
238 );
239 return Ok(());
240 }
241 if !key_matches_unique_index(&unique_indexes, &self.config.write.key) {
242 let available: Vec<String> = unique_indexes
243 .iter()
244 .map(|idx| {
245 let mut cols: Vec<&str> = idx.iter().map(String::as_str).collect();
246 cols.sort_unstable();
247 format!("({})", cols.join(", "))
248 })
249 .collect();
250 return Err(FaucetError::Config(format!(
251 "mysql sink: write_mode {} requires `key` {:?} to exactly match a PRIMARY KEY or \
252 UNIQUE index on table '{}' — MySQL's `ON DUPLICATE KEY UPDATE` resolves on the \
253 table's real unique indexes, so an unmatched key would silently upsert on the \
254 wrong index. Existing unique indexes: {}",
255 self.config.write.write_mode.as_str(),
256 self.config.write.key,
257 self.config.table_name,
258 available.join(", "),
259 )));
260 }
261 Ok(())
262 }
263
264 async fn insert_json(
270 &self,
271 conn: &mut MySqlConnection,
272 records: &[Value],
273 column: &str,
274 ) -> Result<usize, FaucetError> {
275 if records.is_empty() {
276 return Ok(0);
277 }
278
279 let placeholders: Vec<&str> = records.iter().map(|_| "(?)").collect();
281 let insert_sql = format!(
282 "INSERT INTO {} ({}) VALUES {}",
283 quote_ident_mysql(&self.config.table_name),
284 quote_ident_mysql(column),
285 placeholders.join(", ")
286 );
287
288 let mut q = sqlx::query(&insert_sql);
289 for record in records {
290 let json_str = serde_json::to_string(record)
291 .map_err(|e| FaucetError::Sink(format!("failed to serialize record: {e}")))?;
292 q = q.bind(json_str);
293 }
294
295 q.execute(&mut *conn)
296 .await
297 .map_err(|e| FaucetError::Sink(format!("MySQL insert failed: {e}")))?;
298
299 Ok(records.len())
300 }
301
302 async fn insert_auto_map_with_conflict(
316 &self,
317 conn: &mut MySqlConnection,
318 records: &[Value],
319 conflict_key: Option<&[String]>,
320 ) -> Result<usize, FaucetError> {
321 if records.is_empty() {
322 return Ok(0);
323 }
324
325 let columns: Vec<String> = sqlx::query(
327 "SELECT COLUMN_NAME FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME = ? AND TABLE_SCHEMA = DATABASE() ORDER BY ORDINAL_POSITION"
328 )
329 .bind(&self.config.table_name)
330 .fetch_all(&mut *conn)
331 .await
332 .map_err(|e| FaucetError::Sink(format!("failed to query table columns: {e}")))?
333 .iter()
334 .map(|row| row.get::<String, _>("COLUMN_NAME"))
335 .collect();
336
337 if columns.is_empty() {
338 return Err(FaucetError::Sink(format!(
339 "table '{}' has no columns or does not exist",
340 self.config.table_name
341 )));
342 }
343
344 let mut matched_rows: Vec<Vec<(&String, &Value)>> = Vec::with_capacity(records.len());
351 let mut used: std::collections::HashSet<&str> = std::collections::HashSet::new();
352
353 for record in records {
354 let obj = record
355 .as_object()
356 .ok_or_else(|| FaucetError::Sink("AutoMap requires JSON object records".into()))?;
357
358 let matching: Vec<(&String, &Value)> = columns
359 .iter()
360 .filter_map(|col| obj.get(col).map(|v| (col, v)))
361 .collect();
362
363 if matching.is_empty() {
364 tracing::warn!(
365 record_keys = ?obj.keys().collect::<Vec<_>>(),
366 table_columns = ?columns,
367 "record has no keys matching table columns, skipping"
368 );
369 continue;
370 }
371
372 for (c, _) in &matching {
373 used.insert(c.as_str());
374 }
375 matched_rows.push(matching);
376 }
377
378 if matched_rows.is_empty() {
379 return Ok(0);
380 }
381
382 let insert_columns: Vec<String> = columns
384 .iter()
385 .filter(|c| used.contains(c.as_str()))
386 .cloned()
387 .collect();
388
389 let num_cols = insert_columns.len();
390 let num_rows = matched_rows.len();
391 let col_names: Vec<String> = insert_columns
392 .iter()
393 .map(|c| quote_ident_mysql(c))
394 .collect();
395
396 const MAX_MYSQL_PARAMS: usize = 65535;
402 let max_rows_per_insert = (MAX_MYSQL_PARAMS / num_cols).max(1);
403
404 for sub in matched_rows.chunks(max_rows_per_insert) {
405 let row_placeholder = format!("({})", vec!["?"; num_cols].join(", "));
407 let value_tuples: Vec<&str> =
408 (0..sub.len()).map(|_| row_placeholder.as_str()).collect();
409 let base_query = format!(
410 "INSERT INTO {} ({}) VALUES {}",
411 quote_ident_mysql(&self.config.table_name),
412 col_names.join(", "),
413 value_tuples.join(", ")
414 );
415 let query = match conflict_key {
416 Some(key) => format!("{base_query} {}", on_duplicate_clause(key, &insert_columns)),
417 None => base_query,
418 };
419
420 let mut q = sqlx::query(&query);
421 for matched in sub {
422 for col in &insert_columns {
423 let val = matched.iter().find(|(c, _)| *c == col).map(|(_, v)| *v);
424 q = match val {
429 None | Some(Value::Null) => q.bind(None::<String>),
430 Some(Value::Bool(b)) => q.bind(*b),
431 Some(Value::Number(n)) => {
432 if let Some(i) = n.as_i64() {
433 q.bind(i)
434 } else if let Some(f) = n.as_f64() {
435 q.bind(f)
436 } else {
437 q.bind(n.to_string())
439 }
440 }
441 Some(Value::String(s)) => q.bind(s.clone()),
442 Some(v) => q.bind(v.to_string()),
445 };
446 }
447 }
448
449 q.execute(&mut *conn)
450 .await
451 .map_err(|e| FaucetError::Sink(format!("MySQL insert failed: {e}")))?;
452 }
453
454 Ok(num_rows)
455 }
456
457 async fn insert_auto_map(
463 &self,
464 conn: &mut MySqlConnection,
465 records: &[Value],
466 ) -> Result<usize, FaucetError> {
467 self.insert_auto_map_with_conflict(conn, records, None)
468 .await
469 }
470
471 async fn delete_by_keys(
475 &self,
476 conn: &mut MySqlConnection,
477 deletes: &[faucet_core::KeyTuple],
478 ) -> Result<usize, FaucetError> {
479 if deletes.is_empty() {
480 return Ok(0);
481 }
482 let key = &self.config.write.key;
483 let table_ref = quote_ident_mysql(&self.config.table_name);
484 let col_list = key
485 .iter()
486 .map(|k| quote_ident_mysql(k))
487 .collect::<Vec<_>>()
488 .join(", ");
489
490 const MAX_MYSQL_PARAMS: usize = 65535;
491 let per = (MAX_MYSQL_PARAMS / key.len().max(1)).max(1);
492 let mut total = 0usize;
493
494 for chunk in deletes.chunks(per) {
495 let tuples: Vec<String> = chunk
496 .iter()
497 .map(|_| format!("({})", vec!["?"; key.len()].join(", ")))
498 .collect();
499 let sql = format!(
500 "DELETE FROM {table_ref} WHERE ({col_list}) IN ({})",
501 tuples.join(", ")
502 );
503 let mut q = sqlx::query(&sql);
504 for kt in chunk {
505 for (_, v) in &kt.0 {
506 q = match v {
508 Value::Null => q.bind(None::<String>),
509 Value::Bool(b) => q.bind(*b),
510 Value::Number(n) => {
511 if let Some(i) = n.as_i64() {
512 q.bind(i)
513 } else if let Some(f) = n.as_f64() {
514 q.bind(f)
515 } else {
516 q.bind(n.to_string())
517 }
518 }
519 Value::String(s) => q.bind(s.clone()),
520 other => q.bind(other.to_string()),
521 };
522 }
523 }
524 let res = q
525 .execute(&mut *conn)
526 .await
527 .map_err(|e| FaucetError::Sink(format!("MySQL delete failed: {e}")))?;
528 total += res.rows_affected() as usize;
529 }
530 Ok(total)
531 }
532
533 async fn apply_plan(&self, plan: &faucet_core::WritePlan) -> Result<usize, FaucetError> {
537 let mut tx = self
538 .pool
539 .begin()
540 .await
541 .map_err(|e| FaucetError::Sink(format!("MySQL transaction begin failed: {e}")))?;
542
543 let mut affected = 0usize;
544 if !plan.upserts.is_empty() {
545 affected += self
546 .insert_auto_map_with_conflict(&mut tx, &plan.upserts, Some(&self.config.write.key))
547 .await?;
548 }
549 if !plan.deletes.is_empty() {
550 affected += self.delete_by_keys(&mut tx, &plan.deletes).await?;
551 }
552
553 tx.commit()
554 .await
555 .map_err(|e| FaucetError::Sink(format!("MySQL transaction commit failed: {e}")))?;
556 Ok(affected)
557 }
558
559 async fn read_columns(&self) -> Result<Vec<(String, String, bool)>, FaucetError> {
568 let rows = sqlx::query(
572 "SELECT CAST(COLUMN_NAME AS CHAR) AS COLUMN_NAME, \
573 CAST(DATA_TYPE AS CHAR) AS DATA_TYPE, \
574 CAST(IS_NULLABLE AS CHAR) AS IS_NULLABLE \
575 FROM INFORMATION_SCHEMA.COLUMNS \
576 WHERE TABLE_NAME = ? AND TABLE_SCHEMA = DATABASE() ORDER BY ORDINAL_POSITION",
577 )
578 .bind(&self.config.table_name)
579 .fetch_all(&self.pool)
580 .await
581 .map_err(|e| FaucetError::Sink(format!("failed to query table columns: {e}")))?;
582
583 Ok(rows
584 .iter()
585 .map(|row| {
586 (
587 row.get::<String, _>("COLUMN_NAME"),
588 row.get::<String, _>("DATA_TYPE").to_ascii_lowercase(),
591 row.get::<String, _>("IS_NULLABLE")
592 .eq_ignore_ascii_case("YES"),
593 )
594 })
595 .collect())
596 }
597
598 async fn ensure_commit_table(&self) -> Result<(), FaucetError> {
607 let sql = format!(
608 "CREATE TABLE IF NOT EXISTS {t} ({s} VARCHAR({w}) PRIMARY KEY, {k} TEXT NOT NULL, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)",
609 t = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TABLE),
610 s = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_SCOPE_COL),
611 k = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TOKEN_COL),
612 w = SCOPE_COL_WIDTH,
613 );
614 sqlx::query(&sql)
615 .execute(&self.pool)
616 .await
617 .map_err(|e| FaucetError::Sink(format!("MySQL commit-table create failed: {e}")))?;
618 Ok(())
619 }
620}
621
622#[async_trait]
623impl faucet_core::Sink for MysqlSink {
624 fn config_schema(&self) -> serde_json::Value {
625 serde_json::to_value(faucet_core::schema_for!(MysqlSinkConfig))
626 .expect("schema serialization")
627 }
628
629 fn dataset_uri(&self) -> String {
630 format!(
631 "{}?table={}",
632 faucet_core::redact_uri_credentials(&self.config.connection_url),
633 self.config.table_name
634 )
635 }
636
637 async fn check(
643 &self,
644 ctx: &faucet_core::check::CheckContext,
645 ) -> Result<faucet_core::check::CheckReport, FaucetError> {
646 use faucet_core::check::{CheckReport, Probe};
647
648 let started = std::time::Instant::now();
649 let probe =
650 match tokio::time::timeout(ctx.timeout, sqlx::query("SELECT 1").execute(&self.pool))
651 .await
652 {
653 Ok(Ok(_)) => Probe::pass("auth", started.elapsed()),
654 Ok(Err(e)) => Probe::fail_hint(
655 "auth",
656 started.elapsed(),
657 e.to_string(),
658 "check connection_url / credentials / that the database is reachable",
659 ),
660 Err(_) => Probe::fail_hint(
661 "auth",
662 started.elapsed(),
663 "timed out",
664 "check connection_url / credentials / that the database is reachable",
665 ),
666 };
667 Ok(CheckReport::single(probe))
668 }
669
670 fn supported_write_modes(&self) -> &'static [faucet_core::WriteMode] {
671 &[
672 faucet_core::WriteMode::Append,
673 faucet_core::WriteMode::Upsert,
674 faucet_core::WriteMode::Delete,
675 ]
676 }
677
678 fn dedups_by_key(&self) -> bool {
679 self.config.write.dedups_by_key()
680 }
681
682 fn supports_schema_evolution(&self) -> bool {
683 true
684 }
685
686 async fn current_schema(&self) -> Result<Option<serde_json::Value>, FaucetError> {
694 let columns = self.read_columns().await?;
695 if columns.is_empty() {
696 return Ok(None); }
698
699 let mut props = serde_json::Map::new();
700 for (name, data_type, nullable) in columns {
701 props.insert(name, mysql_data_type_to_json_schema(&data_type, nullable));
702 }
703 Ok(Some(
704 serde_json::json!({ "type": "object", "properties": props }),
705 ))
706 }
707
708 async fn evolve_schema(&self, evolution: &SchemaEvolution) -> Result<(), FaucetError> {
718 let table_ref = quote_ident_mysql(&self.config.table_name);
719
720 let current = self.read_columns().await?;
724 let existing: std::collections::HashSet<&str> =
725 current.iter().map(|(n, _, _)| n.as_str()).collect();
726
727 let mut conn = self
728 .pool
729 .acquire()
730 .await
731 .map_err(|e| FaucetError::Sink(format!("MySQL evolve acquire failed: {e}")))?;
732
733 for c in &evolution.additions {
734 if existing.contains(c.name.as_str()) {
737 continue;
738 }
739 let t = json_schema_base_type(&c.to).unwrap_or(SqlBaseType::Text);
740 sqlx::query(&build_add_column_sql(&table_ref, &c.name, t))
741 .execute(&mut *conn)
742 .await
743 .map_err(|e| {
744 FaucetError::Sink(format!("MySQL ADD COLUMN {} failed: {e}", c.name))
745 })?;
746 }
747
748 for c in &evolution.widenings {
749 let t = json_schema_base_type(&c.to).unwrap_or(SqlBaseType::Text);
750 sqlx::query(&build_modify_column_sql(&table_ref, &c.name, t))
751 .execute(&mut *conn)
752 .await
753 .map_err(|e| {
754 FaucetError::Sink(format!("MySQL MODIFY COLUMN {} failed: {e}", c.name))
755 })?;
756 }
757
758 for col in &evolution.relax_nullability {
759 let existing_type = current
763 .iter()
764 .find(|(n, _, _)| n == col)
765 .map(|(_, dt, nullable)| {
766 let fragment = mysql_data_type_to_json_schema(dt, *nullable);
767 json_schema_base_type(&fragment).unwrap_or(SqlBaseType::Text)
768 })
769 .unwrap_or(SqlBaseType::Text);
770 let sql = format!(
771 "ALTER TABLE {table_ref} MODIFY COLUMN {} {} NULL",
772 quote_ident_mysql(col),
773 mysql_keyword(existing_type)
774 );
775 sqlx::query(&sql)
776 .execute(&mut *conn)
777 .await
778 .map_err(|e| FaucetError::Sink(format!("MySQL DROP NOT NULL {col} failed: {e}")))?;
779 }
780
781 Ok(())
782 }
783
784 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
797 if records.is_empty() {
798 return Ok(0);
799 }
800
801 if !matches!(self.config.write.write_mode, faucet_core::WriteMode::Append) {
803 let plan = faucet_core::plan_writes(records, &self.config.write);
804 if let Some((idx, msg)) = plan.failed.first() {
805 return Err(FaucetError::Sink(format!(
806 "mysql {}: row {idx}: {msg}",
807 self.config.write.write_mode.as_str()
808 )));
809 }
810 return self.apply_plan(&plan).await;
811 }
812
813 let mut conn = self
814 .pool
815 .acquire()
816 .await
817 .map_err(|e| FaucetError::Sink(format!("MySQL pool acquire failed: {e}")))?;
818
819 let chunks: Vec<&[Value]> = if self.config.batch_size == 0 {
820 vec![records]
824 } else {
825 records.chunks(self.config.batch_size).collect()
826 };
827
828 let mut total = 0;
829 for chunk in chunks {
830 total += match &self.config.column_mapping {
831 MysqlColumnMapping::Json { column } => {
832 self.insert_json(&mut conn, chunk, column).await?
833 }
834 MysqlColumnMapping::AutoMap => self.insert_auto_map(&mut conn, chunk).await?,
835 };
836 }
837
838 tracing::info!(
839 table = %self.config.table_name,
840 rows = total,
841 "MySQL write complete"
842 );
843 Ok(total)
844 }
845
846 async fn write_batch_partial(
855 &self,
856 records: &[Value],
857 ) -> Result<Vec<faucet_core::RowOutcome>, FaucetError> {
858 if matches!(self.config.write.write_mode, faucet_core::WriteMode::Append) {
859 self.write_batch(records).await?;
860 return Ok(records.iter().map(|_| Ok(())).collect());
861 }
862
863 let plan = faucet_core::plan_writes(records, &self.config.write);
864 self.apply_plan(&plan).await?;
865
866 let mut outcomes: Vec<faucet_core::RowOutcome> = records.iter().map(|_| Ok(())).collect();
867 for (idx, msg) in &plan.failed {
868 outcomes[*idx] = Err(FaucetError::Sink(format!(
869 "mysql {}: {msg}",
870 self.config.write.write_mode.as_str()
871 )));
872 }
873 Ok(outcomes)
874 }
875
876 fn supports_idempotent_writes(&self) -> bool {
877 true
878 }
879
880 async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
881 self.ensure_commit_table().await?;
882 let sql = format!(
883 "SELECT {k} FROM {t} WHERE {s} = ?",
884 t = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TABLE),
885 k = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TOKEN_COL),
886 s = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_SCOPE_COL),
887 );
888 let row = sqlx::query(&sql)
889 .bind(scope_key(scope))
890 .fetch_optional(&self.pool)
891 .await
892 .map_err(|e| FaucetError::Sink(format!("MySQL token read failed: {e}")))?;
893 Ok(row.map(|r| r.get::<String, _>(0)))
894 }
895
896 async fn write_batch_idempotent(
897 &self,
898 records: &[Value],
899 scope: &str,
900 token: &str,
901 ) -> Result<usize, FaucetError> {
902 self.ensure_commit_table().await?;
903
904 let plan = if matches!(self.config.write.write_mode, faucet_core::WriteMode::Append) {
907 None
908 } else {
909 let plan = faucet_core::plan_writes(records, &self.config.write);
910 if let Some((idx, msg)) = plan.failed.first() {
911 return Err(FaucetError::Sink(format!(
912 "mysql {}: row {idx}: {msg}",
913 self.config.write.write_mode.as_str()
914 )));
915 }
916 Some(plan)
917 };
918
919 let mut tx = self
920 .pool
921 .begin()
922 .await
923 .map_err(|e| FaucetError::Sink(format!("MySQL transaction begin failed: {e}")))?;
924
925 let written = match &plan {
930 Some(plan) => {
931 let mut affected = 0usize;
932 if !plan.upserts.is_empty() {
933 affected += self
934 .insert_auto_map_with_conflict(
935 &mut tx,
936 &plan.upserts,
937 Some(&self.config.write.key),
938 )
939 .await?;
940 }
941 if !plan.deletes.is_empty() {
942 affected += self.delete_by_keys(&mut tx, &plan.deletes).await?;
943 }
944 affected
945 }
946 None => match &self.config.column_mapping {
947 MysqlColumnMapping::Json { column } => {
948 self.insert_json(&mut tx, records, column).await?
949 }
950 MysqlColumnMapping::AutoMap => self.insert_auto_map(&mut tx, records).await?,
951 },
952 };
953
954 let upsert = format!(
955 "INSERT INTO {t} ({s}, {k}) VALUES (?, ?) ON DUPLICATE KEY UPDATE {k} = VALUES({k})",
956 t = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TABLE),
957 s = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_SCOPE_COL),
958 k = quote_ident_mysql(faucet_core::idempotency::COMMIT_TOKEN_TOKEN_COL),
959 );
960 sqlx::query(&upsert)
961 .bind(scope_key(scope))
962 .bind(token)
963 .execute(&mut *tx)
964 .await
965 .map_err(|e| FaucetError::Sink(format!("MySQL token upsert failed: {e}")))?;
966
967 tx.commit()
968 .await
969 .map_err(|e| FaucetError::Sink(format!("MySQL transaction commit failed: {e}")))?;
970
971 Ok(written)
972 }
973}
974
975#[cfg(test)]
976mod tests {
977 use super::*;
978
979 #[test]
983 fn commit_token_table_is_the_shared_constant() {
984 assert_eq!(
985 faucet_core::idempotency::COMMIT_TOKEN_TABLE,
986 "_faucet_commit_token"
987 );
988 }
989
990 #[test]
991 fn quote_ident_mysql_simple() {
992 assert_eq!(quote_ident_mysql("my_table"), "`my_table`");
993 }
994
995 #[test]
996 fn quote_ident_mysql_with_backtick() {
997 assert_eq!(quote_ident_mysql("has`tick"), "`has``tick`");
998 }
999
1000 #[test]
1001 fn quote_ident_mysql_empty() {
1002 assert_eq!(quote_ident_mysql(""), "``");
1003 }
1004
1005 #[test]
1006 fn quote_ident_mysql_special_chars() {
1007 assert_eq!(quote_ident_mysql("table; DROP"), "`table; DROP`");
1008 }
1009
1010 #[test]
1011 fn mysql_on_duplicate_clause() {
1012 let clause =
1013 on_duplicate_clause(&["id".to_string()], &["id".to_string(), "name".to_string()]);
1014 assert_eq!(clause, "ON DUPLICATE KEY UPDATE `name` = VALUES(`name`)");
1015 }
1016
1017 #[test]
1018 fn mysql_on_duplicate_all_keys_self_assign() {
1019 let clause = on_duplicate_clause(&["id".to_string()], &["id".to_string()]);
1020 assert_eq!(clause, "ON DUPLICATE KEY UPDATE `id` = `id`");
1021 }
1022
1023 #[test]
1024 fn mysql_on_duplicate_composite_key_partial_update() {
1025 let clause = on_duplicate_clause(
1026 &["a".to_string(), "b".to_string()],
1027 &["a".to_string(), "b".to_string(), "v".to_string()],
1028 );
1029 assert_eq!(clause, "ON DUPLICATE KEY UPDATE `v` = VALUES(`v`)");
1030 }
1031
1032 #[test]
1033 fn mysql_add_column_ddl() {
1034 let sql = build_add_column_sql("`t`", "email", SqlBaseType::Text);
1035 assert_eq!(sql, "ALTER TABLE `t` ADD COLUMN `email` LONGTEXT");
1036
1037 let sql = build_add_column_sql("`t`", "ev`il", SqlBaseType::Integer);
1039 assert_eq!(sql, "ALTER TABLE `t` ADD COLUMN `ev``il` BIGINT");
1040 }
1041
1042 #[test]
1043 fn mysql_modify_column_ddl() {
1044 let sql = build_modify_column_sql("`t`", "score", SqlBaseType::Double);
1045 assert_eq!(sql, "ALTER TABLE `t` MODIFY COLUMN `score` DOUBLE");
1046
1047 let sql = build_modify_column_sql("`t`", "flag", SqlBaseType::Boolean);
1048 assert_eq!(sql, "ALTER TABLE `t` MODIFY COLUMN `flag` TINYINT(1)");
1049 }
1050
1051 #[test]
1052 fn mysql_keyword_mapping() {
1053 assert_eq!(mysql_keyword(SqlBaseType::Integer), "BIGINT");
1054 assert_eq!(mysql_keyword(SqlBaseType::Double), "DOUBLE");
1055 assert_eq!(mysql_keyword(SqlBaseType::Boolean), "TINYINT(1)");
1056 assert_eq!(mysql_keyword(SqlBaseType::Text), "LONGTEXT");
1057 assert_eq!(mysql_keyword(SqlBaseType::Json), "JSON");
1058 }
1059
1060 fn idx(cols: &[&str]) -> std::collections::BTreeSet<String> {
1061 cols.iter().map(|s| s.to_string()).collect()
1062 }
1063
1064 fn keyvec(cols: &[&str]) -> Vec<String> {
1065 cols.iter().map(|s| s.to_string()).collect()
1066 }
1067
1068 #[test]
1069 fn key_matches_single_column_index() {
1070 let indexes = vec![idx(&["id"])];
1071 assert!(key_matches_unique_index(&indexes, &keyvec(&["id"])));
1072 }
1073
1074 #[test]
1075 fn key_matches_composite_index_reordered() {
1076 let indexes = vec![idx(&["a", "b"])];
1078 assert!(key_matches_unique_index(&indexes, &keyvec(&["b", "a"])));
1079 }
1080
1081 #[test]
1082 fn key_does_not_match_subset_of_index() {
1083 let indexes = vec![idx(&["a", "b"])];
1085 assert!(!key_matches_unique_index(&indexes, &keyvec(&["a"])));
1086 }
1087
1088 #[test]
1089 fn key_does_not_match_superset_of_index() {
1090 let indexes = vec![idx(&["a"])];
1092 assert!(!key_matches_unique_index(&indexes, &keyvec(&["a", "b"])));
1093 }
1094
1095 #[test]
1096 fn key_does_not_match_disjoint_index() {
1097 let indexes = vec![idx(&["id"])];
1098 assert!(!key_matches_unique_index(&indexes, &keyvec(&["other"])));
1099 }
1100
1101 #[test]
1102 fn key_matches_one_of_multiple_indexes() {
1103 let indexes = vec![idx(&["id"]), idx(&["email"])];
1105 assert!(key_matches_unique_index(&indexes, &keyvec(&["email"])));
1106 assert!(key_matches_unique_index(&indexes, &keyvec(&["id"])));
1107 assert!(!key_matches_unique_index(&indexes, &keyvec(&["name"])));
1109 assert!(!key_matches_unique_index(
1111 &indexes,
1112 &keyvec(&["id", "email"])
1113 ));
1114 }
1115
1116 #[test]
1117 fn key_does_not_match_empty_index_set() {
1118 let indexes: Vec<std::collections::BTreeSet<String>> = vec![];
1120 assert!(!key_matches_unique_index(&indexes, &keyvec(&["id"])));
1121 }
1122
1123 #[test]
1124 fn empty_key_never_matches() {
1125 let indexes = vec![idx(&["id"])];
1126 assert!(!key_matches_unique_index(&indexes, &keyvec(&[])));
1127 let empty_idx = vec![idx(&[])];
1129 assert!(!key_matches_unique_index(&empty_idx, &keyvec(&[])));
1130 }
1131
1132 #[test]
1133 fn key_matches_composite_index_among_several() {
1134 let indexes = vec![idx(&["id"]), idx(&["tenant", "slug"])];
1135 assert!(key_matches_unique_index(
1136 &indexes,
1137 &keyvec(&["slug", "tenant"])
1138 ));
1139 assert!(!key_matches_unique_index(&indexes, &keyvec(&["tenant"])));
1140 }
1141
1142 #[test]
1143 fn mysql_data_type_round_trips_to_json_schema() {
1144 use serde_json::json;
1145 assert_eq!(
1146 mysql_data_type_to_json_schema("bigint", false),
1147 json!({"type":"integer"})
1148 );
1149 assert_eq!(
1150 mysql_data_type_to_json_schema("int", false),
1151 json!({"type":"integer"})
1152 );
1153 assert_eq!(
1156 mysql_data_type_to_json_schema("tinyint", false),
1157 json!({"type":"integer"})
1158 );
1159 assert_eq!(
1160 mysql_data_type_to_json_schema("double", false),
1161 json!({"type":"number"})
1162 );
1163 assert_eq!(
1164 mysql_data_type_to_json_schema("decimal", false),
1165 json!({"type":"number"})
1166 );
1167 assert_eq!(
1168 mysql_data_type_to_json_schema("json", false),
1169 json!({"type":"object"})
1170 );
1171 assert_eq!(
1172 mysql_data_type_to_json_schema("varchar", false),
1173 json!({"type":"string"})
1174 );
1175 assert_eq!(
1177 mysql_data_type_to_json_schema("datetime", true),
1178 json!({"type":["string","null"]})
1179 );
1180 }
1181}