Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
111 changes: 60 additions & 51 deletions cometx/cli/copy.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand All @@ -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":
Expand Down Expand Up @@ -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"]
Expand All @@ -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":
Expand Down Expand Up @@ -1335,47 +1334,46 @@ 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/"):
model_name = dir_name[len("models/"):]
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(
Comment thread
chasefortier marked this conversation as resolved.
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"]
Expand All @@ -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"]
Expand All @@ -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(
Comment thread
chasefortier marked this conversation as resolved.
model_name=model_name,
Comment thread
chasefortier marked this conversation as resolved.
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,
)

Expand Down Expand Up @@ -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,
Expand Down
51 changes: 36 additions & 15 deletions cometx/framework/comet/download_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -1214,26 +1214,46 @@ 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):
self.summary["assets"] += 1
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")

Expand All @@ -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
Expand Down