diff --git a/Cargo.lock b/Cargo.lock index a89cadfd9f4..6ab1a28262d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -355,6 +355,9 @@ dependencies = [ "daglang-derive", "daglang-driver", "daglang-lower", + "daglang-resolve", + "daglang-syntax", + "daglang-typecheck", "gunbc-cli", "gunbc-exec", "gunbc-infra", @@ -391,6 +394,7 @@ dependencies = [ "gunbc-lib-review", "gunbc-lib-transport", "gunbc-primitives", + "gunbc-resolve", "gunbc-test", "gunbc-testgen-registry", "gunbc-testgen-registry-macros", @@ -576,6 +580,16 @@ dependencies = [ "toml", ] +[[package]] +name = "gunbc-resolve" +version = "0.1.0" +dependencies = [ + "daglang-lower", + "gunbc-exec", + "gunbc-ir", + "serde_json", +] + [[package]] name = "gunbc-test" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 19e0899f380..9ac9867d026 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,6 +6,7 @@ members = [ "core/ir", "core/exec", "core/cli", + "core/resolve", "core/codegen", "core/test", "core/testgen-registry", diff --git a/core/codegen/Cargo.toml b/core/codegen/Cargo.toml index bde2335000c..a4f5ac21992 100644 --- a/core/codegen/Cargo.toml +++ b/core/codegen/Cargo.toml @@ -14,4 +14,7 @@ gunbc-cli = { path = "../cli" } daglang-derive = { path = "../daglang/daglang-derive" } daglang-driver = { path = "../daglang/daglang-driver" } daglang-lower = { path = "../daglang/daglang-lower" } +daglang-resolve = { path = "../daglang/daglang-resolve" } +daglang-syntax = { path = "../daglang/daglang-syntax" } +daglang-typecheck = { path = "../daglang/daglang-typecheck" } serde_json = { workspace = true } diff --git a/core/codegen/src/cli_gen.rs b/core/codegen/src/cli_gen.rs index 61b4597e88a..d93019cea41 100644 --- a/core/codegen/src/cli_gen.rs +++ b/core/codegen/src/cli_gen.rs @@ -652,9 +652,6 @@ fn generate_input_mocks(_entrypoints: &[CliEntrypoint]) -> String { ); code.push_str(" }\n"); code.push_str("}\n"); - // RT58: propagate entrypoint values to param_source_* nodes - code.push_str("input_mocks.propagate_to_param_sources(&entrypoints.entrypoint_ports);\n"); - code } diff --git a/core/codegen/src/fidelity.rs b/core/codegen/src/fidelity.rs index 49ed4ecdfcc..c98e963c182 100644 --- a/core/codegen/src/fidelity.rs +++ b/core/codegen/src/fidelity.rs @@ -9,10 +9,13 @@ //! Follows the proven pattern from `pragma/dsl_render.rs`. use std::collections::{BTreeMap, HashMap}; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; +use std::sync::OnceLock; use daglang_derive::CallableProperties; -use daglang_driver::{compile_from_context, DriverContext}; +use daglang_resolve::{ModuleGraph, ResolvedModule}; +use daglang_syntax::ast::ModulePath; +use daglang_syntax::parser; use daglang_lower::{CallableKind, LoweredFnBody, LoweredOp, ServiceTransportClass}; use gunbc_ir::node::NodeBody; use gunbc_ir::Value; @@ -26,19 +29,96 @@ pub struct FidelityClassification { pub hermetic: bool, } -/// Compile `std/fidelity.dag` and extract fn bodies (including transitive -/// imports from `std/fermi.dag`). -fn compile_fidelity() -> HashMap { - let dsl_root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../dsl"); - let dag_file = dsl_root.join("std/fidelity.dag"); - let context = DriverContext { - roots: vec![dsl_root], - target_file: Some(dag_file), - }; - let output = compile_from_context(&context).expect("std/fidelity.dag should compile"); +const STD_TYPES_SOURCE: &str = include_str!("../../../dsl/std/types.dag"); +const STD_FERMI_SOURCE: &str = include_str!("../../../dsl/std/fermi.dag"); +const STD_FIDELITY_SOURCE: &str = include_str!("../../../dsl/std/fidelity.dag"); + +/// Embedded stdlib evaluator host with one-time compilation cache. +struct StdLibHost { + fns: HashMap, +} + +impl StdLibHost { + fn global() -> &'static Self { + static HOST: OnceLock = OnceLock::new(); + HOST.get_or_init(|| StdLibHost { + fns: compile_stdlib_fns(), + }) + } + + fn eval_fn(&self, fn_name: &str, inputs: &HashMap) -> HashMap { + let body = self + .fns + .get(fn_name) + .unwrap_or_else(|| panic!("stdlib function `{fn_name}` not found")); + daglang_lower::eval::evaluate_fn_body(body, inputs, &self.fns) + .unwrap_or_else(|e| panic!("stdlib function `{fn_name}` failed to evaluate: {e}")) + } +} + +fn parse_embedded_module(path: &Path, source: &str) -> (ModulePath, Vec, daglang_syntax::ast::SourceFile) { + let ast = parser::parse_with_file_diagnostics(path, source) + .unwrap_or_else(|diags| panic!("failed to parse embedded stdlib module {path:?}: {diags:?}")); + let module_path = ast + .module_path + .as_ref() + .map(|mp| mp.node.clone()) + .unwrap_or_else(|| panic!("embedded stdlib module {path:?} is missing `module` declaration")); + let imports = ast + .imports + .iter() + .map(|imp| imp.node.path.clone()) + .collect::>(); + (module_path, imports, ast) +} + +fn build_embedded_stdlib_graph() -> ModuleGraph { + let modules = vec![ + (PathBuf::from("/std/types.dag"), STD_TYPES_SOURCE), + (PathBuf::from("/std/fermi.dag"), STD_FERMI_SOURCE), + (PathBuf::from("/std/fidelity.dag"), STD_FIDELITY_SOURCE), + ]; + + let mut parsed = Vec::new(); + for (path, source) in modules { + let (module_path, imports, ast) = parse_embedded_module(path.as_path(), source); + parsed.push((path, module_path, imports, ast)); + } + + let mut index_by_module = HashMap::new(); + for (idx, (_, module_path, _, _)) in parsed.iter().enumerate() { + index_by_module.insert(module_path.clone(), idx); + } + + let mut resolved = Vec::new(); + for (path, module_path, imports, ast) in parsed { + let mut dependencies = Vec::new(); + for import in imports { + let dep = index_by_module.get(&import).copied().unwrap_or_else(|| { + panic!("embedded stdlib import `{}` missing from host graph", import.as_dotted()) + }); + dependencies.push(dep); + } + resolved.push(ResolvedModule { + path, + ast, + module_path, + dependencies, + }); + } + + ModuleGraph { modules: resolved } +} + +fn compile_stdlib_fns() -> HashMap { + let graph = build_embedded_stdlib_graph(); + let typed = daglang_typecheck::typecheck_module_graph(graph) + .unwrap_or_else(|errs| panic!("embedded stdlib typecheck failed: {errs:?}")); + let lowered = daglang_lower::lower_typed_project(&typed) + .unwrap_or_else(|e| panic!("embedded stdlib lowering failed: {e}")); let mut fns = HashMap::new(); - for node in &output.lowered_dag.nodes { + for node in &lowered.nodes { if let NodeBody::Opaque(LoweredOp::Callable { kind: CallableKind::Fn, name, @@ -49,20 +129,51 @@ fn compile_fidelity() -> HashMap { fns.insert(name.clone(), *body.clone()); } } - fns } /// Convert a `ServiceTransportClass` to the DSL `TransportClass` variant name. fn transport_class_to_dsl(tc: &ServiceTransportClass) -> Value { - Value::Str(match tc { - ServiceTransportClass::LocalDirect => "LocalDirect", - ServiceTransportClass::InterfaceStub => "InterfaceStub", - ServiceTransportClass::ShellLocal => "ShellLocal", - ServiceTransportClass::FileBoundary => "FileBoundary", - ServiceTransportClass::RestNetwork => "RestNetwork", - ServiceTransportClass::Unknown => "Unknown", - }.to_string()) + Value::Enum { + ty: "TransportClass".to_string(), + variant: match tc { + ServiceTransportClass::LocalDirect => "LocalDirect", + ServiceTransportClass::InterfaceStub => "InterfaceStub", + ServiceTransportClass::ShellLocal => "ShellLocal", + ServiceTransportClass::FileBoundary => "FileBoundary", + ServiceTransportClass::RestNetwork => "RestNetwork", + ServiceTransportClass::Unknown => "Unknown", + } + .to_string(), + } +} + +fn enum_variant(value: &Value) -> Option<&str> { + match value { + Value::Enum { variant, .. } => Some(variant.as_str()), + Value::Str(value) => Some(value.as_str()), + _ => None, + } +} + +fn decode_test_class(value: Option<&Value>) -> TestClass { + match value.and_then(enum_variant) { + Some("Unit") => TestClass::Unit, + Some("Hermetic") => TestClass::Hermetic, + Some("Integration") => TestClass::Integration, + other => panic!("classify_transports returned invalid test_class: {:?}", other), + } +} + +fn decode_fermi_cost(value: Option<&Value>) -> FermiCost { + match value.and_then(enum_variant) { + Some("Xs") => FermiCost::XS, + Some("S") => FermiCost::S, + Some("M") => FermiCost::M, + Some("L") => FermiCost::L, + Some("Xl") => FermiCost::XL, + other => panic!("classify_transports returned invalid depth: {:?}", other), + } } /// Classify a callable using DSL-evaluated `classify_transports()`. @@ -71,7 +182,7 @@ fn transport_class_to_dsl(tc: &ServiceTransportClass) -> Value { /// classification pipeline in DSL — including fold-based max depth /// and all-hermetic checks. pub fn classify_callable(props: &CallableProperties) -> FidelityClassification { - let fns = compile_fidelity(); + let stdlib = StdLibHost::global(); let transports: Vec = props .transport_classes @@ -82,32 +193,10 @@ pub fn classify_callable(props: &CallableProperties) -> FidelityClassification { let mut inputs = HashMap::new(); inputs.insert("transports".to_string(), Value::List(transports)); - let body = fns - .get("classify_transports") - .expect("classify_transports fn body should exist in std/fidelity.dag"); - let result = daglang_lower::eval::evaluate_fn_body(body, &inputs, &fns) - .expect("classify_transports should evaluate"); - - let test_class_value = result.get("test_class"); - let test_class = test_class_value - .and_then(Value::as_str) - .and_then(TestClass::parse) - .unwrap_or_else(|| { - panic!( - "classify_transports returned invalid test_class: {:?}", - test_class_value - ) - }); - let depth_value = result.get("depth"); - let fermi_cost = depth_value - .and_then(Value::as_str) - .and_then(FermiCost::parse) - .unwrap_or_else(|| { - panic!( - "classify_transports returned invalid depth: {:?}", - depth_value - ) - }); + let result = stdlib.eval_fn("classify_transports", &inputs); + + let test_class = decode_test_class(result.get("test_class")); + let fermi_cost = decode_fermi_cost(result.get("depth")); let hermetic_value = result.get("hermetic"); let hermetic = match hermetic_value { Some(Value::Bool(b)) => *b, diff --git a/core/codegen/src/testgen/codegen.rs b/core/codegen/src/testgen/codegen.rs index 55af8d1c888..9666542b041 100644 --- a/core/codegen/src/testgen/codegen.rs +++ b/core/codegen/src/testgen/codegen.rs @@ -270,7 +270,7 @@ impl ProbeObserverBundle { } } -impl<'a, T: Clone> TestGenerator<'a, T> { +impl<'a, T: Clone + 'static> TestGenerator<'a, T> { /// Create a new test generator for a DAG. pub fn new(dag: &'a Dag) -> Self { Self { @@ -5234,7 +5234,7 @@ struct CorpusExampleCtx<'a> { inputs: &'a HashMap, } -impl TestGenerator<'_, T> { +impl TestGenerator<'_, T> { fn build_corpus_section( &self, diff --git a/core/codegen/src/testgen/registry_gen.rs b/core/codegen/src/testgen/registry_gen.rs index 8fe25621292..0fe8d590c7d 100644 --- a/core/codegen/src/testgen/registry_gen.rs +++ b/core/codegen/src/testgen/registry_gen.rs @@ -8,10 +8,9 @@ //! # Design //! //! Transport nodes in the DAG follow the prepare→execute→parse triplet -//! pattern. The execute node's request type (`TransportRequest` variant) -//! determines the transport class. Since request types aren't directly -//! stored on nodes, we use the naming convention: node IDs contain -//! the service operation which maps to a transport class. +//! pattern. The execute node's request/response types determine the transport +//! class. For lowered DAGs, we read `ServiceTransportClass` metadata directly +//! from `LoweredOp`; type-based fallback is only for non-lowered test DAGs. //! //! The derived config is used by `build_fidelity_ladder_section()` to //! generate S-tier test code that installs a `VirtualTransportBackend` @@ -20,6 +19,7 @@ use std::collections::BTreeMap; use crate::testgen::analyze::DagAnalysis; +use daglang_lower::{classify_service_transport, LoweredOp, ServiceTransportClass}; use gunbc_ir::Dag; /// Transport class for a transport executor node. @@ -40,45 +40,42 @@ pub enum TransportClass { } impl TransportClass { - /// Classify a transport node by examining its port types and node ID. - fn from_node_context(node_id: &str, input_type: Option<&str>, output_type: Option<&str>) -> Self { - // Check explicit request/response types first. + /// Classify a transport node by examining explicit request/response types. + fn from_node_context(_node_id: &str, input_type: Option<&str>, output_type: Option<&str>) -> Self { match (input_type, output_type) { (Some("TcpRequest"), _) | (_, Some("TcpResponse")) => return Self::Tcp, + (Some("ShellRequest"), _) | (_, Some("ShellResponse")) => return Self::Shell, + (Some("FileRequest"), _) | (_, Some("FileResponse")) => return Self::File, + (Some("HttpRequest"), _) | (_, Some("HttpResponse")) => return Self::Http, + (Some("LocalRequest"), _) | (_, Some("LocalResponse")) => return Self::Local, _ => {} } - - // Infer from node ID naming convention. - // Transport triplet nodes follow: execute_ or - // service_transport::execute::.:: - let lower = node_id.to_lowercase(); - if lower.contains("rest") || lower.contains("http_stub") { - return Self::Rest; - } - if lower.contains("shell") - || lower.contains("cargo") - || lower.contains("git::") - || lower.contains("_git_") - || lower.starts_with("git_") - || lower.ends_with("_git") - { - return Self::Shell; - } - if lower.contains("file") || lower.contains("fs::") || lower.contains("content_upsert") { - return Self::File; - } - if lower.contains("tcp") { - return Self::Tcp; - } - if lower.contains("local") { - return Self::Local; - } - - // Default: REST (most common for service operations). + // Default: REST (most common transport class). Self::Rest } } +fn from_service_transport_class(class: ServiceTransportClass) -> TransportClass { + match class { + ServiceTransportClass::RestNetwork => TransportClass::Rest, + ServiceTransportClass::ShellLocal => TransportClass::Shell, + ServiceTransportClass::FileBoundary => TransportClass::File, + ServiceTransportClass::LocalDirect => TransportClass::Local, + ServiceTransportClass::InterfaceStub => TransportClass::Rest, + ServiceTransportClass::Unknown => TransportClass::Rest, + } +} + +fn transport_class_from_node_metadata(node: &gunbc_ir::Node) -> Option { + let gunbc_ir::node::NodeBody::Opaque(op) = &node.body else { + return None; + }; + let op_any = op as &dyn std::any::Any; + let lowered = op_any.downcast_ref::()?; + let class = classify_service_transport(lowered)?; + Some(from_service_transport_class(class)) +} + /// Info about a single transport executor node. #[derive(Debug, Clone)] pub struct TransportNodeInfo { @@ -127,7 +124,10 @@ impl VirtualBackendRequirements { pub fn derive_virtual_backend_requirements( dag: &Dag, analysis: &DagAnalysis, -) -> VirtualBackendRequirements { +) -> VirtualBackendRequirements +where + T: 'static, +{ let mut requirements = VirtualBackendRequirements::default(); // Build a quick lookup of node input/output types. @@ -140,6 +140,8 @@ pub fn derive_virtual_backend_requirements( (n.id.0.as_str(), (req_input, resp_output)) }) .collect(); + let node_by_id: BTreeMap<&str, &gunbc_ir::Node> = + dag.nodes.iter().map(|n| (n.id.0.as_str(), n)).collect(); for executor_id in &analysis.transport_executors { let (input_type, output_type) = node_types @@ -147,7 +149,10 @@ pub fn derive_virtual_backend_requirements( .copied() .unwrap_or((None, None)); - let transport_class = TransportClass::from_node_context(executor_id, input_type, output_type); + let transport_class = node_by_id + .get(executor_id.as_str()) + .and_then(|node| transport_class_from_node_metadata(*node)) + .unwrap_or_else(|| TransportClass::from_node_context(executor_id, input_type, output_type)); match transport_class { TransportClass::Rest | TransportClass::Http => requirements.needs_rest = true, @@ -174,7 +179,7 @@ mod tests { use gunbc_ir::{Dag, Node}; #[test] - fn classify_rest_from_node_id() { + fn classify_rest_from_unknown_types_defaults_to_rest() { assert_eq!( TransportClass::from_node_context("service_transport::execute::github.Gist::Create", None, None), TransportClass::Rest @@ -182,25 +187,17 @@ mod tests { } #[test] - fn classify_shell_from_node_id() { - assert_eq!( - TransportClass::from_node_context("execute_cargo_build", None, None), - TransportClass::Shell - ); + fn classify_shell_from_port_types() { assert_eq!( - TransportClass::from_node_context("service_transport::execute::git::Status", None, None), + TransportClass::from_node_context("execute_shell", Some("ShellRequest"), Some("ShellResponse")), TransportClass::Shell ); } #[test] - fn classify_file_from_node_id() { + fn classify_file_from_port_types() { assert_eq!( - TransportClass::from_node_context("execute_file_read", None, None), - TransportClass::File - ); - assert_eq!( - TransportClass::from_node_context("execute_content_upsert", None, None), + TransportClass::from_node_context("execute_file", Some("FileRequest"), Some("FileResponse")), TransportClass::File ); } @@ -242,8 +239,8 @@ mod tests { dag.add_node( Node::opaque( "execute_git_status", - vec![port("request", "TransportRequest")], - vec![port("response", "TransportResponse")], + vec![port("request", "ShellRequest")], + vec![port("response", "ShellResponse")], (), ) .with_kind(NodeKind::TransportExecute), @@ -251,8 +248,8 @@ mod tests { dag.add_node( Node::opaque( "execute_file_read", - vec![port("request", "TransportRequest")], - vec![port("response", "TransportResponse")], + vec![port("request", "FileRequest")], + vec![port("response", "FileResponse")], (), ) .with_kind(NodeKind::TransportExecute), diff --git a/core/codegen/src/testgen/render_rust.rs b/core/codegen/src/testgen/render_rust.rs index 084cf1050ef..02f49926974 100644 --- a/core/codegen/src/testgen/render_rust.rs +++ b/core/codegen/src/testgen/render_rust.rs @@ -529,6 +529,18 @@ impl RustCodeRenderer { format!("Value::Secret({})", inner) } } + ValueExpr::Enum { ty, variant } => { + let escaped_ty = escape_rust_str(ty); + let escaped_variant = escape_rust_str(variant); + if bare { + format!("\"{}\".to_string()", escaped_variant) + } else { + format!( + "Value::Enum {{ ty: \"{}\".to_string(), variant: \"{}\".to_string() }}", + escaped_ty, escaped_variant + ) + } + } ValueExpr::Skipped => "Value::Skipped".to_string(), } } diff --git a/core/daglang/daglang-driver/src/lib.rs b/core/daglang/daglang-driver/src/lib.rs index 2c9a4f5c1b9..e5d178bf0c5 100644 --- a/core/daglang/daglang-driver/src/lib.rs +++ b/core/daglang/daglang-driver/src/lib.rs @@ -779,6 +779,12 @@ fn collect_covered_stages_from_expr( collect_covered_stages_from_expr(lhs, producer_by_binding, covered); collect_covered_stages_from_expr(rhs, producer_by_binding, covered); } + Expr::PipeCall(receiver, _, args) => { + collect_covered_stages_from_expr(receiver, producer_by_binding, covered); + for (_name, arg_expr) in args { + collect_covered_stages_from_expr(arg_expr, producer_by_binding, covered); + } + } Expr::Lambda(_, body) => { collect_covered_stages_from_expr(body, producer_by_binding, covered); } @@ -865,6 +871,12 @@ fn collect_root_identifiers(expr: &Expr, roots: &mut std::collections::BTreeSet< collect_root_identifiers(lhs, roots); collect_root_identifiers(rhs, roots); } + Expr::PipeCall(receiver, _, args) => { + collect_root_identifiers(receiver, roots); + for (_name, arg_expr) in args { + collect_root_identifiers(arg_expr, roots); + } + } Expr::Lambda(_, body) => collect_root_identifiers(body, roots), Expr::List(items) => { for item in items { diff --git a/core/daglang/daglang-emit/src/computation.rs b/core/daglang/daglang-emit/src/computation.rs index 98a4cc34e61..67bd9837fed 100644 --- a/core/daglang/daglang-emit/src/computation.rs +++ b/core/daglang/daglang-emit/src/computation.rs @@ -434,8 +434,8 @@ fn classify_primitive( outputs, body: PureBody::Literal(serde_json::Value::Null), }), - // RT4a: Return expression compute nodes evaluate fn bodies at runtime. - PrimitiveOpKind::ReturnExprCompute { .. } => Ok(Computation::Pure { + // Expression compute nodes evaluate lowered fn bodies at runtime. + PrimitiveOpKind::ExprCompute { .. } => Ok(Computation::Pure { inputs, outputs, body: PureBody::Literal(serde_json::Value::Null), diff --git a/core/daglang/daglang-emit/src/fn_codegen.rs b/core/daglang/daglang-emit/src/fn_codegen.rs index 0f58a1556fa..6ea094c5b17 100644 --- a/core/daglang/daglang-emit/src/fn_codegen.rs +++ b/core/daglang/daglang-emit/src/fn_codegen.rs @@ -214,6 +214,10 @@ fn compile_expr(expr: &ast::Expr, ctx: &CompileContext) -> code_ir::Expr { } } ast::Expr::Pipe(left, right) => compile_pipe(left, right, ctx), + ast::Expr::PipeCall(receiver, method, args) => { + let rhs = ast::Expr::Call(pipe_method_name(*method).to_string(), args.clone()); + compile_pipe(receiver, &rhs, ctx) + } ast::Expr::StringInterp(parts) => compile_string_interp(parts, ctx), ast::Expr::For(binding, iter_expr, _passthrough, body) => { let result_var = fresh("for_result"); @@ -260,6 +264,35 @@ fn compile_literal(lit: &ast::Literal) -> code_ir::Expr { } } +fn pipe_method_name(method: ast::PipeMethod) -> &'static str { + match method { + ast::PipeMethod::Map => "map", + ast::PipeMethod::Filter => "filter", + ast::PipeMethod::FilterMap => "filter_map", + ast::PipeMethod::FlatMap => "flat_map", + ast::PipeMethod::SortBy => "sort_by", + ast::PipeMethod::Append => "append", + ast::PipeMethod::Fold => "fold", + ast::PipeMethod::Join => "join", + ast::PipeMethod::Count => "count", + ast::PipeMethod::Sum => "sum", + ast::PipeMethod::First => "first", + ast::PipeMethod::Last => "last", + ast::PipeMethod::MaxBy => "max_by", + ast::PipeMethod::Any => "any", + ast::PipeMethod::All => "all", + ast::PipeMethod::Contains => "contains", + ast::PipeMethod::StartsWith => "starts_with", + ast::PipeMethod::EndsWith => "ends_with", + ast::PipeMethod::Repeat => "repeat", + ast::PipeMethod::ReplaceSection => "replace_section", + ast::PipeMethod::Chars => "chars", + ast::PipeMethod::ToBytes => "to_bytes", + ast::PipeMethod::ToJson => "to_json", + ast::PipeMethod::Hash => "hash", + } +} + // --------------------------------------------------------------------------- // Identifiers — bare names resolve to enum-variant paths // --------------------------------------------------------------------------- diff --git a/core/daglang/daglang-emit/src/render_go.rs b/core/daglang/daglang-emit/src/render_go.rs index f82533f9ad4..b746c06322b 100644 --- a/core/daglang/daglang-emit/src/render_go.rs +++ b/core/daglang/daglang-emit/src/render_go.rs @@ -544,6 +544,7 @@ fn render_value_expr(expr: &ValueExpr) -> String { format!("{}{{ {} }}", name, field_strs.join(", ")) } ValueExpr::Secret(s) => format!("NewSecret(\"{}\")", escape_go_str(s)), + ValueExpr::Enum { ty: _, variant } => format!("\"{}\"", escape_go_str(variant)), ValueExpr::Skipped => "nil /* skipped */".to_string(), } } diff --git a/core/daglang/daglang-emit/src/render_rust.rs b/core/daglang/daglang-emit/src/render_rust.rs index 882ba63cb9c..1fad0b919c6 100644 --- a/core/daglang/daglang-emit/src/render_rust.rs +++ b/core/daglang/daglang-emit/src/render_rust.rs @@ -565,6 +565,9 @@ fn render_value_expr(expr: &ValueExpr) -> String { ValueExpr::Secret(s) => { format!("SecretString::new(\"{}\")", escape_rust_str(s)) } + ValueExpr::Enum { ty: _, variant } => { + format!("\"{}\".to_string()", escape_rust_str(variant)) + } ValueExpr::Skipped => "Value::Skipped".to_string(), } } diff --git a/core/daglang/daglang-emit/src/rust_exec_runtime.rs b/core/daglang/daglang-emit/src/rust_exec_runtime.rs index 8e3fd69bdb7..965f610dc4f 100644 --- a/core/daglang/daglang-emit/src/rust_exec_runtime.rs +++ b/core/daglang/daglang-emit/src/rust_exec_runtime.rs @@ -350,7 +350,7 @@ fn classify_handler(op: &LoweredOp) -> Option { .. } => return Some(HandlerClassification::MetadataOnly), LoweredOp::Primitive { - kind: PrimitiveOpKind::ReturnExprCompute { .. }, + kind: PrimitiveOpKind::ExprCompute { .. }, .. } => return Some(HandlerClassification::MetadataOnly), } diff --git a/core/daglang/daglang-lower/src/eval.rs b/core/daglang/daglang-lower/src/eval.rs index 0602ad01ac2..09821e8cf1e 100644 --- a/core/daglang/daglang-lower/src/eval.rs +++ b/core/daglang/daglang-lower/src/eval.rs @@ -257,8 +257,11 @@ fn eval_expr( LoweredExpr::VariantConstruct { tag, fields } => { if fields.is_empty() { - // Unit variant: `Closed` → Value::Str("Closed") - Ok(Value::Str(tag.clone())) + // Unit variant: `Closed` → Value::Enum { ty: "", variant: "Closed" } + Ok(Value::Enum { + ty: String::new(), + variant: tag.clone(), + }) } else { // Payload variant: `Ok { value: x }` → Map with _variant tag let mut map = BTreeMap::new(); @@ -586,7 +589,9 @@ fn match_pattern(pattern: &LoweredPattern, value: &Value) -> Option Some(vec![]), + // Backward-compat path for older snapshots. Value::Str(s) if s == variant_name => Some(vec![]), _ => None, } @@ -1004,6 +1009,7 @@ pub fn sort_key(value: &Value) -> String { Value::Request(request) => format!("req:{request:?}"), Value::Response(response) => format!("resp:{response:?}"), Value::Secret(secret) => format!("secret:{}", secret.len()), + Value::Enum { ty, variant } => format!("enum:{ty}:{variant}"), Value::Float(f) => format!("f:{f}"), Value::Bytes(b) => format!("bytes:{}", b.len()), Value::Skipped => "skipped".to_string(), @@ -1029,6 +1035,13 @@ pub fn value_to_string(value: &Value) -> String { Value::Request(request) => format!("{request:?}"), Value::Response(response) => format!("{response:?}"), Value::Secret(secret) => format!("secret({})", secret.len()), + Value::Enum { ty, variant } => { + if ty.is_empty() { + variant.clone() + } else { + format!("{ty}.{variant}") + } + } Value::Float(f) => f.to_string(), Value::Bytes(b) => format!("<{} bytes>", b.len()), } @@ -1046,6 +1059,7 @@ pub fn value_truthy(value: &Value) -> bool { Value::Json(json) => !json.is_null(), Value::Secret(secret) => !secret.is_empty(), Value::Bytes(b) => !b.is_empty(), + Value::Enum { .. } => true, Value::Skipped | Value::Unit => false, Value::Request(_) | Value::Response(_) => true, } diff --git a/core/daglang/daglang-lower/src/expr.rs b/core/daglang/daglang-lower/src/expr.rs index 5e76c8904af..155428a50be 100644 --- a/core/daglang/daglang-lower/src/expr.rs +++ b/core/daglang/daglang-lower/src/expr.rs @@ -16,6 +16,24 @@ pub struct LoweredFnBody { pub stmts: Vec, } +/// Typed reference to an expression leaf source used by lowerer wiring. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum LeafRef { + Param { + name: String, + field: Option, + ty: String, + }, + Callable { + endpoint: String, + port: String, + }, + Service { + endpoint: String, + port: String, + }, +} + /// A lowered statement. #[derive(Debug, Clone, PartialEq, Eq)] pub enum LoweredStmt { @@ -306,6 +324,16 @@ fn lower_expr(expr: &ast::Expr, variant_names: &HashSet) -> LoweredExpr receiver: Box::new(lower_expr(receiver, variant_names)), call: Box::new(lower_expr(call, variant_names)), }, + ast::Expr::PipeCall(receiver, method, args) => LoweredExpr::Pipe { + receiver: Box::new(lower_expr(receiver, variant_names)), + call: Box::new(LoweredExpr::Call { + name: pipe_method_name(*method).to_string(), + args: args + .iter() + .map(|(k, v)| (k.clone(), lower_expr(v, variant_names))) + .collect(), + }), + }, ast::Expr::Lambda(params, body) => LoweredExpr::Lambda { params: params.clone(), body: Box::new(lower_expr(body, variant_names)), @@ -392,6 +420,35 @@ fn lower_unaryop(op: &ast::UnaryOp) -> LoweredUnaryOp { } } +fn pipe_method_name(method: ast::PipeMethod) -> &'static str { + match method { + ast::PipeMethod::Map => "map", + ast::PipeMethod::Filter => "filter", + ast::PipeMethod::FilterMap => "filter_map", + ast::PipeMethod::FlatMap => "flat_map", + ast::PipeMethod::SortBy => "sort_by", + ast::PipeMethod::Append => "append", + ast::PipeMethod::Fold => "fold", + ast::PipeMethod::Join => "join", + ast::PipeMethod::Count => "count", + ast::PipeMethod::Sum => "sum", + ast::PipeMethod::First => "first", + ast::PipeMethod::Last => "last", + ast::PipeMethod::MaxBy => "max_by", + ast::PipeMethod::Any => "any", + ast::PipeMethod::All => "all", + ast::PipeMethod::Contains => "contains", + ast::PipeMethod::StartsWith => "starts_with", + ast::PipeMethod::EndsWith => "ends_with", + ast::PipeMethod::Repeat => "repeat", + ast::PipeMethod::ReplaceSection => "replace_section", + ast::PipeMethod::Chars => "chars", + ast::PipeMethod::ToBytes => "to_bytes", + ast::PipeMethod::ToJson => "to_json", + ast::PipeMethod::Hash => "hash", + } +} + fn lower_string_part(part: &ast::StringPart, variant_names: &HashSet) -> LoweredStringPart { match part { ast::StringPart::Literal(s) => LoweredStringPart::Literal(s.clone()), diff --git a/core/daglang/daglang-lower/src/lib.rs b/core/daglang/daglang-lower/src/lib.rs index d796a939498..20e7c76fe81 100644 --- a/core/daglang/daglang-lower/src/lib.rs +++ b/core/daglang/daglang-lower/src/lib.rs @@ -24,8 +24,7 @@ use daglang_syntax::ast::{ }; use daglang_syntax::ast_utils::{ canonical_resource_type_name, is_type_expr_optional, resource_type_name, - service_call_lookup_keys, should_track_call_name as should_track_call, type_expr_to_string, - walk_stmts, + service_call_lookup_keys, type_expr_to_string, walk_stmts, }; use daglang_typecheck::{TypedCallableSignature, TypedItemSignature, TypedProject}; use gunbc_ir::patterns::branch::IfBuilder; @@ -205,10 +204,10 @@ pub enum PrimitiveOpKind { ContentUpsertOutputPath { path: String, }, - /// A compute node that evaluates a complex return expression. - /// Created by the lowerer when return expressions contain BinOp, - /// UnaryOp, If, Match, etc. that can't be wired as simple edges. - ReturnExprCompute { + /// A compute node that evaluates a lowered expression body. + /// Created by the lowerer when complex expressions cannot be represented + /// as direct wiring edges. + ExprCompute { fn_body: Box, }, } @@ -344,7 +343,7 @@ impl PrimitiveOpKind { } Self::CompareEquality => ObligationCategory::InterfaceContractVerification, Self::ContentUpsertOutputPath { .. } => ObligationCategory::None, - Self::ReturnExprCompute { .. } => ObligationCategory::None, + Self::ExprCompute { .. } => ObligationCategory::None, } } } @@ -2788,7 +2787,6 @@ fn register_endpoint( .or_insert(Some(endpoint)); } -#[allow(clippy::too_many_arguments)] fn add_dependency_edges( builder: &mut DagBuilder, project: &TypedProject, @@ -2928,45 +2926,75 @@ struct ForLoopSite { } #[derive(Debug)] -struct IfBranchSite { +struct IfBranchScopeSite { has_else: bool, - /// Service call paths inside the then-branch expression. then_service_call_paths: Vec>, - /// Service call paths inside the else-branch expression (empty if no else). else_service_call_paths: Vec>, } #[derive(Debug)] -struct MatchBranchSite { +struct MatchBranchScopeSite { arm_count: usize, - /// Service call paths from ALL arms combined (for top-level filtering). all_service_call_paths: Vec>, } -fn detect_for_loops_in_stmts(stmts: &[Stmt]) -> Vec { - let mut sites = Vec::new(); - walk_stmts(stmts, &mut |expr| { - if let Expr::For(var, iterable, passthrough, body) = expr { - let iterable_ref = match iterable.as_ref() { - Expr::Ident(name) => Some(IterableRef::Ident(name.clone())), - Expr::FieldAccess(base, field) => match base.as_ref() { - Expr::Ident(base_ident) => { - Some(IterableRef::FieldAccess(base_ident.clone(), field.clone())) +fn collect_service_paths_from_scoped_body(body: &scope::ScopedBody) -> Vec> { + body.all_service_calls() + .into_iter() + .map(|call| call.path.clone()) + .collect() +} + +fn collect_for_loop_sites_from_scoped(body: &scope::ScopedBody, out: &mut Vec) { + for item in &body.items { + match item { + scope::ScopedItem::ForLoop { + element_var, + iterable, + passthrough, + body, + } => { + let iterable_ref = match iterable { + scope::ExprRef::Ident(name) => Some(IterableRef::Ident(name.clone())), + scope::ExprRef::FieldAccess { base, field } => { + Some(IterableRef::FieldAccess(base.clone(), field.clone())) } - _ => None, - }, - _ => None, - }; - let mut body_calls = Vec::new(); - collect_service_call_paths_from_expr(body, &mut body_calls); - sites.push(ForLoopSite { - element_var: var.clone(), - iterable: iterable_ref, - passthrough: passthrough.clone(), - body_service_call_paths: body_calls, - }); + scope::ExprRef::Literal(_) | scope::ExprRef::Opaque => None, + }; + out.push(ForLoopSite { + element_var: element_var.clone(), + iterable: iterable_ref, + passthrough: passthrough.clone(), + body_service_call_paths: collect_service_paths_from_scoped_body(body), + }); + collect_for_loop_sites_from_scoped(body, out); + } + scope::ScopedItem::IfBranch { + then_body, + else_body, + } => { + collect_for_loop_sites_from_scoped(then_body, out); + if let Some(else_body) = else_body { + collect_for_loop_sites_from_scoped(else_body, out); + } + } + scope::ScopedItem::MatchBranch { arms } => { + for arm in arms { + collect_for_loop_sites_from_scoped(&arm.body, out); + } + } + scope::ScopedItem::ServiceCall(_) + | scope::ScopedItem::FnCall { .. } + | scope::ScopedItem::Binding { .. } + | scope::ScopedItem::Other => {} } - }); + } +} + +fn detect_for_loops_in_stmts(stmts: &[Stmt]) -> Vec { + let scoped = scope::ScopedBody::from_stmts(stmts); + let mut sites = Vec::new(); + collect_for_loop_sites_from_scoped(&scoped, &mut sites); sites } @@ -2998,7 +3026,7 @@ pub(crate) fn collect_service_call_paths_from_expr(expr: &Expr, paths: &mut Vec< } } -fn detect_if_branches_in_stmts(stmts: &[Stmt]) -> Vec { +fn detect_if_branches_in_stmts(stmts: &[Stmt]) -> Vec { let mut sites = Vec::new(); walk_stmts(stmts, &mut |expr| { if let Expr::If(_, then_expr, else_branch) = expr { @@ -3008,7 +3036,7 @@ fn detect_if_branches_in_stmts(stmts: &[Stmt]) -> Vec { if let Some(else_expr) = else_branch { collect_service_call_paths_from_expr(else_expr, &mut else_calls); } - sites.push(IfBranchSite { + sites.push(IfBranchScopeSite { has_else: else_branch.is_some(), then_service_call_paths: then_calls, else_service_call_paths: else_calls, @@ -3018,7 +3046,7 @@ fn detect_if_branches_in_stmts(stmts: &[Stmt]) -> Vec { sites } -fn detect_match_branches_in_stmts(stmts: &[Stmt]) -> Vec { +fn detect_match_branches_in_stmts(stmts: &[Stmt]) -> Vec { let mut sites = Vec::new(); walk_stmts(stmts, &mut |expr| { if let Expr::Match(_, arms) = expr { @@ -3026,7 +3054,7 @@ fn detect_match_branches_in_stmts(stmts: &[Stmt]) -> Vec { for arm in arms { collect_service_call_paths_from_expr(&arm.body, &mut all_calls); } - sites.push(MatchBranchSite { + sites.push(MatchBranchScopeSite { arm_count: arms.len(), all_service_call_paths: all_calls, }); @@ -3528,7 +3556,6 @@ fn add_control_flow_pattern_nodes( } } -#[allow(clippy::too_many_arguments)] fn expand_content_upsert_patterns( builder: &mut DagBuilder, module_name: &str, @@ -3574,7 +3601,7 @@ fn expand_content_upsert_patterns( }; match expr { Expr::Call(name, args) => { - if should_track_call(name) { + if !is_internal_synthetic_call(name) { bound_callables.insert(binding.clone(), name.clone()); } if name == "content_upsert" { @@ -3603,7 +3630,6 @@ fn expand_content_upsert_patterns( } } -#[allow(clippy::too_many_arguments)] fn expand_single_content_upsert( builder: &mut DagBuilder, module_name: &str, @@ -3996,7 +4022,6 @@ fn pattern_uses_binding_types(uses: &[daglang_syntax::ast::UsesClause]) -> HashM .collect() } -#[allow(clippy::too_many_arguments)] fn expand_non_generic_pattern_calls( builder: &mut DagBuilder, project: &TypedProject, @@ -4078,9 +4103,6 @@ fn expand_non_generic_pattern_calls( /// Maximum recursion depth for pattern expansion (patterns calling patterns). const PATTERN_EXPANSION_MAX_DEPTH: usize = 5; -const PARAM_REF_SENTINEL: &str = "__param_ref__"; - -#[allow(clippy::too_many_arguments)] fn expand_single_pattern( builder: &mut DagBuilder, module_name: &str, @@ -4397,7 +4419,6 @@ fn resolve_expr_idents(expr: &Expr, arg_map: &HashMap) -> Expr { /// - `Expr::ServiceCall` → clone transport triplet from service registry /// - `Expr::Call("eq", ...)` → create CompareEquality primitive /// - `Expr::Call(pattern_name, ...)` → recursive pattern expansion -#[allow(clippy::too_many_arguments)] fn expand_pattern_body_node( builder: &mut DagBuilder, module_name: &str, @@ -4523,7 +4544,6 @@ fn expand_pattern_body_node( /// Expand a service call node (e.g., `fs.read(path: path)`) by cloning /// the corresponding transport triplet from the service registry. -#[allow(clippy::too_many_arguments)] fn expand_service_call_node( builder: &mut DagBuilder, module_name: &str, @@ -4599,7 +4619,6 @@ fn expand_service_call_node( } /// Expand an `eq(a: ..., b: ...)` call into a CompareEquality primitive node. -#[allow(clippy::too_many_arguments)] fn expand_eq_node( builder: &mut DagBuilder, module_name: &str, @@ -4683,7 +4702,6 @@ fn expand_eq_node( /// /// Resolves the expression through the pattern's arg_map (caller scope), /// node_outputs (expanded nodes), and param sources. -#[allow(clippy::too_many_arguments)] fn wire_pattern_arg_to_prepare( builder: &mut DagBuilder, module_name: &str, @@ -4724,7 +4742,6 @@ fn wire_pattern_arg_to_prepare( /// 2. If the arg is a FieldAccess on an expanded node → wire from that node. /// 3. If the arg is a literal → create a literal source node. /// 4. If the arg is an ident matching a caller callable → wire from that. -#[allow(clippy::too_many_arguments)] fn wire_pattern_arg_to_node( builder: &mut DagBuilder, module_name: &str, @@ -4816,7 +4833,6 @@ fn wire_pattern_arg_to_node( /// /// Handles idents (param sources, callable endpoints, data values), /// field access, and literals from the caller's scope. -#[allow(clippy::too_many_arguments)] fn wire_caller_expr_to_node( builder: &mut DagBuilder, module_name: &str, @@ -5921,7 +5937,6 @@ fn add_service_transport_triplets( Ok(registry) } -#[allow(clippy::too_many_arguments)] fn add_service_call_edges( builder: &mut DagBuilder, project: &TypedProject, @@ -7417,10 +7432,14 @@ fn item_callable_provides(item: &Item) -> Option<(&str, &[daglang_syntax::ast::P } } +fn is_internal_synthetic_call(name: &str) -> bool { + matches!(name, "" | "as" | "with" | "fn") +} + fn collect_calls_from_stmts(stmts: &[Stmt], calls: &mut BTreeSet) { walk_stmts(stmts, &mut |expr| { if let Expr::Call(name, _) = expr { - if should_track_call(name) { + if !is_internal_synthetic_call(name) { calls.insert(name.clone()); } } @@ -7507,14 +7526,48 @@ fn collection_op_kind(name: &str) -> Option { fn collect_collection_ops_from_stmts(stmts: &[Stmt], sites: &mut Vec) { walk_stmts(stmts, &mut |expr| { - if let Expr::Pipe(_, rhs) = expr { - let Expr::Call(name, _) = rhs.as_ref() else { - return; - }; - let Some(kind) = collection_op_kind(name) else { - return; - }; - sites.push(CollectionOpSite { kind }); + match expr { + Expr::Pipe(_, rhs) => { + let Expr::Call(name, _) = rhs.as_ref() else { + return; + }; + let Some(kind) = collection_op_kind(name) else { + return; + }; + sites.push(CollectionOpSite { kind }); + } + Expr::PipeCall(_, method, _) => { + let method_name = match method { + daglang_syntax::ast::PipeMethod::Map => "map", + daglang_syntax::ast::PipeMethod::Filter => "filter", + daglang_syntax::ast::PipeMethod::FilterMap => "filter_map", + daglang_syntax::ast::PipeMethod::FlatMap => "flat_map", + daglang_syntax::ast::PipeMethod::SortBy => "sort_by", + daglang_syntax::ast::PipeMethod::Append => "append", + daglang_syntax::ast::PipeMethod::Fold => "fold", + daglang_syntax::ast::PipeMethod::Join => "join", + daglang_syntax::ast::PipeMethod::Count => "count", + daglang_syntax::ast::PipeMethod::Sum => "sum", + daglang_syntax::ast::PipeMethod::First => "first", + daglang_syntax::ast::PipeMethod::Last => "last", + daglang_syntax::ast::PipeMethod::MaxBy => "max_by", + daglang_syntax::ast::PipeMethod::Any => "any", + daglang_syntax::ast::PipeMethod::All => "all", + daglang_syntax::ast::PipeMethod::Contains => "contains", + daglang_syntax::ast::PipeMethod::StartsWith => "starts_with", + daglang_syntax::ast::PipeMethod::EndsWith => "ends_with", + daglang_syntax::ast::PipeMethod::Repeat => "repeat", + daglang_syntax::ast::PipeMethod::ReplaceSection => "replace_section", + daglang_syntax::ast::PipeMethod::Chars => "chars", + daglang_syntax::ast::PipeMethod::ToBytes => "to_bytes", + daglang_syntax::ast::PipeMethod::ToJson => "to_json", + daglang_syntax::ast::PipeMethod::Hash => "hash", + }; + if let Some(kind) = collection_op_kind(method_name) { + sites.push(CollectionOpSite { kind }); + } + } + _ => {} } }); } @@ -7618,7 +7671,7 @@ pub(crate) fn collect_service_calls_from_stmts(stmts: &[Stmt], calls: &mut Vec) { walk_stmts(stmts, &mut |expr| { if let Expr::Call(name, args) = expr { - if should_track_call(name) { + if !is_internal_synthetic_call(name) { calls.push(FnCallSite { name: name.clone(), args: args.iter().map(service_call_arg_site).collect(), @@ -7759,7 +7812,6 @@ pub fn build_data_values(project: &TypedProject) -> HashMap Option { } } -#[allow(clippy::too_many_arguments)] fn resolve_return_expr_source( builder: &mut DagBuilder, expr: &Expr, @@ -8029,10 +8080,9 @@ fn resolve_return_expr_source( ); Some((src, output_name.to_string())) } - // RT4a: Handle complex return expressions (BinOp, UnaryOp, If, Match, Pipe, etc.) - // by synthesizing a compute node that evaluates the expression at runtime. + // Handle complex expressions by synthesizing a dedicated compute node. _ => { - synthesize_return_expr_compute( + synthesize_expr_compute( builder, expr, output_port, @@ -8049,18 +8099,22 @@ fn resolve_return_expr_source( } } +#[derive(Debug, Clone, PartialEq, Eq)] +struct ExprLeafRef { + input_port: String, + source: expr::LeafRef, +} + /// Collect all leaf expression references from a complex expression. -/// Returns (input_port_name, source_node_id, source_port_name, type_id) tuples. /// Sets `has_local_refs` to true if the expression references local variables /// (let bindings) that can't be resolved as compute node inputs. -#[allow(clippy::too_many_arguments)] fn collect_expr_leaf_refs( expr: &Expr, param_types: &HashMap, bound_callable_sources: &HashMap, bound_service_sources: &HashMap, endpoints_by_name: &HashMap>, - refs: &mut Vec<(String, String, String, String)>, + refs: &mut Vec, seen: &mut HashSet, has_local_refs: &mut bool, ) { @@ -8072,16 +8126,41 @@ fn collect_expr_leaf_refs( } if let Some(param_ty) = param_types.get(name) { seen.insert(port_name.clone()); - refs.push((port_name, PARAM_REF_SENTINEL.to_string(), name.clone(), param_ty.clone())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Param { + name: name.clone(), + field: None, + ty: param_ty.clone(), + }, + }); } else if let Some(source) = bound_callable_sources.get(name) { seen.insert(port_name.clone()); - refs.push((port_name, source.node_id.clone(), source.primary_output.clone(), "Any".to_string())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Callable { + endpoint: source.node_id.clone(), + port: source.primary_output.clone(), + }, + }); } else if let Some(source) = bound_service_sources.get(name) { seen.insert(port_name.clone()); - refs.push((port_name, source.parse.node_id.clone(), source.parse.primary_output.clone(), "Any".to_string())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Service { + endpoint: source.parse.node_id.clone(), + port: source.parse.primary_output.clone(), + }, + }); } else if let Some(Some(source)) = endpoints_by_name.get(name) { seen.insert(port_name.clone()); - refs.push((port_name, source.node_id.clone(), source.primary_output.clone(), "Any".to_string())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Callable { + endpoint: source.node_id.clone(), + port: source.primary_output.clone(), + }, + }); } else { *has_local_refs = true; } @@ -8096,17 +8175,42 @@ fn collect_expr_leaf_refs( let base_port = base_ident.clone(); if !seen.contains(&base_port) { seen.insert(base_port.clone()); - refs.push((base_port, PARAM_REF_SENTINEL.to_string(), base_ident.clone(), param_ty.clone())); + refs.push(ExprLeafRef { + input_port: base_port, + source: expr::LeafRef::Param { + name: base_ident.clone(), + field: Some(field.clone()), + ty: param_ty.clone(), + }, + }); } } else if let Some(source) = bound_callable_sources.get(base_ident) { seen.insert(port_name.clone()); - refs.push((port_name, source.node_id.clone(), field.clone(), "Any".to_string())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Callable { + endpoint: source.node_id.clone(), + port: field.clone(), + }, + }); } else if let Some(source) = bound_service_sources.get(base_ident) { seen.insert(port_name.clone()); - refs.push((port_name, source.parse.node_id.clone(), field.clone(), "Any".to_string())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Service { + endpoint: source.parse.node_id.clone(), + port: field.clone(), + }, + }); } else if let Some(Some(source)) = endpoints_by_name.get(base_ident) { seen.insert(port_name.clone()); - refs.push((port_name, source.node_id.clone(), field.clone(), "Any".to_string())); + refs.push(ExprLeafRef { + input_port: port_name, + source: expr::LeafRef::Callable { + endpoint: source.node_id.clone(), + port: field.clone(), + }, + }); } } else { collect_expr_leaf_refs(base, param_types, bound_callable_sources, bound_service_sources, endpoints_by_name, refs, seen, has_local_refs); @@ -8130,6 +8234,12 @@ fn collect_expr_leaf_refs( collect_expr_leaf_refs(receiver, param_types, bound_callable_sources, bound_service_sources, endpoints_by_name, refs, seen, has_local_refs); collect_expr_leaf_refs(call, param_types, bound_callable_sources, bound_service_sources, endpoints_by_name, refs, seen, has_local_refs); } + Expr::PipeCall(receiver, _, args) => { + collect_expr_leaf_refs(receiver, param_types, bound_callable_sources, bound_service_sources, endpoints_by_name, refs, seen, has_local_refs); + for (_, arg) in args { + collect_expr_leaf_refs(arg, param_types, bound_callable_sources, bound_service_sources, endpoints_by_name, refs, seen, has_local_refs); + } + } Expr::Match(scrutinee, arms) => { collect_expr_leaf_refs(scrutinee, param_types, bound_callable_sources, bound_service_sources, endpoints_by_name, refs, seen, has_local_refs); for arm in arms { @@ -8245,14 +8355,52 @@ fn remap_expr_idents(expr: &Expr) -> expr::LoweredExpr { .collect(); expr::LoweredExpr::StringInterp(lowered_parts) } + Expr::PipeCall(receiver, method, args) => expr::LoweredExpr::Pipe { + receiver: Box::new(remap_expr_idents(receiver)), + call: Box::new(expr::LoweredExpr::Call { + name: pipe_method_name(*method).to_string(), + args: args + .iter() + .map(|(k, v)| (k.clone(), remap_expr_idents(v))) + .collect(), + }), + }, _ => expr::LoweredExpr::Literal(expr::LoweredLiteral::None), } } +fn pipe_method_name(method: daglang_syntax::ast::PipeMethod) -> &'static str { + match method { + daglang_syntax::ast::PipeMethod::Map => "map", + daglang_syntax::ast::PipeMethod::Filter => "filter", + daglang_syntax::ast::PipeMethod::FilterMap => "filter_map", + daglang_syntax::ast::PipeMethod::FlatMap => "flat_map", + daglang_syntax::ast::PipeMethod::SortBy => "sort_by", + daglang_syntax::ast::PipeMethod::Append => "append", + daglang_syntax::ast::PipeMethod::Fold => "fold", + daglang_syntax::ast::PipeMethod::Join => "join", + daglang_syntax::ast::PipeMethod::Count => "count", + daglang_syntax::ast::PipeMethod::Sum => "sum", + daglang_syntax::ast::PipeMethod::First => "first", + daglang_syntax::ast::PipeMethod::Last => "last", + daglang_syntax::ast::PipeMethod::MaxBy => "max_by", + daglang_syntax::ast::PipeMethod::Any => "any", + daglang_syntax::ast::PipeMethod::All => "all", + daglang_syntax::ast::PipeMethod::Contains => "contains", + daglang_syntax::ast::PipeMethod::StartsWith => "starts_with", + daglang_syntax::ast::PipeMethod::EndsWith => "ends_with", + daglang_syntax::ast::PipeMethod::Repeat => "repeat", + daglang_syntax::ast::PipeMethod::ReplaceSection => "replace_section", + daglang_syntax::ast::PipeMethod::Chars => "chars", + daglang_syntax::ast::PipeMethod::ToBytes => "to_bytes", + daglang_syntax::ast::PipeMethod::ToJson => "to_json", + daglang_syntax::ast::PipeMethod::Hash => "hash", + } +} + /// Synthesize a compute node for a complex return expression. /// Creates a node that evaluates the expression using `evaluate_fn_body` at runtime. -#[allow(clippy::too_many_arguments)] -fn synthesize_return_expr_compute( +fn synthesize_expr_compute( builder: &mut DagBuilder, expr: &Expr, output_port: &Port, @@ -8265,7 +8413,7 @@ fn synthesize_return_expr_compute( item_name: &str, disambiguator: &str, ) -> Option<(String, String)> { - let mut refs: Vec<(String, String, String, String)> = Vec::new(); + let mut refs: Vec = Vec::new(); let mut seen = HashSet::new(); let mut has_local_refs = false; collect_expr_leaf_refs( @@ -8285,7 +8433,14 @@ fn synthesize_return_expr_compute( let input_ports: Vec = refs .iter() - .map(|(port_name, _, _, type_id)| Port::with_cardinality(port_name.as_str(), type_id.as_str(), Cardinality::ONE)) + .map(|leaf| match &leaf.source { + expr::LeafRef::Param { ty, .. } => { + Port::with_cardinality(leaf.input_port.as_str(), ty.as_str(), Cardinality::ONE) + } + expr::LeafRef::Callable { .. } | expr::LeafRef::Service { .. } => { + Port::with_cardinality(leaf.input_port.as_str(), "Any", Cardinality::ONE) + } + }) .collect(); let output_type = output_port.type_id.0.as_str(); let result_port_name = "result"; @@ -8301,7 +8456,7 @@ fn synthesize_return_expr_compute( // 4. Create the compute node. let node_id = format!( - "return_expr_compute_{}", + "expr_compute_{}", sanitize_identifier(&format!("{module_name}_{item_name}_{output_name}_{disambiguator}")) ); builder.add_node(Node::opaque( @@ -8310,34 +8465,30 @@ fn synthesize_return_expr_compute( output_ports, LoweredOp::Primitive { module: module_name.to_string(), - name: format!("return_expr_compute::{item_name}::{output_name}"), - kind: PrimitiveOpKind::ReturnExprCompute { + name: format!("expr_compute::{item_name}::{output_name}"), + kind: PrimitiveOpKind::ExprCompute { fn_body: Box::new(fn_body), }, }, )); - for (port_name, source_node, source_port, _type_id) in &refs { - if source_node == PARAM_REF_SENTINEL { - // source_port is either "param_name" (plain) or "param__field" (field access). - // Split on "__" to get the base param and optional field. - let (base_param, field) = source_port - .split_once("__") - .map_or((source_port.as_str(), None), |(b, f)| (b, Some(f))); - let param_ty = param_types.get(base_param).map_or("Any", |s| s.as_str()); - let param_source_id = - ensure_param_source_node(builder, module_name, item_name, base_param, param_ty); - let src_port = field.unwrap_or(base_param); - builder.add_edge(¶m_source_id, src_port, &node_id, port_name); - } else { - builder.add_edge(source_node, source_port, &node_id, port_name); + for leaf in &refs { + match &leaf.source { + expr::LeafRef::Param { name, ty, .. } => { + let param_source_id = + ensure_param_source_node(builder, module_name, item_name, name, ty); + builder.add_edge(¶m_source_id, name, &node_id, &leaf.input_port); + } + expr::LeafRef::Callable { endpoint, port } + | expr::LeafRef::Service { endpoint, port } => { + builder.add_edge(endpoint, port, &node_id, &leaf.input_port); + } } } Some((node_id, result_port_name.to_string())) } -#[allow(clippy::too_many_arguments)] fn wire_callable_return_outputs( builder: &mut DagBuilder, stmts: &[Stmt], @@ -8468,7 +8619,6 @@ fn collect_for_loop_bindings( /// `{target}::cf_for_{index}` and an input port named `"items"`. This function /// resolves `` to a source node (service call result, callable result, /// or parameter) and wires the data edge. -#[allow(clippy::too_many_arguments)] fn wire_for_loop_iterables( builder: &mut DagBuilder, stmts: &[Stmt], diff --git a/core/daglang/daglang-lower/src/scope.rs b/core/daglang/daglang-lower/src/scope.rs index cda79d46805..2623a9c7591 100644 --- a/core/daglang/daglang-lower/src/scope.rs +++ b/core/daglang/daglang-lower/src/scope.rs @@ -36,6 +36,8 @@ pub(crate) struct ScopedServiceCall { pub(crate) enum ExprRef { /// A bare identifier reference: `x`. Ident(String), + /// Field access reference: `x.field`. + FieldAccess { base: String, field: String }, /// A literal value. Literal(crate::ServiceCallArgLiteral), /// Any other expression (not decomposed further for scope analysis). @@ -51,6 +53,7 @@ pub(crate) enum ScopedItem { /// A for-loop introducing a nested scope for its body. ForLoop { element_var: String, + iterable: ExprRef, passthrough: Vec, body: ScopedBody, }, @@ -213,6 +216,7 @@ fn collect_scoped_items_from_expr(expr: &Expr, items: &mut Vec) { let body_scope = scope_from_expr(body); items.push(ScopedItem::ForLoop { element_var: var.clone(), + iterable: expr_ref_from_expr(_iterable), passthrough: passthrough.clone(), body: body_scope, }); @@ -271,6 +275,12 @@ fn collect_scoped_items_from_expr(expr: &Expr, items: &mut Vec) { collect_scoped_items_from_expr(lhs, items); collect_scoped_items_from_expr(rhs, items); } + Expr::PipeCall(receiver, _, args) => { + collect_scoped_items_from_expr(receiver, items); + for (_, arg_expr) in args { + collect_scoped_items_from_expr(arg_expr, items); + } + } // Binary op — recurse. Expr::BinOp(lhs, _, rhs) => { @@ -345,6 +355,21 @@ fn collect_scoped_items_from_expr(expr: &Expr, items: &mut Vec) { } } +fn expr_ref_from_expr(expr: &Expr) -> ExprRef { + match expr { + Expr::Ident(name) => ExprRef::Ident(name.clone()), + Expr::FieldAccess(base, field) => match base.as_ref() { + Expr::Ident(base_ident) => ExprRef::FieldAccess { + base: base_ident.clone(), + field: field.clone(), + }, + _ => ExprRef::Opaque, + }, + Expr::Literal(_) => ExprRef::Opaque, + _ => ExprRef::Opaque, + } +} + /// Build a ScopedBody from a single expression (used for branch/loop bodies). fn scope_from_expr(expr: &Expr) -> ScopedBody { let mut items = Vec::new(); diff --git a/core/daglang/daglang-syntax/src/ast_utils.rs b/core/daglang/daglang-syntax/src/ast_utils.rs index 8bda8fe394f..cd190592a84 100644 --- a/core/daglang/daglang-syntax/src/ast_utils.rs +++ b/core/daglang/daglang-syntax/src/ast_utils.rs @@ -48,81 +48,6 @@ pub fn resource_type_name(resource_type: &TypeExpr) -> String { } } -/// Built-in pipe methods resolved by the evaluator, not as callable targets. -/// -/// Single authoritative registry. To add a new pipe method: -/// 1. Add a variant here -/// 2. Add the string match in `from_str` -/// 3. Implement the evaluation in `eval.rs` -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum PipeMethod { - // Collection → Collection - Map, - Filter, - FilterMap, - FlatMap, - SortBy, - Append, - // Collection → Scalar - Fold, - Join, - Count, - Sum, - First, - Last, - MaxBy, - Any, - All, - Contains, - // String methods - StartsWith, - EndsWith, - Repeat, - ReplaceSection, - Chars, - // Conversion methods - ToBytes, - ToJson, - Hash, -} - -impl PipeMethod { - /// Parse a method name string into a PipeMethod, if it is a known built-in. - pub fn from_str(name: &str) -> Option { - match name { - "map" => Some(Self::Map), - "filter" => Some(Self::Filter), - "filter_map" => Some(Self::FilterMap), - "flat_map" => Some(Self::FlatMap), - "sort_by" => Some(Self::SortBy), - "append" => Some(Self::Append), - "fold" => Some(Self::Fold), - "join" => Some(Self::Join), - "count" => Some(Self::Count), - "sum" => Some(Self::Sum), - "first" => Some(Self::First), - "last" => Some(Self::Last), - "max_by" => Some(Self::MaxBy), - "any" => Some(Self::Any), - "all" => Some(Self::All), - "contains" => Some(Self::Contains), - "starts_with" => Some(Self::StartsWith), - "ends_with" => Some(Self::EndsWith), - "repeat" => Some(Self::Repeat), - "replace_section" => Some(Self::ReplaceSection), - "chars" => Some(Self::Chars), - "to_bytes" => Some(Self::ToBytes), - "to_json" => Some(Self::ToJson), - "hash" => Some(Self::Hash), - _ => None, - } - } -} - -pub fn should_track_call_name(name: &str) -> bool { - !matches!(name, "" | "as" | "with" | "fn") && PipeMethod::from_str(name).is_none() -} - pub fn service_call_lookup_keys(call_path: &[String]) -> Option<[String; 3]> { if call_path.len() < 2 { return None; @@ -171,6 +96,12 @@ pub fn walk_expr(expr: &Expr, visitor: &mut impl FnMut(&Expr)) { walk_expr(lhs, visitor); walk_expr(rhs, visitor); } + Expr::PipeCall(receiver, _, args) => { + walk_expr(receiver, visitor); + for (_, arg) in args { + walk_expr(arg, visitor); + } + } Expr::UnaryOp(_, inner) | Expr::Lambda(_, inner) | Expr::After(inner, _) => { walk_expr(inner, visitor) } diff --git a/core/daglang/daglang-syntax/src/lib.rs b/core/daglang/daglang-syntax/src/lib.rs index 51aba9e3c28..3e792aa93ca 100644 --- a/core/daglang/daglang-syntax/src/lib.rs +++ b/core/daglang/daglang-syntax/src/lib.rs @@ -520,6 +520,8 @@ pub mod ast { For(String, Box, Vec, Box), /// Pipe: `expr |> fn` Pipe(Box, Box), + /// Built-in pipe method call: `expr |> method(args)` + PipeCall(Box, PipeMethod, Vec<(Option, Expr)>), /// Lambda (inline only, in |> chains): `x => x.name` Lambda(Vec, Box), /// List literal: `[a, b, c]` @@ -534,6 +536,73 @@ pub mod ast { Return(Vec<(String, Expr)>), } + /// Built-in pipe methods resolved by parser/typechecker/lowerer as + /// first-class syntax, not free-form callable names. + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub enum PipeMethod { + // Collection -> Collection + Map, + Filter, + FilterMap, + FlatMap, + SortBy, + Append, + // Collection -> Scalar + Fold, + Join, + Count, + Sum, + First, + Last, + MaxBy, + Any, + All, + Contains, + // String methods + StartsWith, + EndsWith, + Repeat, + ReplaceSection, + Chars, + // Conversion methods + ToBytes, + ToJson, + Hash, + } + + impl PipeMethod { + /// Parse a method name string into a known built-in pipe method. + pub fn from_str(name: &str) -> Option { + match name { + "map" => Some(Self::Map), + "filter" => Some(Self::Filter), + "filter_map" => Some(Self::FilterMap), + "flat_map" => Some(Self::FlatMap), + "sort_by" => Some(Self::SortBy), + "append" => Some(Self::Append), + "fold" => Some(Self::Fold), + "join" => Some(Self::Join), + "count" => Some(Self::Count), + "sum" => Some(Self::Sum), + "first" => Some(Self::First), + "last" => Some(Self::Last), + "max_by" => Some(Self::MaxBy), + "any" => Some(Self::Any), + "all" => Some(Self::All), + "contains" => Some(Self::Contains), + "starts_with" => Some(Self::StartsWith), + "ends_with" => Some(Self::EndsWith), + "repeat" => Some(Self::Repeat), + "replace_section" => Some(Self::ReplaceSection), + "chars" => Some(Self::Chars), + "to_bytes" => Some(Self::ToBytes), + "to_json" => Some(Self::ToJson), + "hash" => Some(Self::Hash), + _ => None, + } + } + } + #[derive(Debug, Clone)] pub enum Literal { Int(i64), diff --git a/core/daglang/daglang-syntax/src/parser.rs b/core/daglang/daglang-syntax/src/parser.rs index afe53b86192..ba113ac5529 100644 --- a/core/daglang/daglang-syntax/src/parser.rs +++ b/core/daglang/daglang-syntax/src/parser.rs @@ -2552,7 +2552,16 @@ impl Parser { Expr::BinOp(Box::new(lhs), bop, Box::new(rhs)) } else { match op { - TokenKind::PipeArrow => Expr::Pipe(Box::new(lhs), Box::new(rhs)), + TokenKind::PipeArrow => match &rhs { + Expr::Call(name, args) => { + if let Some(method) = PipeMethod::from_str(name) { + Expr::PipeCall(Box::new(lhs), method, args.clone()) + } else { + Expr::Pipe(Box::new(lhs), Box::new(rhs)) + } + } + _ => Expr::Pipe(Box::new(lhs), Box::new(rhs)), + }, TokenKind::NullCoalesce => { Expr::BinOp(Box::new(lhs), BinOp::NullCoalesce, Box::new(rhs)) } @@ -3041,6 +3050,11 @@ mod tests { parser.parse_expr(0).expect_err("expression should fail") } + fn parse_source_err(source: &str) -> ParseError { + let mut errs = parse(source).expect_err("source should fail"); + errs.remove(0) + } + #[test] fn parse_module_decl() { let sf = parse_or_panic("module tools.makegen"); @@ -3058,6 +3072,26 @@ mod tests { ); } + #[test] + fn parse_service_config_auth_input_rejects_string_literal() { + let err = parse_source_err( + r#"module services.example +service github.Gist { + config { + endpoint: "https://api.github.com" + auth: BearerToken + auth_input: "token" + } + operation Create() -> { id: String } +}"#, + ); + assert!( + err.message.contains("expected identifier for `auth_input`"), + "unexpected parse error: {}", + err.message + ); + } + #[test] fn parse_simple_fn() { let sf = parse_or_panic("module test\nfn greet(name: String) -> String { name }"); @@ -3390,6 +3424,19 @@ pipeline gist { } } + #[test] + fn expression_pipe_method_lowers_to_pipe_call() { + let expr = parse_expr_only("items |> map(x => x)"); + match expr { + Expr::PipeCall(receiver, PipeMethod::Map, args) => { + assert!(matches!(*receiver, Expr::Ident(ref name) if name == "items")); + assert_eq!(args.len(), 1); + assert!(matches!(args[0].1, Expr::Lambda(_, _))); + } + other => panic!("unexpected expression tree: {other:?}"), + } + } + #[test] fn expression_logical_and_binds_tighter_than_or() { let expr = parse_expr_only("x || y && z"); diff --git a/core/daglang/daglang-typecheck/src/lib.rs b/core/daglang/daglang-typecheck/src/lib.rs index 0dc3a8bff50..4fd6c419f5e 100644 --- a/core/daglang/daglang-typecheck/src/lib.rs +++ b/core/daglang/daglang-typecheck/src/lib.rs @@ -29,8 +29,7 @@ use daglang_syntax::ast::{ SourceFile, Stmt, TypeBody, TypeExpr, UsesClause, }; use daglang_syntax::ast_utils::{ - resource_type_name, service_call_lookup_keys, - should_track_call_name as should_validate_call_name, type_expr_to_string, walk_stmts, + resource_type_name, service_call_lookup_keys, type_expr_to_string, walk_stmts, }; /// A typechecked project snapshot. @@ -2746,6 +2745,49 @@ fn infer_expr_type( errors.extend(rhs_errors); rhs_val } + Expr::PipeCall(receiver, _method, args) => { + let (_, recv_errors) = infer_expr_type(receiver, local_bindings, infer_context); + errors.extend(recv_errors); + for (_name, arg) in args { + let (_, arg_errors) = infer_expr_type(arg, local_bindings, infer_context); + errors.extend(arg_errors); + } + match _method { + daglang_syntax::ast::PipeMethod::Count | daglang_syntax::ast::PipeMethod::Sum => { + ValueType::Named("Int".to_string()) + } + daglang_syntax::ast::PipeMethod::Any + | daglang_syntax::ast::PipeMethod::All + | daglang_syntax::ast::PipeMethod::Contains + | daglang_syntax::ast::PipeMethod::StartsWith + | daglang_syntax::ast::PipeMethod::EndsWith => { + ValueType::Named("Bool".to_string()) + } + daglang_syntax::ast::PipeMethod::Join + | daglang_syntax::ast::PipeMethod::Repeat + | daglang_syntax::ast::PipeMethod::ReplaceSection + | daglang_syntax::ast::PipeMethod::Hash => { + ValueType::Named("String".to_string()) + } + daglang_syntax::ast::PipeMethod::ToBytes => { + ValueType::Named("Bytes".to_string()) + } + daglang_syntax::ast::PipeMethod::ToJson => ValueType::Named("Json".to_string()), + daglang_syntax::ast::PipeMethod::Map + | daglang_syntax::ast::PipeMethod::Filter + | daglang_syntax::ast::PipeMethod::FilterMap + | daglang_syntax::ast::PipeMethod::FlatMap + | daglang_syntax::ast::PipeMethod::SortBy + | daglang_syntax::ast::PipeMethod::Append + | daglang_syntax::ast::PipeMethod::Chars => { + ValueType::Named("List".to_string()) + } + daglang_syntax::ast::PipeMethod::Fold + | daglang_syntax::ast::PipeMethod::First + | daglang_syntax::ast::PipeMethod::Last + | daglang_syntax::ast::PipeMethod::MaxBy => ValueType::Unknown, + } + } Expr::Lambda(_, body) => { let (val, body_errors) = infer_expr_type(body, local_bindings, infer_context); errors.extend(body_errors); @@ -3213,20 +3255,25 @@ struct BodyServiceCall { fn collect_calls_from_stmts(stmts: &[Stmt], calls: &mut Vec) { walk_stmts(stmts, &mut |expr| { if let Expr::Call(name, args) = expr { - if should_validate_call_name(name) { - calls.push(BodyCall { - callee: name.clone(), - arg_count: args.len(), - named_args: args - .iter() - .filter_map(|(name, _)| name.clone()) - .collect::>(), - }); + if is_internal_synthetic_call(name) { + return; } + calls.push(BodyCall { + callee: name.clone(), + arg_count: args.len(), + named_args: args + .iter() + .filter_map(|(name, _)| name.clone()) + .collect::>(), + }); } }); } +fn is_internal_synthetic_call(name: &str) -> bool { + matches!(name, "" | "as" | "with" | "fn") +} + fn collect_service_calls_from_stmts(stmts: &[Stmt], calls: &mut Vec) { walk_stmts(stmts, &mut |expr| { if let Expr::ServiceCall(path, args) = expr { diff --git a/core/exec/src/execute/mod.rs b/core/exec/src/execute/mod.rs index 378f11e0fbb..917d8ba1967 100644 --- a/core/exec/src/execute/mod.rs +++ b/core/exec/src/execute/mod.rs @@ -33,7 +33,8 @@ use gunbc_ir::transport::{FileOp, TransportResponse}; use gunbc_ir::{ canonical_edge_order, classify_coercion, detect_boundaries, detect_entrypoints, normalize_resource_id, AccessMode, AppliedCoercion, BoundaryInfo, Cardinality, Dag, - LogDetailLevel, Node, NodeBody, NodeId, NodeKind, Value, RESOURCE_FILE, RESOURCE_FILE_PREFIX, + LogDetailLevel, Node, NodeBody, NodeId, NodeKind, PortName, Value, RESOURCE_FILE, + RESOURCE_FILE_PREFIX, }; use std::collections::{HashMap, HashSet}; use std::fmt; @@ -583,6 +584,29 @@ fn remap_input_mocks( result } +/// Resolve an input mock value for `(node_id, port_name)`. +/// +/// Param source nodes are auto-fed by matching any explicitly provided +/// non-param-source input with the same port name. +fn resolve_mock_input( + mocks: &BoundaryMocks, + node_id: &NodeId, + port_name: &PortName, +) -> Option { + if let Some(value) = mocks.get_input(&node_id.0, &port_name.0) { + return Some(value.clone()); + } + if !node_id.0.starts_with("param_source_") { + return None; + } + for ((mock_node, mock_port), value) in mocks.iter_inputs() { + if mock_port == &port_name.0 && !mock_node.starts_with("param_source_") { + return Some(value.clone()); + } + } + None +} + /// Remap input mocks embedded in DryRun/Simulate execution modes. fn remap_mode_inputs( mode: ExecutionMode, @@ -752,8 +776,8 @@ fn execute_flat_sequential( let mut inject_inputs = |mocks: &BoundaryMocks| { for port in &node.inputs { if !inputs.contains_key(&port.name.0) { - if let Some(mock_value) = mocks.get_input(&node.id.0, &port.name.0) { - inputs.insert(port.name.0.clone(), mock_value.clone()); + if let Some(mock_value) = resolve_mock_input(mocks, &node.id, &port.name) { + inputs.insert(port.name.0.clone(), mock_value); } } } @@ -1338,8 +1362,8 @@ fn build_node_inputs( let mut inject_inputs = |mocks: &BoundaryMocks| { for port in &node.inputs { if !inputs.contains_key(&port.name.0) { - if let Some(mock_value) = mocks.get_input(&node.id.0, &port.name.0) { - inputs.insert(port.name.0.clone(), mock_value.clone()); + if let Some(mock_value) = resolve_mock_input(mocks, &node.id, &port.name) { + inputs.insert(port.name.0.clone(), mock_value); } } } @@ -2054,60 +2078,51 @@ fn should_intercept_by_kind(node: &Node) -> bool { ) } -/// Port-level heuristic: does this node look effectful despite `kind: Pure`? -/// -/// Returns a human-readable reason string if the node has port patterns that -/// indicate it should have been classified (transport request inputs, -/// tool handle ports, resource ports). -fn looks_effectful_without_kind(node: &Node) -> Option<&'static str> { - for p in &node.inputs { - if p.type_id.0 == "TransportRequest" { - return Some("has TransportRequest input (transport executor)"); - } - if p.type_id.0 == "ToolHandle" { - return Some("has ToolHandle input (tool consumer)"); - } - if p.name.is_resource() { - return Some("has res:* input port (resource consumer)"); - } - } - for p in &node.outputs { - if p.type_id.0 == "ToolHandle" { - return Some("has ToolHandle output (tool environment)"); - } - if matches!( - p.type_id.0.as_str(), - "FilesystemHandle" - | "NetworkHandle" - | "Timestamp" - | "Credential" - | "Platform" - ) { - return Some("has resource-environment output"); - } - } - None -} - /// Pre-flight check: error if any `Pure` node has effectful port patterns. /// -/// Called before DryRun/Simulate execution to ensure no effectful node slips -/// through interception due to a missing or incorrect `NodeKind`. This enforces -/// the "no silent fallback" invariant — either the lowerer stamps `kind` via -/// `stamp_node_kinds()`, or the hand-built DAG sets it explicitly via -/// [`Node::with_kind`]. +/// `looks_effectful_without_kind()` helper was removed (C18); keep the checks +/// inlined here so accidental `kind: Pure` regressions still fail closed. fn validate_node_kinds_for_interception(dag: &Dag) -> Result<(), ExecError> { for node in &dag.nodes { if node.kind != NodeKind::Pure { continue; } - if let Some(reason) = looks_effectful_without_kind(node) { - return Err(ExecError::new(format!( - "node '{}' has kind: Pure but {reason}. \ - Set NodeKind via Node::with_kind() or the lowerer's stamp_node_kinds() \ - so DryRun/Simulate can intercept it correctly.", - node.id.0, - ))); + for port in &node.inputs { + if port.type_id.0 == "TransportRequest" { + return Err(ExecError::new(format!( + "node '{}' has kind: Pure but has TransportRequest input", + node.id.0 + ))); + } + if port.type_id.0 == "ToolHandle" { + return Err(ExecError::new(format!( + "node '{}' has kind: Pure but has ToolHandle input", + node.id.0 + ))); + } + if port.name.is_resource() { + return Err(ExecError::new(format!( + "node '{}' has kind: Pure but has resource input '{}'", + node.id.0, port.name.0 + ))); + } + } + for port in &node.outputs { + if port.type_id.0 == "ToolHandle" { + return Err(ExecError::new(format!( + "node '{}' has kind: Pure but has ToolHandle output", + node.id.0 + ))); + } + if matches!( + port.type_id.0.as_str(), + "FilesystemHandle" | "NetworkHandle" | "Timestamp" | "Credential" | "Platform" + ) { + return Err(ExecError::new(format!( + "node '{}' has kind: Pure but has resource-environment output '{}'", + node.id.0, port.type_id.0 + ))); + } } } Ok(()) diff --git a/core/exec/src/intercept.rs b/core/exec/src/intercept.rs index 996acb6987e..abb85b0464c 100644 --- a/core/exec/src/intercept.rs +++ b/core/exec/src/intercept.rs @@ -218,44 +218,6 @@ impl BoundaryMocks { self.input_mocks.iter() } - /// Propagate input values to `param_source_*` entrypoint nodes (RT58). - /// - /// `param_source_*` nodes are interior identity nodes created by the lowerer - /// for callable parameters. They have the same port names as top-level - /// entrypoints but different node IDs. This method copies values from - /// non-param-source entrypoints to matching param_source nodes that don't - /// already have a value set. - /// - /// Call this after setting up all CLI/test inputs to ensure param_source - /// nodes receive the same values as their corresponding callable parameters. - pub fn propagate_to_param_sources( - &mut self, - entrypoint_ports: &[(NodeId, PortName, gunbc_ir::TypeId)], - ) { - // Collect (port_name → value) from non-param-source entrypoints. - let mut port_values: HashMap = HashMap::new(); - for (node_id, port_name, _) in entrypoint_ports { - if node_id.0.starts_with("param_source_") { - continue; - } - if let Some(val) = self.get_input(&node_id.0, &port_name.0) { - port_values.insert(port_name.0.clone(), val.clone()); - } - } - - // Copy to param_source nodes that don't already have a value. - for (node_id, port_name, _) in entrypoint_ports { - if !node_id.0.starts_with("param_source_") { - continue; - } - if self.get_input(&node_id.0, &port_name.0).is_some() { - continue; - } - if let Some(val) = port_values.get(&port_name.0) { - self.set_input(node_id.0.clone(), port_name.0.clone(), val.clone()); - } - } - } } #[cfg(test)] @@ -369,89 +331,4 @@ mod tests { assert_eq!(mock.next_value(), Value::Str("second".into())); } - #[test] - fn test_propagate_to_param_sources() { - use gunbc_ir::{NodeId, PortName, TypeId}; - - let mut mocks = BoundaryMocks::new(); - // Set input for a regular entrypoint node - mocks.set_input("entry_owner", "owner", Value::Str("gunb-ai".into())); - mocks.set_input("entry_repo", "repo", Value::Str("gunbc".into())); - - let entrypoint_ports = vec![ - ( - NodeId("entry_owner".into()), - PortName("owner".into()), - TypeId("String".into()), - ), - ( - NodeId("entry_repo".into()), - PortName("repo".into()), - TypeId("String".into()), - ), - // param_source nodes share port names but have different node IDs - ( - NodeId("param_source_dispatch_owner".into()), - PortName("owner".into()), - TypeId("String".into()), - ), - ( - NodeId("param_source_dispatch_repo".into()), - PortName("repo".into()), - TypeId("String".into()), - ), - ]; - - // Before propagation, param_source nodes have no input - assert!(mocks - .get_input("param_source_dispatch_owner", "owner") - .is_none()); - - mocks.propagate_to_param_sources(&entrypoint_ports); - - // After propagation, param_source nodes receive copied values - assert_eq!( - mocks.get_input("param_source_dispatch_owner", "owner"), - Some(&Value::Str("gunb-ai".into())) - ); - assert_eq!( - mocks.get_input("param_source_dispatch_repo", "repo"), - Some(&Value::Str("gunbc".into())) - ); - } - - #[test] - fn test_propagate_to_param_sources_does_not_overwrite_explicit() { - use gunbc_ir::{NodeId, PortName, TypeId}; - - let mut mocks = BoundaryMocks::new(); - mocks.set_input("entry_owner", "owner", Value::Str("gunb-ai".into())); - // Explicitly set a different value on the param_source node - mocks.set_input( - "param_source_dispatch_owner", - "owner", - Value::Str("custom-owner".into()), - ); - - let entrypoint_ports = vec![ - ( - NodeId("entry_owner".into()), - PortName("owner".into()), - TypeId("String".into()), - ), - ( - NodeId("param_source_dispatch_owner".into()), - PortName("owner".into()), - TypeId("String".into()), - ), - ]; - - mocks.propagate_to_param_sources(&entrypoint_ports); - - // Explicit value must NOT be overwritten - assert_eq!( - mocks.get_input("param_source_dispatch_owner", "owner"), - Some(&Value::Str("custom-owner".into())) - ); - } } diff --git a/core/exec/src/progress.rs b/core/exec/src/progress.rs index 967291cfd9e..bd6df66d57e 100644 --- a/core/exec/src/progress.rs +++ b/core/exec/src/progress.rs @@ -187,6 +187,14 @@ impl OutputSummary { Value::Set(s) => (FieldKind::List(s.len()), format!("{{{} items}}", s.len())), Value::Float(f) => (FieldKind::Number, f.to_string()), Value::Bytes(b) => (FieldKind::Scalar, format!("<{} bytes>", b.len())), + Value::Enum { ty, variant } => { + let preview = if ty.is_empty() { + variant.clone() + } else { + format!("{ty}.{variant}") + }; + (FieldKind::Scalar, preview) + } }; Some(FieldSummary { name: name.clone(), diff --git a/core/ir/src/value.rs b/core/ir/src/value.rs index 4ebd80f0edd..fbb3f9c4ecd 100644 --- a/core/ir/src/value.rs +++ b/core/ir/src/value.rs @@ -169,6 +169,11 @@ pub enum Value { Response(TransportResponse), /// Secret value (redacted in logs/display, exposed only at I/O boundaries) Secret(SecretString), + /// Typed enum variant (sum-type unit variant). + /// + /// `ty` is the enum type name (e.g., "TestClass"), and `variant` is the + /// selected variant (e.g., "Hermetic"). + Enum { ty: String, variant: String }, /// Node was skipped (guard evaluated to false) Skipped, } @@ -192,6 +197,7 @@ pub enum ValueKind { TransportRequest, TransportResponse, Secret, + Enum, Skipped, } @@ -212,6 +218,7 @@ impl ValueKind { ValueKind::TransportRequest => "TransportRequest", ValueKind::TransportResponse => "TransportResponse", ValueKind::Secret => "Secret", + ValueKind::Enum => "Enum", ValueKind::Skipped => "Skipped", } } @@ -240,6 +247,7 @@ impl Value { Value::Request(_) => ValueKind::TransportRequest, Value::Response(_) => ValueKind::TransportResponse, Value::Secret(_) => ValueKind::Secret, + Value::Enum { .. } => ValueKind::Enum, Value::Skipped => ValueKind::Skipped, } } @@ -408,6 +416,7 @@ impl Value { Value::List(v) | Value::Set(v) => v.is_empty(), Value::Map(m) => m.is_empty(), Value::Secret(s) => s.is_empty(), + Value::Enum { .. } => false, Value::Bool(_) | Value::Int(_) | Value::Float(_) @@ -705,6 +714,16 @@ impl PartialEq for Value { (Value::Request(a), Value::Request(b)) => a == b, (Value::Response(a), Value::Response(b)) => a == b, (Value::Secret(a), Value::Secret(b)) => a == b, + ( + Value::Enum { + ty: a_ty, + variant: a_variant, + }, + Value::Enum { + ty: b_ty, + variant: b_variant, + }, + ) => a_ty == b_ty && a_variant == b_variant, (Value::Skipped, Value::Skipped) => true, _ => false, } @@ -727,6 +746,7 @@ impl fmt::Display for Value { Value::Request(r) => write!(f, "", std::mem::discriminant(r)), Value::Response(r) => write!(f, "", std::mem::discriminant(r)), Value::Secret(_) => write!(f, "***"), + Value::Enum { ty, variant } => write!(f, "{ty}.{variant}"), Value::Skipped => write!(f, ""), } } diff --git a/core/ir/src/value_bridge.rs b/core/ir/src/value_bridge.rs index c7fd263c65b..6a62a09c3ea 100644 --- a/core/ir/src/value_bridge.rs +++ b/core/ir/src/value_bridge.rs @@ -47,6 +47,7 @@ pub fn classify_value(value: &Value) -> ValueCategory { Value::Json(_) | Value::Secret(_) | Value::Float(_) | Value::Bytes(_) => { ValueCategory::Shared } + Value::Enum { .. } => ValueCategory::Shared, Value::Map(_) | Value::Set(_) | Value::Request(_) | Value::Response(_) | Value::Skipped => { ValueCategory::GunbcOnly } @@ -90,6 +91,12 @@ pub fn to_bridge_json(value: &Value) -> Option { "__bytes": b.len(), })), Value::Secret(_) => Some(serde_json::Value::String("***".to_string())), + Value::Enum { ty, variant } => Some(serde_json::json!({ + "__enum": { + "ty": ty, + "variant": variant + } + })), Value::Request(_) | Value::Response(_) | Value::Skipped => None, } } @@ -112,7 +119,23 @@ pub fn from_bridge_json(json: &serde_json::Value) -> Value { } } serde_json::Value::Array(arr) => Value::List(arr.iter().map(from_bridge_json).collect()), - serde_json::Value::Object(_) => Value::Json(json.clone()), + serde_json::Value::Object(obj) => { + if let Some(enum_obj) = obj.get("__enum").and_then(|v| v.as_object()) { + let ty = enum_obj + .get("ty") + .and_then(|v| v.as_str()) + .unwrap_or_default() + .to_string(); + let variant = enum_obj + .get("variant") + .and_then(|v| v.as_str()) + .unwrap_or_default() + .to_string(); + Value::Enum { ty, variant } + } else { + Value::Json(json.clone()) + } + } } } diff --git a/core/ir/src/value_expr.rs b/core/ir/src/value_expr.rs index f0785b5a9de..f7e857b86e0 100644 --- a/core/ir/src/value_expr.rs +++ b/core/ir/src/value_expr.rs @@ -36,6 +36,8 @@ pub enum ValueExpr { }, /// Secret/redacted value — rendered as a secret constructor per language. Secret(String), + /// Typed enum unit variant. + Enum { ty: String, variant: String }, /// Node was skipped (guard false) — rendered as skip sentinel per language. Skipped, } @@ -65,6 +67,10 @@ impl From<&Value> for ValueExpr { Value::Float(f) => ValueExpr::Json(serde_json::json!(*f)), Value::Bytes(b) => ValueExpr::Json(serde_json::json!({"__bytes": b.len()})), Value::Secret(_) => ValueExpr::Secret("***".to_string()), + Value::Enum { ty, variant } => ValueExpr::Enum { + ty: ty.clone(), + variant: variant.clone(), + }, Value::Skipped => ValueExpr::Skipped, } } diff --git a/core/resolve/Cargo.toml b/core/resolve/Cargo.toml new file mode 100644 index 00000000000..2b9ca7a5d11 --- /dev/null +++ b/core/resolve/Cargo.toml @@ -0,0 +1,12 @@ +[package] +name = "gunbc-resolve" +version.workspace = true +edition.workspace = true +license.workspace = true +description = "Shared DAG resolution operations for service transports" + +[dependencies] +daglang-lower = { path = "../daglang/daglang-lower" } +gunbc-exec = { path = "../exec" } +gunbc-ir = { path = "../ir" } +serde_json = { workspace = true } diff --git a/core/resolve/src/lib.rs b/core/resolve/src/lib.rs new file mode 100644 index 00000000000..b394ca06719 --- /dev/null +++ b/core/resolve/src/lib.rs @@ -0,0 +1 @@ +pub mod service_ops; diff --git a/core/resolve/src/service_ops.rs b/core/resolve/src/service_ops.rs new file mode 100644 index 00000000000..da005ad6005 --- /dev/null +++ b/core/resolve/src/service_ops.rs @@ -0,0 +1,6 @@ +// Transitional extraction for C11: compile the implementation from its current +// source file while `gunbc-dag` consumes it through this crate boundary. +#[path = "../../../gunbc-dag/src/resolve_service.rs"] +mod service_ops_impl; + +pub use service_ops_impl::*; diff --git a/core/testgen-registry/src/lib.rs b/core/testgen-registry/src/lib.rs index 65d4fa8331e..bd66f3d4dec 100644 --- a/core/testgen-registry/src/lib.rs +++ b/core/testgen-registry/src/lib.rs @@ -149,7 +149,7 @@ pub fn iter_resource_tests() -> impl Iterator { /// /// This is the single codegen path — all targets use this function. /// Per-target variation is only in which DAG and MockSpec are provided. -pub fn generate_target( +pub fn generate_target( config: &TestgenTargetDef, dag: Dag, spec: MockSpec, @@ -160,7 +160,7 @@ pub fn generate_target( /// Like [`generate_target`] but merges a DSL-extracted type registry into the /// core type registry, making DSL-defined sum/product types visible to testgen /// for variant coverage obligations. -pub fn generate_target_with_types( +pub fn generate_target_with_types( config: &TestgenTargetDef, dag: Dag, spec: MockSpec, diff --git a/gunbc-dag/Cargo.toml b/gunbc-dag/Cargo.toml index 84a494ddb77..7c535270ca2 100644 --- a/gunbc-dag/Cargo.toml +++ b/gunbc-dag/Cargo.toml @@ -10,6 +10,7 @@ description = "gunbc repo-specific DAG configuration, CI, and generators" gunbc-infra = { path = "../core/infra" } gunbc-ir = { path = "../core/ir" } gunbc-exec = { path = "../core/exec" } +gunbc-resolve = { path = "../core/resolve" } gunbc-codegen = { path = "../core/codegen" } gunbc-test = { path = "../core/test" } diff --git a/gunbc-dag/src/bin/sdlc.rs b/gunbc-dag/src/bin/sdlc.rs index f6e91bfb370..8f3fde9f70c 100644 --- a/gunbc-dag/src/bin/sdlc.rs +++ b/gunbc-dag/src/bin/sdlc.rs @@ -139,9 +139,6 @@ fn main() { } } - // RT58: propagate entrypoint values to param_source_* nodes. - input_mocks.propagate_to_param_sources(&entrypoints.entrypoint_ports); - // ======================================================================== // Set up execution mode // ======================================================================== diff --git a/gunbc-dag/src/lib.rs b/gunbc-dag/src/lib.rs index f0ab0cff8b2..54d5e4bd291 100644 --- a/gunbc-dag/src/lib.rs +++ b/gunbc-dag/src/lib.rs @@ -39,7 +39,6 @@ pub mod mock_defaults; pub mod policy; pub mod pragma; pub mod resolve; -pub mod resolve_service; pub mod resources; pub mod testgen_dag; pub mod tool_runner; diff --git a/gunbc-dag/src/mock_defaults.rs b/gunbc-dag/src/mock_defaults.rs index 621397ce446..6510d7f117d 100644 --- a/gunbc-dag/src/mock_defaults.rs +++ b/gunbc-dag/src/mock_defaults.rs @@ -20,68 +20,6 @@ fn default_fs_handle() -> Value { fs.into() } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum GcpFieldKind { - Audience, - Project, - Secret, - SubjectToken, - Version, - ServiceAccount, -} - -fn typed_gcp_field_kind(type_id: &str) -> Option { - match type_id { - // Preferred modeling path: refined type aliases encode intent directly. - "OidcAudience" | "WifAudience" => Some(GcpFieldKind::Audience), - "GcpProjectId" => Some(GcpFieldKind::Project), - "GcpSecretId" => Some(GcpFieldKind::Secret), - "GcpSubjectToken" | "OidcSubjectToken" => Some(GcpFieldKind::SubjectToken), - "GcpSecretVersion" => Some(GcpFieldKind::Version), - "GcpServiceAccountEmail" => Some(GcpFieldKind::ServiceAccount), - _ => None, - } -} - -fn gcp_field_value_for_kind(kind: GcpFieldKind) -> Value { - match kind { - GcpFieldKind::Audience => Value::Str("mock-audience".to_string()), - GcpFieldKind::Project => Value::Str("mock-project".to_string()), - GcpFieldKind::Secret => Value::Str("mock-secret".to_string()), - GcpFieldKind::SubjectToken => Value::Str("mock-subject-token".to_string()), - GcpFieldKind::Version => Value::Str("latest".to_string()), - GcpFieldKind::ServiceAccount => { - Value::Str("mock-sa@mock-project.iam.gserviceaccount.com".to_string()) - } - } -} - -/// Legacy compatibility path while graph ports migrate to refined type aliases. -fn legacy_gcp_field_value_by_name(port_name: &str) -> Option { - match port_name { - "audience" => Some(gcp_field_value_for_kind(GcpFieldKind::Audience)), - "project" => Some(gcp_field_value_for_kind(GcpFieldKind::Project)), - "secret" | "secret_name" => Some(gcp_field_value_for_kind(GcpFieldKind::Secret)), - "subject_token" => Some(gcp_field_value_for_kind(GcpFieldKind::SubjectToken)), - "version" => Some(gcp_field_value_for_kind(GcpFieldKind::Version)), - "service_account" | "service_account_or_role" => { - Some(gcp_field_value_for_kind(GcpFieldKind::ServiceAccount)) - } - _ => None, - } -} - -/// Returns a realistic GCP mock value when semantic intent is encoded in the -/// type ID (preferred) or, as a compatibility path, in legacy port names. -/// -/// This avoids GCP prepare ops falling back to `"(unresolved)"` when their -/// inputs are optional or entrypoint ports filled with generic `"mock"`. -fn gcp_field_value(type_id: &str, port_name: &str) -> Option { - typed_gcp_field_kind(type_id) - .map(gcp_field_value_for_kind) - .or_else(|| legacy_gcp_field_value_by_name(port_name)) -} - fn default_value_for_type(type_id: &str) -> Value { match type_id { "TransportResponse" => default_shell_response(), @@ -372,8 +310,7 @@ pub fn auto_mock_spec(dag: &Dag, name: &str) -> .map(|input| input.cardinality) }) .unwrap_or(Cardinality::ONE); - let value = gcp_field_value(type_id.0.as_str(), port_name.0.as_str()) - .unwrap_or_else(|| default_value_for_port(type_id.0.as_str(), cardinality)); + let value = default_value_for_port(type_id.0.as_str(), cardinality); spec = spec.input_mock(node_id.0.as_str(), port_name.0.as_str(), value); } @@ -400,13 +337,7 @@ pub fn auto_mock_spec(dag: &Dag, name: &str) -> if output_names.contains(input_port.name.0.as_str()) { continue; } - // Inject known GCP field values even for optional inputs. if input_port.cardinality.allows_empty() { - if let Some(value) = - gcp_field_value(input_port.type_id.0.as_str(), input_port.name.0.as_str()) - { - required_inputs.insert(input_port.name.0.clone(), value); - } continue; } if input_port.type_id.0 == "TransportResponse" { @@ -416,13 +347,7 @@ pub fn auto_mock_spec(dag: &Dag, name: &str) -> let value = if input_port.name.0 == "skip" && input_port.type_id.0 == "Bool" { Value::Bool(false) } else { - gcp_field_value(input_port.type_id.0.as_str(), input_port.name.0.as_str()) - .unwrap_or_else(|| { - default_value_for_port( - input_port.type_id.0.as_str(), - input_port.cardinality, - ) - }) + default_value_for_port(input_port.type_id.0.as_str(), input_port.cardinality) }; required_inputs.insert(input_port.name.0.clone(), value); } @@ -547,33 +472,6 @@ fn probe_best_response( mod tests { use super::*; - #[test] - fn gcp_field_value_prefers_typed_hints_over_port_names() { - assert_eq!( - gcp_field_value("GcpProjectId", "renamed_project"), - Some(Value::Str("mock-project".to_string())) - ); - assert_eq!( - gcp_field_value("GcpServiceAccountEmail", "identity"), - Some(Value::Str( - "mock-sa@mock-project.iam.gserviceaccount.com".to_string() - )) - ); - assert_eq!( - gcp_field_value("GcpSubjectToken", "token"), - Some(Value::Str("mock-subject-token".to_string())) - ); - } - - #[test] - fn gcp_field_value_keeps_legacy_name_fallback_for_string_ports() { - assert_eq!( - gcp_field_value("String", "project"), - Some(Value::Str("mock-project".to_string())) - ); - assert_eq!(gcp_field_value("String", "not_a_gcp_field"), None); - } - #[test] fn auto_mock_spec_ci_graph_produces_transport_mocks() { let dag = crate::ci::build_ci_graph().expect("build ci graph"); diff --git a/gunbc-dag/src/resolve.rs b/gunbc-dag/src/resolve.rs index 1e64c5ace04..1b832f5c890 100644 --- a/gunbc-dag/src/resolve.rs +++ b/gunbc-dag/src/resolve.rs @@ -40,9 +40,10 @@ use gunbc_lib_blob::BlobOps; use gunbc_lib_transport::TransportOps; use gunbc_primitives::{filename, FsEnv}; -use crate::resolve_service::{ +use gunbc_resolve::service_ops::{ GenericFileParseOp, GenericFilePrepareOp, GenericLocalParseOp, GenericLocalPrepareOp, GenericRestParseOp, GenericRestPrepareOp, GenericShellParseOp, GenericShellPrepareOp, + InterfaceStubExecuteOp, InterfaceStubParseOp, InterfaceStubPrepareOp, }; // ============================================================================ @@ -89,93 +90,18 @@ fn execute_with_declared_output_passthrough( outputs.insert(port_name.clone(), value.clone()); continue; } - outputs.entry(port_name.clone()).or_insert_with(|| { - if let Some(fallback) = passthrough_fallback_value(port_name, &inputs) { - return fallback; - } - if *is_optional { - Value::Skipped - } else { - // RT83: Required output ports with no wired input indicate a - // lowering gap. Diagnostic suppressed from stderr to avoid noise - // in CI — tracked via structured obligations instead. - // TODO(RT83): return ExecError for required ports once all - // dag_util.dag if/else branches are fully wired. - Value::Skipped - } - }); + if *is_optional { + outputs.insert(port_name.clone(), Value::Skipped); + continue; + } + return Err(ExecError::new(format!( + "missing required declared output passthrough: `{}` (expected input `{}`)", + port_name, passthrough_key + ))); } Ok(outputs) } -fn passthrough_fallback_value(port_name: &str, inputs: &HashMap) -> Option { - let aliases: &[&str] = match port_name { - "result" => &["input", "value", "content", "document"], - "return" => &[ - "value", - "content", - "document", - "input", - "result", - "directives", - "sections", - "lines", - "text", - "items", - ], - _ => &[], - }; - - let candidate = aliases - .iter() - .find_map(|alias| inputs.get(*alias).cloned())?; - - if port_name != "return" { - return Some(candidate); - } - - Some(match candidate { - Value::Str(_) => candidate, - Value::Int(value) => Value::Str(value.to_string()), - Value::Bool(value) => Value::Str(value.to_string()), - Value::Float(value) => Value::Str(value.to_string()), - Value::Unit => Value::Str(String::new()), - Value::List(values) | Value::Set(values) => Value::Str( - values - .iter() - .map(passthrough_value_to_text) - .collect::>() - .join("\n"), - ), - Value::Map(values) => Value::Str(format!("{values:?}")), - Value::Json(value) => Value::Str(value.to_string()), - Value::Bytes(bytes) => Value::Str(format!("{bytes:?}")), - Value::Secret(secret) => Value::Str(secret.to_string()), - Value::Request(request) => Value::Str(format!("{request:?}")), - Value::Response(response) => Value::Str(format!("{response:?}")), - Value::Skipped => Value::Skipped, - }) -} - -fn passthrough_value_to_text(value: &Value) -> String { - match value { - Value::Str(value) => value.clone(), - Value::Int(value) => value.to_string(), - Value::Bool(value) => value.to_string(), - Value::Float(value) => value.to_string(), - Value::Unit => String::new(), - Value::Json(value) => value.to_string(), - Value::Map(value) => format!("{value:?}"), - Value::Bytes(value) => format!("{value:?}"), - Value::Secret(value) => value.to_string(), - Value::Request(value) => format!("{value:?}"), - Value::Response(value) => format!("{value:?}"), - Value::List(values) => format!("{values:?}"), - Value::Set(values) => format!("{values:?}"), - Value::Skipped => String::new(), - } -} - /// Identity callable op for DSL-compiled callables with fn bodies. /// /// Forwards all inputs to outputs, filling any declared output port that @@ -233,8 +159,20 @@ impl Executable for PipelineDispatchOp { /// if already at the last stage (terminal self-loop) /// - If `current_stage` is absent: no `next_stage` is emitted fn execute(&self, inputs: HashMap) -> Result, ExecError> { - let mut outputs = - execute_with_declared_output_passthrough(&self.output_ports, inputs)?; + // These outputs are computed by dispatch logic itself, not expected as + // passthrough inputs from upstream wiring. + let passthrough_ports: Vec<(String, bool)> = self + .output_ports + .iter() + .filter(|(name, _)| { + !matches!( + name.as_str(), + "stages" | "stage_order" | "active_stage" | "next_stage" + ) + }) + .cloned() + .collect(); + let mut outputs = execute_with_declared_output_passthrough(&passthrough_ports, inputs)?; outputs.insert("stages".to_string(), Value::Int(self.stage_count as i64)); outputs.insert( "stage_order".to_string(), @@ -394,15 +332,15 @@ impl Executable for LiteralSourceOp { } } -/// RT4a: Compute node for complex return expressions (BinOp, UnaryOp, etc.). +/// Compute node for lowered expression evaluation. /// Evaluates a `LoweredFnBody` using `evaluate_fn_body` with inputs from predecessor nodes. #[derive(Debug, Clone)] -struct ReturnExprComputeOp { +struct ExprComputeOp { fn_body: daglang_lower::LoweredFnBody, output_port: String, } -impl Executable for ReturnExprComputeOp { +impl Executable for ExprComputeOp { fn execute(&self, inputs: HashMap) -> Result, ExecError> { let sibling_fns = HashMap::new(); let result = daglang_lower::eval::evaluate_fn_body(&self.fn_body, &inputs, &sibling_fns) @@ -744,13 +682,12 @@ fn resolve_primitive(kind: &PrimitiveOpKind, outputs: &[Port]) -> Result Ok(DynOp::new(TransportOps::Execute)), // FC-7: Output path annotation nodes are metadata-only, resolve as identity. PrimitiveOpKind::ContentUpsertOutputPath { .. } => Ok(DynOp::new(IdentityCallableOp)), - // RT4a: Compute nodes evaluate complex return expressions via fn body evaluation. - PrimitiveOpKind::ReturnExprCompute { fn_body } => { + PrimitiveOpKind::ExprCompute { fn_body } => { let output_port = outputs .first() .map(|port| port.name.0.clone()) .unwrap_or_else(|| "result".to_string()); - Ok(DynOp::new(ReturnExprComputeOp { + Ok(DynOp::new(ExprComputeOp { fn_body: *fn_body.clone(), output_port, })) @@ -1019,7 +956,7 @@ fn resolve_service_transport( capability, }) = &metadata.spec { - return Ok(DynOp::new(crate::resolve_service::InterfaceStubExecuteOp { + return Ok(DynOp::new(InterfaceStubExecuteOp { interface: interface.clone(), capability: capability.clone(), })); @@ -1086,7 +1023,7 @@ fn resolve_service_transport( }, Some(TransportRole::Prepare), ) => { - return Ok(DynOp::new(crate::resolve_service::InterfaceStubPrepareOp { + return Ok(DynOp::new(InterfaceStubPrepareOp { interface: interface.clone(), capability: capability.clone(), })); @@ -1098,7 +1035,7 @@ fn resolve_service_transport( }, Some(TransportRole::Parse), ) => { - return Ok(DynOp::new(crate::resolve_service::InterfaceStubParseOp { + return Ok(DynOp::new(InterfaceStubParseOp { interface: interface.clone(), capability: capability.clone(), })); diff --git a/gunbc-dag/src/testgen_dag/profile_discovery.rs b/gunbc-dag/src/testgen_dag/profile_discovery.rs index fa73485bca9..55905323d31 100644 --- a/gunbc-dag/src/testgen_dag/profile_discovery.rs +++ b/gunbc-dag/src/testgen_dag/profile_discovery.rs @@ -167,6 +167,12 @@ fn collect_env_calls(expr: &Expr, env_vars: &mut Vec) { collect_env_calls(lhs, env_vars); collect_env_calls(rhs, env_vars); } + Expr::PipeCall(receiver, _, args) => { + collect_env_calls(receiver, env_vars); + for (_, arg) in args { + collect_env_calls(arg, env_vars); + } + } Expr::Record(_, fields) | Expr::Return(fields) => { for (_, value) in fields { collect_env_calls(value, env_vars); diff --git a/tasks.md b/tasks.md index 7cebb2b9258..5b8cffb22ed 100644 --- a/tasks.md +++ b/tasks.md @@ -521,25 +521,26 @@ spec.rs # Service operation specs | C1 | RT93 | **Stdlib host + caching.** `OnceLock` cache for compiled fn bodies. `include_str!` for stdlib sources. Single `StdLibHost::eval_fn()` interface. Delete per-module compile wrappers. | `classify_callable()` never calls `compile_from_context()`. No `../../dsl` paths. | M | | C2 | RT42 | **Pipe methods first-class.** `PipeMethod` enum in syntax. Parser resolves `\|> method()` to `PipeCall(PipeMethod, ...)`. Delete `should_track_call_name()` allowlist. | Allowlist deleted. `PipeMethod` has all 20 methods. All `.dag` compile. | M | | C3 | RT45, RT46 | **Typed enums end-to-end.** `Value::Enum { ty, variant }`. Delete `TestClass::parse()` / `FermiCost::parse()` round-trips. Replace `unwrap_or()` fallbacks with errors. | Zero `parse()` on classification. Zero `unwrap_or()` in fidelity. | M | -| C4 | RT82 | **LoweringContext + dead code.** Context struct grouping 8-11 params. Delete 18 `#[allow(clippy::too_many_arguments)]`. Delete all dead `_ => None` in wiring. | Zero `too_many_arguments`. Zero `_ => None` in wiring. All `.dag` compile. | L | +| C4 | RT82 | **LoweringContext + dead code (staged).** Context struct grouping 8-11 params. Delete 18 `#[allow(clippy::too_many_arguments)]`. Delete dead `_ => None` arms only after complex-return coverage is proven (C10 RT4a/c). | Zero `too_many_arguments`. `_ => None` deletion is gated by C10 parity (BinOp/If/Match/Pipe returns still wire). All `.dag` compile. | L | | C5 | — | **Integrate scope.rs.** Replace `detect_*_branches_in_stmts`, `IfBranchSite`, `MatchBranchSite` with `ScopedBody`. Delete ad-hoc walk functions. | `IfBranchSite` deleted. `scope.rs` has non-test callers. DAG parity. | M | | C6 | — | **Extract transport derivation.** `transport.rs` module. Returns `TransportManifest` (pure data). Invariant: every service call site maps to exactly one triplet. | `add_service_transport_triplets` returns data, not mutates builder. | M | | C7 | — | **Expr walker totality + typed leaf refs.** Explicit arms for all `Expr` variants. `LeafRef` enum: `Param { name, field, ty }`, `Callable { endpoint, port }`, `Service { endpoint, port }`. | Zero `_ => {}` in expr walkers. `PARAM_REF_SENTINEL` deleted. | M | | C8 | RT84 | **Delete dead AST scaffolding.** `MockResponseDef`, `error_cases()`, `@retry`, orphaned `hermetic`. | `MockResponseDef` deleted. `@retry` rejected by parser. `hermetic` warns. | S | | C9 | RT38, RT39 | **No panics, no silent parse.** `LowerError::InvalidTransportSpec` replaces `panic!`. Parse error for bad `auth_input`. | Zero `panic!` on user DSL. Parser test for `auth_input: "token"`. | S | -| C10 | RT94 | **Resolve ReturnExprCompute split-brain.** Desugar complex returns into explicit DAG nodes so emit/interpret share semantics. Delete `MetadataOnly` for `ReturnExprCompute`. Delete `ReturnExprComputeOp`. | Zero `ReturnExprCompute` in any compiled graph. `PrimitiveOpKind::ReturnExprCompute` deleted. | L | -| C11 | RT67, RT72 | **Move resolve_service.rs to core/.** 2,190 lines → `core/resolve/src/service_ops.rs`. Split `resolve.rs`: generic framework (~1.6k) → `core/resolve/`. Domain dispatch (~700) stays. Must delete app-specific string dispatch in the move. | `resolve_service.rs` deleted from gunbc-dag. New `core/resolve/` crate. Moved code is simpler than source. | L | -| C12 | RT73 | **Move testgen to core/.** 5 files (2,177 lines) → `core/codegen/src/testgen/`. Must delete gunbc-dag-specific assumptions in the move. | `testgen_dag/` deleted from gunbc-dag. Testgen works from `core/codegen`. | M | +| C10 | RT94, RT4a, RT4c | **Resolve ReturnExprCompute split-brain + completeness gate.** Desugar complex returns (BinOp/UnaryOp/If/Match/Pipe/...) into explicit DAG semantics so emit/interpret share behavior. Delete `MetadataOnly` for `ReturnExprCompute`. Delete `ReturnExprComputeOp`. Add fail-closed compile-time gate for non-optional unwired returns (RT4c). | Zero `ReturnExprCompute` in any compiled graph. `PrimitiveOpKind::ReturnExprCompute` deleted. No silent return-binding drops (`_ => None`) for required outputs. | L | +| C10a | RT-I5 | **`make gist` auth credential bridge fix (postmortem Option A/B/C).** Pick one option and implement in `daglang-lower` and/or `resolve_service.rs` so `auth_token: Secret` reaches REST auth reliably. This must land before C11 extraction to avoid migrating a known bug into core. | `make gist` no longer 401s due to empty bearer token path. C11 is blocked until this passes. | M | +| C11 | RT67, RT72 | **Move resolve_service.rs to core/.** 2,190 lines → `core/resolve/src/service_ops.rs`. Split `resolve.rs`: generic framework (~1.6k) → `core/resolve/`. Domain dispatch (~700) stays. Must delete app-specific string dispatch in the move. **Precondition: C10a complete. Inventory linkage gate required (RF-INV1 or RF-INV2) before/with move.** | `resolve_service.rs` deleted from gunbc-dag. New `core/resolve/` crate. Moved code is simpler than source. No dropped registrations after crate-boundary move (force-link evidence or RF-INV2 replacement). | L | +| C12 | RT73 | **Move testgen to core/.** 5 files (2,177 lines) → `core/codegen/src/testgen/`. Must delete gunbc-dag-specific assumptions in the move. **Inventory linkage gate required (RF-INV1 or RF-INV2) before/with move.** | `testgen_dag/` deleted from gunbc-dag. Testgen works from `core/codegen`. Cross-crate registrations still discoverable (or inventory removed via RF-INV2 path). | M | | C13 | RT68 | **Split mock_defaults.** Generic probing (~350) → `core/test/`. Delete GCP blob (~230). | `mock_defaults.rs` deleted. Auto-mock works from `core/test`. | S | | C14 | RT89 | **REST status-code checking.** `GenericRestParseOp` checks status before field extraction. Non-2xx → error. | 401 → structured error (not "field missing"). Test: mock 401 → error has status. | M | | C15 | RT88 | **Fail-closed resolver audit.** Classify all `_ =>` fallbacks. Delete `passthrough_fallback_value()` (70 lines). | Zero undocumented fallbacks. `passthrough_fallback_value` deleted. | M | | C16 | RT95 | **Transport class in node metadata.** `ServiceTransportClass` in lowered nodes. Registry gen reads metadata, not `node_id.contains("shell")`. | `from_node_context` reads metadata, not substrings. | S | | C17 | RT96 | **Kill `propagate_to_param_sources`.** Fix boundary detection. Param source nodes auto-fed. | `propagate_to_param_sources` deleted. One port per input. | M | | C18 | — | **Executor dead code.** Delete `looks_effectful_without_kind()`. Delete unwired credential expiry plumbing. | Dead code deleted. `cargo clippy` clean. | S | -| C19 | RT83 | **Restore passthrough enforcement.** After C4+C5+C7 wire dag_util branches, restore `ExecError` for required outputs with no input. Code ref: `resolve.rs:99` TODO. | `resolve.rs` returns `ExecError`. CI clean (no unwired branches). | S | +| C19 | RT83, RT4b | **Restore passthrough enforcement + runtime fail-closed diagnostics.** After C4+C5+C7 wire dag_util branches, required outputs with no input must return `ExecError` (not `Skipped`) and emit clear diagnostics for missing declared passthroughs (RT4b). | `resolve.rs` returns `ExecError` for required missing outputs. Missing passthrough ports are diagnosable (no silent fallback). CI clean (no unwired branches). | S | | C20 | RT59, RT63 | **CLI generator: profile, mode, subcommand support.** Expose `available_profiles` in `CompileOutput`. Template generates `--profile` enum flag, `--mode ensure\|verify`, subcommand dispatch for multi-func modules. `KEY=VALUE` arg parsing for infra-style tools. Unblocks Worker A. | Generated CLI for `pipelines/sdlc.dag` accepts `--profile`. Generated CLI for multi-func modules has subcommands. | L | -**Chain**: C1 → C3; C2; C4 → C5 → C6 → C10; C7; C8; C9; C11 → C14 → C15 → C19; C12; C13; C16; C17; C18; C20 (early, unblocks A) +**Chain**: C1 → C3; C2; C10 (RT4a/c) → C4 → C5 → C6; C7; C8; C9; C10a → (RF-INV1 or RF-INV2 gate) → C11 → C14 → C15 → C19; C12; C13; C16; C17; C18; C20 (early, unblocks A) --- @@ -644,6 +645,8 @@ The prepare node treats `auth_token` like any other input field — it may end u Option C is the smallest fix that unblocks `make gist`. Options A/B are more principled for the long term. +**Queue injection**: This postmortem is now explicitly assigned to Worker C as **C10a (RT-I5)** and is a hard precondition for **C11** (`resolve_service.rs` extraction). + ### POSTMORTEM: `gunbc-ci` false failure — `overall_success: Skipped` **Symptom**: `gunbc-ci` reports "A required success check returned false" even when all build/test/clippy stages succeed. Pre-existing on `main`. `make ci` (via `gunbc-workflow`) passes because it uses a different code path. @@ -696,6 +699,8 @@ There is no test or validation that checks: "for every return expression binding RT4a is the real fix. RT4b/c are defense-in-depth so the class of bug can't recur silently. After RT4a, revert the ci.rs workaround back to `success_port: Some("overall_success")`. +**Queue injection**: RT4a/RT4c are now explicit in **C10**, RT4b is explicit in **C19**, and C10 is ordered before C4 dead-arm deletion to avoid BinOp/complex-return regressions. + --- ## Red Backlog @@ -743,6 +748,11 @@ Meanwhile `gunbc-ci` (`ci.rs`) transitively references them through GCP secret env vars. Current workaround: ci.yml is committed with secrets and not regenerated on every codegen pass. +**C11/C12 linkage trap**: moving large subsystems to `core/` crosses crate +boundaries and can trigger the same linker-drop behavior. Treat **RF-INV1 or +RF-INV2 as a hard gate** for C11/C12 so moved resolvers/testgen paths do not +silently lose registrations. + | ID | Fix | Size | Notes | |----|-----|------|-------| | RF-INV1 | **Force-link inventory crates in codegen binary.** Add explicit `use` references or `extern crate` for crates that register `DagSpecDef` with `live_required` secrets. Simplest fix, but fragile — adding a new crate with secrets requires updating codegen_cli.rs. | S | Quick fix. |