Skip to content
Open
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
94 changes: 94 additions & 0 deletions sdk/python/feast/infra/registry/sql.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from pydantic import StrictInt, StrictStr
from sqlalchemy import ( # type: ignore
BigInteger,
CheckConstraint,
Column,
Index,
Integer,
Expand All @@ -26,6 +27,7 @@
select,
update,
)
from sqlalchemy.dialects import mysql
from sqlalchemy.engine import Engine
from sqlalchemy.exc import IntegrityError

Expand Down Expand Up @@ -277,6 +279,98 @@
)


class ObjectAuditLogOperation(str, Enum):
"""Closed set for object_audit_log.operation."""

CREATE = "create"
UPDATE = "update"
DELETE = "delete"


class ObjectAuditLogObjectType(str, Enum):
"""SQL table names for object_audit_log.object_type. Not proto class names,
not FeastObjectType ('feature view'), not lineage camelCase ('featureView').
Not CHECKed: the set grows with new registry tables and we have no Alembic."""

PROJECTS = "projects"
ENTITIES = "entities"
DATA_SOURCES = "data_sources"
FEATURE_VIEWS = "feature_views"
STREAM_FEATURE_VIEWS = "stream_feature_views"
SORTED_FEATURE_VIEWS = "sorted_feature_views"
ON_DEMAND_FEATURE_VIEWS = "on_demand_feature_views"
FEATURE_SERVICES = "feature_services"
SAVED_DATASETS = "saved_datasets"
VALIDATION_REFERENCES = "validation_references"
MANAGED_INFRA = "managed_infra"
PERMISSIONS = "permissions"


# Append-only audit of registry object create/update/delete.
# Do not queue, outbox, or write after commit — audit
# failure must fail the gRPC write.
#
# before_proto / after_proto are gzip of proto3 wire bytes (null on create /
# delete respectively). Future Feast proto changes must stay additive (new field
# numbers; reserved on deletes; no type/number reuse) so historical rows remain
# FromString-readable. object_type is a kind discriminator (table/kind name), not
# the protobuf message type name and not a schema version. Forensic reads should
# gzip.decompress + FromString, not from_proto().
object_audit_log = Table(
"object_audit_log",
metadata,
Column(
"id",
BigInteger().with_variant(Integer, "sqlite"),
primary_key=True,
autoincrement=True,
),
# project_id is the project name (same as sibling tables). Survives
# delete_project so forensic reads still show the name; a reused name
# gets a new feast_metadata PROJECT_UUID, so incarnation identity is
# project_uuid, not project_id.
Column("project_id", String(255), nullable=False),
Comment thread
piket marked this conversation as resolved.
Column("project_uuid", String(36), nullable=False),
Column("object_type", String(50), nullable=False),
Column("object_name", String(255), nullable=False),
Column("operation", String(20), nullable=False),
CheckConstraint(
"operation IN ("
+ ", ".join(repr(op.value) for op in ObjectAuditLogOperation)
+ ")",
name="ck_object_audit_log_operation",
),
Column("actor", String(255), nullable=True),
Column(
"before_proto",
LargeBinary().with_variant(mysql.LONGBLOB, "mysql"),
nullable=True,
), # gzip(proto); null on create
Column(
"after_proto",
LargeBinary().with_variant(mysql.LONGBLOB, "mysql"),
nullable=True,
), # gzip(proto); null on delete
Column("request_id", String(64), nullable=True),
# UTC epoch seconds, same as last_updated_timestamp and
# materialization_interval_history.recorded_at (int(datetime.timestamp())).
Column("recorded_at", BigInteger, nullable=False),
Comment thread
piket marked this conversation as resolved.
)
Index(
"idx_object_audit_log_project_object_recorded",
object_audit_log.c.project_uuid,
object_audit_log.c.object_type,
object_audit_log.c.object_name,
object_audit_log.c.recorded_at,
)
# Time-range / retention queries ("last N hours", DELETE WHERE recorded_at < …).
# Actor and request_id lookups are follow-up indexes if those filters show up.
Index(
"idx_object_audit_log_recorded_at",
object_audit_log.c.recorded_at,
)


class FeastMetadataKeys(Enum):
LAST_UPDATED_TIMESTAMP = "last_updated_timestamp"
PROJECT_UUID = "project_uuid"
Expand Down
135 changes: 133 additions & 2 deletions sdk/python/tests/unit/infra/registry/test_sql_registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,22 @@
from datetime import timedelta

import pytest
from sqlalchemy import func, inspect, select
from sqlalchemy.exc import IntegrityError

