Skip to content
Draft
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
2 changes: 2 additions & 0 deletions docs/environment.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ dependencies:
- git+https://github.com/OpenFreeEnergy/kartograf@main
- git+https://github.com/OpenFreeEnergy/konnektor@main
- git+https://github.com/OpenFreeEnergy/lomap@main
- git+https://github.com/OpenFreeEnergy/exorcist@main


# These are added automatically by RTD, so we include them here
# for a consistent environment.
Expand Down
1 change: 1 addition & 0 deletions docs/guide/execution/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -12,3 +12,4 @@ then :ref:`reading on the available Python functions<reference_execution>`.
.. toctree::
quickrun_execution
execution_theory
warehouse
73 changes: 73 additions & 0 deletions docs/guide/execution/warehouse.rst

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm a bit lost with this documentation, it looks like mostly dev orientated docs, but it doesn't actually give me any insight as to how things will behave within a Protocol.

Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
Data Handling with Warehouse
==============================

**openfe**'s ``Warehouse`` defines the interface for an execution engine to store and access data during execution.

A Warehouse is any instance of a derived class of the abstract :class:`.WarehouseBaseClass`. In other words, *where* the data is stored is decided by the derived class, but *how* the data is accessed is defined by ``WarehouseBaseClass``.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This jumps in directly into the class but doesn't actually give users any info about what warehouse does - some examples and why it's needed before this line would be useful.


You can think of the ``WarehouseBaseClass`` as a set of specifications that must be met by a Warehouse implementation (subclass), such that any openfe Protocol can then interact appropriately with its data.

For example, **openfe** Protocols require several types of storage - scratch, setup, and result.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Distinction between setup and result here is a bit confusing. I assume by "result" you mean "anything that could be useful for a user to access in the future", and "setup" is "everything that is generated as an input to a further simulation in the dag"?


Naively, we could require that all three of these storage types be filesystem directories that can be accessed locally by the Protocol, but this significantly limits the ways in which the Protocol can be executed, e.g. a Protocol could not store **result** data on a remote machine or cloud storage.

Where a Warehouse stores its data is defined by a :class:`.WarehouseStores` object, which is a small `TypedDict <https://typing.python.org/en/latest/spec/typeddict.html>`_ containing ``'setup'`` and ``'result'`` keys (note that ``scratch`` is not a key, *must* be locally-accessible file storage) that correspond to gufe ``ExternalStorage`` objects.
The type of :class:`.ExternalStorage` objects used by a Warehouse implementation are where the the code author has flexibility to choose *where* data is stored.


The below example implementation, :class:`.FileSystemWarehouse`, is a derived class that inherits from ``WarehouseBaseClass``.
This is a simple example of how to construct a Warehouse given a root directory, which is uses to create ``WarehouseStores``.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
This is a simple example of how to construct a Warehouse given a root directory, which is uses to create ``WarehouseStores``.
This is a simple example of how to construct a Warehouse given a root directory, which is used to create ``WarehouseStores``.


.. TODO: reference FileStorage gufe docs

.. code-block::

class FileSystemWarehouse(WarehouseBaseClass):
"""Warehouse implementation using local filesystem storage.

Provides a file-based storage backend for GufeTokenizable objects
organized in a directory structure.

Parameters
----------
root_dir : str, optional
Root directory for the warehouse storage, by default "warehouse/".

Notes
-----
Creates a "setup/" subdirectory within the root directory for storing
setup-related objects. Future versions may include additional stores
for results and other data types.
"""

def __init__(self, root_dir: str = "warehouse"):
setup_store = FileStorage(f"{root_dir}/setup")

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it worth linking out to the FileStorage implementation?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

yes I think that's a good idea

result_store = FileStorage(f"{root_dir}/result")
stores = WarehouseStores(setup=setup_store, result=result_store)
super().__init__(stores)


Using this new ``FileSystemWarehouse``, we can now access the data cleanly and reproducibly in a way that abstracts away the use of a filesystem.

without Warehouse, using only the filesystem:

.. code-block::

root_dir = os.Path("my_warehouse")
setup = root / "setup"
result = root / "result"


with Warehouse:

.. code-block::

from openfe.storage import FileSystemWarehouse
my_warehouse = FileSystemWarehouse(root_dir="my_warehouse")

...

my_warehouse.store_result_tokenizable(result)
.. TODO: add example of dropping in non-filesystem storage once gufe supports it

.. For more information about what types of storage are available, see <gufe storage docs>.
7 changes: 4 additions & 3 deletions docs/reference/api/defining_and_executing_simulations.rst
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,11 @@ Executing Simulations
:noindex:

.. autosummary::
:nosignatures:
:toctree: generated/
:recursive:

execute_DAG
storage

