From 9d92f95e9d1df572cdef25c4b644d79bacf2a6a0 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Thu, 7 May 2026 10:38:24 +0200 Subject: [PATCH 01/11] Add a default thread for Non-MPI simulations --- .../libmuscle/manager/instance_manager.py | 9 +++++++-- libmuscle/python/libmuscle/planner/planner.py | 20 ++++++++++++++++--- .../libmuscle/planner/test/test_planner.py | 20 +++++++++++++++++++ 3 files changed, 44 insertions(+), 5 deletions(-) diff --git a/libmuscle/python/libmuscle/manager/instance_manager.py b/libmuscle/python/libmuscle/manager/instance_manager.py index 1e7658c6..87b75a48 100644 --- a/libmuscle/python/libmuscle/manager/instance_manager.py +++ b/libmuscle/python/libmuscle/manager/instance_manager.py @@ -6,7 +6,7 @@ from multiprocessing import Queue from queue import Empty -from ymmsl.v0_2 import Configuration, ExecutionModel, Reference +from ymmsl.v0_2 import Configuration, ExecutionModel, Reference, ThreadedResReq from libmuscle.errors import ConfigurationError from libmuscle.manager.instance_registry import InstanceRegistry @@ -139,9 +139,14 @@ def start_all(self) -> None: stdout_path = idir / 'stdout.txt' stderr_path = idir / 'stderr.txt' + res_req = self._configuration.resources.get(component.name) + if res_req is None: + # No resources defined; use default of 1 thread for DIRECT components. + res_req = ThreadedResReq(component.name, 1) + request = InstantiationRequest( instance, program, - self._configuration.resources[component.name], + res_req, resources, idir, workdir, stdout_path, stderr_path) _logger.info(f'Instantiating {instance}') self._requests_out.put(request) diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index 16371512..59012862 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -3,8 +3,8 @@ from typing import Dict, Iterable, List, Mapping, Set, Tuple from ymmsl.v0_2 import ( - Component, Configuration, Model, MPICoresResReq, MPINodesResReq, - Operator, Reference, ResourceRequirements, ThreadedResReq) + Component, Configuration, ExecutionModel, Model, MPICoresResReq, + MPINodesResReq, Operator, Reference, ResourceRequirements, ThreadedResReq) from libmuscle.planner.resources import OnNodeResources, Resources from libmuscle.util import instance_indices @@ -516,13 +516,27 @@ def allocate_all( # Analyse model model = ModelGraph(configuration.root_model()) - requirements = configuration.resources programs = configuration.programs exclusive = { i for c in model.components() for i in c.instances() if (c.implementation and not programs[c.implementation].can_share_resources)} + # Build requirements: filling in defaults for DIRECT components + # that have no resources defined (default: 1 thread). + requirements: Dict[Reference, ResourceRequirements] = dict( + configuration.resources) + for component in model.components(): + if (component.name not in requirements and + component.implementation is not None and + component.implementation in programs and + programs[component.implementation].execution_model + is ExecutionModel.DIRECT): + _logger.debug( + f'No resources defined for {component.name}, using default' + ' of 1 thread.') + requirements[component.name] = ThreadedResReq(component.name, 1) + # Allocate unallocated_instances = [ i for c in model.components() for i in c.instances()] diff --git a/libmuscle/python/libmuscle/planner/test/test_planner.py b/libmuscle/python/libmuscle/planner/test/test_planner.py index 0e164dc7..9e9f648a 100644 --- a/libmuscle/python/libmuscle/planner/test/test_planner.py +++ b/libmuscle/python/libmuscle/planner/test/test_planner.py @@ -277,3 +277,23 @@ def test_impossible_virtual_allocation() -> None: planner = Planner(res) with pytest.raises(InsufficientResourcesAvailable): planner.allocate_all(config, virtual=True) + + +def test_default_resource_for_direct_without_resources() -> None: + """Test that a DIRECT component with no resources gets a default of 1 thread.""" + model = Model( + 'test_model', None, '', None, + [Component('micro', Ports(), '', 'micro')]) + # Program with DIRECT execution model (the default), no resources defined + programs = [Program(Ref('micro'), executable='/usr/bin/micro')] + config = Configuration('config', [], [model], None, None, programs, None) + + res = resources({'node001': [c(1), c(2), c(3), c(4)]}) + + planner = Planner(res) + allocations = planner.allocate_all(config) + + # Should have allocated 1 core (the default) for the micro instance + assert Ref('micro') in allocations + assert len(allocations[Ref('micro')].by_rank) == 1 + assert allocations[Ref('micro')].by_rank[0].total_cores() == 1 From d12d317a78a7cf4bc4f6e7a6a916327ec0457462 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Thu, 7 May 2026 10:49:59 +0200 Subject: [PATCH 02/11] fix error --- libmuscle/python/libmuscle/mcp/tcp_util.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libmuscle/python/libmuscle/mcp/tcp_util.py b/libmuscle/python/libmuscle/mcp/tcp_util.py index d14f51d4..fc3a6b9b 100644 --- a/libmuscle/python/libmuscle/mcp/tcp_util.py +++ b/libmuscle/python/libmuscle/mcp/tcp_util.py @@ -53,7 +53,7 @@ def recv_all(socket: SocketType, length: int) -> bytes: received_count += received_now - return databuf + return bytes(databuf) def send_int64(socket: SocketType, data: int) -> None: From 6657e3d08f13ec3db6c1395cfed68300f0591888 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Fri, 8 May 2026 11:22:04 +0200 Subject: [PATCH 03/11] simplify if statement --- libmuscle/python/libmuscle/planner/planner.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index 59012862..e4061bd7 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -527,11 +527,10 @@ def allocate_all( requirements: Dict[Reference, ResourceRequirements] = dict( configuration.resources) for component in model.components(): + program = programs.get(component.implementation) if (component.name not in requirements and - component.implementation is not None and - component.implementation in programs and - programs[component.implementation].execution_model - is ExecutionModel.DIRECT): + program is not None and + program.execution_model is ExecutionModel.DIRECT): _logger.debug( f'No resources defined for {component.name}, using default' ' of 1 thread.') From f7c20cc5a6758e1118d6f656cd624b82bb464df2 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Fri, 8 May 2026 11:26:24 +0200 Subject: [PATCH 04/11] simplified if statement --- libmuscle/python/libmuscle/planner/planner.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index e4061bd7..517d8557 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -526,11 +526,12 @@ def allocate_all( # that have no resources defined (default: 1 thread). requirements: Dict[Reference, ResourceRequirements] = dict( configuration.resources) + direct_programs = { + name for name, prog in programs.items() + if prog.execution_model is ExecutionModel.DIRECT} for component in model.components(): - program = programs.get(component.implementation) - if (component.name not in requirements and - program is not None and - program.execution_model is ExecutionModel.DIRECT): + if (component.implementation in direct_programs and + component.name not in requirements): _logger.debug( f'No resources defined for {component.name}, using default' ' of 1 thread.') From 3611e9bdea1a480231266cdaa03b7dcae1253fc8 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Fri, 8 May 2026 11:39:30 +0200 Subject: [PATCH 05/11] define and use get_resources --- .../libmuscle/manager/instance_manager.py | 8 ++--- libmuscle/python/libmuscle/planner/planner.py | 19 ++++------ libmuscle/python/libmuscle/util.py | 36 ++++++++++++++++++- 3 files changed, 44 insertions(+), 19 deletions(-) diff --git a/libmuscle/python/libmuscle/manager/instance_manager.py b/libmuscle/python/libmuscle/manager/instance_manager.py index 87b75a48..3dc2f0d2 100644 --- a/libmuscle/python/libmuscle/manager/instance_manager.py +++ b/libmuscle/python/libmuscle/manager/instance_manager.py @@ -6,7 +6,7 @@ from multiprocessing import Queue from queue import Empty -from ymmsl.v0_2 import Configuration, ExecutionModel, Reference, ThreadedResReq +from ymmsl.v0_2 import Configuration, ExecutionModel, Reference from libmuscle.errors import ConfigurationError from libmuscle.manager.instance_registry import InstanceRegistry @@ -18,6 +18,7 @@ from libmuscle.native_instantiator.native_instantiator import NativeInstantiator from libmuscle.planner.planner import Planner, ResourceAssignment from libmuscle.planner.resources import Resources +from libmuscle.util import get_resources _logger = logging.getLogger(__name__) @@ -139,10 +140,7 @@ def start_all(self) -> None: stdout_path = idir / 'stdout.txt' stderr_path = idir / 'stderr.txt' - res_req = self._configuration.resources.get(component.name) - if res_req is None: - # No resources defined; use default of 1 thread for DIRECT components. - res_req = ThreadedResReq(component.name, 1) + res_req = get_resources(self._configuration, component.name) request = InstantiationRequest( instance, program, diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index 517d8557..2e2df8d7 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -3,11 +3,11 @@ from typing import Dict, Iterable, List, Mapping, Set, Tuple from ymmsl.v0_2 import ( - Component, Configuration, ExecutionModel, Model, MPICoresResReq, + Component, Configuration, Model, MPICoresResReq, MPINodesResReq, Operator, Reference, ResourceRequirements, ThreadedResReq) from libmuscle.planner.resources import OnNodeResources, Resources -from libmuscle.util import instance_indices +from libmuscle.util import get_resources, instance_indices _logger = logging.getLogger(__name__) @@ -524,18 +524,11 @@ def allocate_all( # Build requirements: filling in defaults for DIRECT components # that have no resources defined (default: 1 thread). - requirements: Dict[Reference, ResourceRequirements] = dict( - configuration.resources) - direct_programs = { - name for name, prog in programs.items() - if prog.execution_model is ExecutionModel.DIRECT} + requirements: Dict[Reference, ResourceRequirements] = {} for component in model.components(): - if (component.implementation in direct_programs and - component.name not in requirements): - _logger.debug( - f'No resources defined for {component.name}, using default' - ' of 1 thread.') - requirements[component.name] = ThreadedResReq(component.name, 1) + res = get_resources(configuration, component.name) + if res is not None: + requirements[component.name] = res # Allocate unallocated_instances = [ diff --git a/libmuscle/python/libmuscle/util.py b/libmuscle/python/libmuscle/util.py index 9b375a0b..c8c1e499 100644 --- a/libmuscle/python/libmuscle/util.py +++ b/libmuscle/python/libmuscle/util.py @@ -1,10 +1,44 @@ import itertools +import logging from pathlib import Path import sys import time from typing import Generator, List, Optional, cast -from ymmsl.v0_2 import Reference +from ymmsl.v0_2 import ( + Configuration, ExecutionModel, Reference, ResourceRequirements, + ThreadedResReq) + + +_logger = logging.getLogger(__name__) + + +def get_resources( + configuration: Configuration, name: Reference + ) -> ResourceRequirements: + """Get the resource requirements for a component. + + If no resources are defined for this component and + it uses a non-MPI execution model (DIRECT), a default of 1 thread is + returned. + + Args: + configuration: The configuration to look up resources in. + name: The name of the component to get resources for. + + Returns: + The resource requirements for the component. + """ + res_req = configuration.resources.get(name) + if res_req is None: + implementation = configuration.root_model().components[name].implementation + if (implementation is not None and + configuration.programs[implementation].execution_model + is ExecutionModel.DIRECT): + _logger.debug( + f'No resources defined for {name}, using default of 1 thread.') + res_req = ThreadedResReq(name, 1) + return res_req def instance_to_kernel(instance: Reference) -> Reference: From 0dcb94e1033cdaed967d32d61a0f6dee5d1dacb0 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Fri, 8 May 2026 12:03:43 +0200 Subject: [PATCH 06/11] update call in planner --- libmuscle/python/libmuscle/planner/planner.py | 9 ++++----- libmuscle/python/libmuscle/util.py | 10 +++------- 2 files changed, 7 insertions(+), 12 deletions(-) diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index 2e2df8d7..a0108313 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -524,11 +524,10 @@ def allocate_all( # Build requirements: filling in defaults for DIRECT components # that have no resources defined (default: 1 thread). - requirements: Dict[Reference, ResourceRequirements] = {} - for component in model.components(): - res = get_resources(configuration, component.name) - if res is not None: - requirements[component.name] = res + requirements: dict[Reference, ResourceRequirements] = { + component.name: get_resources(configuration, component.name) + for component in model.components() + } # Allocate unallocated_instances = [ diff --git a/libmuscle/python/libmuscle/util.py b/libmuscle/python/libmuscle/util.py index c8c1e499..c7d9279b 100644 --- a/libmuscle/python/libmuscle/util.py +++ b/libmuscle/python/libmuscle/util.py @@ -31,13 +31,9 @@ def get_resources( """ res_req = configuration.resources.get(name) if res_req is None: - implementation = configuration.root_model().components[name].implementation - if (implementation is not None and - configuration.programs[implementation].execution_model - is ExecutionModel.DIRECT): - _logger.debug( - f'No resources defined for {name}, using default of 1 thread.') - res_req = ThreadedResReq(name, 1) + _logger.debug( + f'No resources defined for {name}, using default of 1 thread.') + res_req = ThreadedResReq(name, 1) return res_req From 06d824d3a25eb1df793471fbb0e013b3edf123ff Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Fri, 8 May 2026 12:05:44 +0200 Subject: [PATCH 07/11] remove unused import --- libmuscle/python/libmuscle/util.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/libmuscle/python/libmuscle/util.py b/libmuscle/python/libmuscle/util.py index c7d9279b..6cbc1f2c 100644 --- a/libmuscle/python/libmuscle/util.py +++ b/libmuscle/python/libmuscle/util.py @@ -6,8 +6,8 @@ from typing import Generator, List, Optional, cast from ymmsl.v0_2 import ( - Configuration, ExecutionModel, Reference, ResourceRequirements, - ThreadedResReq) + Configuration, Reference, ResourceRequirements, ThreadedResReq + ) _logger = logging.getLogger(__name__) From 66191dd672967c15f6d461b778ae98196ad4cdcf Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Mon, 11 May 2026 09:21:17 +0200 Subject: [PATCH 08/11] moved get_resources to ymmsl-python --- .../libmuscle/manager/instance_manager.py | 3 +- libmuscle/python/libmuscle/planner/planner.py | 4 +-- .../libmuscle/planner/test/test_planner.py | 20 ------------ libmuscle/python/libmuscle/util.py | 32 +------------------ 4 files changed, 4 insertions(+), 55 deletions(-) diff --git a/libmuscle/python/libmuscle/manager/instance_manager.py b/libmuscle/python/libmuscle/manager/instance_manager.py index 3dc2f0d2..2dab1f99 100644 --- a/libmuscle/python/libmuscle/manager/instance_manager.py +++ b/libmuscle/python/libmuscle/manager/instance_manager.py @@ -18,7 +18,6 @@ from libmuscle.native_instantiator.native_instantiator import NativeInstantiator from libmuscle.planner.planner import Planner, ResourceAssignment from libmuscle.planner.resources import Resources -from libmuscle.util import get_resources _logger = logging.getLogger(__name__) @@ -140,7 +139,7 @@ def start_all(self) -> None: stdout_path = idir / 'stdout.txt' stderr_path = idir / 'stderr.txt' - res_req = get_resources(self._configuration, component.name) + res_req = self._configuration.get_resources(component.name) request = InstantiationRequest( instance, program, diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index a0108313..a5b58be7 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -7,7 +7,7 @@ MPINodesResReq, Operator, Reference, ResourceRequirements, ThreadedResReq) from libmuscle.planner.resources import OnNodeResources, Resources -from libmuscle.util import get_resources, instance_indices +from libmuscle.util import instance_indices _logger = logging.getLogger(__name__) @@ -525,7 +525,7 @@ def allocate_all( # Build requirements: filling in defaults for DIRECT components # that have no resources defined (default: 1 thread). requirements: dict[Reference, ResourceRequirements] = { - component.name: get_resources(configuration, component.name) + component.name: configuration.get_resources(component.name) for component in model.components() } diff --git a/libmuscle/python/libmuscle/planner/test/test_planner.py b/libmuscle/python/libmuscle/planner/test/test_planner.py index 9e9f648a..0e164dc7 100644 --- a/libmuscle/python/libmuscle/planner/test/test_planner.py +++ b/libmuscle/python/libmuscle/planner/test/test_planner.py @@ -277,23 +277,3 @@ def test_impossible_virtual_allocation() -> None: planner = Planner(res) with pytest.raises(InsufficientResourcesAvailable): planner.allocate_all(config, virtual=True) - - -def test_default_resource_for_direct_without_resources() -> None: - """Test that a DIRECT component with no resources gets a default of 1 thread.""" - model = Model( - 'test_model', None, '', None, - [Component('micro', Ports(), '', 'micro')]) - # Program with DIRECT execution model (the default), no resources defined - programs = [Program(Ref('micro'), executable='/usr/bin/micro')] - config = Configuration('config', [], [model], None, None, programs, None) - - res = resources({'node001': [c(1), c(2), c(3), c(4)]}) - - planner = Planner(res) - allocations = planner.allocate_all(config) - - # Should have allocated 1 core (the default) for the micro instance - assert Ref('micro') in allocations - assert len(allocations[Ref('micro')].by_rank) == 1 - assert allocations[Ref('micro')].by_rank[0].total_cores() == 1 diff --git a/libmuscle/python/libmuscle/util.py b/libmuscle/python/libmuscle/util.py index 6cbc1f2c..9b375a0b 100644 --- a/libmuscle/python/libmuscle/util.py +++ b/libmuscle/python/libmuscle/util.py @@ -1,40 +1,10 @@ import itertools -import logging from pathlib import Path import sys import time from typing import Generator, List, Optional, cast -from ymmsl.v0_2 import ( - Configuration, Reference, ResourceRequirements, ThreadedResReq - ) - - -_logger = logging.getLogger(__name__) - - -def get_resources( - configuration: Configuration, name: Reference - ) -> ResourceRequirements: - """Get the resource requirements for a component. - - If no resources are defined for this component and - it uses a non-MPI execution model (DIRECT), a default of 1 thread is - returned. - - Args: - configuration: The configuration to look up resources in. - name: The name of the component to get resources for. - - Returns: - The resource requirements for the component. - """ - res_req = configuration.resources.get(name) - if res_req is None: - _logger.debug( - f'No resources defined for {name}, using default of 1 thread.') - res_req = ThreadedResReq(name, 1) - return res_req +from ymmsl.v0_2 import Reference def instance_to_kernel(instance: Reference) -> Reference: From 89a396ef11ba6eeeb996bbdf3c9c728c60aa0d22 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Mon, 11 May 2026 10:43:26 +0200 Subject: [PATCH 09/11] update pyproject.toml --- pyproject.toml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 8940ea73..77bd7b7e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -32,8 +32,8 @@ dependencies = [ "psutil>=5.0.0", "parsimonious", "numpy>=1.22", - "ymmsl>=0.16.0,<0.17", # Also in examples requirements.txt - # "ymmsl @ git+https://github.com/multiscale/ymmsl-python.git@develop", + # "ymmsl>=0.16.0,<0.17", # Also in examples requirements.txt + "ymmsl @ git+https://github.com/multiscale/ymmsl-python.git@develop", "yatiml>=0.12,<0.13", ] From dbb141b56747175a4ce2dc34f9f09ecca079f48e Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Mon, 11 May 2026 16:19:46 +0200 Subject: [PATCH 10/11] remove bytes changes --- libmuscle/python/libmuscle/mcp/tcp_util.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libmuscle/python/libmuscle/mcp/tcp_util.py b/libmuscle/python/libmuscle/mcp/tcp_util.py index fc3a6b9b..d14f51d4 100644 --- a/libmuscle/python/libmuscle/mcp/tcp_util.py +++ b/libmuscle/python/libmuscle/mcp/tcp_util.py @@ -53,7 +53,7 @@ def recv_all(socket: SocketType, length: int) -> bytes: received_count += received_now - return bytes(databuf) + return databuf def send_int64(socket: SocketType, data: int) -> None: From 408c2116f9177adb2c69a15c325ed22721e8e881 Mon Sep 17 00:00:00 2001 From: Iris van der Werf Date: Mon, 18 May 2026 09:42:37 +0200 Subject: [PATCH 11/11] remove get_resources comments --- libmuscle/python/libmuscle/planner/planner.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/libmuscle/python/libmuscle/planner/planner.py b/libmuscle/python/libmuscle/planner/planner.py index c952d1e1..2fe735b1 100644 --- a/libmuscle/python/libmuscle/planner/planner.py +++ b/libmuscle/python/libmuscle/planner/planner.py @@ -523,8 +523,6 @@ def allocate_all( if (c.implementation and not programs[c.implementation].can_share_resources)} - # Build requirements: filling in defaults for DIRECT components - # that have no resources defined (default: 1 thread). requirements: dict[Reference, ResourceRequirements] = { root_model.name + component.name: configuration.get_resources(root_model.name + component.name)