7 Commits
Author SHA1 Message Date
dsql 6466b797d3 chore: bump to 1.1.0 (logging-discipline audit)
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-08-10 23:00:50 -04:00
dsql 7691c79842 fix: log-XOR-raise — demote the pre-raise log to DEBUG, revise the contract (D1)
every wrapped method did log.exception (ERROR + traceback) immediately before re-raising
the driver error — the same failure reported at two layers (the log AND the raised
exception). that is the log-and-raise violation: a method that re-raises must not also
error-log, because the caller — which alone knows whether the failure is fatal or routine
— is the one that logs. demote all wrapped-method logs to log.debug(..., exc_info=True):
the traceback stays available at DEBUG, and the raised exception is the single loud
terminal signal. docstring updated from "logs via getLogger and re-raises" to the
corrected raise-XOR-log contract. no behavior change beyond log level — the fail-loud
re-raise is unchanged.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-08-09 02:18:08 -04:00
dsql a35dbcbb6e release: 1.0.0
first stable release. pre-1.0.0 verification complete: all surviving MED regressions and
gaps resolved and independently re-fired, tree audited clean across the suite.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-09 18:53:15 -04:00
dsql 71993fea1b docs: add __aenter__/__aexit__ one-liner docstrings for twin parity with redis
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-06 02:55:58 -04:00
dsql e4c61cf6c3 fix: log-and-reraise create_pool failures; lock close() against a racing connect()
connect() built the pool with `await asyncpg.create_pool(...)` outside the
try that wraps SELECT-1 validation - with the default min_size=1 the pool
connects eagerly, so a bad host/auth propagated with no log.exception line,
breaking the documented "every method logs then re-raises" contract.
create_pool() now sits inside the same try/log/re-raise as validation, and
either failure point tears down a partially-built pool before re-raising.

close() didn't take _connect_lock while connect() did, so a close() racing
an in-flight connect() would see _pool is None and no-op as success while
connect() went on to install a live pool - a shutdown handler racing a
reconnect could "close" the instance while real connections stayed open.
close() now takes the same lock via a shared _close_locked() helper.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-03 19:11:52 -04:00
dsql d619d4ca55 refactor: derive __version__ from package metadata (single source)
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-03 17:00:59 -04:00
dsql f07a8da5a4 docs: compress prose/module docstrings, em-dash->hyphen (de-bloat wave 1)
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-03 00:16:11 -04:00
4 changed files with 75 additions and 68 deletions
+3 -3
View File
@@ -10,18 +10,18 @@ a sibling of the `mongo` lib. Class is **`PsqlDB`**.
`requirements.txt`: `requirements.txt`:
``` ```
psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v0.1.5 psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v1.0.0
``` ```
Direct: Direct:
```bash ```bash
pip install "psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v0.1.5" pip install "psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v1.0.0"
``` ```
Pulls `asyncpg`. Pulls `asyncpg`.
Drop the `@v0.1.5` suffix from the line above to install the latest unpinned. Drop the `@v1.0.0` suffix from the line above to install the latest unpinned.
## The two-layer API ## The two-layer API
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
[project] [project]
name = "psql" name = "psql"
version = "0.1.5" version = "1.1.0"
description = "async postgres wrapper over asyncpg: two-layer API (friendly verbs + raw escape hatch), fail-loud, config-free" description = "async postgres wrapper over asyncpg: two-layer API (friendly verbs + raw escape hatch), fail-loud, config-free"
requires-python = ">=3.10" requires-python = ">=3.10"
dependencies = [ dependencies = [
+7
View File
@@ -1,3 +1,10 @@
from importlib.metadata import PackageNotFoundError, version
from .psql import PsqlDB from .psql import PsqlDB
__all__ = ["PsqlDB"] __all__ = ["PsqlDB"]
try:
__version__ = version("psql")
except PackageNotFoundError:
__version__ = "0.0.0+unknown"
+64 -64
View File
@@ -1,5 +1,5 @@
""" """
async postgres wrapper over asyncpg two-layer API (friendly verbs + raw escape hatch) async postgres wrapper over asyncpg - two-layer API (friendly verbs + raw escape hatch)
object pattern (one pool per process), attach to the app: object pattern (one pool per process), attach to the app:
from psql import PsqlDB from psql import PsqlDB
@@ -13,44 +13,32 @@ context manager:
async with PsqlDB(database="app", user="postgres") as db: async with PsqlDB(database="app", user="postgres") as db:
await db.execute("CREATE TABLE ...") await db.execute("CREATE TABLE ...")
lifecycle:
construction is sync, opens no socket. connect() builds the asyncpg pool, validates it
with `SELECT 1` (fail loud on bad host/credentials immediately, not on the first real
op), returns self; concurrent connect() calls are lock-serialized so only one pool is
ever live. close() closes the pool and nulls the reference (a closed instance reports
not-connected).
dsn: dsn:
host/port default to None, not "localhost"/5432 asyncpg only reads a dsn's embedded host/port default to None, not "localhost"/5432 - asyncpg only reads a dsn's embedded
host/port when the host/port kwargs are falsy, so a hardcoded default would silently host/port when the host/port kwargs are falsy, so a hardcoded default would silently
shadow the dsn's server. pass dsn=... in pool_kwargs alone (no host/port) to let the shadow the dsn's server. pass dsn=... in pool_kwargs alone (no host/port) to let the
dsn's own host/port reach asyncpg; the no-dsn path still defaults to localhost:5432. dsn's own host/port reach asyncpg; the no-dsn path still defaults to localhost:5432.
two-layer API: two-layer API:
LAYER 1 friendly, portable verbs for simple single-table CRUD. these hide the LAYER 1 - friendly, portable verbs for simple single-table CRUD, byte-for-byte
dialect and are byte-for-byte identical to the `mysql` lib, so a dev swaps psql<->mysql identical to the `mysql` lib (zero call-site changes swapping psql<->mysql):
with zero call-site changes: create_database, create_table, drop, insert, get, create_database, create_table, drop, insert, get, get_one, delete, exists, upsert.
get_one, delete, exists, upsert. LAYER 2 - raw escape hatch for the complex ~20% (joins, aggregates, CTEs, window
LAYER 2 — raw escape hatch for the complex ~20% (joins, aggregates, CTEs, window functions): execute, fetch, fetchone, fetchval, transaction, using `$1, $2`
functions): execute, fetch, fetchone, fetchval, transaction. you write the SQL with placeholders (asyncpg style). deliberately nothing in between - no query builder/ORM.
`$1, $2` placeholders (asyncpg style) + params; the wrapper still gives pooling,
parameterization, fail-loud errors, and row->dict conversion. raw SQL is psql-specific.
there is deliberately NOTHING in between — no query builder / ORM. a join goes through
raw fetch(), never a chainable .where()/.join().
rows:
layer 1 and fetch/fetchone return plain dicts ({column: value}), not asyncpg Record
objects — identical shape to the mysql lib.
placeholders / injection safety: placeholders / injection safety:
values are ALWAYS parameterized layer 1 builds `$1, $2` internally; layer 2 takes values are ALWAYS parameterized - layer 1 builds `$1, $2` internally; layer 2 takes
your `$1` placeholders + *params. never f-string/format a value into SQL. only your `$1` placeholders + *params. never f-string/format a value into SQL. only
identifiers (table/column names) are interpolated, and they are quoted. identifiers (table/column names) are interpolated, and they are quoted.
errors (FAIL LOUD unlike the mongo lib's swallow-and-default): errors (FAIL LOUD - unlike the mongo lib's swallow-and-default):
every method catches the driver error (asyncpg.PostgresError / InterfaceError, OSError every method catches the driver error (asyncpg.PostgresError / InterfaceError, OSError
on connection loss), logs via getLogger(__name__), and re-raises. a None/[] return is on connection loss) and re-raises it - the raised exception IS the signal, and the
only ever a real result (no row, empty table) — never a swallowed failure. for anything caller (which alone knows fatal-vs-routine) decides and logs. the wrapped method itself
logs only at DEBUG (with the traceback) so a failure isn't reported twice (raise XOR
log). a None/[] return is only ever a real result (no row, empty table) - never a
swallowed failure. for anything
not wrapped, use the raw `.pool` property (the asyncpg.Pool). not wrapped, use the raw `.pool` property (the asyncpg.Pool).
""" """
@@ -68,7 +56,7 @@ _DRIVER_ERRORS = (asyncpg.PostgresError, asyncpg.InterfaceError, OSError)
def _quote_ident(identifier: str) -> str: def _quote_ident(identifier: str) -> str:
"""quote a sql identifier (table/column), escaping embedded double-quotes """quote a sql identifier (table/column), escaping embedded double-quotes
identifiers can't be parameterized, so they are interpolated quoting + doubling any identifiers can't be parameterized, so they are interpolated - quoting + doubling any
embedded quote is the postgres-safe way to do that for caller-supplied names. embedded quote is the postgres-safe way to do that for caller-supplied names.
""" """
return '"' + identifier.replace('"', '""') + '"' return '"' + identifier.replace('"', '""') + '"'
@@ -138,62 +126,75 @@ class PsqlDB:
""" """
async with self._connect_lock: async with self._connect_lock:
if self._pool is not None: if self._pool is not None:
await self.close() await self._close_locked()
pool = await asyncpg.create_pool(**self._config) pool = None
try: try:
pool = await asyncpg.create_pool(**self._config)
await pool.fetchval("SELECT 1") await pool.fetchval("SELECT 1")
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.connect() validation failed") log.debug("psql.connect() failed", exc_info=True)
await pool.close() if pool is not None:
await pool.close()
raise raise
except BaseException: except BaseException:
await pool.close() if pool is not None:
await pool.close()
raise raise
self._pool = pool self._pool = pool
return self return self
async def close(self) -> None: async def close(self) -> None:
"""close the pool on shutdown and null the reference (so pool/connected checks """close the pool on shutdown and null the reference (so pool/connected checks
report not-connected against a dead pool)""" report not-connected against a dead pool)
guarded by the same lock as connect() - a close() racing an in-flight connect()
waits for it rather than no-opping against a not-yet-installed pool and leaving
the just-built one live.
"""
async with self._connect_lock:
await self._close_locked()
async def _close_locked(self) -> None:
"""close the pool and null the reference; caller must hold `_connect_lock`"""
if self._pool is None: if self._pool is None:
return return
try: try:
await self._pool.close() await self._pool.close()
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.close()") log.debug("psql.close()", exc_info=True)
raise raise
finally: finally:
self._pool = None self._pool = None
async def __aenter__(self) -> "PsqlDB": async def __aenter__(self) -> "PsqlDB":
"""enter: connect() and return self"""
return await self.connect() return await self.connect()
async def __aexit__(self, exc_type, exc, tb) -> None: async def __aexit__(self, exc_type, exc, tb) -> None:
"""exit: close(), ignoring exc_type/exc/tb"""
await self.close() await self.close()
@property @property
def pool(self) -> asyncpg.Pool: def pool(self) -> asyncpg.Pool:
"""raw asyncpg.Pool escape hatch; full driver surface, raises """raw asyncpg.Pool escape hatch for copy/prepare/listen-notify/cursors and
anything not wrapped; full driver surface, raises"""
use for copy/prepare/listen-notify/cursors and anything not wrapped.
"""
if self._pool is None: if self._pool is None:
raise RuntimeError("psql: not connected; call await db.connect() first") raise RuntimeError("psql: not connected; call await db.connect() first")
return self._pool return self._pool
# ------------------------------------------------------------------------- # -------------------------------------------------------------------------
# layer 2 raw escape hatch (you write the SQL, $1 placeholders) # layer 2 - raw escape hatch (you write the SQL, $1 placeholders)
async def execute(self, query: str, *params: Any) -> str: async def execute(self, query: str, *params: Any) -> str:
"""run a statement (INSERT/UPDATE/DELETE/DDL); return asyncpg's status string """run a statement (INSERT/UPDATE/DELETE/DDL); return asyncpg's status string
the status string is e.g. "INSERT 0 1" / "UPDATE 3" / "DELETE 2" parse it or use e.g. "INSERT 0 1" / "UPDATE 3" / "DELETE 2" - parse it, or use the layer-1 verbs
the layer-1 verbs (insert/delete) which return structured values instead. (insert/delete) which return structured values instead.
""" """
try: try:
return await self.pool.execute(query, *params) return await self.pool.execute(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.execute(): %s", query) log.debug("psql.execute(): %s", query, exc_info=True)
raise raise
async def fetch(self, query: str, *params: Any) -> List[dict]: async def fetch(self, query: str, *params: Any) -> List[dict]:
@@ -201,20 +202,19 @@ class PsqlDB:
try: try:
rows = await self.pool.fetch(query, *params) rows = await self.pool.fetch(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.fetch(): %s", query) log.debug("psql.fetch(): %s", query, exc_info=True)
raise raise
return [_row_to_dict(r) for r in rows] return [_row_to_dict(r) for r in rows]
async def fetchone(self, query: str, *params: Any) -> Optional[dict]: async def fetchone(self, query: str, *params: Any) -> Optional[dict]:
"""run a query and return the first row as a dict, or None if no rows """run a query and return the first row as a dict, or None if no rows
named fetchone (not asyncpg's fetchrow) to match the mysql lib's layer-2 surface; named fetchone (not asyncpg's fetchrow) to match the mysql lib's layer-2 surface.
maps to the driver's fetchrow internally.
""" """
try: try:
row = await self.pool.fetchrow(query, *params) row = await self.pool.fetchrow(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.fetchone(): %s", query) log.debug("psql.fetchone(): %s", query, exc_info=True)
raise raise
return _row_to_dict(row) if row is not None else None return _row_to_dict(row) if row is not None else None
@@ -223,7 +223,7 @@ class PsqlDB:
try: try:
return await self.pool.fetchval(query, *params) return await self.pool.fetchval(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.fetchval(): %s", query) log.debug("psql.fetchval(): %s", query, exc_info=True)
raise raise
def transaction(self): def transaction(self):
@@ -234,20 +234,20 @@ class PsqlDB:
await conn.execute("INSERT ...", a) await conn.execute("INSERT ...", a)
await conn.execute("UPDATE ...", b) await conn.execute("UPDATE ...", b)
commits on clean exit, rolls back and re-raises on any error. `conn` is a raw commits on clean exit, rolls back and re-raises on any error. `conn` is a raw
asyncpg connection (use its $1-placeholder execute/fetch/... directly). asyncpg connection ($1-placeholder execute/fetch/... directly).
""" """
return _Transaction(self.pool) return _Transaction(self.pool)
# ------------------------------------------------------------------------- # -------------------------------------------------------------------------
# layer 1 friendly portable verbs (identical across psql/mysql) # layer 1 - friendly portable verbs (identical across psql/mysql)
async def create_database(self, name: str) -> None: async def create_database(self, name: str) -> None:
"""CREATE DATABASE name (raises if it already exists postgres has no IF NOT """CREATE DATABASE name (raises if it already exists - postgres has no IF NOT
EXISTS for CREATE DATABASE; catch the duplicate error or check first)""" EXISTS for CREATE DATABASE; catch the duplicate error or check first)"""
try: try:
await self.pool.execute(f"CREATE DATABASE {_quote_ident(name)}") await self.pool.execute(f"CREATE DATABASE {_quote_ident(name)}")
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.create_database(%s)", name) log.debug("psql.create_database(%s)", name, exc_info=True)
raise raise
async def create_table(self, name: str, schema: Dict[str, str]) -> None: async def create_table(self, name: str, schema: Dict[str, str]) -> None:
@@ -262,7 +262,7 @@ class PsqlDB:
try: try:
await self.pool.execute(f"CREATE TABLE IF NOT EXISTS {_quote_ident(name)} ({cols})") await self.pool.execute(f"CREATE TABLE IF NOT EXISTS {_quote_ident(name)} ({cols})")
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.create_table(%s)", name) log.debug("psql.create_table(%s)", name, exc_info=True)
raise raise
async def drop(self, name: str, *, table: bool = True) -> None: async def drop(self, name: str, *, table: bool = True) -> None:
@@ -271,14 +271,14 @@ class PsqlDB:
try: try:
await self.pool.execute(f"DROP {kind} IF EXISTS {_quote_ident(name)}") await self.pool.execute(f"DROP {kind} IF EXISTS {_quote_ident(name)}")
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.drop(%s)", name) log.debug("psql.drop(%s)", name, exc_info=True)
raise raise
async def insert(self, table: str, values: Dict[str, Any]) -> int: async def insert(self, table: str, values: Dict[str, Any]) -> int:
"""INSERT one row from {column: value}; return the inserted rowcount (1) """INSERT one row from {column: value}; return the inserted rowcount (1)
values are parameterized ($1, $2, ...). returns the number of rows inserted (1 on values are parameterized ($1, $2, ...). returns the number of rows inserted (1 on
success) the portable return shape shared with the mysql lib. success) - the portable return shape shared with the mysql lib.
""" """
cols = list(values.keys()) cols = list(values.keys())
placeholders = ", ".join(f"${i + 1}" for i in range(len(cols))) placeholders = ", ".join(f"${i + 1}" for i in range(len(cols)))
@@ -287,14 +287,14 @@ class PsqlDB:
try: try:
status = await self.pool.execute(query, *values.values()) status = await self.pool.execute(query, *values.values())
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.insert(%s)", table) log.debug("psql.insert(%s)", table, exc_info=True)
raise raise
return _status_count(status) return _status_count(status)
async def get(self, table: str, conditions: Optional[Dict[str, Any]] = None) -> List[dict]: async def get(self, table: str, conditions: Optional[Dict[str, Any]] = None) -> List[dict]:
"""SELECT * rows matching equality `conditions` (col = val AND ...) as dicts """SELECT * rows matching equality `conditions` (col = val AND ...) as dicts
conditions=None/{} returns all rows. simple equality only anything more complex conditions=None/{} returns all rows. simple equality only - anything more complex
goes through raw fetch(). goes through raw fetch().
""" """
where, params = _where(conditions) where, params = _where(conditions)
@@ -302,7 +302,7 @@ class PsqlDB:
try: try:
rows = await self.pool.fetch(query, *params) rows = await self.pool.fetch(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.get(%s)", table) log.debug("psql.get(%s)", table, exc_info=True)
raise raise
return [_row_to_dict(r) for r in rows] return [_row_to_dict(r) for r in rows]
@@ -313,7 +313,7 @@ class PsqlDB:
try: try:
row = await self.pool.fetchrow(query, *params) row = await self.pool.fetchrow(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.get_one(%s)", table) log.debug("psql.get_one(%s)", table, exc_info=True)
raise raise
return _row_to_dict(row) if row is not None else None return _row_to_dict(row) if row is not None else None
@@ -327,7 +327,7 @@ class PsqlDB:
try: try:
status = await self.pool.execute(query, *params) status = await self.pool.execute(query, *params)
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.delete(%s)", table) log.debug("psql.delete(%s)", table, exc_info=True)
raise raise
return _status_count(status) return _status_count(status)
@@ -338,11 +338,11 @@ class PsqlDB:
try: try:
return bool(await self.pool.fetchval(query, *params)) return bool(await self.pool.fetchval(query, *params))
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.exists(%s)", table) log.debug("psql.exists(%s)", table, exc_info=True)
raise raise
async def upsert(self, table: str, values: Dict[str, Any], conflict: Sequence[str]) -> int: async def upsert(self, table: str, values: Dict[str, Any], conflict: Sequence[str]) -> int:
"""INSERT ... ON CONFLICT (conflict_cols) DO UPDATE insert or update on key clash """INSERT ... ON CONFLICT (conflict_cols) DO UPDATE - insert or update on key clash
`conflict` is the list of columns forming the unique/pk constraint to upsert on. `conflict` is the list of columns forming the unique/pk constraint to upsert on.
the wrapper emits ON CONFLICT here (the mysql lib emits ON DUPLICATE KEY UPDATE for the wrapper emits ON CONFLICT here (the mysql lib emits ON DUPLICATE KEY UPDATE for
@@ -360,7 +360,7 @@ class PsqlDB:
try: try:
status = await self.pool.execute(query, *values.values()) status = await self.pool.execute(query, *values.values())
except _DRIVER_ERRORS: except _DRIVER_ERRORS:
log.exception("psql.upsert(%s)", table) log.debug("psql.upsert(%s)", table, exc_info=True)
raise raise
return _status_count(status) return _status_count(status)
@@ -380,7 +380,7 @@ class _Transaction:
await self._tx.start() await self._tx.start()
except BaseException: except BaseException:
# start() (or transaction()) failing after acquire would otherwise leak the # start() (or transaction()) failing after acquire would otherwise leak the
# pooled connection __aexit__ is not called when __aenter__ raises. release # pooled connection - __aexit__ is not called when __aenter__ raises. release
# it and reset so a failed transaction start never burns a pool slot. # it and reset so a failed transaction start never burns a pool slot.
await self._pool.release(self._conn) await self._pool.release(self._conn)
self._conn = None self._conn = None