Skip to content

Blocking vs Async Execution

Understanding the execution model is crucial for building reliable components. TinySystems uses a blocking-by-default model with explicit async patterns when needed.

Default Behavior: Blocking

When you call output(), execution blocks until all downstream processing completes:

go
func (c *Component) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    fmt.Println("1. Starting")

    // This blocks until downstream chain completes
    err := output(ctx, "output", msg)

    fmt.Println("2. Downstream complete")
    return err
}

Execution Timeline

Time ------------------------------------------------------------------>

Node A          Node B          Node C          Node D
  |               |               |               |
  | Handle()      |               |               |
  |   |           |               |               |
  |   v           |               |               |
  | output()------|--> Handle()   |               |
  |   | BLOCKED   |     |         |               |
  |   |           |     v         |               |
  |   |           |   output()----|--> Handle()   |
  |   |           |     | BLOCKED |     |         |
  |   |           |     |         |     v         |
  |   |           |     |         |   output()----|--> Handle()
  |   |           |     |         |     | BLOCKED |     |
  |   |           |     |         |     |         |     v
  |   |           |     |         |     |         |   return nil
  |   |           |     |         |     <---------|-----+
  |   |           |     |         |   return nil  |
  |   |           |     <---------|-----+         |
  |   |           |   return nil  |               |
  |   <-----------|-----+         |               |
  | return nil    |               |               |
  v               |               |               |

Why Blocking by Default?

1. Backpressure

Slow consumers naturally slow down producers:

go
// Producer (fast)
func (c *Producer) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    for _, item := range items {
        // Blocks if downstream is slow
        output(ctx, "output", item)
    }
    return nil
}

// Consumer (slow)
func (c *Consumer) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    time.Sleep(100 * time.Millisecond)  // Slow processing
    return nil
}

Without blocking, the producer would overwhelm the system with messages.

2. Rate Limiting

The ticker component demonstrates this:

go
// Ticker emits every N milliseconds
func (t *Ticker) emit(ctx context.Context, handler module.Handler) error {
    timer := time.NewTimer(t.settings.Delay)

    for {
        select {
        case <-timer.C:
            // Blocks until downstream completes
            _ = handler(ctx, OutPort, t.settings.Context)

            // Only THEN reset timer
            timer.Reset(t.settings.Delay)
        case <-ctx.Done():
            return ctx.Err()
        }
    }
}

If processing takes longer than the interval, the next tick waits.

3. Error Propagation

Errors flow back to the source:

go
func (c *Component) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    err := output(ctx, "output", msg)
    if err != nil {
        // We know downstream failed
        log.Error("downstream failed", "error", err)
        return err
    }
    return nil
}

4. Transaction-like Semantics

Know when a unit of work is complete:

go
func (c *BatchProcessor) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    batch := msg.(Batch)

    for _, item := range batch.Items {
        err := output(ctx, "output", item)
        if err != nil {
            return err  // Stop on first error
        }
    }

    // All items processed successfully
    return output(ctx, "complete", BatchComplete{Count: len(batch.Items)})
}

Async Execution

For fire-and-forget scenarios, use goroutines explicitly:

go
func (c *AsyncComponent) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    go func() {
        // Preserve trace context
        asyncCtx := trace.ContextWithSpanContext(
            context.Background(),
            trace.SpanContextFromContext(ctx),
        )

        // This runs independently
        _ = output(asyncCtx, "output", msg)
    }()

    return nil  // Returns immediately
}

When to Use Async

Use CaseBlockingAsync
Sequential processing
Error handling needed
Rate limiting
Fire-and-forget
Parallel fan-out
Long-running background tasks

Async Risks

Goroutine leaks: If downstream blocks forever, goroutines accumulate:

go
// DANGER: Potential goroutine leak
func (c *LeakyComponent) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    go func() {
        // If downstream never responds, this goroutine lives forever
        output(context.Background(), "output", msg)
    }()
    return nil
}

Mitigation: Use context with timeout:

