-
-
Notifications
You must be signed in to change notification settings - Fork 0
/
main.py
245 lines (193 loc) · 8.14 KB
/
main.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
from pathlib import Path
import shutil
from syftbox.lib import Client
import torch
import torch.nn as nn
from torch.utils.data import DataLoader, TensorDataset
import json
from typing_extensions import Optional
API_NAME = "model_aggregator"
TEST_DATASET_NAME = "mnist_dataset.pt"
SAMPLE_TEST_DATASET_PATH = Path("./samples/test_data") / TEST_DATASET_NAME
# Exception name to indicate the state cannot advance
# as there are some pre-requisites that are not met
class StateNotReady(Exception):
pass
class SimpleNN(nn.Module):
def __init__(self):
super(SimpleNN, self).__init__()
self.fc1 = nn.Linear(28 * 28, 128)
self.fc2 = nn.Linear(128, 10)
def forward(self, x):
x = x.view(-1, 28 * 28)
x = torch.relu(self.fc1(x))
x = self.fc2(x)
return x
def get_app_private_data(client: Client, api_name: str) -> Path:
"""
Returns the private data directory of the app
"""
return client.workspace.data_dir / "private" / api_name
def init_aggregator(client: Client) -> None:
"""
Creates the `model_aggregator` api in the `api_data` folder
with the following structure:
```
api_data
└── model_aggregator
└── launch
└── running
└── done
```
"""
model_aggregator = client.api_data(API_NAME)
for folder in ["launch", "running", "done"]:
model_aggregator_folder = model_aggregator / folder
model_aggregator_folder.mkdir(parents=True, exist_ok=True)
# Create the private data directory for the app
# This is where the private test data will be stored
app_pvt_dir = get_app_private_data(client, API_NAME)
app_pvt_dir.mkdir(parents=True, exist_ok=True)
# Copy the test dataset to the private data directory
test_dataset_path = app_pvt_dir / TEST_DATASET_NAME
if not test_dataset_path.is_file():
print(f"No test dataset found. Please copy sample test dataset from \n"
f"{SAMPLE_TEST_DATASET_PATH.resolve()} to {test_dataset_path.resolve()}")
# shutil.copy(SAMPLE_TEST_DATASET_PATH, test_dataset_path)
def launch_aggregator(client: Client) -> None:
"""
Iterates over the launch folder and copies the participants.json file
to the running folder
We look for the participants.json file in the launch folder
"""
launch_folder = client.api_data(API_NAME) / "launch"
running_folder = client.api_data(API_NAME) / "running"
participants_json = launch_folder / "participants.json"
if participants_json.is_file():
print("Copying participants.json to running folder")
shutil.move(participants_json, running_folder)
def get_model_files(path: Path) -> list[Path]:
return list(path.glob("trained_mnist_label_*.pt"))
def aggregate_models(client: Client) -> None:
"""
Iterates over the running folder and tries to advance it
It loads in the participants.json file and aggregates the models
"""
running_folder = client.api_data(API_NAME) / "running"
participants_json = running_folder / "participants.json"
if not participants_json.is_file():
raise StateNotReady("participants.json file not found in the running folder")
with open(participants_json, "r") as f:
participants = json.load(f)["participants"]
model_output_path = running_folder / "global_model.pt"
global_model = SimpleNN()
global_model_state_dict = global_model.state_dict()
aggregated_model_weights = {}
n_peers = len(participants)
aggregated_peers = []
missing_peers = []
for email in participants:
their_public_folder: Path = client.datasites / email / "public"
their_model_files: list[Path] = get_model_files(their_public_folder)
if len(their_model_files) == 0:
print(f"No models found for {email} in '{their_public_folder}'")
missing_peers.append(email)
continue
for model_file in their_model_files:
print(f"Aggregating model '{model_file.name} from {email}")
model_file = their_public_folder / model_file
user_model_state = torch.load(model_file, weights_only=True)
for key in global_model_state_dict.keys():
# If user model has a different architecture than my global model.
# Skip it
if user_model_state.keys() != global_model_state_dict.keys():
print(
f"Model {model_file.name} from {email} has an invalid architecture"
)
continue
if aggregated_model_weights.get(key, None) is None:
aggregated_model_weights[key] = user_model_state[key] * (
1 / n_peers
)
else:
aggregated_model_weights[key] += user_model_state[key] * (
1 / n_peers
)
aggregated_peers.append(email)
if not aggregated_model_weights:
return (None, None)
global_model.load_state_dict(aggregated_model_weights)
torch.save(global_model.state_dict(), model_output_path)
return (participants, missing_peers)
def calculate_model_accuracy(global_model_path: Path, dataset_path: Path) -> float:
model = SimpleNN()
model.load_state_dict(torch.load(global_model_path, weights_only=True))
model.eval()
# load the saved mnist subset
images, labels = torch.load(dataset_path)
# create a tensordataset
dataset = TensorDataset(images, labels)
# create a dataloader for the dataset
data_loader = DataLoader(dataset, batch_size=64, shuffle=True)
correct = 0
total = 0
with torch.no_grad():
for images, labels in data_loader:
outputs = model(images)
_, predicted = torch.max(outputs.data, 1)
total += labels.size(0)
correct += (predicted == labels).sum().item()
accuracy = 100 * correct / total
return accuracy
def evaluate_global_model(
client: Client, participants: Optional[list[str]], missing_peers: Optional[list[str]]
) -> None:
if not participants:
raise StateNotReady("No models aggregated. Skipping evaluation")
running_folder = client.api_data(API_NAME) / "running"
global_model_path = running_folder / "global_model.pt"
if not global_model_path.is_file():
raise StateNotReady(
f"ERROR: global model path ({global_model_path}) does not exist"
)
# Evaluate the global model
test_dataset_path: Path = get_app_private_data(client, API_NAME) / TEST_DATASET_NAME
accuracy: float = calculate_model_accuracy(global_model_path, test_dataset_path)
# Write the accuracy to an results.json file
results = {
"accuracy": accuracy,
"participants": participants,
"missing_peers": missing_peers,
}
print("Accuracy Results:", results)
return results
def save_result(results: dict):
running_folder = client.api_data(API_NAME) / "running"
results_path = running_folder / "results.json"
participants_json = running_folder / "participants.json"
with open(results_path, "w") as f:
json.dump(results, f, indent=4)
# If no missing peers, move the global model and results.json to the done folder
done_folder = client.api_data(API_NAME) / "done"
model_path = running_folder / "global_model.pt"
if not missing_peers:
shutil.move(participants_json, done_folder)
shutil.move(model_path, done_folder)
shutil.move(results_path, done_folder)
if __name__ == "__main__":
client = Client.load()
try:
# Step 1: Init the Aggregator API
init_aggregator(client)
# Step 2: Launch the Aggregator
# Iterates over the launch folder and looks for the participants.json file
launch_aggregator(client)
# Step 3: Aggregate the Models
participants, missing_peers = aggregate_models(client)
# Step 4: Evaluate model
results = evaluate_global_model(client, participants, missing_peers)
# Step 5: Save the results
save_result(results)
except StateNotReady as e:
print(f"StateNotReady: {e}")
exit(0)