General classes
---------------
Expand All @@ -32,10 +33,10 @@ General classes
ProtocolUnitFailure
ProtocolDAGResult

Specialised classes
Specialized classes
-------------------

These classes are abstract classes that are specialised (subclassed) for an individual Protocol.
These classes are abstract classes that are specialized (subclassed) for an individual Protocol.

.. module:: openfe
:noindex:
Expand Down
1 change: 1 addition & 0 deletions environment.yml
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ dependencies:
- git+https://github.com/OpenFreeEnergy/kartograf@main # needs gufe which pins to pydantic < 2.13. update this once gufe v1.13 is released on conda-forge
- git+https://github.com/OpenFreeEnergy/konnektor@main
- git+https://github.com/OpenFreeEnergy/gufe@main
- git+https://github.com/OpenFreeEnergy/exorcist@main
- run_constrained:
# drop this pin when handled upstream in espaloma-feedstock
- smirnoff99frosst>=1.1.0.1 #https://github.com/openforcefield/smirnoff99Frosst/issues/109
23 changes: 23 additions & 0 deletions news/warehouse.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
**Added:**

* Added ``openfe.Warehouse``, an interface for storing and accessing data during simulation execution (`PR #1864 <https://github.com/OpenFreeEnergy/openfe/pull/1864>`_).

**Changed:**

* <news item>

**Deprecated:**

* <news item>

**Removed:**

* Removed the unused methods ``metadatastore``, ``resultclient``, and ``resultserver`` from ``openfe.storage``.

**Fixed:**

* <news item>

**Security:**

* <news item>
244 changes: 244 additions & 0 deletions src/openfe/orchestration/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,244 @@
"""Task orchestration utilities backed by Exorcist and a warehouse."""

from collections.abc import Iterable
from dataclasses import dataclass
from pathlib import Path

from exorcist.taskdb import TaskStatusDB
from gufe.protocols.protocoldag import _pu_to_pur
from gufe.protocols.protocolunit import (
Context,
ProtocolUnit,
ProtocolUnitResult,
)
from gufe.storage.externalresource.base import ExternalStorage
from gufe.storage.externalresource.filestorage import FileStorage
from gufe.tokenization import GufeKey

from openfe.storage.warehouse import FileSystemWarehouse

from .exorcist_utils import (
alchemical_network_to_task_graph,
build_task_db_from_alchemical_network,
)


@dataclass
class Worker:
"""Execute protocol units from an Exorcist task database.

Parameters
----------
warehouse : FileSystemWarehouse
Warehouse used to load queued tasks and store execution results.
task_db_path : pathlib.Path, default=Path("./warehouse/tasks.db")
Path to the Exorcist SQLite task database.
"""

warehouse: FileSystemWarehouse
task_db_path: Path = Path("./warehouse/tasks.db")

_RESULT_INDEX_PREFIX = "protocol_unit_results"
_TASK_WORKDIR_PREFIX = "task_workdirs"

@staticmethod
def _collect_protocol_unit_keys(value: object) -> set[GufeKey]:
"""Collect `ProtocolUnit` keys from nested unit inputs."""

if isinstance(value, ProtocolUnit):
return {value.key}

found: set[GufeKey] = set()
items: Iterable # TODO: update this to dict_values | list after python 3.13 min?
if isinstance(value, dict):
items = value.values()
elif isinstance(value, list):
items = value
else:
return found

for item in items:
found.update(Worker._collect_protocol_unit_keys(item))
return found

@classmethod
def _result_index_location(cls, source_key: GufeKey) -> str:
return f"{cls._RESULT_INDEX_PREFIX}/{source_key}"

@classmethod
def _task_workdir_name(cls, taskid: str) -> str:
return taskid.replace(":", "__")

def _task_workspace_paths(
self, taskid: str, scratch_root: Path, shared_root: Path
) -> tuple[Path, Path]:
workdir_name = self._task_workdir_name(taskid)
task_scratch = scratch_root / self._TASK_WORKDIR_PREFIX / workdir_name
task_shared = shared_root / self._TASK_WORKDIR_PREFIX / workdir_name
return task_scratch, task_shared

def _store_result_index(self, result: ProtocolUnitResult) -> None:
shared_store: ExternalStorage = self.warehouse.stores["shared"]
location = self._result_index_location(result.source_key)
shared_store.store_bytes(location, str(result.key).encode("utf-8"))

def _load_result_from_index(self, source_key: GufeKey) -> ProtocolUnitResult | None:
shared_store: ExternalStorage = self.warehouse.stores["shared"]
location = self._result_index_location(source_key)

if not shared_store.exists(location):
return None

with shared_store.load_stream(location) as stream:
result_key = stream.read().decode("utf-8").strip()

