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
27 changes: 18 additions & 9 deletions sdk/python/feast/api/registry/rest/metrics.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,9 @@
grpc_call,
paginate_and_sort,
)
from feast.errors import FeastObjectNotFoundException
from feast.permissions.action import AuthzedAction
from feast.permissions.security_manager import assert_permissions
from feast.protos.feast.registry import RegistryServer_pb2


Expand Down Expand Up @@ -433,15 +436,21 @@ async def recently_visited(
key = f"recently_visited_{user}"
visits = []
if project:
try:
visits_json = (
server.registry.get_project_metadata(project, key)
if server
else None
)
visits = json.loads(visits_json) if visits_json else []
except Exception:
visits = []
if server:
try:
project_obj = server.registry.get_project(
name=project, allow_cache=True
)
assert_permissions(
resource=project_obj, actions=[AuthzedAction.DESCRIBE]
)
except FeastObjectNotFoundException:
pass
try:
visits_json = server.registry.get_project_metadata(project, key)
visits = json.loads(visits_json) if visits_json else []
except Exception:
visits = []
else:
try:
if server:
Expand Down
51 changes: 29 additions & 22 deletions sdk/python/feast/api/registry/rest/rest_registry_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -235,28 +235,35 @@ async def dispatch(self, request: Request, call_next):
else:
object_type = None
object_name = None
visit = {
"path": path,
"timestamp": _utc_now().isoformat(),
"project": project,
"user": user,
"object": object_type,
"object_name": object_name,
"method": method,
}
try:
visits_json = self.registry.get_project_metadata(project, key)
visits = json.loads(visits_json) if visits_json else []
except Exception:
visits = []
visits.append(visit)
visits = visits[-self.recent_visits_limit :]
try:
self.registry.set_project_metadata(
project, key, json.dumps(visits)
)
except Exception as e:
logger.warning(f"Failed to persist recent visits: {e}")

response = await call_next(request)

if response.status_code < 400:
visit = {
"path": path,
"timestamp": _utc_now().isoformat(),
"project": project,
"user": user,
"object": object_type,
"object_name": object_name,
"method": method,
}
try:
visits_json = self.registry.get_project_metadata(
project, key
)
visits = json.loads(visits_json) if visits_json else []
except Exception:
visits = []
visits.append(visit)
visits = visits[-self.recent_visits_limit :]
try:
self.registry.set_project_metadata(
project, key, json.dumps(visits)
)
except Exception as e:
logger.warning(f"Failed to persist recent visits: {e}")
return response
Comment on lines +264 to +266

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Critical] Unreachable code after return statement

There are three lines of unreachable code after the return statement. The response handling and call_next() execution will never be reached, which could break the middleware functionality.

Current code:

+                    return response
                 response = await call_next(request)
                 return response

Suggested:

Suggested change
except Exception as e:
logger.warning(f"Failed to persist recent visits: {e}")
return response
return response
# Remove the unreachable code below

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The suggestion to remove those lines would break the middleware for all non-GET requests.

response = await call_next(request)
return response

