Loading proto/forecaster.proto +13 −38 Original line number Diff line number Diff line Loading @@ -18,52 +18,27 @@ package forecaster; import "context.proto"; service ForecasterService { rpc ComputeTopologyForecast(ForecastTopology) returns (ListLinkCapacity) {} rpc ComputeLinkForecast (context.LinkId ) returns (LinkCapacity ) {} rpc ForecastLinkCapacity (ForecastLinkCapacityRequest ) returns (ForecastLinkCapacityReply ) {} rpc ForecastTopologyCapacity(ForecastTopologyCapacityRequest) returns (ForecastTopologyCapacityReply) {} } message ForecastRequest { oneof uuid { context.TopologyId topology_id = 1; context.LinkId link_id = 2; } context.TopogyId topology_id = 1; float forecast_window_seconds = 2; } message SingleForecast { context.Timestamp timestamp= 1; double value = 2; } message Forecast { oneof uuid { context.TopologyId topologyId= 1; context.LinkId linkId = 2; } repeated SingleForecast forecast = 3; } enum AvailabilityPredictionEnum { FORECASTED_AVAILABILITY = 0; FORECASTED_UNAVAILABILITY = 1; } message ForecastPrediction { AvailabilityPredictionEnum prediction = 1; message ForecastLinkCapacityRequest { context.LinkId link_id = 1; double forecast_window_seconds = 2; } message LinkCapacity { message ForecastLinkCapacityReply { context.LinkId link_id = 1; float total_capacity_gbps = 2; float current_capacity_gbps = 3; float forecasted_capacity_gbps = 4; } message ListLinkCapacity { repeated LinkCapacity forecasted_link_capacities = 1; message ForecastTopologyCapacityRequest { context.TopologyId topology_id = 1; double forecast_window_seconds = 2; } message ForecastTopologyCapacityReply { repeated ForecastLinkCapacityReply link_capacities = 1; } src/common/Constants.py +2 −0 Original line number Diff line number Diff line Loading @@ -57,6 +57,7 @@ class ServiceNameEnum(Enum): OPTICALATTACKMITIGATOR = 'opticalattackmitigator' CACHING = 'caching' TE = 'te' FORECASTER = 'forecaster' # Used for test and debugging only DLT_GATEWAY = 'dltgateway' Loading @@ -82,6 +83,7 @@ DEFAULT_SERVICE_GRPC_PORTS = { ServiceNameEnum.INTERDOMAIN .value : 10010, ServiceNameEnum.PATHCOMP .value : 10020, ServiceNameEnum.TE .value : 10030, ServiceNameEnum.FORECASTER .value : 10040, # Used for test and debugging only ServiceNameEnum.DLT_GATEWAY .value : 50051, Loading src/forecaster/Dockerfile +5 −5 Original line number Diff line number Diff line Loading @@ -54,17 +54,17 @@ RUN rm *.proto RUN find . -type f -exec sed -i -E 's/(import\ .*)_pb2/from . \1_pb2/g' {} \; # Create component sub-folders, get specific Python packages RUN mkdir -p /var/teraflow/device WORKDIR /var/teraflow/device COPY src/device/requirements.in requirements.in RUN mkdir -p /var/teraflow/forecaster WORKDIR /var/teraflow/forecaster COPY src/forecaster/requirements.in requirements.in RUN pip-compile --quiet --output-file=requirements.txt requirements.in RUN python3 -m pip install -r requirements.txt # Add component files into working directory WORKDIR /var/teraflow COPY src/context/. context/ COPY src/device/. device/ COPY src/forecaster/. forecaster/ COPY src/monitoring/. monitoring/ # Start the service ENTRYPOINT ["python", "-m", "device.service"] ENTRYPOINT ["python", "-m", "forecaster.service"] src/forecaster/client/ForecasterClient.py +16 −35 Original line number Diff line number Diff line Loading @@ -15,9 +15,11 @@ import grpc, logging from common.Constants import ServiceNameEnum from common.Settings import get_service_host, get_service_port_grpc from common.proto.context_pb2 import Device, DeviceConfig, DeviceId, Empty from common.proto.device_pb2 import MonitoringSettings from common.proto.device_pb2_grpc import DeviceServiceStub from common.proto.forecaster_pb2 import ( ForecastLinkCapacityReply, ForecastLinkCapacityRequest, ForecastTopologyCapacityReply, ForecastTopologyCapacityRequest ) from common.proto.forecaster_pb2_grpc import ForecasterServiceStub from common.tools.client.RetryDecorator import retry, delay_exponential from common.tools.grpc.Tools import grpc_message_to_json_string Loading @@ -28,8 +30,8 @@ RETRY_DECORATOR = retry(max_retries=MAX_RETRIES, delay_function=DELAY_FUNCTION, class ForecasterClient: def __init__(self, host=None, port=None): if not host: host = get_service_host(ServiceNameEnum.DEVICE) if not port: port = get_service_port_grpc(ServiceNameEnum.DEVICE) if not host: host = get_service_host(ServiceNameEnum.FORECASTER) if not port: port = get_service_port_grpc(ServiceNameEnum.FORECASTER) self.endpoint = '{:s}:{:s}'.format(str(host), str(port)) LOGGER.debug('Creating channel to {:s}...'.format(str(self.endpoint))) self.channel = None Loading @@ -39,7 +41,7 @@ class ForecasterClient: def connect(self): self.channel = grpc.insecure_channel(self.endpoint) self.stub = DeviceServiceStub(self.channel) self.stub = ForecasterServiceStub(self.channel) def close(self): if self.channel is not None: self.channel.close() Loading @@ -47,36 +49,15 @@ class ForecasterClient: self.stub = None @RETRY_DECORATOR def AddDevice(self, request : Device) -> DeviceId: LOGGER.debug('AddDevice request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.AddDevice(request) LOGGER.debug('AddDevice result: {:s}'.format(grpc_message_to_json_string(response))) def ForecastLinkCapacity(self, request : ForecastLinkCapacityRequest) -> ForecastLinkCapacityReply: LOGGER.debug('ForecastLinkCapacity request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.ForecastLinkCapacity(request) LOGGER.debug('ForecastLinkCapacity result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def ConfigureDevice(self, request : Device) -> DeviceId: LOGGER.debug('ConfigureDevice request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.ConfigureDevice(request) LOGGER.debug('ConfigureDevice result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def DeleteDevice(self, request : DeviceId) -> Empty: LOGGER.debug('DeleteDevice request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.DeleteDevice(request) LOGGER.debug('DeleteDevice result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def GetInitialConfig(self, request : DeviceId) -> DeviceConfig: LOGGER.debug('GetInitialConfig request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.GetInitialConfig(request) LOGGER.debug('GetInitialConfig result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def MonitorDeviceKpi(self, request : MonitoringSettings) -> Empty: LOGGER.debug('MonitorDeviceKpi request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.MonitorDeviceKpi(request) LOGGER.debug('MonitorDeviceKpi result: {:s}'.format(grpc_message_to_json_string(response))) def ForecastTopologyCapacity(self, request : ForecastTopologyCapacityRequest) -> ForecastTopologyCapacityReply: LOGGER.debug('ForecastTopologyCapacity request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.ForecastTopologyCapacity(request) LOGGER.debug('ForecastTopologyCapacity result: {:s}'.format(grpc_message_to_json_string(response))) return response src/forecaster/greet_client.pydeleted 100644 → 0 +0 −23 Original line number Diff line number Diff line import greet_pb2_grpc import greet_pb2 import grpc #from tests import greet_pb2 #from tests import greet_pb2_grpc import datetime as dt def run(): with grpc.insecure_channel('localhost:50051') as channel: stub = greet_pb2_grpc.GreeterStub(channel) print("The stub has been created") forecastRequest = greet_pb2.ForecastTopology(forecast_window_seconds = 36000, topology_id = "1" ) print("The request has been sent to the client") hello_reply = stub.ComputeTopologyForecast(forecastRequest) print("The response has been received") print(hello_reply) if __name__ == "__main__": run() No newline at end of file Loading
proto/forecaster.proto +13 −38 Original line number Diff line number Diff line Loading @@ -18,52 +18,27 @@ package forecaster; import "context.proto"; service ForecasterService { rpc ComputeTopologyForecast(ForecastTopology) returns (ListLinkCapacity) {} rpc ComputeLinkForecast (context.LinkId ) returns (LinkCapacity ) {} rpc ForecastLinkCapacity (ForecastLinkCapacityRequest ) returns (ForecastLinkCapacityReply ) {} rpc ForecastTopologyCapacity(ForecastTopologyCapacityRequest) returns (ForecastTopologyCapacityReply) {} } message ForecastRequest { oneof uuid { context.TopologyId topology_id = 1; context.LinkId link_id = 2; } context.TopogyId topology_id = 1; float forecast_window_seconds = 2; } message SingleForecast { context.Timestamp timestamp= 1; double value = 2; } message Forecast { oneof uuid { context.TopologyId topologyId= 1; context.LinkId linkId = 2; } repeated SingleForecast forecast = 3; } enum AvailabilityPredictionEnum { FORECASTED_AVAILABILITY = 0; FORECASTED_UNAVAILABILITY = 1; } message ForecastPrediction { AvailabilityPredictionEnum prediction = 1; message ForecastLinkCapacityRequest { context.LinkId link_id = 1; double forecast_window_seconds = 2; } message LinkCapacity { message ForecastLinkCapacityReply { context.LinkId link_id = 1; float total_capacity_gbps = 2; float current_capacity_gbps = 3; float forecasted_capacity_gbps = 4; } message ListLinkCapacity { repeated LinkCapacity forecasted_link_capacities = 1; message ForecastTopologyCapacityRequest { context.TopologyId topology_id = 1; double forecast_window_seconds = 2; } message ForecastTopologyCapacityReply { repeated ForecastLinkCapacityReply link_capacities = 1; }
src/common/Constants.py +2 −0 Original line number Diff line number Diff line Loading @@ -57,6 +57,7 @@ class ServiceNameEnum(Enum): OPTICALATTACKMITIGATOR = 'opticalattackmitigator' CACHING = 'caching' TE = 'te' FORECASTER = 'forecaster' # Used for test and debugging only DLT_GATEWAY = 'dltgateway' Loading @@ -82,6 +83,7 @@ DEFAULT_SERVICE_GRPC_PORTS = { ServiceNameEnum.INTERDOMAIN .value : 10010, ServiceNameEnum.PATHCOMP .value : 10020, ServiceNameEnum.TE .value : 10030, ServiceNameEnum.FORECASTER .value : 10040, # Used for test and debugging only ServiceNameEnum.DLT_GATEWAY .value : 50051, Loading
src/forecaster/Dockerfile +5 −5 Original line number Diff line number Diff line Loading @@ -54,17 +54,17 @@ RUN rm *.proto RUN find . -type f -exec sed -i -E 's/(import\ .*)_pb2/from . \1_pb2/g' {} \; # Create component sub-folders, get specific Python packages RUN mkdir -p /var/teraflow/device WORKDIR /var/teraflow/device COPY src/device/requirements.in requirements.in RUN mkdir -p /var/teraflow/forecaster WORKDIR /var/teraflow/forecaster COPY src/forecaster/requirements.in requirements.in RUN pip-compile --quiet --output-file=requirements.txt requirements.in RUN python3 -m pip install -r requirements.txt # Add component files into working directory WORKDIR /var/teraflow COPY src/context/. context/ COPY src/device/. device/ COPY src/forecaster/. forecaster/ COPY src/monitoring/. monitoring/ # Start the service ENTRYPOINT ["python", "-m", "device.service"] ENTRYPOINT ["python", "-m", "forecaster.service"]
src/forecaster/client/ForecasterClient.py +16 −35 Original line number Diff line number Diff line Loading @@ -15,9 +15,11 @@ import grpc, logging from common.Constants import ServiceNameEnum from common.Settings import get_service_host, get_service_port_grpc from common.proto.context_pb2 import Device, DeviceConfig, DeviceId, Empty from common.proto.device_pb2 import MonitoringSettings from common.proto.device_pb2_grpc import DeviceServiceStub from common.proto.forecaster_pb2 import ( ForecastLinkCapacityReply, ForecastLinkCapacityRequest, ForecastTopologyCapacityReply, ForecastTopologyCapacityRequest ) from common.proto.forecaster_pb2_grpc import ForecasterServiceStub from common.tools.client.RetryDecorator import retry, delay_exponential from common.tools.grpc.Tools import grpc_message_to_json_string Loading @@ -28,8 +30,8 @@ RETRY_DECORATOR = retry(max_retries=MAX_RETRIES, delay_function=DELAY_FUNCTION, class ForecasterClient: def __init__(self, host=None, port=None): if not host: host = get_service_host(ServiceNameEnum.DEVICE) if not port: port = get_service_port_grpc(ServiceNameEnum.DEVICE) if not host: host = get_service_host(ServiceNameEnum.FORECASTER) if not port: port = get_service_port_grpc(ServiceNameEnum.FORECASTER) self.endpoint = '{:s}:{:s}'.format(str(host), str(port)) LOGGER.debug('Creating channel to {:s}...'.format(str(self.endpoint))) self.channel = None Loading @@ -39,7 +41,7 @@ class ForecasterClient: def connect(self): self.channel = grpc.insecure_channel(self.endpoint) self.stub = DeviceServiceStub(self.channel) self.stub = ForecasterServiceStub(self.channel) def close(self): if self.channel is not None: self.channel.close() Loading @@ -47,36 +49,15 @@ class ForecasterClient: self.stub = None @RETRY_DECORATOR def AddDevice(self, request : Device) -> DeviceId: LOGGER.debug('AddDevice request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.AddDevice(request) LOGGER.debug('AddDevice result: {:s}'.format(grpc_message_to_json_string(response))) def ForecastLinkCapacity(self, request : ForecastLinkCapacityRequest) -> ForecastLinkCapacityReply: LOGGER.debug('ForecastLinkCapacity request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.ForecastLinkCapacity(request) LOGGER.debug('ForecastLinkCapacity result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def ConfigureDevice(self, request : Device) -> DeviceId: LOGGER.debug('ConfigureDevice request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.ConfigureDevice(request) LOGGER.debug('ConfigureDevice result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def DeleteDevice(self, request : DeviceId) -> Empty: LOGGER.debug('DeleteDevice request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.DeleteDevice(request) LOGGER.debug('DeleteDevice result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def GetInitialConfig(self, request : DeviceId) -> DeviceConfig: LOGGER.debug('GetInitialConfig request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.GetInitialConfig(request) LOGGER.debug('GetInitialConfig result: {:s}'.format(grpc_message_to_json_string(response))) return response @RETRY_DECORATOR def MonitorDeviceKpi(self, request : MonitoringSettings) -> Empty: LOGGER.debug('MonitorDeviceKpi request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.MonitorDeviceKpi(request) LOGGER.debug('MonitorDeviceKpi result: {:s}'.format(grpc_message_to_json_string(response))) def ForecastTopologyCapacity(self, request : ForecastTopologyCapacityRequest) -> ForecastTopologyCapacityReply: LOGGER.debug('ForecastTopologyCapacity request: {:s}'.format(grpc_message_to_json_string(request))) response = self.stub.ForecastTopologyCapacity(request) LOGGER.debug('ForecastTopologyCapacity result: {:s}'.format(grpc_message_to_json_string(response))) return response
src/forecaster/greet_client.pydeleted 100644 → 0 +0 −23 Original line number Diff line number Diff line import greet_pb2_grpc import greet_pb2 import grpc #from tests import greet_pb2 #from tests import greet_pb2_grpc import datetime as dt def run(): with grpc.insecure_channel('localhost:50051') as channel: stub = greet_pb2_grpc.GreeterStub(channel) print("The stub has been created") forecastRequest = greet_pb2.ForecastTopology(forecast_window_seconds = 36000, topology_id = "1" ) print("The request has been sent to the client") hello_reply = stub.ComputeTopologyForecast(forecastRequest) print("The response has been received") print(hello_reply) if __name__ == "__main__": run() No newline at end of file