{ "cells": [ { "cell_type": "markdown", "metadata": {}, "source": [ "# A small compute campaign\n", "\n", "A workspace holds the durable state of a campaign, a job is one input-specific calculation, and a manager drives each job through its runner steps. `collect` turns finished jobs into structured outputs that can be inspected or stored. This notebook is the in-notebook miniature of the [workflow CLI quickstart](https://docs.httk.org/httk-workflow/dev/main/quickstart/), using only a local temporary directory and a mock VASP." ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "import shlex\n", "import sys\n", "import tempfile\n", "from pathlib import Path\n", "\n", "from httk.core import DataRecord\n", "from httk.store import Backend, SqlStore\n", "from httk.workflow import TaskManager, Workspace\n", "from httk.workflow.collecting import collect\n", "from httk.workflow.scaffold import new_job\n", "\n", "tmpdir = Path(tempfile.mkdtemp(prefix=\"httk-campaign-\"))\n", "package = tmpdir / \"campaign\"\n", "package.mkdir()\n", "(package / \"httk_workflow.toml\").write_text(\n", " \"\"\"[workflow]\n", "id = 'notebook.campaign'\n", "\n", "[workflow.runner]\n", "entry = 'run'\n", "steps = ['prepare', 'run', 'publish']\n", "initial_step = 'prepare'\n", "data_mode = 'transactional'\n", "\n", "[workflow.collect]\n", "file = 'collect.py'\n", "\n", "[workflow.outputs.energy]\n", "entry_type = '_httk_records'\n", "role = 'energy'\n", "\"\"\",\n", " encoding=\"utf-8\",\n", ")\n", "\n", "# A deterministic stand-in for VASP: it reads POSCAR and writes familiar result files.\n", "(package / \"mock_vasp.py\").write_text(\n", " \"\"\"from pathlib import Path\n", "\n", "poscar = Path('POSCAR')\n", "label = poscar.read_text(encoding='utf-8').splitlines()[0].strip()\n", "energy = {'Si-A': -10.0, 'Si-B': -10.25, 'Si-C': -10.5}[label]\n", "Path('OUTCAR').write_text(f'mock VASP energy {energy} eV', encoding='utf-8')\n", "Path('OSZICAR').write_text(str(energy), encoding='utf-8')\n", "Path('CONTCAR').write_text(poscar.read_text(encoding='utf-8'), encoding='utf-8')\n", "Path('vasprun.xml').write_text('', encoding='utf-8')\n", "Path('energy.txt').write_text(str(energy), encoding='utf-8')\n", "\"\"\",\n", " encoding=\"utf-8\",\n", ")\n", "\n", "# The runner follows prepare -> run -> publish, as a packaged VASP runner does.\n", "(package / \"run\").write_text(\n", " \"\"\"#!/usr/bin/env python3\n", "import shlex\n", "import shutil\n", "from httk.workflow import Runner\n", "\n", "workflow = Runner('notebook.campaign')\n", "\n", "@workflow.step\n", "def prepare(a):\n", " shutil.copyfile(a.payload / 'files/POSCAR', a.workdir / 'POSCAR')\n", " a.advance('run')\n", "\n", "@workflow.step\n", "def run(a):\n", " result = a.run(shlex.split(str(a.setting('vasp.command'))))\n", " if result.returncode:\n", " a.fail('mock.failed', 'the mock VASP failed')\n", " return\n", " a.advance('publish')\n", "\n", "@workflow.step\n", "def publish(a):\n", " a.put(a.workdir / 'energy.txt', 'campaign/energy.txt')\n", " a.succeed()\n", "\n", "raise SystemExit(workflow.main())\n", "\"\"\",\n", " encoding=\"utf-8\",\n", ")\n", "(package / \"run\").chmod(0o755)\n", "\n", "(package / \"collect.py\").write_text(\n", " \"\"\"from httk.core import DataRecord\n", "\n", "def collect(record):\n", " energy = float((record.workdir / 'energy.txt').read_text(encoding='utf-8'))\n", " return {'energy': DataRecord.from_value('urn:notebook:energy', 'energy', energy)}\n", "\"\"\",\n", " encoding=\"utf-8\",\n", ")" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "workspace = Workspace.initialize(tmpdir / 'workspace')\n", "mock_command = f'{sys.executable} {package / \"mock_vasp.py\"}'\n", "workspace.set_setting('vasp.command', mock_command)\n", "\n", "poscars = []\n", "for label in ('Si-A', 'Si-B', 'Si-C'):\n", " poscar = tmpdir / f'POSCAR.{label}'\n", " poscar.write_text(\n", " f'{label}\\n1.0\\n2 0 0\\n0 2 0\\n0 0 2\\nSi\\n1\\nDirect\\n0 0 0\\n',\n", " encoding='utf-8',\n", " )\n", " poscars.append(poscar)\n", "\n", "jobs = [\n", " new_job(workspace, package, files={'POSCAR': poscar}, tag=f'case-{index}')\n", " for index, poscar in enumerate(poscars)\n", "]\n", "assert len(jobs) == 3" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "with TaskManager(workspace, heartbeat_interval=0.01, maximum_workers=3) as manager:\n", " manager.run_until_idle(timeout=30)\n", "\n", "states = [workspace.find_marker_by_id(job.job_id).kind for job in jobs]\n", "assert states == ['succeeded'] * 3, states\n", "print('job states:', states)" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "collected = list(collect(workspace, allow_job_collector=True))\n", "assert len(collected) == len(jobs)\n", "energies = []\n", "for item in collected:\n", " energy = item.outputs['energy']\n", " assert isinstance(energy, DataRecord)\n", " energies.append(energy.value)\n", "actual_tags = sorted(job.tag for job in jobs)\n", "assert actual_tags == ['case-0', 'case-1', 'case-2']\n", "print('campaign tags:', actual_tags)\n", "print('campaign energies:', sorted(energies))\n", "assert sorted(energies) == [-10.5, -10.25, -10.0]" ] }, { "cell_type": "code", "execution_count": null, "metadata": {}, "outputs": [], "source": [ "# Store selected collected outputs; the complete CLI `collect --into` path also stores provenance.\n", "database = Backend.sqlite(tmpdir / 'results.sqlite')\n", "store = SqlStore(database, entry_records={})\n", "for item in collected:\n", " store.save(item.outputs['energy'])\n", "\n", "search = store.searcher()\n", "record = search.variable(DataRecord)\n", "search.add(record.name == 'energy')\n", "stored_values = sorted(row.record.value for row in search.results(record=record))\n", "assert stored_values == [-10.5, -10.25, -10.0]\n", "print('stored energies:', stored_values)\n", "database.dispose()" ] }, { "cell_type": "markdown", "metadata": {}, "source": [ "At scale, the same workspace/job boundary supports [remote transfer and managers](https://docs.httk.org/httk-workflow/dev/main/workflow_cli/), partitioning a large run into [campaign workspaces](https://docs.httk.org/httk-workflow/dev/main/campaigns/), and [precheck](https://docs.httk.org/httk-workflow/dev/main/taskmanager/) before submitting expensive work. The local mock keeps this example fast; a deployment replaces only the workspace setting and runner details." ] } ], "metadata": { "kernelspec": { "display_name": "Python 3", "language": "python", "name": "python3" }, "language_info": { "codemirror_mode": { "name": "ipython", "version": 3 }, "file_extension": ".py", "mimetype": "text/x-python", "name": "python", "nbconvert_exporter": "python", "pygments_lexer": "ipython3", "version": "3.12.3" } }, "nbformat": 4, "nbformat_minor": 5 }