| import re |
| from copy import deepcopy |
| from typing import Any |
|
|
| from app.models.workflow import ( |
| ContractResponse, |
| DependencyImpactResponse, |
| DocumentationResponse, |
| EnvironmentPromotionResponse, |
| IntentDriftResponse, |
| LineageField, |
| LineageResponse, |
| QualityResponse, |
| ReleasePlanRequest, |
| ReleasePlanResponse, |
| RoiResponse, |
| SelfHealResponse, |
| ValidationIssue, |
| WebhookInspectResponse, |
| WorkflowDocument, |
| WorkflowPackageRequest, |
| WorkflowPackageResponse, |
| ) |
| from app.services.operations import estimate_cost |
| from app.services.optimizer import optimizer |
| from app.services.validation import validator |
|
|
|
|
| FIELD_REF = re.compile(r"\$json\.([A-Za-z_][\w.]*)") |
| ENV_REF = re.compile(r"\$\{ENV\.([A-Z][A-Z0-9_]*)\}") |
| WORD = re.compile(r"[a-z][a-z0-9_-]{2,}") |
| SECRET_NAMES = {"password", "secret", "token", "api_key", "apikey", "authorization"} |
| PERSONAL_NAMES = {"email", "phone", "address", "name", "ssn", "dob"} |
| FINANCIAL_NAMES = {"amount", "invoice", "card", "iban", "account", "price"} |
| STOP_WORDS = {"that", "with", "from", "this", "then", "into", "workflow", "create", "build"} |
|
|
|
|
| def _walk(value: Any, prefix: str = ""): |
| if isinstance(value, dict): |
| for key, nested in value.items(): |
| path = f"{prefix}.{key}" if prefix else str(key) |
| yield path, nested |
| yield from _walk(nested, path) |
| elif isinstance(value, list): |
| for index, nested in enumerate(value): |
| yield from _walk(nested, f"{prefix}[{index}]") |
|
|
|
|
| def _classification(field: str) -> str: |
| parts = {part.lower() for part in re.split(r"[.\[\]_-]+", field) if part} |
| if parts & SECRET_NAMES: |
| return "secret" |
| if parts & FINANCIAL_NAMES: |
| return "financial" |
| if parts & PERSONAL_NAMES: |
| return "personal" |
| return "internal" |
|
|
|
|
| def lineage(workflow: WorkflowDocument) -> LineageResponse: |
| dependencies: dict[str, list[str]] = {node.id: [] for node in workflow.nodes} |
| fields: dict[str, dict[str, set[str] | str]] = {} |
| labels = {node.id: node.data.label for node in workflow.nodes} |
| for edge in workflow.edges: |
| if edge.target in dependencies and edge.source in labels: |
| dependencies[edge.target].append(labels[edge.source]) |
| for node in workflow.nodes: |
| for path, value in _walk(node.data.parameters): |
| candidates = {path} |
| if isinstance(value, str): |
| candidates.update(FIELD_REF.findall(value)) |
| for field in candidates: |
| record = fields.setdefault( |
| field, |
| {"sources": set(), "consumers": set(), "classification": _classification(field)}, |
| ) |
| if field == path: |
| record["sources"].add(node.data.label) |
| else: |
| record["consumers"].add(node.data.label) |
| result = [ |
| LineageField( |
| field=field, |
| source_nodes=sorted(record["sources"]), |
| consumer_nodes=sorted(record["consumers"]), |
| classification=record["classification"], |
| ) |
| for field, record in sorted(fields.items()) |
| ] |
| return LineageResponse( |
| fields=result, |
| node_dependencies=dependencies, |
| sensitive_paths=[field.field for field in result if field.classification in {"secret", "personal", "financial"}], |
| ) |
|
|
|
|
| def _type_name(value: Any) -> str: |
| if value is None: |
| return "null" |
| if isinstance(value, bool): |
| return "boolean" |
| if isinstance(value, (int, float)): |
| return "number" |
| if isinstance(value, dict): |
| return "object" |
| if isinstance(value, list): |
| return "array" |
| return "string" |
|
|
|
|
| def contract_check(sample: dict[str, Any], expected: dict[str, str]) -> ContractResponse: |
| inferred = {path: _type_name(value) for path, value in _walk(sample) if not isinstance(value, (dict, list))} |
| violations = [] |
| for path, expected_type in expected.items(): |
| actual = inferred.get(path) |
| if actual is None: |
| violations.append(f"Missing required field: {path}") |
| elif actual != expected_type: |
| violations.append(f"{path} expected {expected_type}, received {actual}") |
| return ContractResponse(valid=not violations, inferred_schema=inferred, violations=violations) |
|
|
|
|
| def quality(workflow: WorkflowDocument) -> QualityResponse: |
| validation = validator.validate(workflow) |
| findings = list(validation.issues) |
| lineage_result = lineage(workflow) |
| code_nodes = [node for node in workflow.nodes if node.data.type == "n8n-nodes-base.code"] |
| if code_nodes: |
| findings.append(ValidationIssue(code="review_code", severity="warning", message="Code nodes require review and sandboxed testing.")) |
| if lineage_result.sensitive_paths: |
| findings.append(ValidationIssue(code="sensitive_data", severity="warning", message="Sensitive fields require destination and retention review.")) |
| retry_nodes = sum("retry" in str(node.data.parameters).lower() for node in workflow.nodes) |
| scores = { |
| "reliability": max(0, validation.score - (10 if retry_nodes == 0 else 0)), |
| "maintainability": max(0, 100 - len(code_nodes) * 12 - max(0, len(workflow.nodes) - 30)), |
| "security": max(0, 100 - len(lineage_result.sensitive_paths) * 4 - len(code_nodes) * 8), |
| "cost": max(0, 100 - sum("langchain" in node.data.type for node in workflow.nodes) * 7), |
| "observability": 90 if workflow.settings.get("saveExecutionProgress") else 55, |
| } |
| return QualityResponse(overall=round(sum(scores.values()) / len(scores)), scores=scores, findings=findings) |
|
|
|
|
| def intent_drift(workflow: WorkflowDocument, requirement: str) -> IntentDriftResponse: |
| terms = {term for term in WORD.findall(requirement.lower()) if term not in STOP_WORDS} |
| workflow_text = " ".join( |
| f"{node.data.label} {node.data.type} {node.data.subtitle or ''}" for node in workflow.nodes |
| ).lower() |
| covered = sorted(term for term in terms if term in workflow_text) |
| missing = sorted(terms - set(covered)) |
| return IntentDriftResponse( |
| alignment_score=round(len(covered) / max(1, len(terms)) * 100), |
| covered_terms=covered, |
| missing_terms=missing, |
| ) |
|
|
|
|
| def promote(workflow: WorkflowDocument, environment: str, values: dict[str, Any]) -> EnvironmentPromotionResponse: |
| promoted = deepcopy(workflow) |
| replacements = 0 |
| unresolved: set[str] = set() |
|
|
| def replace(value: Any) -> Any: |
| nonlocal replacements |
| if isinstance(value, str): |
| def substitute(match: re.Match[str]) -> str: |
| nonlocal replacements |
| key = match.group(1) |
| if key not in values: |
| unresolved.add(key) |
| return match.group(0) |
| replacements += 1 |
| return str(values[key]) |
| return ENV_REF.sub(substitute, value) |
| if isinstance(value, dict): |
| return {key: replace(nested) for key, nested in value.items()} |
| if isinstance(value, list): |
| return [replace(nested) for nested in value] |
| return value |
|
|
| for node in promoted.nodes: |
| node.data.parameters = replace(node.data.parameters) |
| promoted.meta.model_extra["environment"] = environment |
| return EnvironmentPromotionResponse( |
| workflow=promoted, environment=environment, replacements=replacements, unresolved=sorted(unresolved) |
| ) |
|
|
|
|
| def release_plan(request: ReleasePlanRequest) -> ReleasePlanResponse: |
| validation = validator.validate(request.workflow) |
| blocked = not validation.valid |
| strategy_steps = { |
| "shadow": ["Mirror sanitized input to the candidate version", "Suppress side-effect nodes", "Compare schemas, outputs, and latency"], |
| "canary": [f"Route {request.traffic_percentage}% of eligible traffic to the candidate", "Monitor success rate and latency", "Increase traffic only after approval"], |
| "synthetic": ["Schedule representative test inputs", "Verify credentials and contracts", "Alert on consecutive failures"], |
| } |
| return ReleasePlanResponse( |
| strategy=request.strategy, |
| status="blocked" if blocked else "draft", |
| steps=strategy_steps[request.strategy], |
| rollback_conditions=[ |
| f"Error rate exceeds {request.max_error_rate:.1%}", |
| f"Success rate drops below {request.success_threshold:.1%}", |
| "Output contract or sensitive-data policy fails", |
| ], |
| warnings=[issue.message for issue in validation.issues if issue.severity == "error"], |
| ) |
|
|
|
|
| def package_workflow(request: WorkflowPackageRequest) -> WorkflowPackageResponse: |
| return WorkflowPackageResponse( |
| manifest={ |
| "format": "flowforge.workflow-package/v1", |
| "name": request.workflow.name, |
| "workflowVersion": request.workflow.meta.version or 1, |
| "n8nCompatible": True, |
| "contents": ["workflow", "tests", "contracts", "environments"], |
| }, |
| workflow=request.workflow, |
| tests=request.tests, |
| contracts=request.contracts, |
| environments=request.environments, |
| ) |
|
|
|
|
| def documentation(workflow: WorkflowDocument) -> DocumentationResponse: |
| nodes = "\n".join(f"- **{node.data.label}** (`{node.data.type}`): {node.data.subtitle or 'Configured action'}" for node in workflow.nodes) |
| connections = "\n".join( |
| f"- {next((n.data.label for n in workflow.nodes if n.id == edge.source), edge.source)} -> " |
| f"{next((n.data.label for n in workflow.nodes if n.id == edge.target), edge.target)}" |
| for edge in workflow.edges |
| ) |
| return DocumentationResponse(markdown=f"# {workflow.name}\n\n{workflow.meta.description or 'n8n workflow'}\n\n## Nodes\n\n{nodes}\n\n## Connections\n\n{connections or '- None'}\n") |
|
|
|
|
| def roi(workflow: WorkflowDocument, executions: int, minutes_saved: float, hourly_rate: float, sla_minutes: float) -> RoiResponse: |
| hours = executions * minutes_saved / 60 |
| cost = estimate_cost(workflow, executions).estimated_monthly_usd |
| estimated_duration = max(1, len(workflow.nodes) * 250) |
| return RoiResponse( |
| hours_saved=round(hours, 2), labor_value_usd=round(hours * hourly_rate, 2), |
| estimated_operating_cost_usd=cost, net_value_usd=round(hours * hourly_rate - cost, 2), |
| estimated_duration_ms=estimated_duration, |
| sla_headroom_percent=round(max(0, 1 - estimated_duration / (sla_minutes * 60_000)) * 100, 2), |
| ) |
|
|
|
|
| def inspect_webhook(payload: dict[str, Any], redact: bool) -> WebhookInspectResponse: |
| result = deepcopy(payload) |
| schema = {path: _type_name(value) for path, value in _walk(payload) if not isinstance(value, (dict, list))} |
| redacted = [] |
| if redact: |
| def scrub(value: Any, prefix: str = "") -> None: |
| if not isinstance(value, dict): |
| return |
| for key, nested in value.items(): |
| path = f"{prefix}.{key}" if prefix else key |
| if _classification(path) in {"secret", "personal", "financial"}: |
| value[key] = "[REDACTED]" |
| redacted.append(path) |
| else: |
| scrub(nested, path) |
| scrub(result) |
| return WebhookInspectResponse(payload=result, schema_map=schema, redacted_fields=redacted) |
|
|
|
|
| def dependency_impact(workflow: WorkflowDocument, dependency: str) -> DependencyImpactResponse: |
| query = dependency.lower() |
| affected = [node for node in workflow.nodes if query in f"{node.data.type} {node.data.label} {node.data.parameters}".lower()] |
| affected_ids = {node.id for node in affected} |
| downstream_ids: set[str] = set() |
| frontier = list(affected_ids) |
| while frontier: |
| source = frontier.pop() |
| for edge in workflow.edges: |
| if edge.source == source and edge.target not in downstream_ids: |
| downstream_ids.add(edge.target) |
| frontier.append(edge.target) |
| labels = {node.id: node.data.label for node in workflow.nodes} |
| count = len(affected_ids | downstream_ids) |
| severity = "none" if count == 0 else "low" if count == 1 else "medium" if count < 5 else "high" |
| return DependencyImpactResponse( |
| affected_nodes=[node.data.label for node in affected], |
| downstream_nodes=[labels[node_id] for node_id in downstream_ids if node_id in labels], |
| severity=severity, |
| ) |
|
|
|
|
| def self_heal(workflow: WorkflowDocument, errors: list[str]) -> SelfHealResponse: |
| before = quality(workflow) |
| optimized = optimizer.optimize(workflow) |
| after = quality(optimized.workflow) |
| changes = [suggestion.title for suggestion in optimized.suggestions] |
| if errors: |
| changes.append("Attached execution errors as approval evidence") |
| return SelfHealResponse( |
| proposed_workflow=optimized.workflow, |
| changes=changes, |
| quality_before=before.overall, |
| quality_after=after.overall, |
| ) |
|
|