from feast.entity import Entity
from feast.errors import EntityNotFoundException
from feast.feature_view import MATERIALIZATION_INTERVALS_MAX_LEN, FeatureView
from feast.field import Field
from feast.infra.offline_stores.file_source import FileSource
from feast.infra.registry.sql import SqlRegistry, SqlRegistryConfig
from feast.infra.registry.sql import (
ObjectAuditLogObjectType,
ObjectAuditLogOperation,
SqlRegistry,
SqlRegistryConfig,
object_audit_log,
)
from feast.project import Project
from feast.types import Float32
from feast.utils import _utc_now

Expand Down Expand Up @@ -84,7 +94,7 @@ def test_delete_entity(self, sqlite_registry):

sqlite_registry.delete_entity("test_entity", "test_project")

with pytest.raises(Exception):
with pytest.raises(EntityNotFoundException):
sqlite_registry.get_entity("test_entity", "test_project")

def test_get_project_metadata_model_returns_initialized_metadata(
Expand Down Expand Up @@ -361,3 +371,124 @@ def _race():
"driver_stats", "test_project"
)
assert len(history) == 1


_OBJECT_AUDIT_LOG_COLUMNS = {
"id",
"project_id",
"project_uuid",
"object_type",
"object_name",
"operation",
"actor",
"before_proto",
"after_proto",
"request_id",
"recorded_at",
}

_AUDIT_PROJECT_UUID = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee"
_OBJECT_AUDIT_LOG_NULLABLE = {
"actor",
"before_proto",
"after_proto",
"request_id",
}


def _object_audit_log_row_count(registry: SqlRegistry) -> int:
with registry.write_engine.begin() as conn:
return conn.execute(
select(func.count()).select_from(object_audit_log)
).scalar_one()


class TestObjectAuditLogSchema:
"""Ticket 3: create_all creates an empty object_audit_log; no inserts yet."""

def test_create_all_creates_object_audit_log_columns_and_index(
self, sqlite_registry
):
inspector = inspect(sqlite_registry.write_engine)
assert inspector.has_table("object_audit_log")

columns = {
col["name"]: col for col in inspector.get_columns("object_audit_log")
}
assert set(columns) == _OBJECT_AUDIT_LOG_COLUMNS
for name, col in columns.items():
assert col["nullable"] is (name in _OBJECT_AUDIT_LOG_NULLABLE)

index_names = {idx["name"] for idx in inspector.get_indexes("object_audit_log")}
assert "idx_object_audit_log_project_object_recorded" in index_names
assert "idx_object_audit_log_recorded_at" in index_names

check_names = {
ck["name"] for ck in inspector.get_check_constraints("object_audit_log")
}
assert "ck_object_audit_log_operation" in check_names

@pytest.mark.parametrize(
"operation",
[
ObjectAuditLogOperation.CREATE,
ObjectAuditLogOperation.UPDATE,
ObjectAuditLogOperation.DELETE,
],
)
def test_operation_check_allows_enum_values(self, sqlite_registry, operation):
with sqlite_registry.write_engine.begin() as conn:
conn.execute(
object_audit_log.insert().values(
project_id="p",
project_uuid=_AUDIT_PROJECT_UUID,
object_type=ObjectAuditLogObjectType.ENTITIES.value,
object_name="e",
operation=operation.value,
recorded_at=0,
)
)

@pytest.mark.parametrize("operation", ["CREATE", "upsert", "delete ", ""])
def test_operation_check_rejects_unknown_values(self, sqlite_registry, operation):
with sqlite_registry.write_engine.begin() as conn:
with pytest.raises(IntegrityError):
conn.execute(
object_audit_log.insert().values(
project_id="p",
project_uuid=_AUDIT_PROJECT_UUID,
object_type=ObjectAuditLogObjectType.ENTITIES.value,
object_name="e",
operation=operation,
recorded_at=0,
)
)

def test_delete_project_leaves_object_audit_log_rows(self, sqlite_registry):
sqlite_registry.apply_project(Project(name="reused"))
with sqlite_registry.write_engine.begin() as conn:
conn.execute(
object_audit_log.insert().values(
project_id="reused",
project_uuid=_AUDIT_PROJECT_UUID,
object_type=ObjectAuditLogObjectType.ENTITIES.value,
object_name="e",
operation=ObjectAuditLogOperation.CREATE.value,
recorded_at=0,
)
)

sqlite_registry.delete_project("reused")
assert _object_audit_log_row_count(sqlite_registry) == 1

def test_apply_and_delete_entity_leave_object_audit_log_empty(
self, sqlite_registry
):
entity = Entity(name="test_entity", description="Test entity")
sqlite_registry.apply_entity(entity, "test_project")
sqlite_registry.delete_entity("test_entity", "test_project")

with pytest.raises(EntityNotFoundException):
sqlite_registry.get_entity("test_entity", "test_project")

assert _object_audit_log_row_count(sqlite_registry) == 0
Loading