Expand Down
40 changes: 30 additions & 10 deletions sdk/python/feast/feature_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -584,12 +584,22 @@ async def chat_ui():
@app.post("/materialize", dependencies=[Depends(inject_user_details)])
async def materialize(request: MaterializeRequest) -> None:
with feast_metrics.track_request_latency("/materialize"):
for feature_view in request.feature_views or []:
resource = await _get_feast_object(feature_view, True)
assert_permissions(
resource=resource,
actions=[AuthzedAction.WRITE_ONLINE],
if request.feature_views:
for feature_view in request.feature_views:
resource = await _get_feast_object(feature_view, True)
assert_permissions(
resource=resource,
actions=[AuthzedAction.WRITE_ONLINE],
)
else:
feature_views_to_materialize = store._get_feature_views_to_materialize(
None
)
for fv in feature_views_to_materialize:
assert_permissions(
resource=fv,
actions=[AuthzedAction.WRITE_ONLINE],
)

if request.disable_event_timestamp:
now = datetime.now()
Expand All @@ -615,12 +625,22 @@ async def materialize(request: MaterializeRequest) -> None:
@app.post("/materialize-incremental", dependencies=[Depends(inject_user_details)])
async def materialize_incremental(request: MaterializeIncrementalRequest) -> None:
with feast_metrics.track_request_latency("/materialize-incremental"):
for feature_view in request.feature_views or []:
resource = await _get_feast_object(feature_view, True)
assert_permissions(
resource=resource,
actions=[AuthzedAction.WRITE_ONLINE],
if request.feature_views:
for feature_view in request.feature_views:
resource = await _get_feast_object(feature_view, True)
assert_permissions(
resource=resource,
actions=[AuthzedAction.WRITE_ONLINE],
)
else:
feature_views_to_materialize = store._get_feature_views_to_materialize(
None
)
for fv in feature_views_to_materialize:
assert_permissions(
resource=fv,
actions=[AuthzedAction.WRITE_ONLINE],
)
await run_in_threadpool(
store.materialize_incremental,
utils.make_tzaware(parser.parse(request.end_ts)),
Expand Down
29 changes: 29 additions & 0 deletions sdk/python/feast/registry_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -891,6 +891,13 @@ def DeleteValidationReference(
def ListProjectMetadata(
self, request: RegistryServer_pb2.ListProjectMetadataRequest, context
):
try:
project = self.proxied_registry.get_project(
name=request.project, allow_cache=True
)
assert_permissions(resource=project, actions=[AuthzedAction.DESCRIBE])
except FeastObjectNotFoundException:
pass
return RegistryServer_pb2.ListProjectMetadataResponse(
project_metadata=[
project_metadata.to_proto()
Expand Down Expand Up @@ -923,6 +930,10 @@ def ApplyMaterialization(
return Empty()

def UpdateInfra(self, request: RegistryServer_pb2.UpdateInfraRequest, context):
project = self.proxied_registry.get_project(
name=request.project, allow_cache=True
)
assert_permissions(resource=project, actions=[AuthzedAction.UPDATE])
self.proxied_registry.update_infra(
infra=Infra.from_proto(request.infra),
project=request.project,
Expand All @@ -931,6 +942,10 @@ def UpdateInfra(self, request: RegistryServer_pb2.UpdateInfraRequest, context):
return Empty()

def GetInfra(self, request: RegistryServer_pb2.GetInfraRequest, context):
project = self.proxied_registry.get_project(
name=request.project, allow_cache=True
)
assert_permissions(resource=project, actions=[AuthzedAction.DESCRIBE])
return self.proxied_registry.get_infra(
project=request.project, allow_cache=request.allow_cache
).to_proto()
Expand Down Expand Up @@ -1063,6 +1078,13 @@ def DeleteProject(self, request: RegistryServer_pb2.DeleteProjectRequest, contex
def GetRegistryLineage(
self, request: RegistryServer_pb2.GetRegistryLineageRequest, context
):
try:
project = self.proxied_registry.get_project(
name=request.project, allow_cache=True
)
assert_permissions(resource=project, actions=[AuthzedAction.DESCRIBE])
except FeastObjectNotFoundException:
pass
direct_relationships, indirect_relationships = (
self.proxied_registry.get_registry_lineage(
project=request.project,
Expand Down Expand Up @@ -1101,6 +1123,13 @@ def GetObjectRelationships(
self, request: RegistryServer_pb2.GetObjectRelationshipsRequest, context
):
"""Get relationships for a specific object."""
try:
project = self.proxied_registry.get_project(
name=request.project, allow_cache=True
)
assert_permissions(resource=project, actions=[AuthzedAction.DESCRIBE])
except FeastObjectNotFoundException:
pass
relationships = self.proxied_registry.get_object_relationships(
project=request.project,
object_type=request.object_type,
Expand Down
Loading