1use crate::key::KeySet;
16use crate::model::batch_write_request::MutationGroup as ProtoMutationGroup;
17use crate::model::mutation::Operation;
18use crate::value::Value;
19use rand::seq::IteratorRandom;
20use std::slice::Iter;
21use std::vec::IntoIter;
22
23#[derive(Clone, Debug, PartialEq)]
37pub struct Mutation {
38 pub(crate) inner: InternalMutation,
39}
40
41#[derive(Clone, Debug, PartialEq)]
42pub(crate) enum InternalMutation {
43 Insert(Write),
46 Update(Write),
49 InsertOrUpdate(Write),
52 Replace(Write),
55 Delete(Delete),
57}
58
59#[derive(Clone, Debug, PartialEq)]
61pub(crate) struct Write {
62 pub(crate) table: String,
63 pub(crate) columns: Vec<String>,
64 pub(crate) values: Vec<Value>,
65}
66
67#[derive(Clone, Debug, PartialEq)]
69pub(crate) struct Delete {
70 pub(crate) table: String,
71 pub(crate) key_set: KeySet,
74}
75
76impl Mutation {
77 pub fn new_insert_builder(table: impl Into<String>) -> WriteBuilder {
87 WriteBuilder::new(table, MutationType::Insert)
88 }
89
90 pub fn new_update_builder(table: impl Into<String>) -> WriteBuilder {
101 WriteBuilder::new(table, MutationType::Update)
102 }
103
104 pub fn new_insert_or_update_builder(table: impl Into<String>) -> WriteBuilder {
115 WriteBuilder::new(table, MutationType::InsertOrUpdate)
116 }
117
118 pub fn new_replace_builder(table: impl Into<String>) -> WriteBuilder {
129 WriteBuilder::new(table, MutationType::Replace)
130 }
131
132 pub fn delete(table: impl Into<String>, key_set: KeySet) -> Mutation {
139 Mutation {
140 inner: InternalMutation::Delete(Delete {
141 table: table.into(),
142 key_set,
143 }),
144 }
145 }
146
147 pub(crate) fn build_proto(self) -> crate::model::Mutation {
148 match self.inner {
149 InternalMutation::Insert(write) => {
150 crate::model::Mutation::new().set_insert(write.into_proto())
151 }
152 InternalMutation::Update(write) => {
153 crate::model::Mutation::new().set_update(write.into_proto())
154 }
155 InternalMutation::InsertOrUpdate(write) => {
156 crate::model::Mutation::new().set_insert_or_update(write.into_proto())
157 }
158 InternalMutation::Replace(write) => {
159 crate::model::Mutation::new().set_replace(write.into_proto())
160 }
161 InternalMutation::Delete(delete) => {
162 crate::model::Mutation::new().set_delete(delete.into_proto())
163 }
164 }
165 }
166
167 pub(crate) fn select_mutation_key(
174 mutations: &[crate::model::Mutation],
175 ) -> Option<crate::model::Mutation> {
176 if mutations.is_empty() {
177 return None;
178 }
179
180 let selected_non_insert = mutations
182 .iter()
183 .filter(|m| {
184 m.operation.as_ref().is_some_and(|op| {
185 !matches!(
186 op,
187 Operation::Insert(_) | Operation::Send(_) | Operation::Ack(_)
188 )
189 })
190 })
191 .choose(&mut rand::rng())
192 .cloned();
193
194 if selected_non_insert.is_some() {
195 return selected_non_insert;
196 }
197
198 let max_insert = mutations
200 .iter()
201 .filter_map(|m| match &m.operation {
202 Some(Operation::Insert(write)) => Some((m, write.values.len())),
203 _ => None,
204 })
205 .max_by_key(|&(_, rows)| rows)
206 .map(|(m, _)| m);
207
208 max_insert.cloned().or_else(|| mutations.first().cloned())
209 }
210}
211
212impl Write {
213 fn into_proto(self) -> crate::model::mutation::Write {
214 crate::model::mutation::Write::new()
215 .set_table(self.table)
216 .set_columns(self.columns)
217 .set_values(vec![
218 self.values
219 .into_iter()
220 .map(Value::into_serde_value)
221 .collect::<wkt::ListValue>(),
222 ])
223 }
224}
225
226impl Delete {
227 fn into_proto(self) -> crate::model::mutation::Delete {
228 crate::model::mutation::Delete::new()
229 .set_table(self.table)
230 .set_key_set(self.key_set.into_proto())
231 }
232}
233
234pub struct WriteBuilder {
236 table: String,
237 mutation_type: MutationType,
238 columns: Vec<String>,
239 values: Vec<Value>,
240}
241
242enum MutationType {
243 Insert,
244 Update,
245 InsertOrUpdate,
246 Replace,
247}
248
249impl WriteBuilder {
250 fn new(table: impl Into<String>, mutation_type: MutationType) -> Self {
251 Self {
252 table: table.into(),
253 mutation_type,
254 columns: Vec::new(),
255 values: Vec::new(),
256 }
257 }
258
259 pub fn set(self, column_name: impl Into<String>) -> ValueBinder {
269 ValueBinder {
270 builder: self,
271 column: column_name.into(),
272 }
273 }
274
275 pub fn build(self) -> Mutation {
277 let write = Write {
278 table: self.table,
279 columns: self.columns,
280 values: self.values,
281 };
282 let inner = match self.mutation_type {
283 MutationType::Insert => InternalMutation::Insert(write),
284 MutationType::Update => InternalMutation::Update(write),
285 MutationType::InsertOrUpdate => InternalMutation::InsertOrUpdate(write),
286 MutationType::Replace => InternalMutation::Replace(write),
287 };
288 Mutation { inner }
289 }
290}
291
292pub struct ValueBinder {
294 builder: WriteBuilder,
295 column: String,
296}
297
298impl ValueBinder {
299 pub fn to<T: Into<Value>>(mut self, value: T) -> WriteBuilder {
301 self.builder.columns.push(self.column);
302 self.builder.values.push(value.into());
303 self.builder
304 }
305}
306
307#[derive(Clone, Debug, PartialEq)]
309#[non_exhaustive]
310pub struct MutationGroup {
311 mutations: Vec<Mutation>,
312}
313
314impl MutationGroup {
315 pub fn new(mutations: Vec<Mutation>) -> Self {
317 Self { mutations }
318 }
319
320 pub fn mutations(&self) -> &[Mutation] {
322 &self.mutations
323 }
324
325 #[allow(dead_code)]
326 pub(crate) fn build_proto(self) -> ProtoMutationGroup {
327 ProtoMutationGroup::new().set_mutations(self.mutations.into_iter().map(|m| m.build_proto()))
328 }
329}
330
331impl IntoIterator for MutationGroup {
332 type Item = Mutation;
333 type IntoIter = IntoIter<Mutation>;
334
335 fn into_iter(self) -> Self::IntoIter {
336 self.mutations.into_iter()
337 }
338}
339
340impl<'a> IntoIterator for &'a MutationGroup {
341 type Item = &'a Mutation;
342 type IntoIter = Iter<'a, Mutation>;
343
344 fn into_iter(self) -> Self::IntoIter {
345 self.mutations.iter()
346 }
347}
348
349#[macro_export]
389macro_rules! mutation {
390 (@col $id:ident) => {
392 stringify!($id)
393 };
394 (@col $lit:expr) => {
395 $lit
396 };
397
398 (@build_write $builder:ident, $table:tt, { $($col:tt : $val:expr),* $(,)? }) => {
400 $crate::mutation::Mutation::$builder($table)
401 $(.set($crate::mutation!(@col $col)).to($val))*
402 .build()
403 };
404
405 (delete $table:expr, $key_set:expr) => {
407 $crate::mutation::Mutation::delete($table, $key_set)
408 };
409
410 (insert $table:tt { $($rest:tt)* }) => {
412 $crate::mutation!(@build_write new_insert_builder, $table, { $($rest)* })
413 };
414
415 (update $table:tt { $($rest:tt)* }) => {
417 $crate::mutation!(@build_write new_update_builder, $table, { $($rest)* })
418 };
419
420 (insert_or_update $table:tt { $($rest:tt)* }) => {
422 $crate::mutation!(@build_write new_insert_or_update_builder, $table, { $($rest)* })
423 };
424
425 (replace $table:tt { $($rest:tt)* }) => {
427 $crate::mutation!(@build_write new_replace_builder, $table, { $($rest)* })
428 };
429}
430
431#[cfg(test)]
432mod tests {
433 use super::*;
434 use crate::to_value::ToValue;
435
436 #[test]
437 fn auto_traits() {
438 static_assertions::assert_impl_all!(Mutation: Send, Sync, Clone, std::fmt::Debug);
439 static_assertions::assert_impl_all!(Write: Send, Sync, Clone, std::fmt::Debug);
440 static_assertions::assert_impl_all!(Delete: Send, Sync, Clone, std::fmt::Debug);
441 static_assertions::assert_impl_all!(WriteBuilder: Send, Sync);
442 static_assertions::assert_impl_all!(ValueBinder: Send, Sync);
443 static_assertions::assert_impl_all!(MutationGroup: Send, Sync, Clone, std::fmt::Debug);
444 }
445
446 #[test]
447 fn mutation_group() {
448 let mutation1 = Mutation::new_insert_builder("Users")
449 .set("UserId")
450 .to(1)
451 .build();
452 let mutation2 = Mutation::new_insert_builder("Users")
453 .set("UserId")
454 .to(2)
455 .build();
456 let group = MutationGroup::new(vec![mutation1.clone(), mutation2.clone()]);
457 assert_eq!(group.mutations.len(), 2);
458 assert_eq!(group.mutations[0], mutation1);
459 assert_eq!(group.mutations[1], mutation2);
460 }
461
462 #[test]
463 fn mutation_group_into_iter() {
464 let mutation1 = Mutation::new_insert_builder("Users")
465 .set("UserId")
466 .to(1)
467 .build();
468 let mutation2 = Mutation::new_insert_builder("Users")
469 .set("UserId")
470 .to(2)
471 .build();
472 let group = MutationGroup::new(vec![mutation1.clone(), mutation2.clone()]);
473
474 let mutations: Vec<_> = group.into_iter().collect();
475 assert_eq!(mutations, vec![mutation1, mutation2]);
476 }
477
478 #[test]
479 fn mutation_group_iter_ref() {
480 let mutation1 = Mutation::new_insert_builder("Users")
481 .set("UserId")
482 .to(1)
483 .build();
484 let mutation2 = Mutation::new_insert_builder("Users")
485 .set("UserId")
486 .to(2)
487 .build();
488 let group = MutationGroup::new(vec![mutation1.clone(), mutation2.clone()]);
489
490 let mutations: Vec<_> = (&group).into_iter().collect();
491 assert_eq!(mutations, vec![&mutation1, &mutation2]);
492 }
493
494 #[test]
495 fn value_binder_to_owned_value() {
496 let by_ref = Mutation::new_insert_builder("Users")
497 .set("UserId")
498 .to(1)
499 .build();
500 let by_value = Mutation::new_insert_builder("Users")
501 .set("UserId")
502 .to(1.to_value())
503 .build();
504 assert_eq!(by_ref, by_value);
505 }
506
507 #[test]
508 #[allow(clippy::needless_borrows_for_generic_args)]
509 fn test_value_binder_borrowed_types() {
510 let id_string = String::from("user-123");
511 let age = 42i64;
512 let active = true;
513
514 let mutation = Mutation::new_insert_builder("Users")
515 .set("UserId")
516 .to(&id_string)
517 .set("Age")
518 .to(&age)
519 .set("Active")
520 .to(&active)
521 .set("Role")
522 .to(&"admin")
523 .build();
524
525 match mutation.inner {
526 InternalMutation::Insert(write) => {
527 assert_eq!(write.values[0].as_string(), "user-123");
528 assert_eq!(write.values[1].as_string(), "42");
529 assert!(write.values[2].as_bool());
530 assert_eq!(write.values[3].as_string(), "admin");
531 }
532 _ => panic!("Expected Insert mutation"),
533 }
534 }
535
536 #[test]
537 fn value_binder_to_turbofished_ref_none() {
538 let mut_ref_turbofished = Mutation::new_insert_builder("Users")
539 .set("Age")
540 .to::<&Option<i64>>(&None)
541 .build();
542 let mut_owned_none = Mutation::new_insert_builder("Users")
543 .set("Age")
544 .to::<Option<i64>>(None)
545 .build();
546 assert_eq!(mut_ref_turbofished, mut_owned_none);
547 }
548
549 #[test]
550 fn value_binder_to_untyped_null() {
551 let mut_null = Mutation::new_insert_builder("Users")
552 .set("Age")
553 .to(Value::null())
554 .build();
555 let mut_unit_none = Mutation::new_insert_builder("Users")
556 .set("Age")
557 .to(None::<()>)
558 .build();
559 let mut_value_none = Mutation::new_insert_builder("Users")
560 .set("Age")
561 .to(None::<Value>)
562 .build();
563 assert_eq!(mut_null, mut_unit_none);
564 assert_eq!(mut_null, mut_value_none);
565 }
566
567 #[test]
568 fn insert_builder() {
569 let mutation = Mutation::new_insert_builder("Users")
570 .set("UserId")
571 .to(1)
572 .set("UserName")
573 .to("Alice")
574 .build();
575
576 match mutation.inner {
577 InternalMutation::Insert(write) => {
578 assert_eq!(write.table, "Users");
579 assert_eq!(write.columns, vec!["UserId", "UserName"]);
580 assert_eq!(write.values.len(), 2);
581 assert_eq!(write.values[0].as_string(), "1");
582 assert_eq!(write.values[1].as_string(), "Alice");
583 }
584 _ => panic!("Expected Insert mutation"),
585 }
586 }
587
588 #[test]
589 fn update_builder() {
590 let mutation = Mutation::new_update_builder("Users")
591 .set("UserId")
592 .to(1)
593 .build();
594
595 match mutation.inner {
596 InternalMutation::Update(write) => {
597 assert_eq!(write.table, "Users");
598 assert_eq!(write.columns, vec!["UserId"]);
599 assert_eq!(write.values.len(), 1);
600 assert_eq!(write.values[0].as_string(), "1");
601 }
602 _ => panic!("Expected Update mutation"),
603 }
604 }
605
606 #[test]
607 fn insert_or_update_builder() {
608 let mutation = Mutation::new_insert_or_update_builder("Users")
609 .set("UserId")
610 .to(1)
611 .build();
612
613 match mutation.inner {
614 InternalMutation::InsertOrUpdate(write) => {
615 assert_eq!(write.table, "Users");
616 assert_eq!(write.columns, vec!["UserId"]);
617 assert_eq!(write.values.len(), 1);
618 assert_eq!(write.values[0].as_string(), "1");
619 }
620 _ => panic!("Expected InsertOrUpdate mutation"),
621 }
622 }
623
624 #[test]
625 fn replace_builder() {
626 let mutation = Mutation::new_replace_builder("Users")
627 .set("UserId")
628 .to(1)
629 .build();
630
631 match mutation.inner {
632 InternalMutation::Replace(write) => {
633 assert_eq!(write.table, "Users");
634 assert_eq!(write.columns, vec!["UserId"]);
635 assert_eq!(write.values.len(), 1);
636 assert_eq!(write.values[0].as_string(), "1");
637 }
638 _ => panic!("Expected Replace mutation"),
639 }
640 }
641
642 #[test]
643 fn build_proto_insert() {
644 let mutation = Mutation::new_insert_builder("Users")
645 .set("UserId")
646 .to(1)
647 .set("UserName")
648 .to("Alice")
649 .build();
650 let proto = mutation.build_proto();
651 match proto.operation {
652 Some(Operation::Insert(write)) => {
653 assert_eq!(write.table, "Users");
654 assert_eq!(write.columns, vec!["UserId", "UserName"]);
655 assert_eq!(write.values.len(), 1);
656 assert_eq!(write.values[0].len(), 2);
657 assert_eq!(write.values[0][0], serde_json::json!("1"));
658 assert_eq!(write.values[0][1], serde_json::json!("Alice"));
659 }
660 _ => panic!("Expected Insert operation, got {:?}", proto.operation),
661 }
662 }
663
664 #[test]
665 fn build_proto_update() {
666 let mutation = Mutation::new_update_builder("Users")
667 .set("UserId")
668 .to(1)
669 .build();
670 let proto = mutation.build_proto();
671 match proto.operation {
672 Some(Operation::Update(write)) => {
673 assert_eq!(write.table, "Users");
674 assert_eq!(write.columns, vec!["UserId"]);
675 assert_eq!(write.values.len(), 1);
676 }
677 _ => panic!("Expected Update operation, got {:?}", proto.operation),
678 }
679 }
680
681 #[test]
682 fn build_proto_insert_or_update() {
683 let mutation = Mutation::new_insert_or_update_builder("Users")
684 .set("UserId")
685 .to(1)
686 .build();
687 let proto = mutation.build_proto();
688 match proto.operation {
689 Some(Operation::InsertOrUpdate(write)) => {
690 assert_eq!(write.table, "Users");
691 assert_eq!(write.columns, vec!["UserId"]);
692 assert_eq!(write.values.len(), 1);
693 }
694 _ => panic!(
695 "Expected InsertOrUpdate operation, got {:?}",
696 proto.operation
697 ),
698 }
699 }
700
701 #[test]
702 fn build_proto_replace() {
703 let mutation = Mutation::new_replace_builder("Users")
704 .set("UserId")
705 .to(1)
706 .build();
707 let proto = mutation.build_proto();
708 match proto.operation {
709 Some(Operation::Replace(write)) => {
710 assert_eq!(write.table, "Users");
711 assert_eq!(write.columns, vec!["UserId"]);
712 assert_eq!(write.values.len(), 1);
713 }
714 _ => panic!("Expected Replace operation, got {:?}", proto.operation),
715 }
716 }
717
718 #[test]
719 fn build_proto_delete() {
720 let key_set = crate::key::KeySet::builder().build();
721 let mutation = Mutation::delete("Users", key_set);
722 let proto = mutation.build_proto();
723 match proto.operation {
724 Some(Operation::Delete(delete)) => {
725 assert_eq!(delete.table, "Users");
726 }
727 _ => panic!("Expected Delete operation, got {:?}", proto.operation),
728 }
729 }
730
731 #[test]
732 fn test_select_mutation_key_empty() {
733 let mutations = vec![];
734 let key = Mutation::select_mutation_key(&mutations);
735 assert!(key.is_none());
736 }
737
738 #[test]
739 fn test_select_mutation_key_prefers_insert_or_update_over_insert() {
740 let m1 = Mutation::new_insert_builder("Users")
741 .set("UserId")
742 .to(1)
743 .build()
744 .build_proto();
745 let m2 = Mutation::new_insert_or_update_builder("Users")
746 .set("UserId")
747 .to(2)
748 .build()
749 .build_proto();
750 let mutations = vec![m1.clone(), m2.clone()];
751 let key = Mutation::select_mutation_key(&mutations);
752 assert_eq!(key, Some(m2));
753 }
754
755 #[test]
756 fn test_select_mutation_key_only_insert_prefers_largest() {
757 let m1 = Mutation::new_insert_builder("Users")
758 .set("UserId")
759 .to(1)
760 .build()
761 .build_proto();
762
763 let row1 = vec![serde_json::json!("2")]
765 .into_iter()
766 .collect::<wkt::ListValue>();
767 let row2 = vec![serde_json::json!("3")]
768 .into_iter()
769 .collect::<wkt::ListValue>();
770 let write2 = crate::model::mutation::Write::new()
771 .set_table("Users")
772 .set_columns(vec!["UserId".to_string()])
773 .set_values(vec![row1, row2]);
774 let m2 = crate::model::Mutation::new().set_insert(write2);
775
776 let mutations = vec![m1.clone(), m2.clone()];
777 let key = Mutation::select_mutation_key(&mutations);
778 assert_eq!(key, Some(m2));
779 }
780
781 #[test]
782 fn test_select_mutation_key_mix() {
783 let m1 = Mutation::new_insert_builder("Users")
784 .set("UserId")
785 .to(1)
786 .build()
787 .build_proto();
788 let m2 = Mutation::new_update_builder("Users")
789 .set("UserId")
790 .to(2)
791 .build()
792 .build_proto();
793 let m3 = Mutation::new_insert_or_update_builder("Users")
794 .set("UserId")
795 .to(3)
796 .build()
797 .build_proto();
798 let mutations = vec![m1.clone(), m2.clone(), m3.clone()];
799 let key = Mutation::select_mutation_key(&mutations).expect("Expected a key");
800 assert!(
802 key == m2 || key == m3,
803 "Expected either m2 or m3 to be selected, got {:?}",
804 key
805 );
806 }
807
808 #[test]
809 fn test_select_mutation_key_only_non_insert() {
810 let m1 = Mutation::new_update_builder("Users")
811 .set("UserId")
812 .to(1)
813 .build()
814 .build_proto();
815 let m2 = Mutation::new_replace_builder("Users")
816 .set("UserId")
817 .to(2)
818 .build()
819 .build_proto();
820 let mutations = vec![m1.clone(), m2.clone()];
821 let key = Mutation::select_mutation_key(&mutations).expect("Expected a key");
822 assert!(
824 key == m1 || key == m2,
825 "Expected either m1 or m2 to be selected, got {:?}",
826 key
827 );
828 }
829
830 #[test]
831 fn test_select_mutation_key_operation_none() {
832 let m1 = crate::model::Mutation::default();
833 let m2 = crate::model::Mutation::default();
834 let mutations = vec![m1.clone(), m2.clone()];
835 let key = Mutation::select_mutation_key(&mutations);
836 assert_eq!(key, Some(m1));
837 }
838
839 #[test]
840 fn test_mutation_macro_insert() {
841 let macro_mutation = mutation!(insert "Singers" {
842 SingerId: 1_i64,
843 FirstName: "Marc",
844 LastName: "Richards",
845 });
846
847 let builder_mutation = Mutation::new_insert_builder("Singers")
848 .set("SingerId")
849 .to(1_i64)
850 .set("FirstName")
851 .to("Marc")
852 .set("LastName")
853 .to("Richards")
854 .build();
855
856 assert_eq!(macro_mutation, builder_mutation);
857 }
858
859 #[test]
860 fn test_mutation_macro_update_with_string_literal() {
861 let macro_mutation = mutation!(update "Albums" {
862 SingerId: 1_i64,
863 "AlbumTitle": "New Title",
864 MarketingBudget: 100_000_i64,
865 });
866
867 let builder_mutation = Mutation::new_update_builder("Albums")
868 .set("SingerId")
869 .to(1_i64)
870 .set("AlbumTitle")
871 .to("New Title")
872 .set("MarketingBudget")
873 .to(100_000_i64)
874 .build();
875
876 assert_eq!(macro_mutation, builder_mutation);
877 }
878
879 #[test]
880 fn test_mutation_macro_insert_or_update() {
881 let macro_mutation = mutation!(insert_or_update "Singers" {
882 SingerId: 2_i64,
883 FirstName: "Alice",
884 });
885
886 let builder_mutation = Mutation::new_insert_or_update_builder("Singers")
887 .set("SingerId")
888 .to(2_i64)
889 .set("FirstName")
890 .to("Alice")
891 .build();
892
893 assert_eq!(macro_mutation, builder_mutation);
894 }
895
896 #[test]
897 fn test_mutation_macro_replace() {
898 let macro_mutation = mutation!(replace "Singers" {
899 SingerId: 3_i64,
900 FirstName: "Bob",
901 });
902
903 let builder_mutation = Mutation::new_replace_builder("Singers")
904 .set("SingerId")
905 .to(3_i64)
906 .set("FirstName")
907 .to("Bob")
908 .build();
909
910 assert_eq!(macro_mutation, builder_mutation);
911 }
912
913 #[test]
914 fn test_mutation_macro_delete() {
915 let macro_mutation = mutation!(delete "Singers", KeySet::all());
916 let direct_mutation = Mutation::delete("Singers", KeySet::all());
917
918 assert_eq!(macro_mutation, direct_mutation);
919 }
920
921 #[test]
922 fn test_mutation_macro_complex_table_expressions() {
923 fn get_table() -> &'static str {
924 "Singers"
925 }
926 mod constants {
927 pub const TABLE: &str = "Singers";
928 }
929
930 let macro_mutation1 = mutation!(insert (get_table()) {
931 SingerId: 1_i64,
932 });
933 let macro_mutation2 = mutation!(insert (constants::TABLE) {
934 SingerId: 1_i64,
935 });
936 let delete_mutation = mutation!(delete constants::TABLE, KeySet::all());
937
938 let expected_insert = Mutation::new_insert_builder("Singers")
939 .set("SingerId")
940 .to(1_i64)
941 .build();
942 let expected_delete = Mutation::delete("Singers", KeySet::all());
943
944 assert_eq!(macro_mutation1, expected_insert);
945 assert_eq!(macro_mutation2, expected_insert);
946 assert_eq!(delete_mutation, expected_delete);
947 }
948}