Skip to content
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- `datacontract lint --all-errors`: unknown fields are warnings instead of errors, unless a custom `--schema` rejects them

### Added
- `datacontract import postgres` imports primary and foreign keys for SELECT-only roles with catalog access.
- `datacontract edit`: enable the editor's AI assistant via `DATACONTRACT_EDITOR_AI_*` environment variables (endpoint, API key, model, provider, auth header)

### Fixed
Expand Down
222 changes: 193 additions & 29 deletions datacontract/imports/postgres_importer.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@

Reads table and column metadata from ``information_schema``, the same catalog
``datacontract test`` reads back to verify physical types, so an imported
contract passes on the first run without hand-editing. Comments come from
contract passes on the first run without hand-editing. Primary keys use
PostgreSQL catalog metadata when available, with the portable
information-schema result as a fallback. Comments come from
``pg_description`` via ``obj_description`` / ``col_description``, which
``information_schema`` does not expose.

Expand All @@ -13,9 +15,15 @@

from __future__ import annotations

from typing import Any, Dict, List, Optional
import logging
from typing import Any, Callable, Dict, List, Optional

from open_data_contract_standard.model import OpenDataContractStandard, SchemaObject, SchemaProperty
from open_data_contract_standard.model import (
OpenDataContractStandard,
Relationship,
SchemaObject,
SchemaProperty,
)

from datacontract.config import Config
from datacontract.engines.ibis.native_type import reconstruct_native_type
Expand All @@ -33,6 +41,8 @@
DEFAULT_PORT = 5432
DEFAULT_SCHEMA = "public"

logger = logging.getLogger(__name__)

# to_regclass() yields NULL instead of raising for objects the user cannot see,
# so a missing comment never fails the whole import.
_TABLES_QUERY = """
Expand Down Expand Up @@ -66,17 +76,83 @@
ORDER BY table_name, ordinal_position
"""

