diff --git a/CHANGELOG.md b/CHANGELOG.md index e6a12bd5..42adde16 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added - Server: + - workflow Composer: ordered, resumable task stages reuse existing Docker/SLURM jobs with independently snapshotted resource policies; AlphaFold MSA/features now run CPU-only before GPU model/relax. - version-2 scientific workspaces: modular task-selected RFdiffusion modes, Mol* residue selection (with result controls hidden by default), declarative linked result views, bounded table pages, and EASIFA table-to-structure mapping. Breaking: every task type must declare `input_workspace` explicitly — startup fails closed when a custom registry omits it. - runner: patch RFdiffusion's checkpoint-override parsing (`bool("false")` is True), so binder runs respect `preprocess.sidechain_input=false` instead of crashing in the broken upstream sidechain path. - runner: new alphafold family (official google-deepmind/alphafold @ c77e5d2a) — monomer/pTM/multimer presets, full_dbs MSA, Amber relaxation (best by default); DBs and the 2022-12-06 params release ro-mounted from /mnt/db. @@ -56,6 +57,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed - Server: + - AlphaFold multimer full-database runs now pass the required UniRef30 database path. + - prepared SLURM deploys build staged SIFs from the matching `:next` runner image instead of silently reusing `:latest`. + - create-task parameters: choice controls now serialize their selected value, preserving AlphaFold multimer and every other non-default select option. - Result polling/viewers: terminal task states reload once without overlapping polls; pending result pages keep polling; Mol* teardown completes before iframe removal, stale teardown continuations cannot replace newer previews, and preview loaders remain visible. - AlphaFold stages: drain the stderr translator before wrapper exit and preserve both process statuses so final stage markers cannot be lost. - SLURM stages: `srun -u`, unbuffered AlphaFold stderr, and a Python translator stream wrapper phases live, so `run_stage` records intermediates before job exit. diff --git a/server/README.md b/server/README.md index 7c905643..f1e39874 100644 --- a/server/README.md +++ b/server/README.md @@ -449,6 +449,27 @@ and a watchdog kills work that exceeds the snapshotted runtime. Invalid fields fail closed in the admin API and again at submission/launch rather than being silently discarded. +### Ordered workflows + +A task type may declare an ordered `workflow` whose stages reuse the same +runtime family and immutable task snapshot. The worker acts as a lightweight +Composer: it submits one existing `Job` at a time, validates that allocation's +outputs through the runner contract, persists the stage/job state, and releases +the next stage only after success. Each stage has an independently snapshotted +resource policy in the configuration UI. + +```text +Submission -> Composer -> [CPU: MSA + features] -> [GPU: model + relax] -> Results + | features.pkl | ranked_0.pdb + +---- persisted --------+ +``` + +AlphaFold is the first composed task. `alphafold.features` runs MSA and feature +construction without GPU GRES or Apptainer `--nv`; after `features.pkl` is +validated, `alphafold.model` receives the GPU allocation for inference and +optional relaxation. Restarts cancel only the active allocation and resume at +the first incomplete stage. + ## 5. Authentication The server uses Bearer-token authentication (replaces the old HTTP Basic Auth + `users.txt` model). @@ -1187,7 +1208,8 @@ job is enqueued. The allocation wrapper writes `REVODESIGN_JOB_ID=` as its first stdout line. The worker stores that real SLURM ID in the `slurm_job_id` column of the tasks table so cancellation and restart recovery can address the -allocation directly. +allocation directly. Composed tasks additionally persist all stage handles and +states in `workflow_state`; `slurm_job_id` remains the currently active handle. ### 13.8 Live Output diff --git a/server/config/task_types.yaml b/server/config/task_types.yaml index c675d614..1e84e6f3 100644 --- a/server/config/task_types.yaml +++ b/server/config/task_types.yaml @@ -754,6 +754,17 @@ task_types: featuring: "Building features" modeling: "Folding models" relaxing: "Amber relaxation" + workflow: + - name: features + display_name: "MSA and features" + requires_gpu: false + runner_args: ["-s", "features"] + stage_markers: ["msa_searching", "featuring"] + - name: model + display_name: "Model and relax" + requires_gpu: true + runner_args: ["-s", "model"] + stage_markers: ["modeling", "relaxing"] params: - name: "model_preset" type: "str" diff --git a/server/docker/runners/alphafold/Dockerfile b/server/docker/runners/alphafold/Dockerfile index d8f4ecb5..b539082b 100644 --- a/server/docker/runners/alphafold/Dockerfile +++ b/server/docker/runners/alphafold/Dockerfile @@ -4,9 +4,11 @@ FROM alpine/git:2.47.2 AS alphafold-source ARG ALPHAFOLD_REPO=https://github.com/google-deepmind/alphafold.git ARG ALPHAFOLD_REF=c77e5d2a8961d1a353632c462914ff0a32a950f6 +COPY ./docker/runners/alphafold/staged_pipeline.patch /tmp/staged_pipeline.patch RUN git init /opt/alphafold && git -C /opt/alphafold remote add origin ${ALPHAFOLD_REPO} && \ git -C /opt/alphafold fetch --depth 1 origin ${ALPHAFOLD_REF} && git -C /opt/alphafold checkout --detach FETCH_HEAD && \ - rm -rf /opt/alphafold/.git + git -C /opt/alphafold apply --check /tmp/staged_pipeline.patch && \ + git -C /opt/alphafold apply /tmp/staged_pipeline.patch && rm -rf /opt/alphafold/.git FROM --platform=linux/amd64 python:3.11-slim diff --git a/server/docker/runners/alphafold/run.sh b/server/docker/runners/alphafold/run.sh index 80c6ae09..268ea62b 100644 --- a/server/docker/runners/alphafold/run.sh +++ b/server/docker/runners/alphafold/run.sh @@ -4,9 +4,11 @@ set -e task_context_src="${TASK_CONTEXT_SRC:-/app/revocompute/task_context.sh}" [[ -f "$task_context_src" ]] && source "$task_context_src" -usage() { echo "Usage: $0 -i -o "; exit 1; } -while getopts ":i:o:" opt; do case "${opt}" in i) input_file=$OPTARG ;; o) output_dir=$OPTARG ;; ?) usage ;; esac; done +usage() { echo "Usage: $0 -i -o [-s all|features|model]"; exit 1; } +run_stage=all +while getopts ":i:o:s:" opt; do case "${opt}" in i) input_file=$OPTARG ;; o) output_dir=$OPTARG ;; s) run_stage=$OPTARG ;; ?) usage ;; esac; done [[ -z "${input_file:-}" || -z "${output_dir:-}" ]] && usage +[[ "$run_stage" =~ ^(all|features|model)$ ]] || usage input_file=$(readlink -f "$input_file"); output_dir=$(readlink -f "$output_dir") [[ ! -f "$input_file" ]] && { echo "Task manifest not found: $input_file"; exit 1; } mkdir -p "$output_dir" @@ -18,6 +20,15 @@ NUM_MULTIMER=$(_parse_param num_multimer_predictions_per_model 1) MODELS_TO_RELAX=$(_parse_param models_to_relax best) BENCHMARK=$(_parse_param benchmark false) fasta_path=$(primary_input) +fasta_name=$(basename "$fasta_path") +fasta_name=${fasta_name%.*} +features_path="${output_dir}/${fasta_name}/features.pkl" +features_marker="${output_dir}/.alphafold-features-complete" +use_gpu_relax=true +[[ "$run_stage" == features ]] && use_gpu_relax=false +if [[ "$run_stage" == model ]]; then + [[ -s "$features_path" && -f "$features_marker" ]] || { echo "Validated AlphaFold features are missing" >&2; exit 1; } +fi # A100 memory behaviour: let JAX overcommit via unified memory. export TF_FORCE_UNIFIED_MEMORY=1 @@ -34,13 +45,15 @@ af_args=( "--db_preset=${DB_PRESET}" "--model_preset=${MODEL_PRESET}" "--models_to_relax=${MODELS_TO_RELAX}" - "--use_gpu_relax=true" + "--use_gpu_relax=${use_gpu_relax}" + "--run_stage=${run_stage}" "--benchmark=${BENCHMARK}" "--bfd_database_path=${DB}/bfd/bfd_metaclust_clu_complete_id30_c90_final_seq.sorted_opt" "--mgnify_database_path=${DB}/mgnify/mgy_clusters.fa" "--template_mmcif_dir=${DB}/pdb_mmcif/mmcif_files" "--obsolete_pdbs_path=${DB}/pdb_mmcif/obsolete.dat" "--uniref90_database_path=${DB}/uniref90/uniref90.fasta" + "--uniref30_database_path=${DB}/uniref30_uc30/UniRef30_2022_02/UniRef30_2022_02" ) if [[ "$MODEL_PRESET" == "multimer" ]]; then af_args+=( @@ -51,7 +64,6 @@ if [[ "$MODEL_PRESET" == "multimer" ]]; then else af_args+=( "--pdb70_database_path=${DB}/pdb70/pdb70" - "--uniref30_database_path=${DB}/uniref30_uc30/UniRef30_2022_02/UniRef30_2022_02" ) fi @@ -90,6 +102,13 @@ if [[ ${alphafold_status} -ne 0 || ${translator_status} -ne 0 ]]; then exit "${translator_status}" fi +if [[ "$run_stage" == features ]]; then + [[ -s "$features_path" ]] || { echo "AlphaFold produced no features.pkl" >&2; exit 1; } + touch "$features_marker" + echo "AlphaFold feature construction complete." + exit 0 +fi + [[ -n "$(ls "${output_dir}"/*/ranked_0.pdb 2>/dev/null || true)" ]] || { echo "AlphaFold produced no ranked_0.pdb" >&2; exit 1; } touch "${output_dir}/task_finished" diff --git a/server/docker/runners/alphafold/staged_pipeline.patch b/server/docker/runners/alphafold/staged_pipeline.patch new file mode 100644 index 00000000..f6693615 --- /dev/null +++ b/server/docker/runners/alphafold/staged_pipeline.patch @@ -0,0 +1,116 @@ +diff --git a/run_alphafold.py b/run_alphafold.py +index 3d8a4f4..a41fb00 100644 +--- a/run_alphafold.py ++++ b/run_alphafold.py +@@ -231,6 +231,12 @@ flags.DEFINE_boolean( + 'recommended to enable if possible. GPUs must be available' + ' if this setting is enabled.', + ) ++flags.DEFINE_enum( ++ 'run_stage', ++ 'all', ++ ['all', 'features', 'model'], ++ 'Run the complete pipeline, feature construction only, or modeling only.', ++) + flags.DEFINE_integer( + 'jackhmmer_n_cpu', + # Unfortunately, os.process_cpu_count() is only available in Python 3.13+. +@@ -364,17 +370,22 @@ def predict_structure( + if not os.path.exists(msa_output_dir): + os.makedirs(msa_output_dir) + +- # Get features. +- t_0 = time.time() +- feature_dict = data_pipeline.process( +- input_fasta_path=fasta_path, msa_output_dir=msa_output_dir +- ) +- timings['features'] = time.time() - t_0 +- +- # Write out features as a pickled dictionary. + features_output_path = os.path.join(output_dir, 'features.pkl') +- with open(features_output_path, 'wb') as f: +- pickle.dump(feature_dict, f, protocol=4) ++ if FLAGS.run_stage == 'model': ++ logging.info('Reading precomputed features from %s', features_output_path) ++ with open(features_output_path, 'rb') as f: ++ feature_dict = pickle.load(f) ++ else: ++ t_0 = time.time() ++ feature_dict = data_pipeline.process( ++ input_fasta_path=fasta_path, msa_output_dir=msa_output_dir ++ ) ++ timings['features'] = time.time() - t_0 ++ with open(features_output_path, 'wb') as f: ++ pickle.dump(feature_dict, f, protocol=4) ++ if FLAGS.run_stage == 'features': ++ logging.info('Feature construction complete for %s', fasta_name) ++ return + + unrelaxed_pdbs = {} + unrelaxed_proteins = {} +@@ -667,36 +678,40 @@ def main(argv): + data_pipeline = monomer_data_pipeline + + model_runners = {} +- model_names = config.MODEL_PRESETS[FLAGS.model_preset] +- for model_name in model_names: +- model_config = config.model_config(model_name) +- if run_multimer_system: +- model_config.model.num_ensemble_eval = num_ensemble +- else: +- model_config.data.eval.num_ensemble = num_ensemble +- model_params = data.get_model_haiku_params( +- model_name=model_name, data_dir=FLAGS.data_dir +- ) +- model_runner = model.RunModel(model_config, model_params) +- for i in range(num_predictions_per_model): +- model_runners[f'{model_name}_pred_{i}'] = model_runner ++ amber_relaxer = None ++ if FLAGS.run_stage != 'features': ++ model_names = config.MODEL_PRESETS[FLAGS.model_preset] ++ for model_name in model_names: ++ model_config = config.model_config(model_name) ++ if run_multimer_system: ++ model_config.model.num_ensemble_eval = num_ensemble ++ else: ++ model_config.data.eval.num_ensemble = num_ensemble ++ model_params = data.get_model_haiku_params( ++ model_name=model_name, data_dir=FLAGS.data_dir ++ ) ++ model_runner = model.RunModel(model_config, model_params) ++ for i in range(num_predictions_per_model): ++ model_runners[f'{model_name}_pred_{i}'] = model_runner + +- logging.info( +- 'Have %d models: %s', len(model_runners), list(model_runners.keys()) +- ) ++ logging.info( ++ 'Have %d models: %s', len(model_runners), list(model_runners.keys()) ++ ) + +- amber_relaxer = relax.AmberRelaxation( +- max_iterations=RELAX_MAX_ITERATIONS, +- tolerance=RELAX_ENERGY_TOLERANCE, +- stiffness=RELAX_STIFFNESS, +- exclude_residues=RELAX_EXCLUDE_RESIDUES, +- max_outer_iterations=RELAX_MAX_OUTER_ITERATIONS, +- use_gpu=FLAGS.use_gpu_relax, +- ) ++ amber_relaxer = relax.AmberRelaxation( ++ max_iterations=RELAX_MAX_ITERATIONS, ++ tolerance=RELAX_ENERGY_TOLERANCE, ++ stiffness=RELAX_STIFFNESS, ++ exclude_residues=RELAX_EXCLUDE_RESIDUES, ++ max_outer_iterations=RELAX_MAX_OUTER_ITERATIONS, ++ use_gpu=FLAGS.use_gpu_relax, ++ ) + + random_seed = FLAGS.random_seed + if random_seed is None: +- random_seed = random.randrange(sys.maxsize // len(model_runners)) ++ random_seed = random.randrange( ++ sys.maxsize // max(len(model_runners), len(config.MODEL_PRESETS[FLAGS.model_preset])) ++ ) + logging.info('Using random seed %d for the data pipeline', random_seed) + + # Predict structure for each of the sequences. diff --git a/server/revocompute/app.py b/server/revocompute/app.py index 2b2f77e6..91373f87 100644 --- a/server/revocompute/app.py +++ b/server/revocompute/app.py @@ -198,6 +198,20 @@ def _add_security_headers(response): if manage_db.task_type_get(_tt.name) is None: manage_db.task_type_upsert(_tt.name, enabled=True) _log.info("Seeded task_type_config for %r (enabled=true)", _tt.name) + _parent_resources = manage_db.task_type_get(_tt.name) or {} + for _stage in _tt.workflow: + if manage_db.task_type_get(_stage.name) is None: + _initial = {"enabled": True} + if _stage.requires_gpu: + _initial.update( + { + key: value + for key, value in _parent_resources.items() + if key not in {"tool", "enabled"} and value is not None + } + ) + manage_db.task_type_upsert(_stage.name, **_initial) + _log.info("Seeded workflow resource profile %r", _stage.name) def _is_binary_file(path: str) -> bool: diff --git a/server/revocompute/db.py b/server/revocompute/db.py index 4e6b4bc4..3525f27a 100644 --- a/server/revocompute/db.py +++ b/server/revocompute/db.py @@ -87,6 +87,7 @@ def __init__(self, path: str): Column("input_form", Text), Column("slurm_job_id", String), Column("container_id", String), + Column("workflow_state", Text), ) Index("idx_tasks_uploaded_at", self.tasks_table.c.uploaded_at) self._initialize() @@ -106,7 +107,7 @@ def _initialize(self) -> None: # create_all does not add columns to existing tables — backfill # ones added after a table first shipped (idempotent). existing = {row[1] for row in conn.exec_driver_sql("PRAGMA table_info(tasks)")} - for column in ("container_id",): + for column in ("container_id", "workflow_state"): if column not in existing: conn.exec_driver_sql(f"ALTER TABLE tasks ADD COLUMN {column} VARCHAR") @@ -160,9 +161,9 @@ def upsert_task(self, md5sum: str, **fields) -> None: with self.engine.begin() as conn: conn.execute(stmt) - def update_task(self, md5sum: str, **fields) -> None: + def update_task(self, md5sum: str, **fields) -> bool: if not fields: - return + return False status = fields.get("status") if status: self._ensure_status(status) @@ -174,7 +175,35 @@ def update_task(self, md5sum: str, **fields) -> None: if status is None or (not self._is_deleted_status(status)): stmt = stmt.where(self.tasks_table.c.status.notin_(tuple(self.TERMINAL_STATUSES))) with self.engine.begin() as conn: - conn.execute(stmt) + return conn.execute(stmt).rowcount == 1 + + def claim_task_recovery(self, md5sum: str, *, expected_status: str) -> bool: + """Atomically move one orphaned active task out of recovery scans.""" + if expected_status not in {"queued", "running"}: + return False + stmt = ( + update(self.tasks_table) + .where( + self.tasks_table.c.md5sum == md5sum, + self.tasks_table.c.status == expected_status, + ) + .values(status="pending") + ) + with self.engine.begin() as conn: + return conn.execute(stmt).rowcount == 1 + + def claim_task_cancellation(self, md5sum: str, **fields) -> bool: + """Atomically cancel a task only while it remains active.""" + stmt = ( + update(self.tasks_table) + .where( + self.tasks_table.c.md5sum == md5sum, + self.tasks_table.c.status.in_(("pending", "queued", "running")), + ) + .values(status="cancelled", **fields) + ) + with self.engine.begin() as conn: + return conn.execute(stmt).rowcount == 1 def claim_task_cleanup( self, diff --git a/server/revocompute/resource_audit.py b/server/revocompute/resource_audit.py index b32a3ab6..ab87b42a 100644 --- a/server/revocompute/resource_audit.py +++ b/server/revocompute/resource_audit.py @@ -60,27 +60,32 @@ def main() -> int: failed = False for task_type in list_types(): _, runner = get_task_type(task_type.name) - values = task_values.get(task_type.name, {}) - if values.get("enabled") == 0: + task_config = task_values.get(task_type.name, {}) + if task_config.get("enabled") == 0: print(f"[RESOURCE] {task_type.name}: disabled (not audited)") continue - try: - resolved = resolve_resources( - values.get, - globals_.get, - requires_gpu=task_type.gpus, - allowed_queues=allowed, - default_timeout_seconds=runner.max_runtime_seconds, - ) - accelerator = resolved.gres or "cpu" - partition = resolved.partition or "scheduler-default" - print( - f"[RESOURCE] {task_type.name}: cpus={resolved.cpus} memory={resolved.memory} " - f"time={resolved.slurm_time} accelerator={accelerator} partition={partition}" - ) - except ResourceValidationError as exc: - failed = True - print(f"[RESOURCE] {task_type.name}: INVALID: {exc}", file=sys.stderr) + profiles = [(task_type.name, task_type.gpus)] + if task_type.workflow: + profiles = [(stage.name, stage.requires_gpu) for stage in task_type.workflow] + for profile_name, requires_gpu in profiles: + values = task_values.get(profile_name, task_config if requires_gpu else {}) + try: + resolved = resolve_resources( + values.get, + globals_.get, + requires_gpu=requires_gpu, + allowed_queues=allowed, + default_timeout_seconds=runner.max_runtime_seconds, + ) + accelerator = resolved.gres or "cpu" + partition = resolved.partition or "scheduler-default" + print( + f"[RESOURCE] {profile_name}: cpus={resolved.cpus} memory={resolved.memory} " + f"time={resolved.slurm_time} accelerator={accelerator} partition={partition}" + ) + except ResourceValidationError as exc: + failed = True + print(f"[RESOURCE] {profile_name}: INVALID: {exc}", file=sys.stderr) return 1 if failed else 0 diff --git a/server/revocompute/routes.py b/server/revocompute/routes.py index d5188a92..62afeb5c 100644 --- a/server/revocompute/routes.py +++ b/server/revocompute/routes.py @@ -718,13 +718,24 @@ def upload_file(): # skipcq: PY-R1000 -- route validation branches form one tra if tt.gpus and not g.current_user.get("allow_gpu_use"): return jsonify({"error": "GPU access required for this task type. Contact an administrator."}), 403 resource_policy = None + resource_policies: dict[str, Any] = {} if managedb is not None: try: - resource_policy = managedb.resolve_task_resources( - tt.name, - requires_gpu=tt.gpus, - default_timeout_seconds=runner.max_runtime_seconds, - ) + if tt.workflow: + resource_policies = { + stage.name: managedb.resolve_task_resources( + stage.name, + requires_gpu=stage.requires_gpu, + default_timeout_seconds=runner.max_runtime_seconds, + ) + for stage in tt.workflow + } + else: + resource_policy = managedb.resolve_task_resources( + tt.name, + requires_gpu=tt.gpus, + default_timeout_seconds=runner.max_runtime_seconds, + ) except ResourceValidationError as exc: logging.error("Resource policy rejected submission for %s: %s", task_type, exc) return jsonify({"error": "This task type has an invalid resource policy; contact an administrator."}), 503 @@ -807,6 +818,7 @@ def upload_file(): # skipcq: PY-R1000 -- route validation branches form one tra "submitted_at": datetime.now(tz=timezone.utc).isoformat(), "entities": entities, "resource_policy": resource_policy.public_dict() if resource_policy is not None else None, + "resource_policies": {name: policy.public_dict() for name, policy in resource_policies.items()}, "workspace": workspace_payload, } @@ -1157,9 +1169,20 @@ def cancel_task(md5sum): 400, ) - # SLURM tooling and the Docker socket live only in the worker container; - # the web process delegates the kill there. The task record flips to - # cancelled immediately either way. + now = time.time() + started_at = task.get("started_at") + walltime = (now - started_at) if started_at else None + if not task_store.claim_task_cancellation( + md5sum, + finished_at=now, + walltime=walltime, + error="Task cancelled by user", + ): + return jsonify({"error": "Task state changed before cancellation"}), 409 + task = task_store.get_task(md5sum) or task + + # Claim cancellation in the database before asking the worker to stop + # resources, so a workflow cannot launch its next stage in between. cancel_compute_resources.delay( slurm_job_id=str(task["slurm_job_id"]) if task.get("slurm_job_id") else None, container_id=str(task["container_id"]) if task.get("container_id") else None, @@ -1174,17 +1197,6 @@ def cancel_task(md5sum): logging.warning("Failed to revoke Celery task %s: %s", celery_id, exc) _delete_task_artifacts(task) - - now = time.time() - started_at = task.get("started_at") - walltime = (now - started_at) if started_at else None - task_store.update_task( - md5sum, - status="cancelled", - finished_at=now, - walltime=walltime, - error="Task cancelled by user", - ) return jsonify({"status": "cancelled", "md5sum": md5sum}), 200 @@ -2251,18 +2263,25 @@ def admin_get_config(): return jsonify({"error": "Configuration database not available"}), 500 task_configs = manage_db.task_type_all() type_map = {task_type.name: task_type for task_type in list_types()} + stage_map = { + stage.name: (task_type, stage) for task_type in type_map.values() for stage in task_type.workflow + } for config in task_configs: task_type = type_map.get(config["tool"]) - if task_type is None: + workflow_stage = stage_map.get(config["tool"]) + if task_type is None and workflow_stage is None: continue - config["display_name"] = task_type.display_name - config["requires_gpu"] = task_type.gpus + stage = workflow_stage[1] if workflow_stage else None + task_type = task_type or workflow_stage[0] + config["display_name"] = f"{task_type.display_name} / {stage.display_name}" if stage else task_type.display_name + config["requires_gpu"] = stage.requires_gpu if stage else task_type.gpus config["runtime_family"] = task_type.runtime.name + config["is_workflow_stage"] = stage is not None _, runner = _get_task_type(task_type.name) try: resolved = manage_db.resolve_task_resources( - task_type.name, - requires_gpu=task_type.gpus, + config["tool"], + requires_gpu=config["requires_gpu"], default_timeout_seconds=runner.max_runtime_seconds, ) config["effective_resources"] = resolved.public_dict() @@ -2336,6 +2355,10 @@ def admin_set_config(): known_tools = {entry["tool"] for entry in manage_db.task_type_all()} type_map = {task_type.name: task_type for task_type in list_types()} + profile_gpu = {name: task_type.gpus for name, task_type in type_map.items()} + profile_gpu.update( + {stage.name: stage.requires_gpu for task_type in type_map.values() for stage in task_type.workflow} + ) pending_task_updates: list[tuple[str, dict[str, Any]]] = [] pending_resources: list[tuple[str, Any]] = [] seen_tools: set[str] = set() @@ -2356,8 +2379,7 @@ def admin_set_config(): fields = {field: normalize_resource_value(field, entry[field]) for field in _tt_fields if field in entry} if "enabled" in fields and fields["enabled"] is None: raise ResourceValidationError("enabled cannot be empty") - task_type = type_map.get(tool) - if task_type is not None and not task_type.gpus and fields.get("slurm_gres"): + if not profile_gpu.get(tool, False) and fields.get("slurm_gres"): raise ResourceValidationError(f"CPU-only task {tool!r} cannot request GPU GRES") if fields: pending_task_updates.append((tool, fields)) diff --git a/server/revocompute/static/js/configuration.js b/server/revocompute/static/js/configuration.js index a1e5c139..d89b1d1d 100644 --- a/server/revocompute/static/js/configuration.js +++ b/server/revocompute/static/js/configuration.js @@ -155,8 +155,9 @@ return; } - var enabledCount = taskTypeConfigs.filter(function (c) { return c.enabled !== false; }).length; - taskTypeStatus.textContent = taskTypeConfigs.length + " type(s) · " + enabledCount + " enabled"; + var taskConfigs = taskTypeConfigs.filter(function (c) { return !c.is_workflow_stage; }); + var enabledCount = taskConfigs.filter(function (c) { return c.enabled !== false; }).length; + taskTypeStatus.textContent = taskConfigs.length + " type(s) · " + enabledCount + " enabled"; taskTypeCards.innerHTML = taskTypeConfigs.map(function (config) { var meta = findTypeMeta(config.tool); @@ -232,6 +233,15 @@ '
' + slurmFieldsHtml + '
'; } + var enableControl = config.is_workflow_stage ? 'Workflow stage' : + '' + + (enabled ? "Enabled" : "Disabled") + + '' + + ''; + return ( '
' + '
' + @@ -241,13 +251,7 @@ '' + '
' + '
' + - '' + - (enabled ? "Enabled" : "Disabled") + - '' + - '' + + enableControl + '' + diff --git a/server/revocompute/static/js/input-workspace.js b/server/revocompute/static/js/input-workspace.js index 50554d84..3aac78e4 100644 --- a/server/revocompute/static/js/input-workspace.js +++ b/server/revocompute/static/js/input-workspace.js @@ -107,6 +107,8 @@ control.appendChild(checkboxB); } else if (parameter.choices && parameter.choices.length) { control = element("select", "text-input"); + control.id = "param_" + parameter.name; + control.dataset.paramName = parameter.name; parameter.choices.forEach(function (choice) { var option = element("option", "", String(choice)); option.value = choice; option.selected = choice === parameter.default; control.appendChild(option); diff --git a/server/revocompute/task_runtime.py b/server/revocompute/task_runtime.py index ca99babc..1bec880f 100644 --- a/server/revocompute/task_runtime.py +++ b/server/revocompute/task_runtime.py @@ -17,12 +17,15 @@ import mimetypes import os import re +import signal import shutil import subprocess import threading import time import zipfile +from dataclasses import replace from datetime import datetime +from pathlib import Path from typing import Any import docker @@ -288,6 +291,81 @@ def _run_compute_job( return job.poll() +def _workflow_state(task: dict[str, Any]) -> dict[str, dict[str, Any]]: + raw = task.get("workflow_state") + if not raw: + return {} + try: + parsed = json.loads(raw) + except (TypeError, json.JSONDecodeError): + return {} + return parsed if isinstance(parsed, dict) else {} + + +def _run_compute_workflow( + task_id: str, + task: dict[str, Any], + tt, + runner, + entities: list[dict], + output_dir: str, + resource_policies: dict[str, ResolvedResources], + stage_callback, +) -> JobState: + """Run an ordered workflow through the existing one-allocation Job API.""" + state = _workflow_state(task) + for stage in tt.workflow: + previous = state.get(stage.name, {}) + if previous.get("status") == "completed": + continue + policy = resource_policies.get(stage.name) + if policy is None: + raise ResourceValidationError(f"Workflow stage {stage.name!r} has no resource snapshot") + markers = {name: tt.stage_markers[name] for name in stage.stage_markers} + stage_tt = replace( + tt, + name=stage.name.replace(".", "-"), + runner_args=stage.runner_args, + gpus=stage.requires_gpu, + stage_markers=markers, + workflow=(), + ) + first_marker = next(iter(markers)) + if not task_store.update_task(task_id, status="queued", run_stage=first_marker): + return JobState.CANCELLED + job = _create_job( + task_id, + stage_tt, + runner, + entities, + output_dir, + stage_callback, + username=task.get("username", ""), + resource_policy=policy, + ) + jid = job.submit() + state[stage.name] = {"status": "running", "job_id": jid, "started_at": time.time()} + handles = {"workflow_state": json.dumps(state, sort_keys=True)} + if isinstance(job, SlurmJob): + handles["slurm_job_id"] = jid + elif isinstance(job, DockerJob): + handles["container_id"] = jid + if not task_store.update_task(task_id, **handles): + job.cancel() + return JobState.CANCELLED + result = job.poll() + state[stage.name].update(status=result.value, finished_at=time.time()) + task_store.update_task( + task_id, + workflow_state=json.dumps(state, sort_keys=True), + slurm_job_id=None, + container_id=None, + ) + if result != JobState.COMPLETED: + return result + return JobState.COMPLETED + + # --------------------------------------------------------------------------- # Result finalization and optional archive cache # --------------------------------------------------------------------------- @@ -683,12 +761,17 @@ def _execute_compute_task(md5sum: str, task_type: str = "gremlin", params: dict raw_form = task.get("input_form") entities: list[dict] = [] resource_policy: ResolvedResources | None = None + resource_policies: dict[str, ResolvedResources] = {} if raw_form: try: parsed = json.loads(raw_form) if isinstance(raw_form, str) else raw_form entities = parsed.get("entities", []) if parsed.get("resource_policy"): resource_policy = ResolvedResources.from_snapshot(parsed["resource_policy"]) + raw_policies = parsed.get("resource_policies", {}) + if not isinstance(raw_policies, dict): + raise TypeError("resource_policies must be an object") + resource_policies = {name: ResolvedResources.from_snapshot(policy) for name, policy in raw_policies.items()} except (json.JSONDecodeError, TypeError, ResourceValidationError): logging.warning("Task %s: input_form or resource policy is invalid.", md5sum) _record_failure(md5sum, task, time.time(), "", "Task input or resource policy is invalid") @@ -764,7 +847,19 @@ def _on_stage_change(stage: str) -> None: } if resource_policy is not None: job_kwargs["resource_policy"] = resource_policy - final_state = _run_compute_job(**job_kwargs) + if tt.workflow: + final_state = _run_compute_workflow( + md5sum, + task, + tt, + runner, + entities, + output_dir, + resource_policies, + _on_stage_change, + ) + else: + final_state = _run_compute_job(**job_kwargs) if _task_is_terminal(md5sum): logging.info("Task %s was deleted during execution; skipping result packing and finalization.", md5sum) return @@ -823,6 +918,78 @@ def _on_stage_change(stage: str) -> None: # --------------------------------------------------------------------------- +def _stop_orphaned_workflow_execution(task_id: str, slurm_job_id: str, container_id: str) -> str: + """Stop a workflow allocation before another worker resumes its stage.""" + if slurm_job_id: + if slurm_job_id.isdigit(): + scancel = shutil.which("scancel") + if not scancel: + return f"Cannot resume while SLURM job {slurm_job_id} cannot be cancelled" + try: + subprocess.run([scancel, slurm_job_id], timeout=10, check=True) + except (OSError, subprocess.SubprocessError) as exc: + return f"Could not cancel SLURM job {slurm_job_id}: {exc}" + else: + match = re.fullmatch(r"srun-([1-9][0-9]*)", slurm_job_id) + if not match: + return f"Cannot resume workflow with unknown SLURM handle {slurm_job_id!r}" + pid = int(match.group(1)) + try: + command = Path(f"/proc/{pid}/cmdline").read_bytes().replace(b"\0", b" ") + except FileNotFoundError: + command = b"" + except OSError as exc: + return f"Could not inspect srun process {pid}: {exc}" + if command: + if b"srun" not in command or task_id[:8].encode("ascii") not in command: + return f"Refusing to stop unverified process {pid} for workflow recovery" + try: + os.kill(pid, signal.SIGTERM) + except ProcessLookupError: + pass + except OSError as exc: + return f"Could not stop srun process {pid}: {exc}" + if not _wait_for_process_exit(pid, 10.0): + try: + os.kill(pid, signal.SIGKILL) + except ProcessLookupError: + pass + except OSError as exc: + return f"Could not kill srun process {pid}: {exc}" + if not _wait_for_process_exit(pid, 2.0): + return f"srun process {pid} did not exit after SIGKILL" + + if container_id: + client = None + try: + client = docker.from_env() + client.containers.get(container_id).stop(timeout=10) + except docker.errors.NotFound: + pass + except docker.errors.DockerException as exc: + return f"Could not stop Docker container {container_id}: {exc}" + finally: + if client is not None: + client.close() + return "" + + +def _wait_for_process_exit(pid: int, timeout: float) -> bool: + """Wait for a PID to disappear or become a reaped-ready zombie.""" + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + try: + state = Path(f"/proc/{pid}/stat").read_text(encoding="utf-8").split()[2] + except FileNotFoundError: + return True + except (OSError, IndexError): + state = "" + if state == "Z": + return True + time.sleep(0.1) + return False + + def _recover_orphaned_tasks() -> int: """Resolve compute records whose owning Celery worker disappeared.""" handled = 0 @@ -832,6 +999,44 @@ def _recover_orphaned_tasks() -> int: md5sum = task["md5sum"] slurm_job_id = str(task.get("slurm_job_id") or "").strip() container_id = str(task.get("container_id") or "") + task_type = task.get("task_type", "gremlin") + try: + workflow_task = bool(_get_task_type(task_type)[0].workflow) + except KeyError: + workflow_task = False + if workflow_task: + if not task_store.claim_task_recovery(md5sum, expected_status=str(task.get("status") or "")): + continue + stop_error = _stop_orphaned_workflow_execution(md5sum, slurm_job_id, container_id) + if stop_error: + task_store.update_task(md5sum, status="queued", error=stop_error) + logging.error("Recovery left workflow %s queued: %s", md5sum, stop_error) + handled += 1 + continue + state = _workflow_state(task) + for step in state.values(): + if step.get("status") == "running": + step["status"] = "interrupted" + if not task_store.update_task( + md5sum, + status="pending", + slurm_job_id=None, + container_id=None, + workflow_state=json.dumps(state, sort_keys=True), + error=None, + ): + handled += 1 + continue + try: + resumed = run_compute_task.apply_async(args=[md5sum], kwargs={"task_type": task_type}) + except Exception as exc: # pylint: disable=broad-except + task_store.update_task(md5sum, status="queued", error=f"Workflow recovery enqueue failed: {exc}") + logging.exception("Recovery could not enqueue workflow %s", md5sum) + else: + if not task_store.update_task(md5sum, celery_task_id=resumed.id): + resumed.revoke(terminate=True) + handled += 1 + continue if slurm_job_id: cancellation_error = "" scancel = shutil.which("scancel") @@ -864,7 +1069,6 @@ def _recover_orphaned_tasks() -> int: continue if not container_id: continue - task_type = task.get("task_type", "gremlin") logging.info("Recovery: checking orphaned task %s", md5sum) try: from revocompute.job.runners.docker_runner import DockerJob diff --git a/server/revocompute/task_types/__init__.py b/server/revocompute/task_types/__init__.py index 20ed1f83..7983dfa3 100644 --- a/server/revocompute/task_types/__init__.py +++ b/server/revocompute/task_types/__init__.py @@ -80,6 +80,17 @@ class RuntimeFamily: slurm_image: str = "" +@dataclass(frozen=True) +class WorkflowStage: + """One scheduler allocation in an ordered task workflow.""" + + name: str + display_name: str + requires_gpu: bool + runner_args: tuple[str, ...] = () + stage_markers: tuple[str, ...] = () + + @dataclass(frozen=True) class TaskType: """Portable user-facing task definition. @@ -105,6 +116,7 @@ class TaskType: runner_args: tuple[str, ...] = () gpus: bool = False stage_markers: dict[str, str] = field(default_factory=dict) + workflow: tuple[WorkflowStage, ...] = () params: tuple[TaskParam, ...] = () input_workspace: tuple[InputCapability, ...] = () result_workspace: tuple[ResultView, ...] = () @@ -334,6 +346,50 @@ def _load_result_workspace(raw: Any) -> tuple[ResultView, ...]: return tuple(views) +def _load_workflow(raw: Any, task_name: str, stage_markers: dict[str, str]) -> tuple[WorkflowStage, ...]: + if raw is None: + return () + if not isinstance(raw, list) or len(raw) < 2: + raise ValueError(f"Task type {task_name!r} workflow must contain at least two stages") + stages: list[WorkflowStage] = [] + seen: set[str] = set() + for entry in raw: + if not isinstance(entry, dict) or set(entry) - { + "name", + "display_name", + "requires_gpu", + "runner_args", + "stage_markers", + }: + raise ValueError(f"Task type {task_name!r} has an invalid workflow stage") + name = entry.get("name") + requires_gpu = entry.get("requires_gpu", False) + runner_args = entry.get("runner_args", ()) + raw_markers = entry.get("stage_markers", ()) + if not isinstance(name, str) or not name.replace("_", "").isalnum() or name in seen: + raise ValueError(f"Task type {task_name!r} has an invalid or duplicate workflow stage name") + if not isinstance(requires_gpu, bool): + raise ValueError(f"Workflow stage {task_name}.{name} requires_gpu must be a boolean") + if not isinstance(runner_args, list) or not all(isinstance(arg, str) for arg in runner_args): + raise ValueError(f"Workflow stage {task_name}.{name} runner_args must be a list of strings") + if not isinstance(raw_markers, list) or not all(isinstance(marker, str) for marker in raw_markers): + raise ValueError(f"Workflow stage {task_name}.{name} stage_markers must be a list of strings") + markers = tuple(raw_markers) + if not markers or not set(markers).issubset(stage_markers): + raise ValueError(f"Workflow stage {task_name}.{name} must reference declared stage markers") + seen.add(name) + stages.append( + WorkflowStage( + name=f"{task_name}.{name}", + display_name=str(entry.get("display_name") or name.replace("_", " ").title()), + requires_gpu=requires_gpu, + runner_args=tuple(runner_args), + stage_markers=markers, + ) + ) + return tuple(stages) + + # --------------------------------------------------------------------------- # YAML loader # --------------------------------------------------------------------------- @@ -439,6 +495,7 @@ def load_registry(task_types_yaml: str, runners_dir: str, enabled: set[str]) -> raise ValueError(f"Task type {name!r} min_input_files must be between zero and max_input_files") params = tuple(TaskParam(**{**p, "choices": tuple(p.get("choices", []))}) for p in entry.get("params", [])) + stage_markers = entry.get("stage_markers", {}) tt = TaskType( name=name, display_name=entry["display_name"], @@ -454,7 +511,8 @@ def load_registry(task_types_yaml: str, runners_dir: str, enabled: set[str]) -> allow_multiple_inputs=allow_multiple_inputs, max_input_files=max_input_files, min_input_files=min_input_files, - stage_markers=entry.get("stage_markers", {}), + stage_markers=stage_markers, + workflow=_load_workflow(entry.get("workflow"), name, stage_markers), params=params, input_workspace=_load_input_workspace(entry.get("input_workspace")), result_workspace=_load_result_workspace(entry.get("result_workspace")), diff --git a/server/run/revocompute_ctl/registry.py b/server/run/revocompute_ctl/registry.py index 225241fe..1fc25354 100644 --- a/server/run/revocompute_ctl/registry.py +++ b/server/run/revocompute_ctl/registry.py @@ -13,6 +13,7 @@ import os import re import sys +import tempfile from dataclasses import dataclass from pathlib import Path @@ -201,7 +202,8 @@ def validate_slurm_images(state, families: list[RuntimeFamily]) -> None: if not Path(target).is_file(): print(f"[SLURM] Missing SIF image: {family.slurm_image}", file=sys.stderr) print( - f" Build it: apptainer build --fakeroot {family.slurm_image} {state.server_root()}/{family.definition}", + " Build it: apptainer build --fakeroot " + f"{family.slurm_image} {state.server_root()}/{family.definition}", file=sys.stderr, ) missing += 1 @@ -215,6 +217,30 @@ def validate_slurm_images(state, families: list[RuntimeFamily]) -> None: raise RegistryError +def _docker_tag(image: str, suffix: str = "latest") -> str: + repository = image.rsplit("/", 1)[-1] + return f"{image}:{suffix}" if ":" not in repository and "@" not in image else image + + +def _docker_image_id(state, tag: str) -> str: + return run_cmd( + ["docker", "image", "inspect", "--format", "{{.Id}}", tag], + env=state.exported(), + check=False, + capture=True, + ).stdout.strip() + + +def _sif_source_tag(state, family: RuntimeFamily) -> str: + """Use the prepared runner image when this restart built one.""" + latest = _docker_tag(family.docker_image) + if latest.endswith(":latest"): + prepared = f"{latest[:-len(':latest')]}:next" + if _docker_image_id(state, prepared): + return prepared + return latest + + def sif_stale(state, family: RuntimeFamily) -> bool: """True when the deployed SIF needs a rebuild: it is missing, or the family's docker image was created after the SIF (covers image updates @@ -222,8 +248,10 @@ def sif_stale(state, family: RuntimeFamily) -> bool: unchanged while the SIF still predates it).""" if not Path(family.slurm_image).is_file(): return True - image = family.docker_image - tag = f"{image}:latest" if ":" not in image and "@" not in image else image + latest = _docker_tag(family.docker_image) + tag = _sif_source_tag(state, family) + if tag != latest and _docker_image_id(state, tag) != _docker_image_id(state, latest): + return True created = run_cmd( ["docker", "image", "inspect", "--format", "{{.Created}}", tag], env=state.exported(), @@ -239,6 +267,21 @@ def sif_stale(state, family: RuntimeFamily) -> bool: return image_ts > os.path.getmtime(family.slurm_image) +def _sif_definition_for_tag(def_file: Path, source_tag: str) -> tuple[str, str | None]: + """Return a definition using the prepared Docker tag, plus a temp path to clean up.""" + text = def_file.read_text(encoding="utf-8") + current = _first_directive_value(text, "From:") + if source_tag == current: + return str(def_file), None + updated = text.replace(f"From: {current}", f"From: {source_tag}", 1) + handle = tempfile.NamedTemporaryFile("w", encoding="utf-8", suffix=".def", delete=False) + try: + handle.write(updated) + finally: + handle.close() + return handle.name, handle.name + + def build_slurm_images(state, families: list[RuntimeFamily]) -> int: """Stage SIFs as ``.next`` for missing or stale families only; promotion (promotion.py) moves them into place after down. Returns the @@ -269,9 +312,18 @@ def build_slurm_images(state, families: list[RuntimeFamily]) -> int: # Atomic staging: a killed build must never leave a corrupt .next # that the next run treats as a valid staging. staging = f"{staged}.build" - result = run_cmd( - ["apptainer", "build", "--fakeroot", staging, str(def_file)], env=state.exported(), check=False + build_definition, temporary_definition = _sif_definition_for_tag( + def_file, _sif_source_tag(state, family) ) + try: + result = run_cmd( + ["apptainer", "build", "--fakeroot", staging, build_definition], + env=state.exported(), + check=False, + ) + finally: + if temporary_definition is not None: + os.remove(temporary_definition) if result.returncode != 0: if os.path.isfile(staging): os.remove(staging) diff --git a/server/run/revocompute_ctl/sweep.py b/server/run/revocompute_ctl/sweep.py index 2d901667..5587ad40 100644 --- a/server/run/revocompute_ctl/sweep.py +++ b/server/run/revocompute_ctl/sweep.py @@ -23,10 +23,32 @@ print(job_id) """ -SWEEP_SOURCE = """import time -from revocompute.task_runtime import _record_failure, task_store +SWEEP_SOURCE = """import json +import time +from revocompute.task_runtime import _get_task_type, _record_failure, task_store for task in task_store.list_tasks(): if task.get("status") in {"queued", "running"}: + task_type = task.get("task_type", "gremlin") + try: + is_workflow = bool(getattr(_get_task_type(task_type)[0], "workflow", ())) + except KeyError: + is_workflow = False + if is_workflow: + try: + state = json.loads(task.get("workflow_state") or "{}") + except (TypeError, json.JSONDecodeError): + state = {} + for step in state.values(): + if step.get("status") == "running": + step["status"] = "interrupted" + task_store.update_task( + task["md5sum"], + status="queued", + slurm_job_id=None, + container_id=None, + workflow_state=json.dumps(state, sort_keys=True), + ) + continue _record_failure( task["md5sum"], task, @@ -77,7 +99,7 @@ def pre_stop_sweep_slurm(state, compose_cmd: tuple[str, ...]) -> None: env=state.exported(), check=False, ) - print("Marking in-flight tasks failed before stopping the stack...") + print("Preserving workflows and finalizing other in-flight tasks before stopping the stack...") marked = run_cmd( [*compose_cmd, *compose_args(state), "--env-file", state.env_file, "exec", "-T", "worker", "python3", "-"], env=state.exported(), diff --git a/server/tests/test_race_conditions.py b/server/tests/test_race_conditions.py index da46eb32..7967d418 100644 --- a/server/tests/test_race_conditions.py +++ b/server/tests/test_race_conditions.py @@ -55,6 +55,32 @@ def test_race_cancel_finished_task_rejected(monkeypatch, tmp_path): assert "not pending or running" in resp.json["error"] +def test_race_cancel_loses_atomic_claim_to_completion(monkeypatch, tmp_path): + module = _load_pssm_module(monkeypatch, tmp_path, extra_env={"RUNNER_UID": "1234", "RUNNER_GID": "5678"}) + from revocompute import routes + + client = module.app.test_client() + auth_header = _test_client_auth(module) + md5sum = _insert_pending_task(module, tmp_path / "result") + + def finish_before_claim(task_id, **fields): + assert task_id == md5sum + module.task_store.update_task(md5sum, status="finished") + return False + + monkeypatch.setattr(module.task_store, "claim_task_cancellation", finish_before_claim) + monkeypatch.setattr( + routes.cancel_compute_resources, + "delay", + lambda **kwargs: pytest.fail("completed task resources must not be cancelled"), + ) + + response = client.post(f"/compute/api/cancel/{md5sum}", headers=auth_header) + + assert response.status_code == 409 + assert module.task_store.get_task(md5sum)["status"] == "finished" + + def test_race_cancel_already_cancelled_task(monkeypatch, tmp_path): """Re-cancelling a cancelled task returns 400.""" module = _load_pssm_module(monkeypatch, tmp_path, extra_env={"RUNNER_UID": "1234", "RUNNER_GID": "5678"}) diff --git a/server/tests/test_restart_ctl.py b/server/tests/test_restart_ctl.py index 77530bef..008aebf9 100644 --- a/server/tests/test_restart_ctl.py +++ b/server/tests/test_restart_ctl.py @@ -23,6 +23,7 @@ from pathlib import Path import pytest +import yaml SERVER_DIR = Path(__file__).resolve().parents[1] RUN_DIR = SERVER_DIR / "run" @@ -35,7 +36,7 @@ from revocompute_ctl import stamp as stamp_mod # noqa: E402 from revocompute_ctl import sweep as sweep_mod # noqa: E402 from revocompute_ctl.env import EnvState, parse_env_file # noqa: E402 -from revocompute_ctl.registry import RuntimeFamily, build_slurm_images # noqa: E402 +from revocompute_ctl.registry import RuntimeFamily, _docker_tag, build_slurm_images # noqa: E402 from revocompute_ctl.steps import Step, StepRegistry, run_walk # noqa: E402 RUNNER_IMAGE = "revodesign-revocompute-runner" @@ -323,7 +324,6 @@ def test_sif_staging_builds_missing_skips_unchanged(tmp_path, monkeypatch): state, _log = _shimmed_state(monkeypatch, tmp_path, bin_dir, {}) build_slurm_images(state, [family]) # missing SIF → stage .next assert (sif_dir / "gremlin.sif.next").is_file() - (sif_dir / "gremlin.sif.next").unlink() (sif_dir / "gremlin.sif").touch() build_slurm_images(state, [family]) # image older than SIF → skip @@ -334,6 +334,39 @@ def test_sif_staging_builds_missing_skips_unchanged(tmp_path, monkeypatch): assert (sif_dir / "gremlin.sif.next").is_file() +def test_docker_tag_distinguishes_registry_port_from_image_tag(): + assert _docker_tag("registry.example:5000/team/runner") == "registry.example:5000/team/runner:latest" + assert _docker_tag("registry.example:5000/team/runner:v2") == "registry.example:5000/team/runner:v2" + assert _docker_tag("registry.example:5000/team/runner@sha256:abc") == "registry.example:5000/team/runner@sha256:abc" + + +def test_sif_staging_builds_changed_prepared_image_from_next(tmp_path, monkeypatch): + bin_dir = _write_shims(tmp_path) + sif_dir = tmp_path / "sifs" + sif_dir.mkdir() + (sif_dir / "gremlin.sif").touch() + family = RuntimeFamily( + "gremlin", + RUNNER_IMAGE, + "docker/runners/pssm_gremlin/Dockerfile", + "docker/runners/pssm_gremlin/gremlin.def", + str(sif_dir / "gremlin.sif"), + ) + state, log = _shimmed_state( + monkeypatch, + tmp_path, + bin_dir, + {f"{RUNNER_IMAGE}:latest": "sha256:old", f"{RUNNER_IMAGE}:next": "sha256:new"}, + ) + + build_slurm_images(state, [family]) + + assert (sif_dir / "gremlin.sif.next").is_file() + apptainer_line = next(line for line in log.read_text().splitlines() if line.startswith("build ")) + definition = Path(apptainer_line.rsplit(" ", 1)[-1]) + assert not definition.exists() # temporary prepared-tag definition is cleaned up + + def test_sif_staging_drops_failed_runner_from_enabled_list(tmp_path, monkeypatch): bin_dir = _write_shims(tmp_path) monkeypatch.setenv("APPTAINER_FAIL", "1") @@ -469,6 +502,14 @@ def test_rollback_restores_previous_set_and_config(tmp_path, monkeypatch): bin_dir = _write_shims(tmp_path) config_dir = tmp_path / "config" shutil.copytree(Path(REPO_DIR) / "server" / "config", config_dir) + registry_file = config_dir / "task_types.yaml" + registry = yaml.safe_load(registry_file.read_text(encoding="utf-8")) + sif = tmp_path / "sifs" / "gremlin.sif" + sif.parent.mkdir() + sif.write_text("current", encoding="utf-8") + Path(f"{sif}.previous").write_text("previous", encoding="utf-8") + registry["runtime_families"]["gremlin"]["slurm_image"] = str(sif) + registry_file.write_text(yaml.safe_dump(registry, sort_keys=False), encoding="utf-8") task_dir, _auth_dir, env_file = _deploy_env(tmp_path, config_dir) ids = {f"{RUNNER_IMAGE}:previous": "sha256:prev", f"{RUNNER_IMAGE}:latest": "sha256:bad"} @@ -487,7 +528,6 @@ def test_rollback_restores_previous_set_and_config(tmp_path, monkeypatch): }, ) # Drift the registry in a validation-neutral way. - registry_file = config_dir / "task_types.yaml" registry_file.write_text(registry_file.read_text(encoding="utf-8") + "\n# drifted\n", encoding="utf-8") result = _run_cli(monkeypatch, tmp_path, env_file, bin_dir, "restart", "--rollback") @@ -496,6 +536,7 @@ def test_rollback_restores_previous_set_and_config(tmp_path, monkeypatch): assert not registry_file.read_text(encoding="utf-8").endswith("# drifted\n") # config restored commands = (tmp_path / "docker.log").read_text(encoding="utf-8").splitlines() assert f"tag {RUNNER_IMAGE}:previous {RUNNER_IMAGE}:latest" in commands + assert sif.read_text(encoding="utf-8") == "previous" assert "All prepared deployment services are running." in result.stdout rolled = stamp_mod.load_stamp(state) assert rolled["mode"] == "rollback" diff --git a/server/tests/test_runner_script_static.py b/server/tests/test_runner_script_static.py index 7405d291..dfb4125c 100644 --- a/server/tests/test_runner_script_static.py +++ b/server/tests/test_runner_script_static.py @@ -24,7 +24,7 @@ def test_prepared_activation_audits_resources_before_stopping_services(): assert "-m revocompute.resource_audit" in steps_source -def test_slurm_pre_stop_sweep_preserves_pending_and_finalizes_started_tasks(): +def test_slurm_pre_stop_sweep_preserves_pending_and_resumable_workflows(): sweep = (SERVER_ROOT / "run" / "revocompute_ctl" / "sweep.py").read_text(encoding="utf-8") down = (SERVER_ROOT / "run" / "revocompute_ctl" / "steps.py").read_text(encoding="utf-8") @@ -34,8 +34,11 @@ def test_slurm_pre_stop_sweep_preserves_pending_and_finalizes_started_tasks(): assert '"scancel"' in sweep assert 'if task.get("status") in {"queued", "running"}:' in sweep assert '"pending"' not in sweep - assert "from revocompute.task_runtime import _record_failure, task_store" in sweep + assert "_record_failure" in sweep and "task_store" in sweep assert "_record_failure(" in sweep + assert 'status="queued"' in sweep + assert 'json.loads(task.get("workflow_state") or "{}")' in sweep + assert 'getattr(_get_task_type(task_type)[0], "workflow", ())' in sweep assert "pre_stop_sweep_slurm" in down @@ -57,7 +60,7 @@ def test_runner_script_executes_pipeline_commands_as_arrays(): assert '"${cmd[@]}"' in script -def _run_with_manifest(script, input_file, output_dir, env, params=None): +def _run_with_manifest(script, input_file, output_dir, env, params=None, extra_args=()): """Run a runner script under the v2 protocol: write task.json next to the input, point -i at it, and set TASK_MANIFEST + TASK_CONTEXT_SRC.""" manifest_path = input_file.parent / "task.json" @@ -85,7 +88,7 @@ def _run_with_manifest(script, input_file, output_dir, env, params=None): str(SERVER_ROOT / "docker" / "runners" / "common" / "task_context.sh"), ) return subprocess.run( - ["bash", str(script), "-i", str(manifest_path), "-o", str(output_dir)], + ["bash", str(script), *extra_args, "-i", str(manifest_path), "-o", str(output_dir)], env=env, check=False, capture_output=True, @@ -143,6 +146,75 @@ def test_alphafold_runner_drains_final_stage_before_exit(tmp_path): assert (output_dir / "task_finished").is_file() +def test_alphafold_feature_stage_stops_before_modeling(tmp_path): + input_file = tmp_path / "input.fasta" + output_dir = tmp_path / "outputs" + alphafold_root = tmp_path / "alphafold" + fake_context = tmp_path / "task_context.sh" + fake_python = tmp_path / "fake-python" + fake_args = tmp_path / "alphafold.args" + input_file.write_text(">first\nAAAA\n>second\nBBBB\n", encoding="utf-8") + output_dir.mkdir() + alphafold_root.mkdir() + fake_context.write_text( + '_parse_param() { [[ "$1" == model_preset ]] && printf "multimer\\n" || printf "%s\\n" "$2"; }\n' + 'primary_input() { printf "%s\\n" "$FAKE_PRIMARY_INPUT"; }\n', + encoding="utf-8", + ) + fake_python.write_text( + "#!/bin/bash\n" + "set -e\n" + 'printf "%s\\n" "$@" > "$FAKE_ARGS_FILE"\n' + 'for arg in "$@"; do\n' + ' case "$arg" in --output_dir=*) output_dir=${arg#*=} ;; --run_stage=*) stage=${arg#*=} ;; esac\n' + "done\n" + '[[ "$stage" == features ]]\n' + 'mkdir -p "$output_dir/input"\n' + 'printf "FEATURES\\n" > "$output_dir/input/features.pkl"\n', + encoding="utf-8", + ) + fake_python.chmod(0o755) + env = os.environ.copy() + env.update( + { + "ALPHAFOLD_PATH": str(alphafold_root), + "ALPHAFOLD_PYTHON": str(fake_python), + "FAKE_ARGS_FILE": str(fake_args), + "FAKE_PRIMARY_INPUT": str(input_file), + "TASK_CONTEXT_SRC": str(fake_context), + "ALPHAFOLD_STAGE_TRANSLATOR": str(SERVER_ROOT / "docker" / "runners" / "common" / "stage_translate.py"), + "ALPHAFOLD_STAGE_PATTERNS": str(SERVER_ROOT / "docker" / "runners" / "alphafold" / "alphafold.stages"), + "TMPDIR": str(tmp_path), + } + ) + + completed = _run_with_manifest( + ALPHAFOLD_RUNNER_SCRIPT, + input_file, + output_dir, + env, + params={"model_preset": "multimer"}, + extra_args=("-s", "features"), + ) + + assert completed.returncode == 0, completed.stderr + args = fake_args.read_text(encoding="utf-8") + assert "--model_preset=multimer" in args + assert "--uniref30_database_path=" in args + assert "--uniprot_database_path=" in args + assert "--pdb70_database_path=" not in args + assert (output_dir / ".alphafold-features-complete").is_file() + assert not (output_dir / "task_finished").exists() + + +def test_alphafold_image_applies_staged_pipeline_to_pinned_source(): + dockerfile = (SERVER_ROOT / "docker" / "runners" / "alphafold" / "Dockerfile").read_text() + patch = (SERVER_ROOT / "docker" / "runners" / "alphafold" / "staged_pipeline.patch").read_text() + assert "git -C /opt/alphafold apply --check /tmp/staged_pipeline.patch" in dockerfile + assert "FLAGS.run_stage == 'model'" in patch + assert "FLAGS.run_stage == 'features'" in patch + + def test_opendde_runner_uses_writable_snapshot_copy_and_checks_outputs(): script = OPENDDE_RUNNER_SCRIPT.read_text() diff --git a/server/tests/test_task_type_registry.py b/server/tests/test_task_type_registry.py index 0ce1e11e..889edfee 100644 --- a/server/tests/test_task_type_registry.py +++ b/server/tests/test_task_type_registry.py @@ -91,6 +91,23 @@ def test_shared_tasks_resolve_one_runtime_and_runner_config(): assert rfdiffusion.runner_args == ("rfdiffusion",) +def test_alphafold_declares_cpu_then_gpu_workflow(): + with _preserve_registry(): + task_types.load_registry( + str(SERVER_ROOT / "config" / "task_types.yaml"), + str(SERVER_ROOT / "config" / "runners"), + {"alphafold"}, + ) + alphafold, _ = task_types.get("alphafold") + + assert [(stage.name, stage.requires_gpu) for stage in alphafold.workflow] == [ + ("alphafold.features", False), + ("alphafold.model", True), + ] + assert alphafold.workflow[0].runner_args == ("-s", "features") + assert alphafold.workflow[1].runner_args == ("-s", "model") + + def test_input_workspace_capabilities_cover_simple_and_complex_tasks(): enabled = {"placer-rfdiffusion", "easifa"} with _preserve_registry(): diff --git a/server/tests/test_tasks.py b/server/tests/test_tasks.py index 70c290c4..ccee8b65 100644 --- a/server/tests/test_tasks.py +++ b/server/tests/test_tasks.py @@ -208,6 +208,7 @@ def test_create_task_uses_capability_plugins_with_safe_fallbacks(): assert "workspace.validate()" in orchestrator assert 'formData.append("input_paths"' in orchestrator assert 'formData.append("params[" + name + "]"' in orchestrator + assert 'control.id = "param_" + parameter.name;' in workspace def test_full_stack_smoke_uses_manifest_first_result_contract(): @@ -265,6 +266,43 @@ class _DummyAsyncResult: assert manifest["files"][0]["relative_path"] == "2KL8.fasta" +def test_alphafold_multimer_submission_preserves_selected_preset(monkeypatch, tmp_path): + module = _load_pssm_module( + monkeypatch, + tmp_path, + extra_env={"RUNNER_UID": "1234", "RUNNER_GID": "5678", "ENABLED_TASKRUNNERS": "alphafold"}, + ) + client = module.app.test_client() + auth_header = _test_client_auth(module) + user = module.app.config["user_db"].get_user_by_username("tester") + module.app.config["user_db"].update_user(user["id"], allow_gpu_use=True) + + class _Queued: + id = "queued-alphafold-multimer" + + monkeypatch.setattr(module.run_compute_task, "apply_async", lambda *args, **kwargs: _Queued()) + fasta_path = Path(__file__).resolve().parents[2] / "tests/data/fasta/Sli_S4.fasta" + with fasta_path.open("rb") as handle: + response = client.post( + "/compute/api/post", + headers=auth_header, + data={"task_type": "alphafold", "params[model_preset]": "multimer", "file": (handle, fasta_path.name)}, + content_type="multipart/form-data", + ) + + assert response.status_code == 302, response.get_data(as_text=True)[:300] + task_id = response.headers["Location"].rsplit("/", 1)[-1] + task = module.task_store.get_task(task_id) + manifest = json.loads( + (Path(module.app.config["WORKSPACE_FOLDER"]) / "tester" / task_id / "inputs" / "task.json").read_text() + ) + input_form = json.loads(task["input_form"]) + assert manifest["params"]["model_preset"] == "multimer" + assert set(input_form["resource_policies"]) == {"alphafold.features", "alphafold.model"} + assert input_form["resource_policies"]["alphafold.features"]["requires_gpu"] is False + assert input_form["resource_policies"]["alphafold.model"]["requires_gpu"] is True + + def test_create_task_page_has_categorized_rail_and_validation_panel(): template = (SERVER_PACKAGE / "templates" / "create_task.html").read_text(encoding="utf-8") for marker in ( @@ -945,9 +983,9 @@ def test_task_store_update_ignores_late_non_deleted_updates(monkeypatch, tmp_pat ) # Simulate stale worker writes arriving after a delete request. - module.task_store.update_task(md5sum, status="running", run_stage="blast") - module.task_store.update_task(md5sum, status="finished", walltime=12.3, error=None) - module.task_store.update_task(md5sum, run_stage="hhblits") + assert module.task_store.update_task(md5sum, status="running", run_stage="blast") is False + assert module.task_store.update_task(md5sum, status="finished", walltime=12.3, error=None) is False + assert module.task_store.update_task(md5sum, run_stage="hhblits") is False task = module.task_store.get_task(md5sum) assert task is not None @@ -956,6 +994,33 @@ def test_task_store_update_ignores_late_non_deleted_updates(monkeypatch, tmp_pat assert task["finished_at"] == deleted_at +def test_task_store_recovery_claim_is_atomic(monkeypatch, tmp_path): + module = _load_pssm_module( + monkeypatch, + tmp_path, + extra_env={"RUNNER_UID": "1234", "RUNNER_GID": "5678"}, + ) + md5sum = _insert_pending_task(module, tmp_path / "result") + module.task_store.update_task(md5sum, status="running") + + assert module.task_store.claim_task_recovery(md5sum, expected_status="running") is True + assert module.task_store.claim_task_recovery(md5sum, expected_status="running") is False + assert module.task_store.get_task(md5sum)["status"] == "pending" + + +def test_task_store_cancellation_claim_is_active_only(monkeypatch, tmp_path): + module = _load_pssm_module( + monkeypatch, + tmp_path, + extra_env={"RUNNER_UID": "1234", "RUNNER_GID": "5678"}, + ) + md5sum = _insert_pending_task(module, tmp_path / "result") + + assert module.task_store.claim_task_cancellation(md5sum, error="cancelled") is True + assert module.task_store.claim_task_cancellation(md5sum, error="again") is False + assert module.task_store.get_task(md5sum)["error"] == "cancelled" + + def test_run_compute_task_does_not_resurrect_deleted_task(monkeypatch, tmp_path): module = _load_pssm_module( monkeypatch, diff --git a/server/tests/test_workflow_composer.py b/server/tests/test_workflow_composer.py new file mode 100644 index 00000000..fb6e2c8a --- /dev/null +++ b/server/tests/test_workflow_composer.py @@ -0,0 +1,267 @@ +# Copyright (c) 2026 The REvoDesign Developers. +# Distributed under the terms of the GNU General Public License v3.0. +# SPDX-License-Identifier: GPL-3.0-only + +from __future__ import annotations + +import json +import signal +from pathlib import Path + +import pytest +from revocompute.job import JobState +from revocompute.resource_policy import ResolvedResources +from revocompute.task_types import RuntimeFamily, RunnerConfig, TaskType, WorkflowStage + + +def _policy(requires_gpu: bool) -> ResolvedResources: + return ResolvedResources( + cpus=8, + memory="16G", + max_runtime_seconds=3600, + partition="gpu" if requires_gpu else "cpu", + gres="gpu:1" if requires_gpu else None, + nodes=1, + ntasks=1, + qos=None, + account=None, + constraint=None, + exclusive=False, + requires_gpu=requires_gpu, + sources={}, + ) + + +def test_composer_resumes_after_completed_feature_stage(monkeypatch): + server_root = Path(__file__).resolve().parents[1] + monkeypatch.setenv("SERVER_DIR", str(server_root)) + monkeypatch.setenv("CONFIG_DIR", str(server_root / "config")) + monkeypatch.setenv("ENABLED_TASKRUNNERS", "alphafold") + from revocompute import task_runtime + + runtime = RuntimeFamily("alphafold", "image", ("bash", "run.sh"), "Dockerfile", "runner.def", "image.sif") + stages = ( + WorkflowStage("alphafold.features", "Features", False, ("-s", "features"), ("msa",)), + WorkflowStage("alphafold.model", "Model", True, ("-s", "model"), ("model",)), + ) + task_type = TaskType( + "alphafold", + "AlphaFold2", + runtime, + ".fasta", + "FASTA", + gpus=True, + stage_markers={"msa": "MSA", "model": "Model"}, + workflow=stages, + ) + updates = [] + created = [] + + class _Job: + def submit(self): + return "42" + + def poll(self): + return JobState.COMPLETED + + def _create(*args, **kwargs): + created.append((args[1], kwargs["resource_policy"])) + return _Job() + + monkeypatch.setattr(task_runtime, "_create_job", _create) + + def _update(task_id, **fields): + del task_id + updates.append(fields) + return True + + monkeypatch.setattr(task_runtime.task_store, "update_task", _update) + task = { + "username": "tester", + "workflow_state": json.dumps({"alphafold.features": {"status": "completed", "job_id": "41"}}), + } + + result = task_runtime._run_compute_workflow( + "a" * 32, + task, + task_type, + RunnerConfig(), + [], + "/tmp/results", + {"alphafold.features": _policy(False), "alphafold.model": _policy(True)}, + lambda stage: None, + ) + + assert result == JobState.COMPLETED + assert len(created) == 1 + assert created[0][0].name == "alphafold-model" + assert created[0][0].gpus is True + assert created[0][1].requires_gpu is True + final_state = json.loads(updates[-1]["workflow_state"]) + assert final_state["alphafold.features"]["status"] == "completed" + assert final_state["alphafold.model"]["status"] == "completed" + + +def test_composer_does_not_submit_after_cancellation_claim_fails(monkeypatch): + from revocompute import task_runtime + + runtime = RuntimeFamily("alphafold", "image", ("bash", "run.sh"), "Dockerfile", "runner.def", "image.sif") + stage = WorkflowStage("alphafold.model", "Model", True, ("-s", "model"), ("model",)) + task_type = TaskType( + "alphafold", + "AlphaFold2", + runtime, + ".fasta", + "FASTA", + gpus=True, + stage_markers={"model": "Model"}, + workflow=(stage,), + ) + monkeypatch.setattr(task_runtime.task_store, "update_task", lambda *args, **kwargs: False) + monkeypatch.setattr(task_runtime, "_create_job", lambda *args, **kwargs: pytest.fail("job must not be created")) + + result = task_runtime._run_compute_workflow( + "b" * 32, + {}, + task_type, + RunnerConfig(), + [], + "/tmp/results", + {"alphafold.model": _policy(True)}, + lambda stage_name: None, + ) + + assert result == JobState.CANCELLED + + +def test_composer_cancels_submitted_job_when_handle_cannot_be_persisted(monkeypatch): + from revocompute import task_runtime + + runtime = RuntimeFamily("alphafold", "image", ("bash", "run.sh"), "Dockerfile", "runner.def", "image.sif") + stage = WorkflowStage("alphafold.model", "Model", True, ("-s", "model"), ("model",)) + task_type = TaskType( + "alphafold", + "AlphaFold2", + runtime, + ".fasta", + "FASTA", + gpus=True, + stage_markers={"model": "Model"}, + workflow=(stage,), + ) + updates = iter((True, False)) + cancelled = [] + + class _Job: + def submit(self): + return "42" + + def cancel(self): + cancelled.append(True) + + monkeypatch.setattr(task_runtime.task_store, "update_task", lambda *args, **kwargs: next(updates)) + monkeypatch.setattr(task_runtime, "_create_job", lambda *args, **kwargs: _Job()) + + result = task_runtime._run_compute_workflow( + "c" * 32, + {}, + task_type, + RunnerConfig(), + [], + "/tmp/results", + {"alphafold.model": _policy(True)}, + lambda stage_name: None, + ) + + assert result == JobState.CANCELLED + assert cancelled == [True] + + +def test_workflow_recovery_claims_stops_and_requeues_once(monkeypatch): + from revocompute import task_runtime + + task = { + "md5sum": "d" * 32, + "status": "running", + "task_type": "alphafold", + "slurm_job_id": "1234", + "container_id": None, + "workflow_state": json.dumps({"alphafold.features": {"status": "running"}}), + } + claims = [] + stops = [] + updates = [] + + class _TaskType: + workflow = (object(),) + + class _Queued: + id = "replacement-task" + + monkeypatch.setattr(task_runtime.task_store, "list_tasks", lambda: [task]) + monkeypatch.setattr( + task_runtime.task_store, + "claim_task_recovery", + lambda task_id, expected_status: claims.append((task_id, expected_status)) or True, + ) + monkeypatch.setattr( + task_runtime.task_store, + "update_task", + lambda task_id, **fields: updates.append((task_id, fields)) or True, + ) + monkeypatch.setattr(task_runtime, "_get_task_type", lambda name: (_TaskType(), object())) + monkeypatch.setattr( + task_runtime, + "_stop_orphaned_workflow_execution", + lambda *args: stops.append(args) or "", + ) + monkeypatch.setattr(task_runtime.run_compute_task, "apply_async", lambda *args, **kwargs: _Queued()) + + assert task_runtime._recover_orphaned_tasks() == 1 + assert claims == [("d" * 32, "running")] + assert stops == [("d" * 32, "1234", "")] + assert updates[-1][1]["celery_task_id"] == "replacement-task" + + +def test_workflow_recovery_enqueue_failure_stays_discoverable(monkeypatch): + from revocompute import task_runtime + + task = {"md5sum": "e" * 32, "status": "queued", "task_type": "alphafold"} + updates = [] + + class _TaskType: + workflow = (object(),) + + monkeypatch.setattr(task_runtime.task_store, "list_tasks", lambda: [task]) + monkeypatch.setattr(task_runtime.task_store, "claim_task_recovery", lambda *args, **kwargs: True) + monkeypatch.setattr( + task_runtime.task_store, + "update_task", + lambda task_id, **fields: updates.append(fields) or True, + ) + monkeypatch.setattr(task_runtime, "_get_task_type", lambda name: (_TaskType(), object())) + monkeypatch.setattr(task_runtime, "_stop_orphaned_workflow_execution", lambda *args: "") + monkeypatch.setattr( + task_runtime.run_compute_task, + "apply_async", + lambda *args, **kwargs: (_ for _ in ()).throw(RuntimeError("broker unavailable")), + ) + + assert task_runtime._recover_orphaned_tasks() == 1 + assert updates[-1]["status"] == "queued" + assert "broker unavailable" in updates[-1]["error"] + + +def test_workflow_recovery_escalates_srun_termination(monkeypatch): + from revocompute import task_runtime + + kills = [] + waits = iter((False, True)) + monkeypatch.setattr(Path, "read_bytes", lambda self: b"srun task-1234") + monkeypatch.setattr(task_runtime.os, "kill", lambda pid, sig: kills.append((pid, sig))) + monkeypatch.setattr(task_runtime, "_wait_for_process_exit", lambda pid, timeout: next(waits)) + + error = task_runtime._stop_orphaned_workflow_execution("task-1234", "srun-42", "") + + assert error == "" + assert kills == [(42, signal.SIGTERM), (42, signal.SIGKILL)] diff --git a/tests/data/fasta/Sli_S4.fasta b/tests/data/fasta/Sli_S4.fasta new file mode 100644 index 00000000..8988ef43 --- /dev/null +++ b/tests/data/fasta/Sli_S4.fasta @@ -0,0 +1,5 @@ +>Sli +MDYFLLLPEDCVCDILSFTSPKDVVISSAISRGFNSAAESDVIWVKFLPDDYEDINSRYVSPRIYPSKKELYFSLCDFPVLMDGGKLSFSLDKKTGKKCFMISARELAITWGVDTPWYWEWISHPDSRFSEVAHLKGVSWLDIRGTIGTQILSKRTKYVVYLVFKLSKNHDGLEIANAFVRFVNRVSDKEAEERASVVSLVGKRVRRRKRNVKCPRKRVDGWMEIELGNFINDTGDDGDVEARLMEITQLHGKGGLIVQGIEFRPE + +>S4-nosig +DFDYMQLVLTWPPSFCYPTGTCKRTSNNFTIHGLWPEKNRFRLEFCSGGAAYKKFELQDRIVSDLDRHWIQMKFNEQEAKQKQPLWNHEYKRHGRCCYNLYDQNAYFLLAMRLKDKLDLVTTLRTHGITPGTKHTFDEIKSAIKTVTNQVDPDLKCVEHTKGVQELKEIGICFTPSADSFYPCRQSNTCDEKGTAILFR