Warning, /swf-testbed/agents/snakemake_workflows/example_Snakefile is written in an unsupported language. File is not indexed.
0001 """
0002 Example Snakemake workflow triggered by the Snakemake Agent.
0003
0004 All fields from the incoming ActiveMQ message are available via
0005 config["key"]. For example, when the prompt_processing_agent publishes
0006 an stf_processed message the following keys are available:
0007
0008 config["msg_type"] # "stf_processed"
0009 config["run_id"] # run number
0010 config["panda_task_id"] # PanDA task ID
0011 config["task_status"] # e.g. "done", "failed"
0012 config["processed"] # number of processed files
0013 config["failed"] # number of failed files
0014 config["timed_out"] # "True" / "False"
0015 config["input_dataset"] # input Rucio dataset
0016 config["output_dataset"] # output Rucio dataset
0017 config["execution_id"] # workflow execution ID
0018 config["trigger_timestamp"]# ISO timestamp of when the agent received the message
0019
0020 Note: All values arrive as strings. Cast as needed.
0021 """
0022
0023 # ---------------------------------------------------------------------------
0024 # Configuration with sensible defaults
0025 # ---------------------------------------------------------------------------
0026 RUN_ID = config.get("run_id", "unknown")
0027 TASK_STATUS = config.get("task_status", "unknown")
0028 OUTPUT_DATASET = config.get("output_dataset", "none")
0029 TIMESTAMP = config.get("trigger_timestamp", "no-timestamp")
0030
0031
0032 # ---------------------------------------------------------------------------
0033 # Default target
0034 # ---------------------------------------------------------------------------
0035 rule all:
0036 input:
0037 f"results/run_{RUN_ID}/summary.txt",
0038 f"results/run_{RUN_ID}/status.json",
0039
0040
0041 # ---------------------------------------------------------------------------
0042 # Rules
0043 # ---------------------------------------------------------------------------
0044 rule generate_summary:
0045 """Create a human-readable summary of the processing result."""
0046 output:
0047 f"results/run_{RUN_ID}/summary.txt",
0048 params:
0049 run_id=RUN_ID,
0050 task_status=TASK_STATUS,
0051 output_dataset=OUTPUT_DATASET,
0052 processed=config.get("processed", "0"),
0053 failed=config.get("failed", "0"),
0054 timestamp=TIMESTAMP,
0055 run:
0056 import os
0057 os.makedirs(os.path.dirname(output[0]), exist_ok=True)
0058 with open(output[0], "w") as fh:
0059 fh.write(f"=== Processing Summary ===\n")
0060 fh.write(f"Run ID : {params.run_id}\n")
0061 fh.write(f"Task status : {params.task_status}\n")
0062 fh.write(f"Output dataset : {params.output_dataset}\n")
0063 fh.write(f"Processed : {params.processed}\n")
0064 fh.write(f"Failed : {params.failed}\n")
0065 fh.write(f"Triggered at : {params.timestamp}\n")
0066
0067
0068 rule generate_status_json:
0069 """Create a machine-readable JSON status file."""
0070 output:
0071 f"results/run_{RUN_ID}/status.json",
0072 run:
0073 import json, os
0074 os.makedirs(os.path.dirname(output[0]), exist_ok=True)
0075 status = {k: v for k, v in config.items()}
0076 status["workflow"] = "example_Snakefile"
0077 with open(output[0], "w") as fh:
0078 json.dump(status, fh, indent=2)