_PRIMARY_KEYS_QUERY = """
_INFORMATION_SCHEMA_PRIMARY_KEYS_QUERY = """
SELECT kcu.table_name, kcu.column_name, kcu.ordinal_position
FROM information_schema.table_constraints tc
JOIN information_schema.key_column_usage kcu
ON tc.constraint_name = kcu.constraint_name
AND tc.table_schema = kcu.table_schema
AND tc.table_name = kcu.table_name
WHERE tc.constraint_type = 'PRIMARY KEY'
AND tc.table_schema = %s
ORDER BY kcu.table_name, kcu.ordinal_position
"""

# pg_constraint is the preferred source of primary-key metadata. The portable
# information-schema query remains a fallback when reading the catalog fails.
# conkey retains the declared order of columns in composite primary keys.
_CATALOG_PRIMARY_KEYS_QUERY = """
SELECT relation.relname AS table_name,
attribute.attname AS column_name,
con.oid AS constraint_oid,
key_column.ordinality AS ordinal_position
FROM pg_catalog.pg_constraint AS con
JOIN pg_catalog.pg_class AS relation
ON relation.oid = con.conrelid
JOIN pg_catalog.pg_namespace AS namespace
ON namespace.oid = relation.relnamespace
CROSS JOIN LATERAL unnest(con.conkey)
WITH ORDINALITY AS key_column(attnum, ordinality)
JOIN pg_catalog.pg_attribute AS attribute
ON attribute.attrelid = relation.oid
AND attribute.attnum = key_column.attnum
AND NOT attribute.attisdropped
WHERE con.contype = 'p'
AND namespace.nspname = %s
ORDER BY relation.relname, key_column.ordinality
"""

# Pairing conkey and confkey by ordinality preserves the source-to-target
# column mapping for both single-column and composite foreign keys.
_FOREIGN_KEYS_QUERY = """
SELECT source_relation.relname AS table_name,
source_attribute.attname AS column_name,
target_relation.relname AS foreign_table_name,
target_namespace.nspname AS foreign_table_schema,
target_attribute.attname AS foreign_column_name,
con.oid AS constraint_oid,
key_column.ordinality AS ordinal_position
FROM pg_catalog.pg_constraint AS con
JOIN pg_catalog.pg_class AS source_relation
ON source_relation.oid = con.conrelid
JOIN pg_catalog.pg_namespace AS source_namespace
ON source_namespace.oid = source_relation.relnamespace
JOIN pg_catalog.pg_class AS target_relation
ON target_relation.oid = con.confrelid
JOIN pg_catalog.pg_namespace AS target_namespace
ON target_namespace.oid = target_relation.relnamespace
CROSS JOIN LATERAL unnest(con.conkey, con.confkey)
WITH ORDINALITY AS key_column(source_attnum, target_attnum, ordinality)
JOIN pg_catalog.pg_attribute AS source_attribute
ON source_attribute.attrelid = source_relation.oid
AND source_attribute.attnum = key_column.source_attnum
AND NOT source_attribute.attisdropped
JOIN pg_catalog.pg_attribute AS target_attribute
ON target_attribute.attrelid = target_relation.oid
AND target_attribute.attnum = key_column.target_attnum
AND NOT target_attribute.attisdropped
WHERE con.contype = 'f'
AND source_namespace.nspname = %s
Comment thread
vtulus marked this conversation as resolved.
-- Skip the per-partition clones Postgres adds for a foreign key that references a partitioned table.
AND NOT EXISTS (
SELECT 1
FROM pg_catalog.pg_constraint AS parent_constraint
WHERE parent_constraint.oid = con.conparentid
AND parent_constraint.conrelid = con.conrelid
)
ORDER BY source_relation.relname, con.oid, key_column.ordinality
"""


class PostgresImporter(Importer):
def import_source(self, source: str, import_args: dict, config: "Config | None" = None) -> OpenDataContractStandard:
Expand Down Expand Up @@ -119,9 +195,16 @@ def import_postgres_from_connector(
try:
table_rows = _fetch(connection, _TABLES_QUERY, (schema,))
column_rows = _fetch(connection, _COLUMNS_QUERY, (schema,))
# A user without access to information_schema constraints still gets a
# usable contract, just without primary keys.
primary_key_rows = _fetch(connection, _PRIMARY_KEYS_QUERY, (schema,), optional=True)
information_schema_primary_key_rows = _fetch(
connection, _INFORMATION_SCHEMA_PRIMARY_KEYS_QUERY, (schema,), optional=True
)
catalog_primary_key_rows = _fetch(connection, _CATALOG_PRIMARY_KEYS_QUERY, (schema,), optional=True)
primary_key_rows = (
catalog_primary_key_rows
if catalog_primary_key_rows is not None
else information_schema_primary_key_rows or []
)
foreign_key_rows = _fetch(connection, _FOREIGN_KEYS_QUERY, (schema,), optional=True)
finally:
connection.close()

Expand All @@ -146,8 +229,47 @@ def import_postgres_from_connector(
schema=schema,
)
]
selected_table_names = {table["table_name"] for table in selected}
selected_column_names = {
table_name: {row["column_name"] for row in column_rows if row["table_name"] == table_name}
for table_name in selected_table_names
}
if catalog_primary_key_rows is not None:
primary_key_rows = _visible_constraint_rows(
primary_key_rows,
lambda row: row["column_name"] in selected_column_names.get(row["table_name"], set()),
)
foreign_key_rows_from_selected_tables = [
row for row in foreign_key_rows or [] if row["table_name"] in selected_table_names
]
omitted_foreign_keys = {}
for row in foreign_key_rows_from_selected_tables:
if row["foreign_table_schema"] != schema or row["foreign_table_name"] not in selected_table_names:
omitted_foreign_keys.setdefault(row["constraint_oid"], row)
for row in omitted_foreign_keys.values():
logger.warning(
"Omitting foreign key from %s.%s to %s.%s because the target table is not included in the imported contract.",
schema,
row["table_name"],
row["foreign_table_schema"],
row["foreign_table_name"],
)

selected_foreign_key_rows = [
row
for row in foreign_key_rows_from_selected_tables
if row["foreign_table_schema"] == schema and row["foreign_table_name"] in selected_table_names
]
if foreign_key_rows is not None:
selected_foreign_key_rows = _visible_constraint_rows(
selected_foreign_key_rows,
lambda row: (
row["column_name"] in selected_column_names.get(row["table_name"], set())
and row["foreign_column_name"] in selected_column_names.get(row["foreign_table_name"], set())
),
)
odcs.schema_ = [
_create_schema(table, column_rows, primary_key_rows)
_create_schema(table, column_rows, primary_key_rows, selected_foreign_key_rows)
for table in sorted(selected, key=lambda row: row["table_name"].lower())
]
report_unmapped_types(odcs)
Expand Down Expand Up @@ -178,26 +300,26 @@ def postgres_connection(host: str, port: int, database: str, config: Optional[Co
)


def _fetch(connection, query: str, params: tuple, optional: bool = False) -> List[Dict[str, Any]]:
def _fetch(connection, query: str, params: tuple, optional: bool = False) -> Optional[List[Dict[str, Any]]]:
"""Run a catalog query and return its rows as dicts keyed by column name."""
try:
with connection.cursor() as cursor:
with connection.cursor() as cursor:
try:
cursor.execute(query, params)
columns = [description[0] for description in cursor.description]
return [dict(zip(columns, row)) for row in cursor.fetchall()]
except Exception as e:
if optional:
# Roll back so the failed statement doesn't poison the transaction.
connection.rollback()
return []
raise DataContractException(
type="schema",
result="failed",
name="postgres catalog query failed",
reason=f"Could not read the Postgres catalog: {e}",
engine="datacontract-cli",
original_exception=e,
)
except Exception as e:
if optional:
# Roll back so the failed statement doesn't poison the transaction.
connection.rollback()
return None
raise DataContractException(
type="schema",
result="failed",
name="postgres catalog query failed",
reason=f"Could not read the Postgres catalog: {e}",
engine="datacontract-cli",
original_exception=e,
)


def _select_tables(table_rows: List[Dict[str, Any]], tables: Optional[List[str]]) -> List[Dict[str, Any]]:
Expand All @@ -207,26 +329,67 @@ def _select_tables(table_rows: List[Dict[str, Any]], tables: Optional[List[str]]
return [row for row in table_rows if row["table_name"].lower() in wanted]


def _visible_constraint_rows(
rows: List[Dict[str, Any]], row_is_visible: Callable[[Dict[str, Any]], bool]
) -> List[Dict[str, Any]]:
"""Drop constraints with a column the role can't read, so its name stays out of the contract."""
hidden = {row["constraint_oid"] for row in rows if not row_is_visible(row)}
return [row for row in rows if row["constraint_oid"] not in hidden]


