Skip to content

feat(engine): standardize node outputs with O(1) sourceRefs passport - #86

Merged
Vamsi-o merged 1 commit into
mainfrom
standardizing-node-outputs
Sep 29, 2026
Merged

Vamsi-o merged 1 commit into
mainfrom
standardizing-node-outputs

Conversation

@TejaBudumuru3

Copy link
Copy Markdown
Contributor

Pull Request: Standardize Node Outputs and Implement O(1) Coordinate-Based Data Provenance

📌 Summary of Changes

This pull request completely refactors BuildFlow's workflow execution engine, data structures, and variable resolution logic to standardize data contracts across all nodes.

Previously, node execution on the main branch relied on untyped any signatures, disparate node-specific output wrappers (e.g. raw spreadsheet 2D arrays, isolated status objects), a sequential linear evaluation loop, and column-index heuristics for interpolation.

This PR establishes:

  1. The Universal Data Contract (ExecuteItem): An atomic record model encapsulating business payloads (json) and coordinate-based lineage metadata (sourceRefs).
  2. Unified 2D Output Matrices (ExecuteItem[][]): Enforces a [wireIndex][rowIndex] output topology across all node executors, unlocking true multi-pin routing.
  3. Constant-Time ($O(1)$) Variable Resolution: Replaces linear/heuristic scans with direct coordinate dereferencing against execution context storage.
  4. BFS DAG Engine Traversal: Replaces sequential evaluation with an edge-driven queue supporting complex, branched graphs.
  5. Interface Segregation for OAuth Nodes: Isolates credential-only context (SheetAuthContext) from execution context (NodeExecutionContext), ensuring type safety for UI metadata routes.

🏛️ Architectural Paradigm Comparison

Architectural Dimension main Branch (Legacy System) Current Branch (standardizing-node-outputs)
Data Contract Untyped any. Every node emitted completely arbitrary JSON shapes. Uniform ExecuteItem structure (json payload + sourceRefs coordinate mapping).
Node Output Type output?: any (Could be a single object, primitive, or custom wrapper). output?: ExecuteItem[][] (Strict 2D output matrix: [Pins][Items]).
Input Context inputData?: any (Nodes received untyped data blobs). items: ExecuteItem[] (Strictly typed array of input records).
Data Provenance Non-existent. Nodes had zero record of which upstream node or row generated the data. sourceRefs dictionary: Direct coordinate mapping (Record<NodeId, { wireIndex, rowIndex }>).
Variable Resolution Fragile heuristics: Inspected columns and rows arrays specific to Google Sheets. Constant-time $O(1)$ coordinate dereferencing directly against the execution store.
Workflow Traversal Flat sequential loop: for (const node of nodes). No support for branching or pins. Breadth-First Search (BFS) DAG Queue with pin-based edge slicing.
Looping Architecture Hardcoded loop logic inside executor.ts with retry delays and a custom loopResult object. Handled via array inputs and batch-aware executors emitting standard item matrices.

📂 Detailed File-by-File Technical Changes

1. packages/common/src/index.ts (Core Schema Contract)

  • Changes:
    • Defined and exported ExecuteItemSchema and ExecuteItem type.
  • Technical Details:
    export const ExecuteItemSchema = z.object({
      json: z.record(z.string(), z.any()),
      sourceRefs: z.record(
        z.string(),
        z.object({
          wireIndex: z.number().default(0),
          rowIndex: z.number()
        })
      ).optional()
    });
    
    export type ExecuteItem = z.infer<typeof ExecuteItemSchema>;
  • Rationale: Establishes the core data exchange model. sourceRefs is structured as a dictionary (Record<NodeId, Coordinates>) rather than an array, enabling $O(1)$ hash-map lookups instead of $O(N)$ linear array searching.

2. packages/nodes/src/registry/Execution.config.types.ts (Execution Contracts)

  • Changes:
    • Replaced inputData?: any with items: ExecuteItem[].
    • Replaced output?: any with output?: ExecuteItem[][].
    • Marked nodeId: string as strictly required (removed optional flag ?).
  • Technical Details:
    export interface ExecutionContext {
      nodeId: string;
      userId: string;
      credentialId?: string;
      config: Record<string, any>;
      items: ExecuteItem[];
    }
    
    export interface ExecutionResult {
      success: boolean;
      output?: ExecuteItem[][];
      error?: string;
      metadata?: Record<any, any>;
    }
  • Rationale: Enforces strict compile-time type safety across all node implementations. Eliminates the risk of missing node IDs causing provenance corruption.

3. packages/nodes/src/google-sheets/google-sheets.executor.ts (Auth Isolation & Output Stamping)

  • Changes:
    • Created SheetAuthContext and made NodeExecutionContext extend SheetAuthContext.
    • Updated executeReadRows to transform spreadsheet rows into individual ExecuteItem instances and stamp initial coordinates.
    • Updated executeWriteRows, executeAppendRows, and executeClearRows to preserve sourceRefs: item.sourceRefs and return [outputBoxes].
    • Updated helper methods (getSheets, getSheetTabs, getHeaderRow, ensureSheetService) to accept SheetAuthContext.
  • Technical Details:
    interface SheetAuthContext {
      userId: string;
      credentialId: string;
      authType?: string;
    }
    
    interface NodeExecutionContext extends SheetAuthContext {
      nodeId: string;
      config: any;
      items: ExecuteItem[];
    }
  • Rationale: Applies the Interface Segregation Principle. Metadata dropdown routes (e.g. /getSheets) only need authentication credentials, while workflow step execution strictly requires complete execution parameters.

