Commit effacff8 authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

fix: 修复FileDAG的prc/qc状态取值逻辑,更新IFS配置并新增测试与文档

parent 2b89f6a6
Loading
Loading
Loading
Loading
+8 −4
Original line number Diff line number Diff line
@@ -127,6 +127,8 @@ class FileDAG(BaseDAG):
            if not (force_success or this_task.get("success")):
                continue
            task = dict(this_task.get("task") or {})
            submission_prc_status = data.get("prc_status")
            submission_qc_status = data.get("qc_status")
            run = DagRun(
                dag=dag_group["dag"],
                dag_run_group=dag_group["dag_run_group"],
@@ -146,8 +148,8 @@ class FileDAG(BaseDAG):
                detector=task.get("detector"),
                filter=task.get("filter"),
                custom_id=task.get("custom_id"),
                prc_status=task.get("prc_status"),
                qc_status=task.get("qc_status"),
                prc_status=task.get("prc_status") if submission_prc_status in (None, "") else submission_prc_status,
                qc_status=task.get("qc_status") if submission_qc_status in (None, "") else submission_qc_status,
                n_file_expected=task.get("n_file_expected", -1),
                n_file_found=task.get("n_file_found", -1),
                data_list=(this_task.get("relevant_data_id_list") if return_data_list else []),
@@ -191,6 +193,8 @@ class FileDAG(BaseDAG):
            if not (force_success or this_task.get("success")):
                continue
            task = dict(this_task.get("task") or {})
            submission_prc_status = data.get("prc_status")
            submission_qc_status = data.get("qc_status")
            run = DagRun(
                dag=dag_group["dag"],
                dag_run_group=dag_group["dag_run_group"],
@@ -210,8 +214,8 @@ class FileDAG(BaseDAG):
                detector=task.get("detector"),
                filter=task.get("filter"),
                custom_id=task.get("custom_id"),
                prc_status=task.get("prc_status"),
                qc_status=task.get("qc_status"),
                prc_status=task.get("prc_status") if submission_prc_status in (None, "") else submission_prc_status,
                qc_status=task.get("qc_status") if submission_qc_status in (None, "") else submission_qc_status,
                n_file_expected=task.get("n_file_expected", -1),
                n_file_found=task.get("n_file_found", -1),
                data_list=(this_task.get("relevant_data_id_list") if return_data_list else []),
+6 −6
Original line number Diff line number Diff line
@@ -14,24 +14,24 @@ match:
  detector_group: "EFFECTIVE"

# 默认输入参数(JSON 字符串;DAG 的输入统一从该字符串解析)
default_params_input: '{"dag":"csst-ifs-l2","dag_run":"csst-ifs-l2-inttest","dataset":"test-dataset","instrument":"IFS","obs_type":"SCI","obs_group":"G1","obs_id":"","batch_id":"inttest", "docker_images":{"csst-ifs-l2-cube":"latest", "csst-ifs-l2-map":"latest"}}'
default_params_input: '{"dag":"csst-ifs-l2","dag_run":"csst-ifs-l2-inttest","dataset":"csst-ifs-c11-sim-v1","upstream_batch_id":"inttest","data_model":"csst-ifs-l1-a","instrument":"IFS","obs_type":"SCI","obs_group":"NGC6217-cen","obs_id":"","detector":"","batch_id":"inttest", "docker_images":{"csst-ifs-l2-cube":"latest", "csst-ifs-l2-map":"latest"}}'


# Submission 页面默认值配置(用于表单初始化与约束)
submission:
  data:
    dataset:
      default: test-dataset
      default: csst-ifs-c11-sim-v1
      options: []
      fixed: false
      allow_empty: false
    upstream_batch_id:
      default: ''
      default: inttest
      options: []
      fixed: true
      fixed: false
      allow_empty: false
    data_model:
      default: ''
      default: csst-ifs-l1-a
      options: []
      fixed: true
      allow_empty: false
@@ -49,7 +49,7 @@ submission:
      fixed: false
      allow_empty: false
    obs_group:
      default: G1
      default: NGC6217-cen
      options: []
      fixed: false
      allow_empty: true

task.md

0 → 100644
+123 −0
Original line number Diff line number Diff line
# csst-dag v2 字段矩阵

说明:

- 只写当前真实代码行为。
- 纵向看字段,横向看 `FileDAG / StandaloneDAG / BrickDAG`
- 表尽量少。

来源缩写:

- `group输入``DagRunGroup.trigger(batch_id=..., priority=...)`
- `submission.data`:trigger 的 `data`
- `submission.proc`:trigger 的 `proc`
- `task``FileDAG` 中 dispatcher 生成的中间 `task`
- `运行时`:代码自动生成或计算
- `默认值``DagRun` / `DagRunGroup` 模型默认值
- `固定配置`:DAG YAML `tasks[].image``brick_multipliers`

## 1. Trigger 输入字段

| group | 字段 | 在 DAG YAML `submission` 中声明 | 说明 |
| --- | --- | --- | --- |
| `group` | `batch_id` | 否 | 当前新生成 `DagRunGroup` 的批次 |
| `group` | `priority` | 否 | 当前新生成 `DagRunGroup` 的优先级 |
| `data` | `dataset` | 是 | 统一骨架字段 |
| `data` | `upstream_batch_id` | 是 | 上游检索批次,不等于 `group.batch_id` |
| `data` | `data_model` | 是 | 统一骨架字段 |
| `data` | `instrument` | 是 | 统一骨架字段 |
| `data` | `obs_type` | 是 | 统一骨架字段 |
| `data` | `obs_group` | 是 | 统一骨架字段 |
| `data` | `obs_id` | 是 | 统一骨架字段 |
| `data` | `detector` | 是 | 统一骨架字段 |
| `data` | `filter` | 是 | 统一骨架字段 |
| `data` | `healpix` | 是 | 统一骨架字段 |
| `data` | `custom_id` | 是 | 统一骨架字段 |
| `data` | `prc_status` | 是 | 统一骨架字段 |
| `data` | `qc_status` | 是 | 统一骨架字段 |
| `proc` | `pmapname` | 是 | 处理参数 |
| `proc` | `ref_cat` | 是 | 处理参数 |
| `proc` | `extra_kwargs` | 是 | 处理参数 |
| `docker_images` | `<image_name>: <tag>` | 否 | key 必须来自 DAG YAML `tasks[].image` |

补充:

- DAG YAML `submission:` 只声明 `data``proc`
- `group``docker_images` 是 trigger 顶层参数。
-`StandaloneDAG` 会自动补齐 `submission.data` 的统一骨架。
- `fixed/default/options/allow_empty` 属于表单阶段约束,不是 payload 的单独来源类型。

## 2. 非 `data` 字段矩阵

这一张表把 `group / job / proc / aux / run` 一次看完。真正大的差异不在这里,而在下一张 `data` 主表。

| payload group | 字段 | FileDAG | StandaloneDAG | BrickDAG | 说明 |
| --- | --- | --- | --- | --- | --- |
| `group` | `dag` | 运行时 | 运行时 | 运行时 | 当前 DAG 名 |
| `group` | `dag_run_group` | 运行时 | 运行时 | 运行时 | 自动生成 UUIDv7 |
| `group` | `batch_id` | `group输入` | `group输入` | `group输入` | 当前新生成 group 的批次 |
| `group` | `priority` | `group输入` | `group输入` | `group输入` | 当前新生成 group 的优先级 |
| `group` | `created_at` | 运行时 | 运行时 | 运行时 | 自动生成时间 |
| `group` | `docker_images` | `docker_images` + 固定配置 | `docker_images` + 固定配置 | `docker_images` + 固定配置 | 与 DAG YAML `tasks[].image` 合并并校验 |
| `group` | `n_dag_runs` | 运行时 | 运行时 | 运行时 | 根据最终 run 数量计算 |
| `job` | `dag` | 运行时 | 运行时 | 运行时 | 从当前 DAG 写入 |
| `job` | `dag_run` | 运行时 | 运行时 | 运行时 | 自动生成 UUIDv7 |
| `job` | `dag_run_group` | 运行时 | 运行时 | 运行时 | 从父 group 继承 |
| `job` | `batch_id` | `group输入` | `group输入` | `group输入` | 从父 group 继承 |
| `job` | `priority` | `group输入` | `group输入` | `group输入` | 从父 group 继承 |
| `proc` | `pmapname` | `submission.proc` | `submission.proc` | `submission.proc` | 直接写入 |
| `proc` | `ref_cat` | `submission.proc` | `submission.proc` | `submission.proc` | 直接写入 |
| `proc` | `extra_kwargs` | `submission.proc` | `submission.proc` | `submission.proc` | 直接写入 |
| `aux` | `object` | 默认值 | 默认值 | 默认值 | 当前代码未主动赋值 |
| `aux` | `proposal_id` | 默认值 | `submission.data` | 默认值 | 目前只有 `StandaloneDAG` 直接写 |
| `aux` | `data_list` | 运行时 | 固定空列表 | 固定空列表 | `FileDAG` 可写 relevant data id list |
| `aux` | `n_file_expected` | `task`/运行时 | 固定 `0` | 默认值 | `FileDAG` 来自 dispatcher |
| `aux` | `n_file_found` | `task`/运行时 | 固定 `0` | 默认值 | `FileDAG` 来自 dispatcher |
| `run` | `trial` | 默认值 | 默认值 | 默认值 | 默认 `1` |
| `run` | `reason` | 默认值 | 默认值 | 默认值 | 默认 `initial` |
| `run` | `status` | 默认值 | 默认值 | 默认值 | 默认 `received` |
| `run` | `created_at` | 运行时 | 运行时 | 运行时 | 自动生成 |
| `run` | `started_at` | 默认值 | 默认值 | 默认值 | 初始为空 |
| `run` | `finished_at` | 默认值 | 默认值 | 默认值 | 初始为空 |
| `run` | `elapsed_seconds` | 默认值 | 默认值 | 默认值 | 初始为空 |

## 3. `payload.dag_runs[*].data` 主表

这是最关键的一张表。

| group | 字段 | FileDAG | StandaloneDAG | BrickDAG | 说明 |
| --- | --- | --- | --- | --- | --- |
| `data` | `data_model` | `task` | `submission.data` | `submission.data` | `FileDAG` 从 task 取;后两者直接投影 |
| `data` | `dataset` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `upstream_batch_id` | `submission.data` | `submission.data` | `submission.data` | 与 `group.batch_id` 语义分离 |
| `data` | `instrument` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `obs_type` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `obs_group` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `obs_id` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `detector` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `filter` | `task` | `submission.data` | `submission.data` | `FileDAG` 只有 task 里真的带了才会有值 |
| `data` | `healpix` | 默认值 | `submission.data` | 运行时 | `BrickDAG` 由 unique healpix 展开结果覆盖 |
| `data` | `custom_id` | `task` | `submission.data` | `submission.data` | 同上 |
| `data` | `prc_status` | `submission.data` 优先,否则 `task` | `submission.data` | `submission.data` | `FileDAG` 当前唯一显式特例 |
| `data` | `qc_status` | `submission.data` 优先,否则 `task` | `submission.data` | `submission.data` | `FileDAG` 当前唯一显式特例 |

## 4. `FileDAG.task` 说明

这一段只解释主表里 `FileDAG = task` 到底是什么意思。

| 项 | 当前规则 |
| --- | --- |
| `task` 来源 | `Dispatcher.dispatch_*()``data_basis` 分组后的组级投影 |
| `task` 固定附加字段 | `n_file_expected`, `n_file_found` |
| `task` 基础字段 | 只保留当前 `group_by` 中存在于 `data_basis` 的字段 |
| `dispatch_file` 常见字段 | `dataset`, `instrument`, `obs_type`, `obs_group`, `obs_id`, `detector`,以及可能附带唯一文件 ID |
| `dispatch_obsid` 常见字段 | `dataset`, `instrument`, `obs_type`, `obs_group`, `obs_id` |
| `dispatch_obsgroup` 常见字段 | `dataset`, `instrument`, `obs_type`, `obs_group` |
| `dispatch_obsgroup_detector` 常见字段 | `dataset`, `instrument`, `obs_type`, `obs_group`, `detector` |
| `task` 不保证包含 | `filter`, `custom_id`, `prc_status`, `qc_status` 等非分组键字段 |

## 5. 额外字段

| 字段位置 | 字段 | FileDAG | StandaloneDAG | BrickDAG | 说明 |
| --- | --- | --- | --- | --- | --- |
| 根层级(非标准 group) | multiplier 字段 | 无 | 无 | 固定配置 / `submission.data` | 来自 DAG YAML `brick_multipliers[].field`;若 submission 显式给值则优先用 submission |
+60 −0
Original line number Diff line number Diff line
@@ -61,3 +61,63 @@ def test_v2_file_dag_schedule_flow():
    assert len(payload["dag_runs"]) == 1
    assert payload["dag_runs"][0]["data"]["obs_id"] == "1"
    assert payload["dag_runs"][0]["data"]["upstream_batch_id"] == "src-b1"


def test_v2_file_dag_preserves_zero_status_filters_in_payload():
    sys.path.insert(0, "/nfs/airflow/csst-pipeline-portal/external/csst-dag")

    from csst_dag.v2.dag.base import DependencyBundle
    from csst_dag.v2.dag.file import FileDAG

    class DummyFile(FileDAG):
        pass

    dag = DummyFile("csst-msc-l1-ooc")
    group = dag.build_trigger_payload(
        data={
            "dataset": "d1",
            "instrument": "MSC",
            "obs_type": "WIDE",
            "obs_group": "W5",
            "obs_id": "1",
            "batch_id": "b1",
            "qc_status": 0,
            "prc_status": 0,
        },
        proc={"pmapname": "p1"},
        docker_images={},
        deps=DependencyBundle(
            plan=Table(
                [
                    {
                        "dataset": "d1",
                        "instrument": "MSC",
                        "obs_type": "WIDE",
                        "obs_group": "W5",
                        "obs_id": "1",
                    }
                ]
            ),
            data=Table(
                [
                    {
                        "dataset": "d1",
                        "instrument": "MSC",
                        "obs_type": "WIDE",
                        "obs_group": "W5",
                        "obs_id": "1",
                        "detector": "06",
                        "data_uuid": "u-1",
                        "data_model": "raw",
                    }
                ]
            ),
            upstream=[],
        ),
    )
    payload = group.to_payload()
    assert len(group.dag_runs) == 1
    assert group.dag_runs[0].qc_status == 0
    assert group.dag_runs[0].prc_status == 0
    assert payload["dag_runs"][0]["data"]["qc_status"] == 0
    assert payload["dag_runs"][0]["data"]["prc_status"] == 0