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
17 changes: 16 additions & 1 deletion concore_cli/commands/validate.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,8 @@ def _classify_message(message, bucket_name):
return {"error_type": "missing_edge_source"}
if message.startswith("Edge references non-existent target node:"):
return {"error_type": "missing_edge_target"}
if message.startswith("Edge label '") and "has more than one source" in message:
return {"error_type": "edge_label_multiple_sources"}
if message == "Workflow contains cycles (expected for control loops)":
return {"error_type": "cycle_detected"}
if message.startswith("Invalid port number:"):
Expand Down Expand Up @@ -292,6 +294,7 @@ def finalize():
edge_label_regex = re.compile(r"0x([a-fA-F0-9]+)_(\S+)")
zmq_edges = 0
file_edges = 0
label_sources = {}

for edge in edges:
try:
Expand All @@ -300,13 +303,25 @@ def finalize():
label_tag = edge.find("EdgeLabel")

if label_tag and label_tag.text:
if edge_label_regex.match(label_tag.text.strip()):
edge_label = label_tag.text.strip()
if edge_label_regex.match(edge_label):
zmq_edges += 1
else:
file_edges += 1
label_sources.setdefault(edge_label, set()).add(
edge.get("source")
)
except Exception:
pass

# mkconcore mounts a file edge for its first source only
for edge_label, sources in label_sources.items():
if len(sources) > 1:
errors.append(
f"Edge label '{edge_label}' has more than one source node: "
f"{', '.join(sorted(sources))}"
)

if zmq_edges > 0:
info.append(f"ZMQ-based edges: {zmq_edges}")
if file_edges > 0:
Expand Down
11 changes: 10 additions & 1 deletion mkconcore.py
Original file line number Diff line number Diff line change
Expand Up @@ -382,8 +382,17 @@ def cleanup_script_files():
#Validate edge labels
safe_name(edge_label, f"Edge label '{edge_label}'")

source_label = nodes_dict[edge['source']]
if edge_label not in edges_dict:
edges_dict[edge_label] = [nodes_dict[edge['source']], []]
edges_dict[edge_label] = [source_label, []]
elif edges_dict[edge_label][0] != source_label:
# an edge label is one volume with a single writer, a second
# source would get no out mount and its writes would be lost
logging.error(
f"Edge label '{raw_label}' has more than one source node: "
f"{edges_dict[edge_label][0]} and {source_label}"
)
sys.exit(1)
edges_dict[edge_label][1].append(nodes_dict[edge['target']])
except (IndexError, AttributeError, KeyError):
logging.debug('An edge with no valid properties or missing node encountered and ignored')
Expand Down
33 changes: 33 additions & 0 deletions tests/test_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,39 @@ def test_build_command_missing_source(self):
)
self.assertNotEqual(result.exit_code, 0)

def test_build_command_rejects_edge_label_with_multiple_sources(self):
with self.runner.isolated_filesystem(temp_dir=self.temp_dir):
Path("src").mkdir()
for name in ["a.py", "b.py", "c.py"]:
Path("src", name).write_text("import concore\n")
Path("workflow.graphml").write_text(
'<graphml xmlns:y="http://www.yworks.com/xml/graphml">\n'
'<graph id="G" edgedefault="directed">\n'
'<node id="n0"><data key="d0"><y:NodeLabel>A:a.py</y:NodeLabel></data></node>\n'
'<node id="n1"><data key="d0"><y:NodeLabel>B:b.py</y:NodeLabel></data></node>\n'
'<node id="n2"><data key="d0"><y:NodeLabel>C:c.py</y:NodeLabel></data></node>\n'
'<edge source="n0" target="n2"><data key="d1"><y:EdgeLabel>shared</y:EdgeLabel></data></edge>\n'
'<edge source="n1" target="n2"><data key="d1"><y:EdgeLabel>shared</y:EdgeLabel></data></edge>\n'
"</graph>\n"
"</graphml>\n"
)

result = self.runner.invoke(
cli,
[
"build",
"workflow.graphml",
"--source",
"src",
"--output",
"out",
"--type",
"posix",
],
)
self.assertNotEqual(result.exit_code, 0)
self.assertIn("Edge label 'shared' has more than one source", result.output)

def test_build_command_from_project_dir(self):
with self.runner.isolated_filesystem(temp_dir=self.temp_dir):
result = self.runner.invoke(cli, ["init", "test-project"])
Expand Down
34 changes: 34 additions & 0 deletions tests/test_graph.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import json
import unittest
import tempfile
import shutil
Expand Down Expand Up @@ -188,6 +189,39 @@ def test_validate_with_existing_source_file(self):

self.assertIn("Validation passed", result.output)

def test_validate_edge_label_with_multiple_sources(self):
content = """
<graphml xmlns:y="http://www.yworks.com/xml/graphml">
<graph id="G" edgedefault="directed">
<node id="n0">
<data key="d0"><y:NodeLabel>A:a.py</y:NodeLabel></data>
</node>
<node id="n1">
<data key="d0"><y:NodeLabel>B:b.py</y:NodeLabel></data>
</node>
<node id="n2">
<data key="d0"><y:NodeLabel>C:c.py</y:NodeLabel></data>
</node>
<edge source="n0" target="n2">
<data key="d1"><y:EdgeLabel>shared</y:EdgeLabel></data>
</edge>
<edge source="n1" target="n2">
<data key="d1"><y:EdgeLabel>shared</y:EdgeLabel></data>
</edge>
</graph>
</graphml>
"""
filepath = self.create_graph_file("shared_label.graphml", content)

result = self.runner.invoke(cli, ["validate", filepath])

self.assertNotEqual(result.exit_code, 0)
self.assertIn("Edge label 'shared' has more than one source", result.output)

result = self.runner.invoke(cli, ["validate", filepath, "--format", "json"])
error_types = [e["error_type"] for e in json.loads(result.output)["errors"]]
self.assertIn("edge_label_multiple_sources", error_types)

def test_validate_zmq_port_conflict(self):
content = """
<graphml xmlns:y="http://www.yworks.com/xml/graphml">
Expand Down
Loading