Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Fixed

- PostgreSQL compound index names now fit its 63-byte identifier limit. Long names keep deterministic digests and their index suffix, so typed action indexes are created instead of silently skipped at startup.
- Concurrent first use of a PostgreSQL collection now serializes its schema bootstrap, avoiding races while creating the shared table and base indexes.
- PostgreSQL `$exists` now treats JSON `null` like the in-memory query engine, fixing queries for absent or null nested values.
- File-backed SQLite closes the old `aiosqlite` connection when rebinding across event loops, preventing a leaked worker thread from keeping the process alive.
- Graph transactions isolate their request identity map and invalidate touched parent cache entries after commit or rollback. Index setup is cached per database instance so a second store receives its own indexes.
Expand Down
2 changes: 1 addition & 1 deletion SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -312,7 +312,7 @@ No built-in migration framework. Adapters do not enforce schemas. Adding optiona

### 5.2 Pushdown vs in-memory

- **Postgres**: `translate_query` pushes the whole operator surface into JSONB SQL; `$text` becomes `to_tsvector('simple', …) @@ plainto_tsquery('simple', …)` (requires `$fields`; the GIN from `@fulltext_index` / `attribute(fulltext=True)` serves it when the fields match in order). Per-class indexes are `(entity, <fields>)` (or `WHERE entity = …` with `partial_by_entity`), descending keys `DESC NULLS LAST`, so typed `find(sort=…, limit=…)` walks an index; the whole-document GIN is optional (`JVSPATIAL_PG_GIN_INDEX`). Untranslatable queries fall back to a full scan + `QueryEngine.match`.
- **Postgres**: `translate_query` pushes the whole operator surface into JSONB SQL; `$text` becomes `to_tsvector('simple', …) @@ plainto_tsquery('simple', …)` (requires `$fields`; the GIN from `@fulltext_index` / `attribute(fulltext=True)` serves it when the fields match in order). Per-class indexes are `(entity, <fields>)` (or `WHERE entity = …` with `partial_by_entity`), descending keys `DESC NULLS LAST`, and generated index names stay within PostgreSQL's 63-byte identifier limit using a stable digest when needed, so typed `find(sort=…, limit=…)` can walk an index; first-use collection table and base-index creation is serialized per database instance; the whole-document GIN is optional (`JVSPATIAL_PG_GIN_INDEX`). Untranslatable queries fall back to a full scan + `QueryEngine.match`.
- **MongoDB**: native pushdown; queries run server-side (`$text` uses the collection's text index; `$fields` is stripped).
- **SQLite**: translated to SQL via `SQLiteTranslator` (subset; complex `$or` chains may fall back).
- **DynamoDB**: limited pushdown via `Select=COUNT` and key conditions; remainder filtered client-side.
Expand Down
37 changes: 35 additions & 2 deletions jvspatial/db/postgres.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@
import asyncio
import contextlib
import contextvars
import hashlib
import json
import logging
import re
Expand Down Expand Up @@ -130,6 +131,25 @@ def current_tenant() -> Optional[str]:
# max 63 bytes (Postgres limit). We do NOT support quoted identifiers; pin to
# the safe ASCII subset so we never need to worry about escape edge cases.
_SAFE_IDENT_RE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]{0,62}$")
_POSTGRES_MAX_IDENTIFIER_BYTES = 63


def _postgres_index_name(base: str, suffix: str) -> str:
"""Build a stable index name within PostgreSQL's 63-byte identifier limit.

Keep existing names unchanged when possible. Long compound field paths are
shortened with a digest of the complete name, preserving collision
resistance and the ``idx`` / ``uniq`` / ``fts`` suffix.
"""
full_name = f"{base}_{suffix}"
if len(full_name.encode("utf-8")) <= _POSTGRES_MAX_IDENTIFIER_BYTES:
return full_name

digest = hashlib.sha256(full_name.encode("utf-8")).hexdigest()[:10]
fixed_bytes = len(digest) + len(suffix) + 2 # two underscores
prefix_bytes = _POSTGRES_MAX_IDENTIFIER_BYTES - fixed_bytes
shortened = f"{base[:prefix_bytes]}_{digest}_{suffix}"
return shortened