go
func (c *SafeAsyncComponent) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    go func() {
        // Timeout after 30 seconds
        asyncCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
        defer cancel()

        // Preserve tracing
        asyncCtx = trace.ContextWithSpanContext(asyncCtx, trace.SpanContextFromContext(ctx))

        _ = output(asyncCtx, "output", msg)
    }()
    return nil
}

Parallel Processing

For parallel fan-out, combine async with error collection:

go
func (c *ParallelProcessor) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    items := msg.([]Item)

    var wg sync.WaitGroup
    errors := make(chan error, len(items))

    for _, item := range items {
        wg.Add(1)
        go func(item Item) {
            defer wg.Done()
            if err := output(ctx, "output", item); err != nil {
                errors <- err
            }
        }(item)
    }

    // Wait for all to complete
    wg.Wait()
    close(errors)

    // Check for errors
    for err := range errors {
        return err  // Return first error
    }

    return nil
}

The Split Component Pattern

The split component demonstrates blocking iteration:

go
func (c *Split) Handle(ctx context.Context, handler module.Handler, port string, msg any) any {
    input := msg.(InMessage)

    for _, item := range input.Array {
        // Each item blocks until processed
        if err := handler(ctx, OutPort, OutMessage{
            Context: input.Context,
            Item:    item,
        }); err != nil {
            return err  // Stop on error
        }
    }

    return nil
}

Behavior:

  • Items processed sequentially
  • Error in any item stops iteration
  • Downstream controls pace

Context Handling in Async

Always preserve trace context for observability:

go
import "go.opentelemetry.io/otel/trace"

func (c *AsyncComponent) Handle(ctx context.Context, output module.Handler, port string, msg any) any {
    // Extract span context from incoming request
    spanCtx := trace.SpanContextFromContext(ctx)

    go func() {
        // Create new context with span context (but no deadline/cancellation)
        asyncCtx := trace.ContextWithSpanContext(context.Background(), spanCtx)

        // Now traces will be connected
        output(asyncCtx, "output", msg)
    }()

    return nil
}

Comparison Summary

AspectBlockingAsync
Return timingAfter downstream completesImmediately
Error handlingErrors propagate backErrors lost (unless handled)
BackpressureNaturalNone
Goroutine countBoundedCan grow unbounded
TracingAutomaticManual context propagation
Use caseMost componentsFire-and-forget, parallel

Best Practices

  1. Default to blocking - Use async only when necessary
  2. Handle async errors - Log or handle errors in goroutines
  3. Preserve trace context - Keep observability working
  4. Use timeouts - Prevent goroutine leaks
  5. Consider backpressure - Async can overwhelm downstream

Next Steps

Durable execution and runs

Everything above is about how one component hands a message to the next inside a single request. There's a second layer: how the whole agent executes across pod restarts.

Trigger-driven agents run durable by default. When a trigger fires (a signal, a cron, a webhook that doesn't hold a caller), each hop is written to a persistent work queue before it moves on, and every completed step is recorded to a run ledger in JetStream. If a pod dies mid-run, another pod picks up from the last recorded step. Redelivery is deduplicated by step key, so a step never runs twice.

The result is a run: a durable record you can list, inspect step by step, and retry. Traces show timing; runs show state and recovery.

Request/response agents run classic. An http_server holds a live connection and blocks on its response port, so its chain uses core request/reply delivery — which also load-balances across the module's pods, giving you horizontal scaling for stateless endpoints. A durable, fire-and-forget hop would never return the response the caller is waiting on, so it stays out of that path.

Components declare their transport

You never choose durable versus classic. There is no flow-level setting and no annotation. The mode is derived from the components in the agent.

A component that blocks on a response declares it by implementing the SyncRPC capability:

go
// http_server holds a live HTTP connection until the response arrives,
// so its request→response path must run classic.
func (h *Component) SyncRPC() module.SyncRPCInfo {
    return module.SyncRPCInfo{}
}

At save time the platform reads that declaration from the published component metadata and classifies the graph: every node in a connected subgraph that contains a SyncRPC component runs classic; disconnected, trigger-driven subgraphs in the same agent run durable, independently.

Component authors state one fact about how their component communicates. Nobody — not a user, not a flow, not the agent that built it — picks a mode. It follows from what the agent is made of.

Build Kubernetes workflows with a prompt