diff --git a/cometx/cli/copy.py b/cometx/cli/copy.py index f1c4c00..c8f9ec8 100644 --- a/cometx/cli/copy.py +++ b/cometx/cli/copy.py @@ -1196,18 +1196,17 @@ def update_datagrid_contents( asset_map[old_asset_id] = result["assetId"] def _log_asset( - self, experiment, path, asset_type, log_filename, assets_metadata, asset_map + self, experiment, path, asset_type, asset_data, asset_map ): - log_as_filename = assets_metadata[log_filename].get( - "logAsFileName", - None, - ) - step = assets_metadata[log_filename].get("step") - epoch = assets_metadata[log_filename].get("epoch") - old_asset_id = assets_metadata[log_filename].get("assetId") + log_as_filename = asset_data.get("logAsFileName", None) + original_filename = asset_data["fileName"] + disk_filename = asset_data.get("diskFileName", original_filename) + step = asset_data.get("step") + epoch = asset_data.get("epoch") + old_asset_id = asset_data.get("assetId") if asset_type in self.ignore: return - sanitized_filename = sanitize_filename(log_filename) + sanitized_filename = sanitize_filename(disk_filename) filename = os.path.join(path, asset_type, sanitized_filename) if not os.path.isfile(filename): @@ -1217,7 +1216,7 @@ def _log_asset( print("Missing file %r: unable to copy" % filename) return - metadata = assets_metadata[log_filename].get("metadata") + metadata = asset_data.get("metadata") metadata = json.loads(metadata) if metadata else {} if asset_type == "notebook": @@ -1286,11 +1285,11 @@ def _log_asset( metadata, binary_io, step, - log_as_filename or log_filename, + log_as_filename or original_filename, ) if result is None: print( - f"ERROR: Unable to log {asset_type} asset {log_as_filename or log_filename}; skipping" + f"ERROR: Unable to log {asset_type} asset {log_as_filename or original_filename}; skipping" ) else: asset_map[old_asset_id] = result["assetId"] @@ -1301,7 +1300,7 @@ def _log_asset( metadata, step, filename, - log_as_filename or log_filename, + log_as_filename or original_filename, old_asset_id, ) elif asset_type == "confusion-matrix": @@ -1335,28 +1334,30 @@ def _log_asset( metadata, binary_io, step, - log_as_filename or log_filename, + log_as_filename or original_filename, ) if result is None: print( - f"ERROR: Unable to log {asset_type} asset {log_as_filename or log_filename}; skipping" + f"ERROR: Unable to log {asset_type} asset {log_as_filename or original_filename}; skipping" ) else: asset_map[old_asset_id] = result["assetId"] elif asset_type == "video": - name = os.path.basename(filename) binary_io = open(filename, "rb") result = experiment.log_video( - binary_io, name=log_as_filename or name, step=step, epoch=epoch - ) # done! + binary_io, + name=log_as_filename or original_filename, + step=step, + epoch=epoch, + ) if result is None: print( - f"ERROR: Unable to log {asset_type} asset {log_as_filename or name}; skipping" + f"ERROR: Unable to log {asset_type} asset {log_as_filename or original_filename}; skipping" ) else: asset_map[old_asset_id] = result["assetId"] elif asset_type == "model-element": - dir_name = assets_metadata[log_filename].get("dir", "") + dir_name = asset_data.get("dir", "") # The dir field includes a "models/" prefix added by the # backend; strip it to get the actual model name. if dir_name.startswith("models/"): @@ -1364,18 +1365,15 @@ def _log_asset( else: model_name = dir_name binary_io = open(filename, "rb") - result = experiment._log_asset( - binary_io, - file_name=log_as_filename or log_filename, - copy_to_tmp=True, - asset_type=asset_type, + result = experiment.log_model( + name=model_name, + file_or_folder=binary_io, + file_name=log_as_filename or original_filename, metadata=metadata, - step=step, - grouping_name=model_name, ) if result is None: print( - f"ERROR: Unable to log {asset_type} asset {log_as_filename or log_filename}; skipping" + f"ERROR: Unable to log {asset_type} asset {log_as_filename or original_filename}; skipping" ) else: asset_map[old_asset_id] = result["assetId"] @@ -1386,11 +1384,11 @@ def _log_asset( metadata, filename, step, - log_as_filename or log_filename, + log_as_filename or original_filename, ) if result is None: print( - f"ERROR: Unable to log {asset_type} asset {log_as_filename or log_filename}; skipping" + f"ERROR: Unable to log {asset_type} asset {log_as_filename or original_filename}; skipping" ) else: asset_map[old_asset_id] = result["assetId"] @@ -1402,39 +1400,50 @@ def log_assets(self, experiment, path, assets_metadata): # Create mapping from old asset id to new asset id asset_map = {} # Process all of the non-nested assets first: - for log_filename in assets_metadata: - asset_type = assets_metadata[log_filename].get("type", "asset") or "asset" + for asset_data in assets_metadata: + asset_type = asset_data.get("type", "asset") or "asset" if asset_type not in ["confusion-matrix", "embeddings", "datagrid"]: - if ( - "remote" in assets_metadata[log_filename] - and assets_metadata[log_filename]["remote"] - ): - asset = assets_metadata[log_filename] - experiment.log_remote_asset( - uri=asset["link"], - remote_file_name=asset["fileName"], - step=asset["step"], - metadata=asset["metadata"], - ) + if asset_data.get("remote", False): + if asset_type == "model-element": + dir_name = asset_data.get("dir", "") + if dir_name.startswith("models/"): + model_name = dir_name[len("models/"):] + else: + model_name = dir_name + raw_metadata = asset_data.get("metadata") + metadata = json.loads(raw_metadata) if raw_metadata else None + experiment.log_remote_model( + model_name=model_name, + uri=asset_data["link"], + metadata=metadata, + sync_mode=False, + ) + else: + raw_metadata = asset_data.get("metadata") + metadata = json.loads(raw_metadata) if raw_metadata else None + experiment.log_remote_asset( + uri=asset_data["link"], + remote_file_name=asset_data["fileName"], + step=asset_data["step"], + metadata=metadata, + ) else: self._log_asset( experiment, path, asset_type, - log_filename, - assets_metadata, + asset_data, asset_map, ) # Process all nested assets: - for log_filename in assets_metadata: - asset_type = assets_metadata[log_filename].get("type", "asset") or "asset" + for asset_data in assets_metadata: + asset_type = asset_data.get("type", "asset") or "asset" if asset_type in ["confusion-matrix", "embeddings", "datagrid"]: self._log_asset( experiment, path, asset_type, - log_filename, - assets_metadata, + asset_data, asset_map, ) @@ -1632,11 +1641,11 @@ def log_all(self, experiment, experiment_folder): assets_metadata_filename = os.path.join( experiment_folder, "assets", "assets_metadata.jsonl" ) - assets_metadata = {} + assets_metadata = [] if os.path.exists(assets_metadata_filename): for line in open(assets_metadata_filename): data = json.loads(line) - assets_metadata[data["fileName"]] = data + assets_metadata.append(data) self.log_assets( experiment, diff --git a/cometx/framework/comet/download_manager.py b/cometx/framework/comet/download_manager.py index 1af3655..20b9294 100644 --- a/cometx/framework/comet/download_manager.py +++ b/cometx/framework/comet/download_manager.py @@ -1214,6 +1214,39 @@ def download_assets(self, experiment): self.asset_type if self.asset_type else "all" ) if len(assets) > 0: + filename_counts = {} + for asset in assets: + if asset["type"] == "audio" and asset["step"] is not None: + asset_filename = asset["fileName"] + asset["logAsFileName"] = asset_filename + if "." in asset_filename: + asset_filename, ext = asset_filename.split(".", 1) + else: + asset_filename, ext = asset_filename, "" + asset_filename = "%s-%s.%s" % ( + asset_filename, + asset["step"], + ext, + ) + asset["fileName"] = asset_filename + fn = asset["fileName"] + if fn in filename_counts: + filename_counts[fn] += 1 + if "." in fn: + base, ext = fn.rsplit(".", 1) + asset["diskFileName"] = "%s_%d.%s" % ( + base, + filename_counts[fn], + ext, + ) + else: + asset["diskFileName"] = "%s_%d" % ( + fn, + filename_counts[fn], + ) + else: + filename_counts[fn] = 0 + filename = "assets_metadata.jsonl" filepath = os.path.join(assets_path, filename) if self._should_write(filepath): @@ -1221,19 +1254,6 @@ def download_assets(self, experiment): os.makedirs(assets_path, exist_ok=True) with open(filepath, "w") as f: for asset in assets: - if asset["type"] == "audio" and asset["step"] is not None: - asset_filename = asset["fileName"] - asset["logAsFileName"] = asset_filename - if "." in asset_filename: - asset_filename, ext = asset_filename.split(".", 1) - else: - asset_filename, ext = asset_filename, "" - asset_filename = "%s-%s.%s" % ( - asset_filename, - asset["step"], - ext, - ) - asset["fileName"] = asset_filename f.write(json.dumps(asset)) f.write("\n") @@ -1244,9 +1264,10 @@ def download_assets(self, experiment): path = assets_path else: path = os.path.join(assets_path, asset_type) - filename = sanitize_filename(asset["fileName"]) + filename = sanitize_filename( + asset.get("diskFileName", asset["fileName"]) + ) file_path = os.path.join(path, filename) - # Don't download a filename more than once: if file_path not in filenames and self._should_write(file_path): filenames.add(file_path) self.summary["assets"] += 1