seamm_exec package#

Submodules#

seamm_exec.base module#

Virtual base class for the Exec classes.

class seamm_exec.base.Base(logger)[source]#

Bases: object

exec(config, cmd=[], directory=None, input_data=None, env={}, shell=False, ce={})[source]#

Virtual method for executing a command.

Parameters:
  • cmd ([str]) – The command as a list of words.

  • directory (str or Path) – The directory for the tasks files.

  • input_data (str) – Data to be redirected to the stdin of the process.

  • env ({str: str}) – Dictionary of environment variables to pass to the execution environment

  • shell (bool = False) – Whether to use the shell when launching task

  • ce (dict(str, str or int)) – Description of the computational enviroment

get_mode_info(filename, mode)[source]#

Get the mode information for ‘ls’ like listing

ls_format(path, files)[source]#

Format a list of files as in ‘ls’

property name#

The name of this type of executor.

run(config, cmd=[], directory=None, input_data=None, files=None, env={}, return_files=[], shell=False, in_situ=None, ce={})[source]#

Execute a command directly on the current machine.

Parameters:
  • config (configparser.config) – The configuration options for the code

  • cmd ([str]) – The command as a list of words.

  • directory (str or Path) – The directory for the tasks files.

  • input_data (str) – Data to be redirected to the stdin of the process.

  • files ({str: str or byte}) – Dictionary of file names and data to write before execution.

  • env ({str: str}) – Dictionary of environment variables to pass to the execution environment

  • return_files ([str]) – List of files to return/keep from the calculation

  • shell (bool = False) – Whether to use the shell when launching task

  • in_situ (bool or None = None) – Whether to run in the given directory, versus a temp directory, copying back only return_files afterwards. None (the default) auto-detects: run in a temp directory – which honors $TMPDIR, so it lands on node-local scratch under a scheduler that sets it – when running under a batch scheduler (currently SLURM), since the job directory there is typically NFS-mounted and unsafe for an MPI code’s scratch I/O (molssi-seamm/orca_step#20); otherwise run in place, as before.

  • ce (dict(str, str or int)) – Description of the computational enviroment

  • original (This is the)

  • a (blocking interface. It runs one task through)

:param TaskSet with a one-slot: :param synchronous: :param LocalPool: :param without a manifest: :param so its: :param behaviour and the files it leaves are exactly as before.:

seamm_exec.bundle_runner module#

Run a SEAMM-mode bundle inside an allocation (see task_worker.py).

Each task runs through a LocalPool sized to the allocation, with programs resolved on this machine, and the task’s marker directory receives DONE or FAILED.

seamm_exec.bundle_runner.run_seamm_bundle(bundle)[source]#

Run a SEAMM-mode bundle through a LocalPool. Returns the number failed.

seamm_exec.computational_environment module#

Helper routines for determining the computational environment from queueuing systems.

seamm_exec.computational_environment.computational_environment(limits={})[source]#

Examine where we are running to get the computational environment.

Parameters:

limits (dict(str, any)) – Limits to impose on the number of tasks, gpus, etc.

Returns:

The attributes of the computational enviroment, limited by the imposed limits.

Return type:

dict(str, any)

seamm_exec.computational_environment.expand_hostlist(text)[source]#

SLURM’s compressed hostlist -> host names.

tc[053,059-061],gpu7 -> tc053 tc059 tc060 tc061 gpu7; zero padding is kept; a suffix after the brackets (tc[01-02]-ib) and several bracket groups (r[1-2]n[1-2], the product) are expanded as SLURM does.

seamm_exec.computational_environment.running_scheduler()[source]#

The name of the queueing system whose allocation this is, or None.

seamm_exec.computational_environment.scheduler_job_variables()[source]#

The environment variables that mark a batch job, one per scheduler.

seamm_exec.docker module#

The Docker object executes a task using a Docker container

class seamm_exec.docker.Docker(logger=<Logger seamm-exec (WARNING)>)[source]#

Bases: Base

exec(config, cmd=[], directory=None, input_data=None, env={}, shell=False, ce={})[source]#

Execute a command using a docker container..

Parameters:
  • config (dict(str: any)) – The configuration for the code to run

  • cmd ([str]) – The command as a list of words.

  • directory (str or Path) – The directory for the tasks files.

  • input_data (str) – Data to be redirected to the stdin of the process.

  • env ({str: str}) – Dictionary of environment variables to pass to the execution environment

  • shell (bool = False) – Whether to use the shell when launching task

  • ce (dict(str, str or int)) – Description of the computational enviroment

Returns:

{str – Dictionary with stdout, stderr, returncode, etc.

Return type:

str}

The equivalent commandline commmand is:

docker run --rm -v $PWD:/home -w /home psaxe/mopac [mopac input_file]

For this container the input filename defaults to “mopac.dat” so we do not need to add it.

property name#

The name of this type of executor.

seamm_exec.evaluator module#

Evaluate a model chemistry at many structures, over MDI or as tasks.

A step that needs energies (and gradients, and stress) for many structures submits them to an Evaluator and iterates over the results. The evaluator, not the step or the user, chooses how they are computed:

  • the batch path: each structure is a Task from the program’s get_task, run by a TaskSet (on the job’s target: the local pool or a queue, with restart), and read back by the program’s analyze_task;

  • the MDI path: one warm MDI engine per group of structures with the same elements, charge, multiplicity and periodicity, from the program’s get_mdi_engine_command.

The rule (see choose_path()): the batch path when the job’s tasks go to a queue and the program has get_task; else MDI when the program has an MDI engine, unless it prefers the batch path (options["prefers_batch"], set by ORCA, whose engine runs a subprocess per structure); else the batch path in the local pool. Both paths return the same numbers for the same model chemistry.

The program’s contract (classmethods beside get_model_chemistry_options):

get_task(configuration, model_chemistry, *, key, properties, options, resources)
    -> seamm_exec.Task
analyze_task(result, model_chemistry, configuration, *, properties, options)
    -> {"energy": kJ/mol, "gradients": (n, 3) kJ/mol/Å, "stress": GPa, ...}
can_run_task(configuration, model_chemistry, *, options) -> bool   (optional)

A structure for which can_run_task is False (e.g. a periodic system for ORCA or MOPAC) goes to the program’s MDI engine, if it has one, even when the others run as tasks. If there is no engine, or it cannot start here (a queue target with the code only on the cluster), that structure gets a failed result; so does one whose get_task raises. Neither stops the rest.

analyze_task raises AnalysisError when a required property is missing; it never returns partial numbers. options passes what a consumer needs for a fragment: atom_indices, ghost_atoms, charge, multiplicity, an initial guess.

exception seamm_exec.evaluator.AnalysisError[source]#

Bases: RuntimeError

A task’s results lack a requested property.

class seamm_exec.evaluator.Evaluator(node, model_chemistry=None, *, properties=('energy', 'gradients'), path=None, directory=None, target=None, task_set_options=None, resources=None, name='SEAMM')[source]#

Bases: object

Evaluate a model chemistry at many structures.

Parameters:
  • node (seamm.Node) – The step: its directory, flowchart (plug-ins, executor) and options.

  • model_chemistry (dict, optional) – The _model_chemistry wrapper. Default: the node’s variable.

  • properties ([str]) – “energy”, “gradients”, “stress”.

  • path (str, optional) – “mdi” or “batch”, for tests only; the evaluator chooses otherwise.

  • directory (str or Path, optional) – The batch path’s step directory (<directory>/tasks/...). Default: the node’s.

  • target (seamm_scheduler.TargetSection, optional) – The job’s target. Default: found for the job.

  • task_set_options (dict, optional) – Extra arguments for the TaskSet (bundle_tasks, archive, …).

  • resources (seamm_exec.Resources, optional) – The resources of each calculation on the batch path (ranks, memory per rank), passed to the provider’s get_task. Default: the provider’s.

  • name (str) – A name for the MDI engine.

close()[source]#
property mdi_capable#

Whether the program has an MDI engine for this model chemistry.

results()[source]#

Compute what has been submitted; yield an EvaluatorResult for each, as it is ready.

submit(configuration, key=None, *, options=None)[source]#

Add a structure; returns its key (unique, filesystem-safe).

static topology_key(configuration, options=None)[source]#

What must stay fixed for one MDI engine session.

class seamm_exec.evaluator.EvaluatorResult(key: str, ok: bool, energy: float | None = None, gradients: object = None, stress: object = None, reason: str | None = None, restored: bool = False, path: str = 'batch', data: dict = <factory>, elapsed: float = 0.0)[source]#

Bases: object

One structure’s results.

Variables:
  • key (str)

  • ok (bool)

  • energy (float or None) – kJ/mol

  • gradients (numpy.ndarray or None) – (n, 3), kJ/mol/Å

  • stress (list or None) – GPa, as the program gives it

  • reason (str or None) – Why it failed.

  • restored (bool) – From an earlier run (batch path).

  • path (str) – “mdi” or “batch”

  • data (dict) – Everything the program returned.

  • elapsed (float) – Seconds spent on it in this run (0 if restored).

data: dict#
elapsed: float = 0.0#
energy: float | None = None#
gradients: object = None#
key: str#
ok: bool#
path: str = 'batch'#
reason: str | None = None#
restored: bool = False#
stress: object = None#
class seamm_exec.evaluator.Geometry(atomic_numbers, coordinates, charge=0, multiplicity=1, cell=None, name=None)[source]#

Bases: object

A light stand-in for a molsystem configuration: elements, Cartesian coordinates (Å), charge, multiplicity and an optional cell. For structures a step makes on the fly, such as finite-difference displacements.

property n_atoms#
seamm_exec.evaluator.REQUIRED = ('energy', 'gradients')#

Properties a result must have when requested. Stress is required only for a periodic structure (check_properties(..., periodic=True), and the Evaluator checks it for every task); a molecule has none.

seamm_exec.evaluator.check_properties(data, properties, what, periodic=False)[source]#

Raise AnalysisError unless data has every required property in properties (and the stress, if requested, for a periodic structure). For programs’ analyze_task.

seamm_exec.evaluator.choose_path(model_chemistry, provider, target=None)[source]#

“mdi” or “batch” for a model chemistry, a program and the job’s target.

Raises:

ValueError – If the program offers neither path.

seamm_exec.evaluator.mdi_method_and_basis(model_chemistry)[source]#

The (method, basis) an MDI engine is launched with.

method is the program’s own keyword (options["mdi_method_arg"]), falling back to the model chemistry’s method; basis is options["mdi_basis_arg"] (the user’s basis, for programs that take one) or the model chemistry’s basis, or None for programs that take a method alone (MOPAC, xTB, an MLFF).

seamm_exec.evaluator.structure_data(configuration)[source]#

What a program needs from a configuration (or a Geometry).

Returns:

atomic_numbers, symbols, coordinates ((n, 3) Å), charge, multiplicity, periodicity and cell ((3, 3) Å or None).

Return type:

dict

seamm_exec.exec_flowchart module#

The ExecWorkflow object does what it name implies: it executes, or runs, a given flowchart.

It provides the environment for running the computational tasks locally or remotely, using what is commonly called workflow management system (WMS). The WMS concept, as used here, means tools that run given tasks without knowing anything about chemistry. The chemistry specialization is contained in the Flowchart and the nodes that it contains.

class seamm_exec.exec_flowchart.ExecFlowchart(flowchart=None)[source]#

Bases: object

run(root=None, job_id=None)[source]#
class seamm_exec.exec_flowchart.cd(newPath)[source]#

Bases: object

Context manager for changing the current working directory

seamm_exec.exec_flowchart.credential_sections(root)[source]#

The seammrc sections to look in for datastore credentials, in order.

[Dashboard: <name>] where name is the installation’s root directory name (SEAMM_DEV, SEAMM_NEW, …), then, as before, dev for a root whose path contains “dev” or localhost otherwise, then this host’s name.

seamm_exec.exec_flowchart.get_job_id(filename)[source]#

Get the next job id from the given file.

This uses the fasteners module to provide locking so that only one job at a time can access the file, so that the job ids are unique and monotonically increasing.

seamm_exec.exec_flowchart.open_datastore(root, datastore, timeout=20.0)[source]#

Open the database via the datastore

seamm_exec.exec_flowchart.run(job_id=None, wdir=None, db_path=None, setup_logging=True, in_jobserver=False, cmdline=None)[source]#

The standalone flowchart app

seamm_exec.exec_flowchart.run_from_jobserver()[source]#

Helper routine to run from the JobServer.

Gets the arguments from the command line.

seamm_exec.local module#

The Local object does what it name implies: it executes, or runs, an executable locally.

class seamm_exec.local.Local(logger=<Logger seamm-exec (WARNING)>)[source]#

Bases: Base

exec(config, cmd=[], directory=None, input_data=None, env={}, shell=False, ce={})[source]#

Execute a command directly on the current machine.

Parameters:
  • config (dict(str: any)) – The configuration for the code to run

  • cmd ([str]) – The command as a list of words.

  • directory (str or Path) – The directory for the tasks files.

  • input_data (str) – Data to be redirected to the stdin of the process.

  • env ({str: str}) – Dictionary of environment variables to pass to the execution environment

  • shell (bool = False) – Whether to use the shell when launching task

  • ce (dict(str, str or int)) – Description of the computational enviroment

Returns:

{str – Dictionary with stdout, stderr, returncode, etc.

Return type:

str}

property name#

The name of this type of executor.

seamm_exec.local_pool module#

LocalPool: run tasks concurrently on this machine or in this allocation.

The pool’s capacity is the computational environment – the whole machine, or the SLURM allocation the evaluator runs in – and each task gets a share of it from its Resources. Tasks start, in the order submitted, as soon as enough cores and memory are free; a task larger than the pool is clamped to it and runs alone.

Each task runs through the executor’s _run_task() (the body of the original Base.run()), so the conda, modules and docker handling, in_situ and the return-file contract are exactly those of Base.run().

class seamm_exec.local_pool.LocalPool(executor, *, root=None, ce=None, synchronous=False, max_concurrent=None, resolve_programs=False)[source]#

Bases: object

A back end that runs tasks as local subprocesses.

Parameters:
  • executor (seamm_exec.Base) – Provides _run_task() and exec().

  • root (str or Path, optional) – Where the <program>.ini files are, for tasks without config.

  • ce (dict, optional) – The computational environment to share out. Default: computational_environment().

  • synchronous (bool = False) – Run each task in submit(), in the caller’s thread, with the given ce unchanged, the environment untouched and exceptions propagated: exactly the original Base.run().

  • max_concurrent (int, optional) – At most this many tasks at once.

  • resolve_programs (bool = False) – Kept for compatibility; it changes nothing now. A task without config is always configured here, from <root>/<program>.ini and the program’s resolver (see seamm_exec.resolve); a task with config is run with it exactly as given.

cancel(ids)[source]#
capacity()[source]#

The cores, memory (bytes) and GPUs this pool shares out.

config_for(task)[source]#

The program’s configuration: task.config, else <program>.ini.

Adds code_dir, the directory holding code, for commands that run a code’s companion programs (e.g. ORCA’s orca_2aim), when code is a path rather than a bare name.

fetch(task, backend_id)[source]#
has_program(task)[source]#

Whether the task’s program can run here.

name = 'local'#
reattach(records)[source]#

Tasks from an earlier evaluator: kill any still running here; lost.

A process that is not our child cannot be waited for, so instead of adopting it the pool kills its process group and the task is rerun. A pid is reused once its process has gone, so the group is killed only if its leader is the process the manifest recorded: same start time and working directory. (The host name check means a laptop that changed networks leaves the leftovers alone, the safe direction.)

status(ids)[source]#
submit(tasks, directories, on_start=None)[source]#
wait(ids, timeout=None)[source]#

Block until one of ids is finished, or timeout seconds.

seamm_exec.resolve module#

Resolve a program on the machine that runs it.

A task names a program (“orca”, “mopac”, …). Where it runs, the program’s configuration comes from that machine’s <root>/<program>.ini (the section for the executor, normally [local]), and a plug-in may register a resolver that adjusts the configuration, the command and the environment for that machine and the task’s share of it: ORCA must be invoked by its full path, and a parallel run needs the OpenMPI it was built against on the paths.

Resolvers are entry points in the group org.molssi.seamm.exec.resolvers, named after the program:

[project.entry-points."org.molssi.seamm.exec.resolvers"]
orca = "orca_step.resolver:resolve"

with the signature

def resolve(config, cmd, env, ce, root) -> (config, cmd, env)

config is the program’s ini section as a dict (already a copy), cmd the task’s command template (a list), env its extra environment, ce the task’s computational environment (NTASKS, MEM_PER_CPU, …) and root the SEAMM root holding the ini files. The command is formatted with config and ce afterwards, by the executor, exactly as before.

Task.config is never sent to a remote back end (transport = ssh): it was resolved on the evaluator’s machine and names that machine’s paths, so the TaskSet keeps such a task in the evaluator’s own pool, with a warning. A task without it is resolved where it runs, here. On the local transport Task.config, when given, is used as it is, as in the LocalPool.

seamm_exec.resolve.available(program, root)[source]#

Whether program can run here without an ini file: its resolver says so through an optional available(root) attribute.

seamm_exec.resolve.has_resolver(program)[source]#

Whether program has a registered resolver.

seamm_exec.resolve.read_config(program, root, section='local')[source]#

The [section] of <root>/<program>.ini as a dict, or None.

seamm_exec.resolve.register(program, function)[source]#

Register a resolver by hand (tests, or a program without a plug-in).

seamm_exec.resolve.resolve(program, config, cmd, env, ce, root)[source]#

Apply the program’s resolver, if it has one.

Returns:

The configuration, command and environment to run with.

Return type:

(dict, list, dict)

seamm_exec.resolve.resolvers()[source]#

The registered resolvers, {program: callable}.

Built once, under a lock: pool worker threads may ask at the same time, and a second build must not replace the dict (dropping a register()).

seamm_exec.scheduler_backend module#

SchedulerBackend: run tasks as batch jobs of a queueing system.

Each call to SchedulerBackend.submit() is one bundle: one batch job whose allocation runs the bundle’s tasks through the task worker (python -m seamm_exec.task_worker bundle.json) in SEAMM mode, that is, through a LocalPool inside the allocation, with each program resolved from the <program>.ini files of the machine it runs on. The TaskSet decides the bundles (bundle_tasks, bundle_walltime) and how many it keeps in the queue (room()).

Files: a bundle lives in <step dir>/tasks/_bundles/<bundle>.<n>/ (bundle.json, run.sh and the scheduler’s log), the tasks in their own directories, and the worker writes each task’s DONE or FAILED in <step dir>/tasks/<key>/. When the cluster does not share the evaluator’s filesystem the directories are pushed before submission and pulled back once the bundle’s job has ended, to <remote_root>/<job>/..., the same relative paths under the remote job directory.

Backend ids are <job id>#<bundle>.<n>#<key>, so a restarted evaluator can find a task’s job, bundle and files from its manifest record alone (adopt()) and poll it rather than submit it again.

exception seamm_exec.scheduler_backend.QueueFull(message, maybe_submitted=False)[source]#

Bases: RuntimeError

The queue will not take more jobs now (it is full, or cannot be reached); the TaskSet holds the bundle and tries again later.

maybe_submitted is True when the bundle may have reached the queue (the connection dropped during sbatch): its tasks must then stay recorded as queued, so that a restart looks for the job by name.

class seamm_exec.scheduler_backend.SchedulerBackend(queue, *, name='queue', job_directory, stager=None, remote_job_directory=None, directives=None, setup=None, python=None, root=None, executor='local', accepts_config=True, bundle_walltime=None, max_queued=None, poll_interval=30.0, job_name_prefix='seamm')[source]#

Bases: object

A TaskBackend that submits bundles of tasks to a queueing system.

Parameters:
  • queue (seamm_scheduler.QueueBackend) – Submits, polls and cancels jobs (a scheduler plus a transport).

  • name (str) – The name recorded in the manifest, queue:<target>.

  • job_directory (str or Path) – The job’s directory. Task directories must be inside it when the cluster does not share the filesystem.

  • stager (seamm_scheduler.JobStager, optional) – RsyncStager when the cluster does not share the filesystem. None (or a LocalStager) means shared storage.

  • remote_job_directory (str, optional) – Where the job directory is mirrored on the cluster, if not shared.

  • directives (dict, optional) – Site defaults, in the scheduler’s own spelling (the target section’s directive keys: partition, account, qos, export, …).

  • setup (str, optional) – Shell lines run in the batch script before the worker (module loads).

  • python (str, optional) – A Python with seamm_exec where the tasks run. Default: this one.

  • root (str, optional) – The SEAMM root with the <program>.ini files where the tasks run.

  • accepts_config (bool = True) – Whether Task.config is valid where the tasks run (local transport). On a remote cluster it is not sent.

  • bundle_walltime (float, optional) – Seconds to request for a bundle whose tasks give no walltime.

  • max_queued (int, optional) – The most jobs this user may have queued and running at once.

  • poll_interval (float = 30) – Seconds between polls of the queue.

abandon(records, keep=())[source]#

Cancel the jobs of records (tasks whose inputs have changed) unless a task in keep (ids adopted) still runs in the same job.

adopt(task, directory, marker, record)[source]#

Take back a task an earlier evaluator submitted, from its manifest record. Returns its id, or None if the record is not one of ours.

bundles = True#

The TaskSet submits one call per bundle.

cancel(ids)[source]#
check(task, directory, marker)[source]#

Fail early, before anything is submitted, for a task this back end cannot run (a directory outside the job, without shared storage).

fetch(task, backend_id)[source]#
forget(backend_id)[source]#

Drop a task the TaskSet will submit again.

classmethod from_target(section, *, job_directory, root=None, executor='local')[source]#

The back end a target section with tasks = queue describes.

reason(backend_id)[source]#

Why a task is lost, for the manifest.

relative(path)[source]#

path relative to the job directory.

room()[source]#

How many more bundles may be submitted now, or None for no limit.

property shared#
status(ids)[source]#
submit(tasks, directories, on_start=None, bundle=None, markers=None, on_prepared=None, bundle_walltime=None)[source]#

Submit tasks as one bundle: one batch job. Returns their ids.

bundle_walltime (seconds) is the caller’s limit on a bundle, used to bound its time when the tasks give none.

on_prepared(tasks, info) is called before sbatch with the bundle’s unique job name and directory, so the caller can record them: a restart that finds no job id looks the job up by that name.

wait(ids, timeout=None)[source]#

Sleep until the next poll is due.

where(path)[source]#

The path the cluster sees for a local path.

seamm_exec.scheduler_backend.parse_id(backend_id)[source]#

<job id>#<bundle>.<n>#<key> -> (job id, bundle dir name, key).

seamm_exec.scheduler_backend.remote_name(job_directory)[source]#

The name of a job directory on the cluster: its own name and a hash of its full path. Job numbers are unique only within one datastore, and several installations (~/SEAMM, ~/SEAMM_DEV, ChemAI) may share a remote_root; a hand run’s directory may be called anything.

seamm_exec.seamm_exec module#

Provide the primary functions.

seamm_exec.seamm_exec.canvas(with_attribution=True)[source]#

Placeholder function to show example docstring (NumPy format).

Replace this function and doc string for your own project.

Parameters:

with_attribution (bool, Optional, default: True) – Set whether or not to display who the quote is from.

Returns:

quote – Compiled string including quote and optional attribution.

Return type:

str

seamm_exec.targets module#

Find the job’s target: where its tasks run.

A target is a section of the JobServer’s <root>/<jobserver-name>.ini (see seamm_scheduler.config). The evaluator finds it, in order:

  1. given explicitly (TaskSet(target=...), a name or a TargetSection);

  2. <job dir>/target.json, written by the JobServer when it starts the job, so an evaluator running as a batch job on a cluster, which cannot read the JobServer’s ini file, still knows its target;

  3. the environment variable SEAMM_TARGET, naming a section of <root>/<hostname>.ini (or of the file named by SEAMM_TARGETS), for runs by hand with run_flowchart;

  4. none: tasks run in the evaluator’s own LocalPool, as before targets existed.

seamm_exec.targets.find_target(target=None, *, job_directory=None, root=None)[source]#

The target section for the job, or None.

Parameters:
  • target (str or TargetSection, optional) – A section, or the name of one in <root>/<hostname>.ini.

  • job_directory (str or Path, optional) – Where to look for target.json. Default: the current directory.

  • root (str or Path, optional) – The SEAMM root holding the JobServer’s ini file.

Return type:

seamm_scheduler.TargetSection or None

seamm_exec.targets.write_target(section, job_directory)[source]#

Write <job dir>/target.json for a section (what the JobServer does).

seamm_exec.task_worker module#

Run a bundle of tasks, one after another, inside one allocation.

Usage:

python task_worker.py bundle.json

This is the generic worker that a scheduler back end submits for a bundle of small tasks. It is pure Python with no SEAMM imports, so it runs anywhere a python3 exists, and is run by path rather than imported.

bundle.json:

{
  "tasks": [
    {
      "key": "frag-0001",
      "directory": "/path/to/tasks/frag-0001",
      "command": "orca orca.inp > orca.out",   # fully resolved
      "shell": true,
      "env": {"OMP_NUM_THREADS": "1"},
      "input": null,
      "fingerprint": "sha256:...",
      "success_text": {"orca.out": "ORCA TERMINATED NORMALLY"}
    },
    ...
  ]
}

For each task, in order: if <directory>/DONE exists it is skipped; otherwise the command runs in the directory with its standard output and error in stdout.txt and stderr.txt, the return code goes in returncode, and a successful run – return code 0 and, if given, every success_text file containing its text – writes DONE (JSON). A failure does not stop the bundle. The worker exits 0 if every task succeeded and 1 otherwise.

The SEAMM mode#

A bundle with "mode": "seamm" is run by a Python that has seamm_exec (python -m seamm_exec.task_worker bundle.json). Each task is then given as the task layer describes it – program, the cmd template, the names of its input files (already in its directory), return_files, resources, in_situ – and runs through a LocalPool sized to the allocation, exactly as the evaluator’s own pool would run it: the program resolved from <root>/<program>.ini on this machine (and its resolver, see seamm_exec.resolve), {code}/{NTASKS} from the task’s share of the allocation, scratch in $TMPDIR, and only return_files kept. Tasks that fit side by side in the allocation run concurrently.

Each task’s marker directory (<step dir>/tasks/<key>) receives DONE (JSON, in the task layer’s own form, including the list of returned files) or FAILED (JSON, with the reason). A task whose DONE exists is skipped, so a bundle that is submitted again – after a walltime limit, say – runs only what is left.

seamm_exec.task_worker.main(argv=None)[source]#
seamm_exec.task_worker.run_bundle(bundle)[source]#

Run the tasks of a bundle. Returns the number that failed.

seamm_exec.tasks module#

Tasks: the units of work a step hands to a back end.

A step describes each external calculation as a Task – the program, a command template, the input files, the files to keep and the resources – and runs a group of them with a TaskSet, which submits what is not done to a back end (a TaskBackend, e.g. LocalPool), waits, and yields a TaskResult for each as it finishes.

With a step directory, the TaskSet keeps <step dir>/tasks/manifest.json and writes <step dir>/tasks/<key>/DONE when a task finishes, so a rerun of the step in the same directory never recomputes a finished task.

See docs/developer_guide/campaigns/2026-10-02 for the design.

class seamm_exec.tasks.Manifest(path, save_interval=1.0)[source]#

Bases: object

<step dir>/tasks/manifest.json: what the evaluator knows of its tasks.

One record per key: the back end, its id, the state, the fingerprint, the bundle, timestamps, the attempt count and history, and the archive.

flush(force=False)[source]#

Write the manifest if it changed, at most every save_interval.

The manifest is for reattaching and reporting; whether a task finished is recorded by its DONE file, which is written at once.

get(key)[source]#
save()[source]#
update(key, **values)[source]#

Change a record. It is written by the next flush().

class seamm_exec.tasks.Resources(ntasks: int | None = None, cpus_per_task: int = 1, mem_per_cpu: int | None = None, ngpus: int = 0, walltime: float | None = None, partition: str | None = None, account: str | None = None, qos: str | None = None, nodes: int | None = None)[source]#

Bases: object

What a task needs, in scheduler-neutral terms.

Each back end translates these: the LocalPool into a share of this machine or allocation, a scheduler into its directives.

ntasks = None means “all of the back end’s capacity” (for the LocalPool, the whole machine or allocation, as a code got before tasks). Memory is in bytes and walltime in seconds.

account: str | None = None#
cpus_per_task: int = 1#
mem_per_cpu: int | None = None#
ngpus: int = 0#
nodes: int | None = None#
ntasks: int | None = None#
partition: str | None = None#
qos: str | None = None#
walltime: float | None = None#
class seamm_exec.tasks.Task(key: str, program: str, cmd: list = <factory>, files: dict | None = None, return_files: list = <factory>, resources: Resources = <factory>, env: dict = <factory>, in_situ: bool | None = None, shell: bool = False, input_data: str | None = None, estimated_seconds: float | None = None, target: str | None = None, directory: str | Path | None = None, config: dict | None = None, fingerprint: str | None = None, success_text: dict | None = None)[source]#

Bases: object

One external calculation.

Variables:
  • key (str) – Unique within the step and stable across restarts (e.g. a fragment key). Letters, digits and ._+=@,-; it names the task’s directory.

  • program (str) – The program’s identity (“orca”, “mopac”, …), which names its <root>/<program>.ini.

  • cmd ([str]) – The command template; {code}, {code_dir}, {NTASKS}, … are filled in by the back end from the configuration and the task’s share of the resources.

  • files ({str: str or bytes}) – Input files to write before running.

  • return_files ([str]) – Globs of the files to keep; "@subdir+name" moves a file into a subdirectory, as with Base.run().

  • resources (Resources)

  • env ({str: str}) – Extra environment variables.

  • in_situ (bool or None) – True runs in the task directory, leaving output there to watch; False in a temporary directory, copying back only return_files; None picks per Base.run() (scratch under a scheduler, in place otherwise).

  • shell (bool)

  • input_data (str or None) – Data for the standard input.

  • estimated_seconds (float or None) – The step’s estimate of the cost, used by the inline rule for tiny tasks.

  • target (str or None) – Reserved: None is the job’s target.

  • directory (str or Path or None) – Where the task runs and its results land. None means <step dir>/tasks/<key>/, the default for fan-out; a step running a single calculation passes its own directory so its output stays where it always was.

  • config (dict or None) – A local override of the program’s configuration (the <program>.ini section), for steps that resolve it themselves. Remote back ends ignore it.

  • fingerprint (str or None) – Identifies the inputs for restart. None hashes cmd and files; give one when the inputs carry run-dependent text (core counts, absolute paths).

  • success_text ({str: str or [str]} or None) – For codes whose return code does not show failure (ORCA exits 0 after an error termination): each file must contain its text (or all of its texts), or the task failed, gets no DONE and is tried again on the next run.

cmd: list#
config: dict | None = None#
digest()[source]#

The fingerprint of the inputs, used to detect a changed task.

directory: str | Path | None = None#
env: dict#
estimated_seconds: float | None = None#
files: dict | None = None#
fingerprint: str | None = None#
in_situ: bool | None = None#
input_data: str | None = None#
key: str#
program: str#
resources: Resources#
return_files: list#
shell: bool = False#
success_text: dict | None = None#
target: str | None = None#
class seamm_exec.tasks.TaskBackend(*args, **kwargs)[source]#

Bases: Protocol

What runs tasks: the LocalPool now; a scheduler or TaskServer later.

A back end may also offer wait(ids, timeout) to block until one of ids changes state, reattach(records) -> {key: state} to recover tasks an earlier evaluator submitted, capacity() and has_program(task).

cancel(ids: list) → None[source]#
fetch(task: Task, backend_id: str) → TaskResult[source]#
name: str#
status(ids: list) → dict[source]#
submit(tasks: list, directories: list, on_start=None) → list[source]#

Submit tasks; return their ids. on_start(task, info), if given, is called when a task’s process starts (info is what is needed to find it again, e.g. its process group).

class seamm_exec.tasks.TaskResult(key: str, state: str, returncode: int | None = None, stdout: str = '', stderr: str = '', directory: Path | None = None, files: dict = <factory>, attempts: int = 0, history: list = <factory>, in_situ: bool | None = None, run_directory: str | None = None, restored: bool = False, archive: Path | None = None, reason: str | None = None, raw: dict | None = None)[source]#

Bases: object

The outcome of a task.

Variables:
  • key (str)

  • state (str) – finished | failed | cancelled | lost

  • returncode (int or None)

  • stderr (stdout,)

  • directory (Path or None) – Where the task’s results are (for an archived task, where they were).

  • files ({str: str or bytes}) – The returned files’ contents.

  • attempts (int) – How many times the task has been submitted, over all runs.

  • history ([dict]) – One record per attempt.

  • in_situ (bool or None) – Whether it ran in place.

  • run_directory (str or None) – Where it actually ran (a scratch directory when not in situ).

  • restored (bool) – True when the result came from an earlier run rather than this one.

  • reason (str or None) – Why a task failed: its return code, a failed success check, …

  • archive (Path or None) – The tar holding the task’s directory, once archived.

  • raw (dict or None) – The Base.run()-style dictionary, for this run’s results.

archive: Path | None = None#
attempts: int = 0#
directory: Path | None = None#
files: dict#
classmethod from_raw(key, raw, directory=None)[source]#

Make a result from the dictionary Base._run_task() returns.

history: list#
in_situ: bool | None = None#
key: str#
property ok#
raw: dict | None = None#
reason: str | None = None#
restored: bool = False#
returncode: int | None = None#
run_directory: str | None = None#
state: str#
stderr: str = ''#
stdout: str = ''#
class seamm_exec.tasks.TaskSet(node=None, target=None, *, directory=None, executor=None, backend=None, local=None, root=None, manifest=True, archive=False, bundle_tasks=None, bundle_walltime=None, max_attempts=3, max_lost_retries=2, inline_below=None, poll_interval=1.0)[source]#

Bases: object

What a step uses: add tasks, then iterate over their results.

run() yields finished tasks from earlier runs first (from their DONE markers, never recomputed), then submits the rest and yields each result as it completes. Within a run a task that fails (nonzero return code) is not retried; a lost one is, up to max_lost_retries. Across runs a failed or lost task is tried again until it has had max_attempts attempts in all.

Parameters:
  • node (seamm.Node, optional) – The step. Gives the step directory, the executor, root and the job directory.

  • target (str or seamm_scheduler.TargetSection, optional) – The job’s target. Default: found by seamm_exec.targets.find_target() (<job dir>/target.json, then $SEAMM_TARGET), else the local pool. Ignored when backend is given.

  • directory (str or Path, optional) – The step directory, overriding node.directory.

  • executor (seamm_exec.Base, optional) – Overrides node.flowchart.executor.

  • backend (TaskBackend, optional) – Where tasks go. Default: the target’s back end (a SchedulerBackend for tasks = queue), else a LocalPool for this machine.

  • local (LocalPool, optional) – The evaluator’s own pool, for the inline rule. Defaults to backend when that is a LocalPool, else one is made on demand.

  • root (str or Path, optional) – Where the <program>.ini files are. Default node.global_options["root"].

  • manifest (bool = True) – Keep the manifest and DONE markers. False only for Base.run().

  • archive (bool = False) – Pack each bundle’s task directories into tasks/<bundle>.tar once all its tasks are done, and remove the directories.

  • bundle_tasks (int, optional) – Tasks per bundle, in the order added. Default: the target’s bundle_tasks; else one bundle (one task per bundle on a queue without bundle_walltime).

  • bundle_walltime (float, optional) – Start a new bundle when the tasks’ walltimes (or estimates) would add up to more than this many seconds. Default: the target’s.

  • max_attempts (int = 3) – Attempts per task over all runs.

  • max_lost_retries (int = 2) – Resubmissions of a lost task within one run.

  • inline_below (float, optional) – Tasks estimated to take less than this many seconds run in the local pool rather than on a remote back end, if their program is installed here. Default: the target’s inline_below, else 60.

  • poll_interval (float) – Seconds between status checks of a back end without wait().

add(task)[source]#

Add a task. Keys must be unique and filesystem-safe.

property backend#

The back end for this TaskSet’s tasks (the job’s target).

capacity()[source]#

The back end’s capacity, so a step can size its tasks beforehand.

Returns:

{"cores": int, "memory": bytes, "ngpus": int}

Return type:

dict

property local#

The evaluator’s own LocalPool.

marker_directory(key)[source]#

<step dir>/tasks/<key>, holding the task’s DONE.

route(task)[source]#

The back end for a task: the job’s target, except for tiny tasks.

The inline rule: a task whose estimated_seconds is below inline_below runs in the evaluator’s own pool instead of a remote back end, provided its program is installed here and the task fits the pool: a task asking for more cores than the evaluator has goes to the back end however cheap it is, since its input may already say how many ranks to use (seamm_exec#38).

run() → Iterator[TaskResult][source]#

Submit what is not done, wait, and yield results as they finish.

summary()[source]#

Counts by state, for the Dashboard and the step’s report.

Returns:

{"total": n, "finished": n, "failed": n, ...}

Return type:

dict

task_directory(task)[source]#

Where a task runs and its results land.

property tasks_directory#

<step dir>/tasks

seamm_exec.tasks.WORKER_SCRIPT = PosixPath('/home/runner/work/seamm_exec/seamm_exec/seamm_exec/task_worker.py')#

The pure-Python bundle worker, run by path (python task_worker.py bundle.json)

seamm_exec.tasks.run_task(task, node=None, **kwargs)[source]#

Run one task through a TaskSet and return its result.

The convenience for a step that runs a single calculation: it gets the manifest, restart and the pool’s handling of the process like any other task. kwargs go to TaskSet (e.g. directory= when the task runs somewhere other than node.directory).

Return type:

TaskResult

Module contents#

Classes to execute background codes for SEAMM

seamm_exec.get_executor(executor)[source]#

Return an object of the executor requested.

Parameters:

executor (str) – The name of the executor.

Return type:

instance of executor