Genkit Flows as senro Steps: Extremely Powerful AI Pipelines in Go A developer demonstrated wrapping Genkit flows inside senro pipeline steps, turning each Gemini model call into a graph node with its own log, state, retry policy and cache entry. The approach, tested against the real Gemini API with googleai/gemini-2.5-flash on Genkit Go v1.13.1 and senro v1.4.0, addresses the orchestration scaffolding — worker pools, retry loops, and resumability — that flows alone do not provide when summarizing a corpus of release notes into a digest. A Genkit https://genkit.dev flow is an excellent unit of AI work. It is typed, traced, testable, and genkit.DefineFlow gives you something you can call from anywhere. What a flow is not is a unit of orchestration . The moment you have twelve documents to summarise, a flow whose output feeds two other flows, or a model call that fails on the third of forty items, you end up writing the same scaffolding again: a worker pool, a retry loop, a way to avoid re-running the eleven items that already succeeded, logs you can actually read, and some way to find out what happened at three in the morning. senro https://github.com/xavidop/senro is a pipeline engine where a step already has all of that: retries that distinguish infrastructure from workload, an action cache that skips work it has already done, an append-only event stream, and a live attach protocol. I covered it in the introduction https://xavidop.me/go/2026-09-11-introducing-senro-pipeline-engine-go/ and the monorepo article https://xavidop.me/go/2026-09-11-senro-real-pipeline-monorepo-go/ . This article shows how to combine them. Put a Genkit flow inside a senro step and every model call becomes a graph node with its own log, state, retry policy and cache entry. Every example below ran against the real Gemini API with googleai/gemini-2.5-flash , on Genkit Go v1.13.1 and senro v1.4.0. The summaries are what the model actually returned. go get github.com/firebase/genkit/go go get github.com/xavidop/senro The example is a small corpus of release notes in docs/ .md that gets summarised one document at a time, then merged into a digest. Nothing senro-specific here. This is an ordinary Genkit flow: // SummarizeInput is what the summarize flow takes. type SummarizeInput struct { Path string json:"path" Text string json:"text" } // Summary is what it returns. type Summary struct { Path string json:"path" Summary string json:"summary" } func defineFlows g genkit.Genkit core.Flow SummarizeInput, Summary, struct{} { return genkit.DefineFlow g, "summarize", func ctx context.Context, in SummarizeInput Summary, error { text, err := genkit.GenerateText ctx, g, ai.WithSystem "You summarize release notes for an operator. "+ "Answer in one sentence, no preamble." , ai.WithPrompt fmt.Sprintf "Summarize these release notes:\n\n%s", in.Text if err = nil { return Summary{}, err } return Summary{Path: in.Path, Summary: text}, nil } } Called on its own it does what you expect: The billing reconciliation job was moved to a queue worker, accelerating invoice settlement from up to 24 hours to within two minutes. senro has two kinds of step: a command, and a registered Go function. A Genkit flow is a Go function, so it is the second kind: // summarize is the registered flow. A func step runs in the pipeline's own // process, so a package-level handle is all the wiring the two need. var summarize core.Flow SummarizeInput, Summary, struct{} type SummarizeParams struct { Path string json:"path" Name string json:"name" } func init { senro.RegisterFunc "ai/summarize", RunSummarize } // RunSummarize runs the flow and writes the summary into the one workspace // this step mounts. func RunSummarize ctx senro.Ctx, p SummarizeParams error { dir, ok := ctx.Workspace "doc-" + p.Name if ok { return fmt.Errorf "step %s mounts no doc-%s workspace", ctx.StepID , p.Name } text, err := os.ReadFile filepath.Join string dir , "doc.md" if err = nil { return err } res, err := summarize.Run ctx, SummarizeInput{Path: p.Path, Text: string text } if err = nil { return err } fmt.Fprintf ctx.Stdout , "%s\n", res.Summary return os.WriteFile filepath.Join string dir , "summary.txt" , byte res.Summary+"\n" , 0o644 } Three details make this work smoothly: senro.Ctx embeds context.Context summarize.Run with no adapter. Genkit's tracing and cancellation behave exactly as they always did. ctx.Stdout is the step's real log stream os.Stdout instead reaches your terminal and no log file. ctx.Workspace name hands the function the same path a mount gives a command RegisterFunc registers the function once, from an init . The name is the function's identity: it is what the plan records and what feeds the step's cache key, so renaming it invalidates the cache exactly as renaming a command would. Parameters must be JSON-serializable, and decoding is strict, so a renamed field fails loudly instead of running with a zero value. Because the pipeline is a Go program, fanning out over a corpus is a for loop: func main { ctx := context.Background g := genkit.Init ctx, genkit.WithPlugins &googlegenai.GoogleAI{} , genkit.WithDefaultModel "googleai/gemini-2.5-flash" summarize = defineFlows g docs, err := filepath.Glob "docs/ .md" if err = nil { log.Fatal err } sort.Strings docs p := senro.New "ai" w := p.Workflow "summarize" var names, ids string var wss senro.WorkspaceRef corpus := map string string{} for , doc := range docs { text, err := os.ReadFile doc if err = nil { log.Fatal err } name := strings.TrimSuffix filepath.Base doc , ".md" names = append names, name ids = append ids, "summarize/"+name corpus name = string text // One workspace per document, so one document changing does not // invalidate the other summaries' cache keys. wss = append wss, senro.Workspace "doc-"+name, senro.Scope senro.ScopeRun } seed := w.Step "seed", senro.Func "ai/seed", SeedParams{Docs: corpus} for i, ws := range wss { seed.Mount ws.At "/doc-"+names i , senro.RW } for i, name := range names { // One AI call, one node: its own log, state, retry and cache entry. w.Step ids i , senro.Func "ai/summarize", SummarizeParams{ Path: docs i , Name: name, } . Needs "seed" . Mount wss i .At "/doc", senro.RW . Retry 3, retry.OnInfra . Pure . Inputs artifact.File "doc.md" . Outputs artifact.File "summary.txt" } digestStep := w.Step "digest", senro.Func "ai/digest", DigestParams{Names: names} . Needs ids... for i, ws := range wss { digestStep.Mount ws.At "/doc-"+names i , senro.RO } if err := senro.Run ctx, p ; err = nil { log.Print err os.Exit 1 } } seed writes each document into a workspace of its own. Then one summarize step per document, running in parallel. Then digest , which waits for all of them and merges the results. That loop produces this graph, rendered straight from the run's own plan.json : pipeline: ai │ ├─ wave 1 seed func │ ├─ wave 2 3 parallel summarize/alpha func · retry 3/infra │ summarize/beta func · retry 3/infra │ summarize/gamma func · retry 3/infra │ └─ wave 3 digest needs summarize/alpha, summarize/beta, summarize/gamma · func Three documents give three parallel nodes. Thirty would give thirty, from the same for loop, with no change to the code. Run the pipeline binary directly, the way any Go program runs: go run . Or hand the package to the senro CLI, which builds it, execs it and attaches automatically. You get a terminal UI on a TTY and plain streaming lines anywhere else: senro run . php seed - succeeded summarize/beta - succeeded summarize/alpha - succeeded summarize/gamma - succeeded digest - succeeded - alpha: The billing reconciliation job now uses a queue worker for invoice settlement within two minutes, replacing the nightly cron. - beta: The search index now rebuilds incrementally in under a second without blocking writes, replacing the old full rebuild and requiring operators to delete the rebuild-search cron job. - gamma: Service gamma's public API now features per-tenant rate limiting set at 100 requests per second, which may cause HTTP 429 errors for existing bursting integrations. Three real model calls in parallel, each with its own log file, and a merge step that ran once they were all done. While that is running, a second terminal can watch it, or steer it: senro attach terminal UI senro ui browser view, prints a one-time link senro attach --run