A pure MoonBit PostgreSQL client, wire protocol from scratch.
Dependencies
moon add jaredzhou/moonpglet conn = @moonpg.connect("postgres://user:pw@localhost:5432/db")
// Execute DDL / DML
conn.execute("CREATE TABLE users (id SERIAL PRIMARY KEY, name TEXT, email TEXT)") |> ignore
conn.execute(
"INSERT INTO users (name, email) VALUES ($1, $2)",
params=["alice", "a@b.com"],
) |> ignore
// fetch — typed array of rows (auto-closes)
let names : Array[String] = &QueryExecutor::fetch(conn, "SELECT name FROM users ORDER BY id")
let count : Int = &QueryExecutor::fetch_one(conn, "SELECT COUNT(*) FROM users")
// fetch with tuples
let users : Array[(Int, String)] = &QueryExecutor::fetch(conn, "SELECT id, name FROM users")
let (id, name) = &QueryExecutor::fetch_one(conn, "SELECT id, name FROM users WHERE id = $1", params=[1])
// query — manual iteration (MUST close rows in try-catch)
let rows = conn.query("SELECT id, name FROM users")
try {
while rows.has_next() {
let row = rows.get_row()
let id : Int = row.get(0)
let name : String = row.get_by_name("name")
}
rows.close()
} catch {
e => { rows.close(); raise e }
}let conn = @moonpg.connect(conninfo)
// Scalar
let count : Int = &QueryExecutor::fetch_one(conn, "SELECT COUNT(*) FROM users")
// Nullable
let email : String? = &QueryExecutor::fetch_one(conn, "SELECT email FROM users WHERE id = $1", params=[1])
// Tuples — no struct needed
let pairs : Array[(Int, String)] = &QueryExecutor::fetch(conn, "SELECT id, name FROM users")
let (id, name, email) : (Int, String, String) = &QueryExecutor::fetch_one(conn,
"SELECT id, name, email FROM users WHERE id = $1", params=[1],
)
// Custom FromRow struct
impl FromRow for User with fn from_row(r : Row) -> User raise PgError {
User::{ id: r.get(0), name: r.get(1), email: r.get(2) }
}
let users : Array[User] = &QueryExecutor::fetch(conn, "SELECT id, name, email FROM users")
// Or cast to &QueryExecutor for dot-syntax:
let q : &QueryExecutor = conn
let count : Int = q.fetch_one("SELECT COUNT(*) FROM users")
let names : Array[String] = q.fetch("SELECT name FROM users")// query — iterate rows manually
let rows = conn.query("SELECT id, name FROM users")
try {
while rows.has_next() {
let row = rows.get_row()
let id : Int = row.get(0)
let name : String = row.get_by_name("name")
println("\{id}: \{name}")
}
rows.close()
} catch {
e => { rows.close(); raise e }
}
// query_one — single row
let row = conn.query_one("SELECT id FROM users WHERE name = $1", params=["alice"])
let id : Int = row.get(0)conn.execute(
"UPDATE users SET email = $1 WHERE id = $2",
params=[null, 42], // null → SQL NULL
) |> ignorelet pool = Pool::new(PoolConfig::new(
"postgres://user:pw@localhost:5432/db",
max_conns=10,
min_idle=2,
))
// Auto-acquire + auto-release — fetch/rows.close() returns conn to pool
let names : Array[String] = &QueryExecutor::fetch(pool, "SELECT name FROM users")
pool.execute("INSERT INTO users (name) VALUES ($1)", params=["bob"]) |> ignore
// Manual query with try-catch
let rows = pool.query("SELECT id FROM users")
try {
while rows.has_next() { ... }
rows.close()
} catch {
e => { rows.close(); raise e }
}
// Explicit acquire
let pc = pool.acquire()
pc.execute("DELETE FROM users WHERE id = $1", params=[1]) |> ignore
pc.release()
// Inspect pool
let stats = pool.stats()
println("active=\{stats.active_connections} idle=\{stats.idle_connections}")
// Health check + background maintenance
let pool2 = Pool::new(PoolConfig::new(
conninfo,
max_conns=10,
min_idle=2,
max_idle_sec=300,
max_lifetime_sec=3600,
health_check=true,
maintenance_interval_sec=60,
))
pool2.start_maintenance()let conn = @moonpg.connect(conninfo)
// begin_func — auto-commit on success, auto-rollback on error
let result = begin_func(conn, async fn(tx) {
tx.execute("INSERT INTO users (name) VALUES ($1)", params=["alice"]) |> ignore
&QueryExecutor::fetch_one(tx, "SELECT id FROM users WHERE name = $1", params=["alice"])
})
// Manual transaction — try { commit } catch { rollback }
let tx = conn.begin_tx()
try {
tx.execute("UPDATE users SET name = $1 WHERE id = $2", params=["bob", 1]) |> ignore
let name : String = &QueryExecutor::fetch_one(tx, "SELECT name FROM users WHERE id = $1", params=[1])
tx.commit()
} catch {
e => { tx.rollback(); raise e }
}
// Pooled transaction — connection auto-returns to pool on commit/rollback
let pool = Pool::new(PoolConfig::new(conninfo, max_conns=4))
let tx2 = pool.begin_tx()
try {
tx2.execute("DELETE FROM users WHERE id = $1", params=[99]) |> ignore
tx2.commit()
} catch {
e => { tx2.rollback(); raise e }
}// ToValue — MoonBit → PostgreSQL
// Built-in impls: Int, Int64, Double, Bool, String, Bytes, Json,
// Timestamp, Decimal, UUID, Option<T>, Array<T>
let params = [42, 3.14, true, "hello", null] // null → SQL NULL
conn.execute("INSERT INTO t (a, b, c, d, e) VALUES ($1, $2, $3, $4, $5)", params=params)
// FromValue — PostgreSQL → MoonBit
let row = conn.query_one("SELECT a, b, c, d, e FROM t")
let a : Int = row.get(0) // strict: raises on NULL
let b : String? = row.get(1) // nullable: NULL → None
let c : Bool = row.get_by_name("c") // by column name
let d : Json = row.get(3) // jsonb → Json
let e : Timestamp = row.get(4) // timestamptz → Unix µs
// Arrays — PostgreSQL array columns
conn.execute(
"INSERT INTO items (tags) VALUES ($1)",
params=[["red", "green", "blue"]],
) |> ignore
let row2 = conn.query_one("SELECT tags FROM items")
let tags : Array[String] = row2.get(0) // strict: no NULL elements
let tags2 : Array[String?] = row2.get(0) // nullable: NULL elements → None
// Custom impl
impl ToValue for MyType with fn to_value(self) -> Value {
Value::String(self.to_json())
}
impl FromValue for MyType with fn from_value(v : Value) -> MyType raise ValueError {
match v { Value::String(s) => MyType::from_json(s); _ => raise ... }
}// Single-column rows: built-in impls for all FromValue types
let count : Int = &QueryExecutor::fetch_one(conn, "SELECT COUNT(*) FROM users")
let name : String? = &QueryExecutor::fetch_one(conn, "SELECT name FROM users WHERE id = $1", params=[1])
// Custom struct
impl FromRow for User with fn from_row(r : Row) -> User raise PgError {
User::{ id: r.get(0), name: r.get(1), email: r.get(2) }
}
let users : Array[User] = &QueryExecutor::fetch(conn, "SELECT id, name, email FROM users")
// Tuple impls (2–10)
let pairs : Array[(Int, String)] = &QueryExecutor::fetch(conn, "SELECT id, name FROM users")
let (id, name, email) : (Int, String, String) = &QueryExecutor::fetch_one(conn,
"SELECT id, name, email FROM users WHERE id = $1", params=[1],
)let pid = conn.backend_pid() // server process ID
let ver = conn.param("server_version") // e.g. Some("16.4")
let tz = conn.param("TimeZone") // e.g. Some("UTC")
if conn.is_closed() { ... }// System CA
@moonpg.connect("postgres://user:pw@host/db?sslmode=require")
// Custom CA
@moonpg.connect("postgres://user:pw@host/db?sslmode=verify-ca&sslrootcert=/etc/ca.pem")
// Client certificate
@moonpg.connect(
"postgres://user:pw@host/db?sslmode=require&sslcert=/etc/certs/client.pem&sslkey=/etc/certs/client.key",
)@async.with_task_group() group => {
let listener = conn.listen("events", group)
group.spawn_bg() () => {
let c2 = @moonpg.connect(conninfo)
c2.notify("events", payload="hello")
}
let notif = listener.recv()
println("\{notif.channel}: \{notif.payload}")
}// Bulk insert from an iterator — one row in memory at a time
conn.copy_in("COPY users (name, age) FROM STDIN", ["alice\t30\n", "bob\t25\n"].iter())
// Streaming COPY writer
let w = conn.begin_copy("users", ["name", "age"])
w.write_row(["alice", 30])
w.write_row(["bob", 25])
let result = w.finish()// TCP connect timeout (seconds)
@moonpg.connect("postgres://host/db?connect_timeout=5")
// Server-side statement timeout (milliseconds)
@moonpg.connect("postgres://host/db?statement_timeout=30000")@moonpg.connect("postgres://host/db?target_session_attrs=read-write")# Default connection
moon test --target native
# Custom connection
PGCONN="postgres://user:pass@localhost:5432/mydb" moon test --target native| env var | example connstr | auth method |
|---|---|---|
| PG_PLAIN_CONN | postgres://moonpg_plain:plain_pass@localhost:5432/moonpg_test | password |
| PG_MD5_CONN | postgres://moonpg_md5:md5_pass@localhost:5432/moonpg_test | md5 |
| PG_SCRAM_CONN | postgres://moonpg_scram:scram_pass@localhost:5432/moonpg_test | scram-sha-256 |
pub(open) trait Closer {
fn close(Self) -> Unit
}impl FromRow for User with fn from_row(r : Row) -> User raise PgError {
User::{ id: r.get(0), name: r.get(1), email: r.get(2) }
}
let users : Array[User] = pool.fetch("SELECT id, name, email FROM users")impl FromValue for MyType with fn from_value(v : Value) -> MyType raise ValueError {
match v {
Value::String(s) => parse_my_type(s)
_ => raise ValueError::ValueError("expected Value::String, got ...")
}
}while rows.has_next() {
let row = rows.get_row()
let v : Int = row.get(0)
}
rows.close()impl ToValue for MyType with fn to_value(self) -> Value {
...
}pub(open) trait Tx : QueryExecutor {
async fn commit(Self) -> Unit raise PgError
async fn rollback(Self) -> Unit raise PgError
}impl Closer for Connectionimpl QueryExecutor for Connectionasync fn execute(self : Connection, sql : String, params? : Array[&ToValue]) -> ExecResult raise PgErrorimpl TxBeginner for Connectionasync fn Connection::begin_copy(self : Connection, table : String, columns : Array[String]) -> CopyWriter raise PgErrorasync fn Connection::copy_in(self : Connection, sql : String, rows : Iter[String]) -> ExecResult raise PgErrorconn.copy_in("COPY t FROM STDIN", my_rows.iter())async fn[X] Connection::listen(self : Connection, channel : String, group : TaskGroup[X]) -> Listener raise PgError@async.with_task_group() group => {
let listener = conn.listen("events", group)
for ;; {
let notif = listener.recv()
group.spawn_bg() () => { handle(notif) }
}
}async fn Connection::notify(self : Connection, channel : String, payload? : String) -> Unit raise PgErrorlet w = conn.begin_copy("users", ["name", "age"])
w.write_row(["Alice", 30])
w.write_row(["Bob", 25])
let result = w.finish()pub(all) enum IsolationLevel {
ReadCommitted
RepeatableRead
Serializable
}pub(all) struct Pool {
queue : Queue[IdleConn]
conninfo : String
max_conns : Int
min_idle : Int
max_idle_sec : Int
max_lifetime_sec : Int
health_check : Bool
maintenance_interval_sec : Int
count : Ref[Int]
idle_count : Ref[Int]
running : Ref[Bool]
acquire_count : Ref[Int64]
acquire_wait_count : Ref[Int64]
acquire_wait_duration : Ref[Int64]
}impl QueryExecutor for Poolimpl TxBeginner for Poolpub(all) struct PoolConfig {
conninfo : String
max_conns : Int
min_idle : Int
max_idle_sec : Int
max_lifetime_sec : Int
health_check : Bool
maintenance_interval_sec : Int
}let pool = Pool::new(PoolConfig::new("postgres://...", max_conns=10))
// Option A: implicit acquire — Pool auto-manages lifecycle
let rows = pool.query("SELECT 1")
while rows.has_next() { let row = rows.get_row() }
rows.close() // PoolRows.close() returns conn to pool
let row = pool.query_one("SELECT 42")
pool.execute("INSERT ...") |> ignore
let tx = pool.begin_tx()
tx.commit() // PoolDbTx.commit() returns conn to pool
// Option B: explicit acquire — caller manages lifecycle
let pc = pool.acquire()
defer pc.release()
let rows = pc.query("SELECT 1")
rows.close() // PoolRows.close() returns conn to pool
let row = pc.query_one("SELECT 42")
pc.release() // caller releasesfn PoolConfig::new(conninfo : String, max_conns? : Int, min_idle? : Int, max_idle_sec? : Int, max_lifetime_sec? : Int, health_check? : Bool, maintenance_interval_sec? : Int) -> PoolConfigimpl QueryExecutor for PoolConnasync fn execute(self : PoolConn, sql : String, params? : Array[&ToValue]) -> ExecResult raise PgErrorimpl TxBeginner for PoolConnimpl QueryExecutor for PoolDbTxasync fn execute(self : PoolDbTx, sql : String, params? : Array[&ToValue]) -> ExecResult raise PgErrorpub(all) struct PoolStats {
total_connections : Int
idle_connections : Int
active_connections : Int
acquire_count : Int64
acquire_wait_count : Int64
acquire_wait_duration_ms : Int64
}let new_id = begin_func(conn, fn(tx) {
let row = tx.query_one("INSERT INTO users (name) VALUES ($1) RETURNING id", params=["alice"])
row.get(0)
})conn.execute("UPDATE t SET email = $1 WHERE id = $2", params=[null, 42])A pure MoonBit PostgreSQL client, wire protocol from scratch.
Dependencies