4. packages/nodes/src/gmail/gmail.executor.ts (Consumer Node Output Matrix)

  • Changes:
    • Replaced single execution returning an untyped summary object with an input iteration loop.
    • Emits an array of ExecuteItem instances containing execution response data while forwarding incoming sourceRefs: item.sourceRefs.
    • Returns output: [outputBoxes].

5. packages/nodes/src/filter/filter.executor.ts (Multi-Pin Output Partitioning)

  • Changes:
    • Restructured the return shape from a custom { filteredData, discardedData, metadata } object into a true 2D output matrix:
    const wire0 = filteredData.map(item => ({ json: item }));
    const wire1 = discardedData.map(item => ({ json: item }));
    
    return {
      success: true,
      output: [wire0, wire1], // Wire 0 = Kept records, Wire 1 = Discarded records
      metadata: {
        operation_used: operation,
        items_kept: filteredData.length,
        items_discard: discardedData.length
      }
    };
  • Rationale: Enables multi-pin routing. Downstream edges connected to Pin 0 receive matching records, while edges connected to Pin 1 receive discarded records. Telemetry is cleanly separated into top-level metadata.

6. packages/common/src/interpolation.ts (O(1) Variable Resolution)

  • Changes:
    • Extended resolveVariable, interpolateString, and resolveConfigVariables to accept sourceRefs?: Record<string, { wireIndex: number; rowIndex: number }>.
    • Replaced Google Sheets-specific column indexing with coordinate-based dereferencing:
    if (sourceRefs && sourceRefs[nodeName]) {
      const { wireIndex, rowIndex } = sourceRefs[nodeName];
      const item = nodeData[wireIndex]?.[rowIndex];
      const value = getNestedValue(item?.json, path);
      return value !== undefined ? value : `{{${variable}}}`;
    }
    • Added fallback resolution for trigger/webhook nodes (wire 0, row 0) and graceful preservation of unresolved placeholders {{variable}}.

7. apps/worker/src/engine/executor.ts (Engine Normalization & BFS Integration)

  • Changes:
    • Replaced sequential for loop with a Breadth-First Search (BFS) DAG Queue (while (queue.length > 0)).
    • Added defensive dimensional normalization for currentInputData:
    let itemsToProcess: any[] = [];
    if (Array.isArray(currentInputData)) {
      if (currentInputData.length > 0 && Array.isArray(currentInputData[0])) {
        itemsToProcess = currentInputData[0]; // Normalizes 2D output matrix to wire 0
      } else {
        itemsToProcess = currentInputData;    // Preserves 1D batch array
      }
    } else {
      itemsToProcess = currentInputData ? [{ json: currentInputData }] : [];
    }
    • Moved itemsToProcess extraction prior to configuration resolution, supplying itemsToProcess[0]?.sourceRefs to resolveConfigVariables.
    • Slices output pins dynamically when enqueueing downstream tasks:
    const branchData = execute.output?.[outputPinIndex] || [];
    queue.push({ nodeId: edge.target, inputData: branchData });

8. apps/http-backend/src/routes/userRoutes/executionRoutes.ts (Single-Node Testing Compliance)

  • Changes:
    • Injected nodeId: nodeData.id and items: [] into the test context object.
    • Corrected failure logging from error: executionResult.output to error: executionResult.error.
  • Rationale: Ensures single-node testing in the UI strictly conforms to the ExecutionContext contract and prevents node provenance from recording sourceRefs["undefined"].

🧪 Verification & Testing Results

  1. Monorepo Build Compilation:

    • Executed turbo build across all 11 packages (@repo/common, @repo/nodes, @repo/worker, @repo/db, @repo/processor, hooks, http-backend, web).
    • Result: Exit code 0 (Zero TypeScript or build errors).
  2. Runtime Verification Suite:

    • 17 out of 17 automated runtime tests passed:
      • Schema validation with optional and default wireIndex coordinates.
      • Coordinate-based multi-row template resolution (Alice, Bob, Charlie).
      • Fallback resolution for trigger nodes lacking initial provenance metadata.
      • Unresolved placeholder preservation.
      • Dimensional normalization across 2D arrays, 1D arrays, and raw payloads.
  3. Multi-Pin Executor Verification:

    • Validated FilterExecutor with test dataset:
      • Wire 0 (Unique Records): 2 items
      • Wire 1 (Discarded Records): 1 item
      • Metadata: Correctly emitted operation metrics.
  4. Live Execution Test:

    • Ran local development environment (pnpm dev) with live Kafka broker.
    • End-to-end execution completed seamlessly across the worker engine pipeline.

🚀 Next Steps

  • Update the Frontend Web UI (Variable Panel and Execution History components) to consume the new 2D matrix structure and cleanly expose .json properties to the canvas.
  • Connect multi-pin handles on the frontend canvas to wire handles 0 and 1 for the Splitter/Filter nodes.

@coderabbitai

coderabbitai Bot commented Sep 28, 2026

Copy link
Copy Markdown

Important

  • 🔍 Trigger review

This repository does not receive automatic reviews because it has fewer than 10 stars.

⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Advanced

Run ID: 438fbcd1-b711-47d7-b36d-98e455c4eaf4


Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@Vamsi-o Vamsi-o left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Amazing work

@Vamsi-o
Vamsi-o merged commit f55ed37 into main Sep 29, 2026
2 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants