The sidecar CLI (crates/postjevsql-sidecar, contract Deployment,
"Launching it") against a throwaway target and sidecar: converge
declares exactly the configured server, sync plans without changing
anything, sync --apply changes only what differs, and a change a
view depends on is refused with nothing applied.
10use support::postgres::{Instance, TestPostgres}; 11use tokio_postgres::Client; 12 13async fn connect(pg: &Instance) -> Client { 14 let (client, conn) = tokio_postgres::connect(&pg.conn_str(), tokio_postgres::NoTls).await.expect("connect"); 15 tokio::spawn(conn); 16 client.batch_execute("SET statement_timeout = '60s'").await.unwrap(); 17 client 18} 19 20fn host(pg: &Instance) -> String { 21 let conn = pg.conn_str(); 22 conn.split_whitespace().find_map(|kv| kv.strip_prefix("host=")).expect("a host").to_owned() 23} 24 25struct Setup { 26 _target_pg: Instance, 27 _sidecar_pg: Instance, 28 target: Client, 29 sidecar: Client, 30 sidecar_conn: String, 31 dir: tempfile::TempDir, 32} 33 34async fn setup() -> Setup { 35 // Both slots at once (`support/postgres.rs`). `converge` creates the 36 // extension, so the sidecar serves it. 37 let ext = std::env::var("POSTJEVSQL_EXT").expect("POSTJEVSQL_EXT"); 38 let sidecar = TestPostgres::new().with_extension(Path::new(&ext)); 39 let [target_pg, sidecar_pg] = TestPostgres::start_all([TestPostgres::new(), sidecar]).await; 40 let target = connect(&target_pg).await; 41 target 42 .batch_execute( 43 "CREATE TABLE people (id int NOT NULL, old text, name text); 44 CREATE TABLE same (id int); 45 INSERT INTO people VALUES (1, 'x', 'Ada');", 46 ) 47 .await 48 .unwrap(); 49 let dir = tempfile::tempdir().unwrap(); 50 std::fs::write(dir.path().join("key"), "not-a-real-key\n").unwrap(); 51 std::fs::write(dir.path().join("pw"), "unused-under-trust\n").unwrap(); 52 std::fs::write( 53 dir.path().join("config.toml"), 54 format!( 55 "[target]\nhost = \"{}\"\ndbname = \"postgres\"\nsslmode = \"disable\"\nschemas = [\"public\"]\n\ 56 [[mapping]]\nlocal_user = \"postgres\"\nremote_user = \"postgres\"\npassword_file = \"{}\"\n\ 57 [jev]\napi_key_file = \"{}\"\n", 58 host(&target_pg), 59 dir.path().join("pw").display(), 60 dir.path().join("key").display() 61 ), 62 ) 63 .unwrap(); 64 let sidecar = connect(&sidecar_pg).await; 65 Setup { sidecar_conn: sidecar_pg.conn_str(), _target_pg: target_pg, _sidecar_pg: sidecar_pg, target, sidecar, dir } 66} 67 68fn cli(s: &Setup, config: &Path, args: &[&str], env: &[(&str, &str)]) -> Output { 69 Command::new(std::env::var("POSTJEVSQL_SIDECAR_CLI").expect("POSTJEVSQL_SIDECAR_CLI")) 70 .env_clear() 71 .envs(env.iter().copied()) 72 .arg("--sidecar") 73 .arg(&s.sidecar_conn) 74 .arg(config) 75 .args(args) 76 .output() 77 .expect("the CLI runs") 78} 79 80fn run(s: &Setup, args: &[&str]) -> String { 81 let out = cli(s, &s.dir.path().join("config.toml"), args, &[]); 82 let stdout = String::from_utf8_lossy(&out.stdout).into_owned(); 83 assert!(out.status.success(), "{args:?} failed:\n{stdout}\n{}", String::from_utf8_lossy(&out.stderr)); 84 stdout 85}
schema.table → its columns, in attnum order, and its oid.
88async fn catalog(c: &Client) -> Vec<(String, u32, Vec<String>)> { 89 c.query( 90 "SELECT r.relname::text, r.oid, 91 array_agg(a.attname::text ORDER BY a.attnum) 92 FROM pg_foreign_table f JOIN pg_class r ON r.oid = f.ftrelid 93 JOIN pg_attribute a ON a.attrelid = r.oid AND a.attnum > 0 AND NOT a.attisdropped 94 GROUP BY 1, 2 ORDER BY 1", 95 &[], 96 ) 97 .await 98 .unwrap() 99 .iter() 100 .map(|r| (r.get(0), r.get(1), r.get(2))) 101 .collect() 102}
104async fn server_options(c: &Client) -> Vec<String> { 105 c.query_one("SELECT srvoptions FROM pg_foreign_server WHERE srvname = 'target'", &[]) 106 .await 107 .unwrap() 108 .get(0) 109} 110 111#[tokio::test(flavor = "multi_thread")] 112async fn converge_declares_exactly_the_configured_options() { 113 let s = setup().await; 114 run(&s, &["converge"]); 115 s.sidecar.batch_execute("ALTER SERVER target OPTIONS (ADD extensions 'postjevsql')").await.unwrap(); 116 run(&s, &["converge"]); 117 let options = server_options(&s.sidecar).await; 118 assert!(!options.iter().any(|o| o.starts_with("extensions=")), "{options:?}"); 119 assert!(options.contains(&"sslmode=disable".to_owned()), "{options:?}"); 120 // A second run is a no-op. 121 run(&s, &["converge"]); 122 assert_eq!(server_options(&s.sidecar).await, options); 123} 124 125#[tokio::test(flavor = "multi_thread")] 126async fn sync_plans_then_applies_only_what_differs() { 127 let s = setup().await; 128 run(&s, &["converge"]); 129 run(&s, &["sync", "--apply"]); 130 let before = catalog(&s.sidecar).await; 131 assert_eq!(before.len(), 2, "{before:?}"); 132 let rows = s.sidecar.query("SELECT name FROM people", &[]).await.unwrap(); 133 assert_eq!(rows[0].get::<_, String>(0), "Ada"); 134 135 s.target 136 .batch_execute("ALTER TABLE people DROP COLUMN old, ADD COLUMN email text; CREATE TABLE fresh (id bigint)") 137 .await 138 .unwrap(); 139 140 // The plan is printed and changes nothing. 141 let plan = run(&s, &["sync"]); 142 assert!(plan.contains("DROP COLUMN \"old\""), "{plan}"); 143 assert!(plan.contains("ADD COLUMN \"email\""), "{plan}"); 144 assert!(plan.contains("CREATE FOREIGN TABLE \"public\".\"fresh\""), "{plan}"); 145 assert!(!plan.contains("\"same\""), "{plan}"); 146 assert_eq!(catalog(&s.sidecar).await, before); 147 148 run(&s, &["sync", "--apply"]); 149 let after = catalog(&s.sidecar).await; 150 let find = |name: &str| after.iter().find(|t| t.0 == name).cloned().expect(name); 151 assert_eq!(find("people").2, ["id", "name", "email"]); 152 assert_eq!(find("fresh").2, ["id"]); 153 // Changed in place, never dropped and re-imported. 154 assert_eq!(find("people").1, before[0].1); 155 assert_eq!(find("same").1, before[1].1); 156 assert!(run(&s, &["sync"]).contains("nothing to change")); 157} 158 159#[tokio::test(flavor = "multi_thread")] 160async fn a_dependent_view_refuses_the_whole_sync() { 161 let s = setup().await; 162 run(&s, &["converge"]); 163 run(&s, &["sync", "--apply"]); 164 s.sidecar.batch_execute("CREATE VIEW olds AS SELECT old FROM people").await.unwrap(); 165 s.target.batch_execute("ALTER TABLE people DROP COLUMN old; CREATE TABLE fresh (id int)").await.unwrap(); 166 let before = catalog(&s.sidecar).await; 167 168 let out = cli(&s, &s.dir.path().join("config.toml"), &["sync", "--apply"], &[]); 169 let stderr = String::from_utf8_lossy(&out.stderr); 170 assert!(!out.status.success()); 171 assert!(stderr.contains("view olds"), "{stderr}"); 172 assert!(stderr.contains("DROP COLUMN \"old\""), "{stderr}"); 173 // Nothing applied, the new table included. 174 assert_eq!(catalog(&s.sidecar).await, before); 175} 176 177#[tokio::test(flavor = "multi_thread")] 178async fn a_secret_set_twice_or_nowhere_is_refused() { 179 let s = setup().await; 180 let config = s.dir.path().join("config.toml"); 181 let key = s.dir.path().join("key"); 182 183 let out = cli(&s, &config, &["converge"], &[("TYPESAFE_API_KEY", "sk-secret-value")]); 184 let stderr = String::from_utf8_lossy(&out.stderr); 185 assert!(!out.status.success()); 186 assert!(stderr.contains("set twice"), "{stderr}"); 187 assert!(stderr.contains(&key.display().to_string()) && stderr.contains("$TYPESAFE_API_KEY"), "{stderr}"); 188 assert!(!stderr.contains("sk-secret-value") && !stderr.contains("not-a-real-key"), "{stderr}"); 189 190 let text = std::fs::read_to_string(&config).unwrap(); 191 let neither = s.dir.path().join("neither.toml"); 192 std::fs::write(&neither, text.split("[jev]").next().unwrap()).unwrap(); 193 let out = cli(&s, &neither, &["converge"], &[]); 194 let stderr = String::from_utf8_lossy(&out.stderr); 195 assert!(!out.status.success()); 196 assert!(stderr.contains("not set") && stderr.contains("$TYPESAFE_API_KEY"), "{stderr}"); 197 // Refused before converging anything. 198 let servers: i64 = 199 s.sidecar.query_one("SELECT count(*) FROM pg_foreign_server", &[]).await.map(|r| r.get(0)).unwrap_or(0); 200 assert_eq!(servers, 0); 201}