From f4ca1ce9670980c041902a7051899c7ee47d41c0 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Mon, 3 Aug 2026 12:14:12 +0000
Subject: [PATCH 1/4] Initial plan
From 7b22f1a8876ee73141b1d4e74d948dadcdd3e4ec Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Mon, 3 Aug 2026 12:20:53 +0000
Subject: [PATCH 2/4] chore: outline dataflow workflow fix plan
Co-authored-by: pelikhan <4175913+pelikhan@users.noreply.github.com>
---
.github/skills/agentic-workflows/SKILL.md | 1 +
1 file changed, 1 insertion(+)
diff --git a/.github/skills/agentic-workflows/SKILL.md b/.github/skills/agentic-workflows/SKILL.md
index 6fb19019416..141445630d5 100644
--- a/.github/skills/agentic-workflows/SKILL.md
+++ b/.github/skills/agentic-workflows/SKILL.md
@@ -71,6 +71,7 @@ Load these files from `github/gh-aw` (they are not available locally).
- `.github/aw/test-coverage.md`
- `.github/aw/test-expression.md`
- `.github/aw/token-optimization-caching-budgets.md`
+- `.github/aw/token-optimization-observability.md`
- `.github/aw/token-optimization.md`
- `.github/aw/triggers.md`
- `.github/aw/update-agentic-workflow.md`
From 653de726d246d2e5662062ec2dead48b88f193ec Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Mon, 3 Aug 2026 12:36:51 +0000
Subject: [PATCH 3/4] fix: bound dataflow dataset workflow runtime and
reporting
Co-authored-by: pelikhan <4175913+pelikhan@users.noreply.github.com>
---
.../dataflow-pr-discussion-dataset.lock.yml | 263 ++++------
.../dataflow-pr-discussion-dataset.md | 463 ++++++++++--------
2 files changed, 357 insertions(+), 369 deletions(-)
diff --git a/.github/workflows/dataflow-pr-discussion-dataset.lock.yml b/.github/workflows/dataflow-pr-discussion-dataset.lock.yml
index 5ff7ad395b6..22df9af951e 100644
--- a/.github/workflows/dataflow-pr-discussion-dataset.lock.yml
+++ b/.github/workflows/dataflow-pr-discussion-dataset.lock.yml
@@ -1,5 +1,5 @@
-# gh-aw-metadata: {"schema_version":"v4","frontmatter_hash":"0442a494289950abde81e098ef655239847438fd679784ffe31e5290f1468528","body_hash":"421b26ff01056f750ef268a35325e6b4a64a8992f21253b54d9f1e8eea4bcba2","strict":true,"agent_id":"copilot","engine_versions":{"copilot":"1.0.77"}}
-# gh-aw-manifest: {"version":1,"secrets":["COPILOT_GITHUB_TOKEN","GH_AW_GITHUB_MCP_SERVER_TOKEN","GH_AW_GITHUB_TOKEN","GH_AW_OTEL_GRAFANA_AUTHORIZATION","GH_AW_OTEL_GRAFANA_ENDPOINT","GH_AW_OTEL_SENTRY_AUTHORIZATION","GH_AW_OTEL_SENTRY_ENDPOINT","GITHUB_TOKEN"],"actions":[{"repo":"actions/cache/restore","sha":"55cc8345863c7cc4c66a329aec7e433d2d1c52a9","version":"v6.1.0"},{"repo":"actions/cache/save","sha":"55cc8345863c7cc4c66a329aec7e433d2d1c52a9","version":"v6.1.0"},{"repo":"actions/checkout","sha":"3d3c42e5aac5ba805825da76410c181273ba90b1","version":"v7.0.1"},{"repo":"actions/download-artifact","sha":"3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c","version":"v8.0.1"},{"repo":"actions/github-script","sha":"3a2844b7e9c422d3c10d287c895573f7108da1b3","version":"v9.0.0"},{"repo":"actions/setup-node","sha":"820762786026740c76f36085b0efc47a31fe5020","version":"v7.0.0"},{"repo":"actions/upload-artifact","sha":"043fb46d1a93c77aae656e7c1c64a875d1fc6a0a","version":"v7.0.1"},{"repo":"safedep/pmg","sha":"5ac0f275b83d9d5a9342c6aae977ec32fa330daa","version":"v1"}],"containers":[{"image":"ghcr.io/github/gh-aw-firewall/agent:0.27.43","digest":"sha256:04e2d1987a565000a8f114b89d806ae7a3864dd4f944be65275b28c93d8690e6","pinned_image":"ghcr.io/github/gh-aw-firewall/agent:0.27.43@sha256:04e2d1987a565000a8f114b89d806ae7a3864dd4f944be65275b28c93d8690e6"},{"image":"ghcr.io/github/gh-aw-firewall/api-proxy:0.27.43","digest":"sha256:d85f57975af5ea23af4996e41ed73fbc8f5b4a47402472bfe82e508f352cb0c1","pinned_image":"ghcr.io/github/gh-aw-firewall/api-proxy:0.27.43@sha256:d85f57975af5ea23af4996e41ed73fbc8f5b4a47402472bfe82e508f352cb0c1"},{"image":"ghcr.io/github/gh-aw-firewall/cli-proxy:0.27.43","digest":"sha256:65c45ea2967984d0024f3df61bc71335658a77ede96c8d9665da7a5f33a795ab","pinned_image":"ghcr.io/github/gh-aw-firewall/cli-proxy:0.27.43@sha256:65c45ea2967984d0024f3df61bc71335658a77ede96c8d9665da7a5f33a795ab"},{"image":"ghcr.io/github/gh-aw-firewall/squid:0.27.43","digest":"sha256:26be5e0b8c8f4c41c8a59126b29bb5d80b07253597472ded2a16bdd75abcbf9d","pinned_image":"ghcr.io/github/gh-aw-firewall/squid:0.27.43@sha256:26be5e0b8c8f4c41c8a59126b29bb5d80b07253597472ded2a16bdd75abcbf9d"},{"image":"ghcr.io/github/gh-aw-mcpg:v0.4.7","digest":"sha256:7545220a9aca134b71e51193ee0eaf4c50756ebf8fbd25a63ae7556e62815c00","pinned_image":"ghcr.io/github/gh-aw-mcpg:v0.4.7@sha256:7545220a9aca134b71e51193ee0eaf4c50756ebf8fbd25a63ae7556e62815c00"},{"image":"ghcr.io/github/gh-aw-node","digest":"sha256:0d9f1fb5fd6610c0ac1f5194a38e45a8a1e81f8a390d5142d8e4e6f26a4b3196","pinned_image":"ghcr.io/github/gh-aw-node@sha256:0d9f1fb5fd6610c0ac1f5194a38e45a8a1e81f8a390d5142d8e4e6f26a4b3196"},{"image":"ghcr.io/github/github-mcp-server:v1.8.0","digest":"sha256:d5a18c04b92714c309eb46a2305087e91a4dbd80420f6e462656699f95093520","pinned_image":"ghcr.io/github/github-mcp-server:v1.8.0@sha256:d5a18c04b92714c309eb46a2305087e91a4dbd80420f6e462656699f95093520"}]}
+# gh-aw-metadata: {"schema_version":"v4","frontmatter_hash":"d3b50c56cd0267b90b93602f0056ffc5b1f8cf26c02f8595f7cda068afa2d5b8","body_hash":"9b06defff555f6d3244af112d46036857781188f0283d999c6e2b0ea27959d9b","strict":true,"agent_id":"copilot","engine_versions":{"copilot":"1.0.77"}}
+# gh-aw-manifest: {"version":1,"secrets":["COPILOT_GITHUB_TOKEN","GH_AW_GITHUB_MCP_SERVER_TOKEN","GH_AW_GITHUB_TOKEN","GH_AW_OTEL_GRAFANA_AUTHORIZATION","GH_AW_OTEL_GRAFANA_ENDPOINT","GH_AW_OTEL_SENTRY_AUTHORIZATION","GH_AW_OTEL_SENTRY_ENDPOINT","GITHUB_TOKEN"],"actions":[{"repo":"actions/cache/restore","sha":"55cc8345863c7cc4c66a329aec7e433d2d1c52a9","version":"v6.1.0"},{"repo":"actions/cache/save","sha":"55cc8345863c7cc4c66a329aec7e433d2d1c52a9","version":"v6.1.0"},{"repo":"actions/checkout","sha":"3d3c42e5aac5ba805825da76410c181273ba90b1","version":"v7.0.1"},{"repo":"actions/download-artifact","sha":"3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c","version":"v8.0.1"},{"repo":"actions/github-script","sha":"3a2844b7e9c422d3c10d287c895573f7108da1b3","version":"v9.0.0"},{"repo":"actions/setup-node","sha":"820762786026740c76f36085b0efc47a31fe5020","version":"v7.0.0"},{"repo":"actions/setup-python","sha":"5fda3b95a4ea91299a34e894583c3862153e4b97","version":"v7.0.0"},{"repo":"actions/upload-artifact","sha":"043fb46d1a93c77aae656e7c1c64a875d1fc6a0a","version":"v7.0.1"},{"repo":"safedep/pmg","sha":"5ac0f275b83d9d5a9342c6aae977ec32fa330daa","version":"v1"}],"containers":[{"image":"ghcr.io/github/gh-aw-firewall/agent:0.27.43","digest":"sha256:04e2d1987a565000a8f114b89d806ae7a3864dd4f944be65275b28c93d8690e6","pinned_image":"ghcr.io/github/gh-aw-firewall/agent:0.27.43@sha256:04e2d1987a565000a8f114b89d806ae7a3864dd4f944be65275b28c93d8690e6"},{"image":"ghcr.io/github/gh-aw-firewall/api-proxy:0.27.43","digest":"sha256:d85f57975af5ea23af4996e41ed73fbc8f5b4a47402472bfe82e508f352cb0c1","pinned_image":"ghcr.io/github/gh-aw-firewall/api-proxy:0.27.43@sha256:d85f57975af5ea23af4996e41ed73fbc8f5b4a47402472bfe82e508f352cb0c1"},{"image":"ghcr.io/github/gh-aw-firewall/cli-proxy:0.27.43","digest":"sha256:65c45ea2967984d0024f3df61bc71335658a77ede96c8d9665da7a5f33a795ab","pinned_image":"ghcr.io/github/gh-aw-firewall/cli-proxy:0.27.43@sha256:65c45ea2967984d0024f3df61bc71335658a77ede96c8d9665da7a5f33a795ab"},{"image":"ghcr.io/github/gh-aw-firewall/squid:0.27.43","digest":"sha256:26be5e0b8c8f4c41c8a59126b29bb5d80b07253597472ded2a16bdd75abcbf9d","pinned_image":"ghcr.io/github/gh-aw-firewall/squid:0.27.43@sha256:26be5e0b8c8f4c41c8a59126b29bb5d80b07253597472ded2a16bdd75abcbf9d"},{"image":"ghcr.io/github/gh-aw-mcpg:v0.4.7","digest":"sha256:7545220a9aca134b71e51193ee0eaf4c50756ebf8fbd25a63ae7556e62815c00","pinned_image":"ghcr.io/github/gh-aw-mcpg:v0.4.7@sha256:7545220a9aca134b71e51193ee0eaf4c50756ebf8fbd25a63ae7556e62815c00"},{"image":"ghcr.io/github/gh-aw-node","digest":"sha256:0d9f1fb5fd6610c0ac1f5194a38e45a8a1e81f8a390d5142d8e4e6f26a4b3196","pinned_image":"ghcr.io/github/gh-aw-node@sha256:0d9f1fb5fd6610c0ac1f5194a38e45a8a1e81f8a390d5142d8e4e6f26a4b3196"},{"image":"ghcr.io/github/github-mcp-server:v1.8.0","digest":"sha256:d5a18c04b92714c309eb46a2305087e91a4dbd80420f6e462656699f95093520","pinned_image":"ghcr.io/github/github-mcp-server:v1.8.0@sha256:d5a18c04b92714c309eb46a2305087e91a4dbd80420f6e462656699f95093520"}]}
# This file was automatically generated by gh-aw. DO NOT EDIT. To debug this workflow, load the skill at https://github.com/github/gh-aw/blob/main/debug.md
#
# ___ _ _
@@ -30,7 +30,6 @@
# - shared/discussions-data-fetch.md
# - shared/otlp.md
# - shared/pmg.md
-# - shared/repo-memory-standard.md
# - shared/reporting.md
#
# Secrets used:
@@ -51,6 +50,7 @@
# - actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
# - actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0 (source v9)
# - actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0
+# - actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0
# - actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
# - safedep/pmg@5ac0f275b83d9d5a9342c6aae977ec32fa330daa # v1
#
@@ -307,27 +307,25 @@ jobs:
GH_AW_GITHUB_RUN_ID: ${{ github.run_id }}
GH_AW_GITHUB_SERVER_URL: ${{ github.server_url }}
GH_AW_GITHUB_WORKSPACE: ${{ github.workspace }}
- GH_AW_WIKI_NOTE: ${{ '' }}
# poutine:ignore untrusted_checkout_exec
run: |
bash "${RUNNER_TEMP}/gh-aw/actions/create_prompt_first.sh"
{
- cat << 'GH_AW_PROMPT_4618c422275e1a44_EOF'
+ cat << 'GH_AW_PROMPT_2fab087bd0af0f8a_EOF'
- GH_AW_PROMPT_4618c422275e1a44_EOF
+ GH_AW_PROMPT_2fab087bd0af0f8a_EOF
cat "${RUNNER_TEMP}/gh-aw/prompts/xpia.md"
cat "${RUNNER_TEMP}/gh-aw/prompts/temp_folder_prompt.md"
cat "${RUNNER_TEMP}/gh-aw/prompts/markdown.md"
cat "${RUNNER_TEMP}/gh-aw/prompts/cache_memory_prompt.md"
- cat "${RUNNER_TEMP}/gh-aw/prompts/repo_memory_prompt.md"
cat "${RUNNER_TEMP}/gh-aw/prompts/safe_outputs_prompt.md"
- cat << 'GH_AW_PROMPT_4618c422275e1a44_EOF'
+ cat << 'GH_AW_PROMPT_2fab087bd0af0f8a_EOF'
Tools: create_discussion, missing_tool, missing_data, noop
- GH_AW_PROMPT_4618c422275e1a44_EOF
+ GH_AW_PROMPT_2fab087bd0af0f8a_EOF
cat "${RUNNER_TEMP}/gh-aw/prompts/mcp_cli_tools_prompt.md"
- cat << 'GH_AW_PROMPT_4618c422275e1a44_EOF'
+ cat << 'GH_AW_PROMPT_2fab087bd0af0f8a_EOF'
The following GitHub context information is available for this workflow:
{{#if github.actor}}
@@ -356,16 +354,16 @@ jobs:
{{/if}}
- GH_AW_PROMPT_4618c422275e1a44_EOF
+ GH_AW_PROMPT_2fab087bd0af0f8a_EOF
cat "${RUNNER_TEMP}/gh-aw/prompts/cli_proxy_with_safeoutputs_prompt.md"
- cat << 'GH_AW_PROMPT_4618c422275e1a44_EOF'
+ cat << 'GH_AW_PROMPT_2fab087bd0af0f8a_EOF'
{{#runtime-import .github/workflows/shared/pmg.md}}
{{#runtime-import .github/workflows/shared/discussions-data-fetch.md}}
{{#runtime-import .github/workflows/shared/reporting.md}}
{{#runtime-import .github/workflows/shared/otlp.md}}
{{#runtime-import .github/workflows/dataflow-pr-discussion-dataset.md}}
- GH_AW_PROMPT_4618c422275e1a44_EOF
+ GH_AW_PROMPT_2fab087bd0af0f8a_EOF
} > "$GH_AW_PROMPT"
- name: Interpolate variables and render templates
uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
@@ -400,12 +398,6 @@ jobs:
GH_AW_GITHUB_SERVER_URL: ${{ github.server_url }}
GH_AW_GITHUB_WORKSPACE: ${{ github.workspace }}
GH_AW_MCP_CLI_SERVERS_LIST: '- `safeoutputs` — run `safeoutputs --help` to see available tools'
- GH_AW_MEMORY_BRANCH_NAME: 'memory/dataflow-dataset'
- GH_AW_MEMORY_CONSTRAINTS: "\n\n**Constraints:**\n- **Allowed Files**: Only files matching patterns: *.json, *.jsonl, *.csv, *.md\n- **Max File Size**: 102400 bytes (0.10 MB) per file\n- **Max File Count**: 100 files per commit\n- **Max Patch Size**: 10240 bytes (10 KB) total per push (max: 1024 KB)\n"
- GH_AW_MEMORY_DESCRIPTION: ' Tracks dataset build statistics and run metadata for the DataFlow pipeline'
- GH_AW_MEMORY_DIR: '/tmp/gh-aw/repo-memory/default/'
- GH_AW_MEMORY_TARGET_REPO: ' of the current repository'
- GH_AW_WIKI_NOTE: ''
with:
script: |
const { setupGlobals } = require('${{ runner.temp }}/gh-aw/actions/setup_globals.cjs');
@@ -430,13 +422,7 @@ jobs:
GH_AW_GITHUB_RUN_ID: process.env.GH_AW_GITHUB_RUN_ID,
GH_AW_GITHUB_SERVER_URL: process.env.GH_AW_GITHUB_SERVER_URL,
GH_AW_GITHUB_WORKSPACE: process.env.GH_AW_GITHUB_WORKSPACE,
- GH_AW_MCP_CLI_SERVERS_LIST: process.env.GH_AW_MCP_CLI_SERVERS_LIST,
- GH_AW_MEMORY_BRANCH_NAME: process.env.GH_AW_MEMORY_BRANCH_NAME,
- GH_AW_MEMORY_CONSTRAINTS: process.env.GH_AW_MEMORY_CONSTRAINTS,
- GH_AW_MEMORY_DESCRIPTION: process.env.GH_AW_MEMORY_DESCRIPTION,
- GH_AW_MEMORY_DIR: process.env.GH_AW_MEMORY_DIR,
- GH_AW_MEMORY_TARGET_REPO: process.env.GH_AW_MEMORY_TARGET_REPO,
- GH_AW_WIKI_NOTE: process.env.GH_AW_WIKI_NOTE
+ GH_AW_MCP_CLI_SERVERS_LIST: process.env.GH_AW_MCP_CLI_SERVERS_LIST
}
});
- name: Validate prompt placeholders
@@ -554,6 +540,10 @@ jobs:
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
with:
persist-credentials: false
+ - name: Setup Python
+ uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0
+ with:
+ python-version: '3.12'
- name: Create gh-aw temp directory
run: bash "${RUNNER_TEMP}/gh-aw/actions/create_gh_aw_tmp_dir.sh"
- name: Configure gh CLI for GitHub Enterprise
@@ -596,16 +586,6 @@ jobs:
GH_AW_CACHE_DIR: /tmp/gh-aw/cache-memory
GH_AW_MIN_INTEGRITY: approved
run: bash "${RUNNER_TEMP}/gh-aw/actions/setup_cache_memory_git.sh"
- # Repo memory git-based storage configuration from frontmatter processed below
- - name: Clone repo-memory branch (default)
- env:
- GH_TOKEN: ${{ github.token }}
- GITHUB_SERVER_URL: ${{ github.server_url }}
- BRANCH_NAME: memory/dataflow-dataset
- TARGET_REPO: ${{ github.repository }}
- MEMORY_DIR: /tmp/gh-aw/repo-memory/default
- CREATE_ORPHAN: true
- run: bash "${RUNNER_TEMP}/gh-aw/actions/clone_repo_memory_branch.sh"
- name: Install gh CLI
run: |
bash "${RUNNER_TEMP}/gh-aw/actions/install_gh_cli.sh"
@@ -761,18 +741,89 @@ jobs:
REPO_OWNER: ${{ github.repository_owner }}
- name: Install DataFlow
run: |
+ mkdir -p /tmp/gh-aw/agent/dataflow/{input,output,pipeline,reports}
mkdir -p /tmp/gh-aw/python
python3 -m venv /tmp/gh-aw/python/venv
- /tmp/gh-aw/python/venv/bin/pip install --quiet open-dataflow
- /tmp/gh-aw/python/venv/bin/python3 -c "
- import dataflow
- print('DataFlow', getattr(dataflow, '__version__', 'installed'), 'ready')
- # Print available operators for reference
- import pkgutil, dataflow.operators as ops
- available = [m.name for m in pkgutil.iter_modules(ops.__path__)]
- print('Operator modules:', available)
- "
- mkdir -p /tmp/gh-aw/agent/dataflow/{input,output,pipeline,reports}
+
+ VENV=/tmp/gh-aw/python/venv
+ STATUS=/tmp/gh-aw/agent/dataflow/output/dataflow_runtime.json
+ LOG=/tmp/gh-aw/agent/dataflow/output/dataflow_install.log
+ : > "$LOG"
+
+ dataflow_ready=false
+ if timeout 5m "$VENV/bin/pip" install --quiet --disable-pip-version-check --retries 3 --timeout 60 "uv==0.8.3" >>"$LOG" 2>&1 &&
+ timeout 20m env UV_HTTP_TIMEOUT=60 UV_HTTP_RETRIES=3 \
+ "$VENV/bin/uv" pip install --python "$VENV/bin/python3" "open-dataflow==1.0.10" >>"$LOG" 2>&1 &&
+ timeout 2m "$VENV/bin/python3" - <<'PY' >>"$LOG" 2>&1
+ import json
+ from pathlib import Path
+ from dataflow.utils.storage import FileStorage
+ from dataflow.operators.general_text import (
+ CharNumberFilter,
+ HashDeduplicateFilter,
+ MinHashDeduplicateFilter,
+ )
+
+ fixture_dir = Path("/tmp/gh-aw/agent/dataflow/smoke")
+ fixture_dir.mkdir(parents=True, exist_ok=True)
+ fixture = fixture_dir / "fixture.jsonl"
+ fixture.write_text(
+ "\n".join(
+ json.dumps(record)
+ for record in [
+ {"id": "valid-1", "text": "Alpha words " * 8},
+ {"id": "short", "text": "tiny"},
+ {"id": "valid-1-dup", "text": "Alpha words " * 8},
+ ]
+ )
+ + "\n"
+ )
+
+ storage = FileStorage(first_entry_file_name=str(fixture), cache_path=str(fixture_dir / "cache"))
+ storage.step()
+ CharNumberFilter(threshold=50).run(storage=storage, input_key="text")
+ after_length = storage.step().read("dict")
+ if len(after_length) != 2:
+ raise RuntimeError(f"unexpected CharNumberFilter result: {len(after_length)}")
+
+ _ = MinHashDeduplicateFilter(threshold=0.85)
+ storage.write(after_length)
+ storage.step()
+ HashDeduplicateFilter().run(storage=storage, input_key="text")
+ after_dedup = storage.step().read("dict")
+ if len(after_dedup) != 1:
+ raise RuntimeError(f"unexpected HashDeduplicateFilter result: {len(after_dedup)}")
+
+ print("DataFlow API smoke test passed")
+ PY
+ then
+ dataflow_ready=true
+ else
+ echo "::warning::DataFlow installation or API validation failed; the workflow will use the pure-Python or mixed fallback path."
+ fi
+
+ "$VENV/bin/python3" - < "${RUNNER_TEMP}/gh-aw/safeoutputs/config.json" << 'GH_AW_SAFE_OUTPUTS_CONFIG_66ea2c9eeeebad01_EOF'
- {"create_discussion":{"category":"audits","close_older_discussions":true,"expires":168,"fallback_to_issue":true,"max":1,"title_prefix":"[dataflow-dataset] "},"create_report_incomplete_issue":{},"missing_data":{},"missing_tool":{},"noop":{"max":1,"report-as-issue":"true"},"push_repo_memory":{"memories":[{"dir":"/tmp/gh-aw/repo-memory/default","id":"default","max_file_count":100,"max_file_size":102400,"max_patch_size":10240}]},"report_incomplete":{},"upload_artifact":{"max-size-bytes":104857600,"max-uploads":3,"retention-days":30,"skip-archive":false}}
- GH_AW_SAFE_OUTPUTS_CONFIG_66ea2c9eeeebad01_EOF
+ cat > "${RUNNER_TEMP}/gh-aw/safeoutputs/config.json" << 'GH_AW_SAFE_OUTPUTS_CONFIG_85dfe4c486f56ca8_EOF'
+ {"create_discussion":{"category":"audits","close_older_discussions":true,"expires":168,"fallback_to_issue":true,"max":1,"title_prefix":"[dataflow-dataset] "},"create_report_incomplete_issue":{},"missing_data":{},"missing_tool":{},"noop":{"max":1,"report-as-issue":"true"},"report_incomplete":{},"upload_artifact":{"max-size-bytes":104857600,"max-uploads":3,"retention-days":30,"skip-archive":false}}
+ GH_AW_SAFE_OUTPUTS_CONFIG_85dfe4c486f56ca8_EOF
- name: Generate Safe Outputs Tools
env:
GH_AW_TOOLS_META_JSON: |
@@ -1339,21 +1390,6 @@ jobs:
if [ ! -f /tmp/gh-aw/agent_output.json ]; then
echo '{"items":[]}' > /tmp/gh-aw/agent_output.json
fi
- # Upload repo memory as artifacts for push job
- - name: Sanitize repo-memory filenames (default)
- if: always()
- continue-on-error: true
- env:
- MEMORY_DIR: /tmp/gh-aw/repo-memory/default
- run: bash "${RUNNER_TEMP}/gh-aw/actions/sanitize_repo_memory_filenames.sh"
- - name: Upload repo-memory artifact (default)
- if: always()
- uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1
- with:
- name: repo-memory-default
- path: /tmp/gh-aw/repo-memory/default
- retention-days: 1
- if-no-files-found: ignore
- name: Commit cache-memory changes
if: always()
env:
@@ -1419,7 +1455,6 @@ jobs:
- evals
- push_evals_state
- push_experiments_state
- - push_repo_memory
- safe_outputs
- update_cache_memory
if: >
@@ -1671,10 +1706,6 @@ jobs:
GH_AW_DAILY_AI_CREDITS_TOTAL_EFFECTIVE_TOKENS: ${{ needs.activation.outputs.daily_ai_credits_total_effective_tokens }}
GH_AW_DAILY_AI_CREDITS_THRESHOLD: ${{ needs.activation.outputs.daily_ai_credits_threshold }}
GH_AW_SAFE_OUTPUT_MESSAGES: "{\"footer\":\"\\u003e 🌊 *Dataset built by [{workflow_name}]({run_url})*{ai_credits_suffix}{history_link}\",\"runStarted\":\"🌊 DataFlow Dataset Builder starting! [{workflow_name}]({run_url}) is processing discussions and PRs with OpenDCAI/DataFlow...\",\"runSuccess\":\"✅ DataFlow dataset ready! [{workflow_name}]({run_url}) produced a cleaned, deduplicated dataset. Artifacts uploaded. 📊\",\"runFailure\":\"⚠️ DataFlow pipeline failed! [{workflow_name}]({run_url}) {status}. Check the run logs.\"}"
- GH_AW_PUSH_REPO_MEMORY_RESULT: ${{ needs.push_repo_memory.result }}
- GH_AW_REPO_MEMORY_VALIDATION_FAILED_default: ${{ needs.push_repo_memory.outputs.validation_failed_default }}
- GH_AW_REPO_MEMORY_VALIDATION_ERROR_default: ${{ needs.push_repo_memory.outputs.validation_error_default }}
- GH_AW_REPO_MEMORY_PATCH_SIZE_EXCEEDED_default: ${{ needs.push_repo_memory.outputs.patch_size_exceeded_default }}
GH_AW_GROUP_REPORTS: "false"
GH_AW_FAILURE_REPORT_AS_ISSUE: "true"
GH_AW_MISSING_TOOL_REPORT_AS_FAILURE: "true"
@@ -2326,100 +2357,6 @@ jobs:
clean: false
persist-credentials: false
- push_repo_memory:
- needs:
- - activation
- - agent
- - detection
- if: >
- always() && (!cancelled()) && (needs.detection.result == 'success' || needs.detection.result == 'skipped') &&
- needs.agent.result == 'success'
- runs-on: ubuntu-slim
- permissions:
- contents: write
- concurrency:
- group: "push-repo-memory-${{ github.repository }}|memory/dataflow-dataset"
- cancel-in-progress: false
- env:
- GH_AW_RUNTIME_FEATURES: ${{ vars.GH_AW_RUNTIME_FEATURES }}
- outputs:
- patch_size_exceeded_default: ${{ steps.push_repo_memory_default.outputs.patch_size_exceeded }}
- validation_error_default: ${{ steps.push_repo_memory_default.outputs.validation_error }}
- validation_failed_default: ${{ steps.push_repo_memory_default.outputs.validation_failed }}
- steps:
- - name: Checkout actions folder
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- with:
- repository: github/gh-aw
- sparse-checkout: |
- actions
- clean: false
- persist-credentials: false
- - name: Setup Scripts
- id: setup
- uses: ./actions/setup
- with:
- destination: ${{ runner.temp }}/gh-aw/actions
- job-name: ${{ github.job }}
- trace-id: ${{ needs.activation.outputs.setup-trace-id }}
- parent-span-id: ${{ needs.activation.outputs.setup-parent-span-id || needs.activation.outputs.setup-span-id }}
- env:
- GH_AW_SETUP_WORKFLOW_NAME: "DataFlow PR & Discussion Dataset Builder"
- GH_AW_CURRENT_WORKFLOW_REF: ${{ github.repository }}/.github/workflows/dataflow-pr-discussion-dataset.lock.yml@${{ github.ref }}
- GH_AW_INFO_VERSION: "1.0.77"
- GH_AW_INFO_AWF_VERSION: "v0.27.43"
- GH_AW_INFO_ENGINE_ID: "copilot"
- - name: Checkout repository
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- with:
- persist-credentials: false
- sparse-checkout: .
- - name: Configure Git credentials
- env:
- GITHUB_REPOSITORY: ${{ github.repository }}
- GITHUB_SERVER_URL: ${{ github.server_url }}
- GITHUB_TOKEN: ${{ github.token }}
- run: bash "${RUNNER_TEMP}/gh-aw/actions/configure_git_credentials.sh"
- - name: Download repo-memory artifact (default)
- uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8.0.1
- continue-on-error: true
- with:
- name: repo-memory-default
- path: /tmp/gh-aw/repo-memory/default
- - name: Push repo-memory changes (default)
- id: push_repo_memory_default
- if: always()
- uses: actions/github-script@3a2844b7e9c422d3c10d287c895573f7108da1b3 # v9.0.0
- env:
- GH_TOKEN: ${{ github.token }}
- GITHUB_RUN_ID: ${{ github.run_id }}
- GITHUB_SERVER_URL: ${{ github.server_url }}
- ARTIFACT_DIR: /tmp/gh-aw/repo-memory/default
- MEMORY_ID: default
- TARGET_REPO: ${{ github.repository }}
- BRANCH_NAME: memory/dataflow-dataset
- MAX_FILE_SIZE: 102400
- MAX_FILE_COUNT: 100
- MAX_PATCH_SIZE: 10240
- ALLOWED_EXTENSIONS: '[]'
- FILE_GLOB_FILTER: "*.json *.jsonl *.csv *.md"
- with:
- script: |
- const { setupGlobals } = require('${{ runner.temp }}/gh-aw/actions/setup_globals.cjs');
- setupGlobals(core, github, context, exec, io, getOctokit);
- const { main } = require('${{ runner.temp }}/gh-aw/actions/push_repo_memory.cjs');
- await main();
- - name: Restore actions folder
- if: always()
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- with:
- repository: github/gh-aw
- sparse-checkout: |
- actions/setup
- sparse-checkout-cone-mode: true
- clean: false
- persist-credentials: false
-
safe_outputs:
needs:
- activation
diff --git a/.github/workflows/dataflow-pr-discussion-dataset.md b/.github/workflows/dataflow-pr-discussion-dataset.md
index 551ff40e32a..e122a767c9f 100644
--- a/.github/workflows/dataflow-pr-discussion-dataset.md
+++ b/.github/workflows/dataflow-pr-discussion-dataset.md
@@ -24,10 +24,6 @@ network:
imports:
- shared/pmg.md
- uses: shared/discussions-data-fetch.md
- - uses: shared/repo-memory-standard.md
- with:
- branch-name: memory/dataflow-dataset
- description: Tracks dataset build statistics and run metadata for the DataFlow pipeline
- shared/reporting.md
- shared/otlp.md
tools:
@@ -42,18 +38,89 @@ tools:
steps:
- name: Install DataFlow
run: |
+ mkdir -p /tmp/gh-aw/agent/dataflow/{input,output,pipeline,reports}
mkdir -p /tmp/gh-aw/python
python3 -m venv /tmp/gh-aw/python/venv
- /tmp/gh-aw/python/venv/bin/pip install --quiet open-dataflow
- /tmp/gh-aw/python/venv/bin/python3 -c "
- import dataflow
- print('DataFlow', getattr(dataflow, '__version__', 'installed'), 'ready')
- # Print available operators for reference
- import pkgutil, dataflow.operators as ops
- available = [m.name for m in pkgutil.iter_modules(ops.__path__)]
- print('Operator modules:', available)
- "
- mkdir -p /tmp/gh-aw/agent/dataflow/{input,output,pipeline,reports}
+
+ VENV=/tmp/gh-aw/python/venv
+ STATUS=/tmp/gh-aw/agent/dataflow/output/dataflow_runtime.json
+ LOG=/tmp/gh-aw/agent/dataflow/output/dataflow_install.log
+ : > "$LOG"
+
+ dataflow_ready=false
+ if timeout 5m "$VENV/bin/pip" install --quiet --disable-pip-version-check --retries 3 --timeout 60 "uv==0.8.3" >>"$LOG" 2>&1 &&
+ timeout 20m env UV_HTTP_TIMEOUT=60 UV_HTTP_RETRIES=3 \
+ "$VENV/bin/uv" pip install --python "$VENV/bin/python3" "open-dataflow==1.0.10" >>"$LOG" 2>&1 &&
+ timeout 2m "$VENV/bin/python3" - <<'PY' >>"$LOG" 2>&1
+ import json
+ from pathlib import Path
+ from dataflow.utils.storage import FileStorage
+ from dataflow.operators.general_text import (
+ CharNumberFilter,
+ HashDeduplicateFilter,
+ MinHashDeduplicateFilter,
+ )
+
+ fixture_dir = Path("/tmp/gh-aw/agent/dataflow/smoke")
+ fixture_dir.mkdir(parents=True, exist_ok=True)
+ fixture = fixture_dir / "fixture.jsonl"
+ fixture.write_text(
+ "\n".join(
+ json.dumps(record)
+ for record in [
+ {"id": "valid-1", "text": "Alpha words " * 8},
+ {"id": "short", "text": "tiny"},
+ {"id": "valid-1-dup", "text": "Alpha words " * 8},
+ ]
+ )
+ + "\n"
+ )
+
+ storage = FileStorage(first_entry_file_name=str(fixture), cache_path=str(fixture_dir / "cache"))
+ storage.step()
+ CharNumberFilter(threshold=50).run(storage=storage, input_key="text")
+ after_length = storage.step().read("dict")
+ if len(after_length) != 2:
+ raise RuntimeError(f"unexpected CharNumberFilter result: {len(after_length)}")
+
+ _ = MinHashDeduplicateFilter(threshold=0.85)
+ storage.write(after_length)
+ storage.step()
+ HashDeduplicateFilter().run(storage=storage, input_key="text")
+ after_dedup = storage.step().read("dict")
+ if len(after_dedup) != 1:
+ raise RuntimeError(f"unexpected HashDeduplicateFilter result: {len(after_dedup)}")
+
+ print("DataFlow API smoke test passed")
+ PY
+ then
+ dataflow_ready=true
+ else
+ echo "::warning::DataFlow installation or API validation failed; the workflow will use the pure-Python or mixed fallback path."
+ fi
+
+ "$VENV/bin/python3" - <0.25.
-Deduplicate with MinHash threshold=0.85, or exact-hash fallback.
-If DataFlow is unavailable, use pure Python fallback.
-Compute retention_rate = output/input, upload `dataset_clean.jsonl` as artifact, post a stats table (input count, output count, retention rate, operators used) to a GitHub Discussion in category `audits`, and update the memory branch with run metadata.
+Use the runtime status from `/tmp/gh-aw/agent/dataflow/output/dataflow_runtime.json`.
+Keep the 50–100,000 character bound and 0.25 alphabetic-character ratio.
+Use current DataFlow operators from `dataflow.operators.general_text` for supported stages, skip `AlphaWordsFilter`, and fall back to Python where needed so the report can distinguish `mixed` from `fallback`.
+Compute retention_rate = output/input, upload `dataset_clean.jsonl` as artifact, and post a stats table (including execution mode and operators actually used) to a GitHub Discussion in category `audits`.
{{else}}
# DataFlow PR & Discussion Dataset Builder
@@ -159,13 +226,13 @@ GitHub Discussions + PRs
│
▼
┌─────────────┐
- │ DataFlow │ Text length filter → alpha-ratio filter → stop-word filter
+ │ Mixed │ DataFlow char-length filter → Python alpha-ratio filter
│ Filters │
└──────┬──────┘
│
▼
┌─────────────┐
- │ DataFlow │ Near-duplicate removal (MinHash or exact-hash)
+ │ DataFlow │ Near-duplicate removal (MinHash or exact-hash fallback)
│ Dedup │
└──────┬──────┘
│
@@ -177,38 +244,17 @@ GitHub Discussions + PRs
## Step-by-Step Instructions
-### Step 1: Inspect Available DataFlow Operators
-
-Before building the pipeline, discover which operators are installed:
-
-```bash
-/tmp/gh-aw/python/venv/bin/python3 -c "
-import pkgutil, dataflow.operators as ops
-for m in pkgutil.iter_modules(ops.__path__):
- print(m.name)
-"
-```
+### Step 1: Inspect the DataFlow Runtime Status
-Then list classes in the `filter` and `dedup` sub-modules (if present):
-
-```bash
-/tmp/gh-aw/python/venv/bin/python3 -c "
-import inspect
-try:
- import dataflow.operators.filter as f
- print('filter:', [n for n, _ in inspect.getmembers(f, inspect.isclass)])
-except Exception as e:
- print('filter module error:', e)
-
-try:
- import dataflow.operators.dedup as d
- print('dedup:', [n for n, _ in inspect.getmembers(d, inspect.isclass)])
-except Exception as e:
- print('dedup module error:', e)
-"
-```
+Read `/tmp/gh-aw/agent/dataflow/output/dataflow_runtime.json` first.
-Use the discovered class names throughout the pipeline below.
+- If `dataflow_ready` is `true`, use the smoke-tested current API:
+ - `dataflow.utils.storage.FileStorage`
+ - `dataflow.operators.general_text.CharNumberFilter`
+ - `dataflow.operators.general_text.MinHashDeduplicateFilter`
+ - `dataflow.operators.general_text.HashDeduplicateFilter`
+- If `dataflow_ready` is `false`, print the recorded warning and use the pure-Python fallback pipeline.
+- Never use `AlphaWordsFilter`; it can trigger an implicit NLTK download.
### Step 2: Normalise Raw Data into JSONL
@@ -296,31 +342,32 @@ Write `/tmp/gh-aw/agent/dataflow/pipeline/02_pipeline.py`:
```python
#!/usr/bin/env python3
"""
-DataFlow text processing pipeline:
- 1. Load JSONL into FileStorage
- 2. Text length filter (heuristic — no LLM required)
- 3. Alpha-ratio filter (heuristic — no LLM required)
- 4. Near-duplicate removal (MinHash or exact-hash — no LLM required)
- 5. Save clean output
+Dataset cleaning pipeline:
+ 1. Load JSONL records and the runtime status file
+ 2. DataFlow char-length filter when the validated API is available
+ 3. Deterministic Python alpha-ratio filter (to avoid NLTK downloads)
+ 4. DataFlow near-deduplication (MinHash, then exact-hash fallback) when available
+ 5. Pure-Python exact-hash fallback when DataFlow is unavailable
+ 6. Save clean output and honest execution-mode statistics
"""
-import json, sys, inspect, traceback
+import json
from pathlib import Path
INPUT = "/tmp/gh-aw/agent/dataflow/input/combined_raw.jsonl"
OUTPUT = "/tmp/gh-aw/agent/dataflow/output/dataset_clean.jsonl"
STATS = "/tmp/gh-aw/agent/dataflow/output/pipeline_stats.json"
+RUNTIME = "/tmp/gh-aw/agent/dataflow/output/dataflow_runtime.json"
Path("/tmp/gh-aw/agent/dataflow/output").mkdir(parents=True, exist_ok=True)
-# ── Load DataFlow storage ─────────────────────────────────────────────────────
-try:
- from dataflow.utils.storage import FileStorage
- storage = FileStorage(first_entry_file_name=INPUT)
- print(f"FileStorage loaded with {len(storage)} records")
-except Exception as e:
- print(f"FileStorage error: {e} — falling back to raw JSONL processing")
- storage = None
+runtime = {}
+if Path(RUNTIME).exists():
+ runtime = json.loads(Path(RUNTIME).read_text())
+
+dataflow_ready = bool(runtime.get("dataflow_ready"))
+if not dataflow_ready and runtime.get("warning"):
+ print(f"Runtime warning: {runtime['warning']}")
stats = {
"input_count": 0,
@@ -328,117 +375,110 @@ stats = {
"after_alpha_filter": 0,
"after_dedup": 0,
"operators_used": [],
- "fallback_mode": storage is None,
+ "execution_mode": "fallback",
+ "fallback_mode": True,
+ "dataflow_ready": dataflow_ready,
+ "dataflow_version": runtime.get("dataflow_version", ""),
+ "warnings": [],
}
+if runtime.get("warning"):
+ stats["warnings"].append(runtime["warning"])
-# ── Helper: count records in a JSONL file ────────────────────────────────────
-def count_jsonl(path: str) -> int:
- try:
- return sum(1 for _ in open(path))
- except FileNotFoundError:
- return 0
+def alpha_ratio(text: str) -> float:
+ letters = sum(1 for c in text if c.isalpha())
+ return letters / max(len(text), 1)
-# ── If FileStorage is available, attempt DataFlow operators ──────────────────
-if storage is not None:
- stats["input_count"] = len(storage)
- current_file = INPUT
+def exact_hash_dedup(records):
+ import hashlib
- # ── 1. Text length filter ─────────────────────────────────────────────────
- try:
- from dataflow.operators.filter import TextLengthFilter
- op = TextLengthFilter(min_len=50, max_len=100_000)
- op.run(storage=storage.step(), input_key="text")
- stats["after_length_filter"] = len(storage)
- stats["operators_used"].append("TextLengthFilter")
- print(f"After TextLengthFilter: {stats['after_length_filter']} records")
- except Exception as e:
- print(f"TextLengthFilter skipped: {e}")
- stats["after_length_filter"] = len(storage)
+ seen_hashes = set()
+ kept = []
+ for record in records:
+ digest = hashlib.sha256(record.get("text", "").encode("utf-8")).hexdigest()
+ if digest in seen_hashes:
+ continue
+ seen_hashes.add(digest)
+ kept.append(record)
+ return kept
- # ── 2. Alpha-ratio filter ─────────────────────────────────────────────────
- try:
- from dataflow.operators.filter import AlphaRatioFilter
- op = AlphaRatioFilter(min_ratio=0.25)
- op.run(storage=storage.step(), input_key="text")
- stats["after_alpha_filter"] = len(storage)
- stats["operators_used"].append("AlphaRatioFilter")
- print(f"After AlphaRatioFilter: {stats['after_alpha_filter']} records")
- except Exception as e:
- print(f"AlphaRatioFilter skipped: {e}")
- stats["after_alpha_filter"] = len(storage)
-
- # ── 3. Near-duplicate removal ─────────────────────────────────────────────
- for module_path, class_name, kwargs in [
- ("dataflow.operators.dedup", "MinHashDeduplicator", {"threshold": 0.85}),
- ("dataflow.operators.dedup", "HashDeduplicator", {}),
- ]:
- try:
- mod = __import__(module_path, fromlist=[class_name])
- cls = getattr(mod, class_name)
- op = cls(**kwargs)
- op.run(storage=storage.step(), input_key="text")
- stats["after_dedup"] = len(storage)
- stats["operators_used"].append(class_name)
- print(f"After {class_name}: {stats['after_dedup']} records")
- break
- except Exception as e:
- print(f"{class_name} skipped: {e}")
- else:
- stats["after_dedup"] = len(storage)
-
- # ── Save output ───────────────────────────────────────────────────────────
+with open(INPUT) as fh:
+ raw_records = [json.loads(line) for line in fh if line.strip()]
+
+stats["input_count"] = len(raw_records)
+records_after_length = raw_records
+dataflow_ops_used = []
+python_ops_used = []
+
+if dataflow_ready:
try:
- storage.save(OUTPUT)
- print(f"DataFlow output saved to {OUTPUT}")
+ from dataflow.utils.storage import FileStorage
+ from dataflow.operators.general_text import CharNumberFilter, HashDeduplicateFilter, MinHashDeduplicateFilter
+
+ storage = FileStorage(
+ first_entry_file_name=INPUT,
+ cache_path="/tmp/gh-aw/agent/dataflow/cache",
+ )
+ storage.step() # step 0 = INPUT
+ CharNumberFilter(threshold=50).run(storage=storage, input_key="text")
+ storage.step() # step 1 = length-filter output
+ records_after_length = storage.read("dict")
+ dataflow_ops_used.append("CharNumberFilter")
+ print(f"After CharNumberFilter: {len(records_after_length)} records")
except Exception as e:
- # Fallback: write current data directly
- print(f"storage.save failed ({e}), writing manually")
- records = list(storage)
- with open(OUTPUT, "w") as fh:
- for r in records:
- fh.write(json.dumps(r, ensure_ascii=False) + "\n")
-
-# ── Fallback: lightweight Python pipeline (no DataFlow operators) ─────────────
+ stats["warnings"].append(f"CharNumberFilter failed: {e}")
+ print(f"CharNumberFilter failed: {e} — falling back to Python length filter")
+ records_after_length = [r for r in raw_records if 50 <= len(r.get('text', '')) <= 100_000]
+ python_ops_used.append("python_length_filter")
else:
- print("Running fallback pipeline (pure Python)...")
- import unicodedata, re
+ records_after_length = [r for r in raw_records if 50 <= len(r.get('text', '')) <= 100_000]
+ python_ops_used.append("python_length_filter")
- def alpha_ratio(text: str) -> float:
- letters = sum(1 for c in text if c.isalpha())
- return letters / max(len(text), 1)
+stats["after_length_filter"] = len(records_after_length)
- seen_hashes: set = set()
- kept: list = []
+records_after_alpha = [r for r in records_after_length if alpha_ratio(r.get("text", "")) >= 0.25]
+stats["after_alpha_filter"] = len(records_after_alpha)
+python_ops_used.append("python_alpha_ratio_filter")
- with open(INPUT) as fh:
- raw_records = [json.loads(line) for line in fh if line.strip()]
+records_after_dedup = records_after_alpha
+if dataflow_ready:
+ try:
+ storage.write(records_after_alpha) # writes step 2
+ storage.step() # step 2 = alpha-filter output
+ try:
+ MinHashDeduplicateFilter(threshold=0.85).run(storage=storage, input_key="text")
+ dataflow_ops_used.append("MinHashDeduplicateFilter")
+ except Exception as minhash_error:
+ stats["warnings"].append(f"MinHashDeduplicateFilter failed: {minhash_error}")
+ print(f"MinHashDeduplicateFilter failed: {minhash_error} — trying HashDeduplicateFilter")
+ HashDeduplicateFilter().run(storage=storage, input_key="text")
+ dataflow_ops_used.append("HashDeduplicateFilter")
+ storage.step() # step 3 = dedup output
+ records_after_dedup = storage.read("dict")
+ print(f"After DataFlow deduplication: {len(records_after_dedup)} records")
+ except Exception as e:
+ stats["warnings"].append(f"DataFlow dedup failed: {e}")
+ print(f"DataFlow dedup failed: {e} — falling back to Python exact-hash dedup")
+ records_after_dedup = exact_hash_dedup(records_after_alpha)
+ python_ops_used.append("python_exact_hash_dedup")
+else:
+ records_after_dedup = exact_hash_dedup(records_after_alpha)
+ python_ops_used.append("python_exact_hash_dedup")
- stats["input_count"] = len(raw_records)
+stats["after_dedup"] = len(records_after_dedup)
+stats["operators_used"] = dataflow_ops_used + python_ops_used
- for r in raw_records:
- text = r.get("text", "")
- # Length filter
- if not (50 <= len(text) <= 100_000):
- continue
- stats["after_length_filter"] = stats.get("after_length_filter", 0) + 1
- # Alpha-ratio filter
- if alpha_ratio(text) < 0.25:
- continue
- stats["after_alpha_filter"] = stats.get("after_alpha_filter", 0) + 1
- # Exact dedup
- h = hash(text[:500])
- if h in seen_hashes:
- continue
- seen_hashes.add(h)
- kept.append(r)
-
- stats["after_dedup"] = len(kept)
- stats["operators_used"].append("fallback_python_pipeline")
+if dataflow_ops_used and python_ops_used:
+ stats["execution_mode"] = "mixed"
+elif dataflow_ops_used:
+ stats["execution_mode"] = "dataflow"
+else:
+ stats["execution_mode"] = "fallback"
+stats["fallback_mode"] = stats["execution_mode"] == "fallback"
- with open(OUTPUT, "w") as fh:
- for r in kept:
- fh.write(json.dumps(r, ensure_ascii=False) + "\n")
- print(f"Fallback pipeline output: {len(kept)} records → {OUTPUT}")
+with open(OUTPUT, "w") as fh:
+ for record in records_after_dedup:
+ fh.write(json.dumps(record, ensure_ascii=False) + "\n")
+print(f"{stats['execution_mode']} pipeline output: {len(records_after_dedup)} records → {OUTPUT}")
# ── Write stats ───────────────────────────────────────────────────────────────
Path(STATS).write_text(json.dumps(stats, indent=2))
@@ -454,8 +494,23 @@ Run it:
Verify output:
```bash
-echo "Output records: $(wc -l < /tmp/gh-aw/agent/dataflow/output/dataset_clean.jsonl)"
-cat /tmp/gh-aw/agent/dataflow/output/pipeline_stats.json
+/tmp/gh-aw/python/venv/bin/python3 - << 'EOF'
+import json
+from pathlib import Path
+
+output = Path("/tmp/gh-aw/agent/dataflow/output/dataset_clean.jsonl")
+stats = json.loads(Path("/tmp/gh-aw/agent/dataflow/output/pipeline_stats.json").read_text())
+
+if not output.exists():
+ raise SystemExit("data-processing-error: dataset_clean.jsonl was not created")
+
+records = [json.loads(line) for line in output.read_text().splitlines() if line.strip()]
+print(f"Output records: {len(records)}")
+print(json.dumps(stats, indent=2))
+
+if stats.get("input_count", 0) > 0 and not records:
+ raise SystemExit("data-processing-error: filtering and deduplication produced an empty dataset")
+EOF
```
### Step 4: Upload Dataset Artifact
@@ -480,33 +535,10 @@ Then call the `upload_artifact` safe-output tool:
Record the returned artifact URL for use in the discussion report.
-### Step 5: Update Repo-Memory
-
-Save pipeline statistics for trend tracking across runs:
+### Step 5: Keep Repo-Memory Out of the Success Path
-```bash
-DATE=$(date '+%Y-%m-%d')
-RUN_ID="${GITHUB_RUN_ID}"
-STATS=$(cat /tmp/gh-aw/agent/dataflow/output/pipeline_stats.json)
-
-# Load existing history (or start fresh)
-HISTORY_FILE="/tmp/gh-aw/repo-memory/default/dataflow-runs.jsonl"
-mkdir -p "$(dirname "$HISTORY_FILE")"
-touch "$HISTORY_FILE"
-
-# Append this run
-python3 -c "
-import json, sys, os
-entry = {
- 'date': '$DATE',
- 'run_id': '$RUN_ID',
- **json.loads('''$STATS''')
-}
-with open('$HISTORY_FILE', 'a') as fh:
- fh.write(json.dumps(entry) + '\n')
-print('Run appended to history')
-"
-```
+Do not add repo-memory initialization or signed-commit setup to this workflow run.
+Report the dataset result directly from the generated artifact and `pipeline_stats.json`.
### Step 6: Compute Quality Breakdown
@@ -530,6 +562,10 @@ report = {
"avg_text_length_chars": round(avg_len, 1),
"operators_used": stats.get("operators_used", []),
"input_count": stats.get("input_count", 0),
+ "execution_mode": stats.get("execution_mode", "fallback"),
+ "fallback_mode": stats.get("fallback_mode", True),
+ "dataflow_version": stats.get("dataflow_version", ""),
+ "warnings": stats.get("warnings", []),
"retention_rate_pct": round(len(records) / max(stats.get("input_count", 1), 1) * 100, 1),
}
@@ -572,19 +608,34 @@ artifact_section = ""
if artifact_url:
artifact_section = f"\n### 📦 Dataset Artifact\n\n[Download dataset_clean.jsonl]({artifact_url})\n"
+execution_mode = quality.get("execution_mode", "fallback")
+mode_summary = {
+ "dataflow": "All filtering and deduplication stages used the current DataFlow operators.",
+ "mixed": "DataFlow handled the supported stages while deterministic Python fallbacks preserved the existing alpha-ratio contract.",
+ "fallback": "The workflow used the pure-Python fallback because DataFlow installation or API validation did not succeed.",
+}.get(execution_mode, "Execution mode unavailable.")
+
+warnings = quality.get("warnings", [])
+warning_section = ""
+if warnings:
+ warning_section = "### Runtime Warnings\n\n" + "\n".join(f"- {warning}" for warning in warnings) + "\n\n"
+
operators_str = ", ".join(quality.get("operators_used", ["none"])) or "none"
body = f"""### Summary
-Built a cleaned, deduplicated dataset from GitHub discussions and PRs using [OpenDCAI/DataFlow](https://github.com/OpenDCAI/DataFlow).
+Built a cleaned, deduplicated dataset from GitHub discussions and PRs.
| Metric | Value |
|--------|-------|
+| Execution mode | `{execution_mode}` |
| Input records | {quality.get("input_count", 0):,} |
| Output (clean) records | {quality.get("total_clean_records", 0):,} |
| Retention rate | {quality.get("retention_rate_pct", 0)}% |
| Average text length | {quality.get("avg_text_length_chars", 0):,.0f} chars |
-| DataFlow operators | `{operators_str}` |
+| Operators used | `{operators_str}` |
+
+{mode_summary}
### Records by Source
@@ -592,19 +643,20 @@ Built a cleaned, deduplicated dataset from GitHub discussions and PRs using [Ope
|--------|-------|
{by_source_rows}
-### DataFlow Pipeline
+{warning_section}### Pipeline
The pipeline applied these processing stages:
1. **Normalise** — Merged discussions and PRs into unified JSONL (`title + body` → `text`)
2. **Text length filter** — Kept records with 50–100,000 characters
-3. **Alpha-ratio filter** — Removed records with fewer than 25% alphabetic characters
-4. **Near-duplicate removal** — Eliminated near-identical records using MinHash or exact-hash
+3. **Alpha-ratio filter** — Removed records with fewer than 25% alphabetic characters using a deterministic Python stage
+4. **Near-duplicate removal** — Eliminated near-identical records using DataFlow MinHash or exact-hash fallback
{artifact_section}
### Pipeline Configuration
```yaml
-DataFlow package: open-dataflow (PyPI)
+Execution mode: {execution_mode}
+DataFlow package: open-dataflow=={quality.get("dataflow_version", "not-installed") or "not-installed"}
Source: https://github.com/OpenDCAI/DataFlow
Input: GitHub Discussions + merged PRs from {repo}
Output: JSONL — one record per item, text field for LLM use
@@ -635,37 +687,35 @@ Then emit:
```json
{
"type": "create_discussion",
- "title": "🌊 DataFlow Dataset Build Report — [DATE]",
+ "title": "DataFlow Dataset Build Report ([MODE]) — [DATE]",
"body": "[contents of /tmp/gh-aw/agent/discussion_body.md]",
"category": "audits"
}
```
-Replace `[DATE]` with today's ISO date and `[contents of ...]` with the actual text read
-from `/tmp/gh-aw/agent/discussion_body.md`.
+Replace `[MODE]` with the actual `execution_mode` value (`dataflow`, `mixed`, or `fallback`), replace `[DATE]` with today's ISO date, and replace `[contents of ...]` with the actual text read from `/tmp/gh-aw/agent/discussion_body.md`.
## Success Criteria
- ✅ Both discussions and PRs loaded and normalised into unified JSONL
-- ✅ DataFlow pipeline applied (with graceful fallback if operator API changes)
+- ✅ Execution mode reported honestly as `dataflow`, `mixed`, or `fallback`
- ✅ Clean JSONL artifact uploaded
- ✅ Quality breakdown computed (input/output counts, retention rate, source split)
- ✅ Discussion posted in `audits` category with full pipeline stats
-- ✅ Run metadata appended to repo-memory for trend tracking
## Edge Cases
### DataFlow API Changes
-DataFlow is actively developed — operator class names may change between releases.
-Always `import inspect` and list class names before use; adapt the import paths at runtime.
+Use the pinned, smoke-tested current API from `dataflow.operators.general_text`.
+Do not fall back to obsolete modules such as `dataflow.operators.filter` or `dataflow.operators.dedup`.
### No Data Available
If both input files are empty or missing, post a short discussion noting:
"No discussions or PRs found to process in this run."
### Insufficient Text
-If all records are filtered out by the length or alpha-ratio filter, relax thresholds:
-`min_len=20`, `min_ratio=0.1`, and skip dedup. Note the relaxed thresholds in the report.
+Do not relax the 50–100,000 character or 0.25 alphabetic-character thresholds.
+If those thresholds remove every record, fail with a `data-processing-error` instead of silently widening the filter.
## Usage
@@ -678,4 +728,5 @@ gh aw run dataflow-pr-discussion-dataset
The resulting JSONL dataset is suitable for:
- **LLM fine-tuning**: Supervised fine-tuning (SFT) datasets from real developer discussions
- **RAG indexing**: Embedding-ready clean text chunks from the repository's knowledge base
-- **Analytics**: Deduplicated corpus for topic modelling, sentiment analysis, clustering{{/if}}
\ No newline at end of file
+- **Analytics**: Deduplicated corpus for topic modelling, sentiment analysis, clustering
+{{/if}}
\ No newline at end of file
From 8a6f6ee9b21d9a9637efc4a2721cafae51bc1563 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Mon, 3 Aug 2026 13:25:26 +0000
Subject: [PATCH 4/4] =?UTF-8?q?fix:=20address=20review=20issues=20in=20dat?=
=?UTF-8?q?aflow=20pipeline=20=E2=80=94=20mode=20reporting,=20schema,=20up?=
=?UTF-8?q?per-bound=20filter,=20MinHash?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
- Fix mode reporting: exclude intentional Python alpha-ratio stage from
python_ops_used; mode is now computed from DataFlow-eligible stages only
(CharNumberFilter + HashDeduplicateFilter). Pure DataFlow success reports
as `dataflow`, partial fallback reports as `mixed`.
- Strip DataFlow internal label fields (char_number_filter_label,
minhash_deduplicated_label, hash_deduplicated_label) before writing
output so the dataset schema is uniform regardless of execution path.
- Add explicit upper-bound cap (100,000 chars) after DataFlow
CharNumberFilter output; CharNumberFilter(threshold=50) only enforces
the lower bound.
- Replace MinHashDeduplicateFilter with HashDeduplicateFilter in the
pipeline. The 1.0.10 MinHash implementation queries its LSH index
while an insertion_session is still open (datasketch inconsistency);
HashDeduplicateFilter is deterministic and already validated by the
smoke test. Remove MinHash import, instantiation, and reference from
the smoke test and selected_operators list.
Co-authored-by: gh-aw-bot <259018956+gh-aw-bot@users.noreply.github.com>
---
.../dataflow-pr-discussion-dataset.lock.yml | 5 +--
.../dataflow-pr-discussion-dataset.md | 44 +++++++++----------
2 files changed, 21 insertions(+), 28 deletions(-)
diff --git a/.github/workflows/dataflow-pr-discussion-dataset.lock.yml b/.github/workflows/dataflow-pr-discussion-dataset.lock.yml
index 22df9af951e..17e1cdd0bc0 100644
--- a/.github/workflows/dataflow-pr-discussion-dataset.lock.yml
+++ b/.github/workflows/dataflow-pr-discussion-dataset.lock.yml
@@ -1,4 +1,4 @@
-# gh-aw-metadata: {"schema_version":"v4","frontmatter_hash":"d3b50c56cd0267b90b93602f0056ffc5b1f8cf26c02f8595f7cda068afa2d5b8","body_hash":"9b06defff555f6d3244af112d46036857781188f0283d999c6e2b0ea27959d9b","strict":true,"agent_id":"copilot","engine_versions":{"copilot":"1.0.77"}}
+# gh-aw-metadata: {"schema_version":"v4","frontmatter_hash":"905bd34278b42598e78925c37617ee3eeaf01bb54a55cbf77711c0b0d84f818d","body_hash":"53eab465d720bf4387a7432ded7958daddbb5da3332d9c36f65737fcfff87673","strict":true,"agent_id":"copilot","engine_versions":{"copilot":"1.0.77"}}
# gh-aw-manifest: {"version":1,"secrets":["COPILOT_GITHUB_TOKEN","GH_AW_GITHUB_MCP_SERVER_TOKEN","GH_AW_GITHUB_TOKEN","GH_AW_OTEL_GRAFANA_AUTHORIZATION","GH_AW_OTEL_GRAFANA_ENDPOINT","GH_AW_OTEL_SENTRY_AUTHORIZATION","GH_AW_OTEL_SENTRY_ENDPOINT","GITHUB_TOKEN"],"actions":[{"repo":"actions/cache/restore","sha":"55cc8345863c7cc4c66a329aec7e433d2d1c52a9","version":"v6.1.0"},{"repo":"actions/cache/save","sha":"55cc8345863c7cc4c66a329aec7e433d2d1c52a9","version":"v6.1.0"},{"repo":"actions/checkout","sha":"3d3c42e5aac5ba805825da76410c181273ba90b1","version":"v7.0.1"},{"repo":"actions/download-artifact","sha":"3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c","version":"v8.0.1"},{"repo":"actions/github-script","sha":"3a2844b7e9c422d3c10d287c895573f7108da1b3","version":"v9.0.0"},{"repo":"actions/setup-node","sha":"820762786026740c76f36085b0efc47a31fe5020","version":"v7.0.0"},{"repo":"actions/setup-python","sha":"5fda3b95a4ea91299a34e894583c3862153e4b97","version":"v7.0.0"},{"repo":"actions/upload-artifact","sha":"043fb46d1a93c77aae656e7c1c64a875d1fc6a0a","version":"v7.0.1"},{"repo":"safedep/pmg","sha":"5ac0f275b83d9d5a9342c6aae977ec32fa330daa","version":"v1"}],"containers":[{"image":"ghcr.io/github/gh-aw-firewall/agent:0.27.43","digest":"sha256:04e2d1987a565000a8f114b89d806ae7a3864dd4f944be65275b28c93d8690e6","pinned_image":"ghcr.io/github/gh-aw-firewall/agent:0.27.43@sha256:04e2d1987a565000a8f114b89d806ae7a3864dd4f944be65275b28c93d8690e6"},{"image":"ghcr.io/github/gh-aw-firewall/api-proxy:0.27.43","digest":"sha256:d85f57975af5ea23af4996e41ed73fbc8f5b4a47402472bfe82e508f352cb0c1","pinned_image":"ghcr.io/github/gh-aw-firewall/api-proxy:0.27.43@sha256:d85f57975af5ea23af4996e41ed73fbc8f5b4a47402472bfe82e508f352cb0c1"},{"image":"ghcr.io/github/gh-aw-firewall/cli-proxy:0.27.43","digest":"sha256:65c45ea2967984d0024f3df61bc71335658a77ede96c8d9665da7a5f33a795ab","pinned_image":"ghcr.io/github/gh-aw-firewall/cli-proxy:0.27.43@sha256:65c45ea2967984d0024f3df61bc71335658a77ede96c8d9665da7a5f33a795ab"},{"image":"ghcr.io/github/gh-aw-firewall/squid:0.27.43","digest":"sha256:26be5e0b8c8f4c41c8a59126b29bb5d80b07253597472ded2a16bdd75abcbf9d","pinned_image":"ghcr.io/github/gh-aw-firewall/squid:0.27.43@sha256:26be5e0b8c8f4c41c8a59126b29bb5d80b07253597472ded2a16bdd75abcbf9d"},{"image":"ghcr.io/github/gh-aw-mcpg:v0.4.7","digest":"sha256:7545220a9aca134b71e51193ee0eaf4c50756ebf8fbd25a63ae7556e62815c00","pinned_image":"ghcr.io/github/gh-aw-mcpg:v0.4.7@sha256:7545220a9aca134b71e51193ee0eaf4c50756ebf8fbd25a63ae7556e62815c00"},{"image":"ghcr.io/github/gh-aw-node","digest":"sha256:0d9f1fb5fd6610c0ac1f5194a38e45a8a1e81f8a390d5142d8e4e6f26a4b3196","pinned_image":"ghcr.io/github/gh-aw-node@sha256:0d9f1fb5fd6610c0ac1f5194a38e45a8a1e81f8a390d5142d8e4e6f26a4b3196"},{"image":"ghcr.io/github/github-mcp-server:v1.8.0","digest":"sha256:d5a18c04b92714c309eb46a2305087e91a4dbd80420f6e462656699f95093520","pinned_image":"ghcr.io/github/github-mcp-server:v1.8.0@sha256:d5a18c04b92714c309eb46a2305087e91a4dbd80420f6e462656699f95093520"}]}
# This file was automatically generated by gh-aw. DO NOT EDIT. To debug this workflow, load the skill at https://github.com/github/gh-aw/blob/main/debug.md
#
@@ -761,7 +761,6 @@ jobs:
from dataflow.operators.general_text import (
CharNumberFilter,
HashDeduplicateFilter,
- MinHashDeduplicateFilter,
)
fixture_dir = Path("/tmp/gh-aw/agent/dataflow/smoke")
@@ -786,7 +785,6 @@ jobs:
if len(after_length) != 2:
raise RuntimeError(f"unexpected CharNumberFilter result: {len(after_length)}")
- _ = MinHashDeduplicateFilter(threshold=0.85)
storage.write(after_length)
storage.step()
HashDeduplicateFilter().run(storage=storage, input_key="text")
@@ -809,7 +807,6 @@ jobs:
"dataflow_version": "",
"selected_operators": [
"CharNumberFilter",
- "MinHashDeduplicateFilter",
"HashDeduplicateFilter",
],
"install_log": "$LOG",
diff --git a/.github/workflows/dataflow-pr-discussion-dataset.md b/.github/workflows/dataflow-pr-discussion-dataset.md
index e122a767c9f..971da22434a 100644
--- a/.github/workflows/dataflow-pr-discussion-dataset.md
+++ b/.github/workflows/dataflow-pr-discussion-dataset.md
@@ -58,7 +58,6 @@ steps:
from dataflow.operators.general_text import (
CharNumberFilter,
HashDeduplicateFilter,
- MinHashDeduplicateFilter,
)
fixture_dir = Path("/tmp/gh-aw/agent/dataflow/smoke")
@@ -83,7 +82,6 @@ steps:
if len(after_length) != 2:
raise RuntimeError(f"unexpected CharNumberFilter result: {len(after_length)}")
- _ = MinHashDeduplicateFilter(threshold=0.85)
storage.write(after_length)
storage.step()
HashDeduplicateFilter().run(storage=storage, input_key="text")
@@ -106,7 +104,6 @@ steps:
"dataflow_version": "",
"selected_operators": [
"CharNumberFilter",
- "MinHashDeduplicateFilter",
"HashDeduplicateFilter",
],
"install_log": "$LOG",
@@ -232,7 +229,7 @@ GitHub Discussions + PRs
│
▼
┌─────────────┐
- │ DataFlow │ Near-duplicate removal (MinHash or exact-hash fallback)
+ │ DataFlow │ Exact-hash deduplication (HashDeduplicateFilter)
│ Dedup │
└──────┬──────┘
│
@@ -251,7 +248,6 @@ Read `/tmp/gh-aw/agent/dataflow/output/dataflow_runtime.json` first.
- If `dataflow_ready` is `true`, use the smoke-tested current API:
- `dataflow.utils.storage.FileStorage`
- `dataflow.operators.general_text.CharNumberFilter`
- - `dataflow.operators.general_text.MinHashDeduplicateFilter`
- `dataflow.operators.general_text.HashDeduplicateFilter`
- If `dataflow_ready` is `false`, print the recorded warning and use the pure-Python fallback pipeline.
- Never use `AlphaWordsFilter`; it can trigger an implicit NLTK download.
@@ -346,9 +342,8 @@ Dataset cleaning pipeline:
1. Load JSONL records and the runtime status file
2. DataFlow char-length filter when the validated API is available
3. Deterministic Python alpha-ratio filter (to avoid NLTK downloads)
- 4. DataFlow near-deduplication (MinHash, then exact-hash fallback) when available
- 5. Pure-Python exact-hash fallback when DataFlow is unavailable
- 6. Save clean output and honest execution-mode statistics
+ 4. DataFlow exact-hash deduplication when available; pure-Python fallback otherwise
+ 5. Save clean output and honest execution-mode statistics
"""
import json
@@ -412,7 +407,7 @@ python_ops_used = []
if dataflow_ready:
try:
from dataflow.utils.storage import FileStorage
- from dataflow.operators.general_text import CharNumberFilter, HashDeduplicateFilter, MinHashDeduplicateFilter
+ from dataflow.operators.general_text import CharNumberFilter, HashDeduplicateFilter
storage = FileStorage(
first_entry_file_name=INPUT,
@@ -422,6 +417,7 @@ if dataflow_ready:
CharNumberFilter(threshold=50).run(storage=storage, input_key="text")
storage.step() # step 1 = length-filter output
records_after_length = storage.read("dict")
+ records_after_length = [r for r in records_after_length if len(r.get('text', '')) <= 100_000]
dataflow_ops_used.append("CharNumberFilter")
print(f"After CharNumberFilter: {len(records_after_length)} records")
except Exception as e:
@@ -437,23 +433,16 @@ stats["after_length_filter"] = len(records_after_length)
records_after_alpha = [r for r in records_after_length if alpha_ratio(r.get("text", "")) >= 0.25]
stats["after_alpha_filter"] = len(records_after_alpha)
-python_ops_used.append("python_alpha_ratio_filter")
records_after_dedup = records_after_alpha
if dataflow_ready:
try:
storage.write(records_after_alpha) # writes step 2
storage.step() # step 2 = alpha-filter output
- try:
- MinHashDeduplicateFilter(threshold=0.85).run(storage=storage, input_key="text")
- dataflow_ops_used.append("MinHashDeduplicateFilter")
- except Exception as minhash_error:
- stats["warnings"].append(f"MinHashDeduplicateFilter failed: {minhash_error}")
- print(f"MinHashDeduplicateFilter failed: {minhash_error} — trying HashDeduplicateFilter")
- HashDeduplicateFilter().run(storage=storage, input_key="text")
- dataflow_ops_used.append("HashDeduplicateFilter")
+ HashDeduplicateFilter().run(storage=storage, input_key="text")
storage.step() # step 3 = dedup output
records_after_dedup = storage.read("dict")
+ dataflow_ops_used.append("HashDeduplicateFilter")
print(f"After DataFlow deduplication: {len(records_after_dedup)} records")
except Exception as e:
stats["warnings"].append(f"DataFlow dedup failed: {e}")
@@ -465,19 +454,26 @@ else:
python_ops_used.append("python_exact_hash_dedup")
stats["after_dedup"] = len(records_after_dedup)
-stats["operators_used"] = dataflow_ops_used + python_ops_used
+stats["operators_used"] = dataflow_ops_used + ["python_alpha_ratio_filter"] + python_ops_used
-if dataflow_ops_used and python_ops_used:
- stats["execution_mode"] = "mixed"
-elif dataflow_ops_used:
+if dataflow_ops_used and not python_ops_used:
stats["execution_mode"] = "dataflow"
+elif dataflow_ops_used and python_ops_used:
+ stats["execution_mode"] = "mixed"
else:
stats["execution_mode"] = "fallback"
stats["fallback_mode"] = stats["execution_mode"] == "fallback"
+_DATAFLOW_INTERNAL_KEYS = frozenset({
+ "char_number_filter_label",
+ "minhash_deduplicated_label",
+ "hash_deduplicated_label",
+})
+
with open(OUTPUT, "w") as fh:
for record in records_after_dedup:
- fh.write(json.dumps(record, ensure_ascii=False) + "\n")
+ clean = {k: v for k, v in record.items() if k not in _DATAFLOW_INTERNAL_KEYS}
+ fh.write(json.dumps(clean, ensure_ascii=False) + "\n")
print(f"{stats['execution_mode']} pipeline output: {len(records_after_dedup)} records → {OUTPUT}")
# ── Write stats ───────────────────────────────────────────────────────────────
@@ -650,7 +646,7 @@ The pipeline applied these processing stages:
1. **Normalise** — Merged discussions and PRs into unified JSONL (`title + body` → `text`)
2. **Text length filter** — Kept records with 50–100,000 characters
3. **Alpha-ratio filter** — Removed records with fewer than 25% alphabetic characters using a deterministic Python stage
-4. **Near-duplicate removal** — Eliminated near-identical records using DataFlow MinHash or exact-hash fallback
+4. **Near-duplicate removal** — Eliminated near-identical records using DataFlow HashDeduplicateFilter or exact-hash fallback
{artifact_section}
### Pipeline Configuration