def _create_schema(
table: Dict[str, Any],
column_rows: List[Dict[str, Any]],
primary_key_rows: List[Dict[str, Any]],
foreign_key_rows: List[Dict[str, Any]],
) -> SchemaObject:
table_name = table["table_name"]
primary_keys = {
row["column_name"]: index + 1
for index, row in enumerate(row for row in primary_key_rows if row["table_name"] == table_name)
row["column_name"]: row["ordinal_position"] for row in primary_key_rows if row["table_name"] == table_name
}
properties = [_create_property(row, primary_keys) for row in column_rows if row["table_name"] == table_name]
return create_schema_object(
relationships = {}
schema_relationships = []
foreign_keys_by_constraint = {}
for row in foreign_key_rows:
if row["table_name"] == table_name:
foreign_keys_by_constraint.setdefault(row["constraint_oid"], []).append(row)
for constraint_rows in foreign_keys_by_constraint.values():
if len(constraint_rows) == 1:
row = constraint_rows[0]
relationships.setdefault(row["column_name"], []).append(
Relationship(type="foreignKey", to=f"{row['foreign_table_name']}.{row['foreign_column_name']}")
)
else:
schema_relationships.append(
Relationship(
type="foreignKey",
**{
"from": [f"{table_name}.{row['column_name']}" for row in constraint_rows],
"to": [f"{row['foreign_table_name']}.{row['foreign_column_name']}" for row in constraint_rows],
},
)
)
properties = [
_create_property(row, primary_keys, relationships.get(row["column_name"]))
for row in column_rows
if row["table_name"] == table_name
]
schema = create_schema_object(
name=table_name,
physical_type=table.get("table_type") or "table",
description=_clean(table.get("remarks")),
properties=properties or None,
)
if schema_relationships:
schema.relationships = schema_relationships
return schema


def _create_property(row: Dict[str, Any], primary_keys: Dict[str, int]) -> SchemaProperty:
def _create_property(
row: Dict[str, Any],
primary_keys: Dict[str, int],
relationships: Optional[List[Relationship]] = None,
) -> SchemaProperty:
name = row["column_name"]
max_length = row.get("character_maximum_length")
precision = row.get("numeric_precision")
Expand All @@ -252,6 +415,7 @@ def _create_property(row: Dict[str, Any], primary_keys: Dict[str, int]) -> Schem
format=format,
dimensions=dimensions,
element_type=element_type,
relationships=relationships,
)


Expand Down
6 changes: 5 additions & 1 deletion docs/docs/imports/postgres.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ description: "Create a data contract from a Postgres schema."

# <img className="page-icon" src="/img/icons/postgres.svg" alt="" /> Import: Postgres

Creates a data contract from a Postgres schema by reading table metadata from `information_schema` — including column types with length and precision, nullability, primary keys, and the comments stored in `pg_description`. Works with Postgres and Postgres-compatible databases (e.g. RisingWave).
Creates a data contract from a Postgres schema. Tables and columns are read from `information_schema`, comments from `pg_description`, and key metadata from `pg_catalog`. Works with Postgres and Postgres-compatible databases (e.g. RisingWave).

```bash
datacontract import postgres \
Expand All @@ -22,6 +22,10 @@ The generated contract includes a ready-to-test `servers` block, so you can run

Credentials are provided as environment variables and are the same ones `datacontract test` uses: `DATACONTRACT_POSTGRES_USERNAME` and `DATACONTRACT_POSTGRES_PASSWORD` — see the [Postgres Reference](../reference/postgres.md).

## Key metadata and privileges

Primary keys and foreign keys are read from the Postgres catalog, so a role with `SELECT` on the tables and `USAGE` on the schema is enough. A foreign key is included only when both of its tables are imported from the same schema.

Comment thread
vtulus marked this conversation as resolved.
Working from a DDL file instead of a live database? Use [`datacontract import sql --dialect postgres`](./sql.md).

All options: **[`datacontract import postgres`](../commands/import/postgres.md)**.
Loading
Loading