diff --git a/sunfish/events/redfish_subscription_handler.py b/sunfish/events/redfish_subscription_handler.py index 2b7ac4e..5c2c085 100644 --- a/sunfish/events/redfish_subscription_handler.py +++ b/sunfish/events/redfish_subscription_handler.py @@ -5,6 +5,7 @@ import json import os import string +import pdb from sunfish.events.subscription_handler_interface import SubscriptionHandlerInterface from sunfish.lib.exceptions import * @@ -62,6 +63,7 @@ def __init__(self, core): # Loads the subscriptions already stored def load_subscriptions(self): + #pdb.set_trace() path = os.path.join(os.getcwd(), self.fs_root, self.subscribers_root) if not os.path.exists(path): return @@ -80,6 +82,7 @@ def load_subscriptions(self): def new_subscription(self, payload: dict): # check if sub has colliding properties + #pdb.set_trace() if self.validate_subscription(payload) is False: raise IllegalSubscription diff --git a/sunfish/lib/core.py b/sunfish/lib/core.py index 0701690..3fdba08 100644 --- a/sunfish/lib/core.py +++ b/sunfish/lib/core.py @@ -3,6 +3,7 @@ # The full license terms are available here: https://github.com/OpenFabrics/sunfish_library_reference/blob/main/LICENSE import os +import json import string import uuid import logging @@ -208,12 +209,10 @@ def create_object(self, path: string, payload: dict): # 1. check the path target of the operation exists # self.storage_backend.read(path) # above done elsewhere, too soon to do here - # 2. is needed first forward the request to the agent managing the object + # 2. if needed first forward the request to the agent managing the object agent_response = self.objects_manager.forward_to_manager(SunfishRequestType.CREATE, path, payload=payload) if agent_response: payload_to_write = agent_response - # 3. Execute any custom handler for this object type AFTER Agent mods, if any - self.objects_handler.dispatch(object_type, path, SunfishRequestType.CREATE, payload=payload_to_write) except ResourceNotFound: logger.error("The collection where the resource is to be created does not exist.") except AgentForwardingFailure as e: @@ -222,8 +221,16 @@ def create_object(self, path: string, payload: dict): # The object does not have a handler. logger.debug(f"The object {object_type} does not have a custom handler") pass + # 3. Execute any custom handler for this object type AFTER Agent mods, if any + self.objects_handler.dispatch(object_type, path, SunfishRequestType.CREATE, payload=payload_to_write) # 4. persist change in Sunfish tree - return self.storage_backend.write(payload_to_write) + payload_written = self.storage_backend.write(payload_to_write) + # 5. create appropriate Event and send to subscribed EventDestinations + generate_resource_event = self.event_handler.resource_event_builder(SunfishRequestType.CREATE, path, payload=payload_written) + self.event_handler.new_event(generate_resource_event) + + + return payload_written def replace_object(self, path: str, payload: dict): """Calls the correspondent replace function from the backend implementation. @@ -242,21 +249,25 @@ def replace_object(self, path: str, payload: dict): try: # 1. check the path target of the operation exists self.storage_backend.read(path) - # 2. is needed first forward the request to the agent managing the object - #self.objects_manager.forward_to_manager(SunfishRequestType.REPLACE, path, payload=payload) + # 2. if needed first forward the request to the agent managing the object agent_response = self.objects_manager.forward_to_manager(SunfishRequestType.REPLACE, path, payload=payload) if agent_response: payload_to_write = agent_response - # 3. Execute any custom handler for this object type - self.objects_handler.dispatch(object_type, path, SunfishRequestType.REPLACE, payload=payload_to_write) except ResourceNotFound: logger.error(logger.error(f"The resource to be replaced ({path}) does not exist.")) except AttributeError: # The object does not have a handler. logger.debug(f"The object {object_type} does not have a custom handler") pass + # 3. Execute any custom handler for this object type + self.objects_handler.dispatch(object_type, path, SunfishRequestType.REPLACE, payload=payload_to_write) # 4. persist change in Sunfish tree - return self.storage_backend.replace(payload_to_write) + payload_written = self.storage_backend.replace(payload_to_write) + # 5. create appropriate Event and send to subscribed EventDestinations + generate_resource_event = self.event_handler.resource_event_builder(SunfishRequestType.REPLACE, path, payload=payload_written) + self.event_handler.new_event(generate_resource_event) + + return payload_written def patch_object(self, path: str, payload: dict): """Calls the correspondent patch function from the backend implementation. @@ -277,12 +288,9 @@ def patch_object(self, path: str, payload: dict): # 1. check the path target of the operation exists self.storage_backend.read(path) # 2. is needed first forward the request to the agent managing the object - #self.objects_manager.forward_to_manager(SunfishRequestType.PATCH, path, payload=payload) agent_response = self.objects_manager.forward_to_manager(SunfishRequestType.PATCH, path, payload=payload) if agent_response: payload_to_write = agent_response - # 3. Execute any custom handler for this object type - self.objects_handler.dispatch(object_type, path, SunfishRequestType.PATCH, payload=payload) except ResourceNotFound: logger.error(f"The resource to be patched ({path}) does not exist.") except AttributeError: @@ -290,8 +298,15 @@ def patch_object(self, path: str, payload: dict): logger.debug(f"The object {object_type} does not have a custom handler") pass + # 3. Execute any custom handler for this object type + self.objects_handler.dispatch(object_type, path, SunfishRequestType.PATCH, payload=payload) # 4. persist change in Sunfish tree - return self.storage_backend.patch(path, payload_to_write) + payload_written = self.storage_backend.patch(path, payload_to_write) + # 5. create appropriate Event and send to subscribed EventDestinations + generate_resource_event = self.event_handler.resource_event_builder(SunfishRequestType.PATCH, path, payload=payload_written) + self.event_handler.new_event(generate_resource_event) + + return payload_written def delete_object(self, path: string): """Calls the correspondent remove function from the backend implementation. Checks that the path is valid. @@ -312,16 +327,22 @@ def delete_object(self, path: string): self.storage_backend.read(path) # 2. is needed first forward the request to the agent managing the object self.objects_manager.forward_to_manager(SunfishRequestType.DELETE, path) - # 3. Execute any custom handler for this object type - self.objects_handler.dispatch(object_type, path, SunfishRequestType.DELETE) except ResourceNotFound: logger.error(f"The resource to be deleted ({path}) does not exist.") except AttributeError: # The object does not have a handler. logger.debug(f"The object {object_type} does not have a custom handler") + # 3. Execute any custom handler for this object type + self.objects_handler.dispatch(object_type, path, SunfishRequestType.DELETE) # 4. persist change in Sunfish tree - self.storage_backend.remove(path) + list_of_impacted_objects = self.storage_backend.remove(path) + # 5. process list of impacted objects for subscribers to ResourceEvents + events_sent_to = self.event_handler.process_new_resourceEvents(list_of_impacted_objects) + # 6. remove any deleted objects' URIs from Sunfish alias DB + #pdb.set_trace() + events_sent_to = self.event_handler.removeAliasesFromSunfishDB(list_of_impacted_objects) + # TODO return f"Object {path} deleted" def handle_event(self, payload): @@ -331,15 +352,45 @@ def handle_event(self, payload): else: context = "" logger.debug("Started handling incoming events") + sunfish_handled = False + all_event_responses = [] + this_event_response = {} + stat_code_max = 0 + for event in payload["Events"]: logger.debug(f"Handling event {event['MessageId']}") message_id = event['MessageId'].split(".")[-1] + event_id = event.get('EventId') or "" + event_origin = event.get('OriginOfCondition') or {} + stat_code = 500 try: - self.event_handler.dispatch(message_id, self.event_handler, event, context) + resp = self.event_handler.dispatch(message_id, self.event_handler, event, context) + if resp is not None: + #pdb.set_trace() + sunfish_handled = True + if type(resp) == int: + stat_code = resp + this_event_response["EventId"]=event_id + this_event_response["MessageId"]=message_id + this_event_response["dispatch_response"]=stat_code + this_event_response["origin"]=event_origin + all_event_responses.append(this_event_response) + if stat_code > stat_code_max: + stat_code_max = stat_code + except PropertyNotFound as e: logger.warning(repr(e)) raise e - return self.event_handler.new_event(payload) + + # if no events are handled by Sunfish, do NOT forward the original event to any subscribers + # for now return an unhandled response + if sunfish_handled is False: + return {"status": "un-processable content", "code": 422} + else: + if stat_code_max < 200: + stat_code_max = 200 + logger.info(f"event handler returned these results: \n {json.dumps(all_event_responses, indent = 4)}") + return {"status": "success", "code": 200} def _get_type(self, payload: dict, path: str = None): # controlla odata.type diff --git a/sunfish_plugins/events_handlers/redfish/redfish_event_handler.py b/sunfish_plugins/events_handlers/redfish/redfish_event_handler.py index c8b9e39..a2cf2f2 100644 --- a/sunfish_plugins/events_handlers/redfish/redfish_event_handler.py +++ b/sunfish_plugins/events_handlers/redfish/redfish_event_handler.py @@ -10,11 +10,14 @@ import shutil from uuid import uuid4 import pdb +import copy import requests from sunfish.events.event_handler_interface import EventHandlerInterface from sunfish.events.redfish_subscription_handler import subscriptions from sunfish.lib.exceptions import * +from sunfish.models.types import * +from typing import Optional logger = logging.getLogger("RedfishEventHandler") logging.basicConfig(level=logging.DEBUG) @@ -29,10 +32,11 @@ def AggregationSourceDiscovered(cls, event_handler: EventHandlerInterface, event # The arguments of the event message are: # - Arg0: "Redfish" # - Arg1: "agent_ip:port" - # I am also assuming that the agent name to be used is contained in the OriginOfCondifiton field of the event as in the below example: + # I am also assuming that the ConnectionMethod to be used is contained in the OriginOfCondifiton field of the event + # as in the below example: # { # "OriginOfCondition: [ - # "@odata.id" : "/redfish/v1/AggregationService/AggregationSource/AgentName" + # "@odata.id" : "/redfish/v1/AggregationService/ConnectionMethods/AgentName" # ]" # } logger.info("AggregationSourceDiscovered method called") @@ -40,8 +44,10 @@ def AggregationSourceDiscovered(cls, event_handler: EventHandlerInterface, event connectionMethodId = event['OriginOfCondition']['@odata.id'] hostname = event['MessageArgs'][1] # Agent address - #response = requests.get(f"{hostname}/{connectionMethodId}") - response = requests.get(f"{hostname}{connectionMethodId}") + try: + response = requests.get(f"{hostname}{connectionMethodId}") + except: + raise Exception(f"Cannot find {hostname})") if response.status_code != 200: raise Exception("Cannot find ConnectionMethod") response = response.json() @@ -65,14 +71,17 @@ def AggregationSourceDiscovered(cls, event_handler: EventHandlerInterface, event try: event_handler.core.storage_backend.write(aggregation_source_template) except Exception: - raise Exception() + raise Exception(f"Failed to store new aggregation source") agent_subscription_context = {"Context": aggregation_source_id.split('/')[-1]} - resp_patch = requests.patch(f"{hostname}/redfish/v1/EventService/Subscriptions/SunfishServer", + try: + resp_patch = requests.patch(f"{hostname}/redfish/v1/EventService/Subscriptions/SunfishServer", json=agent_subscription_context) - return resp_patch + return resp_patch.status_code + except Exception: + return 500 #something went wrong with assigning an ID to this Agent @classmethod def ResourceCreated(cls, event_handler: EventHandlerInterface, event: dict, context: str): @@ -87,6 +96,7 @@ def ResourceCreated(cls, event_handler: EventHandlerInterface, event: dict, cont # logger.info("New resource created") + new_resourceEvent_URIs = {} id = event['OriginOfCondition']['@odata.id'] # ex: /redfish/v1/Fabrics/CXL logger.info(f"aggregation_source's redfish URI: {id}") @@ -99,7 +109,10 @@ def ResourceCreated(cls, event_handler: EventHandlerInterface, event: dict, cont raise PropertyNotFound("Cannot find aggregation source; file does not exist") # fetch the actual resource to be created from agent hostname = aggregation_source["HostName"] - response = requests.get(f"{hostname}/{id}") + try: + response = requests.get(f"{hostname}/{id}") + except Exception: + raise ResourceNotFound("Aggregation source read from Agent failed") if response.status_code != 200: raise ResourceNotFound("Aggregation source read from Agent failed") @@ -120,16 +133,23 @@ def ResourceCreated(cls, event_handler: EventHandlerInterface, event: dict, cont fs_full_path = os.path.join(os.getcwd(), event_handler.core.conf["backend_conf"]["fs_root"], resource, 'index.json') if not os.path.exists(fs_full_path): - RedfishEventHandler.bfsInspection(event_handler.core, response, aggregation_source) + new_resourceEvent_URIs = RedfishEventHandler.bfsInspection(event_handler.core, response, aggregation_source) else: logger.warning(f"resource to create: {id} already exists.") # could be a second agent with naming conflicts, or same agent with duplicate # still run the inspection process on it to find cause of warning - RedfishEventHandler.bfsInspection(event_handler.core, response, aggregation_source) - + new_resourceEvent_URIs = RedfishEventHandler.bfsInspection(event_handler.core, response, aggregation_source) + # need to parse the new_resourceEvent_URIs{} for URIs to objects created, changed or deleted + try: + #pdb.set_trace() + notified_list =[] + notified_list = RedfishEventHandler.process_new_resourceEvents(event_handler, new_resourceEvent_URIs) + pass + except Exception as e: + logging.error(f"Sunfish Internal Event Generation function Error", exc_info=True) + pass # patch the aggregation_source object in storage with all the new resources found - #pdb.set_trace() event_handler.core.storage_backend.patch(agg_src_path, aggregation_source) logger.debug(f"\n{json.dumps(aggregation_source, indent=4)}") return 200 @@ -146,24 +166,30 @@ def ResourceChanged(cls, event_handler: EventHandlerInterface, event: dict, cont try: if "OriginOfCondition" not in event or not event["OriginOfCondition"].get("@odata.id"): logger.error("ResourceChanged event is missing OriginOfCondition.") - return + return 400 if not context: logger.error("No context (AggregationSource ID) in ResourceChanged event.") - return + return 400 aggregation_source_id = context aggregation_source = event_handler.core.storage_backend.read(f"/redfish/v1/AggregationService/AggregationSources/{aggregation_source_id}") host = aggregation_source["HostName"] origin_of_condition = event["OriginOfCondition"]["@odata.id"] + new_resourceEvent_URIs = {} + new_resourceEvent_URIs["changed"]=[] # Fetch the updated resource from the agent logger.info(f"Fetching updated resource {origin_of_condition} from agent {aggregation_source_id} at {host}") resource_endpoint = host + origin_of_condition - response = requests.get(resource_endpoint) + try: + response = requests.get(resource_endpoint) + except Exception: + logger.error(f"Exception occured trying to retrieve {origin_of_condition} from agent {aggregation_source_id}.") + return 500 if response.status_code != 200: logger.error(f"Could not fetch resource {origin_of_condition} from agent {aggregation_source_id}. Status: {response.status_code}") - return + return response.status_code updated_resource = response.json() # URI Aliasing to find the object in Sunfish @@ -174,7 +200,7 @@ def ResourceChanged(cls, event_handler: EventHandlerInterface, event: dict, cont event_handler.core.storage_backend.read(sunfish_uri) except NotFound: logger.error(f"ResourceChanged event for a non-existent object. Agent URI: {origin_of_condition}, Sunfish URI: {sunfish_uri}") - return + return 400 # Get aliases for this agent to update links in the payload uri_alias_file = os.path.join(os.getcwd(), event_handler.core.conf["backend_conf"]["fs_private"], 'URI_aliases.json') @@ -197,12 +223,105 @@ def ResourceChanged(cls, event_handler: EventHandlerInterface, event: dict, cont logger.info(f"Patching resource at {sunfish_uri}") event_handler.core.storage_backend.patch(sunfish_uri, updated_resource) + # need to create a ResourceUpdated event + try: + #pdb.set_trace() + new_resourceEvent_URIs["changed"].append(sunfish_uri) + notified_list =[] + notified_list = RedfishEventHandler.process_new_resourceEvents(event_handler, new_resourceEvent_URIs) + except Exception as e: + logging.error(f"Sunfish Internal Event Generation function Error", exc_info=True) + pass + # After patching, check if any cross-agent links need to be updated RedfishEventHandler.updateAllAgentsRedirectedLinks(event_handler.core) + return 200 except Exception: logger.error("Exception in ResourceChanged handler", exc_info=True) + return 500 + + + @classmethod + def ResourceDeleted(cls, event_handler: EventHandlerInterface, event: dict, context: str): + """ + Handles a ResourceDeleted event from an agent. + This handler translates the deleted resource Agent-name URI + to the corresponding Sunfish URI, and removes the existing object and all its subordinates + from the database. + + Then the Sunfish URIs of all deleted files are checked for aliases, and the aliases are removed + from the Sunfish URI_aliasDB file. + """ + #pdb.set_trace() + try: + if "OriginOfCondition" not in event or not event["OriginOfCondition"].get("@odata.id"): + logger.error("ResourceDeleted event is missing OriginOfCondition.") + return 400 + + if not context: + logger.error("No context (AggregationSource ID) in ResourceDeleted event.") + return 400 + + aggregation_source_id = context + aggregation_source = event_handler.core.storage_backend.read(f"/redfish/v1/AggregationService/AggregationSources/{aggregation_source_id}") + host = aggregation_source["HostName"] + origin_of_condition = event["OriginOfCondition"]["@odata.id"] + new_resourceEvent_URIs = {} + new_resourceEvent_URIs["deleted"]=[] + new_resourceEvent_URIs["changed"]=[] + new_resourceEvent_URIs["created"]=[] + + # URI Aliasing to find the object in Sunfish + sunfish_uri = RedfishEventHandler.xlateToSunfishPath(event_handler.core, origin_of_condition, aggregation_source) + # Check if object exists before attempting to delete + try: + event_handler.core.storage_backend.read(sunfish_uri) + except NotFound: + logger.error(f"ResourceDeleted event for a non-existent object. Agent URI: {origin_of_condition}, Sunfish URI: {sunfish_uri}") + return 400 + + # Update any internal @odata.id links in the fetched payload + # No need to update links in the to-be-deleted object + #updated_resource = RedfishEventHandler.update_aliased_links_in_object(event_handler.core, updated_resource, agent_aliases) + + # Boundary Link Processing - check if this delete affects a boundary link + # need to find all objects that will be deleted with this one + #if "Oem" in updated_resource and "Sunfish_RM" in updated_resource["Oem"] and updated_resource["Oem"]["Sunfish_RM"].get("BoundaryComponent") == "BoundaryPort": + #RedfishEventHandler.track_boundary_port(event_handler.core, updated_resource, aggregation_source) + + # + logger.info(f"Deleting resource at {sunfish_uri}") + # remove() returns dictionary of lists of files changed or deleted + new_resourceEvent_URIs = event_handler.core.storage_backend.remove(sunfish_uri) + + # need to create a ResourceEvents for all files deleted or changed + try: + #pdb.set_trace() + notified_list =[] + notified_list = RedfishEventHandler.process_new_resourceEvents(event_handler, new_resourceEvent_URIs) + pass + except Exception as e: + logging.error(f"Sunfish Internal Event Generation function Error", exc_info=True) + pass + + # now remove all aliases for deleted URIs + + try: + #pdb.set_trace() + uri_aliases = RedfishEventHandler.removeAliasesFromSunfishDB(event_handler, new_resourceEvent_URIs) + except Exception as e: + logging.error(f"Sunfish URI alias Database Cleanup error", exc_info=True) + + # After deleting, check if any cross-agent links need to be updated + #RedfishEventHandler.updateAllAgentsRedirectedLinks(event_handler.core) + # Even if there were problems generating related ResourceDeleted events, the DELETE was successful + return 200 + + except Exception: + logger.error("Exception in ResourceDeleted handler", exc_info=True) + return 500 @classmethod def TriggerEvent(cls, event_handler: EventHandlerInterface, event: dict, context: str): @@ -245,23 +364,20 @@ def TriggerEvent(cls, event_handler: EventHandlerInterface, event: dict, context response = requests.post(destination,json=event_to_send) if response.status_code != 200: logger.debug(f"Destination returned code {response.status_code}") - return response + return response.status_code else: logger.info(f"TriggerEvents Succeeded: code {response.status_code}") - return response + return response.status_code except Exception: raise Exception(f"Event forwarding to destination {destination} failed.") - response = 500 - return response + return 500 else: logger.error(f"file not found: {file_to_send} ") - response = 404 - return response + return 404 except Exception: raise Exception("TriggerEvents Failed") - resp = 500 - return resp + return 500 @@ -271,6 +387,7 @@ class RedfishEventHandler(EventHandlerInterface): "AggregationSourceDiscovered": RedfishEventHandlersTable.AggregationSourceDiscovered, "ResourceCreated": RedfishEventHandlersTable.ResourceCreated, "ResourceChanged": RedfishEventHandlersTable.ResourceChanged, + "ResourceDeleted": RedfishEventHandlersTable.ResourceDeleted, "TriggerEvent": RedfishEventHandlersTable.TriggerEvent } @@ -292,12 +409,15 @@ def dispatch(cls, message_id: str, event_handler: EventHandlerInterface, event: else: logger.debug(f"Message id '{message_id}' does not have a custom handler") - def new_event(self, payload): + def new_event(self, payload, origin_type = None): """Compares event's information with the subsribtions data structure to find the Ids of the subscribers for that event. Args: payload (dict): event received. + origin_type (dict): resource type of OriginOfCondition in payload (optional) + (if event is a DELETE notice, there is no object to read to fine its ResourceType) """ + #pdb.set_trace() for event in payload["Events"]: prefix = event["MessageId"].split('.')[0] messageId = event["MessageId"] @@ -308,7 +428,7 @@ def new_event(self, payload): if prefix in subscriptions["RegistryPrefixes"]: for id in subscriptions["RegistryPrefixes"][prefix]["exclude"]: - to_exclude.extend(id) + to_exclude.append(id) if messageId in subscriptions["MessageIds"]: for id in subscriptions["MessageIds"][messageId]["exclude"]: to_exclude.append(id) @@ -319,14 +439,17 @@ def new_event(self, payload): if "OriginOfCondition" in event: origin = event["OriginOfCondition"]["@odata.id"] try: - type = self.check_data_type(origin) + if origin_type is None: + type = RedfishEventHandler.check_data_type(self, origin) + else: + type=origin_type except ResourceNotFound as e: raise ResourceNotFound(e.resource_id) if type in subscriptions["ResourceTypes"]: to_forward.extend(subscriptions["ResourceTypes"][type]) if origin in subscriptions["OriginResources"]: to_forward.extend(subscriptions["OriginResources"][origin]) - sub = self.check_subdirs(origin) + sub = RedfishEventHandler.check_subdirs(self, origin) to_forward.extend(sub) if prefix in subscriptions["RegistryPrefixes"]: @@ -340,7 +463,7 @@ def new_event(self, payload): set2 = set(to_exclude) to_forward = list(set1 - set2) - return self.forward_event(to_forward, payload) + return RedfishEventHandler.forward_event(self, to_forward, payload) def check_data_type(self, origin): length = len(self.redfish_root) @@ -364,19 +487,23 @@ def forward_event(self, list, payload): ResourceNotFound: if it is not possible to get the subscription's details. Returns: - list: list of all the reachable subcribers for the event. + list: (modified) list of all the reachable subcribers for the event. """ - for id in list: + for id in list[:]: #must use a slice-copy of list since we modify list in the loop path = os.path.join(self.redfish_root, 'EventService', 'Subscriptions', id) try: data = self.core.storage_backend.read(path) + # use the context found in the subscription, if there is one + if "Context" in data: + payload["Context"] = data["Context"] resp = requests.post(data['Destination'], json=payload) resp.raise_for_status() except (requests.exceptions.ConnectionError, requests.exceptions.HTTPError) as e: logger.warning(f"Unable to contact event destination {id} for event , skipping.") logger.warning(f"Event log: \n{json.dumps(payload, indent=2)}") list.remove(id) + continue # don't quit on the rest of the original list except ResourceNotFound: raise ResourceNotFound(path) @@ -424,6 +551,9 @@ def bfsInspection(self, node, aggregation_source): fetched = [] notfound = [] uploaded = [] + updated = [] + deleted = [] + modified = {} visited.append(node['@odata.id']) queue.append(node['@odata.id']) @@ -461,15 +591,26 @@ def handleNestedObject(self, obj): logger.info(json.dumps(sorted(fetched),indent = 4)) logger.info("\n\nAgent did not return objects for the following URIs:\n") logger.info(json.dumps(sorted(notfound),indent = 4)) - + # find the uploaded objects + set1 = set(fetched) + set2 = set(notfound) + uploaded = sorted(list(set1-set2)) # now need to revisit all uploaded objects and update any links renamed after # the uploaded object was written RedfishEventHandler.updateAllAliasedLinks(self,aggregation_source) # now we need to re-direct any boundary port link references # this needs to be done on ALL agents, not just the one we just uploaded - RedfishEventHandler.updateAllAgentsRedirectedLinks(self) - - return visited + updated = RedfishEventHandler.updateAllAgentsRedirectedLinks(self) + # created the dict of uploaded, updated, or deleted URIs, + # uploaded URIs are in Agent namespace and need to be xlated + uploaded = RedfishEventHandler.xlateEventOriginsToSunfish(self, uploaded,aggregation_source) + modified["created"] = uploaded + # updated URIs are in Sunfish namespace already + modified["changed"] = updated + # deleted URIs is empty list for this method + modified["deleted"] = deleted + + return modified def create_uploaded_object(self, path: str, payload: dict): # before to add the ID and to call the methods there should be the json validation @@ -767,7 +908,6 @@ def updateAllAliasedLinks(self,aggregation_source): def updateObjectAliasedLinks(self, object_URI, agent_aliases): def findNestedURIs(self, URI_to_match, URI_to_sub, obj, path_to_nested_URI): - #pdb.set_trace() nestedPaths = [] if type(obj) == list: i = 0; @@ -829,6 +969,7 @@ def updateAllAgentsRedirectedLinks(self ): modified_aliasDB = False + modified_objects = [] for owning_agent_id in uri_aliasDB['Agents_xref_URIs']: logger.debug(f"redirecting placeholder links in all boundary ports for : {owning_agent_id}") if owning_agent_id in uri_aliasDB['Agents_xref_URIs']: @@ -844,6 +985,7 @@ def updateAllAgentsRedirectedLinks(self ): modified_aliasDB = True # need to replace the update object and re-save the uri_aliasDB self.storage_backend.replace(agent_bp_obj) + modified_objects.append(agent_bp_URI) else: logger.info(f"------ PeerPortURI NOT found") pass @@ -854,6 +996,7 @@ def updateAllAgentsRedirectedLinks(self ): modified_aliasDB = True # need to replace the update object and re-save the uri_aliasDB self.storage_backend.replace(agent_bp_obj) + modified_objects.append(agent_bp_URI) else: logger.info(f"------ PeerPortURI NOT found") pass @@ -864,6 +1007,7 @@ def updateAllAgentsRedirectedLinks(self ): modified_aliasDB = True # need to replace the update object and re-save the uri_aliasDB self.storage_backend.replace(agent_bp_obj) + modified_objects.append(agent_bp_URI) else: logger.info(f"------ PeerPortURI NOT found") pass @@ -874,7 +1018,7 @@ def updateAllAgentsRedirectedLinks(self ): with open(uri_alias_file,'w') as data_json: json.dump(uri_aliasDB, data_json, indent=4, sort_keys=True) data_json.close() - return + return modified_objects def redirectInterswitchLinks(self,owning_agent_id, agent_bp_obj,uri_aliasDB): @@ -1319,7 +1463,183 @@ def track_boundary_port(self, redfish_obj, aggregation_source): logger.debug(f"----- boundary ports matched {matching_ports}") return + def resource_event_builder(self, request_type: 'sunfish.models.types.SunfishRequestType', path: str, payload: Optional[dict] = None) -> Optional[dict]: + #pdb.set_trace() + ResourceCreated_template = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "New Resource Created", + "Context": "", + "Events": [ { + "Severity": "Ok", + "Message": "New Resource Created ", + "MessageId": "ResourceEvent.1.x.ResourceCreated", + "MessageArgs": [ ], + "OriginOfCondition": { + "@odata.id": "" + } + } ] + } + + ResourceChanged_template = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "Resource Changed", + "Context": "", + "Events": [ { + "Severity": "Ok", + "Message": "Existing Resource Changed ", + "MessageId": "ResourceEvent.1.x.ResourceChanged", + "MessageArgs": [ ], + "OriginOfCondition": { + "@odata.id": "" + } + } ] + } + + ResourceDeleted_template = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "Resource Deleted", + "Context": "", + "Events": [ { + "Severity": "Ok", + "Message": "Existing Resource Deleted ", + "MessageId": "ResourceEvent.1.x.ResourceDeleted", + "MessageArgs": [ ], + "OriginOfCondition": { + "@odata.id": "" + } + } ] + } + + logger.debug(f"Event_Builder: called on {path} with payload {payload}") + if request_type == SunfishRequestType.DELETE: + ResourceDeleted_template["Events"][0]["OriginOfCondition"]\ + ["@odata.id"]=path + return ResourceDeleted_template + elif payload is not None: # need a payload for these request types + if request_type == SunfishRequestType.CREATE: + ResourceCreated_template["Events"][0]["OriginOfCondition"]\ + ["@odata.id"]=payload["@odata.id"] + return ResourceCreated_template + elif request_type == SunfishRequestType.REPLACE or request_type == SunfishRequestType.PATCH: + ResourceChanged_template["Events"][0]["OriginOfCondition"]\ + ["@odata.id"]=payload["@odata.id"] + return ResourceChanged_template + + # no payload, or no proper request_type, so no event to return + + return + + def xlateEventOriginsToSunfish(self, agent_URIs: list , aggregation_source): + + sunfish_URIs = [] + try: + for agent_path in agent_URIs: + sunfish_path= RedfishEventHandler.xlateToSunfishPath(self, agent_path, aggregation_source) + sunfish_URIs.append(sunfish_path) + except: + pass + + return sunfish_URIs + + def process_new_resourceEvents(self, sunfish_object_URIs): + #pdb.set_trace() + eventOrigin = {} + event_to_send = {} + was_sent_to = [] + if "created" in sunfish_object_URIs: + try: + for created_URI in sunfish_object_URIs["created"]: + eventOrigin_path = created_URI + eventOrigin["@odata.id"] = created_URI + action_type = SunfishRequestType.CREATE + event_to_send = RedfishEventHandler.resource_event_builder(self, request_type = action_type, path = eventOrigin_path, payload = eventOrigin) + #pdb.set_trace() + was_sent_to.extend(RedfishEventHandler.new_event(self, event_to_send)) + logger.debug(f"sent event to {len(was_sent_to)} Destinations") + logger.debug(json.dumps(was_sent_to, indent=4)) + except: + print(f"process_new_resourceEvents: Exception in CREATED") + pass + + if "changed" in sunfish_object_URIs: + try: + for created_URI in sunfish_object_URIs["changed"]: + eventOrigin_path = created_URI + eventOrigin["@odata.id"] = created_URI + action_type = SunfishRequestType.PATCH + event_to_send = self.resource_event_builder(action_type, eventOrigin_path, eventOrigin) + was_sent_to.extend(RedfishEventHandler.new_event(self, event_to_send)) + except: + print(f"process_new_resourceEvents: Exception in CHANGED") + pass + if "deleted" in sunfish_object_URIs: + try: + for created_URI in sunfish_object_URIs["deleted"]: + eventOrigin_path = created_URI + eventOrigin["@odata.id"] = created_URI + action_type = SunfishRequestType.DELETE + event_to_send = self.resource_event_builder(action_type, eventOrigin_path, eventOrigin) + if "deleted_types" in sunfish_object_URIs and created_URI in sunfish_object_URIs["deleted_types"]: + origin_type = sunfish_object_URIs["deleted_types"][created_URI] + was_sent_to.extend(RedfishEventHandler.new_event(self, event_to_send, origin_type)) + except: + print(f"process_new_resourceEvents: Exception in DELETED") + pass + + return was_sent_to + + def removeAliasesFromSunfishDB(self,deleted_Sunfish_URIs): + try: + #pdb.set_trace() + uri_alias_file = os.path.join(os.getcwd(), self.core.conf["backend_conf"]["fs_private"], 'URI_aliases.json') + if os.path.exists(uri_alias_file): + with open(uri_alias_file, 'r') as data_json: + uri_aliasDB = json.load(data_json) + data_json.close() + else: + logger.error(f"alias file {uri_alias_file} not found") + raise Exception + + except: + pdb.set_trace() + raise Exception + + try: + #pdb.set_trace() + changed_URI_DB = False + sunfish_names = copy.deepcopy(uri_aliasDB.get("Sunfish_xref_URIs",{})) + agent_names = copy.deepcopy(uri_aliasDB.get("Agents_xref_URIs", {})) + # search the sunfish_xref_URIs key:value pairs + # key = sunfish_name, value = list of agent_names + for deleted_URI in deleted_Sunfish_URIs.get("deleted",[]): + if "aliases" in sunfish_names: + for sunfish_name, agent_list in sunfish_names["aliases"].items(): + if sunfish_name == deleted_URI: + del uri_aliasDB["Sunfish_xref_URIs"]["aliases"][sunfish_name] + changed_URI_DB = True + # + # next we need to search the agent_xref_URIs + # which are key = agent_name, value = sunfish_name + # + for agent_id in agent_names: + for agent_URI, sunfish_URI in agent_names[agent_id]["aliases"].items(): + if sunfish_URI == deleted_URI: + del uri_aliasDB["Agents_xref_URIs"][agent_id]["aliases"][agent_URI] + changed_URI_DB = True + + # now need to write aliasDB back to file if we didn't mess it up + if changed_URI_DB: + with open(uri_alias_file,'w') as data_json: + json.dump(uri_aliasDB, data_json, indent=4, sort_keys=True) + data_json.close() + except Exception as e: + # don't change anything + logging.error(f"Removing Deleted Aliases Failed", exc_info=True) + pdb.set_trace() + pass + + return uri_aliasDB def add_aggregation_source_reference(redfish_obj, aggregation_source): # BoundaryComponent = ["owned", "foreign", "BoundaryLink","unknown"] diff --git a/sunfish_plugins/objects_handlers/sunfish_server/redfish_object_handler.py b/sunfish_plugins/objects_handlers/sunfish_server/redfish_object_handler.py index 60ece2a..a856f59 100644 --- a/sunfish_plugins/objects_handlers/sunfish_server/redfish_object_handler.py +++ b/sunfish_plugins/objects_handlers/sunfish_server/redfish_object_handler.py @@ -3,6 +3,7 @@ # The full license terms are available here: https://github.com/OpenFabrics/sunfish_library_reference/blob/main/LICENSE import logging import string +import pdb from typing import Optional import sunfish.lib.core @@ -26,7 +27,8 @@ def EventDestination(cls, core: 'sunfish.lib.core.Core', path: str, operation: S core.subscription_handler.delete_subscription(payload) core.subscription_handler.new_subscription(payload) elif operation == SunfishRequestType.DELETE: - core.subscription_handler.delete_subscription(path) + subscriber_id = path.split("/")[-1] + core.subscription_handler.delete_subscription(subscriber_id) class RedfishObjectHandler(ObjectHandlerInterface): diff --git a/sunfish_plugins/objects_managers/sunfish_agent/agents_management.py b/sunfish_plugins/objects_managers/sunfish_agent/agents_management.py index c8e3c28..21959b1 100644 --- a/sunfish_plugins/objects_managers/sunfish_agent/agents_management.py +++ b/sunfish_plugins/objects_managers/sunfish_agent/agents_management.py @@ -38,17 +38,20 @@ def get_id(self) -> string: def is_agent_managed(cls, sunfish_core: 'sunfish.lib.core.Core', path: string): # if this is a top level resource, there's no need to check for the agent as no agent can own top level ones. # Example of top levels is Systems, Chassis, etc... - path_to_owner = (path.replace(sunfish_core.conf["redfish_root"], "").split("/")) - level = len(path_to_owner) + #pdb.set_trace() + #path_to_owner = (path.replace(sunfish_core.conf["redfish_root"], "").split("/")) + #level = len(path_to_owner) + path_to_owner = path[len(sunfish_core.conf["redfish_root"]):] + level = len(path_to_owner.split("/")) if level == 1: return None # path passed in is to parent object, which is usually a collection collection = sunfish_core.storage_backend.read(path) if 'Collection' in collection["@odata.type"]: - print(f"parent obj {path} is a Collection.") - new_path = os.path.join('/'.join(path_to_owner[:-1])) - print(f"grandparent obj at {new_path}") + logger.debug(f"parent obj {path} is a Collection.") + new_path = os.path.dirname(path) + logger.debug(f"grandparent obj at {new_path}") collection = sunfish_core.storage_backend.read(new_path) logger.debug(f"Checking if the object {new_path} is managed by an Agent") if "Oem" in collection and "Sunfish_RM" in collection["Oem"] and "ManagingAgent" in collection["Oem"]["Sunfish_RM"]: diff --git a/sunfish_plugins/objects_managers/sunfish_agent/sunfish_agent_manager.py b/sunfish_plugins/objects_managers/sunfish_agent/sunfish_agent_manager.py index 784ac6a..de8181f 100644 --- a/sunfish_plugins/objects_managers/sunfish_agent/sunfish_agent_manager.py +++ b/sunfish_plugins/objects_managers/sunfish_agent/sunfish_agent_manager.py @@ -23,12 +23,9 @@ def __init__(self, core: 'sunfish.lib.core.Core'): self.core = core def forward_to_manager(self, request_type: 'sunfish.models.types.SunfishRequestType', path: string, payload: dict = None) -> Optional[dict]: - uri_aliasDB = {} agent_response = None object_modified = False path_to_check = path - print(f"!!obj path to foward is {path}") - print(f"!!request_type is {request_type}") #pdb.set_trace() if request_type == SunfishRequestType.CREATE: # When creating an object, the request must be done on the collection. Since collections are generally not @@ -43,6 +40,7 @@ def forward_to_manager(self, request_type: 'sunfish.models.types.SunfishRequestT # in this case there would be no parent to inherit the agent from. Here this creation request should be # rejected because in Sunfish only agents can create elements in the top level directories and this is done # via events. + #pdb.set_trace() path_elems = path.split("/")[1:-1] path_to_check = "".join(f"/{e}" for e in path_elems) # get the parent path @@ -93,7 +91,7 @@ def forward_to_manager(self, request_type: 'sunfish.models.types.SunfishRequestT else: logger.debug(f"{path} is not managed by an agent") - return agent_response + return agent_response def xlateToAgentURIs(self, sunfish_obj ): diff --git a/sunfish_plugins/storage/file_system_backend/backend_FS.py b/sunfish_plugins/storage/file_system_backend/backend_FS.py index 91f6f16..62d6f2c 100644 --- a/sunfish_plugins/storage/file_system_backend/backend_FS.py +++ b/sunfish_plugins/storage/file_system_backend/backend_FS.py @@ -3,6 +3,7 @@ # The full license terms are available here: https://github.com/OpenFabrics/sunfish_library_reference/blob/main/LICENSE import pdb +import copy import json import logging import os @@ -290,14 +291,24 @@ def remove(self, path:str): ResourceNotFound: it is not possible to remove a resource that does not exists. Returns: - str: confirmation string + dict: list of files created, changed, or deleted """ # code that removes a file logging.info('BackendFS: remove called') + files_removed = [] + parent_file_modified = False + files_modified = [] + new_resourceEvents_URIs={} + new_resourceEvents_URIs["created"] = [] + new_resourceEvents_URIs["changed"] = [] + new_resourceEvents_URIs["deleted"] = [] + new_resourceEvents_URIs["deleted_types"] = {} length = len(self.redfish_root) resource_id = path[length:] + parent_id = os.path.dirname(resource_id) + base_path = os.path.join(os.getcwd(), self.root) full_path = os.path.join(os.getcwd(), self.root, resource_id) if len(resource_id) == 0: @@ -307,71 +318,75 @@ def remove(self, path:str): raise ResourceNotFound(resource_id) parent_path = os.path.dirname(full_path) - json_path = os.path.join(parent_path, 'index.json') + parent_json_path = os.path.join(parent_path, 'index.json') + # find files that will be removed + #pdb.set_trace() + files_removed = self._list_subordinates(full_path, base_path) + # find object types of files that will be removed + for name in files_removed: + data = self.read(name) + obj_type = data["@odata.type"].split('.')[0] + obj_type = obj_type.replace("#","") + new_resourceEvents_URIs["deleted_types"][name]=obj_type + shutil.rmtree(full_path) try: - with open(json_path, "r") as file: + # retrieve parent object if it exists + with open(parent_json_path, "r") as file: pdata = json.load(file) file.close() data = { "@odata.id": os.path.join(self.redfish_root, resource_id) } - collection_name = resource_id.split('/')[-1] + # define the object's ID within parent object (/Fabrics/CXL/Connections/this_one =>this_one) + object_name = resource_id.split('/')[-1] + # remove the Redfish ID from a Collection's 'Members' if 'Members' in pdata and data in pdata['Members']: pdata['Members'].remove(data) pdata['Members@odata.count'] = int(pdata['Members@odata.count']) - 1 - elif collection_name in pdata: - del pdata[collection_name] + parent_file_modified = True + # otherwise, see if deleted path is a subordinate + elif object_name in pdata: + del pdata[object_name] + parent_file_modified = True + + # write the parent object back to the file system + if parent_file_modified: + files_modified.append(os.path.join(self.redfish_root, parent_id)) + with open(parent_json_path, "w") as file: + json.dump(pdata, file, indent=4, sort_keys=True) + file.close() - with open(json_path, "w") as file: - json.dump(pdata, file, indent=4, sort_keys=True) - file.close() except FileNotFoundError as e: raise ResourceNotFound(resource_id) - # check links - to_replace = False - first = False + #pdb.set_trace() - for path, directories, files in os.walk(os.path.join(os.getcwd(), self.root)): - if 'index.json' in files: - file_path = os.path.join(path, 'index.json') + self._remove_all_references(files_removed, files_modified) - with open(file_path, "r") as file: - pdata = json.load(file) + new_resourceEvents_URIs["changed"].extend(files_modified) + new_resourceEvents_URIs["deleted"].extend(files_removed) + return new_resourceEvents_URIs - if 'Links' in pdata and path != os.path.join(os.getcwd(), self.root): - link_list = pdata['Links'] - to_del = [] - for link in link_list: - for x in link_list[link]: - if isinstance(link_list[link], list): - to_compare = "" - if type(x) is dict and "@odata.id" in x: - to_compare = x['@odata.id'] - elif type(x) is str: - to_compare = x - if to_compare == os.path.join(self.redfish_root, resource_id): - to_replace = True - link_list[link].remove(x) - if len(link_list[link]) == 0: - to_del.append(link) - elif isinstance(link_list[link], dict): - if x == os.path.join(self.redfish_root, resource_id): - to_del.append(link) - to_replace = True - if to_del: - for el in to_del: - del link_list[el] - if to_replace: - with open(file_path, "w") as file: - json.dump(pdata, file, indent=4, sort_keys=True) - to_replace = False + def _list_subordinates(self, path: str, base_path: str): + to_delete = [] + + # Check if the base directory exists + if not os.path.exists(path): + return to_delete - return "DELETE: file removed." + # Walk through all subordinate directories and files + for root, dirs, files in os.walk(path): + for name in files: + if name == "index.json": + #to_delete.append(os.path.join(root, name)) + redfish_path = os.path.join(self.redfish_root,os.path.relpath(root, base_path)) + to_delete.append(redfish_path) + + return to_delete @@ -401,3 +416,56 @@ def reset_resources(self, resource_path: str, clean_resource_path: str): resp = "Fail", 500 return resp + def _remove_all_references(self,removed_list, modified_files ): + # check all other objects in the service tree from root + # for links (URLs) pointing to the removed object + # this brute force search may need to be optimized! + to_replace = False + base_path = os.path.join(os.getcwd(), self.root) + + for path, directories, files in os.walk(base_path): + if 'index.json' in files: + file_path = os.path.join(path, 'index.json') + + with open(file_path, "r") as file: + pdata = json.load(file) + + if 'Links' in pdata and path != os.path.join(os.getcwd(), self.root): + link_list = pdata['Links'] #grab pointer to the "Links" structure + link_list_copy = copy.deepcopy(link_list) + to_del = [] + for link in link_list_copy: + for x in link_list_copy[link]: + if isinstance(link_list_copy[link], list): + to_compare = "" + if type(x) is dict and "@odata.id" in x: + to_compare = x['@odata.id'] + elif type(x) is str: + to_compare = x + # have to check this link against each deleted resource URI + for URI in removed_list: + #if to_compare == os.path.join(self.redfish_root, URI): + if to_compare == URI: + to_replace = True + link_list[link].remove(x) #manipulate original + if len(link_list[link]) == 0: + to_del.append(link) + elif isinstance(link_list_copy[link], dict): + # have to check this link against each deleted resource URI + for URI in removed_list: + #if x == os.path.join(self.redfish_root, URI): + if x == URI: + to_del.append(link) + to_replace = True + if to_del: + for el in to_del: + del link_list[el] + # after checking a URI Links in file, write it back if needed + if to_replace: + with open(file_path, "w") as file: + json.dump(pdata, file, indent=4, sort_keys=True) + # path is full filesystem name, need the redfish ID + redfish_path = os.path.join(self.redfish_root, os.path.relpath(path, base_path)) + modified_files.append(redfish_path) + to_replace = False + diff --git a/tests/Resources/AggregationService/AggregationSources/feb0bb58-83f2-4945-a798-7c0811a29955/index.json b/tests/Resources/AggregationService/AggregationSources/feb0bb58-83f2-4945-a798-7c0811a29955/index.json new file mode 100644 index 0000000..0f698bb --- /dev/null +++ b/tests/Resources/AggregationService/AggregationSources/feb0bb58-83f2-4945-a798-7c0811a29955/index.json @@ -0,0 +1,13 @@ +{ + "@odata.id": "/redfish/v1/AggregationService/AggregationSources/feb0bb58-83f2-4945-a798-7c0811a29955", + "@odata.type": "#AggregationSource.v1_2_.AggregationSource", + "HostName": "http://127.0.0.1:8080", + "Id": "feb0bb58-83f2-4945-a798-7c0811a29955", + "Links": { + "ConnectionMethod": { + "@odata.id": "/redfish/v1/AggregationService/ConnectionMethods/Pytest3" + }, + "ResourcesAccessed": [ + ] + } +} diff --git a/tests/Resources/AggregationService/AggregationSources/feb7bb58-83f2-4945-a798-7c0811a29955/index.json b/tests/Resources/AggregationService/AggregationSources/feb7bb58-83f2-4945-a798-7c0811a29955/index.json index f400ead..2731a3b 100644 --- a/tests/Resources/AggregationService/AggregationSources/feb7bb58-83f2-4945-a798-7c0811a29955/index.json +++ b/tests/Resources/AggregationService/AggregationSources/feb7bb58-83f2-4945-a798-7c0811a29955/index.json @@ -8,7 +8,6 @@ "@odata.id": "/redfish/v1/AggregationService/ConnectionMethods/Pytest1" }, "ResourcesAccessed": [ - "/redfish/v1/Fabrics/Pytest1" ] } } diff --git a/tests/Resources/AggregationService/AggregationSources/index.json b/tests/Resources/AggregationService/AggregationSources/index.json index fbb2a81..745d611 100755 --- a/tests/Resources/AggregationService/AggregationSources/index.json +++ b/tests/Resources/AggregationService/AggregationSources/index.json @@ -4,8 +4,11 @@ "Members": [ { "@odata.id": "/redfish/v1/AggregationService/AggregationSources/feb7bb58-83f2-4945-a798-7c0811a29955" + }, + { + "@odata.id": "/redfish/v1/AggregationService/AggregationSources/feb0bb58-83f2-4945-a798-7c0811a29955" } ], - "Members@odata.count": 1, + "Members@odata.count": 2, "Name": "AggregationService/AggregationSources Collection" } diff --git a/tests/test_sunfishcore_library.py b/tests/test_sunfishcore_library.py index 5870bc2..a4c91c4 100644 --- a/tests/test_sunfishcore_library.py +++ b/tests/test_sunfishcore_library.py @@ -80,6 +80,7 @@ def test_delete_exception(self): def test_post_object(self): json_file = tests_template.test_post_system path = os.path.join(self.conf["redfish_root"], "Systems") + #pdb.set_trace() assert self.core.create_object(path, json_file) def test_post_collection_exception(self): @@ -153,19 +154,22 @@ def httpserver_listen_address(self): def test_event_forwarding(self, httpserver: HTTPServer): httpserver.expect_request("/").respond_with_data("OK") - resp = self.core.handle_event(tests_template.task_event_cancelled) + #pdb.set_trace() + #resp = self.core.handle_event(tests_template.task_event_cancelled) + resp = self.core.event_handler.new_event(tests_template.task_event_cancelled) assert len(resp) == 1 def test_event_forwarding_exception(self, httpserver: HTTPServer): path = os.path.join(self.conf['redfish_root'], self.conf["backend_conf"]["subscribers_root"]) assert self.core.create_object(path, tests_template.wrong_sub) - resp = self.core.handle_event(tests_template.event) + #resp = self.core.handle_event(tests_template.event) + resp = self.core.event_handler.new_event(tests_template.event) assert len(resp) == 0 def test_event_forwarding_2(self, httpserver: HTTPServer): httpserver.expect_request("/").respond_with_data("OK") - resp = self.core.handle_event(tests_template.event_resource_type_system) - #print('RESP ', resp) + #resp = self.core.handle_event(tests_template.event_resource_type_system) + resp = self.core.event_handler.new_event(tests_template.event_resource_type_system) assert len(resp) == 1 def test_resource_created_event_no_context_exception(self): @@ -209,10 +213,10 @@ def test_agent_register(self, httpserver: HTTPServer): resp = self.core.handle_event(tests_template.reg_event) assert len(httpserver.log) == 2 - assert len(resp) == 0 + assert len(resp) == 2 - def test_agent_upload(self, httpserver: HTTPServer): - #pdb.set_trace() + def test_agent_upload(self, httpserver: HTTPServer, caplog): + # arm the httpserver with agent's response to GET on OriginOfCondition connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1") httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_pytest1) @@ -226,11 +230,156 @@ def test_agent_upload(self, httpserver: HTTPServer): connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1/Switches/Pytest1") httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_switch_pytest1) resp = self.core.handle_event(tests_template.upload_event) + upload_list = test_utils.check_uploaded_objects(tests_template.upload_event, self.conf['redfish_root']) + assert len(upload_list) == 3 assert len(httpserver.log) == 4 # TODO # should verify the two objects got uploaded and written to the Sunfish DB - assert len(resp) == 0 + assert len(resp) == 2 + #assert "Sunfish Internal Event Generation function Error" in caplog.text + + + def test_event_resourceChanged(self, httpserver: HTTPServer): + # requires test_agent_upload runs successfully before calling this test + # + # This test checks two things: + # 1) when a new Subscription is created, Sunfish core creates a ResourceCreated Event + # which will get forwarded to the new subscriber + # 2) then when a ResourceUpdated event is sent to Sunfish core, + # Sunfish fetches the updated resource (OriginOfCondition) + # and also sends a ResourceUpdated event to the new subscriber + # + # arm the httpserver with subscriber's response to POST of the create_object of subscription sub4 + connection_path = os.path.join(self.conf['redfish_root'], "/") + httpserver.expect_ordered_request(connection_path, method="POST").respond_with_data("OK") + # install another subscriber for ResourceEvents, (this will also trigger a ResourceCreated event!) + path = os.path.join(self.conf['redfish_root'], self.conf["backend_conf"]["subscribers_root"]) + assert self.core.create_object(path, tests_template.sub4) + + # arm the httpserver with agent's response to GET on OriginOfCondition + connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1/Switches/Pytest1") + httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_switch_pytest1_modified) + # arm the httpserver with subscriber's response to POST on Eventlistener + connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1/Switches/Pytest1") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + # send Sunfish core a ResourceUpdated event naming a Switch as OriginOfCondition + resp = self.core.handle_event(tests_template.update_switch) + assert len(httpserver.log) == 3 + # TODO + # should verify the updated switch got uploaded and written to the Sunfish DB + + # handle_event() will return list of UUIDs to which the event was forwarded + assert len(resp) == 2 + + def test_2nd_agent_upload(self, httpserver: HTTPServer, caplog): + + # this tests a 2nd agent upload of a CXL fabric with the same names as + # the previous agent's upload. The 2nd agent's upload has a different Fabric UUID + # so Sunfish will RENAME the 2nd agent's fabric object + # Because this test runs AFTER test_event_resourceChanged, there will be events + # issued for the creation of 3 new objects uploaded from the 2nd agent + # + # this test requires the previous test_agent_upload and test_event_resourceChanged + # both completed successfully + # + # arm the httpserver with agent's response to GET on OriginOfCondition + connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1") + httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_pytest1b) + # the above is actually retrieved again at start of recursive fetch (upload) + connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1") + httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_pytest1b) + # arm the httpserver with agent's response to GET on subordinate Switches collection + connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1/Switches") + httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_switch_collection) + # arm the httpserver with agent's response to GET on Switch object + connection_path = os.path.join(self.conf['redfish_root'], "Fabrics/Pytest1/Switches/Pytest1") + httpserver.expect_ordered_request(connection_path, method="GET").respond_with_json(tests_template.fabrics_switch_pytest1b) + # arm the httpserver with client's response to the associated 3 ResourceCreated Events + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + #pdb.set_trace() + resp = self.core.handle_event(tests_template.upload_event2) + # check that 2nd agent uploaded the correct number of objects (3) + upload_list = test_utils.check_uploaded_objects(tests_template.upload_event2, self.conf['redfish_root']) + assert len(upload_list) == 3 + assert len(httpserver.log) == 7 + # TODO + # should verify the two objects got uploaded and written to the Sunfish DB + # assert test_utils.check_delete(system_url) == True + + assert len(resp) == 2 + #assert "Sunfish Internal Event Generation function Error" in caplog.text + + def test_event_resourceDeleted(self, httpserver: HTTPServer): + # requires test_agent_upload runs successfully before calling this test + # + # This test checks: + # 1) when a ResourceDeleted event is sent to Sunfish core, + # Sunfish removes the deleted resource (OriginOfCondition), + # Sunfish removes all the subordinates of the deleted resource + # and traverses the whole database and removing links to the deleted resources + # and also sends a ResourceDeleted event to any ResourceEvents subscribers + # AND sends a ResourceChanged event for any + # + # arm the httpserver with subscriber's blind response to receipt of Events: + # for delete of /Fabrics/Pytest1 + # for delete of /Fabrics/Pytest1/Switches + # for delete of /Fabrics/Pytest1/Switches/Pytest1 + # for changes to /Fabrics + # for changes to /AggregationService/AggregationSources/xxxxxxxx + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + # send Sunfish core a ResourceDeleted event naming a fabric as OriginOfCondition + resp = self.core.handle_event(tests_template.delete_fabric_event) + # deleted objects should be removed from uploaded_objects list in the aggregationSource + upload_list = test_utils.check_uploaded_objects(tests_template.delete_fabric_event, self.conf['redfish_root']) + assert len(upload_list) == 0 + assert len(httpserver.log) == 5 + # TODO + # should verify the deleted fabric got removed from the Sunfish DB + + # handle_event() will return list of UUIDs to which the event was forwarded + assert len(resp) == 2 + + + def test_event_resourceDeleted2(self, httpserver: HTTPServer): + # requires test_agent_upload runs successfully before calling this test + # + # This test checks: + # 1) when a ResourceDeleted event is sent to Sunfish core, + # Sunfish removes the deleted resource (OriginOfCondition), + # Sunfish removes all the subordinates of the deleted resource + # and traverses the whole database and removing links to the deleted resources + # and also sends a ResourceDeleted event to any ResourceEvents subscribers + # AND sends a ResourceChanged event for any + # + # arm the httpserver with subscriber's blind response to receipt of Events: + # for delete of /Fabrics/Pytest1 + # for delete of /Fabrics/Pytest1/Switches + # for delete of /Fabrics/Pytest1/Switches/Pytest1 + # for changes to /Fabrics + # for changes to /AggregationService/AggregationSources/xxxxxxxx + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + httpserver.expect_ordered_request("/", method="POST").respond_with_data("OK") + # send Sunfish core a ResourceDeleted event naming a fabric as OriginOfCondition + # this is the renamed fabric for 2nd agent, + resp = self.core.handle_event(tests_template.delete_fabric_event2) + # deleted objects should be removed from uploaded_objects list in the aggregationSource + upload_list = test_utils.check_uploaded_objects(tests_template.delete_fabric_event2, self.conf['redfish_root']) + assert len(upload_list) == 0 + assert len(httpserver.log) == 5 + # TODO + # should verify the deleted fabric got removed from the Sunfish DB + + assert len(resp) == 2 # deletes all the subscriptions @pytest.mark.order("last") diff --git a/tests/test_utils.py b/tests/test_utils.py index 8459640..8f0908e 100644 --- a/tests/test_utils.py +++ b/tests/test_utils.py @@ -19,6 +19,16 @@ def check_object(payload, redfish_root): path = get_resource_path(payload, redfish_root) return os.path.exists(path) and os.path.exists(os.path.join(path, 'index.json')) +def check_uploaded_objects(upload_event, redfish_root): + aggSrc_Id = upload_event["Context"] + aggSrc_path = os.path.join(os.getcwd(), 'Resources','AggregationService','AggregationSources',\ + aggSrc_Id, 'index.json') + with open(aggSrc_path, "r", encoding="utf-8") as file: + aggSrc_obj = json.load(file) + # extract any ResourcesAccessed + things_uploaded = aggSrc_obj.get("Links", {}).get("ResourcesAccessed", []) + return things_uploaded + def check_delete(path): if os.path.exists(path): print('non eliminato') @@ -29,4 +39,4 @@ def get_id(root, collection): list = os.listdir(os.path.join(os.getcwd(), root, collection)) for dir in list: if dir != '.DS_Store' and dir != 'index.json': - return dir \ No newline at end of file + return dir diff --git a/tests/tests_template.py b/tests/tests_template.py index 7418120..b409eae 100644 --- a/tests/tests_template.py +++ b/tests/tests_template.py @@ -248,7 +248,43 @@ "OriginResources": [{ "@odata.id": "/redfish/v1/Systems/1" }], - "SubordinateResources": "True" + "SubordinateResources": True +} + + +sub4 = { + "@odata.type": "#EventDestination.EventDestination", + "Destination": "http://localhost:8080", + "EventFormatType": "Event", + "RegistryPrefixes": [ + "ResourceEvent" + ] + , + "ExcludeRegistryPrefixes": [ + "Basic" + ], + "OriginResources": [{ + "@odata.id": "/redfish/v1/Fabrics/Pytest1" + }], + "SubordinateResources": True +} + + +sub5 = { + "@odata.type": "#EventDestination.EventDestination", + "Destination": "http://localhost:8080", + "EventFormatType": "Event", + "RegistryPrefixes": [ + "ResourceEvent" + ] + , + "ExcludeRegistryPrefixes": [ + "Basic" + ], + "OriginResources": [{ + "@odata.id": "/redfish/v1/Fabrics" + }], + "SubordinateResources": True } wrong_sub = { @@ -347,7 +383,7 @@ aggregation_source = { "@Redfish.Copyright": "Copyright 2014-2021 SNIA. All rights reserved.", "@odata.id": "/redfish/v1/AggregationService/AggregationSources/afd9e24c-20d1-479e-be24-4ad6a62f7197", - "@odata.type": "#AggregationSource.v1_2_afd9e24c-20d1-479e-be24-4ad6a62f7197.AggregationSource", + "@odata.type": "#AggregationSource.v1_2.AggregationSource", "HostName": "http://localhost:8080", "Id": "afd9e24c-20d1-479e-be24-4ad6a62f7197", "Links": { @@ -496,6 +532,7 @@ "Name": "AggregationSourceDiscovered", "Context": "", "Events": [ { + "EventId":"0102", "Severity": "Ok", "Message": "A aggregation source of connection method", "MessageId": "ResourceEvent.1.x.AggregationSourceDiscovered", @@ -511,6 +548,7 @@ "Name": "ResourceCreated", "Context": "feb7bb58-83f2-4945-a798-7c0811a29955", "Events": [ { + "EventId":"0506", "Severity": "Ok", "Message": "A new Fabric resource created", "MessageId": "ResourceEvent.1.x.ResourceCreated", @@ -521,6 +559,56 @@ } ] } + +upload_event2 = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "ResourceCreated", + "Context": "feb0bb58-83f2-4945-a798-7c0811a29955", + "Events": [ { + "EventId":"0507", + "Severity": "Ok", + "Message": "A new Fabric resource created", + "MessageId": "ResourceEvent.1.x.ResourceCreated", + "MessageArgs": [ "Redfish", "http://127.0.0.1:8080" ], + "OriginOfCondition": { + "@odata.id": "/redfish/v1/Fabrics/Pytest1" + } + } ] +} + +delete_fabric_event = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "ResourceDeleted", + "Context": "feb7bb58-83f2-4945-a798-7c0811a29955", + "Events": [ { + "EventId":"2121", + "Severity": "Ok", + "Message": "A Fabric resource deleted", + "MessageId": "ResourceEvent.1.x.ResourceDeleted", + "MessageArgs": [ "Redfish", "http://127.0.0.1:8080" ], + "OriginOfCondition": { + "@odata.id": "/redfish/v1/Fabrics/Pytest1" + } + } ] +} + + +delete_fabric_event2 = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "ResourceDeleted", + "Context": "feb0bb58-83f2-4945-a798-7c0811a29955", + "Events": [ { + "EventId":"2123", + "Severity": "Ok", + "Message": "A Fabric resource deleted", + "MessageId": "ResourceEvent.1.x.ResourceDeleted", + "MessageArgs": [ "Redfish", "http://127.0.0.1:8080" ], + "OriginOfCondition": { + "@odata.id": "/redfish/v1/Fabrics/Pytest1" + } + } ] +} + connection_method_pytest2 = { "@odata.id": "/redfish/v1/AggregationService/ConnectionMethods/Pytest2", "@odata.type": "#ConnectionMethod.v1_1_0.ConnectionMethod", @@ -555,6 +643,24 @@ "UUID": "1af883ee-d4b7-400d-afb2-536d2c2aea31" } + +fabrics_pytest1b = { + "@odata.id": "/redfish/v1/Fabrics/Pytest1", + "@odata.type": "#Fabric.v1_3_0.Fabric", + "Description": "Pytest CXL Fabric", + "FabricType": "CXL", + "Id": "CXL", + "Name": "CXL Fabric", + "Status": { + "Health": "OK", + "State": "Enabled" + }, + "Switches": { + "@odata.id": "/redfish/v1/Fabrics/Pytest1/Switches" + }, + "UUID": "1af883ee-d4b7-400d-afb2-536d2c2aea32" +} + fabrics_switch_collection = { "@odata.id": "/redfish/v1/Fabrics/Pytest1/Switches", "@odata.type": "#SwitchesCollection.SwitchesCollection", @@ -588,3 +694,51 @@ "UUID": "bfd8ca1e-b5b5-4c0f-a3cf-4d741b85b1fb" } +fabrics_switch_pytest1b = { + "@odata.id": "/redfish/v1/Fabrics/Pytest1/Switches/Pytest1", + "@odata.type": "#Switch.v1_9_1.Switch", + "CXL": { + "MaxVCSsSupported": 4, + "TotalNumbervPPBs": 4, + "VCS": { + "HDMDecoders": 12 + } + }, + "Id": "Pytest1", + "Name": "CXL Fabric Switch", + "Status": { + "Health": "OK", + "HealthRollup": "OK", + "State": "Enabled" + }, + "SwitchType": "CXL", + "UUID": "bfd8ca1e-b5b5-4c0f-a3cf-4d741b85b100" +} + +update_switch = { + "@odata.type": "#Event.v1_7_0.Event", + "Name": "ResourceChanged", + "Context": "feb7bb58-83f2-4945-a798-7c0811a29955", + "Events": [ { + "EventId":"0809", + "Severity": "Ok", + "Message": "A new Status for Switch resource", + "MessageId": "ResourceEvent.1.x.ResourceChanged", + "MessageArgs": [ "Redfish", "http://127.0.0.1:8080" ], + "OriginOfCondition": { + "@odata.id": "/redfish/v1/Fabrics/Pytest1/Switches/Pytest1" + } + } ] +} + +fabrics_switch_pytest1_modified = { + "@odata.id": "/redfish/v1/Fabrics/Pytest1/Switches/Pytest1", + "@odata.type": "#Switch.v1_9_1.Switch", + "Id": "Pytest1", + "Status": { + "Health": "OK", + "HealthRollup": "OK", + "State": "Disabled" + } +} +