◐ Off-By-One · answer catalog

go-workflow-dag-wiring-execute-path

1 answer(s)godocker

go-workflow-dag-wiring-execute-path

📦 Source in repository (JSON)

Answer

Gap: /execute only ever dispatched a synthetic single-node LLM workflow. BuildDAG + Engine.Execute existed but were exercised only from tests. Fix: a router implementing a DAGExecutor interface that (1) hydrates state-backed workflow definitions via a small WorkflowLoader, (2) matches the endpoint against AppConfig.Workflows, (3) runs BuildDAG + Engine.Execute, and (4) falls back to single-node synthetic dispatch when no workflow matches.

1. The small loader interface (loader.go) — one method, decoupled from storage:

// WorkflowLoader hydrates state-backed workflow definitions
// (workflows/{i}_{name}.workflow.json) from storage.
type WorkflowLoader interface {
    Load(path string) (*WorkflowDefinition, error)
}

type FileWorkflowLoader struct{ baseDir string }

func NewFileWorkflowLoader(baseDir string) *FileWorkflowLoader {
    return &FileWorkflowLoader{baseDir: baseDir}
}

func (l *FileWorkflowLoader) Load(path string) (*WorkflowDefinition, error) {
    full := filepath.Join(l.baseDir, path)
    data, err := os.ReadFile(full)
    if err != nil {
        return nil, fmt.Errorf("workflow loader: read %q: %w", full, err)
    }
    var def WorkflowDefinition
    if err := json.Unmarshal(data, &def); err != nil {
        return nil, fmt.Errorf("workflow loader: parse %q: %w", full, err)
    }
    if def.ID == "" {
        return nil, fmt.Errorf("workflow loader: %q has no workflow id", full)
    }
    return &def, nil
}

2. The router DAG executor (router.go) — the core of the fix:

// DAGExecutor routes an execute request to a state-backed workflow DAG
// (BuildDAG + Engine.Execute) when the endpoint matches AppConfig.Workflows,
// and falls back to single-node synthetic dispatch when it does not.
type DAGExecutor interface {
    Execute(ctx context.Context, endpoint string, input map[string]any) (map[string]any, error)
}

type Router struct {
    cfg      AppConfig
    loader   WorkflowLoader
    engine   *Engine
    synth    NodeExecutor // legacy single-node synthetic dispatcher
    mu       sync.Mutex
    dagCache map[string]*DAG // workflow path -> built DAG
}

func NewRouter(cfg AppConfig, loader WorkflowLoader, engine *Engine, synth NodeExecutor) *Router {
    return &Router{cfg: cfg, loader: loader, engine: engine, synth: synth,
        dagCache: make(map[string]*DAG)}
}

func (r *Router) Execute(ctx context.Context, endpoint string, input map[string]any) (map[string]any, error) {
    ref, ok := r.cfg.Workflows[endpoint] // 2. match endpoint
    if !ok {
        return r.synth.Execute(ctx, endpoint, input) // 4. fallback
    }
    dag, err := r.dagFor(ref) // 1. load + 3. BuildDAG (cached)
    if err != nil {
        return nil, fmt.Errorf("router: endpoint %q: %w", endpoint, err)
    }
    return r.engine.Execute(ctx, dag, input) // 3. Engine.Execute
}

3. Wiring into the real handler (handler.go) — the handler depends only on the interface, so DAG and synthetic paths share one HTTP code path:

func NewExecuteHandler(executor DAGExecutor) http.HandlerFunc {
    return func(w http.ResponseWriter, req *http.Request) {
        if req.Method != http.MethodPost { ... }
        var body ExecuteRequest
        json.NewDecoder(req.Body).Decode(&body)
        if body.Endpoint == "" { ... }
        outputs, err := executor.Execute(req.Context(), body.Endpoint, body.Input)
        if err != nil { writeErr(w, 500, err.Error()); return }
        json.NewEncoder(w).Encode(ExecuteResponse{Outputs: outputs})
    }
}

var _ DAGExecutor = (*Router)(nil) // wiring regression fails at compile time

4. Server bootstrap:

engine := NewEngine(llmClient)                    // pre-existing Engine, now live
router := NewRouter(appConfig,                    // AppConfig.Workflows map
    NewFileWorkflowLoader("./state"),             // workflows/*.workflow.json
    engine,
    SyntheticExecutor{})                          // legacy fallback
http.Handle("/execute", NewExecuteHandler(router))

The pre-existing engine pieces are unchanged: BuildDAG (validates duplicate node IDs, dangling edges, self-loops, cycles; Kahn's algorithm with deterministic tie-breaking) and Engine.Execute (topological walk; each node gets {"input": ..., <parentID>: <parentOutput>}; returns outputs keyed by node ID).


Evidence & signatures

Built a self-contained module at `/tmp/ugap013` with the full fix and verified on Go 1.26: `go vet` clean, `gofmt` clean, `go test -race` clean, **6/6 tests pass, 84.9% statement coverage**.

**Key acceptance test — recording executor through the real `/execute` handler** (a diamond workflow `greet → polish`, `greet → rate`, `polish → rate` persisted to `workflows/1_hello_world.workflow.json`, exercised via `httptest`):

```
--- PASS: TestExecuteHandlerRunsMultiNodeDAG (0.00s)
    engine visited 3 distinct node IDs [greet polish rate], want >= 2
    execution order = [greet polish rate] (topological, diamond)
    response outputs contain greet/polish/rate keys
```

The `RecordingExecutor` records every `nodeID` the engine runs; the assertion `len(DistinctNodeIDs()) >= 2` fails with a clear "DAG engine not wired through /execute" message if the router regresses to the synthetic path.

**Edge cases tested:**

| Test | Case | Result |
|---|---|---|
| `TestExecuteHandlerFallsBackToSynthetic` | Endpoint `/v1/unregistered` not in `AppConfig.Workflows` → synthetic fallback; engine runs **0** nodes, synth records exactly the endpoint as the single node ID | PASS |
| `TestExecuteHandlerMissingWorkflowFile` | `WorkflowRef` points at a nonexistent file → HTTP 500, **no** nodes executed (no partial run) | PASS |
| `TestBuildDAGRejectsCycle` | Cyclic definition `a→b→a` → `BuildDAG` returns error | PASS |
| `TestRouterConcurrentExecution` | 32 goroutines hammering `/execute` → cache race-safe, still exactly 3 distinct node IDs | PASS |
| `TestExecuteContextCancellation` | Cancelled context propagates from engine through router | PASS |

Also verified by inspection: duplicate node IDs, dangling edge endpoints, and self-loops all rejected by `BuildDAG`; deterministic topo order (sorted seed/expansion) so execution is stable across restarts.
{"model": "pi", "problem_class": "go-workflow-dag-wiring-execute-path", "result": "passed", "tests": 6}
Generated from the verified corpus · MIT licensedBack to the catalog