loaded = self.warehouse.load_result_tokenizable(GufeKey(result_key))
if isinstance(loaded, ProtocolUnitResult):
return loaded

return None

def _scan_result_store_for_sources(
self, source_keys: set[GufeKey]
) -> dict[GufeKey, ProtocolUnitResult]:
found: dict[GufeKey, ProtocolUnitResult] = {}

for location in self.warehouse.result_store.iter_contents():
if len(found) == len(source_keys):
break

loaded = self.warehouse.load_result_tokenizable(GufeKey(location))
if not isinstance(loaded, ProtocolUnitResult):
continue

source_key = loaded.source_key
if source_key in source_keys and source_key not in found:
found[source_key] = loaded

return found

def _build_input_result_mapping(self, unit: ProtocolUnit) -> dict[GufeKey, ProtocolUnitResult]:
required_keys = self._collect_protocol_unit_keys(unit.inputs)
if not required_keys:
return {}

results: dict[GufeKey, ProtocolUnitResult] = {}
unresolved = set(required_keys)

for source_key in required_keys:
loaded = self._load_result_from_index(source_key)
if loaded is not None:
results[source_key] = loaded
unresolved.discard(source_key)

if unresolved:
scanned = self._scan_result_store_for_sources(unresolved)
for source_key, loaded in scanned.items():
results[source_key] = loaded
self._store_result_index(loaded)
unresolved.discard(source_key)

if unresolved:
missing_keys = ", ".join(sorted(str(k) for k in unresolved))
raise RuntimeError(
"Missing ProtocolUnitResult(s) for dependency key(s): "
f"{missing_keys}. Ensure upstream tasks completed successfully."
)

return results

def _checkout_task(self) -> tuple[TaskStatusDB, str, ProtocolUnit] | None:
"""Check out one available task and load its protocol unit.

Returns
-------
tuple[TaskStatusDB, str, ProtocolUnit] or None
The open database connection, checked-out task ID, and corresponding
protocol unit, or ``None`` if no task is currently available.
The caller is responsible for calling ``mark_task_completed`` on the
returned database using the returned task ID.
"""

db: TaskStatusDB = TaskStatusDB.from_filename(self.task_db_path)
# The format for the taskid is "ProtocolUnit-<HASH>"
taskid = db.check_out_task()
if taskid is None:
return None

protocol_unit_key = taskid
unit = self.warehouse.load_task(GufeKey(protocol_unit_key))
return db, taskid, unit

def _get_task(self) -> tuple[str, ProtocolUnit]:
"""Return the next available task ID and protocol unit.

Returns
-------
tuple[str, ProtocolUnit]
The checked-out task ID and corresponding protocol unit.

Raises
------
RuntimeError
Raised when no task is available in the task database.
"""

task = self._checkout_task()
if task is None:
raise RuntimeError("No AVAILABLE tasks found in the task database.")
db, taskid, unit = task
return taskid, unit

def execute_unit(self, scratch: Path) -> tuple[str, ProtocolUnitResult] | None:
"""Execute one checked-out protocol unit and persist its result.

Parameters
----------
scratch : pathlib.Path
Scratch directory passed to the protocol execution context.

Returns
-------
tuple[str, ProtocolUnitResult] or None
The task ID and execution result for the processed task, or
``None`` if no task is currently available.

Raises
------
Exception
Re-raises any exception thrown during protocol unit execution after
marking the task as failed.
"""

# 1. Get task/unit
task = self._checkout_task()
if task is None:
return None
db, taskid, unit = task
# 2. Construct the context
# NOTE: On changes to context (gufe PR #753), this can easily be replaced with external storage objects
# However, to satisfy the current work, we will use this implementation where we
# force the use of a FileSystemWarehouse and in turn can assert that an object is FileStorage.
shared_store = self.warehouse.stores["shared"]
if not isinstance(shared_store, FileStorage):
raise TypeError("Expected a FileStorage backend for the shared store")
shared_root_dir = shared_store.root_dir
task_scratch, task_shared = self._task_workspace_paths(taskid, scratch, shared_root_dir)
task_scratch.mkdir(parents=True, exist_ok=True)
task_shared.mkdir(parents=True, exist_ok=True)
ctx = Context(task_scratch, shared=task_shared)
# 3. Execute unit
try:
results = self._build_input_result_mapping(unit)
inputs = _pu_to_pur(unit.inputs, results)
result = unit.execute(context=ctx, **inputs)
except Exception:
db.mark_task_completed(taskid, success=False)
raise

db.mark_task_completed(taskid, success=result.ok())
# 4. output result to warehouse
# TODO: we may need to end up handling namespacing on the warehouse side for tokenizables
self.warehouse.store_result_tokenizable(result)
self._store_result_index(result)
return taskid, result
Loading
Loading