_PARAM_RE = re.compile(r"\$(\d+)")
Expand Down Expand Up @@ -390,6 +410,7 @@ def __init__(
# Collections we've already created the table + base indexes for.
# Avoids running CREATE TABLE IF NOT EXISTS on the hot path.
self._collections_bootstrapped: Set[str] = set()
self._collection_bootstrap_locks: Dict[str, asyncio.Lock] = {}

# Vector columns configured per collection. Map collection ->
# {field_name: dim}. Populated by :meth:`enable_vector_column`;
Expand Down Expand Up @@ -428,6 +449,7 @@ def _discard_pool_from_dead_loop(self) -> None:
# ``asyncio.Lock`` binds to its loop too, so it has to go as well.
self._pool_lock = asyncio.Lock()
self._collections_bootstrapped.clear()
self._collection_bootstrap_locks.clear()

async def _ensure_pool(self) -> "Pool":
"""Lazily create the asyncpg pool. Idempotent + concurrency-safe."""
Expand Down Expand Up @@ -467,6 +489,7 @@ async def close(self) -> None:
self._pool = None
self._pool_loop = None
self._collections_bootstrapped.clear()
self._collection_bootstrap_locks.clear()

# ---- tenant scope (C6) -------------------------------------------------

Expand Down Expand Up @@ -613,6 +636,13 @@ async def _bootstrap_collection(self, collection: str) -> None:
"""
if collection in self._collections_bootstrapped:
return
lock = self._collection_bootstrap_locks.setdefault(collection, asyncio.Lock())
async with lock:
if collection in self._collections_bootstrapped:
return
await self._create_collection_schema(collection)

async def _create_collection_schema(self, collection: str) -> None:
col = _safe_collection(collection)
schema = _safe_collection(self.schema_name)
gin_sql = (
Expand Down Expand Up @@ -1446,6 +1476,8 @@ async def create_index(
(``<col>_<entity>_<fields>_idx``; smaller, one per class).
``drop_legacy=True`` also drops the unscoped pre-0.0.18
``<col>_<fields>_idx`` / ``_uniq`` index it replaces.
Generated names stay within PostgreSQL's 63-byte identifier limit;
overlong compound names retain a stable digest and index suffix.
* ``fulltext=True`` — a GIN index over ``to_tsvector('simple', ...)``
of the fields in order (``<col>_<entity>_<fields>_fts``, scoped to
``entity`` when given), matching
Expand Down Expand Up @@ -1562,7 +1594,7 @@ async def create_index(
scope = "entity_"
else:
scope = ""
index_name = f"{col}_{scope}{field_names}_{suffix}"
index_name = _postgres_index_name(f"{col}_{scope}{field_names}", suffix)
unique_sql = "UNIQUE " if unique and not fulltext else ""
body = f"ON {schema}.{col} USING {method} ({', '.join(keys)}){where_clause}"

Expand All @@ -1588,7 +1620,8 @@ async def create_index(
f"CREATE {unique_sql}INDEX IF NOT EXISTS {index_name} {body}"
)
if kwargs.get("drop_legacy") and scope:
legacy = f"{col}_{field_names}_{'uniq' if unique else 'idx'}"
legacy_suffix = "uniq" if unique else "idx"
legacy = _postgres_index_name(f"{col}_{field_names}", legacy_suffix)
await conn.execute(f"DROP INDEX IF EXISTS {schema}.{legacy}")

# ---- atomic compound ops (C4) ------------------------------------------
Expand Down
38 changes: 38 additions & 0 deletions tests/db/test_postgres_indexes_text.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
from jvspatial.db._postgres_translate import translate_query, translate_sort
from jvspatial.db.jsondb import JsonDB
from jvspatial.db.mongodb import _native_query
from jvspatial.db.postgres import _postgres_index_name
from jvspatial.db.query import QueryEngine
from jvspatial.db.sqlite import SQLiteDB
from jvspatial.exceptions import QueryError
Expand Down Expand Up @@ -121,6 +122,43 @@ async def _explain(admin: Any, sql: str, params: List[Any]) -> List[Dict[str, An
return _plan_nodes(plan[0]["Plan"])


async def test_long_index_names_are_shortened_deterministically():
base = "node_entity_context_agent_id_context_namespace_context_label"
name = _postgres_index_name(base, "uniq")

assert name == _postgres_index_name(base, "uniq")
assert len(name.encode("utf-8")) <= 63
assert name.endswith("_uniq")
assert name != _postgres_index_name(f"{base}_other", "uniq")


async def test_long_compound_index_is_created_and_reused():
async with _pg() as (_ctx, db, admin, schema):
fields = [
("context.agent_id", 1),
("context.namespace", 1),
("context.label", 1),
]
await db.create_index(
"node", fields, unique=True, entity="LongAction", entity_leading=True
)
await db.create_index(
"node", fields, unique=True, entity="LongAction", entity_leading=True
)

indexes = await _indexes(admin, schema, "node")
matching = [
(name, definition)
for name, definition in indexes.items()
if name.endswith("_uniq") and "context,label" in definition
]
assert len(matching) == 1
name, definition = matching[0]
assert len(name.encode("utf-8")) <= 63
assert "CREATE UNIQUE INDEX" in definition
assert "entity" in definition


# ---- entity-scoped per-class indexes ------------------------------------------


Expand Down
Loading