4 Commits
Author SHA1 Message Date
dsql f90c18ed64 fix: close() nulls pool + connect() lock guards against concurrent races (psql-7, psql-8)
- close() sets self._pool = None (in a finally, even on driver-error close) so a closed
  instance reports not-connected instead of masquerading as still-connected against a
  dead pool (the pool property now raises RuntimeError post-close as intended).
- connect() is now guarded by an internal asyncio.Lock: concurrent connect() calls
  serialize instead of each racing to build + orphan its own live pool. Twin fix with
  mysql at the same Layer-1 signature level.
- docstrings tightened (module header lifecycle/dsn notes, __init__ dsn note) with zero
  behavior change; re-verified against a real embedded postgres.

verified against pixeltable-pgserver (test-only): 8 concurrent connect() calls now leave
0 orphaned pools + exactly 1 live pool (old: 7 orphaned, 8 live simultaneously); close()
now flips the connected-check to False (old: stayed True against a dead pool). full
layer-1/layer-2/transaction/concurrency regression suite still green. fresh-venv install
confirms runtime purity (asyncpg only, no pixeltable-pgserver). bump v0.1.4 -> v0.1.5

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-02 23:25:52 -04:00
dsql 1ae280e5c3 fix: host/port defaults no longer shadow a dsn's embedded host/port (psql-5)
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-02 17:30:10 -04:00
dsql c3205f0614 fix: _where() renders None conditions as IS NULL instead of = NULL
col = $n bound to NULL never matches in sql, so get/get_one/exists/delete
silently missed every row filtered on a None value despite insert() writing
NULL fine. Kept in lockstep with the mysql lib's identical fix.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-02 16:44:27 -04:00
dsql 4e75a800f1 fix: connect() never orphans a pool (psql-3)
connect() now closes an existing pool before re-connecting (no orphan on double-connect),
and tears down the freshly-built pool if the SELECT-1 validation fails before re-raising
(no leaked pool on a failed connect). verified vs embedded postgres. bump v0.1.1 -> v0.1.2

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-01 00:28:27 -04:00
3 changed files with 73 additions and 22 deletions
+9 -3
View File
@@ -10,18 +10,18 @@ a sibling of the `mongo` lib. Class is **`PsqlDB`**.
`requirements.txt`:
```
psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v0.1.1
psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v0.1.5
```
Direct:
```bash
pip install "psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v0.1.1"
pip install "psql @ git+ssh://git@git.rethinkstudios.io/rethink-public/psql.git@v0.1.5"
```
Pulls `asyncpg`.
Drop the `@v0.1.1` suffix from the line above to install the latest unpinned.
Drop the `@v0.1.5` suffix from the line above to install the latest unpinned.
## The two-layer API
@@ -63,6 +63,12 @@ async with PsqlDB(database="app", user="postgres") as db:
await db.insert("events", {"kind": "login"})
```
`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 passing `dsn=...` in
`pool_kwargs` (with no `host`/`port` of your own) lets the dsn's server reach asyncpg
instead of being silently overridden. The no-dsn path above still defaults to
`localhost:5432` when you don't pass `host`.
### Layer 2 — raw SQL for the complex queries
```python
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
[project]
name = "psql"
version = "0.1.1"
version = "0.1.5"
description = "async postgres wrapper over asyncpg: two-layer API (friendly verbs + raw escape hatch), fail-loud, config-free"
requires-python = ">=3.10"
dependencies = [
+63 -18
View File
@@ -14,9 +14,17 @@ context manager:
await db.execute("CREATE TABLE ...")
lifecycle:
construction is sync and opens no socket. connect() builds the asyncpg pool and
validates it with `SELECT 1` so a bad host/credentials fails loud immediately rather
than on the first real op, and returns self. close() closes the pool.
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:
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
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.
two-layer API:
LAYER 1 — friendly, portable verbs for simple single-table CRUD. these hide the
@@ -46,6 +54,7 @@ errors (FAIL LOUD — unlike the mongo lib's swallow-and-default):
not wrapped, use the raw `.pool` property (the asyncpg.Pool).
"""
import asyncio
import logging
from typing import Any, Dict, List, Optional, Sequence
@@ -75,8 +84,8 @@ class PsqlDB:
def __init__(
self,
host: str = "localhost",
port: int = 5432,
host: Optional[str] = None,
port: Optional[int] = None,
database: Optional[str] = None,
user: Optional[str] = None,
password: Optional[str] = None,
@@ -91,7 +100,18 @@ class PsqlDB:
host/port/database/user/password/min_size/max_size/command_timeout are injected by
the caller. extra pool_kwargs pass through to asyncpg.create_pool (ssl, server_
settings, dsn, etc). `host` may be a unix socket directory as well as a hostname.
host/port default to None here (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 shadow the dsn's server. no dsn -> host/port fall back to
localhost:5432; dsn passed -> host/port stay None (letting the dsn's own reach
asyncpg) unless the caller explicitly overrides them.
"""
if "dsn" not in pool_kwargs:
if host is None:
host = "localhost"
if port is None:
port = 5432
self._config = dict(
host=host,
port=port,
@@ -104,22 +124,37 @@ class PsqlDB:
**pool_kwargs,
)
self._pool: Optional[asyncpg.Pool] = None
self._connect_lock = asyncio.Lock()
async def connect(self) -> "PsqlDB":
"""build the pool and validate it with SELECT 1; fail loud on bad config
returns self so callers can write `db = await PsqlDB(...).connect()`.
returns self so callers can write `db = await PsqlDB(...).connect()`. if called
again on an already-connected instance the previous pool is closed first (no
orphaned pool); if the SELECT-1 validation fails the freshly-built pool is torn
down before re-raising, so a failed connect() never leaks a live pool. guarded by
an internal lock so concurrent connect() calls build exactly one pool instead of
each racing to create + orphan their own.
"""
try:
self._pool = await asyncpg.create_pool(**self._config)
await self._pool.fetchval("SELECT 1")
except _DRIVER_ERRORS:
log.exception("psql.connect() failed")
raise
return self
async with self._connect_lock:
if self._pool is not None:
await self.close()
pool = await asyncpg.create_pool(**self._config)
try:
await pool.fetchval("SELECT 1")
except _DRIVER_ERRORS:
log.exception("psql.connect() validation failed")
await pool.close()
raise
except BaseException:
await pool.close()
raise
self._pool = pool
return self
async def close(self) -> None:
"""close the pool on shutdown"""
"""close the pool on shutdown and null the reference (so pool/connected checks
report not-connected against a dead pool)"""
if self._pool is None:
return
try:
@@ -127,6 +162,8 @@ class PsqlDB:
except _DRIVER_ERRORS:
log.exception("psql.close()")
raise
finally:
self._pool = None
async def __aenter__(self) -> "PsqlDB":
return await self.connect()
@@ -363,13 +400,21 @@ class _Transaction:
def _where(conditions: Optional[Dict[str, Any]]) -> tuple:
"""build a parameterized `WHERE col = $1 AND ...` clause + the params list
returns ("", []) when there are no conditions. equality only.
returns ("", []) when there are no conditions. equality only. a None value renders as
`col IS NULL` (not `col = $n` bound to NULL, which sql never matches) and does not
consume a placeholder.
"""
if not conditions:
return "", []
cols = list(conditions.keys())
clause = " AND ".join(f"{_quote_ident(c)} = ${i + 1}" for i, c in enumerate(cols))
return f" WHERE {clause}", list(conditions.values())
parts = []
params = []
for col, val in conditions.items():
if val is None:
parts.append(f"{_quote_ident(col)} IS NULL")
else:
params.append(val)
parts.append(f"{_quote_ident(col)} = ${len(params)}")
return f" WHERE {' AND '.join(parts)}", params
def _status_count(status: str) -> int: