Back to home page

EIC code displayed by LXR

 
 

    


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)