From 100d870a17b47fc70beb39cd0af7b1cccd5bbcbc Mon Sep 17 00:00:00 2001 From: 2233admin <2233admin@users.noreply.github.com> Date: Sat, 25 Jul 2026 14:17:32 +0800 Subject: [PATCH 1/4] feat: internalize DataFlow cleaning pipelines --- backend/api/v1/workflows.py | 7 +- backend/schemas/workflow.py | 4 +- backend/workflow/capability_projection.py | 80 + backend/workflow/compiler.py | 114 ++ backend/workflow/data_operators.py | 1056 +++++++++++++ backend/workflow/dataflow_compat.py | 1106 +++++++++++++ backend/workflow/demand_assembler.py | 190 ++- backend/workflow/external_importer.py | 86 + backend/workflow/node_registry.py | 4 + backend/workflow/opencli_hda_tracer.py | 96 +- backend/workflow/runtime_contracts.py | 15 + backend/workflow/runtime_registry.py | 105 ++ docs/dataflow-compatibility-matrix.md | 54 + docs/dataflow-operator-packs.md | 93 ++ frontend/components/flow/inspector.tsx | 109 +- frontend/lib/flow/store.ts | 95 +- frontend/lib/flow/types.ts | 4 +- frontend/lib/workflow/capabilities.ts | 43 + frontend/lib/workflow/node-catalog.ts | 68 +- frontend/lib/workflow/node-contracts.ts | 29 + frontend/lib/workflow/node-internals.ts | 28 + frontend/lib/workflow/parameter-interface.ts | 165 +- frontend/lib/workflow/schema.ts | 4 +- frontend/scripts/test-data-operator-nodes.mjs | 311 ++++ .../dataflow/test_pinned_compatibility.py | 589 +++++++ tests/compat/dataflow/test_upstream_oracle.py | 300 ++++ .../dataflow/pinned_f62aa134_golden.json | 606 +++++++ .../pinned_f62aa134_phase2_golden.json | 533 +++++++ .../test_dataflow_operator_pipeline_api.py | 1400 +++++++++++++++++ tests/integration/test_workflow_patch_api.py | 114 ++ tests/unit/test_data_operators.py | 370 +++++ tests/unit/test_dataflow_compat.py | 301 ++++ 32 files changed, 8048 insertions(+), 31 deletions(-) create mode 100644 backend/workflow/data_operators.py create mode 100644 backend/workflow/dataflow_compat.py create mode 100644 docs/dataflow-compatibility-matrix.md create mode 100644 docs/dataflow-operator-packs.md create mode 100644 frontend/scripts/test-data-operator-nodes.mjs create mode 100644 tests/compat/dataflow/test_pinned_compatibility.py create mode 100644 tests/compat/dataflow/test_upstream_oracle.py create mode 100644 tests/fixtures/dataflow/pinned_f62aa134_golden.json create mode 100644 tests/fixtures/dataflow/pinned_f62aa134_phase2_golden.json create mode 100644 tests/integration/test_dataflow_operator_pipeline_api.py create mode 100644 tests/unit/test_data_operators.py create mode 100644 tests/unit/test_dataflow_compat.py diff --git a/backend/api/v1/workflows.py b/backend/api/v1/workflows.py index 3228613..4ea4c58 100644 --- a/backend/api/v1/workflows.py +++ b/backend/api/v1/workflows.py @@ -169,9 +169,12 @@ async def draft_demand_workflow( async def import_external_runtime_workflow( body: workflow_schemas.WorkflowExternalImportRequest, ) -> ApiResponse[workflow_schemas.WorkflowPatchResponse]: - """Import LangGraph/LangChain graphs as OpenCLI Admin native nodes.""" + """Import supported external graphs as reviewable OpenCLI Admin native nodes.""" - return ApiResponse.ok(import_external_workflow(body)) + try: + return ApiResponse.ok(import_external_workflow(body)) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc @router.post( diff --git a/backend/schemas/workflow.py b/backend/schemas/workflow.py index 3886e0a..a0d21da 100644 --- a/backend/schemas/workflow.py +++ b/backend/schemas/workflow.py @@ -126,7 +126,7 @@ class WorkflowParameterInterfaceField(BaseModel): id: str = Field(..., min_length=1) label: str = Field(..., min_length=1) groupId: str = Field(..., min_length=1) - type: Literal["text", "textarea", "number", "slider", "select", "boolean", "tokens"] + type: Literal["text", "textarea", "json", "number", "slider", "select", "boolean", "tokens"] binding: WorkflowParameterBinding description: Optional[str] = None order: Optional[float] = None @@ -272,7 +272,7 @@ class WorkflowDemandDraftRequest(BaseModel): locale: Optional[str] = None -ExternalWorkflowRuntime = Literal["langgraph", "langchain"] +ExternalWorkflowRuntime = Literal["langgraph", "langchain", "dataflow"] class WorkflowExternalImportRequest(BaseModel): diff --git a/backend/workflow/capability_projection.py b/backend/workflow/capability_projection.py index 3c9590c..bdc039d 100644 --- a/backend/workflow/capability_projection.py +++ b/backend/workflow/capability_projection.py @@ -17,11 +17,13 @@ WorkflowNodeKind, WorkflowRuntimeCapability, ) +from backend.workflow.data_operators import list_data_operator_specs from backend.workflow.node_registry import WORKFLOW_PRIMITIVE_IDS from backend.workflow.opencli_adapter_nodes import get_opencli_adapter_node_summary from backend.workflow.runtime_contracts import runtime_io_contract_manifest from backend.workflow.runtime_registry import ( COLLECTION_OUTPUT_BINDING_ID, + DATA_OPERATOR_CATALOG_BINDINGS, DEMAND_DRAFT_BINDING_ID, EXTERNAL_TOOL_BINDING_ID, MERGE_BINDING_ID, @@ -268,6 +270,7 @@ def _catalog_capabilities() -> list[WorkflowRuntimeCapability]: probes=["typed_port_contract_registered"], ), ), + *_data_operator_capabilities(), _blocked_catalog( "intelligence.agent.summary", "LLM Summary", @@ -538,6 +541,83 @@ def _catalog_capabilities() -> list[WorkflowRuntimeCapability]: ] +def _data_operator_capabilities() -> list[WorkflowRuntimeCapability]: + specs_by_kind: dict[str, list[object]] = {} + for spec in list_data_operator_specs(): + specs_by_kind.setdefault(spec.kind, []).append(spec) + + rows: list[WorkflowRuntimeCapability] = [] + for catalog_id, binding_id in DATA_OPERATOR_CATALOG_BINDINGS.items(): + operator_kind = catalog_id.rsplit(".", 1)[-1] + specs = sorted( + specs_by_kind.get(operator_kind, []), + key=lambda spec: spec.operator_id, + ) + if not specs: + continue + operators = [ + { + "id": spec.operator_id, + "operatorId": spec.operator_id, + "kind": spec.kind, + "label": spec.label, + "description": spec.description, + "pack": spec.pack_id, + "packId": spec.pack_id, + "version": spec.pack_version, + "packVersion": spec.pack_version, + "status": "runnable", + "readiness": "ready", + "configKeys": list(spec.config_keys), + } + for spec in specs + ] + rows.append( + _capability( + id=catalog_id, + label=f"Data {operator_kind.title()}", + surface="catalog", + status="runnable", + backend_available=True, + kind="agent", + capability="normalize", + provider="workflow", + runtime_binding=binding_id, + reason="Registered versioned data operators execute on Record Candidates.", + tags=["data", "operator", operator_kind], + source="backend.workflow.data_operators", + manifest={ + **_manifest( + schema=f"capability.data.{operator_kind}.v1", + input_ports=[_port("in", "recordCandidate[]")], + output_ports=[_port("out", "recordCandidate[]")], + runtime_binding=binding_id, + trace_events=[ + "partial:outputItemCount", + "completed", + "failed", + ], + probes=["data_operator_registry"], + ), + "operatorIds": [operator["id"] for operator in operators], + "operators": operators, + "packs": sorted({operator["packId"] for operator in operators}), + "params": list( + dict.fromkeys( + key for spec in specs for key in spec.config_keys + ) + ), + "artifacts": [ + "recordCandidate[]", + "metrics", + "rejectedCandidateIds", + ], + }, + ) + ) + return rows + + def _manifest( *, schema: str, diff --git a/backend/workflow/compiler.py b/backend/workflow/compiler.py index bee753d..0939363 100644 --- a/backend/workflow/compiler.py +++ b/backend/workflow/compiler.py @@ -20,6 +20,7 @@ WorkflowProjectNode, WorkflowRuntimePreview, ) +from backend.workflow.data_operators import resolve_data_operator from backend.workflow.hda_templates import materialize_hda_templates from backend.workflow.node_registry import ( forbidden_node_definition_keys, @@ -29,6 +30,7 @@ INTERNAL_ID_SEPARATOR = "::" MAX_NODE_PATH_DEPTH = 4 +_LEGACY_DATA_OPERATOR_PACK_VERSION = "1.0.0" @dataclass(frozen=True) @@ -64,6 +66,22 @@ class _PortContract: [_PortContract("in", "input", "recordCandidate[]")], [_PortContract("out", "output", "recordCandidate[]")], ), + "intelligence.data.generate": ( + [_PortContract("in", "input", "recordCandidate[]")], + [_PortContract("out", "output", "recordCandidate[]")], + ), + "intelligence.data.filter": ( + [_PortContract("in", "input", "recordCandidate[]")], + [_PortContract("out", "output", "recordCandidate[]")], + ), + "intelligence.data.evaluate": ( + [_PortContract("in", "input", "recordCandidate[]")], + [_PortContract("out", "output", "recordCandidate[]")], + ), + "intelligence.data.refine": ( + [_PortContract("in", "input", "recordCandidate[]")], + [_PortContract("out", "output", "recordCandidate[]")], + ), "intelligence.flow.merge": ( [ _PortContract("in1", "input", "recordCandidate[]"), @@ -316,6 +334,7 @@ def _validate_project(project: WorkflowProject) -> list[WorkflowCompileError]: ) errors.extend(_validate_node_origin(node, ["nodes", node.id])) + errors.extend(_validate_data_operator_node(node, ["nodes", node.id])) errors.extend(_validate_typed_edges(project.nodes, project.edges, path_prefix=["edges"])) errors.extend(_cycle_errors(project)) @@ -366,6 +385,100 @@ def _validate_node_origin( return errors +def _validate_data_operator_node( + node: WorkflowProjectNode, + path_prefix: list[str], +) -> list[WorkflowCompileError]: + catalog_id = _read_string((node.ui or {}).get("catalogId")) + prefix = "intelligence.data." + expected_kind = ( + catalog_id.removeprefix(prefix) + if catalog_id and catalog_id.startswith(prefix) + else None + ) + operator_id = _read_string(node.params.get("operatorId")) + if expected_kind not in {"generate", "filter", "evaluate", "refine"}: + if "operatorId" not in node.params: + return [] + return [ + WorkflowCompileError( + code="data_operator_catalog_required", + message=( + f'Workflow node "{node.id}" declares params.operatorId but does ' + "not reference a registered intelligence.data catalog node" + ), + node_id=node.id, + path=[*path_prefix, "ui", "catalogId"], + ) + ] + if operator_id is None: + return [ + WorkflowCompileError( + code="missing_data_operator_id", + message=f'Workflow data node "{node.id}" requires params.operatorId', + node_id=node.id, + path=[*path_prefix, "params", "operatorId"], + ) + ] + + pack_version_provided = "packVersion" in node.params + requested_pack_version = _read_string(node.params.get("packVersion")) + if pack_version_provided and requested_pack_version is None: + return [ + WorkflowCompileError( + code="unsupported_data_operator_version", + message=( + f'Workflow data node "{node.id}" requires params.packVersion ' + "to be a non-empty string when provided" + ), + node_id=node.id, + path=[*path_prefix, "params", "packVersion"], + ) + ] + resolved_pack_version = ( + requested_pack_version or _LEGACY_DATA_OPERATOR_PACK_VERSION + ) + spec = resolve_data_operator(operator_id, resolved_pack_version) + if spec is None: + if resolve_data_operator(operator_id) is not None: + return [ + WorkflowCompileError( + code="unsupported_data_operator_version", + message=( + f'Workflow data node "{node.id}" references unsupported ' + f'version "{resolved_pack_version}" of operator ' + f'"{operator_id}"' + ), + node_id=node.id, + path=[*path_prefix, "params", "packVersion"], + ) + ] + return [ + WorkflowCompileError( + code="unknown_data_operator", + message=( + f'Workflow data node "{node.id}" references unknown operator ' + f'"{operator_id}"' + ), + node_id=node.id, + path=[*path_prefix, "params", "operatorId"], + ) + ] + if spec.kind != expected_kind: + return [ + WorkflowCompileError( + code="data_operator_kind_mismatch", + message=( + f'Data operator "{operator_id}" has kind "{spec.kind}", but node ' + f'"{node.id}" requires "{expected_kind}"' + ), + node_id=node.id, + path=[*path_prefix, "params", "operatorId"], + ) + ] + return [] + + def _validate_typed_edges( nodes: list[WorkflowProjectNode], edges: list[WorkflowProjectEdge], @@ -1004,6 +1117,7 @@ def _validate_package_internals( internal_path_prefix, ) ) + errors.extend(_validate_data_operator_node(internal_node, internal_path_prefix)) if _is_structural_container(internal_node): errors.extend( _validate_package_internals( diff --git a/backend/workflow/data_operators.py b/backend/workflow/data_operators.py new file mode 100644 index 0000000..9d8367e --- /dev/null +++ b/backend/workflow/data_operators.py @@ -0,0 +1,1056 @@ +"""Versioned deterministic data-preparation operator registry.""" + +from __future__ import annotations + +import copy +import hashlib +import html +import json +import re +import string +import unicodedata +from collections.abc import Callable +from dataclasses import dataclass +from html.parser import HTMLParser +from typing import Any, Literal + +from backend.workflow.dataflow_compat import ( + COMPAT_EXECUTORS, + COMPAT_OPERATOR_DEFINITIONS, +) + +DataOperatorKind = Literal["generate", "filter", "evaluate", "refine"] +_Executor = Callable[ + [list[dict[str, Any]], dict[str, Any]], + tuple[list[dict[str, Any]], dict[str, Any], list[str]], +] + + +@dataclass(frozen=True) +class DataOperatorSpec: + id: str + kind: DataOperatorKind + pack_id: str + pack_version: str + label: str + description: str + config_keys: tuple[str, ...] = () + + @property + def operator_id(self) -> str: + return self.id + + def to_manifest(self) -> dict[str, object]: + return { + "operatorId": self.id, + "kind": self.kind, + "packId": self.pack_id, + "packVersion": self.pack_version, + "label": self.label, + "description": self.description, + "configKeys": list(self.config_keys), + "inputPort": "recordCandidate[]", + "outputPort": "recordCandidate[]", + "deterministic": True, + } + + +@dataclass(frozen=True) +class DataOperatorPack: + id: str + version: str + operators: tuple[DataOperatorSpec, ...] + + @property + def pack_id(self) -> str: + return self.id + + def to_manifest(self) -> dict[str, object]: + return { + "packId": self.id, + "version": self.version, + "operators": [spec.to_manifest() for spec in self.operators], + } + + +@dataclass(frozen=True) +class DataOperatorResult: + operator_id: str + pack_id: str + pack_version: str + items: list[dict[str, Any]] + metrics: dict[str, Any] + rejected_count: int + rejected_candidate_ids: list[str] + + def to_details(self) -> dict[str, object]: + return { + "operatorId": self.operator_id, + "packId": self.pack_id, + "packVersion": self.pack_version, + "inputItemCount": self.metrics["inputItemCount"], + "outputItemCount": self.metrics["outputItemCount"], + "metrics": dict(self.metrics), + "rejectedCandidateIds": list(self.rejected_candidate_ids), + } + + +_CORE_PACK_ID = "builtin.core-data" +_TEXT_PACK_ID = "builtin.text-cleaning" +_DATASET_PACK_ID = "builtin.dataset-preparation" +_LEGACY_VERSION = "1.0.0" + +_LEGACY_SPECS = ( + DataOperatorSpec( + "core.generate.instruction-pairs", + "generate", + _CORE_PACK_ID, + _LEGACY_VERSION, + "Instruction pairs", + "Build instruction/output records from normalized candidates.", + ("instructionField", "responseField", "instructionTemplate"), + ), + DataOperatorSpec( + "core.filter.quality", + "filter", + _CORE_PACK_ID, + _LEGACY_VERSION, + "Quality filter", + "Reject records below deterministic text-quality thresholds.", + ( + "fields", + "requiredFields", + "textField", + "minChars", + "maxChars", + "minLength", + "maxLength", + "minQuality", + "blocklist", + ), + ), + DataOperatorSpec( + "core.evaluate.quality", + "evaluate", + _CORE_PACK_ID, + _LEGACY_VERSION, + "Quality evaluation", + "Attach deterministic quality scores and signals.", + ("fields", "minLength", "maxLength"), + ), + DataOperatorSpec( + "core.refine.text", + "refine", + _CORE_PACK_ID, + _LEGACY_VERSION, + "Text refine", + "Normalize whitespace in selected text fields.", + ("fields", "lowercase", "unicodeForm", "redactEmail", "redactPhone"), + ), + DataOperatorSpec( + "text.clean", + "refine", + _TEXT_PACK_ID, + _LEGACY_VERSION, + "Text clean", + "Apply a configurable sequence of deterministic text cleaners.", + ("fields", "operations", "replacement"), + ), + DataOperatorSpec( + "text.rule-filter", + "filter", + _TEXT_PACK_ID, + _LEGACY_VERSION, + "Text rule filter", + "Filter text by length, vocabulary, symbols, and blocklist rules.", + ( + "fields", + "minChars", + "maxChars", + "minWords", + "maxWords", + "minSentences", + "maxSymbolRatio", + "minUniqueWordRatio", + "blocklist", + ), + ), + DataOperatorSpec( + "text.deduplicate", + "filter", + _TEXT_PACK_ID, + _LEGACY_VERSION, + "Text deduplicate", + "Keep the first exact or near-duplicate record.", + ("fields", "mode", "maxHammingDistance"), + ), + DataOperatorSpec( + "text.statistics", + "evaluate", + _TEXT_PACK_ID, + _LEGACY_VERSION, + "Text statistics", + "Attach character, word, sentence, and lexical-diversity statistics.", + ("fields", "outputField"), + ), + DataOperatorSpec( + "data.project", + "refine", + _DATASET_PACK_ID, + _LEGACY_VERSION, + "Data project", + "Select, rename, coalesce, and cast normalized fields.", + ("select", "rename", "coalesce", "casts"), + ), + DataOperatorSpec( + "data.chunk", + "generate", + _DATASET_PACK_ID, + _LEGACY_VERSION, + "Data chunk", + "Split text into deterministic overlapping character chunks.", + ("field", "chunkSize", "overlap"), + ), + DataOperatorSpec( + "data.qa-extract", + "generate", + _DATASET_PACK_ID, + _LEGACY_VERSION, + "QA extract", + "Expand embedded question-answer pairs into grounded candidates.", + ("pairsField", "contextField"), + ), + DataOperatorSpec( + "data.training-format", + "refine", + _DATASET_PACK_ID, + _LEGACY_VERSION, + "Training format", + "Project question-answer candidates to Alpaca or ShareGPT records.", + ("format", "instructionField", "inputField", "outputField", "resultField"), + ), +) + + +def _compat_spec(definition: dict[str, Any]) -> DataOperatorSpec: + def read(camel: str, snake: str | None = None) -> Any: + return definition.get(camel, definition.get(snake or camel)) + + kind = read("kind") + if kind not in {"generate", "filter", "evaluate", "refine"}: + raise ValueError("Compatibility data operator kind is invalid") + config_keys = read("configKeys", "config_keys") or () + if not isinstance(config_keys, (list, tuple)) or any( + not isinstance(key, str) or not key for key in config_keys + ): + raise ValueError("Compatibility data operator configKeys are invalid") + values = { + "id": read("operatorId", "operator_id"), + "pack_id": read("packId", "pack_id"), + "pack_version": read("packVersion", "pack_version"), + "label": read("label"), + "description": read("description"), + } + if any(not isinstance(value, str) or not value for value in values.values()): + raise ValueError("Compatibility data operator definition is incomplete") + return DataOperatorSpec( + values["id"], + kind, + values["pack_id"], + values["pack_version"], + values["label"], + values["description"], + tuple(config_keys), + ) + + +_SPECS = (*_LEGACY_SPECS, *tuple(_compat_spec(item) for item in COMPAT_OPERATOR_DEFINITIONS)) +_SPEC_BY_KEY = {(spec.operator_id, spec.pack_version): spec for spec in _SPECS} +if len(_SPEC_BY_KEY) != len(_SPECS): + raise ValueError("Duplicate data operator id and packVersion") + + +def _packs() -> tuple[DataOperatorPack, ...]: + keys = dict.fromkeys((spec.pack_id, spec.pack_version) for spec in _SPECS) + return tuple( + DataOperatorPack( + pack_id, + version, + tuple( + spec + for spec in _SPECS + if spec.pack_id == pack_id and spec.pack_version == version + ), + ) + for pack_id, version in keys + ) + + +_PACKS = _packs() + + +def resolve_data_operator( + operator_id: str, pack_version: str | None = None +) -> DataOperatorSpec | None: + if pack_version is not None: + return _SPEC_BY_KEY.get((operator_id, pack_version)) + return _SPEC_BY_KEY.get((operator_id, _LEGACY_VERSION)) + + +def list_data_operator_specs() -> tuple[DataOperatorSpec, ...]: + return _SPECS + + +def list_data_operator_packs() -> tuple[DataOperatorPack, ...]: + return _PACKS + + +def get_data_operator_pack( + pack_id: str, pack_version: str | None = None +) -> DataOperatorPack | None: + version = pack_version or _LEGACY_VERSION + return next( + (pack for pack in _PACKS if pack.id == pack_id and pack.version == version), + None, + ) + + +def execute_data_operator( + operator_id: str, + items: list[dict[str, Any]], + config: dict[str, Any] | None = None, + *, + pack_version: str | None = None, +) -> DataOperatorResult: + spec = resolve_data_operator(operator_id, pack_version) + if spec is None: + suffix = f" at packVersion {pack_version}" if pack_version else "" + raise ValueError(f"Unknown data operator: {operator_id}{suffix}") + if not isinstance(items, list) or any(not isinstance(item, dict) for item in items): + raise ValueError("Data operator items must be a list of objects") + resolved_config = dict(config or {}) + unknown = sorted(set(resolved_config) - set(spec.config_keys)) + if unknown: + raise ValueError(f"Unsupported config for {operator_id}: {', '.join(unknown)}") + executor = _EXECUTORS.get((operator_id, spec.pack_version)) + if executor is None: + raise ValueError( + f"Data operator executor is unavailable: {operator_id} at {spec.pack_version}" + ) + output, operator_metrics, rejected = executor( + copy.deepcopy(items), resolved_config + ) + rejected_count = int(operator_metrics.pop("rejectedInputCount", len(rejected))) + metrics = { + "inputItemCount": len(items), + "outputItemCount": len(output), + "rejectedItemCount": rejected_count, + "inputCount": len(items), + "outputCount": len(output), + "rejectedCount": rejected_count, + **operator_metrics, + } + return DataOperatorResult( + operator_id=operator_id, + pack_id=spec.pack_id, + pack_version=spec.pack_version, + items=output, + metrics=metrics, + rejected_count=rejected_count, + rejected_candidate_ids=rejected, + ) + + +def _instruction_pairs(items, config): + instruction_field = _string_config(config, "instructionField", "title") + response_field = _string_config(config, "responseField", "content") + template = _string_config(config, "instructionTemplate", "{title}") + output, rejected = [], [] + for index, item in enumerate(items): + response = _field_text(item, response_field) + if not response: + _record_rejection(rejected, item) + continue + title = _field_text(item, instruction_field) or "the source content" + try: + instruction = template.format(title=title, content=response) + except (IndexError, KeyError, ValueError) as error: + raise ValueError( + "instructionTemplate may only use {title} and {content}" + ) from error + normalized = _normalized(item) + normalized.update( + { + "instruction": instruction, + "input": "", + "output": response, + "response": response, + } + ) + output.append(_with_normalized(item, normalized, index=index)) + return output, { + "generatedPairCount": len(output), + "rejectedInputCount": len(items) - len(output), + }, rejected + + +def _quality_filter(items, config): + text_field = _string_config(config, "textField", "content") + fields = _fields_config(config) if "fields" in config else [text_field] + minimum = _aliased_length(config, "minChars", "minLength", default=1) + maximum = _aliased_optional_length(config, "maxChars", "maxLength") + if maximum is not None and maximum < minimum: + raise ValueError("maxChars must be greater than or equal to minChars") + min_quality = _number_config(config, "minQuality", 0.0, 0.0, 1.0) + required = _string_list(config.get("requiredFields", []), "requiredFields") + blocklist = [ + term.casefold() + for term in _string_list(config.get("blocklist", []), "blocklist") + ] + output, rejected = [], [] + for item in items: + text = _combined_text(item, fields) + score, _ = _quality(text) + data = _normalized(item) + reject = any(data.get(field) in (None, "") for field in required) + reject = reject or any(term in text.casefold() for term in blocklist) + reject = reject or len(text) < minimum + reject = reject or (maximum is not None and len(text) > maximum) + reject = reject or score < min_quality + if reject: + _record_rejection(rejected, item) + else: + output.append(item) + return output, { + "minimumQuality": min_quality, + "rejectedInputCount": len(items) - len(output), + }, rejected + + +def _quality_evaluate(items, config): + fields = _fields_config(config) + minimum = _int_config(config, "minLength", 20, 0) + maximum = _optional_int_config(config, "maxLength", 1) + if maximum is not None and maximum < minimum: + raise ValueError("maxLength must be greater than or equal to minLength") + output, scores = [], [] + for index, item in enumerate(items): + text = _combined_text(item, fields) + score, signals = _quality(text) + signals["lengthWithinBounds"] = len(text) >= minimum and ( + maximum is None or len(text) <= maximum + ) + normalized = _normalized(item) + normalized.update({"qualityScore": score, "qualitySignals": signals}) + output.append(_with_normalized(item, normalized, index=index)) + scores.append(score) + average = round(sum(scores) / len(scores), 4) if scores else 0.0 + return output, {"averageQuality": average}, [] + + +def _text_refine(items, config): + fields = _fields_config(config) + lowercase = _bool_config(config, "lowercase", False) + form = _string_config(config, "unicodeForm", "NFKC") + if form not in {"NFC", "NFD", "NFKC", "NFKD"}: + raise ValueError("unicodeForm must be one of NFC, NFD, NFKC, or NFKD") + redact_email = _bool_config(config, "redactEmail", False) + redact_phone = _bool_config(config, "redactPhone", False) + output, changed = [], 0 + for index, item in enumerate(items): + normalized = _normalized(item) + for field in fields: + value = normalized.get(field) + if not isinstance(value, str): + continue + refined = _collapse_whitespace(unicodedata.normalize(form, value)) + refined = refined.lower() if lowercase else refined + if redact_email: + refined = re.sub( + r"[\w.+-]+@[\w.-]+\.[A-Za-z]{2,}", + "[REDACTED_EMAIL]", + refined, + ) + if redact_phone: + refined = re.sub( + r"(? max_chars), + ("minWords", stats["wordCount"] < min_words), + ("maxWords", max_words is not None and stats["wordCount"] > max_words), + ("minSentences", stats["sentenceCount"] < min_sentences), + ("maxSymbolRatio", stats["symbolRatio"] > max_symbol), + ("minUniqueWordRatio", stats["uniqueWordRatio"] < min_unique), + ("blocklist", any(value in text.casefold() for value in blocklist)), + ) + failed = [name for name, condition in checks if condition] + if failed: + _record_rejection(rejected, item) + for name in failed: + hits[name] = hits.get(name, 0) + 1 + else: + output.append(item) + return output, { + "ruleHits": hits, + "rejectedInputCount": len(items) - len(output), + }, rejected + + +def _deduplicate(items, config): + fields = _fields_config(config) + mode = _string_config(config, "mode", "exact") + if mode not in {"exact", "simhash"}: + raise ValueError("text.deduplicate mode must be exact or simhash") + distance = _int_config(config, "maxHammingDistance", 3, 0) + if distance > 64: + raise ValueError("maxHammingDistance must be between 0 and 64") + output, rejected, exact_seen, fingerprints = [], [], set(), [] + for item in items: + text = _collapse_whitespace(_combined_text(item, fields)).casefold() + if mode == "exact": + duplicate = text in exact_seen + exact_seen.add(text) + else: + fingerprint = _simhash(text) + duplicate = any( + (fingerprint ^ previous).bit_count() <= distance + for previous in fingerprints + ) + fingerprints.append(fingerprint) + if duplicate: + _record_rejection(rejected, item) + else: + output.append(item) + return output, { + "duplicateCount": len(items) - len(output), + "mode": mode, + "rejectedInputCount": len(items) - len(output), + }, rejected + + +def _text_statistics(items, config): + fields = _fields_config(config) + output_field = _string_config(config, "outputField", "dataflowStatistics") + output, total_words = [], 0 + for index, item in enumerate(items): + stats = _statistics(_combined_text(item, fields)) + normalized = _normalized(item) + normalized[output_field] = stats + output.append(_with_normalized(item, normalized, index=index)) + total_words += stats["wordCount"] + average = round(total_words / len(items), 4) if items else 0.0 + return output, {"averageWordCount": average}, [] + + +def _project(items, config): + selected = _string_list(config.get("select", []), "select") + rename = _string_mapping(config.get("rename", {}), "rename") + coalesce = _string_list_mapping(config.get("coalesce", {}), "coalesce") + casts = _string_mapping(config.get("casts", {}), "casts") + output = [] + for index, item in enumerate(items): + source = _normalized(item) + projected = ( + {key: source[key] for key in selected if key in source} + if selected else source + ) + for target, candidates in coalesce.items(): + for candidate in candidates: + if source.get(candidate) not in (None, ""): + projected[target] = source[candidate] + break + for old, new in rename.items(): + if old in projected: + projected[new] = projected.pop(old) + for field, cast in casts.items(): + if field in projected: + projected[field] = _cast(projected[field], cast) + output.append(_with_normalized(item, projected, index=index)) + return output, { + "projectedFieldCount": sum(len(_normalized(item)) for item in output) + }, [] + + +def _chunk(items, config): + field = _string_config(config, "field", "content") + size = _int_config(config, "chunkSize", 1000, 1) + overlap = _int_config(config, "overlap", 0, 0) + if overlap >= size: + raise ValueError("data.chunk overlap must be smaller than chunkSize") + output, rejected, step = [], [], size - overlap + for index, item in enumerate(items): + text = _field_text(item, field) + if not text: + _record_rejection(rejected, item) + continue + chunks = [] + for start in range(0, len(text), step): + chunks.append(text[start : start + size]) + if start + size >= len(text): + break + source_id = _candidate_id(item, index) + for chunk_index, value in enumerate(chunks): + normalized = _normalized(item) + normalized.update({ + field: value, + "chunkIndex": chunk_index, + "chunkCount": len(chunks), + "sourceCandidateId": source_id, + }) + output.append(_with_normalized( + item, normalized, index=index, derived_key=f"chunk:{chunk_index}" + )) + return output, { + "generatedChunkCount": len(output), + "rejectedInputCount": sum(not _field_text(item, field) for item in items), + }, rejected + + +def _qa_extract(items, config): + pairs_field = _string_config(config, "pairsField", "qaPairs") + context_field = _string_config(config, "contextField", "content") + output, rejected, rejected_count = [], [], 0 + for index, item in enumerate(items): + pairs = _field_value(item, pairs_field) + pairs = _field_value(item, f"extra_{pairs_field}") if pairs is None else pairs + if not isinstance(pairs, list): + _record_rejection(rejected, item) + rejected_count += 1 + continue + generated = 0 + source_id = _candidate_id(item, index) + context = _field_text(item, context_field) + for pair_index, pair in enumerate(pairs): + if not isinstance(pair, dict): + continue + question = _first_text(pair, ("question", "q", "instruction")) + answer = _first_text(pair, ("answer", "a", "output", "response")) + if not question or not answer: + continue + normalized = _normalized(item) + normalized.update({ + "question": question, "answer": answer, "context": context, + "sourceCandidateId": source_id, "sourceRefs": _source_refs(item), + "citations": list(pair.get("citations", [])) + if isinstance(pair.get("citations"), list) else [], + }) + output.append(_with_normalized( + item, normalized, index=index, derived_key=f"qa:{pair_index}" + )) + generated += 1 + if not generated: + _record_rejection(rejected, item) + rejected_count += 1 + return output, { + "extractedPairCount": len(output), + "rejectedInputCount": rejected_count, + }, rejected + + +def _training_format(items, config): + format_name = _string_config(config, "format", "alpaca").casefold() + if format_name not in {"alpaca", "sharegpt"}: + raise ValueError("data.training-format format must be alpaca or sharegpt") + instruction_field = _string_config(config, "instructionField", "question") + input_field = _string_config(config, "inputField", "context") + output_field = _string_config(config, "outputField", "answer") + result_field = _string_config(config, "resultField", "trainingData") + output, rejected = [], [] + for index, item in enumerate(items): + instruction = _field_text(item, instruction_field) + response = _field_text(item, output_field) + if not instruction or not response: + _record_rejection(rejected, item) + continue + input_text = _field_text(item, input_field) + formatted = ( + {"instruction": instruction, "input": input_text, "output": response} + if format_name == "alpaca" + else [ + {"from": "human", "value": f"{instruction}\n\n{input_text}".strip()}, + {"from": "gpt", "value": response}, + ] + ) + normalized = _normalized(item) + normalized[result_field] = formatted + output.append(_with_normalized(item, normalized, index=index)) + return output, { + "formattedRecordCount": len(output), + "format": format_name, + "rejectedInputCount": len(items) - len(output), + }, rejected + + +def _quality(text): + stats = _statistics(text) + nonempty = bool(text.strip()) + length_score = min(stats["characterCount"] / 200, 1.0) + symbol_score = 1.0 - stats["symbolRatio"] + score = round( + (0.1 if nonempty else 0) + + 0.4 * length_score + + 0.3 * stats["uniqueWordRatio"] + + 0.2 * symbol_score, + 4, + ) + return min(score, 1.0), { + "nonempty": nonempty, + "lengthScore": round(length_score, 4), + "lexicalDiversity": stats["uniqueWordRatio"], + "symbolScore": round(symbol_score, 4), + } + + +def _statistics(text): + words = re.findall(r"\b[\w'-]+\b", text.casefold(), re.UNICODE) + non_space = [c for c in text if not c.isspace()] + symbols = [c for c in non_space if not c.isalnum() and not c.isalpha()] + sentences = len(re.findall(r"[.!?。!?]+(?:\s|$)", text)) + sentences = 1 if text.strip() and not sentences else sentences + return { + "characterCount": len(text), + "wordCount": len(words), + "sentenceCount": sentences, + "uniqueWordRatio": round(len(set(words)) / len(words), 4) if words else 0.0, + "symbolRatio": round(len(symbols) / len(non_space), 4) if non_space else 0.0, + } + + +def _simhash(text): + tokens = re.findall(r"\w+", text.casefold(), re.UNICODE) or [text] + vector = [0] * 64 + for token in tokens: + digest = int.from_bytes(hashlib.sha256(token.encode()).digest()[:8], "big") + for bit in range(64): + vector[bit] += 1 if digest & (1 << bit) else -1 + return sum(1 << bit for bit, weight in enumerate(vector) if weight >= 0) + + +def _normalized(item): + value = item.get("normalizedData") + return copy.deepcopy(value) if isinstance(value, dict) else {} + + +def _with_normalized(item, normalized, *, index, derived_key=None): + updated = copy.deepcopy(item) + updated["normalizedData"] = normalized + updated["contentHash"] = _stable_hash(normalized) + updated["candidateId"] = ( + _candidate_id(item, index) + if derived_key is None + else _derived_candidate_id(item, index, derived_key, normalized) + ) + return updated + + +def _candidate_id(item, index): + value = item.get("candidateId") + if isinstance(value, str) and value: + return value + identity = { + "contentHash": item.get("contentHash"), + "normalizedData": item.get("normalizedData"), + "raw": item.get("raw"), + "index": index, + } + return f"candidate:{_stable_hash(identity)[:24]}" + + +def _record_rejection(rejected, item): + value = item.get("candidateId", item.get("id")) + if isinstance(value, str) and value: + rejected.append(value) + elif isinstance(value, (int, float)) and not isinstance(value, bool): + rejected.append(str(value)) + + +def _derived_candidate_id(item, index, key, normalized): + value = f"{_candidate_id(item, index)}|{key}|{_stable_hash(normalized)}" + return f"candidate:{hashlib.sha256(value.encode()).hexdigest()[:24]}" + + +def _stable_hash(value): + payload = json.dumps( + value, sort_keys=True, ensure_ascii=False, separators=(",", ":"), default=str + ) + return hashlib.sha256(payload.encode()).hexdigest() + + +def _field_value(item, field): + for value in (item.get("normalizedData"), item.get("raw"), item): + if isinstance(value, dict) and field in value: + return value[field] + return None + + +def _field_text(item, field): + value = _field_value(item, field) + return value.strip() if isinstance(value, str) else "" + + +def _combined_text(item, fields): + return "\n".join(value for field in fields if (value := _field_text(item, field))) + + +def _source_refs(item): + refs = [copy.deepcopy(value) for value in item.get("lineage", []) if isinstance(value, dict)] + url = _field_text(item, "url") + return [*refs, *([{"url": url}] if url else [])] + + +def _first_text(value, keys): + return next( + ( + value[key].strip() + for key in keys + if isinstance(value.get(key), str) and value[key].strip() + ), + "", + ) + + +def _collapse_whitespace(value): + return re.sub(r"\s+", " ", value).strip() + + +def _fields_config(config): + return _string_list(config.get("fields", ["title", "content"]), "fields") + + +def _string_config(config, key, default): + value = config.get(key, default) + if not isinstance(value, str) or not value: + raise ValueError(f"{key} must be a non-empty string") + return value + + +def _bool_config(config, key, default): + value = config.get(key, default) + if not isinstance(value, bool): + raise ValueError(f"{key} must be a boolean") + return value + + +def _int_config(config, key, default, minimum): + value = config.get(key, default) + if not isinstance(value, int) or isinstance(value, bool) or value < minimum: + raise ValueError(f"{key} must be an integer >= {minimum}") + return value + + +def _optional_int_config(config, key, minimum): + return None if key not in config or config[key] is None else _int_config( + config, key, minimum, minimum + ) + + +def _aliased_length(config, primary, alias, *, default): + if primary in config and alias in config and config[primary] != config[alias]: + raise ValueError(f"{primary} and {alias} must match when both are provided") + return _int_config(config, primary if primary in config else alias, default, 0) + + +def _aliased_optional_length(config, primary, alias): + if primary in config and alias in config and config[primary] != config[alias]: + raise ValueError(f"{primary} and {alias} must match when both are provided") + return _optional_int_config(config, primary if primary in config else alias, 1) + + +def _number_config(config, key, default, minimum, maximum): + value = config.get(key, default) + if not isinstance(value, (int, float)) or isinstance(value, bool): + raise ValueError(f"{key} must be a number") + resolved = float(value) + if not minimum <= resolved <= maximum: + raise ValueError(f"{key} must be between {minimum} and {maximum}") + return resolved + + +def _string_list(value, name): + if not isinstance(value, list) or any( + not isinstance(item, str) or not item for item in value + ): + raise ValueError(f"{name} must be a list of non-empty strings") + return list(value) + + +def _string_mapping(value, name): + if not isinstance(value, dict) or any( + not isinstance(key, str) + or not key + or not isinstance(item, str) + or not item + for key, item in value.items() + ): + raise ValueError(f"{name} must map non-empty strings to non-empty strings") + return dict(value) + + +def _string_list_mapping(value, name): + if not isinstance(value, dict): + raise ValueError(f"{name} must be an object") + return {key: _string_list(item, f"{name}.{key}") for key, item in value.items()} + + +def _cast(value, cast): + try: + if cast == "string": + return str(value) + if cast == "integer": + return int(value) + if cast == "number": + return float(value) + if cast == "boolean": + if isinstance(value, bool): + return value + if isinstance(value, str) and value.casefold() in {"true", "1", "yes"}: + return True + if isinstance(value, str) and value.casefold() in {"false", "0", "no"}: + return False + raise ValueError + if cast == "json": + return json.loads(value) if isinstance(value, str) else value + except (TypeError, ValueError) as error: + raise ValueError(f"Cannot cast value to {cast}") from error + raise ValueError(f"Unsupported cast: {cast}") + + +_LEGACY_EXECUTORS: dict[str, _Executor] = { + "core.generate.instruction-pairs": _instruction_pairs, + "core.filter.quality": _quality_filter, + "core.evaluate.quality": _quality_evaluate, + "core.refine.text": _text_refine, + "text.clean": _text_clean, + "text.rule-filter": _rule_filter, + "text.deduplicate": _deduplicate, + "text.statistics": _text_statistics, + "data.project": _project, + "data.chunk": _chunk, + "data.qa-extract": _qa_extract, + "data.training-format": _training_format, +} +_EXECUTORS: dict[tuple[str, str], _Executor] = { + **{ + (operator_id, _LEGACY_VERSION): executor + for operator_id, executor in _LEGACY_EXECUTORS.items() + }, + **COMPAT_EXECUTORS, +} + +__all__ = [ + "DataOperatorKind", + "DataOperatorPack", + "DataOperatorResult", + "DataOperatorSpec", + "execute_data_operator", + "get_data_operator_pack", + "list_data_operator_packs", + "list_data_operator_specs", + "resolve_data_operator", +] diff --git a/backend/workflow/dataflow_compat.py b/backend/workflow/dataflow_compat.py new file mode 100644 index 0000000..8fc4bf2 --- /dev/null +++ b/backend/workflow/dataflow_compat.py @@ -0,0 +1,1106 @@ +"""Pinned, dependency-free compatibility for a small DataFlow operator subset.""" + +from __future__ import annotations + +import base64 +import copy +import hashlib +import re +import string +import unicodedata +import zlib +from collections.abc import Callable, Mapping +from dataclasses import dataclass +from datetime import datetime +from functools import cache +from typing import Any + +DATAFLOW_COMPAT_SHA = "f62aa1349e0ff14cb737a4cbda1945d04fde85bb" +COMPAT_PACK_ID = "builtin.text-cleaning" +COMPAT_PACK_VERSION = "1.1.0" + +_Executor = Callable[ + [list[dict[str, Any]], dict[str, Any]], + tuple[list[dict[str, Any]], dict[str, Any], list[str]], +] + + +@dataclass(frozen=True) +class DataFlowInvocation: + operator_id: str + pack_id: str + pack_version: str + kind: str + config: dict[str, Any] + source_id: str + + def to_params(self) -> dict[str, Any]: + return { + "operatorId": self.operator_id, + "packId": self.pack_id, + "packVersion": self.pack_version, + "config": copy.deepcopy(self.config), + } + + +COMPAT_OPERATOR_DEFINITIONS: tuple[dict[str, object], ...] = ( + { + "operatorId": "text.clean", + "kind": "refine", + "packId": COMPAT_PACK_ID, + "packVersion": COMPAT_PACK_VERSION, + "label": "Text clean", + "description": "Apply pinned DataFlow-compatible deterministic text cleaners.", + "configKeys": ["fields", "operations", "htmlEntities"], + "inputPort": "recordCandidate[]", + "outputPort": "recordCandidate[]", + "deterministic": True, + }, + { + "operatorId": "text.rule-filter", + "kind": "filter", + "packId": COMPAT_PACK_ID, + "packVersion": COMPAT_PACK_VERSION, + "label": "Text rule filter", + "description": "Apply pinned DataFlow-compatible deterministic text rules.", + "configKeys": ["fields", "rules"], + "inputPort": "recordCandidate[]", + "outputPort": "recordCandidate[]", + "deterministic": True, + }, + { + "operatorId": "text.deduplicate", + "kind": "filter", + "packId": COMPAT_PACK_ID, + "packVersion": COMPAT_PACK_VERSION, + "label": "Text deduplicate", + "description": "Keep the first record for each pinned DataFlow-compatible hash.", + "configKeys": [ + "fields", + "hashFunction", + "mode", + "nGram", + "diffSize", + "outputKey", + ], + "inputPort": "recordCandidate[]", + "outputPort": "recordCandidate[]", + "deterministic": True, + }, +) + +_SOURCE_PREFIX = f"dataflow@{DATAFLOW_COMPAT_SHA}::" +_ALIASES = { + "RemoveExtraSpaces": ( + "dataflow.operators.general_text.refine.remove_extra_spaces_refiner." + "RemoveExtraSpacesRefiner" + ), + "Lowercase": ( + "dataflow.operators.general_text.refine.lowercase_refiner.LowercaseRefiner" + ), + "HtmlUrlRemover": ( + "dataflow.operators.general_text.refine.html_url_remover_refiner." + "HtmlUrlRemoverRefiner" + ), + "HtmlEntity": ( + "dataflow.operators.general_text.refine.html_entity_refiner.HtmlEntityRefiner" + ), + "RemoveEmoji": ( + "dataflow.operators.general_text.refine.remove_emoji_refiner.RemoveEmojiRefiner" + ), + "RemoveNumber": ( + "dataflow.operators.general_text.refine.remove_number_refiner.RemoveNumberRefiner" + ), + "RemovePunctuation": ( + "dataflow.operators.general_text.refine.remove_punctuation_refiner." + "RemovePunctuationRefiner" + ), + "RemoveRepetitionsPunctuation": ( + "dataflow.operators.general_text.refine." + "remove_repetitions_punctuation_refiner.RemoveRepetitionsPunctuationRefiner" + ), + "ContentNull": ( + "dataflow.operators.general_text.filter.rule_based_filter.ContentNullFilter" + ), + "WordNumber": ( + "dataflow.operators.general_text.filter.word_number_filter.WordNumberFilter" + ), + "SentenceNumber": ( + "dataflow.operators.general_text.filter.rule_based_filter.SentenceNumberFilter" + ), + "CharNumber": ( + "dataflow.operators.general_text.filter.rule_based_filter.CharNumberFilter" + ), + "UniqueWords": ( + "dataflow.operators.general_text.filter.rule_based_filter.UniqueWordsFilter" + ), + "HashDeduplicate": ( + "dataflow.operators.general_text.filter.hash_deduplicate_filter." + "HashDeduplicateFilter" + ), + "ColonEnd": ( + "dataflow.operators.general_text.filter.rule_based_filter.ColonEndFilter" + ), + "LineEndWithEllipsis": ( + "dataflow.operators.general_text.filter.rule_based_filter." + "LineEndWithEllipsisFilter" + ), + "SymbolWordRatio": ( + "dataflow.operators.general_text.filter.rule_based_filter." + "SymbolWordRatioFilter" + ), + "AlphaWords": ( + "dataflow.operators.general_text.filter.rule_based_filter.AlphaWordsFilter" + ), + "HtmlEntityFilter": ( + "dataflow.operators.general_text.filter.rule_based_filter.HtmlEntityFilter" + ), + "IDCard": ( + "dataflow.operators.general_text.filter.rule_based_filter.IDCardFilter" + ), + "NoPunc": ( + "dataflow.operators.general_text.filter.rule_based_filter.NoPuncFilter" + ), + "SpecialCharacter": ( + "dataflow.operators.general_text.filter.rule_based_filter." + "SpecialCharacterFilter" + ), + "Watermark": ( + "dataflow.operators.general_text.filter.rule_based_filter.WatermarkFilter" + ), + "MeanWordLength": ( + "dataflow.operators.general_text.filter.rule_based_filter." + "MeanWordLengthFilter" + ), + "StopWord": ( + "dataflow.operators.general_text.filter.rule_based_filter.StopWordFilter" + ), + "CurlyBracket": ( + "dataflow.operators.general_text.filter.rule_based_filter.CurlyBracketFilter" + ), + "CapitalWords": ( + "dataflow.operators.general_text.filter.rule_based_filter.CapitalWordsFilter" + ), + "LoremIpsum": ( + "dataflow.operators.general_text.filter.rule_based_filter.LoremIpsumFilter" + ), + "LineStartWithBulletpoint": ( + "dataflow.operators.general_text.filter.rule_based_filter." + "LineStartWithBulletpointFilter" + ), + "LineWithJavascript": ( + "dataflow.operators.general_text.filter.rule_based_filter." + "LineWithJavascriptFilter" + ), + "TextNormalization": ( + "dataflow.operators.general_text.refine.text_normalization_refiner." + "TextNormalizationRefiner" + ), + "RemoveImageRefs": ( + "dataflow.operators.general_text.refine.remove_image_ref_refiner." + "RemoveImageRefsRefiner" + ), + "Blocklist": ( + "dataflow.operators.general_text.filter.blocklist_filter.BlocklistFilter" + ), + "NgramHashDeduplicate": ( + "dataflow.operators.general_text.filter.ngramhash_deduplicate_filter." + "NgramHashDeduplicateFilter" + ), +} +DATAFLOW_ALIAS_SOURCE_IDS = { + alias: _SOURCE_PREFIX + dotted_name for alias, dotted_name in _ALIASES.items() +} +_ALIAS_BY_SOURCE_ID = { + source_id: alias for alias, source_id in DATAFLOW_ALIAS_SOURCE_IDS.items() +} + +_CLEAN_OPERATIONS = { + "RemoveExtraSpaces": "removeExtraSpaces", + "Lowercase": "lowercase", + "HtmlUrlRemover": "htmlUrlRemover", + "HtmlEntity": "htmlEntity", + "RemoveEmoji": "removeEmoji", + "RemoveNumber": "removeNumber", + "RemovePunctuation": "removePunctuation", + "RemoveRepetitionsPunctuation": "removeRepetitionsPunctuation", + "TextNormalization": "textNormalization", + "RemoveImageRefs": "removeImageRefs", +} +_DEFAULT_HTML_ENTITIES = [ + "nbsp", + "lt", + "gt", + "amp", + "quot", + "apos", + "hellip", + "ndash", + "mdash", + "lsquo", + "rsquo", + "ldquo", + "rdquo", +] + +# Pinned OpenDCAI/DataFlow blocklists, Apache-2.0 with the upstream repository. +# The compressed bytes are the exact git blobs at DATAFLOW_COMPAT_SHA. +_BLOCKLIST_ASSETS = { + "en": ( + "af851ecef1d5f212caba17339b12ac39cc2fef7d78c74876f67237644fcee8bd", + "eNplV02S9agR3HMKVt55MRPhAyGpJNFCwPDz1Op7jLe+oo/gzELvfWM7orsqQSAVRdbP+337bTa/282XUO1vdu7ZuLmkls6Udx+8My64yZ3O7qnZnOZDGqfq4aLNPkvwUYyLLkBgfdx6BaLIUnePxaXKngIW1cr/Nz57nHfjektW8EE/K37g5CYJLi4Et51ca1IG/up+FsAQ7Oa2BxT3ugc8/HzAiDEIfx1UNx8P6u/puE0l1QE6dZFpLCsSbhtkw7k4iO4QGlObKx+dFPjoiKbLTEs9zSQu0lZVlfolxc5B9XsYfNZH2O2Cb7BdfkG/2SmoFUSFq6riI9L9ZQya19lzov2+LCHpjgafqhRM8y121gchxUWsm5tP8T1K0f7f/PWVJtX2A+7UYXFyOHfoYi+3rrhB4DPDjWZKgd/G5/Aqtwk1T5q4H4IPUrvtDNcbuPuKFhy49BylR8HVfr7fjwOOhg5huR8gzb78NLASaupxs8oiogdUeq63pgKnFxoEqI9nd0qwIBcRqU5dQ2/U154KH5SMLykpYf0YfkZ7mlNwMBRskakv1cy8LaqCq/2Schjol3LW1ib4ILbhPlWkwqUhDSa4M3MEk6WsoKLRG5rVh3PKBV/aIgJP8ROEcypxHCXBTxRYC2qc2WOun/w/yWloeLe9NVb1+IlL4GYWVw5sWvQ4CFJFA4hk2/aSXFP8gXH5GLIgpCACKL/gpUEmKeX+C/Yg3uILLjx7EON6j6qjL3/MkrbNC3x04zRj8Avff5m/P/AzGWawBRondUqYBxb/TXjy0ykegm9rcllSn4JYzG1vnAWMK8/u/CbeUm67g89meS3OiAMRbst8JfO8eyNfbu5hbHoSlCqEu1RcTTPSY8cHV2Qk/G+4gFVAeMjAaQm6m6Dp+HSwpf7R4RreGiZgv1k93MEhhBRmpQeOyTrWIprsKvj6rphxusKa8ZB8orAkP1myarJ71FjA3MUHywZ/OKaUMfjghnxevNmYFmkEgX0Qrki+zSbRNxeq2byLbeQYBhZyylsj3rIZdWVDRrAXSGM2EGr+VlUFarGLO7ElHUenCiDbkx4wSktOfElKWqAACjIyvL6Bj0LZ8zDn7zXD51svyewMQTqFYOiywEIEOdEAEpvzZkeVe+5zV95AJboAKSFygJM5fFfBTubDMpwLVSaENz57Wbilw5srwoU8oqd9nMFsKETXjBSK834xH6d1BfABDkV6IBoAHLlJ1U5cnnV+c8ie1L8AmKHo58d89W2r5vBIlbjdyR2d+mBmUXAb1IxpojkoQA25zBYYVRys+szoeNubpXnC+RO3lx2i1gT/Ejo4JKQxZ5jCTqdEghJ7Chx7ivkfOp8oWR38bfIL/tc8A+lEUYXLnG0lIT2ePqwQteKZQzjmVP1Yx+g9E00leXGyM3VccVrtS9honMXyqgH6utrFs84SAvGj0Z0ohFCXQ/kwUUAbSHzlByM6ViU2QdmYsMOfGdRwTMLR54zcM1Q1sa6XCutPWI+JvggFKzeurkk08eb2R423pBnB0Gu9DRhX1A6GBLot265B7qozW78t9ZcrJhXlOtRtspNFU7AAHh4iNmZaajzFQ3QTRcGzSkYqQb6Dc2EIXsqLRL2Yhb7TUpo9SZ/pda6FRtYeQHVw95SwAFSpvcDgnZTOibd9aqTkFG8uA0grRVTRmCyQhF3UJxk7O9hA+EYlqkhDooXLO9YWho11AaUE7aZH6KhlbZ/xuoknVjf+0cWtlELpT1PctgtaFGjPZKU9iFY1CL6CCrYXmRvKZUGtRkSCvZdmluJPpgooraIo82CgC+cvZJlWGC7/sNUzwKop7DpI35OfM9UtrAgodCj4+EadkQ3qvAcSuM7wadJEXgXUN7wM/KMBoUoUN0VH1RgqvMfkVt3ROC529I/v0fBF3bWcQPuJaVtvlmIKILIifUFCANcdR1R/1EMQ6xUNS5Mbm9kN1b/BmJMARRV1qsZ0sWPW9YnF9kcGwPsyu9cZWkVICa3fo+2Z0Cdh9EyiUVnYTFfgHg/DjJNZKhS8dWYfDeTBltAno1Xfaitg2LKrwCvwA8CjZx1BU3toyBZXUp9etOnAMS8tm6Yx0cFRzFsY4NpQ2A3aMyTfER/ocUQqc9jodsZcP6cfpHhnEGSLxW+ipgnygiuxEO7V7htCg7Cpg1FtN7THPlpsS8wWeJJqpRmoZ0HpiXNGhAQOOQ22tD6N2gbw6Io4aBe403CKY0g0bJ9EYZ/UYZAFkc01dHuuB5KvQZDiHM5qSPZP1/ZiVEAxX1pNn+blHWIOakI3lApAYqt98UfXCzmqyHtvuqWXR10yEXUc6tXDyxmsP8wlTX82QduFXalRTyGHs46PBvtKJSwk/cWrZozyxi5ExBGEPbMrK1ua72/8fZvbJW9uYQf5+blw+3W9zU96n+rf//rnn+Y/EFJuDw==", + ), + "zh": ( + "a1d9aa037c8b039ef3b40148b3364ce2ca62ce4a955b7082a16ad99f6cbd1bc0", + "eNpNlslyIjkQhu96mIno6EufPYc5zWleZo7s+9bgMTsGYxZjm8Vgs5uHabTUrR9h/szKAkdQ9f2pUkmpVKaKb9//UN++28hWnTdpuxvZdBIqS0Y0QsqtFsDEhcL6tWQbMXXets1mAyzpOu8eCfpxRTCpE2P54JIh8z5ReOzF8+q8T7heFbjTgw/8ropH9I1xLTCGKcGvWNoWsjr+7JpZadKtuaig97hG8J5euOHYhbc6sjezWza+PtHxsZ5+sLpOTsb2Sx+anBE00PgMne/pwTEwyiNW0+iXrtPodSBTm3mb3kW9kpKnOtywqYUSLwEvdFA6kfMifaCphxWlk/D2mVF8ZaxXBJeBlWro3ZBgaxPslNK5uG1Mlc6Had+wlzp/d0M3Hh5YY6T8O13YRKWLOb2c68+p9/AIo097CLjIDHjWu4zSpaheFICC3mG+0tZOMUKvwzF8CRQC4yszy4uyqY0oSphN6mpwn8eYHiyBO6+G1TyOAodYmWWPFILGa0WWXHbJJjJsUN/Bhi4MCmwlIXzFob3sJRpu6IZ0CMjhYMUOsRG8N0xJg1gAv8ZbQxsyrLDvw4pXfwCWX+aEcVV+GkP5ixYjcGpUo1LRo74bnADs4rgWeCZvPmV4nQwd33PbM9L8wJhics45Xvcl+/QsYVYdpecxjsw8Zt8+gaLUJlQQa1Yc63kxyJ950UXS9q1Myps02BtWbjViJXvCirJaL8K2dSC4CFlwYJFTej2mY0Fv3/7l5Ni+nTdRTmMocotBFSZKOl0TxDfSM1J+WAPld72EVQyElZSJjKWDae1MYc7qdf/nPyL++luEeAGFQJoWDYNqOew4qxjUw4RGdFFdALo+Jphq10WPpKh2TTTu0m/A8rzJ8DvJBHkBh0XBH5Ms/6AbhYBQZJj/6sqkSufdDvhpKnmqmwF2tooCRtfzvq9MGi+lgYqNTAlev4OlKZM5mmJZmSwfEpA8dbEsSzPlsp6flD+mQA49MfzjVwxy1lfk7BdHxKAavvSRuMNA3AXXsWWzRPnDQV2Hk40hhewCqLDM7afL95SpDuQgJ0WeM3hAVuTI5WngSHVg0ltlWim7LwI7jjKARKXJ8Hud4KdMO2JH6P8+4Q1Y83XeLgh6mCSYfp4xXxHs6EhgzwAUqtmevFARdaYTcddHVpxCSABlw2MJHj6WVCoALcummspmq6gOAhwl2OxBqtkW0qb9rq5fNltp6XGUqpPUU4JHrLRoeVKhtn6k7zGhumLUe0qCa0dhymT7Muax3j5pIATpvMsR9M+lkpABMh6U1LmvyHkcAbofJiAcPAapMsEe2kARV5OKwkUGFBUUAEXBZRbmHnPg64houiWer9p4CVjc0I0yHeA5VgfZWPfevO4lDCrm2sk3+lI/7jilwnSntcvPFEUWl+mulZcs0Cnk1ZCBa4IrFYCVvk0STGrN6G4YsybB5biLF64womPlNTuULF67q3NVAs0O0EcMcNEcIz5X3tOHfPpEwTlf8Sus+C1W8um7GNzHj7aEGuDVvBx4zE3vB90ogQn44BNKIwYCIMcvQJ9/OYK9z3sdO9E+kcIZ/3u8LCuvPneFO/WrU8Vpvae2jPKWT+p3b9mB9YGjp9ygZiT8PbrfTYI/Z7q8QS17xXv8BWkEitvOhw9pM9kC5ZmeDZFnaRfHobko6NNRuV3b1hfKX69XTYii3LvGwPZeXKKr/gdO9p9p", + ), +} + +# NLTK's English stopwords corpus is an unpinned upstream runtime dependency. +# Freeze the observed 198-word corpus locally so compatibility never downloads data. +_NLTK_STOPWORDS_ASSET = ( + "f6d005956f407dbc6ea32e5ff0c7e8e6f71488d3239b9023efdc7fc139d6375b", + "eNo1U1uSxCAI/OciOZeZmJGtKFs+1vL2240zVdAQRJ4mSDhtdOJflHD3WCW8g5aNDSfUn0dClgCtXOAloUZycTjg1SR0OSPoFUajxOEZb6vUtbyBj01gn35kPcmJzOeSFwK/bDzXVyDgJZc6FwdaDBSbQ9noVsa+jJb9PaGOSmsMryR3nIIq5K6W5R61J7SYwkUuDriVUD+4OPg3pkEoG7eJMRPs8bgIGIrHipW22hzic/MMH5rJ26D4RvMqiqt6Q+CyHlkwXC3dBA6K9OrZlUTH7m4d0bQ7M5geqO1nYDU4zJKD5ChZ36mXj0CIzLlng1OGZ9lI+wJ5mBIj+neEvRiogqlOMbEbBC6gVwQ8S+BgPj4bzs0hPn+RWhc8Idgxf+Ru0gLqagm7JSBHS/zm8NqeXvNBtcSdf0T5Svenxm6bgRhtYKFdOoMCXPMZoSqy1o2NIjt8yoPKK1wVsDkuB5TjYkdZx/ZYTNu5t57M3auNN5IbyGSUC62O0vWR8SvwRetLJp7R5DOa+xnNCEKC6e0CKw0fKBvdjdkmG5qsc3qdMym6BT7UjZwByKIMpvh/JtYz/dnPPbz5nd2SZYOM7ETcWNwadcRevr2113d/FQ6KDijmH2pQY+Y=", +) + + +def translate_dataflow_alias( + source_id: str, + init_config: Mapping[str, Any] | None = None, + run_config: Mapping[str, Any] | None = None, +) -> DataFlowInvocation: + """Translate one exact, SHA-locked upstream class into a native invocation.""" + + alias = _ALIAS_BY_SOURCE_ID.get(source_id) + if alias is None: + _unsupported("source id is not in the pinned compatibility allowlist") + init = _mapping(init_config, "init_config") + run = _mapping(run_config, "run_config") + + if alias in _CLEAN_OPERATIONS: + allowed_init = {"html_entities"} if alias == "HtmlEntity" else set() + _only_keys(init, allowed_init, alias) + _only_keys(run, {"input_key"}, alias) + field = _required_string(run, "input_key") + config: dict[str, Any] = { + "fields": [field], + "operations": [_CLEAN_OPERATIONS[alias]], + } + if alias == "HtmlEntity": + config["htmlEntities"] = _string_list( + init.get("html_entities", _DEFAULT_HTML_ENTITIES), "html_entities" + ) + return _invocation("text.clean", "refine", config, source_id) + + if alias == "HashDeduplicate": + _only_keys(init, {"hash_func"}, alias) + _only_keys(run, {"input_key", "input_keys", "output_key"}, alias) + hash_function = init.get("hash_func", "md5") + if hash_function not in {"md5", "sha256", "xxh3"}: + _unsupported("HashDeduplicate hash_func must be md5, sha256, or xxh3") + has_one = run.get("input_key") is not None + has_many = run.get("input_keys") is not None + if has_one == has_many: + _unsupported("HashDeduplicate requires exactly one of input_key or input_keys") + if has_many: + fields = _string_list(run["input_keys"], "input_keys") + if len(fields) < 2: + _unsupported("HashDeduplicate input_keys requires at least two fields") + else: + fields = [_required_string(run, "input_key")] + return _invocation( + "text.deduplicate", + "filter", + { + "fields": fields, + "hashFunction": hash_function, + "outputKey": _optional_string( + run, "output_key", "minhash_deduplicated_label" + ), + }, + source_id, + ) + + if alias == "NgramHashDeduplicate": + _only_keys(init, {"n_gram", "hash_func", "diff_size"}, alias) + _only_keys(run, {"input_key", "input_keys", "output_key"}, alias) + hash_function = init.get("hash_func", "md5") + if hash_function not in {"md5", "sha256"}: + _unsupported( + "NgramHashDeduplicate hash_func must be md5 or sha256; " + "xxh3 is unavailable without an explicit dependency" + ) + has_one = run.get("input_key") is not None + has_many = run.get("input_keys") is not None + if has_one == has_many: + _unsupported( + "NgramHashDeduplicate requires exactly one of input_key or input_keys" + ) + if has_many: + fields = _string_list(run["input_keys"], "input_keys") + if len(fields) < 2: + _unsupported("NgramHashDeduplicate input_keys requires at least two fields") + else: + fields = [_required_string(run, "input_key")] + return _invocation( + "text.deduplicate", + "filter", + { + "fields": fields, + "mode": "ngramHash", + "nGram": _integer(init.get("n_gram", 3), "n_gram"), + "hashFunction": hash_function, + "diffSize": _integer(init.get("diff_size", 1), "diff_size"), + "outputKey": _optional_string( + run, "output_key", "minhash_deduplicated_label" + ), + }, + source_id, + ) + + _only_keys(run, {"input_key", "output_key"}, alias) + field = _required_string(run, "input_key") + if alias == "ContentNull": + _only_keys(init, set(), alias) + rule = { + "type": "contentNull", + "outputKey": _optional_string( + run, "output_key", "content_null_filter_label" + ), + } + elif alias == "WordNumber": + _only_keys(init, {"min_words", "max_words"}, alias) + rule = { + "type": "wordNumber", + "min": _integer(init.get("min_words", 20), "min_words"), + "max": _integer(init.get("max_words", 100000), "max_words"), + "outputKey": _optional_string( + run, "output_key", "word_number_filter_label" + ), + } + elif alias == "SentenceNumber": + _only_keys(init, {"min_sentences", "max_sentences"}, alias) + rule = { + "type": "sentenceNumber", + "min": _integer(init.get("min_sentences", 3), "min_sentences"), + "max": _integer(init.get("max_sentences", 7500), "max_sentences"), + "outputKey": _optional_string( + run, "output_key", "sentence_number_filter_label" + ), + } + elif alias == "CharNumber": + _only_keys(init, {"threshold"}, alias) + rule = { + "type": "charNumber", + "threshold": _integer(init.get("threshold", 100), "threshold"), + "outputKey": _optional_string( + run, "output_key", "char_number_filter_label" + ), + } + elif alias == "UniqueWords": + _only_keys(init, {"threshold"}, alias) + rule = { + "type": "uniqueWords", + "threshold": _number(init.get("threshold", 0.1), "threshold"), + "outputKey": _optional_string(run, "output_key", "unique_words_filter"), + } + else: + rule = _translate_phase2_rule(alias, init, run) + return _invocation( + "text.rule-filter", "filter", {"fields": [field], "rules": [rule]}, source_id + ) + + +def _translate_phase2_rule( + alias: str, init: dict[str, Any], run: dict[str, Any] +) -> dict[str, Any]: + output_defaults = { + "ColonEnd": "colonendfilter_label", + "LineEndWithEllipsis": "line_end_with_ellipsis_filter_label", + "SymbolWordRatio": "symbol_word_ratio_filter_label", + "AlphaWords": "alpha_words_filter_label", + "HtmlEntityFilter": "html_entity_filter_label", + "IDCard": "id_card_filter_label", + "NoPunc": "no_punc_filter_label", + "SpecialCharacter": "special_character_filter_label", + "Watermark": "watermark_filter_label", + "MeanWordLength": "mean_word_length_filter_label", + "StopWord": "stop_word_filter_label", + "CurlyBracket": "curly_bracket_filter_label", + "CapitalWords": "capital_words_filter", + "LoremIpsum": "loremipsum_filter_label", + "LineStartWithBulletpoint": "line_start_with_bullet_point_filter_label", + "LineWithJavascript": "line_with_javascript_filter_label", + "Blocklist": "blocklist_filter_label", + } + rule_types = { + "ColonEnd": "colonEnd", + "LineEndWithEllipsis": "lineEndWithEllipsis", + "SymbolWordRatio": "symbolWordRatio", + "AlphaWords": "alphaWords", + "HtmlEntityFilter": "htmlEntity", + "IDCard": "idCard", + "NoPunc": "noPunc", + "SpecialCharacter": "specialCharacter", + "Watermark": "watermark", + "MeanWordLength": "meanWordLength", + "StopWord": "stopWord", + "CurlyBracket": "curlyBracket", + "CapitalWords": "capitalWords", + "LoremIpsum": "loremIpsum", + "LineStartWithBulletpoint": "lineStartWithBulletpoint", + "LineWithJavascript": "lineWithJavascript", + "Blocklist": "blocklist", + } + if alias not in rule_types: + _unsupported("alias has no pinned compatibility translation") + output_key = _optional_string(run, "output_key", output_defaults[alias]) + rule: dict[str, Any] = {"type": rule_types[alias], "outputKey": output_key} + + if alias in {"ColonEnd", "HtmlEntityFilter", "SpecialCharacter"}: + _only_keys(init, set(), alias) + elif alias in { + "LineEndWithEllipsis", + "SymbolWordRatio", + "IDCard", + "NoPunc", + "CurlyBracket", + "LoremIpsum", + "LineStartWithBulletpoint", + "LineWithJavascript", + }: + _only_keys(init, {"threshold"}, alias) + defaults = { + "LineEndWithEllipsis": 0.3, + "SymbolWordRatio": 0.4, + "IDCard": 3, + "NoPunc": 112, + "CurlyBracket": 0.025, + "LoremIpsum": 3e-8, + "LineStartWithBulletpoint": 0.9, + "LineWithJavascript": 3, + } + default = defaults[alias] + rule["threshold"] = ( + _integer(init.get("threshold", default), "threshold") + if isinstance(default, int) + else _number(init.get("threshold", default), "threshold") + ) + elif alias in {"AlphaWords", "StopWord"}: + _only_keys(init, {"threshold", "use_tokenizer"}, alias) + if "threshold" not in init or "use_tokenizer" not in init: + _unsupported(f"{alias} requires threshold and use_tokenizer") + rule["threshold"] = _number(init["threshold"], "threshold") + rule["useTokenizer"] = _boolean(init["use_tokenizer"], "use_tokenizer") + elif alias == "Watermark": + _only_keys(init, {"watermarks"}, alias) + rule["watermarks"] = _string_list( + init.get("watermarks", ["Copyright", "Watermark", "Confidential"]), + "watermarks", + ) + elif alias == "MeanWordLength": + _only_keys(init, {"min_length", "max_length"}, alias) + rule["minLength"] = _number(init.get("min_length", 3), "min_length") + rule["maxLength"] = _number(init.get("max_length", 10), "max_length") + elif alias == "CapitalWords": + _only_keys(init, {"threshold", "use_tokenizer"}, alias) + rule["threshold"] = _number(init.get("threshold", 0.2), "threshold") + rule["useTokenizer"] = _boolean( + init.get("use_tokenizer", False), "use_tokenizer" + ) + elif alias == "Blocklist": + _only_keys(init, {"language", "threshold", "use_tokenizer"}, alias) + language = init.get("language", "en") + if language not in _BLOCKLIST_ASSETS: + _unsupported("Blocklist language must be en or zh") + rule["language"] = language + rule["threshold"] = _integer(init.get("threshold", 1), "threshold") + rule["useTokenizer"] = _boolean( + init.get("use_tokenizer", False), "use_tokenizer" + ) + return rule + + +def _text_clean( + items: list[dict[str, Any]], config: dict[str, Any] +) -> tuple[list[dict[str, Any]], dict[str, Any], list[str]]: + _validate_items(items) + _only_keys(config, {"fields", "operations", "htmlEntities"}, "text.clean") + fields = _string_list(config.get("fields", ["title", "content"]), "fields") + operations = _string_list( + config.get("operations", ["removeExtraSpaces"]), "operations" + ) + unknown = set(operations) - set(_CLEAN_OPERATIONS.values()) + if unknown: + _unsupported(f"text.clean operation is unsupported: {sorted(unknown)[0]}") + entities = _string_list( + config.get("htmlEntities", _DEFAULT_HTML_ENTITIES), "htmlEntities" + ) + output = copy.deepcopy(items) + changed = 0 + for item in output: + normalized = _normalized(item) + for field in fields: + value = normalized.get(field) + if not isinstance(value, str): + continue + refined = value + for operation in operations: + refined = _clean_value(refined, operation, entities) + changed += refined != value + normalized[field] = refined + return output, {"changedFieldCount": changed}, [] + + +def _text_rule_filter( + items: list[dict[str, Any]], config: dict[str, Any] +) -> tuple[list[dict[str, Any]], dict[str, Any], list[str]]: + _validate_items(items) + _only_keys(config, {"fields", "rules"}, "text.rule-filter") + fields = _string_list(config.get("fields", ["content"]), "fields") + rules = config.get( + "rules", + [{"type": "contentNull", "outputKey": "content_null_filter_label"}], + ) + if not isinstance(rules, list) or not rules or any( + not isinstance(rule, dict) for rule in rules + ): + _unsupported("text.rule-filter rules must be a non-empty list of objects") + output: list[dict[str, Any]] = [] + rejected: list[str] = [] + rule_hits: dict[str, int] = {} + for item in copy.deepcopy(items): + normalized = _normalized(item) + accepted = True + for rule in rules: + passed, result = _apply_rule(normalized, fields[0], rule) + if passed: + normalized[_rule_output_key(rule)] = result + else: + accepted = False + rule_type = str(rule.get("type")) + rule_hits[rule_type] = rule_hits.get(rule_type, 0) + 1 + if accepted: + output.append(item) + else: + _record_rejection(rejected, item) + return output, { + "ruleHits": rule_hits, + "rejectedInputCount": len(items) - len(output), + }, rejected + + +def _text_deduplicate( + items: list[dict[str, Any]], config: dict[str, Any] +) -> tuple[list[dict[str, Any]], dict[str, Any], list[str]]: + _validate_items(items) + _only_keys( + config, + {"fields", "hashFunction", "mode", "nGram", "diffSize", "outputKey"}, + "text.deduplicate", + ) + fields = _string_list(config.get("fields", ["title", "content"]), "fields") + mode = config.get("mode", "exact") + if mode not in {"exact", "ngramHash"}: + _unsupported("text.deduplicate mode must be exact or ngramHash") + hash_function = config.get("hashFunction", "md5") + output_key = _config_string(config, "outputKey", "minhash_deduplicated_label") + if hash_function not in {"md5", "sha256", "xxh3"}: + _unsupported("text.deduplicate hashFunction is unsupported") + if mode == "ngramHash" and hash_function == "xxh3": + _unsupported( + "ngramHash xxh3 is unavailable without an explicit dependency" + ) + if hash_function == "xxh3": + try: + from xxhash import xxh3_128 # type: ignore[import-not-found] + except ImportError: + _unsupported("xxh3 is unavailable in this installation") + hasher: Callable[[bytes], Any] = xxh3_128 + else: + hasher = getattr(hashlib, hash_function) + + output: list[dict[str, Any]] = [] + rejected: list[str] = [] + seen: set[str] = set() + seen_ngrams: list[set[str]] = [] + n_gram = _integer(config.get("nGram", 3), "nGram") + diff_size = _integer(config.get("diffSize", 1), "diffSize") + if mode == "ngramHash" and n_gram <= 0: + _unsupported("text.deduplicate nGram must be greater than zero") + for item in copy.deepcopy(items): + normalized = _normalized(item) + if len(fields) > 1: + text = "\n".join(f"{field}:\n{normalized[field]}" for field in fields) + else: + text = normalized[fields[0]] + if mode == "ngramHash": + gram_length = len(text) // n_gram + digests = { + hasher( + text[index * gram_length : (index + 1) * gram_length].encode( + "utf-8" + ) + ).hexdigest() + for index in range(n_gram) + } + if any( + len(digests.intersection(previous)) >= diff_size + for previous in seen_ngrams + ): + _record_rejection(rejected, item) + continue + seen_ngrams.append(digests) + else: + digest = hasher(text.encode("utf-8")).hexdigest() + if digest in seen: + _record_rejection(rejected, item) + continue + seen.add(digest) + normalized[output_key] = 1 + output.append(item) + return output, { + "duplicateCount": len(items) - len(output), + "rejectedInputCount": len(items) - len(output), + }, rejected + + +COMPAT_EXECUTORS: dict[tuple[str, str], _Executor] = { + ("text.clean", COMPAT_PACK_VERSION): _text_clean, + ("text.rule-filter", COMPAT_PACK_VERSION): _text_rule_filter, + ("text.deduplicate", COMPAT_PACK_VERSION): _text_deduplicate, +} + + +def _clean_value(value: str, operation: str, entities: list[str]) -> str: + if operation == "removeExtraSpaces": + return " ".join(value.split()) + if operation == "lowercase": + return value.lower() + if operation == "htmlUrlRemover": + value = re.sub(r"https?:\/\/\S+[\r\n]*", "", value, flags=re.MULTILINE) + return re.sub(r"<.*?>", "", value) + if operation == "htmlEntity": + patterns = [ + pattern + for entity in entities + for pattern in ( + f"&{entity};", + f"&{entity};", + f"&{entity};", + f"&{entity};", + ) + ] + return re.sub("|".join(patterns), "", value) + if operation == "removeEmoji": + return re.sub( + "[" + "\U0001f600-\U0001f64f" + "\U0001f300-\U0001f5ff" + "\U0001f680-\U0001f6ff" + "\U0001f1e0-\U0001f1ff" + "\u2702-\u27b0" + "]+", + "", + value, + ) + if operation == "removeNumber": + return "".join(character for character in value if not character.isdigit()) + if operation == "removePunctuation": + return value.translate(str.maketrans("", "", string.punctuation)) + if operation == "removeRepetitionsPunctuation": + return re.sub(r"([^\w\s_])\1+|(_)\2+", r"\1\2", value) + if operation == "textNormalization": + refined = re.sub( + r"(\d{1,2})[/.](\d{1,2})[/.](\d{2,4})", r"\3-\2-\1", value + ) + date_patterns = ( + (r"\b(\w+)\s+(\d{1,2}),\s+(\d{4})\b", "%B %d, %Y"), + (r"\b(\d{1,2})\s+(\w+)\s+(\d{4})\b", "%d %B %Y"), + ) + for pattern, date_format in date_patterns: + match = re.search(pattern, refined) + if match is None: + continue + date_text = match.group(0) + try: + normalized_date = datetime.strptime( + date_text, date_format + ).strftime("%Y-%m-%d") + except ValueError: + continue + refined = refined.replace(date_text, normalized_date) + return re.sub(r"\$\s?(\d+)", r"\1 USD", refined) + if operation == "removeImageRefs": + return re.sub( + r"!\[\]\(images\/[0-9a-fA-F]\.jpg\)|" + r"[a-fA-F0-9]+\.[a-zA-Z]{3,4}\)|" + r"!\[\]\(images\/[a-f0-9]|" + r"图\s+\d+-\d+:[\u4e00-\u9fa5a-zA-Z0-9]+|" + r"(?:[0-9a-zA-Z]+){7,}|" + r"(?:[一二三四五六七八九十零壹贰叁肆伍陆柒捌玖拾佰仟万亿]+){5,}|" + r"u200e|" + r"÷|\? :|" + r"[�□]|\{\/U\}|" + r"U\+26[0-F][0-D]|U\+273[3-4]|U\+1F[3-6][0-4][0-F]|" + r"U\+1F6[8-F][0-F]", + "", + value, + ) + raise AssertionError(operation) + + +def _apply_rule( + normalized: dict[str, Any], field: str, rule: dict[str, Any] +) -> tuple[bool, int]: + rule_type = rule.get("type") + allowed = { + "contentNull": {"type", "outputKey"}, + "wordNumber": {"type", "min", "max", "outputKey"}, + "sentenceNumber": {"type", "min", "max", "outputKey"}, + "charNumber": {"type", "threshold", "outputKey"}, + "uniqueWords": {"type", "threshold", "outputKey"}, + "colonEnd": {"type", "outputKey"}, + "lineEndWithEllipsis": {"type", "threshold", "outputKey"}, + "symbolWordRatio": {"type", "threshold", "outputKey"}, + "alphaWords": {"type", "threshold", "useTokenizer", "outputKey"}, + "htmlEntity": {"type", "outputKey"}, + "idCard": {"type", "threshold", "outputKey"}, + "noPunc": {"type", "threshold", "outputKey"}, + "specialCharacter": {"type", "outputKey"}, + "watermark": {"type", "watermarks", "outputKey"}, + "meanWordLength": {"type", "minLength", "maxLength", "outputKey"}, + "stopWord": {"type", "threshold", "useTokenizer", "outputKey"}, + "curlyBracket": {"type", "threshold", "outputKey"}, + "capitalWords": {"type", "threshold", "useTokenizer", "outputKey"}, + "loremIpsum": {"type", "threshold", "outputKey"}, + "lineStartWithBulletpoint": {"type", "threshold", "outputKey"}, + "lineWithJavascript": {"type", "threshold", "outputKey"}, + "blocklist": { + "type", + "language", + "threshold", + "useTokenizer", + "outputKey", + }, + } + if rule_type not in allowed: + _unsupported("text.rule-filter rule type is unsupported") + _only_keys(rule, allowed[rule_type], f"text.rule-filter {rule_type}") + text = normalized.get(field) + if rule_type == "contentNull": + return bool(text is not None and text.strip() != ""), 1 + if not text: + return False, 0 + if rule_type == "wordNumber": + count = len(tuple(text.split())) + return ( + _integer(rule.get("min", 20), "min") + <= count + < _integer(rule.get("max", 100000), "max") + ), count + if rule_type == "sentenceNumber": + count = len(re.findall(r"\b[^.!?\n]+[.!?]*", text, flags=re.UNICODE)) + return ( + _integer(rule.get("min", 3), "min") + <= count + <= _integer(rule.get("max", 7500), "max") + ), 1 + if rule_type == "charNumber": + count = len(text.strip().replace(" ", "").replace("\n", "").replace("\t", "")) + return count >= _integer(rule.get("threshold", 100), "threshold"), 1 + if rule_type == "uniqueWords": + words = tuple(text.lower().split()) + ratio = len(set(words)) / len(words) if words else 0.0 + return ratio > _number(rule.get("threshold", 0.1), "threshold"), 1 + if rule_type == "colonEnd": + return not text.endswith(":"), 1 + if rule_type == "lineEndWithEllipsis": + paragraphs = _paragraphs(text) + if not paragraphs: + return False, 1 + occurrences = sum( + paragraph.rstrip().endswith(("...", "…")) for paragraph in paragraphs + ) + ratio = occurrences / len(paragraphs) + return ratio < _number(rule.get("threshold", 0.3), "threshold"), 1 + if rule_type == "symbolWordRatio": + tokens = re.findall(r"\w+|[^\w\s]+", text, flags=re.UNICODE) + if not tokens: + return False, 1 + num_symbols = float(sum(text.count(symbol) for symbol in ("#", "...", "…"))) + return ( + num_symbols / len(tokens) + < _number(rule.get("threshold", 0.4), "threshold") + ), 1 + if rule_type == "alphaWords": + words = _tokenize(text, _boolean(rule.get("useTokenizer"), "useTokenizer")) + if not words: + return False, 1 + ratio = sum(bool(re.search(r"[a-zA-Z]", word)) for word in words) / len( + words + ) + return ratio > _number(rule.get("threshold"), "threshold"), 1 + if rule_type == "htmlEntity": + patterns = ( + pattern + for entity in _DEFAULT_HTML_ENTITIES + for pattern in ( + f"&{entity};", + f"&{entity};", + f"&{entity};", + f"&{entity};", + f"&{entity}", + f"&{entity}", + ) + ) + return not any(pattern in text for pattern in patterns), 1 + if rule_type == "idCard": + pattern = ( + r"(身\s{0,10}份|id\s{0,10}number\s{0,10}|identification|identity|" + r"\s{0,10}ID\s{0,10}No\s{0,10}|id\s{0,10}card\s{0,10}|" + r"NRIC\s{0,10}number\s{0,10}|IC\s{0,10}number\s{0,10}|" + r"resident\s{0,10}registration\s{0,10}|" + r"I.D.\s{0,10}Number\s{0,10})" + ) + return len(re.findall(pattern, text, re.IGNORECASE)) < _integer( + rule.get("threshold", 3), "threshold" + ), 1 + if rule_type == "noPunc": + paragraphs = tuple(line for line in text.split("\n") if line.strip()) + longest = max( + len(sentence.split()) + for paragraph in paragraphs + for sentence in re.split(r"[–.!?,;•/|…]", paragraph) + ) if paragraphs else 0 + return longest <= _integer(rule.get("threshold", 112), "threshold"), 1 + if rule_type == "specialCharacter": + patterns = ( + r"u200e", + r"÷|\? :", + r"[�□]|\{\/U\}", + r"U\+26[0-F][0-D]|U\+273[3-4]|U\+1F[3-6][0-4][0-F]|" + r"U\+1F6[8-F][0-F]", + ) + return not any(re.search(pattern, text) for pattern in patterns), 1 + if rule_type == "watermark": + watermarks = _string_list( + rule.get("watermarks", ["Copyright", "Watermark", "Confidential"]), + "watermarks", + ) + return re.search("|".join(watermarks), text) is None, 1 + if rule_type == "meanWordLength": + words = text.split() + if not words: + return False, 1 + mean_length = round(sum(map(len, words)) / len(words), 2) + return ( + _number(rule.get("minLength", 3), "minLength") + <= mean_length + < _number(rule.get("maxLength", 10), "maxLength") + ), 1 + if rule_type == "stopWord": + words = _tokenize( + text.lower(), _boolean(rule.get("useTokenizer"), "useTokenizer") + ) + count = sum(word in _stopwords() for word in words) + ratio = count / len(words) if words else 0 + return ( + ratio > _number(rule.get("threshold"), "threshold") + and count > 2 + ), 1 + if rule_type == "curlyBracket": + ratio = (text.count("{") + text.count("}")) / len(text) + return ratio < _number(rule.get("threshold", 0.025), "threshold"), 1 + if rule_type == "capitalWords": + words = _tokenize( + text, _boolean(rule.get("useTokenizer", False), "useTokenizer") + ) + ratio = sum(map(str.isupper, words)) / len(words) if words else 0 + return ratio <= _number(rule.get("threshold", 0.2), "threshold"), 1 + if rule_type == "loremIpsum": + ratio = len(re.findall("lorem ipsum", text, re.IGNORECASE)) / len( + text.lower() + ) + return ratio <= _number(rule.get("threshold", 3e-8), "threshold"), 1 + if rule_type == "lineStartWithBulletpoint": + paragraphs = _paragraphs(text) + bullets = ("•", "‣", "▶", "◀", "◦", "■", "□", "▪", "▫", "–") + ratio = sum(line.lstrip().startswith(bullets) for line in paragraphs) / len( + paragraphs + ) + return ratio <= _number(rule.get("threshold", 0.9), "threshold"), 1 + if rule_type == "lineWithJavascript": + paragraphs = _paragraphs(text, normalize=True) + occurrences = sum("javascript" in line for line in paragraphs) + not_javascript = len(paragraphs) - occurrences + return ( + len(paragraphs) <= 3 + or not_javascript >= _integer(rule.get("threshold", 3), "threshold") + ), 1 + if rule_type == "blocklist": + words = _tokenize( + text.lower(), _boolean(rule.get("useTokenizer", False), "useTokenizer") + ) + count = sum(word in _blocklist(rule.get("language", "en")) for word in words) + return count <= _integer(rule.get("threshold", 1), "threshold"), 1 + raise AssertionError(rule_type) + + +def _invocation( + operator_id: str, kind: str, config: dict[str, Any], source_id: str +) -> DataFlowInvocation: + return DataFlowInvocation( + operator_id=operator_id, + pack_id=COMPAT_PACK_ID, + pack_version=COMPAT_PACK_VERSION, + kind=kind, + config=copy.deepcopy(config), + source_id=source_id, + ) + + +def _paragraphs(text: str, *, normalize: bool = False) -> tuple[str, ...]: + paragraphs = tuple( + match + for match in re.findall(r"([^\n]*\n|[^\n]+$)", text) + if match.strip() + ) + if not normalize: + return paragraphs + return tuple(_normalize_rule_line(paragraph) for paragraph in paragraphs) + + +def _normalize_rule_line(value: str) -> str: + punctuation = string.punctuation.replace("_", "") + value = value.translate(str.maketrans("", "", punctuation)) + value = re.sub(r"\s+", " ", value.lower().strip()) + return unicodedata.normalize("NFD", value) + + +def _tokenize(text: str, use_tokenizer: bool) -> tuple[str, ...]: + if not use_tokenizer: + return tuple(text.split()) + try: + from nltk.tokenize import word_tokenize + + return tuple(word_tokenize(text)) + except (ImportError, LookupError): + _unsupported("NLTK word tokenizer data is unavailable in this installation") + + +@cache +def _decoded_asset(sha256: str, encoded: str) -> frozenset[str]: + raw = zlib.decompress(base64.b64decode(encoded)) + if hashlib.sha256(raw).hexdigest() != sha256: + _unsupported("embedded compatibility asset failed its SHA-256 check") + return frozenset(line.strip().lower() for line in raw.decode().splitlines() if line) + + +def _stopwords() -> frozenset[str]: + return _decoded_asset(*_NLTK_STOPWORDS_ASSET) + + +def _blocklist(language: Any) -> frozenset[str]: + if not isinstance(language, str) or language not in _BLOCKLIST_ASSETS: + _unsupported("Blocklist language must be en or zh") + return _decoded_asset(*_BLOCKLIST_ASSETS[language]) + + +def _normalized(item: dict[str, Any]) -> dict[str, Any]: + value = item.get("normalizedData") + if not isinstance(value, dict): + value = {} + item["normalizedData"] = value + return value + + +def _record_rejection(rejected: list[str], item: dict[str, Any]) -> None: + candidate_id = item.get("candidateId") + if isinstance(candidate_id, str): + rejected.append(candidate_id) + elif candidate_id is not None: + rejected.append(str(candidate_id)) + + +def _rule_output_key(rule: dict[str, Any]) -> str: + return _config_string(rule, "outputKey", "filter_label") + + +def _validate_items(items: object) -> None: + if not isinstance(items, list) or any(not isinstance(item, dict) for item in items): + _unsupported("items must be a list of candidate objects") + + +def _mapping(value: Mapping[str, Any] | None, name: str) -> dict[str, Any]: + if value is None: + return {} + if not isinstance(value, Mapping): + _unsupported(f"{name} must be an object") + return copy.deepcopy(dict(value)) + + +def _only_keys(value: Mapping[str, Any], allowed: set[str], name: str) -> None: + unknown = sorted(set(value) - allowed) + if unknown: + _unsupported(f"{name} config key is unsupported: {unknown[0]}") + + +def _required_string(value: Mapping[str, Any], key: str) -> str: + result = value.get(key) + if not isinstance(result, str) or not result: + _unsupported(f"{key} must be a non-empty string") + return result + + +def _optional_string(value: Mapping[str, Any], key: str, default: str) -> str: + result = value.get(key, default) + if not isinstance(result, str) or not result: + _unsupported(f"{key} must be a non-empty string") + return result + + +def _config_string(value: Mapping[str, Any], key: str, default: str) -> str: + return _optional_string(value, key, default) + + +def _string_list(value: Any, name: str) -> list[str]: + if not isinstance(value, list) or any( + not isinstance(item, str) or not item for item in value + ): + _unsupported(f"{name} must be a list of non-empty strings") + return list(value) + + +def _integer(value: Any, name: str) -> int: + if not isinstance(value, int) or isinstance(value, bool): + _unsupported(f"{name} must be an integer") + return value + + +def _number(value: Any, name: str) -> float: + if not isinstance(value, int | float) or isinstance(value, bool): + _unsupported(f"{name} must be a number") + return float(value) + + +def _boolean(value: Any, name: str) -> bool: + if not isinstance(value, bool): + _unsupported(f"{name} must be a boolean") + return value + + +def _unsupported(reason: str) -> Any: + raise ValueError(f"dataflow_operator_unsupported: {reason}") + + +__all__ = [ + "COMPAT_EXECUTORS", + "COMPAT_OPERATOR_DEFINITIONS", + "COMPAT_PACK_ID", + "COMPAT_PACK_VERSION", + "DATAFLOW_ALIAS_SOURCE_IDS", + "DATAFLOW_COMPAT_SHA", + "DataFlowInvocation", + "translate_dataflow_alias", +] diff --git a/backend/workflow/demand_assembler.py b/backend/workflow/demand_assembler.py index f81b45a..429063e 100644 --- a/backend/workflow/demand_assembler.py +++ b/backend/workflow/demand_assembler.py @@ -41,7 +41,13 @@ def draft_workflow_demand(body: WorkflowDemandDraftRequest) -> WorkflowPatchResp ], ) - operations = _native_first_loop_operations(body.project, sources, body.text, body.locale) + operations = _native_first_loop_operations( + body.project, + sources, + body.text, + body.locale, + data_operators=_data_operators_for_need(body.text), + ) return preview_workflow_patch(body.project, operations) @@ -50,6 +56,8 @@ def _native_first_loop_operations( sources: list[dict[str, Any]], demand_text: str, locale: str | None, + *, + data_operators: list[dict[str, Any]], ) -> list[WorkflowPatchOperation]: operations: list[WorkflowPatchOperation] = [] used_node_ids = {node.id for node in project.nodes} @@ -132,8 +140,13 @@ def _native_first_loop_operations( ) merge_id = _unique_id(used_node_ids, "merge-candidates") + operator_node_ids = [ + _unique_id(used_node_ids, f"{operator['id']}-data") + for operator in data_operators + ] accept_id = _unique_id(used_node_ids, "accept-records") sink_id = _unique_id(used_node_ids, "record-sink") + accept_x = 960 + len(operator_node_ids) * 260 operations.extend( [ WorkflowPatchOperation( @@ -171,7 +184,7 @@ def _native_first_loop_operations( ui={ "catalogId": "intelligence.control.record-acceptance", "label": "Record Acceptance", - "position": {"x": 960, "y": 240}, + "position": {"x": accept_x, "y": 240}, }, ), ), @@ -189,12 +202,35 @@ def _native_first_loop_operations( ui={ "catalogId": "intelligence.sink.records", "label": "Records", - "position": {"x": 1220, "y": 240}, + "position": {"x": accept_x + 260, "y": 240}, }, ), ), ] ) + for index, (operator, operator_node_id) in enumerate( + zip(data_operators, operator_node_ids, strict=True) + ): + operations.append( + WorkflowPatchOperation( + op="add_node", + node=WorkflowProjectNode( + id=operator_node_id, + kind="agent", + capability="normalize", + params={ + "operatorId": operator["operatorId"], + "packVersion": operator["packVersion"], + "config": operator.get("config", {}), + }, + ui={ + "catalogId": operator["catalogId"], + "label": operator["label"], + "position": {"x": 960 + index * 260, "y": 240}, + }, + ), + ) + ) for index, normalize_id in enumerate(normalize_ids, start=1): operations.append( WorkflowPatchOperation( @@ -208,13 +244,46 @@ def _native_first_loop_operations( ), ) ) + terminal_id = operator_node_ids[-1] if operator_node_ids else merge_id + if operator_node_ids: + operations.append( + WorkflowPatchOperation( + op="connect_nodes", + edge=WorkflowProjectEdge( + id=_unique_id( + used_edge_ids, + f"e-{merge_id}-{operator_node_ids[0]}", + ), + source=merge_id, + target=operator_node_ids[0], + sourcePort="out", + targetPort="in", + ), + ) + ) + for source_id, target_id in zip( + operator_node_ids, + operator_node_ids[1:], + ): + operations.append( + WorkflowPatchOperation( + op="connect_nodes", + edge=WorkflowProjectEdge( + id=_unique_id(used_edge_ids, f"e-{source_id}-{target_id}"), + source=source_id, + target=target_id, + sourcePort="out", + targetPort="in", + ), + ) + ) operations.extend( [ WorkflowPatchOperation( op="connect_nodes", edge=WorkflowProjectEdge( - id=_unique_id(used_edge_ids, f"e-{merge_id}-{accept_id}"), - source=merge_id, + id=_unique_id(used_edge_ids, f"e-{terminal_id}-{accept_id}"), + source=terminal_id, target=accept_id, sourcePort="out", targetPort="candidates", @@ -235,6 +304,117 @@ def _native_first_loop_operations( return operations +def _data_operators_for_need(text: str) -> list[dict[str, Any]]: + normalized = text.lower() + dataflow_compat = "dataflow" in normalized + training_data = any( + token in normalized + for token in ("训练数据", "sft", "instruction", "instruction data", "微调数据") + ) + quality_work = training_data or any( + token in normalized + for token in ( + "dataflow", + "数据准备", + "数据清洗", + "清洗", + "quality", + "filter", + "evaluate", + "refine", + "质量", + "过滤", + "评估", + "筛选", + ) + ) + if not quality_work: + return [] + + operators = [ + { + "id": "chunk", + "catalogId": "intelligence.data.generate", + "operatorId": "data.chunk", + "packVersion": "1.0.0", + "label": "Chunk Data", + }, + { + "id": "clean", + "catalogId": "intelligence.data.refine", + "operatorId": "text.clean", + "packVersion": "1.1.0" if dataflow_compat else "1.0.0", + "config": ( + { + "fields": ["content"], + "operations": [ + "removeEmoji", + "htmlUrlRemover", + "removeExtraSpaces", + ], + } + if dataflow_compat + else {} + ), + "label": "Clean Text", + }, + { + "id": "deduplicate", + "catalogId": "intelligence.data.filter", + "operatorId": "text.deduplicate", + "packVersion": "1.1.0" if dataflow_compat else "1.0.0", + "config": ( + { + "fields": ["content"], + "mode": "exact", + "hashFunction": "md5", + } + if dataflow_compat + else {} + ), + "label": "Deduplicate Text", + }, + { + "id": "rule-filter", + "catalogId": "intelligence.data.filter", + "operatorId": "text.rule-filter", + "packVersion": "1.1.0" if dataflow_compat else "1.0.0", + "config": ( + { + "fields": ["content"], + "rules": [ + { + "type": "contentNull", + "outputKey": "content_null_filter_label", + } + ], + } + if dataflow_compat + else {} + ), + "label": "Filter Text Rules", + }, + { + "id": "statistics", + "catalogId": "intelligence.data.evaluate", + "operatorId": "text.statistics", + "packVersion": "1.0.0", + "label": "Text Statistics", + }, + ] + if training_data: + operators.append( + { + "id": "generate", + "catalogId": "intelligence.data.generate", + "operatorId": "core.generate.instruction-pairs", + "packVersion": "1.0.0", + "label": "Generate Instruction Pairs", + } + ) + return operators + + def _source_slots_for_need(text: str) -> list[dict[str, Any]]: normalized = text.lower() slots: list[dict[str, Any]] = [] diff --git a/backend/workflow/external_importer.py b/backend/workflow/external_importer.py index 443585e..255d74f 100644 --- a/backend/workflow/external_importer.py +++ b/backend/workflow/external_importer.py @@ -12,6 +12,10 @@ WorkflowProjectEdge, WorkflowProjectNode, ) +from backend.workflow.dataflow_compat import ( + DATAFLOW_COMPAT_SHA, + translate_dataflow_alias, +) from backend.workflow.patcher import preview_workflow_patch EXTERNAL_TOOL_CAPABILITY_ID = "external.tool.capability" @@ -24,6 +28,13 @@ def import_external_workflow(body: WorkflowExternalImportRequest) -> WorkflowPat OpenCLI Admin catalog capabilities and never carry external executors. """ + if body.runtime == "dataflow": + source_sha = _read_string(body.graph.get("sourceSha")) + if source_sha != DATAFLOW_COMPAT_SHA: + raise ValueError( + f"DataFlow graph.sourceSha must equal {DATAFLOW_COMPAT_SHA}" + ) + external_nodes = _extract_nodes(body.graph) external_edges = _extract_edges(body.graph) for edge in external_edges: @@ -95,6 +106,14 @@ def _to_workflow_node( graph_name: str | None, index: int, ) -> WorkflowProjectNode: + if runtime == "dataflow": + return _to_dataflow_node( + external_node, + node_id=node_id, + graph_name=graph_name, + index=index, + ) + external_id = _node_id(external_node) external_type = _node_type(external_node) catalog_id, kind, capability, params = _native_capability_for_external_node( @@ -129,6 +148,73 @@ def _to_workflow_node( ) +def _to_dataflow_node( + external_node: dict[str, Any], + *, + node_id: str, + graph_name: str | None, + index: int, +) -> WorkflowProjectNode: + external_id = _node_id(external_node) + source_id = _dataflow_source_id(external_node) + if "config" in external_node: + raise ValueError( + f'DataFlow node "{external_id}" must use initConfig/runConfig, not config' + ) + invocation = translate_dataflow_alias( + source_id, + external_node.get("initConfig"), + external_node.get("runConfig"), + ) + provenance = { + "runtime": "dataflow", + "sourceSha": DATAFLOW_COMPAT_SHA, + "sourceId": invocation.source_id, + } + return WorkflowProjectNode( + id=node_id, + kind="agent", + capability="normalize", + params={ + "operatorId": invocation.operator_id, + "packVersion": invocation.pack_version, + "config": invocation.config, + }, + ui={ + "catalogId": f"intelligence.data.{invocation.kind}", + "label": _node_label(external_node), + "position": {"x": 180 + (index % 4) * 260, "y": 180 + (index // 4) * 140}, + "externalWorkflow": { + "runtime": "dataflow", + "graphName": graph_name, + "nodeId": external_id, + **provenance, + }, + }, + ) + + +def _dataflow_source_id(external_node: dict[str, Any]) -> str: + source_id = _read_string(external_node.get("sourceId")) + if source_id is None: + module = _read_string(external_node.get("module")) + class_name = _read_string(external_node.get("class")) + if module is None or class_name is None: + raise ValueError( + f'DataFlow node "{_node_id(external_node)}" requires sourceId or module+class' + ) + source_id = ( + f"dataflow@{DATAFLOW_COMPAT_SHA}::" + f"dataflow.operators.{module}.{class_name}" + ) + if not source_id.startswith(f"dataflow@{DATAFLOW_COMPAT_SHA}::"): + raise ValueError( + f'DataFlow node "{_node_id(external_node)}" sourceId must use ' + f"DataFlow SHA {DATAFLOW_COMPAT_SHA}" + ) + return source_id + + def _tool_capability_params(external_node: dict[str, Any]) -> dict[str, Any]: tool_capability = external_node.get("toolCapability") if isinstance(tool_capability, dict): diff --git a/backend/workflow/node_registry.py b/backend/workflow/node_registry.py index 067c70a..5f41451 100644 --- a/backend/workflow/node_registry.py +++ b/backend/workflow/node_registry.py @@ -24,6 +24,10 @@ "intelligence.source.opencli-slot", "intelligence.processing.normalize", "intelligence.processing.dedupe", + "intelligence.data.generate", + "intelligence.data.filter", + "intelligence.data.evaluate", + "intelligence.data.refine", "intelligence.flow.merge", "intelligence.agent.summary", "intelligence.agent.score", diff --git a/backend/workflow/opencli_hda_tracer.py b/backend/workflow/opencli_hda_tracer.py index 998b6af..4504f6e 100644 --- a/backend/workflow/opencli_hda_tracer.py +++ b/backend/workflow/opencli_hda_tracer.py @@ -44,6 +44,7 @@ SOURCE_OUTPUT_REQUIRED, ) from backend.workflow.compiler import INTERNAL_ID_SEPARATOR, compile_workflow_project +from backend.workflow.data_operators import execute_data_operator from backend.workflow.event_mirror import publish_workflow_run_event_mirror from backend.workflow.fleet_inventory import match_workflow_fleet_capability from backend.workflow.joyai_vl_executor import ( @@ -58,6 +59,7 @@ execute_okx_market_ticker_snapshot, ) from backend.workflow.runtime_registry import ( + DATA_OPERATOR_CATALOG_BINDINGS, EXTERNAL_TOOL_BINDING_ID, INBOX_STORE_BINDING_ID, MERGE_BINDING_ID, @@ -92,6 +94,8 @@ class _StoredWorkflowRun: _RUNS: dict[str, _StoredWorkflowRun] = {} +_DATA_OPERATOR_BINDING_IDS = set(DATA_OPERATOR_CATALOG_BINDINGS.values()) +_LEGACY_DATA_OPERATOR_PACK_VERSION = "1.0.0" def build_opencli_hda_trace( @@ -491,6 +495,7 @@ async def start_workflow_run( continue if _is_first_loop_native_node(node): + emitter.emit(node, "started", message=_native_node_started_message(node)) try: details, output_items = await _execute_native_node( node, @@ -516,8 +521,31 @@ async def start_workflow_run( details=reason.details, ) continue + except Exception as exc: + if _binding_id(node) not in _DATA_OPERATOR_BINDING_IDS: + raise + binding_input = _read_dict( + _read_dict(node.runtime.get("binding")).get("input") + ) + reason = WorkflowRunBlockReason( + code="data_operator_execution_failed", + message="Data operator execution failed", + source="data_operator_runtime", + details={ + "bindingId": _binding_id(node), + "operatorId": binding_input.get("operatorId"), + "errorType": type(exc).__name__, + }, + ) + emitter.emit( + node, + "failed", + message=reason.message, + block_reason=reason, + details=reason.details, + ) + continue outputs_by_node[node.id] = output_items - emitter.emit(node, "started", message=_native_node_started_message(node)) if _binding_id(node) == EXTERNAL_TOOL_BINDING_ID: emitter.emit( node, @@ -1384,6 +1412,7 @@ def _is_first_loop_native_node(node: CompiledWorkflowNode) -> bool: return False return binding.get("binding_id") in { NORMALIZE_BINDING_ID, + *_DATA_OPERATOR_BINDING_IDS, MERGE_BINDING_ID, ROUTER_ROUTE_BINDING_ID, RECORD_ACCEPTANCE_BINDING_ID, @@ -1523,6 +1552,55 @@ async def _execute_native_node( }, candidates, ) + if binding_id in _DATA_OPERATOR_BINDING_IDS: + binding = _read_dict(node.runtime.get("binding")) + binding_input = _read_dict(binding.get("input")) + operator_id = _read_string(binding_input.get("operatorId")) + if operator_id is None: + raise ValueError(f"Data operator binding {binding_id} is missing operatorId") + pack_version = ( + _read_string(binding_input.get("packVersion")) + or _LEGACY_DATA_OPERATOR_PACK_VERSION + ) + config = binding_input.get("config") + if not isinstance(config, dict): + raise ValueError("Data operator config must be an object") + result = execute_data_operator( + operator_id, + input_items, + config, + pack_version=pack_version, + ) + output_items = [ + _append_data_operator_lineage( + item, + node, + operator_id=result.operator_id, + run_id=run_id, + ) + for item in result.items + ] + result_details = result.to_details() + rejected_candidate_ids = list(result.rejected_candidate_ids) + return ( + { + **result_details, + "packVersion": result.pack_version, + "bindingId": binding_id, + "inputPort": binding_input.get("inputPort", "recordCandidate[]"), + "outputPort": binding_input.get("outputPort", "recordCandidate[]"), + "inputItemCount": len(input_items), + "outputItemCount": len(output_items), + "rejectedCount": result.metrics.get( + "rejectedCount", len(rejected_candidate_ids) + ), + "rejectedCandidateIds": rejected_candidate_ids[:100], + "rejectedCandidateIdsTruncated": len(rejected_candidate_ids) > 100, + "metrics": result.metrics, + "lineage": _lineage_pointer(node), + }, + output_items, + ) if binding_id == MERGE_BINDING_ID: binding = _read_dict(node.runtime.get("binding")) binding_input = _read_dict(binding.get("input")) @@ -2057,6 +2135,16 @@ def _append_lineage( return updated +def _append_data_operator_lineage( + item: dict[str, Any], + node: CompiledWorkflowNode, + *, + operator_id: str, + run_id: str, +) -> dict[str, Any]: + return _append_lineage(item, node, step=operator_id, run_id=run_id) + + def _accept_candidate( item: dict[str, Any], node: CompiledWorkflowNode, @@ -2333,6 +2421,8 @@ def _notify_send_block_reason( def _native_node_started_message(node: CompiledWorkflowNode) -> str: binding_id = _binding_id(node) + if binding_id in _DATA_OPERATOR_BINDING_IDS: + return "Data operator started" if binding_id == NORMALIZE_BINDING_ID: return "Normalize transform started" if binding_id == MERGE_BINDING_ID: @@ -2356,6 +2446,8 @@ def _native_node_started_message(node: CompiledWorkflowNode) -> str: def _native_node_partial_message(node: CompiledWorkflowNode) -> str: binding_id = _binding_id(node) + if binding_id in _DATA_OPERATOR_BINDING_IDS: + return "Data operator emitted items" if binding_id == NORMALIZE_BINDING_ID: return "Record Candidates projected" if binding_id == MERGE_BINDING_ID: @@ -2379,6 +2471,8 @@ def _native_node_partial_message(node: CompiledWorkflowNode) -> str: def _native_node_completed_message(node: CompiledWorkflowNode) -> str: binding_id = _binding_id(node) + if binding_id in _DATA_OPERATOR_BINDING_IDS: + return "Data operator completed" if binding_id == NORMALIZE_BINDING_ID: return "Normalize transform completed" if binding_id == MERGE_BINDING_ID: diff --git a/backend/workflow/runtime_contracts.py b/backend/workflow/runtime_contracts.py index 7a01c56..18945bc 100644 --- a/backend/workflow/runtime_contracts.py +++ b/backend/workflow/runtime_contracts.py @@ -163,6 +163,21 @@ def to_manifest(self) -> dict[str, object]: event_shape=("partial:recordCandidateCount", "completed"), fixture_coverage=("happy-path", "sse-parity", "odp-redis-mirror"), ), + **{ + f"workflow.data.{operation}": RuntimeIOContract( + binding_id=f"workflow.data.{operation}", + status="executable", + input_ports=(("in", "recordCandidate[]"),), + output_ports=(("out", "recordCandidate[]"),), + input_params=("operatorId", "packVersion", "config"), + output_artifacts=("recordCandidate[]", "metrics", "rejectedCandidateIds"), + permission_gate=(), + config_gate=("data_operator_registry",), + event_shape=("partial:outputItemCount", "completed", "failed"), + fixture_coverage=("data-operator-unit", "workflow-data-operator-e2e"), + ) + for operation in ("generate", "filter", "evaluate", "refine") + }, "workflow.flow.merge": RuntimeIOContract( binding_id="workflow.flow.merge", status="executable", diff --git a/backend/workflow/runtime_registry.py b/backend/workflow/runtime_registry.py index b9921d4..21666d4 100644 --- a/backend/workflow/runtime_registry.py +++ b/backend/workflow/runtime_registry.py @@ -20,6 +20,7 @@ MISSING_TURBOPUSH_CONTENT_TYPE, MISSING_TURBOPUSH_SERVICE, ) +from backend.workflow.data_operators import resolve_data_operator from backend.workflow.runtime_contracts import runtime_io_contract_manifest from backend.workflow.tool_capabilities import resolve_workflow_tool_capability from backend.workflow.turbopush_runtime import ( @@ -43,6 +44,12 @@ SOURCE_POOL_BINDING_ID = "workflow.source-pool.parallel-fanout" COLLECTION_OUTPUT_BINDING_ID = "workflow.collection-output.items" NORMALIZE_BINDING_ID = "workflow.transform.normalize" +DATA_OPERATOR_CATALOG_BINDINGS = { + "intelligence.data.generate": "workflow.data.generate", + "intelligence.data.filter": "workflow.data.filter", + "intelligence.data.evaluate": "workflow.data.evaluate", + "intelligence.data.refine": "workflow.data.refine", +} MERGE_BINDING_ID = "workflow.flow.merge" ROUTER_ROUTE_BINDING_ID = "workflow.router.route" RECORD_ACCEPTANCE_BINDING_ID = "workflow.gate.record-acceptance" @@ -52,6 +59,7 @@ NOTIFY_SEND_BINDING_ID = "workflow.notify.send" EXTERNAL_TOOL_BINDING_ID = "workflow.external-tool.capability" SUPPORTED_TOOL_EXECUTOR_MODES = {"fixture", "okx_market_ticker_snapshot", "joyai_vl_interaction"} +_LEGACY_DATA_OPERATOR_PACK_VERSION = "1.0.0" class WorkflowRuntimeBinding(BaseModel): @@ -95,6 +103,8 @@ def resolve_runtime_metadata( metadata = _resolve_source_pool(node, node_id=resolved_node_id) elif _is_collection_output(node): metadata = _resolve_collection_output(node, node_id=resolved_node_id) + elif _is_data_operator_node(node): + metadata = _resolve_data_operator_node(node, node_id=resolved_node_id) elif _is_normalize_node(node): metadata = _resolve_normalize_node(node, node_id=resolved_node_id) elif _is_merge_node(node): @@ -332,6 +342,91 @@ def _resolve_normalize_node(node: WorkflowProjectNode, *, node_id: str) -> dict[ } +def _resolve_data_operator_node( + node: WorkflowProjectNode, + *, + node_id: str, +) -> dict[str, Any]: + operator_id = _read_string(node.params.get("operatorId")) + pack_version_provided = "packVersion" in node.params + requested_pack_version = _read_string(node.params.get("packVersion")) + resolved_pack_version = ( + requested_pack_version or _LEGACY_DATA_OPERATOR_PACK_VERSION + ) + spec = ( + resolve_data_operator(operator_id, resolved_pack_version) + if operator_id and (requested_pack_version or not pack_version_provided) + else None + ) + known_spec = resolve_data_operator(operator_id) if operator_id else None + catalog_id = _read_string((node.ui or {}).get("catalogId")) + binding_id = DATA_OPERATOR_CATALOG_BINDINGS.get(catalog_id or "") + catalog_kind = binding_id.rsplit(".", 1)[-1] if binding_id else None + if spec is None or spec.kind != catalog_kind: + code = ( + "missing_data_operator_id" + if operator_id is None + else "unsupported_data_operator_version" + if pack_version_provided and known_spec is not None + else "unknown_data_operator" + if spec is None + else "data_operator_kind_mismatch" + ) + return { + "missing_runtime": _dump_missing_runtime( + WorkflowMissingRuntime( + code=code, + node_id=node_id, + kind=node.kind, + capability=node.capability, + required_params=["operatorId"], + message=( + "Data operator node requires a registered operatorId matching " + f'catalog kind "{catalog_kind or "unknown"}".' + ), + ) + ) + } + + flat_config = { + key: value + for key, value in node.params.items() + if key not in {"operatorId", "packVersion", "config"} + } + nested_config = node.params.get("config") + config: Any = ( + {**flat_config, **nested_config} + if isinstance(nested_config, dict) + else flat_config + if nested_config is None + else nested_config + ) + return { + "binding": { + "status": "bound", + "binding_id": binding_id, + "runtime": "workflow", + "channel": "data-operator", + "input": { + "operatorId": operator_id, + "operatorKind": spec.kind, + "packId": spec.pack_id, + "packVersion": spec.pack_version, + "config": config, + "inputPort": "recordCandidate[]", + "outputPort": "recordCandidate[]", + }, + }, + "data_operator": { + "node_id": node_id, + "operator_id": operator_id, + "operator_kind": spec.kind, + "pack_id": spec.pack_id, + "pack_version": spec.pack_version, + }, + } + + def _resolve_merge_node(node: WorkflowProjectNode, *, node_id: str) -> dict[str, Any]: strategy = _read_string(node.params.get("strategy")) or "concat" lineage = node.params.get("preserveLineage") @@ -798,6 +893,16 @@ def _is_normalize_node(node: WorkflowProjectNode) -> bool: ) +def _is_data_operator_node(node: WorkflowProjectNode) -> bool: + return _data_operator_kind(node) is not None + + +def _data_operator_kind(node: WorkflowProjectNode) -> str | None: + catalog_id = _read_string((node.ui or {}).get("catalogId")) + binding_id = DATA_OPERATOR_CATALOG_BINDINGS.get(catalog_id or "") + return binding_id.rsplit(".", 1)[-1] if binding_id else None + + def _is_merge_node(node: WorkflowProjectNode) -> bool: return _read_string((node.ui or {}).get("catalogId")) == "intelligence.flow.merge" or ( node.kind == "flow" and node.capability == "merge" diff --git a/docs/dataflow-compatibility-matrix.md b/docs/dataflow-compatibility-matrix.md new file mode 100644 index 0000000..e49078a --- /dev/null +++ b/docs/dataflow-compatibility-matrix.md @@ -0,0 +1,54 @@ +# DataFlow compatibility matrix + +Baseline: `OpenDCAI/DataFlow@f62aa1349e0ff14cb737a4cbda1945d04fde85bb`. + +OpenCLI Admin preserves DataFlow source identity at import time and compiles it +to a versioned native Data Operator invocation. The compatibility unit is an +upstream class plus its pinned source SHA, not an unversioned class name. + +## Runnable exact compatibility + +`builtin.text-cleaning@1.1.0` currently provides exact, dependency-free +compatibility for 34 upstream classes behind three native operators. + +| Native operator | Pinned DataFlow classes | +| --- | --- | +| `text.clean` | `RemoveExtraSpacesRefiner`, `LowercaseRefiner`, `HtmlUrlRemoverRefiner`, `HtmlEntityRefiner`, `RemoveEmojiRefiner`, `RemoveNumberRefiner`, `RemovePunctuationRefiner`, `RemoveRepetitionsPunctuationRefiner`, `TextNormalizationRefiner`, `RemoveImageRefsRefiner` | +| `text.rule-filter` | `ContentNullFilter`, `WordNumberFilter`, `SentenceNumberFilter`, `CharNumberFilter`, `UniqueWordsFilter`, `ColonEndFilter`, `LineEndWithEllipsisFilter`, `SymbolWordRatioFilter`, `AlphaWordsFilter`, `HtmlEntityFilter`, `IDCardFilter`, `NoPuncFilter`, `SpecialCharacterFilter`, `WatermarkFilter`, `MeanWordLengthFilter`, `StopWordFilter`, `CurlyBracketFilter`, `CapitalWordsFilter`, `LoremIpsumFilter`, `LineStartWithBulletpointFilter`, `LineWithJavascriptFilter`, `BlocklistFilter` | +| `text.deduplicate` | `HashDeduplicateFilter`, `NgramHashDeduplicateFilter` | + +The authoritative alias list and config translation live in +`backend/workflow/dataflow_compat.py`. Golden outputs and upstream source or +asset digests live in `tests/fixtures/dataflow/`. + +## Inventoried but unavailable until their real dependency exists + +| Family | Examples | Required adapter or asset | +| --- | --- | --- | +| Tokenizer or statistical text | `MinHashDeduplicateFilter`, `SimHashDeduplicateFilter`, `NgramFilter`, stop-word/stemming/lemmatization refiners | pinned tokenizer/library pack | +| Local language and quality models | `LanguageFilter`, `SemDeduplicateFilter`, Text-PT and Text-SFT quality filters | model artifact plus model digest | +| PII and entity processing | `NERRefiner`, `PIIAnonymizeRefiner`, `PresidioFilter` | NER/Presidio model adapter | +| Prompted cleaning and evaluation | `PromptedFilter`, `PromptedRefiner`, KBC text cleaner, model-judge filters | configured LLM provider | +| Document extraction | MinerU, PDF-to-Markdown, OCR and VQA extraction | document conversion or modality adapter | +| Retrieval and synthesis | embeddings, retrieval, reranking, multi-hop QA and RAG | embedding, vector-index and model adapters | +| Code and SQL execution | sandbox evaluators and SQL execution filters | isolated execution adapter | + +Unavailable classes fail closed during DataFlow import. They are never mapped +to a merely similar deterministic implementation. + +## Pipeline reproduction + +The native import endpoint accepts `runtime: "dataflow"` with the pinned +`graph.sourceSha`. Each graph node supplies either the complete source identity +or `module + class`, plus constructor and run configuration. The importer +produces ordinary typed Workflow nodes with: + +- canonical `operatorId`; +- exact `packVersion`; +- translated native config; +- safe DataFlow provenance. + +Canvas still uses four reusable execution kinds: `generate`, `filter`, +`evaluate`, and `refine`. Imported upstream classes remain separately visible +as workflow nodes and provenance entries without creating 194 shallow runtime +implementations. diff --git a/docs/dataflow-operator-packs.md b/docs/dataflow-operator-packs.md new file mode 100644 index 0000000..492a131 --- /dev/null +++ b/docs/dataflow-operator-packs.md @@ -0,0 +1,93 @@ +# DataFlow operator packs + +OpenCLI Admin internalizes DataFlow as versioned workflow capabilities, not as +a copy of every upstream Python class. One native operator may cover a family +of upstream fine-grained operators when their runtime and port contracts are +the same. + +Upstream baseline: `OpenDCAI/DataFlow@f62aa1349e0ff14cb737a4cbda1945d04fde85bb`. + +## Runnable v1 scope + +All v1 operators use `recordCandidate[] -> recordCandidate[]`, preserve source +lineage, do not mutate their input, and return metrics plus stable rejected +candidate IDs. + +| Pack | Native operator | DataFlow capability family | +| --- | --- | --- | +| `builtin.core-data@1.0.0` | `core.generate.instruction-pairs` | deterministic instruction-pair shaping | +| | `core.filter.quality` | basic quality gate | +| | `core.evaluate.quality` | basic quality scoring | +| | `core.refine.text` | basic text normalization | +| `builtin.text-cleaning@1.0.0` | `text.clean` | HTML/entity/URL/emoji/whitespace and configurable text cleanup | +| | `text.rule-filter` | non-empty, length, ratio and blocklist rules | +| | `text.deduplicate` | exact and SimHash duplicate removal | +| | `text.statistics` | character, word, sentence and lexical-diversity metrics | +| `builtin.dataset-preparation@1.0.0` | `data.project` | select, rename, coalesce and scalar cast | +| | `data.chunk` | deterministic chunking with overlap | +| | `data.qa-extract` | source-grounded existing QA-pair extraction | +| | `data.training-format` | deterministic Alpaca/ShareGPT-style shaping | + +`builtin.text-cleaning@1.1.0` is the pinned DataFlow compatibility profile. +It keeps the three deep native operators `text.clean`, `text.rule-filter`, and +`text.deduplicate`, while reproducing 34 SHA-locked upstream classes: + +- `RemoveExtraSpacesRefiner`, `LowercaseRefiner`, `HtmlUrlRemoverRefiner`, + `HtmlEntityRefiner`, `RemoveEmojiRefiner`, `RemoveNumberRefiner`, + `RemovePunctuationRefiner`, and `RemoveRepetitionsPunctuationRefiner`; +- `ContentNullFilter`, `WordNumberFilter`, `SentenceNumberFilter`, + `CharNumberFilter`, and `UniqueWordsFilter`; +- the remaining deterministic CPU rule filters, including colon/ellipsis, + symbol and alpha ratios, HTML entity, ID-card, punctuation, watermark, + word-length, stop-word, bracket, capitalization, Lorem Ipsum, bullet-point, + and JavaScript-line checks; +- `TextNormalizationRefiner`, `RemoveImageRefsRefiner`, and `BlocklistFilter`; +- `HashDeduplicateFilter` and `NgramHashDeduplicateFilter` (`md5`/`sha256`). + +DataFlow import identities use the complete form +`dataflow@::dataflow.operators..`. Short or unversioned +aliases are rejected. Imported workflows persist the native `operatorId`, +`packVersion`, translated config, and safe upstream provenance. The runtime +never imports DataFlow or pandas. + +The Canvas keeps four execution primitives (`generate`, `filter`, `evaluate`, +`refine`) and selects an operator from the backend capability manifest. This +keeps the node catalog usable while retaining a versioned, queryable operator +identity in every compiled workflow. + +## Resource-dependent packs + +The following DataFlow families are inventoried but are not declared runnable +until their real resource adapter exists: + +- prompted generation/refinement/evaluation, multi-hop QA, reasoning and + Text2SQL: model-provider adapter; +- embedding, retrieval, reranking and vector index writes: embedding/vector + resource contracts; +- PDF/URL-to-Markdown and MinerU transforms: document conversion adapter; +- OCR, speech and vision operators: modality-specific runtime adapters; +- Perspective, Presidio, LangKit and heavyweight semantic evaluators: explicit + dependency/service adapters. + +These capabilities must fail closed as unavailable. A deterministic placeholder +must never be reported as RAG synthesis, OCR, embedding, or semantic evaluation. + +## Acceptance contract + +Each runnable operator must pass: + +1. direct happy-path, empty-batch, invalid-config, determinism and input + immutability checks; +2. compile-time rejection of missing, unknown, or kind-mismatched `operatorId`; + explicit unsupported pack versions and unpinned DataFlow aliases also fail + closed; +3. runtime events carrying pack/version, lineage, metrics and bounded rejected + IDs without raw source bodies; +4. an end-to-end workflow from fixture source through normalize, preparation, + cleaning, acceptance and record sink; +5. demand-draft projection and Canvas selection without hand-editing the graph. + +Pinned compatibility additionally locks upstream source digests and golden +outputs under `tests/fixtures/dataflow/`. A dependency-capable job may run the +isolated upstream oracle; normal project tests validate the same golden +contract without adding DataFlow to application dependencies. diff --git a/frontend/components/flow/inspector.tsx b/frontend/components/flow/inspector.tsx index 265de46..a497372 100644 --- a/frontend/components/flow/inspector.tsx +++ b/frontend/components/flow/inspector.tsx @@ -3,7 +3,7 @@ import { useState } from "react" import Link from "next/link" import { AlertTriangle, PlugZap } from "lucide-react" -import { useFlowStore } from "@/lib/flow/store" +import { clearParameterDraftEntry, useFlowStore } from "@/lib/flow/store" import type { WorkflowNodeData, FieldConfig } from "@/lib/flow/types" import { Input } from "@/components/ui/input" import { Label } from "@/components/ui/label" @@ -18,7 +18,11 @@ import { } from "@/components/ui/select" import type { NodeInternalStatus } from "@/lib/workflow/node-internals" import { getNodeTemplate } from "@/lib/workflow/node-templates" -import { buildParameterInterfaceView, type ParameterInterfaceViewField } from "@/lib/workflow/parameter-interface" +import { + buildParameterInterfaceView, + parseJsonParameterValue, + type ParameterInterfaceViewField, +} from "@/lib/workflow/parameter-interface" import { blockedActionViewForRuntime } from "@/lib/workflow/capabilities" import { buildCanonicalNodeViewContract } from "@/lib/workflow/canonical-node-contract" import { MonoRow, PanelShell, SectionCaption } from "./inspector-shell" @@ -111,6 +115,8 @@ export function Inspector() { const onEdgesChange = useFlowStore((s) => s.onEdgesChange) const [nodeTab, setNodeTab] = useState<"config" | "prompt" | "run" | "trace">("config") const [parameterGroupTab, setParameterGroupTab] = useState("") + const [jsonDrafts, setJsonDrafts] = useState>({}) + const [jsonErrors, setJsonErrors] = useState>({}) const selected = nodes.filter((n) => n.selected) const selectedEdges = edges.filter((e) => e.selected) @@ -229,6 +235,15 @@ export function Inspector() { const updateParameterField = (field: ParameterInterfaceViewField, value: unknown) => { if (field.readonly) return + if (field.binding.source === "params" && field.binding.fieldId === "operatorId") { + const configField = parameterInterfaceView?.fields.find( + (candidate) => candidate.binding.source === "params" && candidate.binding.fieldId === "config", + ) + if (configField) { + setJsonDrafts((drafts) => clearParameterDraftEntry(drafts, node.id, configField.id)) + setJsonErrors((errors) => clearParameterDraftEntry(errors, node.id, configField.id)) + } + } if (parameterInterfaceView?.mode === "template") { if (field.binding.source === "adapter") { if (field.binding.fieldId === "mode") { @@ -277,6 +292,47 @@ export function Inspector() { ) + if (field.type === "json") { + const draftKey = `${node.id}:${field.id}` + const value = jsonDrafts[draftKey] ?? formatJsonParameterValue(raw) + const error = jsonErrors[draftKey] + return row( +
+