forked from NVIDIA/NVFlare
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Fixes # . ### Description This PR adds following features. - A base class (EdgeTaskExecutor) for developing executors for processing edge tasks. - The EdgeTaskDispatcher (ETD) that is to be installed on CP and is responsible for dispatching edge requests to the appropriate CJ. The ETD determines the right CJ based on the job's "edge_method" meta property against device's capabilities. - Added some test components for determining active devices from simulated device activities. - Added a more elegant way for registering and handling events. ### Types of changes <!--- Put an `x` in all the boxes that apply, and remove the not applicable items --> - [x] Non-breaking change (fix or new feature that would not break existing functionality). - [ ] Breaking change (fix or new feature that would cause existing functionality to change). - [ ] New tests added to cover the changes. - [ ] Quick tests passed locally by running `./runtest.sh`. - [ ] In-line docstrings updated. - [ ] Documentation updated.
- Loading branch information
1 parent
0a8d273
commit 4bb222c
Showing
29 changed files
with
939 additions
and
86 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,13 @@ | ||
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved. | ||
# | ||
# Licensed under the Apache License, Version 2.0 (the "License"); | ||
# you may not use this file except in compliance with the License. | ||
# You may obtain a copy of the License at | ||
# | ||
# http://www.apache.org/licenses/LICENSE-2.0 | ||
# | ||
# Unless required by applicable law or agreed to in writing, software | ||
# distributed under the License is distributed on an "AS IS" BASIS, | ||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
# See the License for the specific language governing permissions and | ||
# limitations under the License. |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,13 @@ | ||
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved. | ||
# | ||
# Licensed under the Apache License, Version 2.0 (the "License"); | ||
# you may not use this file except in compliance with the License. | ||
# You may obtain a copy of the License at | ||
# | ||
# http://www.apache.org/licenses/LICENSE-2.0 | ||
# | ||
# Unless required by applicable law or agreed to in writing, software | ||
# distributed under the License is distributed on an "AS IS" BASIS, | ||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
# See the License for the specific language governing permissions and | ||
# limitations under the License. |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,36 @@ | ||
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved. | ||
# | ||
# Licensed under the Apache License, Version 2.0 (the "License"); | ||
# you may not use this file except in compliance with the License. | ||
# You may obtain a copy of the License at | ||
# | ||
# http://www.apache.org/licenses/LICENSE-2.0 | ||
# | ||
# Unless required by applicable law or agreed to in writing, software | ||
# distributed under the License is distributed on an "AS IS" BASIS, | ||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
# See the License for the specific language governing permissions and | ||
# limitations under the License. | ||
from nvflare.apis.fl_context import FLContext | ||
from nvflare.apis.shareable import Shareable | ||
from nvflare.app_common.abstract.aggregator import Aggregator | ||
|
||
|
||
class EdgeSurveyAggregator(Aggregator): | ||
def __init__(self): | ||
Aggregator.__init__(self) | ||
self.num_devices = 0 | ||
|
||
def accept(self, shareable: Shareable, fl_ctx: FLContext) -> bool: | ||
self.log_info(fl_ctx, f"accepting: {shareable}") | ||
num_devices = shareable.get("num_devices") | ||
if num_devices: | ||
self.num_devices += num_devices | ||
return True | ||
|
||
def reset(self, fl_ctx: FLContext): | ||
self.num_devices = 0 | ||
|
||
def aggregate(self, fl_ctx: FLContext) -> Shareable: | ||
self.log_info(fl_ctx, f"aggregating final result: {self.num_devices}") | ||
return Shareable({"num_devices": self.num_devices}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,36 @@ | ||
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved. | ||
# | ||
# Licensed under the Apache License, Version 2.0 (the "License"); | ||
# you may not use this file except in compliance with the License. | ||
# You may obtain a copy of the License at | ||
# | ||
# http://www.apache.org/licenses/LICENSE-2.0 | ||
# | ||
# Unless required by applicable law or agreed to in writing, software | ||
# distributed under the License is distributed on an "AS IS" BASIS, | ||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
# See the License for the specific language governing permissions and | ||
# limitations under the License. | ||
from nvflare.fuel.f3.cellnet.defs import ReturnCode as CellReturnCode | ||
|
||
|
||
class Status(CellReturnCode): | ||
NO_TASK = "no_task" | ||
NO_JOB = "no_job" | ||
|
||
|
||
class EdgeProtoKey: | ||
STATUS = "status" | ||
DATA = "data" | ||
|
||
|
||
class EdgeContextKey: | ||
JOB_ID = "__edge_job_id__" | ||
EDGE_CAPABILITIES = "__edge_capabilities__" | ||
REQUEST_FROM_EDGE = "__request_from_edge__" | ||
REPLY_TO_EDGE = "__reply_to_edge__" | ||
|
||
|
||
class EventType: | ||
EDGE_REQUEST_RECEIVED = "_edge_request_received" | ||
EDGE_JOB_REQUEST_RECEIVED = "_edge_job_request_received" |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,13 @@ | ||
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved. | ||
# | ||
# Licensed under the Apache License, Version 2.0 (the "License"); | ||
# you may not use this file except in compliance with the License. | ||
# You may obtain a copy of the License at | ||
# | ||
# http://www.apache.org/licenses/LICENSE-2.0 | ||
# | ||
# Unless required by applicable law or agreed to in writing, software | ||
# distributed under the License is distributed on an "AS IS" BASIS, | ||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
# See the License for the specific language governing permissions and | ||
# limitations under the License. |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,59 @@ | ||
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved. | ||
# | ||
# Licensed under the Apache License, Version 2.0 (the "License"); | ||
# you may not use this file except in compliance with the License. | ||
# You may obtain a copy of the License at | ||
# | ||
# http://www.apache.org/licenses/LICENSE-2.0 | ||
# | ||
# Unless required by applicable law or agreed to in writing, software | ||
# distributed under the License is distributed on an "AS IS" BASIS, | ||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
# See the License for the specific language governing permissions and | ||
# limitations under the License. | ||
from nvflare.apis.controller_spec import ClientTask, Task | ||
from nvflare.apis.fl_context import FLContext | ||
from nvflare.apis.impl.controller import Controller | ||
from nvflare.apis.shareable import Shareable | ||
from nvflare.apis.signal import Signal | ||
|
||
|
||
class EdgeSurveyController(Controller): | ||
def __init__(self, num_rounds: int, timeout: int): | ||
Controller.__init__(self) | ||
self.num_rounds = num_rounds | ||
self.timeout = timeout | ||
|
||
def start_controller(self, fl_ctx: FLContext): | ||
pass | ||
|
||
def stop_controller(self, fl_ctx: FLContext): | ||
pass | ||
|
||
def control_flow(self, abort_signal: Signal, fl_ctx: FLContext): | ||
for r in range(self.num_rounds): | ||
task = Task( | ||
name="survey", | ||
data=Shareable(), | ||
timeout=self.timeout, | ||
) | ||
|
||
self.broadcast_and_wait( | ||
task=task, | ||
min_responses=2, | ||
wait_time_after_min_received=0, | ||
fl_ctx=fl_ctx, | ||
abort_signal=abort_signal, | ||
) | ||
|
||
total_devices = 0 | ||
for ct in task.client_tasks: | ||
assert isinstance(ct, ClientTask) | ||
result = ct.result | ||
assert isinstance(result, Shareable) | ||
self.log_info(fl_ctx, f"result from client {ct.client.name}: {result}") | ||
count = result.get("num_devices") | ||
if count: | ||
total_devices += count | ||
|
||
self.log_info(fl_ctx, f"total devices in round {r}: {total_devices}") |
Oops, something went wrong.