|
| 1 | +//! Write-event tag correctness: the `WriteOp` emitted to the Event Plane |
| 2 | +//! must reflect whether the mutation created a new row or replaced an |
| 3 | +//! existing one. Downstream consumers (AFTER triggers, CDC) branch on |
| 4 | +//! `Insert` vs `Update`; a wrong tag silently delivers the wrong event. |
| 5 | +//! |
| 6 | +//! These tests drive real SQL through the harness and observe which |
| 7 | +//! AFTER trigger fires. Sync triggers are used so the audit table is |
| 8 | +//! populated by the time the outer statement returns. |
| 9 | +
|
| 10 | +mod common; |
| 11 | + |
| 12 | +use common::pgwire_harness::TestServer; |
| 13 | + |
| 14 | +/// UPSERT onto an existing primary key must emit `WriteOp::Update`. |
| 15 | +/// An `AFTER UPDATE` trigger must fire; an `AFTER INSERT` trigger must not. |
| 16 | +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] |
| 17 | +async fn upsert_on_existing_row_emits_update() { |
| 18 | + let server = TestServer::start().await; |
| 19 | + |
| 20 | + server.exec("CREATE COLLECTION src").await.unwrap(); |
| 21 | + server.exec("CREATE COLLECTION ins_log").await.unwrap(); |
| 22 | + server.exec("CREATE COLLECTION upd_log").await.unwrap(); |
| 23 | + |
| 24 | + server |
| 25 | + .exec( |
| 26 | + "CREATE SYNC TRIGGER on_ins AFTER INSERT ON src FOR EACH ROW \ |
| 27 | + BEGIN INSERT INTO ins_log (id, src_id) VALUES (NEW.id || '_i', NEW.id); END;", |
| 28 | + ) |
| 29 | + .await |
| 30 | + .unwrap(); |
| 31 | + server |
| 32 | + .exec( |
| 33 | + "CREATE SYNC TRIGGER on_upd AFTER UPDATE ON src FOR EACH ROW \ |
| 34 | + BEGIN INSERT INTO upd_log (id, src_id) VALUES (NEW.id || '_u', NEW.id); END;", |
| 35 | + ) |
| 36 | + .await |
| 37 | + .unwrap(); |
| 38 | + |
| 39 | + // Seed the row (fires AFTER INSERT — expected). |
| 40 | + server |
| 41 | + .exec("UPSERT INTO src (id, v) VALUES ('a', 1)") |
| 42 | + .await |
| 43 | + .unwrap(); |
| 44 | + // Overwrite the row (must fire AFTER UPDATE, not AFTER INSERT). |
| 45 | + server |
| 46 | + .exec("UPSERT INTO src (id, v) VALUES ('a', 2)") |
| 47 | + .await |
| 48 | + .unwrap(); |
| 49 | + |
| 50 | + let ins_rows = server |
| 51 | + .query_text("SELECT id FROM ins_log ORDER BY id") |
| 52 | + .await |
| 53 | + .unwrap(); |
| 54 | + let upd_rows = server |
| 55 | + .query_text("SELECT id FROM upd_log ORDER BY id") |
| 56 | + .await |
| 57 | + .unwrap(); |
| 58 | + |
| 59 | + assert_eq!( |
| 60 | + ins_rows.len(), |
| 61 | + 1, |
| 62 | + "AFTER INSERT must fire once (seed only); got: {ins_rows:?}" |
| 63 | + ); |
| 64 | + assert_eq!( |
| 65 | + upd_rows.len(), |
| 66 | + 1, |
| 67 | + "AFTER UPDATE must fire on the overwrite; got: {upd_rows:?}" |
| 68 | + ); |
| 69 | +} |
| 70 | + |
| 71 | +/// UPSERT onto a fresh primary key must emit `WriteOp::Insert`. |
| 72 | +/// An `AFTER INSERT` trigger must fire. |
| 73 | +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] |
| 74 | +async fn upsert_on_new_row_emits_insert() { |
| 75 | + let server = TestServer::start().await; |
| 76 | + |
| 77 | + server.exec("CREATE COLLECTION src").await.unwrap(); |
| 78 | + server.exec("CREATE COLLECTION ins_log").await.unwrap(); |
| 79 | + |
| 80 | + server |
| 81 | + .exec( |
| 82 | + "CREATE SYNC TRIGGER on_ins AFTER INSERT ON src FOR EACH ROW \ |
| 83 | + BEGIN INSERT INTO ins_log (id, src_id) VALUES (NEW.id || '_i', NEW.id); END;", |
| 84 | + ) |
| 85 | + .await |
| 86 | + .unwrap(); |
| 87 | + |
| 88 | + server |
| 89 | + .exec("UPSERT INTO src (id, v) VALUES ('a', 1)") |
| 90 | + .await |
| 91 | + .unwrap(); |
| 92 | + |
| 93 | + let rows = server.query_text("SELECT id FROM ins_log").await.unwrap(); |
| 94 | + assert_eq!(rows.len(), 1, "AFTER INSERT must fire; got: {rows:?}"); |
| 95 | +} |
| 96 | + |
| 97 | +/// `INSERT ... ON CONFLICT (id) DO UPDATE SET ...` onto an existing row |
| 98 | +/// must emit `WriteOp::Update`. AFTER UPDATE fires, AFTER INSERT does not. |
| 99 | +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] |
| 100 | +async fn insert_on_conflict_do_update_emits_update_on_overwrite() { |
| 101 | + let server = TestServer::start().await; |
| 102 | + |
| 103 | + server.exec("CREATE COLLECTION src").await.unwrap(); |
| 104 | + server.exec("CREATE COLLECTION ins_log").await.unwrap(); |
| 105 | + server.exec("CREATE COLLECTION upd_log").await.unwrap(); |
| 106 | + |
| 107 | + server |
| 108 | + .exec( |
| 109 | + "CREATE SYNC TRIGGER on_ins AFTER INSERT ON src FOR EACH ROW \ |
| 110 | + BEGIN INSERT INTO ins_log (id, src_id) VALUES (NEW.id || '_i', NEW.id); END;", |
| 111 | + ) |
| 112 | + .await |
| 113 | + .unwrap(); |
| 114 | + server |
| 115 | + .exec( |
| 116 | + "CREATE SYNC TRIGGER on_upd AFTER UPDATE ON src FOR EACH ROW \ |
| 117 | + BEGIN INSERT INTO upd_log (id, src_id) VALUES (NEW.id || '_u', NEW.id); END;", |
| 118 | + ) |
| 119 | + .await |
| 120 | + .unwrap(); |
| 121 | + |
| 122 | + server |
| 123 | + .exec("INSERT INTO src (id, v) VALUES ('a', 1)") |
| 124 | + .await |
| 125 | + .unwrap(); |
| 126 | + server |
| 127 | + .exec( |
| 128 | + "INSERT INTO src (id, v) VALUES ('a', 2) \ |
| 129 | + ON CONFLICT (id) DO UPDATE SET v = EXCLUDED.v", |
| 130 | + ) |
| 131 | + .await |
| 132 | + .unwrap(); |
| 133 | + |
| 134 | + let ins_rows = server |
| 135 | + .query_text("SELECT id FROM ins_log ORDER BY id") |
| 136 | + .await |
| 137 | + .unwrap(); |
| 138 | + let upd_rows = server |
| 139 | + .query_text("SELECT id FROM upd_log ORDER BY id") |
| 140 | + .await |
| 141 | + .unwrap(); |
| 142 | + |
| 143 | + assert_eq!( |
| 144 | + ins_rows.len(), |
| 145 | + 1, |
| 146 | + "AFTER INSERT must fire for the first (new-row) insert only; got: {ins_rows:?}" |
| 147 | + ); |
| 148 | + assert_eq!( |
| 149 | + upd_rows.len(), |
| 150 | + 1, |
| 151 | + "AFTER UPDATE must fire on the conflict path; got: {upd_rows:?}" |
| 152 | + ); |
| 153 | +} |
| 154 | + |
| 155 | +/// `INSERT ... ON CONFLICT DO NOTHING` on an existing row must emit no |
| 156 | +/// event — the write is a silent no-op, so neither AFTER INSERT nor |
| 157 | +/// AFTER UPDATE should fire. Regression guard: the bug class includes |
| 158 | +/// both "wrong tag" and "emit when no mutation occurred". |
| 159 | +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] |
| 160 | +async fn insert_on_conflict_do_nothing_emits_no_event() { |
| 161 | + let server = TestServer::start().await; |
| 162 | + |
| 163 | + server.exec("CREATE COLLECTION src").await.unwrap(); |
| 164 | + server.exec("CREATE COLLECTION ins_log").await.unwrap(); |
| 165 | + server.exec("CREATE COLLECTION upd_log").await.unwrap(); |
| 166 | + |
| 167 | + server |
| 168 | + .exec( |
| 169 | + "CREATE SYNC TRIGGER on_ins AFTER INSERT ON src FOR EACH ROW \ |
| 170 | + BEGIN INSERT INTO ins_log (id, src_id) VALUES (NEW.id || '_i', NEW.id); END;", |
| 171 | + ) |
| 172 | + .await |
| 173 | + .unwrap(); |
| 174 | + server |
| 175 | + .exec( |
| 176 | + "CREATE SYNC TRIGGER on_upd AFTER UPDATE ON src FOR EACH ROW \ |
| 177 | + BEGIN INSERT INTO upd_log (id, src_id) VALUES (NEW.id || '_u', NEW.id); END;", |
| 178 | + ) |
| 179 | + .await |
| 180 | + .unwrap(); |
| 181 | + |
| 182 | + server |
| 183 | + .exec("INSERT INTO src (id, v) VALUES ('a', 1)") |
| 184 | + .await |
| 185 | + .unwrap(); |
| 186 | + server |
| 187 | + .exec("INSERT INTO src (id, v) VALUES ('a', 2) ON CONFLICT DO NOTHING") |
| 188 | + .await |
| 189 | + .unwrap(); |
| 190 | + |
| 191 | + let ins_rows = server |
| 192 | + .query_text("SELECT id FROM ins_log ORDER BY id") |
| 193 | + .await |
| 194 | + .unwrap(); |
| 195 | + let upd_rows = server |
| 196 | + .query_text("SELECT id FROM upd_log ORDER BY id") |
| 197 | + .await |
| 198 | + .unwrap(); |
| 199 | + |
| 200 | + assert_eq!( |
| 201 | + ins_rows.len(), |
| 202 | + 1, |
| 203 | + "only the first insert should fire AFTER INSERT; got: {ins_rows:?}" |
| 204 | + ); |
| 205 | + assert!( |
| 206 | + upd_rows.is_empty(), |
| 207 | + "DO NOTHING must not fire AFTER UPDATE; got: {upd_rows:?}" |
| 208 | + ); |
| 209 | +} |
0 commit comments