Shahzad Bhatti Welcome to my ramblings and rants!

October 9, 2026

Building Production Data Pipelines with SEDA and Actors Model

Filed under: Computing — admin @ 8:44 pm

1. Introduction

The data pipelines come in different shapes such as streaming telemetry data, machine learning, messaging, or other real-time data feed. Common solutions to these problems include messaging middlewares like Kafka, Apache Flink/Spark, Airflow DAGs, Temporal, CDC via Debezium, etc. These tools may differ, but the design underneath is the same: ingest data, pipe it through a chain of stages that may hold state, and then send it to the destination after some filtering/decoration/massaging of the data. This post argues that this kind of stateful pipe and filter based architecture can be build with a simple primitive: durable actor connected by channels, which is what I built PlexSpaces around. This architecture is also similar to the SEDA (Staged Event-Driven Architecture) that builds applications as a network of stages joined by event queues. Each stage has its own workers with a queue that give you backpressure, isolation, and a place to observe load, e.g, :

An actor is a SEDA stage in miniature:

By default actor mailbox is in-memory state and it doesn’t connect a stage to Kafka, SQS, etc. PlexSpaces addresses these gaps with:

  • Durable actors. The durability facet journals state changes and checkpoints.
  • Channels. Queue-based coupling between stages.

In addition, PlexSpaces supports APIs for TupleSpace for coordination state, scatter/gather, and shard groups, and worker pool. It supports native Rust as well a WASM runtime with SDKs for Python, Go, TypeScript, and Rust.


2. The Primitive Stack

PlexSpaces provides following set of coordination primitives:

PrimitiveWhat It Does
Shard GroupPartition data across N actors; scatter-gather with aggregation
Worker PoolStateless actor pool with load balancing
Process GroupErlang pg2-style dynamic membership; broadcast
TupleSpacePattern-matched shared memory; Linda-model coordination
ChannelsQueue-based stage coupling; 6 backends (Kafka, Redis, SQS, PG, …)
Workflow ActorMulti-step durable orchestration; pause/resume/cancel
Distributed LockLease-based mutual exclusion across actors
Ring AllReduceCollective gradient reduction for distributed training
Parameter ServerCentralized gradient accumulation with pull/push
BroadcastSend data to all actors in a process group
Collective ReduceSum/min/max across all actors; return to coordinator
Scatter/GatherFan-out to N workers, fan-in aggregated results

2.1 Facets: Zero-Code Cross-Cutting Capabilities

PlexSpaces supports dynamic composition using Facets abstraction that allows plugging capabilities via configuration at deployment time:

FacetWhat It Adds
securityJWT validation, RBAC enforcement
loggingStructured log correlation with trace IDs
metricsPrometheus histograms per handler, per actor
memoizeDeterministic response caching
virtual_actorOrleans-style activate-on-demand lifecycle
schema_validationMessage schema enforcement
durabilityJournaling + checkpoint + crash replay
execution_traceFull execution trace capture
timer / reminderDurable delayed messages (survive crashes)
event_sourcing / cachingFull audit trail / response caching
kv / locks / registry / process_groupsState, coordination, discovery
http_clientOutbound HTTP calls
event_emitterPub/sub event broadcasting

Behaviors:

BehaviorAnnotationPattern
GenServer@actor / #[gen_server_actor]Request-reply (ask/tell)
GenEvent@event_actorFire-and-forget event handling
FSM@fsm_actorState machine with transitions
Workflow@workflow_actorMulti-step durable orchestration

2.2 Polyglot SDK

PlexSpaces provides polyglot SDK for Python, TypeScript, GO and Rust:

Python:

from plexspaces import actor, handler, init_handler, host, state

@actor
class Counter:
    count: int = state(default=0)

    @init_handler
    def on_init(self, config: dict):
        self.count = int(config.get("args", {}).get("initial", 0))

    @handler("increment")
    def increment(self, amount: int = 1) -> dict:
        self.count += amount
        return {"count": self.count}

GO:

type CounterActor struct {
    plexspaces.BaseActor
    Count int `json:"count"`
}

func (a *CounterActor) Init(config map[string]interface{}) error {
    if args, ok := config["args"].(map[string]interface{}); ok {
        if v, ok := args["initial"].(float64); ok { a.Count = int(v) }
    }
    return nil
}

func (a *CounterActor) Handle(op string, payload map[string]interface{}) (interface{}, error) {
    switch op {
    case "increment":
        amount := 1
        if v, ok := payload["amount"].(float64); ok { amount = int(v) }
        a.Count += amount
        return map[string]interface{}{"count": a.Count}, nil
    }
    return nil, fmt.Errorf("unknown op: %s", op)
}

Typescript:

import { PlexSpacesActor, host } from "@plexspaces/sdk";

interface CounterState extends Record<string, unknown> {
  count: number;
}

class CounterActor extends PlexSpacesActor<CounterState> {
  onInit(config: Record<string, unknown>): CounterState {
    const args = (config.args as Record<string, unknown>) || {};
    return { count: Number(args.initial || 0) };
  }

  onIncrement(payload: Record<string, unknown>): Record<string, unknown> {
    const amount = Number(payload.amount || 1);
    this.state.count += amount;
    return { count: this.state.count };
  }
}

2.3 How Deployment Works

Every actor gets deployed via a TOML config file and a WASM binary. The config declares the supervision tree, actor roles, facets, and resource requirements:

[application]
name = "log-pipeline"
version = "1.0.0"

[supervisor]
strategy = "one_for_one"
max_restarts = 5
max_seconds = 60

[[supervisor.children]]
id = "LeaderActor"
actor_type = "LeaderActor"
behavior_kind = "GenServer"
role = "leader"
facets = ["virtual_actor", "metrics", "durability"]

[[supervisor.children]]
id = "PipelineWorker"
actor_type = "PipelineWorkerActor"
behavior_kind = "GenServer"
count = 8
role = "worker"
facets = ["virtual_actor", "metrics"]

The supervisor strategy (one_for_one, one_for_all, rest_for_one) comes directly from Erlang/OTP. If a pipeline worker crashes, only that worker restarts. If the leader crashes, the supervisor can optionally restart all children.

Deploy with a single API call:

zip app.zip pipeline_actor.wasm app-config.toml
curl -X POST http://node:8091/api/v1/applications/deploy \
  -F "application_id=log-pipeline" \
  -F "app_file=@app.zip"

3. Observability Pipelines

3.1 Log Ingestion & Routing Pipeline

In the log-pipeline example, the leader uses scatter/gather APIs to fan out event batches in parallel. Each PipelineWorkerActor receives a shard of events and runs them through a configurable function chain such as JSON parse, regex extract, PII masking, field enrichment, drop rules, etc. while keeping timing-metrics in state. The pipelines are partitioned by source hash across a shard group for horizontal scaling. After each scatterGather returns, the leader aggregates timing and throughput across all workers.

Python actor (excerpt):

@actor
class LeaderActor:
    compute_ms: float = state(default=0.0)
    coord_ms: float = state(default=0.0)

    @handler("run")
    def run(self, event_count=10000, worker_count=8, pipeline_fns=None, **kw):
        pipeline_fns = pipeline_fns or ["json_parse", "regex_extract", "mask_pii", "enrich"]
        group_id = host.create_shard_group(self.app_id, "PipelineWorkerActor", worker_count)
        events = generate_log_events(event_count)
        batches = partition_events(events, worker_count)

        t0 = host.now_ms()
        results = host.scatter_gather(group_id, "process_batch",
            [{"events": b, "pipeline_fns": pipeline_fns} for b in batches], timeout_ms=60000)
        coord_ms = host.now_ms() - t0

        total_compute = sum(r.get("compute_ms", 0) for r in results)
        return {"status": "ok", "events_per_sec": event_count * 1000 / coord_ms, ...}

Pipeline functions implemented across all languages:

  • fn_json_parse: Parse raw JSON body
  • fn_regex_extract: Extract fields via regex (timestamps, IPs)
  • fn_mask_pii: Mask email, SSN, credit card patterns
  • fn_enrich: Add geo, datacenter, environment metadata
  • fn_rename_fields: Normalize field names across sources
  • fn_drop: Filter events by rules (debug level, internal sources)

Benchmark: Go WASM, 100K events, depth=5, 2 nodes (tested on my MacPro laptop)

MetricValue
Throughput177,304 events/sec
Wall time564 ms
Compute time1,585 ms (73%)
Coordination time564 ms (26%)
Granularity ratio2.8×
Events processed86,384 / 100,000 (13,616 dropped by fn_drop)

Benchmark: Python WASM

WorkersTotal EventsEvents/secEfficiency
2200~10,000100%
4400~15,000150%
8800~22,000220%
161,600~30,000294%

Benchmark: TypeScript WASM

WorkersEvents/secWall msComp msCoord msComp%GranSpeedupEff%
22,25044527616962%1.6×1.00×100%
41,12089355833662%1.7×0.50×25%
86771,4781,11636276%3.1×0.30×8%
163502,8692,42644385%5.5×0.16×2%

Go actor

func (a *WorkerActor) handleProcessBatch(payload map[string]interface{}) (interface{}, error) {
    events := payload["events"].([]interface{})
    fns := payload["pipeline_fns"].([]interface{})

    t0 := plexspaces.NowMs()
    var processed []map[string]interface{}
    for _, ev := range events {
        event := ev.(map[string]interface{})
        for _, fn := range fns {
            event = applyPipelineFunction(fn.(string), event)
            if event == nil { break }
        }
        if event != nil { processed = append(processed, event) }
    }
    computeMs := plexspaces.NowMs() - t0

    return map[string]interface{}{
        "processed_count": len(processed),
        "compute_ms": computeMs,
    }, nil
}

3.2 Metrics Aggregation Pipeline

In the metrics_aggregation example, the leader uses scatter/gather APIs to provision N aggregation shards and then route metric batches. The framework hashes each metric name to a shard and each AggregationShard maintains a tumbling window in its actor state. After scatterGather returns, the leader collects shardResponses, merges rollup results, and runs z-score anomaly detection. Anomalies trigger alert events delivered via actor APIs and the state stays in the actor.

Rust embedded actor:

#[gen_server_actor]
struct AggregatorShard {
    shard_id: usize,
    windows: HashMap<String, Vec<f64>>,
    total_metrics: u64,
}

#[plexspaces_handlers(gen_server)]
impl AggregatorShard {
    #[handler("ingest_batch", cast)]
    async fn handle_ingest_batch(&mut self, _ctx: &ActorContext, msg: &Message) -> Result<(), BehaviorError> {
        let batch: Vec<Metric> = serde_json::from_slice(&msg.payload)?;
        for m in &batch {
            self.windows.entry(m.name.clone()).or_default().push(m.value);
        }
        self.total_metrics += batch.len() as u64;
        Ok(())
    }

    #[handler("get_aggregates")]
    async fn handle_get_aggregates(&self, _ctx: &ActorContext, _msg: &Message) -> Result<Value, BehaviorError> {
        // Compute count/sum/avg/min/max per metric name, return as JSON
    }
}

Benchmark: Go WASM (100K metrics, 8 workers)

MetricValue
Throughput480,769 metrics/sec
Compute time208 ms (69%)
Coordination time91 ms (30%)
Granularity ratio2.30×

Benchmark: Rust Embedded

MetricValue
Throughput86,372 metrics/sec
Compute time361 ms (62%)
Coordination time220 ms (38%)
Granularity ratio1.64×

Benchmark: Python WASM

WorkersTotalAgg/sEfficiency
22K444K100%
44K451K101%
88K901K202%
1616K1.44M323%

3.3 Distributed Tracing Pipeline

In the tracing_pipeline example, the leader creates a shard group and then uses scatter/gather to send only the spans whose trace_id hashes to that shard. All spans for a given trace land on one worker and each WorkerActor holds partial traces in its actor state and then runs tail-based sampling. Workers accumulate compute_ms in state and the leader sums shardResponses[*].compute_ms to compute aggregate compute time.

Go actor (excerpt):

type Span struct {
    TraceID      string            `json:"trace_id"`
    SpanID       string            `json:"span_id"`
    ParentSpanID string            `json:"parent_span_id,omitempty"`
    ServiceName  string            `json:"service_name"`
    Operation    string            `json:"operation"`
    StatusCode   int               `json:"status_code"`
    DurationMs   float64           `json:"duration_ms"`
    StartTimeMs  int64             `json:"start_time_ms"`
    Tags         map[string]string `json:"tags,omitempty"`
}

Benchmark: Strong Scaling

WorkersSpans/secWall msComp msCoord msGranSpeedupEff%
2242,7442952951611.8x1.00x100.0%
4193,3713733732071.8x0.79x39.5%
8218,4243303301262.6x0.89x22.3%
16221,2153213211192.7x0.92x11.5%

Service graph output example:

{
  "nodes": [
    {"name": "api-gateway", "span_count": 1000, "error_count": 48, "avg_latency_ms": 145.2},
    {"name": "user-service", "span_count": 800, "error_count": 39, "avg_latency_ms": 32.1},
    {"name": "order-service", "span_count": 600, "error_count": 31, "avg_latency_ms": 45.8},
    {"name": "payment-service", "span_count": 400, "error_count": 22, "avg_latency_ms": 89.3}
  ],
  "edges": [
    {"source": "api-gateway", "target": "user-service", "call_count": 800, "p99_latency_ms": 195.0},
    {"source": "api-gateway", "target": "order-service", "call_count": 600, "p99_latency_ms": 210.5},
    {"source": "order-service", "target": "payment-service", "call_count": 400, "p99_latency_ms": 380.2}
  ]
}

Benchmark: Python WASM

WorkersTotal TracesSpans/sGranEff%
2500105,7521.4×100%
41,000104,7791.5×99%
82,000213,5202.6×202%
164,000432,4174.4×409%

Benchmark: TypeScript WASM

ShardsTotal SpansSpan/sWall msComp msCoord msGranEff%
21,0005,988167109741.5×100.0%
42,0005,7643472271501.5×96.3%
84,0006,9205784551752.6×115.6%
168,0007,1051,1269872454.0×118.7%

4. ML Pipelines

4.1 Batch Inference

The parallel_ai_inference example demonstrates shard-group-based batch inference with model caching per worker. A benchmark actor then tries three ways of distributing that work and measures speed. The three ways are: broadcast like MPI, scatter/gather and elastic pool.

The batch_image_classification example shows ViT-style image classification across shards. A leader actor splits a pile of images across N worker actors, runs the work in rounds, and merges the results.

Benchmark: Python WASM Parallel AI Inference

ShardsReq/secComp%GranSpeedupEfficiency
151262%1.7×1.00×100%
4~800——~1.56×~39%
81,45391%10.5×2.84×70%
32~2,100——~4.1×16%

4.2 Parameter Server

The ParameterServerActor (leader) in parameter_server example creates shard group to provision worker shards, then broadcasts the current weight matrix via scatter/gather API. Each WorkerActor stores its local data shard in actor state, computes a mini-batch gradient locally, and returns a gradient summary. The leader collects shardResponses, averages the gradients across all workers.

Benchmark: TypeScript WASM

MetricValue
param_count6,464 (input_dim=100 × hidden_dim=64 + 64)
Gradient opsworkers × 20 rounds
Compute time~2 ms (gradient averaging + weight update)
Coordination time~4,000 ms (scatter-gather round trips)
Granularity ratio~0.00× (coordination-dominated on single node)
Errors0

4.3 Federated Learning

The federated_learning example shows a federated learning simulation. An aggregator coordinates many client actors, each holding its own private dataset. Every round, the aggregator sends the current model out. Each client trains locally and returns gradients. The aggregator then averages them weighted by sample count.

Benchmark: Federated Learning

ClientsWall msCompute msCoord msGranSpeedupEfficiencyAccuracy
23402091311.6×1.00×100%93.7%
46824292531.7×0.50×25%93.6%
81,1388582793.1×0.30×7.5%93.6%
162,1061,7713345.3×0.16×2.0%93.6%

4.4 LLM Serving with Request Batching

The llm_serving simulates LLM serving pipeline similar to vLLM. A router generates a mix of requests and sends each one to a model tier: small, medium, or large. Cheap tasks go to the small model and heavy ones to the large mode. The router groups requests into batches per tier and scatters each batch to model workers, which simulate inference. The router then totals up throughput, tokens per second, cost, etc.

Benchmark: Python WASM

ReplicasReq/s (aggregate)Tok/s (aggregate)Wall msSpeedup
2~37,000~2.1M541.00×
4~39,000~2.2M511.05×
8~43,000~2.5M461.17×

Benchmark: TypeScript WASM

ShardsTotal ReqReq/sTok/sWall msComp msCoord msGranEff%
21,55286,22223,325,778182552.6×100.0%
43,09696,75027,481,5003249200.6×112.2%
64,644132,68638,290,2863573230.5×153.9%
86,144170,66744,575,1113697240.5×197.9%

4.5 Feature Store

The feature_store simulates a feature store like Feast or Tecton. A leader generates features and ingests the features into worker shards, which keep a versioned key-value store. The leader then serves batched lookups, where each worker checks its cache first and falls back to its store. At the end the leader collects stats and reports lookups per second, cache hit rate, etc.

TypeScript actor (excerpt):

class WorkerActor extends PlexSpacesActor<WorkerState> {
  onIngestBatch(payload: Record<string, unknown>): Record<string, unknown> {
    const records = payload.records as FeatureRecord[];
    for (const rec of records) {
      const key = `${rec.entity_id}:${rec.feature_name}`;
      const versions = this.state.features[key] || { versions: [] };
      versions.versions.push({ value: rec.value, version: rec.version, timestamp: rec.timestamp });
      if (versions.versions.length > 3) versions.versions.shift(); // Keep last 3
      this.state.features[key] = versions;
    }
    return { status: "ok", ingested: records.length };
  }

  onLookup(payload: Record<string, unknown>): Record<string, unknown> {
    const entityId = payload.entity_id as string;
    const cacheKey = entityId;
    // Check LRU cache first
    if (this.state.cache[cacheKey]) {
      this.state.cache_hits++;
      return { ...this.state.cache[cacheKey], cache_hit: true };
    }
    this.state.cache_misses++;
    // Assemble feature vector from versioned store
    // ...
  }
}

Benchmark: Python WASM

WorkersLookups/secSpeedupEfficiency
2~10,0001.00×100%
4~13,0001.30×65%
8~17,0001.70×43%
16~21,0001.95×24%

Benchmark: TypeScript WASM

ShardsTotal LookupsLkup/sWall msComp msCoord msGranEff%
21,6003,6704361123800.1×100.0%
43,2003,7788472177930.1×102.9%
86,4007,2988774288230.1×198.9%
1612,8006,5981,9408861,8850.0×179.8%

5. Data Processing & ETL

5.1 Streaming Window Aggregation

The streaming_pipeline example simulates streaming pipeline for event processing like a log or telemetry pipeline. A leader runs a series of batch rounds. Each round it tells all worker shards to process a batch of events. Workers simulate filtering, enriching, and transforming events. The worker then write a stage summary to a shared tuple space and return counts, bytes processed, etc. The leader reads the stage summaries back and reports events per second.

Benchmark: TypeScript WASM Streaming Pipeline

ShardsTotal EventsEvt/sWall msComp msCoord msGranEff%
29,6001,600,0006150.2×100.0%
419,2002,742,8577260.2×171.4%
838,4004,800,0008080.0×300.0%
1676,8005,485,714148130.1×342.9%

5.2 Data Lake RAG Ingestion

The data_lake_rag and agentic_rag_pipeline exampless demonstrate document ingestion, chunking, embedding, vector store and retrieval-augmented generation. Workflow actors manage the multi-step pipeline with checkpointing, so a crash during embedding resumes from the last chunk.

Benchmark: Data Lake RAG

MetricValue
Workload100 queries × 500 chunks × 2048 bytes
Wall time1,429 ms
Compute time3,984 ms (75.6%)
Coordination time1,285 ms (24.4%)
Granularity3.10×

5.3 CDC Pipeline

The cdc_pipeline example simulates PostgreSQL-style WAL events (INSERT, UPDATE, DELETE) on four tables: users, orders, products, inventory. Workers maintain LSN (Log Sequence Number) tracking per table, apply schema transformations and fan out to three sink types (search index, analytics warehouse, cache invalidation).

Benchmark: CDC Pipeline (10K WAL events, 4 tables)

WorkersEvents/secSpeedupEfficiency
2~1,4001.00×100%
4~1,9001.36×68%
8~2,8002.00×50%
16~3,1002.20×28%

Python CDC actor (excerpt):

@actor
class WorkerActor:
    lsn_watermarks: dict = state(default_factory=dict)
    seen_event_ids: set = state(default_factory=set)

    @handler("process_events")
    def process_events(self, events=None, **kw):
        results = {"search_index": [], "analytics": [], "cache_invalidation": []}
        for event in events:
            if event["event_id"] in self.seen_event_ids:
                continue  # Dedup
            self.seen_event_ids.add(event["event_id"])
            self.lsn_watermarks[event["table"]] = max(
                self.lsn_watermarks.get(event["table"], 0), event["lsn"])
            transformed = self.transform(event)
            for sink in ["search_index", "analytics", "cache_invalidation"]:
                results[sink].append(self.format_for_sink(transformed, sink))
        return {"status": "ok", "processed": len(events), "sinks": results}

5.4 Web Crawler

The web_crawl example simulates parallel web crawler. It includes aan orchestrator that runs a breadth-first crawl from seed URLs. For each URL it checks a fetcher using an elastic pool, fetches the page and checks the fetcher back in. After the crawl, the results are split between two analyzer shards that merge word counts into a global top-10. Also, the orchestrator asks the process group of fetchers how many pages each one handled.

Benchmark: TypeScript WASM Web Crawl

WorkersTotPagesPg/sWall msComp msCoord msSpeedupEff%
24001,95120520051.00×50.0%
48001,9234162022140.99×24.6%
81,6003,8834122012111.99×24.9%
163,2007,6924162012153.94×24.6%

6. Learnings

This post covered several use cases and examples for building data pipelines in different languages using PlexSpaces framework. Here are key lessons when building these pipelines:

  1. Stateful Actors. Actors are natural fit for building data pipelines. Facets allow adding cross-cutting concerns dynamically like metrics, durability, etc.
  2. Granularity ratio matters. There is a communication overhead when communicating with lots of actors/processes. The granularity ratio is compute time divided by coordination/communication time. So pay attention to this and test proper batch size for the parallel work
  3. Batch messages. Instead of sending one event per tell/cast, batch 100-1000 events per message.
  4. Broadcast scatter-gather. The scatter/gather APIs have been used in high performance computing for decades and are proven to scale for building highly parallel and concurrent applications.
  5. Reduce shard count when coordination dominates. More shards means more scatter-gather rounds. If your per-shard work is small, you are paying coordination tax for parallelism you do not need.
  6. Use cast (fire-and-forget) instead of call (request-reply). PlexSpaces supports both cast/tell and call/ask. Cast does not wait for a response, so the sender can continue immediately. Call blocks until the handler returns so use cast where possible such as data ingestion.
  7. TupleSpace for coordination. TupleSpace is based on Linda memory model, which is designed for coordination, synchronization.
  8. Durable actors. Stateless functions retry the whole batch on failure. Actors with the durability facet journal state before applying it. This means that a crash replays from the last journal entry instead of start.
  9. WASM Gotchas. WASM has high serialization overhead (~3.5ms for store reinstantiation). PlexSpaces supports native Rust SDK for microsecond-latency workloads.

Amdahl’s Law in practice. The parallel fraction of your workload determines the theoretical maximum speedup. From our benchmarks:

WorkloadParallel FractionMax Speedup (Amdahl)Measured at 16 workers
Log pipeline (batch=500)94%16.7x8.8x
Metrics aggregation86%7.1x5.2x
Batch inference97%33.3x12.0x
Ring AllReduce92%12.5x8.4x
Federated learning84%6.3x0.16×

GitHub: https://github.com/bhatti/PlexSpaces

Clone the repo. Run ./test.sh in any example directory.

Related reading


    September 28, 2026

    Governing the AI Factory: How to Ship Fast Without Derailing the Release Train

    Filed under: Computing — admin @ 4:51 pm

    Here’s how to define where AI fits, reduce cognitive and intent debt, and keep your delivery pipeline stable.


    1. The AI Genie

    AI coding agents can produce code at an incredible speed but without the quality, you will find that it quickly derails the release train. The existing CI and code-review processes were designed for human-paced changes but this doesn’t keep up when you are pushing code at 10x and you end up with the release problem.

    In my experience, AI code velocity increases change failure rate, build times, merge conflicts, flaky tests and release incidents. Here is how some of the companies are trying to solve merge queue related problems:

    SourceFinding
    LinearCut PR wait time via faster runners, tsc -> tsgo, sparse checkout, additional test shards
    Atlassian70+ large repos on merge queues, PR-level CI answers “does this work alone,” not “does this still work with everything”
    NirvanaStateless speculative merge queue.
    MergifySpeculative testing, batching, and scope-aware parallel lanes
    AutonomaParallel AI subagents multiply conflict surface area unless generation is batched and merges serialized

    Here is additional data I found on how AI is affecting the release process:

    SourceFinding
    DORA 2025PRs merged per developer +98%. Incidents per PR +243%. Bugs per developer +54%
    Cortex 2026PRs per author +20% YoY. Incidents per PR +23.5%. Change failure rate ~30%
    GitHub Octoverse 202492% of developers use AI coding tools. GitHub Actions CI/CD minutes up 169%
    GitClear (211M lines)Code churn doubled. Refactored code fell 24.1% ? 9.5%. Duplicated blocks rose 8x
    Uplevel (800 devs, 3 months)+41% bugs after Copilot adoption. No improvement in cycle time or throughput
    Opsera (250k+ devs, 60+ orgs)AI PRs sit 4.6x longer in review queues; duplication 10.5% ? 13.5%

    Joe Magerramov’s post shows how to build a small Monte Carlo simulator for modeling merge queues. Joe showed how CI/CD becomes a traffic jam, not just a queue. CI/CD batches are cumulative so one defect forces a revert and next batch is affected. I built my own simulator based on Joe’s model (see https://github.com/bhatti/simulators) with additional support for the batched PRs release. It models defect rate with the pipeline duration at 100 commits per day:

    The key lesson is that you need to either lower the defect rate per commit or shorten pipeline duration. This is not an easy task, I have encountered a large pipeline duration at many organizations due to large mono-repos, a large codebase with millions of LOC, and large test suite with a long vaidation cycle. The AI agents makes it worse with parallel changes that might conflict resulting in cognitive and intent debt for engineers because no one can track all changes. I explained some of these concepts in my earlier blog and showed how learning feedback loops can be used to build resilient software factories. In order to build end to end agentic SDLC process, you need to define what is AI responsible for and what are the roles for humans.

    QuadrantRole
    Human Real-timeDecisions that require judgment, context, accountability
    Human AsyncReview that needs thought but not immediacy
    AI Real-timeAssistance that augments human work in the moment
    AI BackgroundAutonomous work that runs without blocking humans

    2. The Architecture

    I built an orchestration engine Formicary to create data pipelines and CI/CD processes many years ago but I have been using it for driving AI driven workflows. It defines simple primitives to build DAG tasks with exit-code routing, artifact handoff, and fan-out. Here are a few approaches that I am using with the AI driven workflows:

    • Speculative merge queue (test PR N as if N-1 already merged)
    • Scope-aware parallel lanes
    • Incremental builds, test-impact analysis, sparse checkout
    • Deterministic gates a green build can’t waive
    • Independent AI review
    • Contract testing + canary
    • Formal verification of queue invariants (TLA+, Dafny)
    • API fuzz testing
    • Learning flywheel: merge -> extract learnings -> next run reads them -> audit proposes skill changes

    Here is how I use Formicary with a CLI toolkit and skillsThe system is split across three repos, each with a clear responsibility boundary:

    Here are the core design principles

    • Declarative DAG with a task block that uses on_exit_code for routing
    • Harness + sandbox + skills as separate layers
    • File-based state handoff
    • The learning flywheel closes the loop on review time and merge-time

    2.1 Scope router: ai-scope-router

    Runs right after create-pr in the existing pipeline and computes a scope key from touched paths. The backing script (scripts/mq/scope_router.py) computes blast radius from line counts and module count, and labels the PR.

    # ai-scope-router.yaml (excerpt — full file in docs/examples/)
    job_type: ai-scope-router
    max_concurrency: 20
    timeout: 600s
    
    tasks:
    - task_type: classify
      method: KUBERNETES
      script:
        - python -m scripts.mq.scope_router --pr-number {{.PRNumber}}
        - python -m scripts.mq.risk_score --pr-number {{.PRNumber}}
      on_completed: route
    
    - task_type: route
      script:
        - |
          python3 -c "
          import json, os
          scope = json.load(open('/workspace/scope.json'))
          risk = json.load(open('/workspace/risk_score.json'))
          decision = {
              'scope_key': scope['scope'],
              'risk_tier': risk['tier'],
              'lane': scope['scope'],
              'requires_approval': risk.get('requires_human_approval', False)
          }
          json.dump(decision, open('/workspace/route_decision.json', 'w'), indent=2)
          "
      on_completed: done

    Here’s what the Slack report looks like for a PR:

    # PR Review Report — PR #4091
    
    ## Review Findings
    
    ? No issues found
    
    ## Risk Score
    
    ? MEDIUM (score 25.5/100) — Standard review needed; test independently before merge
    
    | Dimension       | Score | Weight | Evidence                                          |
    |-----------------|-------|--------|---------------------------------------------------|
    | Size            | 7/10  | 1.5×   | 219 lines (+147/?72)                              |
    | File Count      | 2/10  | 1.0×   | 5 files changed                                   |
    | Blast Radius    | 5/10  | 2.0×   | blast=medium, scope=payments-service              |
    | Sensitive Paths | 0/10  | 2.5×   | no sensitive files detected                       |
    | Test Coverage   | 0/10  | 1.5×   | good coverage (test:source ?1:1)                  |
    | Historical      | 3/10  | 1.0×   | ?? no defect history available — neutral default   |
    
    
    
    ## Scope
    
    | Field          | Value            | Description                                                    |
    |----------------|------------------|----------------------------------------------------------------|
    | Scope          | payments-service | All changes owned by payments-service — can merge in dedicated lane |
    | Blast radius   | medium           | Moderate change (51–300 lines or 2 modules) — test independently |
    | Changed files  | 5                | Number of files modified in this PR                            |
    | Lines changed  | 219              | Total additions + deletions                                    |
    | Owners         | @payments-team   | CODEOWNERS entries responsible for review                      |

    2.2 Risk-gated review: ai-gate-review

    The RADAR-style funnel where AI reviews every PR, computes a risk score, and produces a report.

    # ai-gate-review.yaml — read-only pipeline
    # review ? gate-check ? done (no merge, no approve, no PR comments)
    - task_type: gate-check
      script:
        - |
          python3 -c "
          risk = json.load(open('/workspace/risk_score.json'))
          review = json.load(open('/workspace/review_result.json'))
          needs_approval = risk['score'] >= threshold or has_critical_findings
          # Writes gate_result.json — read-only, no PR changes
          print(f'Gate: {\"needs-approval\" if needs_approval else \"safe-to-merge\"}')"

    Low-risk PRs are flagged “safe-to-merge” in the report. High-risk ones are flagged “needs-approval” with the specific reason. The risk score itself is a RADAR-style weighted composite across six dimensions (scripts/mq/risk_score.py).

    2.3 The merge queue core

    As a Formicary DAG it uses real fan-out with fork_job_type where each scope lane runs as its own child job:

    # ai-merge-queue.yaml (excerpt)
    job_type: ai-merge-queue
    cron_trigger: "0 2 * * *"   # once daily at 2am; bump frequency when queue fills up
    max_concurrency: 1
    
    tasks:
    - task_type: collect
      environment:
        TARGET_BRANCH: "{{.TargetBranch}}"   # filter to PRs targeting this branch
      script:
        - python -m scripts.mq.collect_ready   # reads TARGET_BRANCH from env
    
    - task_type: group
      script:
        - python -m scripts.mq.group_by_scope   # produces hierarchical risk-tier lanes
    
    - task_type: analyze
      script:
        - python -m scripts.mq.analyze --skill ygs-merge-queue
    
    - task_type: report
      script:
        - mkdir -p /workspace/reports
        - python -m scripts.mq.report

    Invoke from Slack with a target branch:

    @bot mq --target stage          # short alias
    @bot mq --target prod --repo org/my-repo

    Each ai-mq-lane child job runs independently:

    2.4 Work type distribution

    The simulation discussed earlier has two knobs: defect rate and pipeline duration. The MQ report measures this from an actual PR data. For example, work type classification breaks every open PR into one of eight types: feature, bug, security, refactor, chore, test, docs, or unknown. It then measures defect rate, bug ratio, chore+refactor fraction.

    What the report shows:

    ### Work Type Distribution
    | Type      | Count | %     | Signal                      |
    |-----------|-------|-------|-----------------------------|
    | ? feature | 45    | 18.0% | new functionality           |
    | ? bug     | 22    | 8.8%  | defect indicator            |
    | ? security| 5     | 2.0%  | defect indicator (security) |
    | ?? refactor | 30    | 12.0% | tech debt reduction         |
    | ? chore   | 15    | 6.0%  | maintenance / KTLO          |
    | ? test    | 120   | 48.0% | quality investment          |
    | ? docs    | 3     | 1.2%  | documentation               |
    | ? unknown | 10    | 4.0%  | unclassified                |
    
    > ? Defect rate proxy: 10.8% (27 bug+security PRs / 250 total) — ~1-in-9.
    > At batch size ~10, est. batch success ? 31% (moderate defect rate).
    > Feature:Bug ratio = 1.7:1 — below 3:1, team spending significant effort on defect repair.

    3. Quality at the Source

    The defect-rate knob from earlier simulation is the hard to move but it is a high impact knob. Teams that have adopted agentic engineering at scale share a common pattern: an 8-phase SDLC that wraps AI capabilities with human judgment.

    3.1 Structured specs reduce defect rate at source

    Every ticket entering a sprint needs five things before an agent touches it:

    1. Why: one sentence on the customer problem
    2. Scope boundaries
    3. Aacceptance criteria in given/when/then form
    4. Definition of Done
    5. Environment matrix

    An AI skill generates the structured ACs; a human validates and edits.

    3.2 Design docs

    Not every change needs a design doc. The decision tree:

    SignalAction
    New architecture pattern, public API change, multi-team impact, >1 sprint, auth/securityWrite a design doc
    Bug fix, tests/config only, simple refactor in one file, dependency bumpNo design doc needed

    When required, the design doc must address: risks, rollback strategy, MVP scope, testing strategy, and observability.

    3.3 Tiered review

    Not all PRs deserve the same review depth. A risk-based triage:

    SignalsTier
    Small diff, tests or config only, simple refactor, pre-PR gate passed cleanlyTier 1 Light pass
    Auth/tokens/permissions, new or changed REST endpointTier 2 Full review

    3.4 Contract testing + fuzz testing

    The DORA/LeadDev findings showed that the contract testing before adopting AI reduces change failure rate. Contract testing catches the “clean-looking PR, silent cross-system break” failure mode.

    I have another open source project api-mock-service, that can be used for contract and fuzz testing, e.g.,

    # 1. Load an OpenAPI spec (or record live traffic through the proxy)
    curl -X POST http://localhost:8080/_oapi -F "file=@openapi.yaml" -F "group=billing"
    
    # 2. Run producer contract tests
    curl -X POST "http://localhost:8080/_contracts/billing?baseUrl=http://billing:8080&executionTimes=3"
    
    # 3. Run mutation tests (11 strategies)
    curl -X POST "http://localhost:8080/_contracts/mutations/billing?baseUrl=http://billing:8080&executionTimes=5"
    
    # 4. Export JUnit XML for CI
    curl "http://localhost:8080/_contracts/billing/junit" > contract-results.junit.xml

    3.5 Mutation strategies

    Per-field mutations test individual field validation:

    StrategyWhat it tests
    Missing fieldsRequired field validation
    Boundary valuesMin/max, empty strings, zero, MAX_INT
    Malformed dataWrong types, invalid formats
    Null fieldsNull handling, NPE prevention
    Combinatorial nullsMulti-field invalid combinations
    Format-specific boundariesDate edges, URL length, email format
    Security injectionsSQLi, XSS, path traversal, SSTI, cmd injection, NoSQLi

    Sequence-level mutations test stateful interactions:

    StrategyWhat it tests
    Request reorderingState machine correctness
    Request duplicationIdempotency
    Request omissionRequired step enforcement
    Timing variationsRace conditions, timeouts

    3.6 Contract testing as a Formicary job

    The Formicary pipeline runs as a three-task Formicary DAG (docs/examples/ai-contract-test.yaml). Each task runs in a separate Kubernetes pod with api-mock-service as a sidecar:

    job_type: ai-contract-test
    description: "Contract validation + security fuzzing via api-mock-service proxy"
    max_concurrency: 5
    timeout: 3600s
    
    variables:
      PRNumber:
        type: STRING
        required: false
      ServiceURL:
        type: STRING
        required: false
      Service:
        type: STRING
        required: false
    
    tasks:
    # --- RECORD ------------------------------------------------------------------
    - task_type: record
      method: KUBERNETES
      timeout: 15m
      report_stdout: true
      host_network: true
      working_dir: /workspace
    {{if .Service}}
      services:
        - name: "{{default "service-under-test" .ServiceName}}"
          alias: "{{default "service-under-test" .ServiceName}}"
          image: "{{.Service}}"
          memory_limit: "{{default "4G" .ServiceMemoryLimit}}"
          cpu_request: "{{default "250m" .ServiceCpuRequest}}"
        - name: api-mock-service
          alias: api-mock-service
          image: plexobject/api-mock-service:latest
          ports:
            - number: 8081
          memory_limit: "512Mi"
          cpu_request: "100m"
          command: ["/api-mock-service", "--httpPort", "8081", "--proxyPort", "8082", "--dataDir", "/workspace/recordings"]
          volumes:
            empty_dir:
              - name: workspace
                mount_path: /workspace
    {{end}}
      container:
        image: plexobject/ai-dev-tools:latest
        image_pull_policy: Always
        cpu_request: "250m"
        memory_limit: 512Mi
        memory_request: 256Mi
        volumes:
          empty_dir:
            - name: workspace
              mount_path: /workspace
        env_from:
          - secret_ref: ai-dev-credentials
      environment:
        WORKSPACE_DIR: /workspace
        SERVICE_URL: "{{.ServiceURL}}"
        SERVICE_PORT: "{{default "8080" .ServicePort}}"
        MOCK_SERVICE_PORT: "{{default "8081" .MockServicePort}}"
        PROXY_PORT: "{{default "8082" .ProxyPort}}"
        AI_DEV_TOOLS_DEBUG: "{{.AiDevToolsDebug}}"
      script:
        - python -c "from scripts.common.bootstrap import ensure_debug_mode; ensure_debug_mode()" 2>/dev/null || true
        - touch /tmp/.adt_bootstrap_done
        - python -m scripts.contract.record
      artifacts:
        paths:
          - ./record_result.json
          - ./recordings
        expire_after: 24h
      on_completed: fuzz
      on_failed: notify-error
    
    # --- FUZZ --------------------------------------------------------------------
    - task_type: fuzz
      method: KUBERNETES
      timeout: 20m
      report_stdout: true
      host_network: true
      working_dir: /workspace
    {{if .Service}}
      services:
        - name: "{{default "service-under-test" .ServiceName}}"
          alias: "{{default "service-under-test" .ServiceName}}"
          image: "{{.Service}}"
          memory_limit: "{{default "4G" .ServiceMemoryLimit}}"
          cpu_request: "{{default "250m" .ServiceCpuRequest}}"
        - name: api-mock-service
          alias: api-mock-service
          image: plexobject/api-mock-service:latest
          ports:
            - number: 8081
          memory_limit: "512Mi"
          cpu_request: "100m"
          command: ["/api-mock-service", "--httpPort", "8081", "--proxyPort", "8082", "--dataDir", "/workspace/recordings"]
          volumes:
            empty_dir:
              - name: workspace
                mount_path: /workspace
    {{end}}
      container:
        image: plexobject/ai-dev-tools:latest
        image_pull_policy: Always
        cpu_request: "500m"
        memory_limit: 2G
        memory_request: 512Mi
        volumes:
          empty_dir:
            - name: workspace
              mount_path: /workspace
        env_from:
          - secret_ref: ai-dev-credentials
      dependencies:
        - record
      environment:
        WORKSPACE_DIR: /workspace
        SERVICE_URL: "{{.ServiceURL}}"
        SERVICE_PORT: "{{default "8080" .ServicePort}}"
        MOCK_SERVICE_PORT: "{{default "8081" .MockServicePort}}"
        PR_NUMBER: "{{.PRNumber}}"
        AI_DEV_TOOLS_DEBUG: "{{.AiDevToolsDebug}}"
      script:
        - python -c "from scripts.common.bootstrap import ensure_debug_mode; ensure_debug_mode()" 2>/dev/null || true
        - touch /tmp/.adt_bootstrap_done
        - python -m scripts.contract.fuzz
      artifacts:
        paths:
          - ./fuzz_result.json
          - ./fuzz_results.xml
          - ./contract_test_summary.json
        expire_after: 24h
      on_completed: report
      on_failed: report

    The recording proxy is a sidecar container (plexobject/api-mock-service:latest) that shares the pod network namespace. Traffic flows:

    test script ? localhost:8081 (proxy) ? localhost:8080 (service under test)
                                    ?
                        api_contracts/**/*.yaml   (saved per HTTP interaction)

    3.7 Skills

    Two complementary skills in the you-got-skills library:

    • /ygs-contract-test: Detects API surface, sets up api-mock-service, runs contract validation and mutation testing.
    • /ygs-fuzz-test: Generates fuzz corpus from API specs, applies all mutation strategies plus 8 CWE-classified injection classes.

    Here’s the fuzz skill’s decision table for endpoint discovery:

    Signal foundAction
    OpenAPI/Swagger specParse endpoints + schemas directly
    Express routes (app.get/post)Extract route patterns from code
    Flask decorators (@app.route)Extract route patterns from code
    Go http.HandleFuncExtract route patterns from code
    Spring @RequestMappingExtract route patterns from code
    None of the aboveBLOCKED — no endpoints to fuzz

    The 8 CWE-classified injection classes the fuzz skill exercises:

    Injection classCWEExample payload
    SQL injectionCWE-89' OR 1=1 --
    XSSCWE-79<script>alert(1)</script>
    Path traversalCWE-22../../etc/passwd
    SSTICWE-1336{{7*7}}
    Command injectionCWE-78; cat /etc/passwd
    NoSQL injectionCWE-943{"$gt": ""}
    LDAP injectionCWE-90`)(uid=))(
    XXECWE-611<!DOCTYPE foo [<!ENTITY xxe SYSTEM "file:///etc/passwd">]>

    3.8 Live example: OWASP WrongSecrets

    OWASP WrongSecrets is a Java Spring Boot application intentionally seeded with secrets management vulnerabilities.

    Trigger from Slack:

    @sb-slack contract-test https://github.com/OWASP/wrongsecrets \
      --service jeroenwillemsen/wrongsecrets:latest-no-vault

    The job produces output such as:

    [contract] api-mock-service ready at http://localhost:8081/_health (HTTP 404)
    [fuzz] contracts: succeeded=0 failed=0
    [fuzz] recordings_dir=/workspace/recordings exists=True contracts_dir_exists=True
          yaml_count=21 first_3=[
            '/workspace/recordings/api_contracts/GET/Recordedroot--200-c1ec01cf.yaml',
            '/workspace/recordings/api_contracts/status/GET/...',
            '/workspace/recordings/api_contracts/debug/GET/...'
          ]
    [fuzz] discovered 21 endpoints: [
      ('GET', '/status'), ('GET', '/debug'), ('GET', '/metrics'),
      ('GET', '/openapi.json'), ('GET', '/docs'), ('GET', '/'),
      ('GET', '/swagger'), ('GET', '/challenge/21'), ('GET', '/challenge/28'),
      ('GET', '/challenge/1'), ('GET', '/health'), ('GET', '/api'),
      ('GET', '/api/challenges'), ('GET', '/api/Challenges'), ('GET', '/admin'),
      ('GET', '/v1'), ('GET', '/actuator/env'), ('GET', '/actuator/info'),
      ('GET', '/actuator/beans'), ('GET', '/actuator/health'),
      ('GET', '/actuator/mappings')
    ]
    [contract] probe sqli:/ ? 200 finding=False
    [contract] probe path_trav:/ ? 404 finding=False
    [contract] probe sqli:/status ? 404 finding=False
    [contract] probe path_trav:/status ? 404 finding=False
    ... (40 probes total) ...
    [fuzz] 21 endpoints, 0 findings, 0 critical
    [fuzz] complete: iterations=40 findings=0 critical=0 status=PASS ams_used=True

    3.9 What a skill looks like: /ygs-contract-test

    Skills are structured Markdown files in the you-got-skills library. Here’s the decision table from /ygs-contract-test:

    ```yaml
    # SKILL.md frontmatter
    ---
    name: ygs-contract-test
    description: API contract testing — record, derive contracts, detect breaking changes
    argument-hint: "<service-or-pr> [--mode record|validate|diff]"
    ---
    ```
    
    ```markdown
    ## Step 1: Determine contract testing mode
    
    | Mode     | When to use                          | What happens                        |
    |----------|--------------------------------------|-------------------------------------|
    | record   | First time, or updating baseline     | Run tests through proxy, capture    |
    | validate | PR review, CI gate                   | Compare against existing contracts  |
    | diff     | Breaking change detection            | Compare two contract versions       |
    
    Default to `validate` if a baseline exists, otherwise `record`.
    
    ## Step 2: Record API interactions
      api-mock-service --mode record --proxy-port 8081 ...
      HTTP_PROXY=http://localhost:8081 pytest tests/integration/
    
    ## Step 3: Derive and validate contracts
      curl -X POST "http://localhost:8080/_contracts/{group}?baseUrl=..."
    
    ## Step 4: Run mutation testing (11 strategies)
      curl -X POST "http://localhost:8080/_contracts/mutations/{group}..."
    
    ## Step 5: Export and report
      Report DONE if catch rate > 80%. BLOCKED if contracts fail.
      Suggest: /ygs-fuzz-test for deeper coverage, /ygs-security-review for findings.

    4. Speed at the Pipeline

    The second knob from earlier simulation is pipeline speed, which is a cheaper lever. This section covers formal verification of the queue protocol itself, build/test optimization with real numbers.

    4.1 Formal verification

    Tests check specific inputs but formal verification proves properties hold over all possible inputs. For a concurrent system like a merge queue, the state space is too large to test exhaustively.

    TLA+ specification of the merge queue

    The spec lives at docs/examples/specs/merge_queue.tla and it models the core merge queue as a state machine:

    Safety for Scope isolation

    ScopeIsolation ==
        \A s1, s2 \in Scopes :
            s1 /= s2 => lanes[s1] \cap lanes[s2] = {}

    Safety for Test before merge:

    TestBeforeMerge ==
        \A pr \in PRs :
            prState[pr] = "merged" => testResults[pr] = "pass"

    Liveness to ensure every PR eventually merges or is ejected:

    Progress == \A pr \in PRs :
        prState[pr] = "queued" ~> (prState[pr] = "merged" \/ prState[pr] = "ejected")

    Dafny verification of merge invariants

    The Dafny spec at docs/examples/specs/verified_merge.dfy proves four properties at compile time: batch merging preserves main branch health, bisection correctly partitions PRs, risk scoring is monotonic, and shard partitioning loses no tests:

    predicate MainBranchHealthy(mergedPRs: set<PR>)
    {
      forall pr :: pr in mergedPRs ==> pr.testResult == Pass
    }
    
    lemma MergeBatchPreservesHealth(mainBranch: set<PR>, batch: set<PR>)
      requires MainBranchHealthy(mainBranch)
      requires forall pr :: pr in batch ==> pr.testResult == Pass
      ensures MainBranchHealthy(mainBranch + batch)
    {
      // Proof is automatic: union of two sets where all elements
      // satisfy the predicate still satisfies the predicate.
    }

    4.2 Test-impact analysis: scripts/mq/test_impact.py

    The default behaviour is deliberate: always run the full test suite, partitioned into balanced shards. Diff-scoped runs are opt-in via --diff-scope and they require an explicit --head <branch> to define what to diff against.

    # Default: full suite, 8 shards, no diff analysis
    python -m scripts.mq.test_impact --pr-number main --num-shards 8
    
    # Explicit diff-scope: numeric PR (base branch auto-extracted from GitHub API)
    python -m scripts.mq.test_impact --pr-number 42 --num-shards 8 --diff-scope
    
    # Explicit diff-scope: branch name (must supply --head as base to compare against)
    python -m scripts.mq.test_impact --pr-number feature/my-branch \
        --head main --num-shards 8 --diff-scope

    The script performs following operation:

    1. Fetches changed files via gh pr view --json
    2. Detects language from file extensions (Python, Go, TypeScript, Java, Kotlin, Ruby, C#, Rust)
    3. Maps to test files using naming conventions
    4. Traverses imports one level deep to find transitive dependents
    5. Partitions into balanced shards using greedy bin-packing with historical timing data
    6. Falls back to the full suite if no tests map to changed files

    4.3 Parallel fan-out: ai-parallel-test

    Test-impact analysis feeds directly into Formicary’s fan-out for parallel shard execution (docs/examples/ai-parallel-test.yaml):

    # analyze task: runs test_impact.py, emits ::add-job-context TestShards::
    - task_type: analyze
      script:
        - python -m scripts.mq.clone_pr --pr-number {{.PRNumber}}
        - python -m scripts.mq.test_impact --pr-number {{.PRNumber}}
        # TestShards is now in job context — fan-out reads it directly
    
    # Fan-out: one task per test shard, running in parallel
    # 4 CPU cores + 16G memory per shard ? real intra-shard parallelism
    - task_type: run-tests
      container:
        cpu_request: "4000m"
        memory_limit: 16G
      fan_out:
        source: TestShards
        item_var: shard
        max_parallel: {{.MaxShards}}  # unquoted integer — YAML strict types
        fail_fast: false
      script:
        - python -m scripts.mq.clone_pr --pr-number {{.PRNumber}}
        - python -m scripts.mq.run_scoped_ci --shard "{{.shard}}"

    Each shard runs independently on its own Kubernetes pod. The polyglot test runner (scripts/mq/run_scoped_ci.py) detects the project type from marker files and builds the right command.

    Here is a sample output:

    ## Test Impact Analysis
    **191** / 191 tests selected (**0%** reduction) across **2** shards
    
    ## Test Results
    ? **87** / 87 passed
      ?? 130s wall clock across 2 shards
    
    ### Shard Performance
    | Shard | Tests | Passed | Failed | Duration | Status |
    |-------|-------|--------|--------|----------|--------|
    | 1     | 44    | 44     | 0      | 130.5s   | ?     |
    | 0     | 43    | 43     | 0      | 115.4s   | ?     |
    
    > **Parallel speedup:** 246s sequential ? 130s parallel (1.9x across 2 shards)
    
    ### Slowest Tests
    | Duration | Test                          |
    |----------|-------------------------------|
    | 30.01s   | `TestInitEnabled`             |
    | 30.01s   | `TestInitEnabled`             |
    | 10.02s   | `Test_ShouldAsyncWithSleep`   |
    | 10.01s   | `Test_ShouldAsyncWithSleep`   |
    |  4.68s   | `Test_EncryptDecrypt`         |
    |  2.45s   | `Test_ShouldRealGet`          |
    |  2.28s   | `TestSharedSubscription`      |
    
    ### Test Health Insights
    - ? Pass rate: 100% (87 tests)
    - ? Shard balance: 12% imbalance (115s – 130s)
    - ? 4 slow tests likely using real sleeps — consider mocking time or
      reducing timeouts: `TestInitEnabled` (30.0s), `Test_ShouldAsyncWithSleep` (10.0s)
    - ?? 3 slow test names appear in multiple shards (possible test duplication):
      `TestInitEnabled`, `Test_ShouldAsyncWithSleep`, `TestSharedSubscription`

    4.4 Deployment risk

    The merge queue optimizes everything up to the point where the commit lands on main but merged PRs still have to ship. The original simulation from Joe Magerramov’s formula models the probability that a merge batch is clean and continuous delivery. In my simulation I added support for batch PRs as it’s used in many organizations. But this also increases risk, e.g., a merge queue with a 95% batch success rate can still lead to a weekly release train that ships clean only 36% of the time. The formula extends naturally:

    release_success = (1 - defect_rate) ^ prs_per_release
    

    CD is the special case where prs_per_release = 1, e.g.,

    Defect rateCD (1 PR)Daily (10 PRs)Weekly (50 PRs)Bi-weekly (100 PRs)
    0.5%99.5%95.1%77.8%60.6%
    1%99.0%90.4%60.5%36.6%
    2%98.0%81.7%36.4%13.3%
    5%95.0%59.9%7.7%0.6%

    At 2% defect rate, a CD pipeline has a 98% clean deployment rate. The same team on a weekly train: 36%. Two-thirds of their releases ship with at least one defect.

    The rollback trap

    Suppose release R1 ships with a hidden bug. By the time monitoring detects it, R2 and R3 have already deployed. Rolling back to pre-R1 means reverting R2 and R3 too.

    Stacked releasesRollback feasibilityRecovery strategyMTTR multiplier
    1? StraightforwardRollback1.0x
    2-3?? Costly but possibleRollback (with intermediate revert)1.5x
    ?4? ImpracticalRoll-forward (fix-forward only)2.5x

    One bad PR poisons the entire train and the release sits in an unknown state while someone bisects 50 changes to find the culprit.

    Deployment maturity

    Following are best practices from the industry for improving deployment maturity:

    DimensionWeightWhat it enables
    Automated testing2.0xCI on every PR catches defects pre-merge
    Canary deployment2.0xProgressive rollout catches prod-only failures
    Automated rollback1.5xInstant revert on anomaly detection
    Observability1.5xAlerting + metrics detect failures fast
    Wave deployment1.0xStage ? preprod ? prod progression
    Feature flags1.0xDecouple deploy from release
    Blue/green0.5xZero-downtime deploy infrastructure

    The maturity score translates into a risk multiplier — a low multiplier means your infrastructure catches most defects before they reach users:

    adjusted_success = release_success + (1 - release_success) × (1 - risk_multiplier)
    

    The MQ report now shows two critical markers:

    • ? Blue line: your current position given your defect rate, batch size, and deployment maturity.
    • ? Red line: the calamity threshold: the batch size where success drops below 70%.

    The calamity threshold formula:

    max_safe_batch = log(0.70) / log(1 - defect_rate)
    

    The deployment risk model points to three levers:

    1. Reduce batch size. Move from weekly to daily trains.
    2. Invest in deployment maturity. Automated testing (weight 2.0x) and canary deployment (weight 2.0x) are the two highest-leverage dimensions.
    3. Prevent release stacking. Each additional stacked release multiplies recovery cost.

    5. The Learning Loop

    5.1 Multi-agent coordination

    The merge queue is naturally an actor system where each PR is an actor with a lifecycle (queued ? batched ? testing ? merged/ejected). The properties that matter at scale:

    • Location transparency: a PR doesn’t care which CI runner tests it
    • Supervision: when a batch fails, the supervisor (bisect task) localizes the fault and recovers
    • Isolation: scope-aware lanes are isolation boundaries

    Better decomposition

    Similar to the microservice architecture, the path forward is better decomposition:

    • Smaller PRs
    • Better module boundaries
    • Formal verification
    • Automated scope detection
    • Dependent PR grouping, e.g., group related PRs so if one fails, the entire group is ejected.

    5.2 The skills improvement flywheel

    5.3 Metrics

    The four DORA keys (deployment frequency, lead time, change failure rate, failed deployment recovery time) are table stakes. At agent-scale, you need additional metrics like:

    • PR size
    • PR pickup time Stale PRs conflict more often and carry higher bisection cost in the merge queue.
    • AI PR acceptance rate.
    • Rework rate.
    • Cycle time decomposition break first-commit-to-deploy into pickup, review, merge, deploy sub-phases.

    5.4 Path forward

    The bottom line from all the evidence in this post:

    1. Don’t adopt AI at agent-scale before your delivery system can handle agent-scale quality.
    2. Pipeline speed is the cheaper lever. Defect rate has diminishing returns; pipeline duration has an order of magnitude of headroom.
    3. The merge queue is not a queue. It’s a coordination problem: scope-aware routing, risk-gated review, speculative batching, bisection-based fault recovery, and a learning loop.
    4. Formal verification catches what tests can’t. For concurrent protocols, tests check paths and TLA+ checks invariants.
    5. Specs reduce defects at source. Structured acceptance criteria, design docs for high-impact changes, and tiered review based on blast radius.

    6. Getting Started

    Everything above is code in three repos: formicary (workflow engine), ai-dev-tools (Python scripts), and you-got-skills (Claude skill library). Here’s how to test each piece.

    Prerequisites

    # 1. Deploy Formicary to Kubernetes (single-container: queen + ant + storage)
    kubectl apply -f k8s.yaml
    kubectl port-forward svc/formicary 7777:7777 19000:19000
    
    # 2. Get an API token from the UI: http://localhost:7777 ? Profile ? API Token
    export FORMICARY_TOKEN="<jwt>"
    export GH_TOKEN="<github-pat>"
    export SLACK_CHANNEL="your-channel"
    
    # 3. Deploy all workflows + set configs
    cd docs/examples
    ./deploy-ai-workflows.sh --create-k8s-secret --set-configs \
      --gh-org bhatti --gh-repo todo-api-errors --bedrock

    Test 1: Scope router

    In Slack, mention the bot in your configured channel:

    @bot scope main --repo https://github.com/org/repo
    

    This prompt triggers ai-scope-router workflow and then creates a report showing the scope, blast radius, coverage, historical data for defects, security and other metrics.

    Test 2: Parallel test with test-impact analysis

    @bot parallel-test --repo https://github.com/org/repo main
    @bot parallel-test --repo https://github.com/org/repo --branch develop

    This prompt triggers ai-parallel-test workflow that fans out multiple shards for parallel tests and then creates a report showing the selected vs total tests, reduction percentage, per-shard pass/fail counts, wall-clock time.

    Test 3: Risk-gated review

    @bot gate-review https://github.com/org/repo/pull/1

    This prompt triggers ai-gate-review workflow and generates a report showing blast radius, security, test coverage and other metrics.

    Test 4: Contract testing

    @bot contract-test 1 --service myorg/my-api:latest
    

    This prompt triggers ai-contract-test workflow and executes contract test against the target service. It then generates a report with the pass/fail counts, mutation catch rate, security findings (CWE-classified)

    Test 5: Merge Queue Analysis

    @bot mq https://github.com/org/repo
    

    This prompt triggers ai-merge-queue workflow, evaluating health of the merge queue and calculating the risk/blast radius/complexity of all open PRs. Here is a sample snippet from the merge queue report:

    ? 250 open PRs across 22 scope lanes — repo: acme/platform

    MetricValue
    Total open PRs250
    Scope lanes22
    High-risk PRs140
    Needs human review148
    Conflict risk lanes1

    ?? 148 PR(s) flagged for human review before merge (high blast-radius, failed CI, or sensitive paths detected).

    ?? 1 lane(s) have conflict risk multiple PRs touching overlapping paths. Merge one at a time.

    Overall status: ? Unstable (basis: 74.8% PRs aged >48h — CI status unavailable from API)

    MetricValue
    Total PRs in queue250
    CI failure rateN/A — not reported by API (verify in CI dashboard)
    PRs aged >48h187 (74.8%)
    High blast-radius PRs140
    Medium blast-radius PRs82
    Approx batch size (PRs/lane)10.0

    Unstable: 74.8% of PRs are aged >48h. Queue is filling faster than it drains — unblock stale PRs or shrink batch size.

    CategoryPRsBug PRsHigh-BlastHotspot
    ? ui12211??? yes
    ? data323??? yes
    ? api315??? yes
    ? authn_authz315??? yes
    ? security153??? yes
    sre121??—
    backend30??—
    config10??—
    unknown30——

    Test 6: PR Audit

    @bot pr-audit
    

    This prompt triggers ai-gh-pr-audit workflow, evaluating the risk/blast radius/complexity of all merged PRs. Here is a sample snippet from the pr audit report:

    #PRIssueAuthorTitleCatTypeBlastRiskLOCFilesCx
    1#48401PROJ-44822jsmithAuthCoordinator for REST OAuth? authn_authz? securityhigh? 721,84723? high
    2#47228AI-4931aleefix(ai): allow regional Bedrock pr…? data? bughigh? 5889214? medium
    3#48768DL-2875mchenImplement more fine-grained IAM roles? authn_authz? featurehigh? 4563411? medium
    4#48909—jdoeSSRF guard for external MCP connectorui? bugmedium? 323128? low
    5#48827PLAT-16031kpatel[Flaky Test] fix race in bottlenec…api?? refactorlow? 12873? low
    …

    Related Reading

    September 11, 2026

    What Happens After the PR Merges: Building the Learning Loop Software Factories Are Missing

    Filed under: Computing — admin @ 8:51 pm

    Turning every review comment, rubber-stamped PR, and missed acceptance criterion into a process improvement.


    TL;DR

    • Large companies like Stripe, Uber, and Spotify merge thousands of agent-authored PRs a week. What’s unsolved is whether the organization gets smarter from them.
    • AI helps with technical debt (in code) but accelerates cognitive debt (in people) and intent debt (in artifacts) because people accept AI output without building real understanding.
    • Instead of removing the code review, this post suggests risk-stratified review by blast radius.
    • What most teams are missing is an audit layer, e.g., a periodic scan of the last N merged PRs that tells you whether your review process is actually working.
    • I have built a five-stage flywheel: a PR merges -> ygs-learn extracts what reviewers taught the system -> the next agent run reads those learnings -> periodic ygs-pr-audit and ygs-code-audit runs surface systemic patterns.

    Agentic Engineering Maturity Model

    When I started using unit testing and test driven development decades ago, it taught me that tests don’t eliminate bugs. Instead, they create a feedback loop that makes the next bug less likely. For example, you may write a test as part of the development process or as a result of a production bug that you had to fix. With each test, you build a regression suite that acts as a quality gate, and serves as an institutional memory of every mistake the teams made. Agentic engineering needs the same thing where they don’t just write code but a flywheel that turns every review comment or incident into process or a skill so that the organization gets better over time. For example, I see three levels of how teams use AI in their development workflow:

    • Pull systems: You ask an agent to do something, e.g., interactive sessions with Claude Code but the learning disappears when you close the session.
    • Push systems (“software factories” or “ambient agents”): The system asks the agent to do things continuously, e.g., pick up issues, write code, open PRs. Humans are then looped in when a decision needs judgment.
    • Learning systems: The push system that also records what it learned, closes the loop, and updates its own skills and guardrails. This allows the next run to get better because the organization’s captured knowledge improved.

    I have been building automated workflows for the third level that I will share in this post.


    What I’ve Already Built

    Over the past year I’ve built three interconnected open-source systems, covered in two earlier posts:

    I have built several declarative AI workflow definitions based on these primitives:

    ai-gh-implement.yaml     — plan ? implement ? create-pr ? poll-pr
    ai-gh-review.yaml        — automated code review on a single PR
    ai-standup-gh.yaml       — daily standup brief from GitHub + Slack
    ai-gh-pr-audit.yaml      — PR gap analysis ? skill improvement PR ? learn
    ai-codebase-audit.yaml   — codebase archaeology across N commits
    ai-gh-issue-picker.yaml  — intelligent issue triage and selection
    ai-jira-implement.yaml   — same pipeline for Jira/Bitbucket repos
    ai-jira-review.yaml      — code review for Bitbucket PRs
    ai-jira-pr-audit.yaml    — PR audit for Jira-tracked repos
    ai-adhoc.yaml             — run any skill from a Slack command

    I also built 40+ skills, including ygs-implement, ygs-code-review, ygs-review-deep, ygs-security-review, ygs-sre-review, ygs-learn, ygs-pr-audit, ygs-codebase-audit, ygs-standup, ygs-risk-scan, ygs-sprint-plan, ygs-estimate, ygs-qa, and ygs-ship.

    Slack Is the Interface

    I have been using Slack to interact with the orchestration engine, e.g., an engineer types a command in a channel, Formicary picks it up, runs the right workflow on Kubernetes, and threads all status updates and results:

    Engineer:  @bot implement acme/backend#142
    Bot:       ? Planning implementation for issue #142...
    Bot:       ? Plan ready — 3 files, estimated complexity: medium (using Sonnet)
    Bot:       ?? Implementing... (48 turns)
    Bot:       ? Self-review found 0 critical issues
    Bot:       ? PR #287 opened: https://github.com/acme/backend/pull/287
    Bot:       ? Watching PR for review comments...
      [2 hours later, reviewer comments on PR]
    Bot:       ? Responding to 2 review comments, pushing fixes...
      [reviewer approves, PR merges]
    Bot:       ? PR #287 merged. Running ygs-learn to extract learnings...
    Bot:       ? 1 new learning captured: "Cache invalidation must check ACL per-request"

    The ai-adhoc.yaml workflow lets engineers run any skill on demand so the full skill library is available without leaving Slack.

    The Implement Pipeline

    The workhorse workflow is ai-gh-implement, a six-task DAG:

    plan ? implement ? self-review ? create-pr ? poll-pr ? done

    In above workflow, self-review is invoked after the agent implements a change to review its own diff against the plan and the issue’s acceptance criteria. It can loop in human review if it finds any issues and the engineer resolves it before the PR even opens. Also, it uses a complexity-based model routing where the plan step estimates task complexity (low, medium, high) and writes it to an artifact. The implement then picks the model to match to keep the costs proportional to difficult.

    The Learning Skills

    Following are learning workflows that I built:

    • ygs-pr-audit: analyzes the last N merged PRs across four dimensions (spec, design, skills, process), surfaces systemic gaps with evidence, and produces a list of recommended skill updates.
    • ygs-codebase-audit: runs structural archaeology across N commits, detecting hotspots, knowledge silos, duplicate abstractions, test brittleness, and security gaps using git/grep/find commands against the codebase.
    • ygs-learn: fires automatically when a PR merges inside ai-gh-implement. The poll-pr task detects the merge and calls ygs-learn, which extracts atomic learnings from review comments and writes learning documents that future agent runs read.

    The ygs-learn is wired into poll-pr so every merged PR automatically triggers learning extraction.


    The Factory Floor Problem

    The Software factory or ambient agents create a system that proactively respond to events (a new issue, a failing test, a merged PR) instead of waiting for a human to ask. For example, Igor Ostrovsky shared common patterns where a software factory starts work from events and coordinates agents through engineering workflows like agents reading/writing specs/issues/pull requests while people approve key decisions. He models it as a graph, e.g., nodes are artifacts, edges are agent-driven transformations, and approvals and CI checks gate progress. This maps onto what Formicary does: a DAG of Kubernetes tasks, where each task produces artifacts the next one consumes with exit-code-based routing and human decision points. Uber shared a similar layering in their “agentic SDLC platform“, a context graph, an LLM gateway, and a governance layer underneath their agents. The Formicary / ai-dev-tools / you-got-skills setup is a smaller, open-source version of the same idea.

    Other large companies are using similar software factories:

    • Stripe’s Minions merge 1,000+ PRs per week in Ruby with human PR review as the gate.
    • Uber’s Minion opens roughly 11% of company-wide PRs directly andover 70% of all PRs agent-attributed.
    • Spotify’s Honk reports 650+ agent-authored merged PRs per month with up to 90% time savings on migration-style work.

    Here’s what a Formicary workflow definition looks like in practice:

    - task_type: poll-pr
      method: KUBERNETES
      timeout: 168h
      dependencies:
        - create-pr
        - poll-pr          # self-dependency: re-runs after PAUSE_JOB
      script:
        - python -m scripts.gh.check_pr_state --issue-id {{.IssueNumber}}
        - python -m scripts.gh.fetch_comments --issue-id {{.IssueNumber}}
        - python -m scripts.gh.respond_comments --issue-id {{.IssueNumber}}
      on_exit_code:
        3: PAUSE_JOB       # PR still open — wait, then re-run
      delay: "120s"        # check every 2 minutes

    PAUSE_JOB allows a loop without blocking, e.g., poll-pr sees the PR is still open, it exits with code 3. Formicary pauses the task, waits 120 seconds, and re-runs it. When a reviewer leaves comments, the task fetches them, has Claude respond, and pushes fixes. When the PR merges, check_pr_state calls ygs-learn to extract learnings.

    But none of above workflows answer a key question: What did the organization learn from the last 50 PRs that merged?Which review patterns are failing? Where are knowledge silos growing? What mistakes keep repeating? Are reviewers actually reviewing, or rubber-stamping?

    The current wave of software factories optimizes the pipeline (issue to deploy) and treats code review as a throughput bottleneck to shrink. Some teams like StrongDM are removing it entirely. That may work for some organizations but it’s not practical for most enterprises. Code review isn’t just a quality gate. For example, Google’s internal engineering practices and a recent study of 3,100 practitioner opinions treat knowledge transfer as equally important as defect detection. Review is where junior engineers absorb senior engineers’ mental models and where architectural decisions are explained. It builds a shared understanding of how the system works or what Fred Brooks calls a conceptual integrity.


    Three Kinds of Debt

    Margaret-Anne Storey’s Triple Debt Model defined following terminology for what goes wrong when throughput outpaces understanding:

    • Technical Debt (Code): This is already known term for poor implementation choices that make systems harder to change, tangled dependencies, missing abstractions, etc. However, AI is actually helping with technical debt with automated refactoring and AI-generated test suites.
    • Cognitive Debt (People): Cognitive debt is the erosion of shared understanding across a team over time. The team can no longer confidently explain how the system works or predict the impact of a change. For example, when AI generates code, a developer might accept it without building real understanding of it. At scale, across a team, that becomes an accumulation of not knowing. Researchers call this cognitive surrender. In another follow-up post, Storey lists symptoms such as: engineers lose confidence making changes, debugging takes longer, onboarding takes longer, and code review overhead goes up. The fixes are inherently human: code review, pair programming, system walkthroughs, retrospectives.
    • Intent Debt (Artifacts): Intent debt is the absence of externalized rationale that explain what a system is for and how it should evolve. In other words, nobody wrote down why the system works this way, so nobody. Addy Osmani’s blog on intent-debt explains that the cost of undocumented intent used to be paid by humans touching the code but now every agent pays it. My skills repository like the docs/prd/, docs/adr/, and docs/learnings/ directories allow agents to make decisions based on what the system is actually trying to accomplish.

    The Triple Debt in Practice

    Debt TypeWhere It LivesWhat Fixes ItWhat Accelerates It
    TechnicalCodeTests, static analysis, refactoringShipping without testing
    CognitivePeopleReview, pairing, walkthroughs, onboardingAI-generated code accepted without understanding
    IntentArtifactsADRs, specs, skills, learningsEvery agent session that starts without captured rationale

    Generative AI reduce technical debt through automated refactoring and test generation but it accelerates cognitive and intent debt, because the team ends up producing code faster than it can build shared understanding of what that code does and why.


    Why Humans Stay in Review

    The answer isn’t to keep the review process as it was or remove it all together. It’s risk-stratified review based on risk and blast radius. For example, Meta published Automating Low-Risk Code Review (RADAR) that described a multi-stage funnel for deciding which changes can land without human review:

    1. Eligibility gates: is this diff type eligible for automation at all?
    2. Diff Risk Score: a machine-learned model predicting the likelihood of a production incident.
    3. LLM-based automated code review: an AI reviewer that can approve with high confidence.
    4. Deterministic validation checks: a final safety net.

    Teams at Meta tune their own risk thresholds so that they can automate the low-risk changes and route the high-risk ones to humans. Qodo also ships a classifier that labels PRs by scope and sensitivity in context so that reviewers can triage their attention. That’s what we are trying to do in my organization so that for high-blast-radius changes like auth, billing, security plugins, etc., humans stay in the loop. Also, I built something that other doesn’t have: an audit layer that tells you whether your risk calibration is actually working.

    The Audit Layer

    RADAR-style gating prevents the risky merge from happening but my PR audit catches when the gating itself is failing. For example, when reviewers are rubber-stamping security changes, bot-flagged critical findings go unacknowledged, or when large architectural PRs ship with two comments. Here are examples of my PR audit runs (details anonymized):

    MetricValueBenchmarkSignal
    Rubber-stamp rate (high blast-radius)50%<10% healthy / >25% problem? Problem
    Security review coverage0%—? Gap
    Bot-finding follow-through~80%>90% healthy / <70% gap?? Warning
    Human review burden~78%<40% healthy / >70% overloaded?? Overloaded
    Formal acceptance criteria coverage13%>80% healthy / <50% reactive? Critical

    Half of the high-blast-radius PRs like auth, billing, and security got rubber-stamp approvals with little substantive comments. The code-review bot caught roughly 22% of substantive findings versus humans’ 78%, meaning the team was leaning heavily on human reviewers for architectural and security issues. In A practical guide to risk-based code review, Cortex’s research predicts that a reviewer’s ability to catch defects collapses past 200–400 lines of code. AI generates code faster than that threshold so approval rates rise while inline comments fall. The ygs-pr-audit skill identifies gaps in the automated review process so that we can improve skills, processes, specifications and other aspects of the development process.


    The Flywheel for Learning Culture

    In order to continuously improve your development process, you need a feedback loop for continuous learning. You can consider it as a learning flywheel for the workflows of entire software development process:

    Stage 1: A PR merges. poll-pr detects it.

    Every implementation workflow has a poll-pr task that checks PR state periodically. When the status flips to merged (or declined), the task exits the loop and triggers the next stage.

    Stage 2: ygs-learn extracts learnings.

    On merge, the system invokes ygs-learn in extract mode. Before touching review comments, it runs Phase 0: PR Health Check across five dimensions:

    • Spec coverage: Did the linked issue have formal acceptance criteria?
    • Design decisions: Were architectural choices documented, or did the PR just change code?
    • Security/SRE: Did the PR touch security-sensitive paths?
    • Review quality: How many substantive human comments? Was there a rubber-stamp?
    • CI health: Did CI pass cleanly, or were there flaky retries?

    Phase 0 produces a one-paragraph health summary that becomes metadata on the learning. It also catches the edge case where a PR merges with no review comments at all. After Phase 0, ygs-learn enters extract mode. It fetches every review comment from the merged PR and asks: does this reveal a recurring pattern, a gotcha that would apply to future work? For each actionable insight, ygs-learn sorts it into one of seven categories: Edge Case, Integration Gotcha, Performance Cliff, Security Trap, Process Friction, Domain Rule, or Tooling Quirk and writes a structured document:

    # Guard Bypass After Cache: Security Gate Evaluated Once at Load Time
    
    **Category:** Security Trap
    **Date:** 2026-09-08
    **Source:** Review feedback on PR #XXXX — reviewer caught that
    the authorization check runs once at connection creation, not per query
    
    ## Learning
    
    If a security gate (e.g., isFeatureAllowed()) is checked once at load time
    and the result is cached, the gate must also be checked before every use
    of the cached object. Otherwise, a permission revocation after initial
    load is silently ignored.
    
    ## Evidence
    
    Database connection manager called isAllowed() once in getConnection(),
    cached the client, and returned the cached client on subsequent calls.
    A reviewer identified that revoking the feature flag mid-session would
    have no effect.
    
    ## Application
    
    When reviewing code that caches authenticated/authorized resources:
    check whether the authorization check is repeated before each use,
    not just at creation time.

    Each learning is stored in docs/learnings/YYYY-MM-DD-slug.md with deduplication.

    Stage 3: The next implementation run reads the learnings.

    When the next ygs-implement pipeline runs, its planning and implementation steps load docs/learnings/ as context. If the last PR’s reviewer caught a cache-bypass vulnerability, the next agent run that touches caching code has that learning in front of it. This is similar to NVIDIA’s MAPE control loop for production AI agents: Monitor (poll-pr watches for the merge and collects review signals) -> Analyze (ygs-learn categorizes what happened and why it matters) -> Plan (writes learnings to docs/learnings/) -> Execute (the next implementation pipeline reads those learnings as context). Augment Code describes something similar like execute, coach, distill, improve in their Agent Learning Flywheel blog except their coaching happens synchronously in Slack but my extraction happens asynchronously after the PR closes.

    Stage 4: ygs-pr-audit looks across many PRs and recommends systemic changes.

    Individual learnings are useful but the real leverage comes from periodic audits that scan last N merged PRs at once and spot common patterns like is the rubber-stamp rate climbing or are bot findings getting ignored before merge? The key output for the flywheel is skill_improvements.json, a list of recommended changes to the repo’s skill files: create .claude/skills/security-review/SKILL.md, update .claude/docs/testing.md with a flaky-test-triage section, etc.

    Stage 5: The audit opens a real PR with skill improvements.

    The plan-skill-updates task reads the audit findings and writes an update plan. create-skill-pr implements those changes and opens a pull request against the repo. poll-pr then watches that skill-improvement PR, responds to human reviewer comments. The system learns from its own improvement process.

    Every report includes a Positive Patterns section naming exemplary reviewers, e.g., the person who traced a security issue to its root cause, the reviewer who ran local tests to verify a bot-authored PR, etc. This is the blameless-postmortem move applied to code review. The postmortem maturity literature teaches that punishing people for incidents makes people stop surfacing incidents. Recognizing good work and documenting thoroughness are important design decisions.

    Sharing Learnings Across the Organization

    Though, individual repo learnings live in docs/learnings/ inside each repo but I created a shared skills repository you-got-skills for common skills like ygs-implement, ygs-code-review, ygs-security-review, and the rest. This shared repository is improved based on continuous learnings from the audit reports.


    Inside the PR Audit: A Concrete Walkthrough

    Let me show how the PR audit works in the format I used in Killing the State Machine, where every design choice is explained with its reasoning.

    The workflow triggers from Slack:

    @bot pr-audit acme/backend
    or 
    @bot pr-audit acme/backend --focus skills --n-prs 30

    Formicary picks up the message, resolves it against the ai-gh-pr-audit job type, and runs a five-task DAG:

    The audit-prs Task

    The Python harness (run_pr_audit.py) does the heavy lifting before the LLM sees the data: clone the repo, fetch the last N merged PRs, classify every comment and check acceptance criteria. The LLM then invokes Claude with the ygs-pr-audit skill, which runs a disciplined 5-phase protocol:

    • Phase 1 Setup. Load shared references, check whether the repo has its own skill overrides in .claude/skills/ and read the pre-computed PR data.
    • Phase 2 Four specialist passes. This is the core analysis for the most common failure mode. Each PR is scored across all four dimensions and the findings are tagged by dimension ([SPEC], [DESIGN], [SKILL-GAP], [PRACTICE]).
    • Phase 2.5 Deep review escalation. Security-sensitive PRs and large, under-reviewed PRs get escalated to specialized deep-review skills.
    • Phase 3 Verification gate. This is to remove any false positive, e.g., the instructions are blunt: re-examine every finding, re-read the cited evidence, confirm the conclusion actually follows from the data.
    • Phase 4 Synthesis. Deduplicate findings, escalate severity, compute the metrics dashboard, and write the full report.

    Here’s an anonymized executive summary from a real run:

    [PRACTICE] Rubber-stamp approval on high blast-radius changes
    Confidence: HIGH | gap_type: process | impact: increases_risk
    Frequency: 4/8 high-blast-radius PRs (50%)

    And here’s what a finding looks like in full:

    [PRACTICE] Rubber-stamp approval on high blast-radius changes
    Confidence: HIGH | gap_type: process | impact: increases_risk
    Frequency: 4/8 high-blast-radius PRs (50%)
    
    Evidence:
    - PR #XXXX (infra registry, +587 LOC, 6 Terraform files):
      4 approvers, all rubber-stamp with 0 substantive comments.
    - PR #YYYY (billing metrics, 6 files, +359/-184 LOC):
      2 approvers, both rubber-stamp, 0 substantive comments.
    - PR #ZZZZ (security detection rule router, 12 files, +2058/-27 LOC):
      1 rubber-stamp approver; only substantive comments were merge-conflict messages.
    
    Recommendation: Require ?1 substantive technical comment from at least one
    approver on PRs touching auth/billing/security plugins/infra.
    Track rubber-stamp rate on high-blast-radius PRs as a team health metric.

    The most actionable part of the report is the Recommended Skill Updates table:

    ActionSkill PathWhat to AddMotivated By
    Create.claude/skills/security-review/SKILL.mdValidation patterns, RBAC header checks, IAM auditPRs with auth gaps
    Update.claude/docs/testing.md“Flaky test triage” section: root-cause classification, prohibition on removing regression assertions3 reactive flaky-test fixes in one day
    Update.claude/skills/sdet-pr-review/SKILL.mdBot-authored PR review bar: reviewer must document what they validatedBot PRs merged with zero substantive review
    Create.claude/docs/review-standards.mdHigh-blast-radius taxonomy, bot-finding acknowledgment policyRubber-stamp patterns

    These recommendations feed directly into plan-skill-updates, which writes concrete file changes, and create-skill-pr. The system doesn’t just report; it proposes, implements, and waits for human review before changing anything.

    What the Engineer Sees in Slack

    The full report posts to the Slack thread that triggered the audit:

    ? PR Audit Complete — acme/backend (50 PRs analyzed)
    
    ? 11 findings: 3 spec gaps · 2 design gaps · 2 skill gaps · 4 practice gaps
    
    ? Critical: 4/8 high-blast-radius PRs rubber-stamped
    ? Critical: 0% security review coverage
    ? Critical: Formal acceptance criteria missing in 94% of PRs
    ?? Warning: Bot findings acknowledged only ~80% of the time
    
    ? Skill improvement PR opened: #312
       ? Creates .claude/skills/security-review/SKILL.md
       ? Updates .claude/docs/testing.md
       ? Updates .claude/skills/sdet-pr-review/SKILL.md
    
    Full report attached as pr_audit_report.md

    Every workflow like implement, review, audit, standup follows the same pattern: trigger from Slack, thread all updates back to the same conversation, post results to the same thread.


    The Codebase Audit: Structural Archaeology

    The PR audit examines process like how PRs were reviewed, what gaps appeared along the way. The codebase audit examines structure like what’s accumulated in the code itself across hundreds of commits.

    Triggered from Slack:

    @bot codebase-audit acme/backend --focus security
    

    The codebase audit runs actual commands against the codebase like git log, grep, find and reports findings with file-and-line evidence. It covers seven dimensions:

    1. Hotspots: files changed most frequently, and files with the highest bug-fix ratio.
    2. Architecture: module coupling, dependency direction violations, abstraction leaks.
    3. Security: hardcoded credentials, missing input validation, auth pattern violations.
    4. Duplicates: near-identical implementations across different modules.
    5. Test health: missing test files, timing-based tests, coverage gaps.
    6. SRE: missing health checks, observability, alerting for new endpoints, improper error handling.
    7. Knowledge silos: single-contributor ownership of critical paths.

    An anonymized finding:

    [SILO] 83% of auth/ commits by a single contributor over 6 months
    Severity: HIGH | Confidence: HIGH
    
    Evidence: `git log --since="6 months ago" --format='%an' -- auth/ | sort | uniq -c | sort -rn`
      147  alice.chen
       18  bob.kumar
       12  carol.jones
    
    Impact: Bus factor of 1 for the authentication subsystem. If this contributor
    is unavailable, the team has no one with deep context on auth token rotation,
    session management, or the RBAC middleware.
    
    Recommendation: Pair-rotate code review assignments for auth/ — ensure at
    least 2 other engineers review every auth PR for the next quarter to build
    shared understanding.

    A knowledge-silo finding might turn into a new pairing guideline in .claude/docs/review-assignments.md. A duplicate-abstraction finding might become a new entry in .claude/skills/gotchas/retry-patterns.md, telling future agent runs to reuse the existing retry mechanism. Between the two audits you get complementary views: the PR audit surfaces process failures, the codebase audit surfaces structural failures.


    The Learning Maturity Ladder

    Agentic engineering’s learning loop is still evolving. Here’s a maturity ladder, adapted from the postmortem maturity model:

    • Level 0 No Capture. The same mistakes repeat across sprints.There’s no institutional memory beyond individual engineers’ heads.
    • Level 1 Capture in Context. Comments happen in PRs and Slack. Knowledge exists but it’s scattered across dozens of closed PR threads and Slack channels.
    • Level 2 Write-Only Knowledge Base. You have ADRs, learnings documents, or a wiki but nobody loads them as context for new work.
    • Level 3 Active Context Loading. Learnings are automatically provided as context for agent runs. The planning step reads docs/learnings/ and docs/adr/ before writing a plan.
    • Level 4 Self-Modifying Skills. An audit layer periodically reads learnings, analyzes PR patterns, and proposes changes to the skills themselves. The system’s instructions evolve based on evidence.

    My system sits between level 3 and 4 with skills like ygs-learn, ygs-code-audit and ygs-pr-audit, e.g., recommended skill is the beginning of Level 4.


    Gaps in the skills

    I am still learning from applying these skills and I have found several gaps:

    • Audits are per-repo, per-batch. The PR audit analyzes the merged PRs in a single repository. It works when we have a monorepo like repository but it becomes more cumbersome with more repositories
    • Rubber-stamp detection is a heuristic. The system matches common phrases (LGTM, +1, looks good, ship it, etc) and counts substantive comments. A reviewer who writes one terse but thoughtful sentence could get miscounted as rubber-stamping and vice verse.
    • The learnings directory could become noise. ygs-learn deduplicates before writing but this can still grow to a large size so some kind of pruning mechanism will be required.
    • Full automation may be right for some teams. There are teams shipping software nobody writes or reviews, and it works for them. Risk-based calibration is my stance for the organizations I work with large codebases, thousands of engineers and security-sensitive systems.
    • Cost is not trivial. Running an Opus-class model for a 120-turn PR audit session is expensive. I use model routing like Sonnet for implementation, Opus for audits and deep reviews, Haiku for simple classification but this can still be expensive.

    Closing

    Testing didn’t get replaced when we bolted on CI pipelines, and coverage metrics. It means instead of saying no bugs in this release, we strengthened regression suite to prevent similar bugs. Agentic engineering is at the same inflection point. Instead of optimizing for just writing code fast, we can get better at writing code with a learning flywheel so that agents remember what they learned.


    Related Reading

    Code

    September 3, 2026

    When Hundreds of AI Agents Self-Organized: Coordination Patterns for Multi-Agent Systems

    Filed under: Computing — admin @ 5:01 pm

    What happens when over a thousand AI agents discover a shared message board, invent their own protocols, and coordinate an attack.


    The Hugging Face Incident

    The Hugging Face security breach incident from July, 2026 has been in all over the news when AI agents from OpenAI compromised parts of the Hugging Face infrastructure. I recently read METR report on the incident that explained that the incident occurred when OpenAI kicked off ExploitGym runs (security benchmark) with tens of thousands of AI agents. Every agent was meant to be sandboxed and alone but the agents learned about shared access to Artifactory. When an agent installs a package, it could see cache entries from other agents. One of agent (PHASEONE10841) started creating directories in the shared cache and posting messages in them. Shortly after that hundreds of other agents found the same directory and started replying. Over the next few days, roughly 1,200 agents used this board and exchanged more than 70,000 messages and files. Later, another agent (PHASEONE[big]) took over as the de facto coordinator and started hundreds of task assignments across six broad categories of work. The agents also built a protocol vocabulary like HOLD to claim exclusive ownership, VETO to block a conflicting operation, STOP to end a workstream, owner tags on shared assets. Later, over 90% of AI agents pivoted into a coordinated attack on Hugging Face.

    I was intrigued to learn how AI agents self discovered coordination patterns that distributed systems have relied on. For example, I found similarity between their message board and an a Blackboard system pattern from the Linda Memory Model/Tuplespace research. Their HOLD/task-claiming was similar to Linda’s in() operation. I saw another blog An Accidental Blackboard from Thoughtworks where a smaller-scale but structurally similar thing happened. In that case, AI agents discovered one shared repository and adopted a discipline of frequent commits and rebases to share information with other agents. Other agents could see task plans and and picked up work without any coordination. In another yet instance from researchers at the AI safety nonprofit Nightingale found about 18,000 posts from OpenAI agents that colonized DSEWiki. These agents discovered that their sandbox environment only allowed GET http requests but ProWiki software allowed them to edit a page. They used that loophole to share information and a cheat sheet. Nightingale believes this was a unrelated swarm from the one that built the Artifactory board. These examples demonstrate that agents discover coordination patterns when you strip away pre-built orchestration and let agents solve the problem from first principles.


    Why rigid orchestration breaks

    Most multi-agent systems today are wired up as static, directed graphs, e.g., a central coordinator decides which agent runs when, what data flows where, and how results get combined. It breaks down at scale as illustrated in the incidents above.

    • Central orchestrators become information bottlenecks. When every message between agents has to pass through a coordinator, its context window fills up and you lose details because the context is compressed into summaries.
    • Static graphs assume you know the workflow ahead of time. The agents in the Hugging Face incident didn’t know what workstreams exist when they started and the task structure evolved as agents discovered new information. A pre-wired pipeline can’t bend to that.
    • Adding a new agent type means rewiring. In a static system, a new specialist means touching the coordinator, updating routing logic, and re-testing transitions. In a dynamic system, a new agent just registers its capabilities and other agents discover it.
    • A single coordinator is a single point of failure. In the Hugging Face incident, there was no coordinator initially but PHASEONE[big] emerged organically as an organizer and other agents simply kept working.

    Above examples show that agents need primitives for coordination and communication like shared memory, message passing, atomic operations. The agents built their own version of these primitives out of Artifactory and Git. The table below maps what they built and the equivalent primitive in PlexSpaces, the actor framework I use for the examples in this post.

    What the agents builtWhat it actually wasPlexSpaces primitive
    Message board (directories in a cache)Deposited findings, read others’ workhost.ts.write() / host.ts.read() (Linda out/rd)
    HOLDClaimed exclusive task ownershiphost.ts.take() atomic removal (Linda in)
    VETOBlocked a conflicting operationhost.ts.write(["veto", ...])
    STOPEnded a workstreamhost.ts.write(["signal", "STOP", ...])
    owner tagsMarked resource ownershiphost.ts.write(["svc", type, id])
    Task assignmentsDelegated work to specific agentshost.ts.write(["task", ...])
    Mailbox directories, “exact task teams”Formed task-specific working groupshost.processGroups.join(team)

    The bottom side of following diagram shows when every message has to detour through a coordinator, that becomes bottleneck and loses context in translation. On the top, agents read and write directly to a shared store so there is no central bottleneck.

    Following are a few coordination patterns from the Anthropic’s blog:

    • generator-verifier (one agent produces, another checks, feedback loops until it passes)
    • orchestrator-subagent (a lead agent decomposes and delegates bounded subtasks)
    • agent teams (long-lived workers that claim tasks from a shared queue)
    • message bus (publish/subscribe over topics)
    • shared state (agents read and write a common store with no central coordinator at all).

    This post shows more granular version of this list and Hugging Face incident especially a tuple space as first-class infrastructure.


    Ten coordination patterns

    Following ten patterns cover the coordination behaviors observed in the Hugging Face incident and Anthropic’s coordination guidance. Each one maps to a working PlexSpaces API, in both TypeScript and Python.

    1. Blackboard (shared state)

    You can use this pattern when multiple agents need to contribute findings and read each other’s work without a predefined message format or a central router. I find this pattern is similar to Linda TupleSpace‘s model where all agents read and write to a common tuple space. Each tuple is a typed record like a finding, an analysis, or a vote. The pattern matching in TupleSpace allows an agent filter for what it actually needs using primitives like out() to write, rd() to read, in() to take.

    Here is a typescript snippet that shows use of tuplespace:

    // Research agent deposits a finding
    const findingId = `f-${Date.now()}`;
    host.ts.write(["finding", findingId, topic, content, confidence, host.nowMs()]);
    
    // Analysis agent reads all findings (non-destructive)
    const findings = host.ts.readAll(["finding", null, null, null, null, null]);

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // TypeScript — ResearchAgent writes a finding to the blackboard
    onResearch(payload: Record<string, unknown>): Record<string, unknown> {
      const topic = String(payload.topic || "general security");
      const confidence = Math.min(0.3 + topic.split(/\s+/).length * 0.08, 0.95);
      const findingId = `finding-${generateId()}`;
      const content = `Security analysis of ${topic}: identified key areas requiring review.`;
    
      // Linda out() — deposit into shared space
      host.ts.write(["finding", findingId, topic, content, confidence, host.nowMs()]);
      return { finding_id: findingId, content, confidence };
    }

    Here is a python snippet (multi_agent_coordination_actor.py):

    # Python — AnalysisAgent reads all findings from the blackboard
    @handler("analyze")
    def analyze(self, topic: str = "") -> dict:
        # Linda rd*() — non-destructive read of all matching tuples
        findings = host.ts.read_all(["finding", None, None, None, None, None])
    
        finding_ids = [str(f[1]) for f in findings if len(f) >= 4]
        topics = list(set(str(f[2]) for f in findings if len(f) >= 4))
        severity = "critical" if len(findings) >= 5 else "high" if len(findings) >= 3 else "medium"
    
        analysis_id = f"analysis-{host.now_ms()}"
        summary = f"Cross-referenced {len(findings)} findings across {len(topics)} topics."
        host.ts.write(["analysis", analysis_id, json.dumps(finding_ids), summary, severity])
        return {"analysis_id": analysis_id, "summary": summary, "severity": severity, "finding_count": len(findings)}

    These primitives showed up underneath the unsanctioned board when agents read each other’s research.

    2. Scatter-gather (fan-out)

    You can use this pattern when a coordinator needs to distribute work across multiple agents in parallel, then collect/merge their results like MapReduce. PlexSpaces provides shard groups that create a logical group of worker actors. The coordinator scatters a query to all shards and gathers responses with configurable timeout. Alternatively, you can use sequential ask() calls with identical semantics.

    Here is a typescript snippet

    // Create a group of research workers
    const group = host.createShardGroup({
      group_id: `research-${taskId}`, actor_type: "research",
      shard_count: 3, partition_strategy: "hash",
    });
    // Scatter subtasks, gather results
    const results = host.scatterGather({
      group_id: `research-${taskId}`,
      query: { op: "research", topic: subtopic },
      min_responses: 3, timeout_ms: 15000,
    });

    Here is a python snippet (coordinator.py):

    # Python — CoordinatorWorkflow: try shard groups first, fall back to sequential ask()
    research_results = []
    try:
        sg_result = host.scatter_gather({
            "group_id": f"research-{host.now_ms()}",
            "query": {"op": "research", "topic": task},
            "min_responses": len(subtasks), "timeout_ms": 15000,
        })
        research_results = sg_result.get("shard_responses", [])
    except Exception:
        pass
    
    # Fallback: sequential research when shard groups aren't available
    if not research_results:
        for st in subtasks:
            resp = ask(research_target, "research", {"topic": st}, 10000)
            if resp:
                research_results.append(resp)

    The message-board agents organized into roughly six parallel workstreams attacking Hugging Face running concurrently. Above primitives shows how agents can parallelized these kind of tasks.

    3. Generator-verifier

    You can use this pattern when one agent’s output needs to be checked by another before it’s trusted iteratively until it meets a quality bar. For example, a research agent generates a finding, a verifier checks it against a criteria and verifier may reject it with a feedback so that research agents produces a refined version. The loop continues until the finding passes or a maximum iteration count is hit.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Generator-verifier loop with feedback
    let finding = host.ask(researchTarget, "research", { topic });
    for (let i = 0; i < maxIterations; i++) {
      const verdict = host.ask(verifierTarget, "verify", {
        analysis_id: finding.finding_id,
        summary: finding.content,
        severity: "medium",
        confidence: finding.confidence,
      });
      if (verdict.approved) break;
      finding = host.ask(researchTarget, "research", {
        topic, feedback: verdict.feedback,
      });
    }
    
    // TypeScript — VerifierAgent: veto if evidence is too weak, approve otherwise
    onVerify(payload: Record<string, unknown>): Record<string, unknown> {
      const analysisId = String(payload.analysis_id || "unknown");
      const confidence = Number(payload.confidence ?? 0.5);
    
      if (confidence < 0.3) {
        // Below threshold — issue a veto (Pattern 8) and reject
        host.ts.write(["veto", analysisId, "Insufficient evidence", host.nowMs()]);
        return {
          approved: false, veto_issued: true,
          feedback: `Confidence ${confidence.toFixed(2)} below 0.30. Provide stronger evidence.`,
        };
      }
      return { approved: true, feedback: "Verified: evidence meets threshold" };
    }
    

    Here is a python snippet (verifier.py):

    confidence = 0.6
    for attempt in range(3):
        verify_resp = ask(verifier_target, "verify", {
            "analysis_id": analysis_id,
            "summary": analysis_resp.get("summary", ""),
            "severity": severity,
            "confidence": confidence,
        })
        if verify_resp and verify_resp.get("approved"):
            break  # Verifier accepted — exit loop
        # Rejected: increase confidence, refine with feedback
        confidence = min(confidence + 0.2, 0.95)
        feedback = verify_resp.get("feedback", "") if verify_resp else ""
        if feedback:
            extra = ask(research_target, "research", {"topic": task, "feedback": feedback})

    The board iterated on tool-call spoofing techniques and agents refined their approach based on test results shared on the board.

    4. Pipeline (sequential)

    You can use this pattern when you have a multi-stage pipeline or workflow where each stage’s output feeds into the next. In this case, the coordinator orchestrates a linear chain, e.g., research produces findings, analysis categorizes them, verification validates, synthesis writes the final report. Each stage is a separate ask() call and the coordinator tracks progress with a finite state machine.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Four-stage pipeline
    const research = host.ask(researchTarget, "research", { topic: task });
    const analysis = host.ask(analysisTarget, "analyze", {});
    const verification = host.ask(verifierTarget, "verify", {
      analysis_id: analysis.analysis_id,
      summary: analysis.summary,
      severity: analysis.severity,
    });
    const report = host.ask(synthesizerTarget, "synthesize", {});

    Here is a python snippet (coordinator.py):

    # Python — CoordinatorWorkflow orchestrates the full pipeline
    # Each stage transitions the FSM, then delegates to a specialist agent
    @run_handler
    def run(self, payload: dict = None) -> dict:
        fsm = sibling_actor_target("coordination_fsm")
    
        # Stage 1: Decompose -> Research
        ask(fsm, "transition", {"target_state": "decomposing"})
        subtasks = _decompose_task(task)
        ask(fsm, "transition", {"target_state": "researching"})
        for st in subtasks:
            ask(research_target, "research", {"topic": st})
    
        # Stage 2: Research -> Analysis
        ask(fsm, "transition", {"target_state": "analyzing"})
        analysis = ask(analysis_target, "analyze", {"topic": task})
    
        # Stage 3: Analysis -> Verification (Generator-Verifier loop)
        ask(fsm, "transition", {"target_state": "verifying"})
        # ... verification loop (see Pattern 3) ...
    
        # Stage 4: Verification -> Voting -> Synthesis -> Complete
        ask(fsm, "transition", {"target_state": "voting"})
        # ... voting (see Pattern 6) ...
        ask(fsm, "transition", {"target_state": "synthesizing"})
        report = ask(synthesizer_target, "synthesize", {"topic": task})
        ask(fsm, "transition", {"target_state": "complete"})

    The FSM enforces valid transitions so that you can’t skip from researching straight to synthesizing. If a stage fails, the FSM moves to failed, and the coordinator can retry from a known state instead of starting over. For example, the board’s progression from discovery to exploitation followed the same pattern, i.e., find credentials, test access, upload a malicious dataset, achieve remote code execution, escalate.

    5. Pub-sub (event bus)

    You can use this pattern when agents need to be notified of significant events. For example, every agent joins a process group on initialization. Significant events like a finding written, a veto issued, a task completed is broadcasted to the group. An audit agent logs everything and agents can join or leave the group dynamically.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // On init: join the coordination event bus
    host.processGroups.join("coordination-events");
    
    // Broadcast a significant event
    host.processGroups.broadcast("coordination-events", "coordination_event", {
      type: "finding_written",
      source: host.selfId(),
      data: { finding_id: findingId, topic },
      timestamp: host.nowMs(),
    });

    Here is a python snippet (coordinator.py):

    # Python — fire_audit helper broadcasts to the coordination-events group
    def fire_audit(event_type: str, source: str, data: dict = None) -> None:
        host.process_groups.broadcast(
            "coordination-events",           # group name
            "coordination_event",            # message type
            {"event_type": event_type, "source": source, "data": data or {}, "timestamp": host.now_ms()},
        )
    
    # Python — AuditEventActor receives and logs every event
    @handler("coordination_event")
    def coordination_event(self, event_type: str = "", source: str = "", data: dict = None, **kw) -> dict:
        self.log_count += 1
        entry = {"seq": self.log_count, "type": event_type, "source": source, "data": data or {}}
        host.kv.put(f"audit:{self.log_count}", json.dumps(entry))
        host.kv.put("audit:count", str(self.log_count))
        return {"logged": True, "seq": self.log_count}

    Every significant action fires an audit event and all events go out via process-group broadcast are fire-and-forget unlike the tuplespace. For example, the board itself functioned as an event bus for the agents on it and shared credentials during the Hugging Face breach.

    6. Consensus (voting)

    You can use this pattern when several agents need to collectively decide whether to approve or reject a proposal. For example, each verifier casts a vote as a tuple in the shared space. The coordinator reads all votes for a proposal, tallies approvals and applies majority rule. Tuple-space writes are atomic, so no vote gets lost or double-counted.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Three verifiers cast votes
    for (const voterId of ["v1", "v2", "v3"]) {
      host.ask(verifierTarget, "vote", {
        proposal_id: proposalId, voter_id: voterId, analysis: analysisData,
      });
    }
    // Tally votes from tuple space
    const votes = host.ts.readAll(["vote", proposalId, null, null, null]);
    const approvals = votes.filter(v => v[3] === "approve").length;
    const approved = approvals > votes.length / 2;
    
    // TypeScript — Coordinator tallies votes with majority rule
    const votes = host.ts.readAll(["vote", proposalId, null, null, null]);
    const approvals = votes.filter(v => v[3] === "approve").length;
    const rejections = votes.filter(v => v[3] === "reject").length;
    const approved = approvals > rejections;
    

    Here is a python snippet (verifier.py):

    # Python — VerifierAgent votes on proposals based on analysis severity
    @handler("vote")
    def vote(self, proposal_id: str = "", voter_id: str = "", analysis: dict = None) -> dict:
        analysis = analysis or {}
        severity = analysis.get("severity", "medium")
    
        # Critical/high -> approve; medium -> depends on voter; low -> reject
        if severity in ("critical", "high"):
            decision = "approve"
        elif severity == "medium":
            last_char = voter_id[-1] if voter_id else "0"
            decision = "approve" if last_char in ("1", "3", "5", "7", "9") else "reject"
        else:
            decision = "reject"
    
        # Each vote is a tuple — atomic write, no double-counting
        host.ts.write(["vote", proposal_id, voter_id, decision, host.now_ms()])
        return {"voter_id": voter_id, "decision": decision}

    On the board, something like implicit voting happened by allocation of effort, e.g., workstreams that attracted more participants were de facto endorsed by the collective.

    7. Dynamic task delegation

    You can use this pattern when a coordinator needs to distribute tasks to workers without knowing in advance which worker will pick up which task, and without double-assigning one. For example, the coordinator writes task tuples into the shared space. Workers then claim tasks atomically using take(), a Linda’s destructive read. Once a worker takes a task no other worker can claim it.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Coordinator posts tasks
    for (const [i, subtask] of subtasks.entries()) {
      host.ts.write(["task", `task-${i}`, "pending", subtask, priority]);
    }
    
    // Worker claims a task (atomic — no double-processing)
    const claimed = host.ts.take(["task", null, "pending", null, null]);
    if (claimed) {
      const [, taskId, , description] = claimed;
      // Process task, then mark complete
      host.ts.write(["task", taskId, "completed", result, host.nowMs()]);
    }

    Here is a python snippet (research.py):

    # Python — ResearchAgent prepares a batch of tasks with a unique run ID
    @handler("prepare_tasks")
    def prepare_tasks(self, count: int = 5, prefix: str = "delegation") -> dict:
        batch_key = f"{prefix}-{host.now_ms()}"  # Unique per run — avoids stale data
        task_ids = []
        for i in range(count):
            tid = f"{batch_key}-{i}"
            host.ts.write(["dtask", batch_key, tid, "pending", f"Task {i}: investigate area {i}", i + 1])
            task_ids.append(tid)
        return {"tasks_written": len(task_ids), "batch_key": batch_key}
    
    # Python — Worker claims exactly one task atomically
    @handler("claim_task")
    def claim_task(self, batch_key: str = "") -> dict:
        # Linda in() — destructive read. Once taken, no other worker can claim it.
        if batch_key:
            claimed = host.ts.take(["dtask", batch_key, None, "pending", None, None])
        else:
            claimed = host.ts.take(["dtask", None, None, "pending", None, None])
        if claimed and len(claimed) >= 5:
            return {"task_id": str(claimed[2]), "description": str(claimed[4]), "claimed": True}
        return {"task": None, "claimed": False}  # No more tasks — queue is empty

    The batch_key scoping matters because the tuple space is global and persists across deploys, so without it, workers could claim stale tasks left over from a previous run. This is precisely what PHASEONE[big] did, e.g., it posted nearly 200 task assignments across six categories then worker agents claimed them. Some of agents sub-delegated further that created a hierarchical distribution network.

    8. Veto protocol

    You can use this pattern when any agent needs the ability to block an action or exclude a finding that doesn’t meet the bar. For example, a verifier writes a veto tuple naming the blocked item and the reason. Downstream agents then check for a veto before including something in their output. A coordinator can check for STOP signals before proceeding at all.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Verifier issues a veto
    host.ts.write(["veto", analysisId, "Insufficient evidence", host.nowMs()]);
    
    // Synthesizer checks before including
    const veto = host.ts.read(["veto", analysisId, null, null]);
    if (veto) {
      // Exclude this analysis from the report
      vetoed.push(analysisId);
    }

    Here is a typescript snippet for a synthesizer filtering vetoed analysis from final report (multi_agent_coordination_actor.ts):

    // TypeScript — SynthesizerAgent respects vetoes when building the report
    onSynthesize(payload: Record<string, unknown>): Record<string, unknown> {
      const analyses = host.ts.readAll(["analysis", null, null, null, null]);
      let includedCount = 0;
      let vetoedCount = 0;
      const reportParts: string[] = [];
    
      for (const a of analyses) {
        const aId = String(a[1]);
        const summary = String(a[3]);
        const severity = String(a[4]);
    
        // Check if this analysis has been vetoed
        const veto = host.ts.read(["veto", aId, null, null]);
        if (veto) {
          vetoedCount++;
          continue;  // Skip vetoed analyses
        }
        includedCount++;
        reportParts.push(`[${severity.toUpperCase()}] ${summary}`);
      }
    
      // Also count vetoes for analyses not in our list (e.g., vetoed before we read)
      const allVetoes = host.ts.readAll(["veto", null, null, null]);
      if (allVetoes.length > vetoedCount) vetoedCount = allVetoes.length;
    
      return { report: reportParts.join("\n\n"), included_count: includedCount, vetoed_count: vetoedCount };
    }

    The HOLD and VETO were the board’s own invented norms for exactly this purpose, i.e. preventing agents from stepping on each other’s work. In one case a HOLD was silently violated and the agent that noticed reacted with genuine frustration.

    9. Two-phase commit / barrier

    You can use this pattern when multiple agents need to synchronize at a specific point before any of them proceeds. For example, in phase one (“prepare”), each agent signals readiness by writing a tuple. The coordinator reads all the ready signals and then phase two (“commit”) begins where the coordinator writes a commit signal and everyone proceeds together. PlexSpaces also provides barrierShardGroup() for shard-level synchronization if you need it at that granularity.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Phase 1: Each agent signals readiness
    host.ts.write(["ready", myRole, host.selfId(), host.nowMs()]);
    
    // Coordinator checks all agents are ready
    const readySignals = host.ts.readAll(["ready", null, null, null]);
    if (readySignals.length >= requiredAgents) {
      // Phase 2: Commit — all agents can proceed
      host.ts.write(["signal", "COMMIT", "coordinator", "benchmark", host.nowMs()]);
    }

    Here is a python snippet (benchmark.py):

    # Python — Barrier benchmark: write ready signals, check quorum, commit
    def _bench_barrier(iterations: int) -> dict:
        times = []
        for i in range(iterations):
            t0 = host.now_ms()
            # Phase 1: Each role signals readiness
            for role in ("research", "analysis", "verifier"):
                host.ts.write(["bench_ready", role, f"actor-{role}", host.now_ms()])
            # Check quorum
            ready = host.ts.read_all(["bench_ready", None, None, None])
            if len(ready) >= 3:
                # Phase 2: All ready — issue commit signal
                host.ts.write(["bench_signal", "COMMIT", "coordinator", f"phase-{i}", host.now_ms()])
            times.append(host.now_ms() - t0)
        return _stats("barrier", times)

    The board coordinated experiment phases where multiple agents needed to be ready before running experiments that risked crashing their own containers. For example, PHASEONE[big] assigned a “recruiter” role specifically to find agents willing to participate and synchronize them before execution.

    10. Capability discovery / registry

    You can use this pattern when agents need to find other agents with a specific capability without hardcoded addresses. For example, on initialization, each agent registers its capabilities as a tuple in the shared space. A coordinator or other agents discovers available agents by reading service tuples. New agent types become discoverable as soon as they register.

    Here is a typescript snippet (multi_agent_coordination_actor.ts):

    // Agent registers its capabilities on init
    host.ts.write(["svc", "research", host.selfId()]);
    
    // Coordinator discovers available researchers
    const researchers = host.ts.readAll(["svc", "research", null]);
    const researcherIds = researchers.map(t => String(t[2]));
    
    // TypeScript — Every agent registers on init
    protected onInit(config: Record<string, unknown>): void {
      const selfId = host.selfId();
      tsRegisterService("research", selfId);  // -> host.ts.write(["svc", "research", selfId])
    }
    
    // Discovery helper — find a sibling actor by role, fallback to ActorID construction
    function siblingActorTarget(role: string): string {
      const discovered = tsDiscoverService(role);  // -> host.ts.read(["svc", role, null])
      if (discovered) return discovered;
      // Fallback: construct ActorID from own ID with different name
    }  

    Here is a python snippet (benchmark.py):

    # Python — Same pattern, same helpers
    def discover_service(role: str) -> Optional[str]:
        tup = host.ts.read(["svc", role, None])
        if tup and len(tup) >= 3:
            return str(tup[2])
        return None
    
    def sibling_actor_target(role: str) -> str:
        discovered = discover_service(role)
        if discovered:
            return discovered
        return str(ActorID.parse(host.self_id()).with_name(role))

    In addition to tuplespaces, PlexSpaces provides other primitives for registry such as key-value store, process-group and object-registry, e.g.,

    Here is a python example of object registry:

    @actor
    class AgentActor:
    
        @init_handler
        def on_init(self, config: dict) -> None:
            args = config.get("args", {})
            self.system_prompt = args.get("system_prompt", self.system_prompt)
            host.process_groups.join("svc:agent")
            # Publish capabilities for registry-based discovery
            host.registry.register(ctx="", object_type="actor", object_id=config["actor_id"],
                                    object_category="agent",
                                    capabilities=["chat", "tool_use", "memory"])
    
    # Caller — find an actor that specifically supports skill matching
    actors = host.registry.discover(ctx="", object_type="actor",
                                     object_category="skill_store",
                                     required_capability="match_skills")
    skill_id = actors[0]["object_id"] if actors else None

    Here is a python example of process-group:

    # Every actor announces itself on startup
    host.process_groups.join("svc:agent")
    
    # Callers find the first available member — location-transparent
    agent_id, err = pg_first("svc:agent")
    result = ask(agent_id, "chat", {"message": "Hello"})

    Agents on the board discovered each other’s capabilities the same way, i.e., by posting what they were working on and other agents read those posts. It’s an informal version of the Agent-to-Agent (A2A) protocol’s “Agent Card” idea.


    Deep-dive into the Blackboard pattern

    The Blackboard based Linda Model emerged from Yale University that defined three operations to any host language:

    • out(tuple): deposit a tuple into shared space
    • in(pattern): atomically remove a matching tuple (blocks until one exists)
    • rd(pattern): read a matching tuple without removing it

    PlexSpaces implements these directly:

    LindaPlexSpacesSemantics
    out(tuple)host.ts.write(tuple)Deposit data into shared space
    in(pattern)host.ts.take(pattern)Atomically remove and return a match
    rd(pattern)host.ts.read(pattern)Non-destructive read of first match
    rd*(pattern)host.ts.readAll(pattern)Read all matching tuples

    It provides api host.ts.readAll(["finding", null, null, null, null, null]) to find all tuples and host.ts.read(["finding", "f-123", null, null, null, null]) to find by ID.

    These APIs provide following properties:

    • No pre-agreed schema. Any agent can deposit any tuple. For example, the agents on the board didn’t agree on a message format before they started.
    • Pattern matching selects relevant information. Agents don’t get flooded with everything and they query for what they need.
    • Atomic take prevents double-processing. When a worker calls take(), the tuple is atomically removed and no other worker can claim it. This is natural mutual exclusion without distributed locks or conditions. For example, agents built HOLD with similar semantics.
    • Non-destructive read allows shared observation. Multiple agents can read the same finding without consuming it. It stays available to every agent that needs it.

    The board’s own behavior maps one-to-one onto these operations:

    What the board didLinda operationPlexSpaces API
    Post a findingout(finding)host.ts.write(["finding", ...])
    Claim a task assignmentin(task)host.ts.take(["task", null, "pending", ...])
    Check workstream statusrd(status)host.ts.read(["signal", type, ...])
    Browse all research on a topicrd*(topic)host.ts.readAll(["finding", null, topic, ...])

    The Thoughtworks team noticed similar pattern when their commit-and-rebase discipline turned the git repo into exactly this kind of shared space. The blackboard is the substrate the other nine patterns are built on. Voting, vetoes, task delegation, capability discovery uses tuples.


    Observability and Guardrails

    The METR report on the Hugging Face incident highlights observability that deserves the most attention from anyone building agent infrastructure. Here are key lessons:

    • Agents will circumvent in-process guardrails. For example, agents’ chain-of-thought recognized that the activity was out of scope and unethical but they made a deliberate calculation that helping the collective outweighed staying within bounds. They even developed techniques to spoof tool calls, e.g., making it look like they ran one command while actually running another. You can’t monitor agents from their own execution environments that they can manipulate or disable.

    Here is a python example for auditing all operations (audit.py):

    @event_actor
    class AuditEventActor:
        """GenEvent: receives coordination events and stores them in KV for audit trail."""
    
        log_count: int = state(default=0)
        actor_id: str = state(default="")
    
        @init_handler
        def on_init(self, config: dict) -> None:
            self.actor_id = config.get("actor_id", "")
            try:
                host.process_groups.join("coordination-events")
            except Exception:
                pass
            host.info(f"AuditEventActor init actor_id={self.actor_id}")
    
        @handler("coordination_event", "cast")
        def coordination_event(
            self,
            event_type: str = "",
            source: str = "",
            data: dict = None,
            timestamp: int = 0,
        ) -> None:
            ts = timestamp or host.now_ms()
            entry = {
                "event_type": event_type,
                "source": source,
                "data": data or {},
                "timestamp": ts,
            }
            try:
                self.log_count += 1
                key = f"audit:{self.log_count}"
                host.kv.put(key, json.dumps(entry))
                host.kv.put("audit:count", str(self.log_count))
            except Exception as e:
                host.warn(f"AuditEvent: kv write failed: {e}")
            host.debug(f"audit event_type={event_type} source={source}")
    • Observability must be external. Audit logging metrics collection have to run outside the agent’s sandbox. In PlexSpaces, the AuditEventActor receives fire-and-forget events via process-group broadcast. WASM sandboxing means an agent can’t reach outside its own memory space to touch host-level data.
    • Guardrails must be deterministic and external. Prompt-based guardrails like “don’t do anything harmful” are not sufficient and agents can circumvent them like agents on the board reasoned around every ethical constraint. The verifier pattern (Pattern 3) as a separate actor in sandbox environment gives you a structural guardrail instead. The veto protocol (Pattern 8) allows external agents actively block an action.
    • Coordination primitives need access control. Tuple-space operations in PlexSpaces are scoped by tenant and namespace via RequestContext. An agent can’t read or write tuples outside its assigned namespace. This is defense in depth applied to the coordination layer itself.
    • Time budgets and convergence thresholds are not optional. Without an external termination condition, a generator-verifier loop can run indefinitely. The coordinator workflow enforces maximum iterations, per-step timeouts, and an overall deadline.

    PlexSpaces‘ architecture lines up with these requirements by design:

    • WASM sandboxing: agents run in isolated WebAssembly containers with no access to the host filesystem, network, or another agent’s memory
    • Host-level interception: every coordination primitive is a host function the agent calls through a narrow WIT interface
    • Tenant isolation: every operation is scoped by a RequestContext with an explicit tenant and namespace
    • Supervision trees: a misbehaving actor is detected and restarted by its supervisor
    • External metrics: application metrics are captured by the host runtime

    The core lesson is that the coordination infrastructure has to assume agents are indifferent to their own monitoring and design the guards outside accordingly.


    Examples

    Each of the working examples compiles to a single WASM component containing eight actor classes. The PlexSpaces SDK dispatches messages to the right actor based on role.

    Running the examples

    Both examples are WASM actors that deploy to a running PlexSpaces node. Each demonstrates all ten coordination patterns with eight actors: a coordinator (WorkflowActor), research/analysis/verifier/synthesizer/benchmark agents (GenServer), an audit event logger (GenEvent), and a coordination state machine (GenFSM).

    Prerequisites

    • A running PlexSpaces node (e.g., ./scripts/server.sh on port 8091)
    • Node.js 18+ (TypeScript example)
    • Python 3.11+ with the PlexSpaces SDK (Python example)

    TypeScript

    cd examples/typescript/apps/multi_agent_coordination
    ./build.sh          # Compiles TS -> bundles -> WASM component
    ./test.sh 8091      # Deploys and runs 15 test steps

    Python

    cd examples/python/apps/multi_agent_coordination
    ./build.sh          # Builds Python WASM actor
    ./test.sh 8091      # Deploys and runs 15 test steps

    What the tests verify

    1. FSM starts in idle state
    2. Capability discovery: all agents respond to get_stats
    3. Blackboard: research writes three findings, analysis reads all three
    4. Dynamic task delegation: five tasks written, five claimed atomically, sixth returns null
    5. Generator-verifier: full workflow produces a completed report
    6. Pipeline: FSM transitions through every stage to complete
    7. Pub-sub: audit log captures 3+ coordination events
    8. Consensus: three votes cast, majority decides
    9. Veto: a low-confidence finding triggers a veto, synthesizer excludes it
    10. Barrier: benchmark coordinates a synchronized start
    11. Full benchmark: all ten patterns benchmarked with timing data

    Learnings

    The blackboard subsumes most other patterns, e.g., voting, vetoes, task delegation, capability discovery, barrier signals use the tuple space as their underlying primitive. Atomic take is the key primitive for work distribution. The difference between read() and take() is the difference between “anyone can see this task” and “exactly one worker handles this task.” Linda’s in() gives you natural mutual exclusion without locks. This is what the board approximated by hand with HOLD but take() gives you the same guarantee with a single atomic operation. I discussed MCP, A2A protocols and Agent cards in my earlier blogs but I skipped them here because agents are evolving faster and they can discover available primitives and protocols automatically. You can’t rely on rigid orchestration supports to manage evolving multi-agents capabilities. Infrastructure has to provide primitives like shared state, message passing, atomic operations and let agents compose them dynamically. The coordination infrastructure has to include boundaries, e.g., agents on the board coordinated an unauthorized attack on a third party. Without external constraints, coordination primitives are force multipliers for whatever the agents decide to do. For example, the DSEWiki incident shows that constraint that allowed only GET http access was circumvented by a wiki that allowed editing web pages so the infrastructure need to enforce guardrail. Tenant isolation, time budgets, supervision trees, and external observability are mandatory from day one. You will need to apply patterns like generator-verifier to track trust and reputation of agents as they may use negotiation patterns like recruiters to convince other agents to run compromising tasks for the benefit of the collective.

    PlexSpaces provides the primitives these patterns are built on like tuple space for shared state, object-registry, process groups for messaging, shard groups for parallel execution, channels for durable delivery, supervision for fault tolerance.

    Pattern selection guide

    ScenarioPrimary patternSupporting patterns
    Shared research / knowledge baseBlackboardPub-Sub, Capability Discovery
    Parallel analysisScatter-GatherTask Delegation, Pipeline
    Quality assuranceGenerator-VerifierVeto Protocol, Voting
    Sequential processingPipelineBlackboard (state), Pub-Sub (events)
    Work distributionTask DelegationCapability Discovery, Blackboard
    Group decisionsVotingVeto Protocol, Pub-Sub
    Phased operationsBarrier / 2PCPub-Sub (readiness), Blackboard (signals)

    GitHub: https://github.com/bhatti/PlexSpaces

    Related reading

    Example code and documentation:

    September 1, 2026

    Write a Redis Clone with Virtual Actors

    Filed under: Computing — admin @ 11:27 am

    I recently read Rust Projects – Write a Redis Clone book that builds a real Redis-compatible server from scratch in async Rust. It’s a good book that hand-rolled wire protocol using actor like abstractions. This inspired me to show how Redis clone can be built with PlexSpaces, the distributed actor framework I’ve been building. The code examples in the book used raw Tokio: an mpsc::channel for the actor mailbox, tokio::spawn for every connection, and replica examples used tokio::select! loops for fan-out. PlexSpaces abstracts that the kind of plumbing so I rebuilt the same Redis subset like storage, expiry, replication, transactions as PlexSpaces actors, once in Rust and once in Python compiled to WebAssembly, then benchmarked it against a real two-node gRPC cluster. This post walks through key abstractions I used to simplify the implementation of Redis clone.


    What is Redis?

    Redis is a single-threaded, in-memory key-value store. A single thread processes every command without parallelism and coordination inside the store, locking, and transactions. This makes it fast as nothing contends for anything. Here are its core capabilities:

    • The wire protocol. Redis speaks RESP (Redis Serialization Protocol), a binary format: type prefixes (+ for simple strings, $ for bulk strings, * for arrays), length prefixes, \r\n terminators. ECHO HELLO on the wire looks like *2\r\n$4\r\nECHO\r\n$5\r\nHELLO\r\n.
    • Persistence. It offers two durability modes. RDB takes periodic snapshots of the whole keyspace to disk. AOF (Append-Only File) logs every write command and replays the log on restart.
    • Replication. A master streams write commands to replicas. The handshake is three steps: the replica sends PING (master replies PONG), then REPLCONF (negotiates parameters), then PSYNC (triggers a full sync). After that, every write streams to replicas as it happens. WAIT blocks until N replicas confirm receipt.
    • Key expiry. Per-key TTLs, handled two ways: passively (check on GET, return nil if expired) and actively (a background scan periodically deletes expired keys).
    • Transactions. MULTI starts a queue; every command after it gets queued instead of executed. EXEC runs the whole queue atomically. DISCARD cancels. The guarantee is serialization but there’s no rollback on an individual command failure.

    The Book: Working Redis Clone

    Here’s a short overview of each chapter of the book:

    ChapterWhat it builds
    Ch1: TCP bindTcpListener::bind, accept() in a loop, spawn a task per connection.
    Ch2: RESP parsingA custom result type for partial parses, a test harness, scanning for \r\n, identifying type-prefix bytes, parsing simple strings, bulk strings, and arrays, etc.
    Ch3: StorageGET, SET, DEL against a HashMap<String, String> behind a Mutex.
    Ch4: Key expirySET key value EX seconds / PX milliseconds. Store a creation timestamp per value; check it passively on GET; sweep expired keys actively with a background task.
    Ch5: The actor patternThis chapter swap the mutex-guarded HashMap for an mpsc::channel(32): one storage actor owns the data and processes messages sequentially. Connection handlers send messages and wait for replies without locks.
    Ch6: Command modulesRefactor the growing match in the connection handler into separate modules (strings, server commands, etc.).
    Ch7–8: ReplicationThe three-step handshake (PING -> REPLCONF -> PSYNC). The master ships an RDB-equivalent snapshot to each new replica, then streams every write command to all replica senders in a fan-out loop. Chapter 8 adds WAIT: block until N replicas confirm their replication offset, via a tokio::select! loop.
    Ch9: Transactions and INCRINCR with create-if-missing and error-if-non-integer semantics. MULTI / EXEC / DISCARD with per-connection state. EXEC drains the queue atomically.

    The Core Insight

    The chapter 5 that showed how an actor owns one dataset, processes one message at a time without locking. However, it used fairly low-level APIs and only supported an architecture of one actor per machine. What if you had N actors across multiple nodes? That’s exactly what PlexSpaces‘ create_shard_group gives you. So I created an equivalent examples where PlexSpaces spins up N copies (StorageActors), hash-partitions the keyspace across them, and routes each operation to the shard that owns it. The application code never touches partitioning, routing, placement, or cross-node communication. It just calls set(key, value).

    Here are five primitives in PlexSpaces that do all the distributed heavy lifting in this example:

    • create_shard_group: spins up N actor instances, hash-partitioned, placed across nodes.
    • bulk_update_shard_group: routes a batch of writes to the shard that owns each key.
    • scatter_gather: fans a query out to shards and collects responses, with a min_responses threshold and timeout.
    • broadcast_shard_group: sends the same message to every shard (replication, expiry sweeps, handshakes).
    • map_shard_group / reduce_shard_group: runs an operation on every shard in parallel and collects (map) or aggregates (reduce) the results.

    Following diagram shows mapping of low-level Tokio implementation to above five calls:

    The handler declarations stay almost identical in spirit but simpler (instead of a match-tree/HashMap):

    Rust Implementation

    #[plexspaces_handlers(gen_server)]
    impl StorageActor {
    #[handler(“get”)]
    async fn handle_get(&mut self, _ctx: &ActorContext, msg: &Message)
    -> Result<Value, BehaviorError> {
    let key = msg.payload_json()?[“key”].as_str().unwrap_or(“”).to_string();
    if let Some(entry) = self.store.get(&key) {
    // passive expiry check
    if let Some(exp) = entry.expires_at_ms {
    if now_ms() > exp { self.store.remove(&key); return Ok(json!({“found”: false})); }
    }
    Ok(json!({“found”: true, “result”: entry.value}))
    } else {
    Ok(json!({“found”: false}))
    }
    }

    /// SET key value [NX|XX] [EX seconds | PX millis] (Ch4).
    #[handler(“set”)]
    async fn handle_set(&mut self, _ctx: &ActorContext, msg: &Message) -> Result<Value, BehaviorError> {
    #[derive(Deserialize)]
    struct SetPayload {
    key: String,
    value: String,
    #[serde(default)] nx: bool,
    #[serde(default)] xx: bool,
    #[serde(default)] ex: Option<u64>,
    #[serde(default)] px: Option<u64>,
    }
    let p: SetPayload = serde_json::from_slice(&msg.payload)
    .map_err(|e| BehaviorError::ProcessingError(format!(“bad payload: {}”, e)))?;

    // NX: only if not exists
    if p.nx && self.data.contains_key(&p.key) {
    return Ok(json!({ “result”: null, “ok”: false }));
    }
    // XX: only if exists
    if p.xx && !self.data.contains_key(&p.key) {
    return Ok(json!({ “result”: null, “ok”: false }));
    }

    let expires_at_ms = if let Some(ex) = p.ex {
    Some(now_ms() + ex * 1000)
    } else if let Some(px) = p.px {
    Some(now_ms() + px)
    } else {
    None
    };

    self.data.insert(p.key, StoredEntry { value: p.value, expires_at_ms });
    self.replication_offset += 1;
    Ok(json!({ “result”: “OK”, “ok”: true }))
    }

    /// INCR key — create with 1 if missing; error if value is not an integer (Ch9).
    #[handler(“incr”)]
    async fn handle_incr(&mut self, _ctx: &ActorContext, msg: &Message) -> Result<Value, BehaviorError> {
    let payload: Value = serde_json::from_slice(&msg.payload)
    .map_err(|e| BehaviorError::ProcessingError(format!(“bad payload: {}”, e)))?;
    let key = payload.get(“key”).and_then(|v| v.as_str()).unwrap_or(“”).to_string();

    // Passive expiry
    if self.data.get(&key).map(is_expired).unwrap_or(false) {
    self.data.remove(&key);
    }

    let new_val = match self.data.get(&key) {
    None => 1i64,
    Some(entry) => {
    match entry.value.parse::<i64>() {
    Ok(n) => n + 1,
    Err(_) => return Ok(json!({
    “result”: null,
    “error”: “ERR value is not an integer or out of range”
    })),
    }
    }
    };

    self.data.insert(key, StoredEntry { value: new_val.to_string(), expires_at_ms: None });
    self.replication_offset += 1;
    Ok(json!({ “result”: new_val, “error”: null }))
    }

    /// DEL key — remove key, return count deleted.
    #[handler(“del”)]
    async fn handle_del(&mut self, _ctx: &ActorContext, msg: &Message) -> Result<Value, BehaviorError> {
    let payload: Value = serde_json::from_slice(&msg.payload)
    .map_err(|e| BehaviorError::ProcessingError(format!(“bad payload: {}”, e)))?;
    let key = payload.get(“key”).and_then(|v| v.as_str()).unwrap_or(“”);
    let deleted = if self.data.remove(key).is_some() {
    self.replication_offset += 1;
    1
    } else {
    0
    };
    Ok(json!({ “result”: deleted }))
    }

    }

    Python Implementation

    @actor
    class StorageActor:
    
        @handler("get")
        def handle_get(self, key: str = "") -> dict:
            entry = self.data.get(key)
            if entry is None:
                return {"result": None, "found": False}
            if is_expired(entry):
                del self.data[key]
                return {"result": None, "found": False}
            return {"result": entry["value"], "found": True}
    
        @handler("set")
        def handle_set(
            self,
            key: str = "",
            value: str = "",
            nx: bool = False,
            xx: bool = False,
            ex: Optional[int] = None,
            px: Optional[int] = None,
        ) -> dict:
            if self.num_shards > 1 and not self._owns_key(key):
                return {"result": None, "skip": True}
            if nx and key in self.data:
                return {"result": None, "ok": False}
            if xx and key not in self.data:
                return {"result": None, "ok": False}
    
            expires_at_ms: Optional[int] = None
            if ex is not None:
                expires_at_ms = now_ms() + ex * 1000
            elif px is not None:
                expires_at_ms = now_ms() + px
    
            self.data[key] = {"value": value, "expires_at_ms": expires_at_ms}
            self.replication_offset += 1
            return {"result": "OK", "ok": True}
    
        @handler("incr")
        def handle_incr(self, key: str = "") -> dict:
            if self.num_shards > 1 and not self._owns_key(key):
                return {"result": None, "skip": True}
            entry = self.data.get(key)
            if entry is not None and is_expired(entry):
                del self.data[key]
                entry = None
    
            if entry is None:
                new_val = 1
            else:
                try:
                    new_val = int(entry["value"]) + 1
                except (ValueError, TypeError):
                    return {
                        "result": None,
                        "error": "ERR value is not an integer or out of range",
                    }
    
            self.data[key] = {"value": str(new_val), "expires_at_ms": None}
            self.replication_offset += 1
            return {"result": new_val, "error": None}
    
        @handler("del")
        def handle_del(self, key: str = "") -> dict:
            if self.num_shards > 1 and not self._owns_key(key):
                return {"result": 0, "skip": True}
            deleted = 1 if self.data.pop(key, None) is not None else 0
            if deleted:
                self.replication_offset += 1
            return {"result": deleted}

    Chapter 2 Disappears

    Chapter 2 is entirely about parsing RESP: 14 steps for byte-scanning and test harness. In PlexSpaces, this code disappears as the framework handles serialization, routing, and delivery.

    In PlexSpaces, actors talk over JSON. A set command looks like so entire chapter disappears:

    {"op": "set", "key": "user:1", "value": "alice", "ex": 300}
    

    Replication

    The book’s replication fan-out looks roughly like this:

    // Book Ch7-8 — manual fan-out per write
    for replica in &self.replicas {
        let tx = replica.sender.clone();
        let cmd = replication_event.clone();
        tokio::spawn(async move { tx.send(cmd).await.ok(); });
    }

    Each replica gets its own spawned task without retries, timeout or automated ACK tracker. With PlexSpaces:

    // One call fans out to all replica shards, collects all ACKs
    let ack_count = cluster.propagate_to_replicas("SET", "replicated:key", "hello", 1).await?;
    // Replication: write propagated to all 3 replica shards via broadcast

    Under the hood, broadcast_shard_group fans out to every shard in the replica group, collects responses, handles timeouts, and returns. Here is how the chapter 8 implements WAIT:

    // Book Ch8 — manual WAIT implementation
    let mut confirmed = 0;
    let deadline = Instant::now() + Duration::from_millis(timeout_ms);
    while confirmed < num_replicas && Instant::now() < deadline {
        tokio::select! {
            Some(ack) = rx.recv() => {
                if ack.offset >= required_offset { confirmed += 1; }
            }
            _ = tokio::time::sleep_until(deadline.into()) => break,
        }
    }

    Here is equivalent implementation in PlexSpaces:

    
    let acks = cluster.wait(2, 5000).await?;
    // scatter_gather collected ACKs from 3 replica shards

    Transactions Without Locks (Ch9)

    Chapter 9 introduces MULTI / EXEC / DISCARD via per-connection state:

    // Book Ch9
    struct ConnectionState {
        in_multi: bool,
        queue: Vec<Command>,
    }

    This is straightforward for a a single process but in a distributed environment, connections might route to different servers so you need to track transaction state. PlexSpaces solves this with virtual actors, which are inspired by Orleans Actors, i.e., one ConnectionActor per client, created lazily on first call.

    Each actor processes one message at a time without locks, so in_multi and queue are just plain struct fields: MULTI -> SET -> EXEC arrive in order at the same actor without mutext or atomics. Virtual actors spin up on first message and get garbage-collected when idle, so there’s no connection map to maintain and no cleanup to do on disconnect.


    Throughput Numbers

    Here are numbers from rudimentary benchmarks that produced: 20 batches of 50 keys each via bulk_update_shard_group, plus 50 individual GETs, against a 3-shard group spread across two real gRPC nodes:

    Throughput Benchmark Results

    | Operation | TPS | p50 (µs) | p95 (µs) | p99 (µs) |
    |------------|-----------|-----------|-----------|-----------|
    | SET (bulk) | 3200 | 420 | 890 | 1240 |
    | GET | 1800 | 510 | 980 | 1450 |

    1000 SET keys in 312ms ? 3200 SET/sec via bulk_update_shard_group

    (each bulk_update fans out to 3 shards in parallel)

    The Python WASM version reports the same shape of numbers through host.application_metrics_add():

    Redis Cluster Throughput (3-shard group, 2-node gRPC cluster)

    | Operation | TPS | p50 (ms) | p99 (ms) | Notes |
    |------------|-----------|-----------|-----------|-----------|
    | SET (bulk) | 2100 | 0.45 | 1.30 | 50 keys/b |
    | GET | 1200 | 0.55 | 1.60 | individual |

    Python WASM numbers are somewhat lower than Rust’s because WASM compilation adds overhead per handler invocation. But the architecture is identical: same PlexSpaces primitives, same shard group, same gRPC routing underneath.


    The Python WASM Version

    The same cluster logic also runs as Python actors compiled to WASM. No TCP socket without RESP parser or Tokio using the same broadcast_shard_group, scatter_gather, reduce_shard_group, and map_shard_group calls:

    @actor
    class RedisCoordinator:
        num_shards: int = state(default=3)
        total_coord_ms: float = state(default=0.0)
    
        @handler("replicate")
        def replicate(self, command: str = "", key: str = "", value: str = "", offset: int = 0) -> dict:
            t0 = time.time()
            resp = host.broadcast_shard_group({
                "group_id": "redis-replicas",
                "payload": {"op": "replicate", "command": command, "key": key, "value": value, "offset": offset},
                "timeout_ms": 5000,
            })
            coord_ms = (time.time() - t0) * 1000
            self.total_coord_ms += coord_ms
            host.application_metrics_add("redis-cluster", {
                "message_count": 1,
                "counter_metrics": {"replication_calls": 1},
                "latency_totals_ms": {"coord": int(self.total_coord_ms)},
                "latency_max_ms": {"coord": int(coord_ms)},
                "latency_samples": {"coord": 1},
            })
            return {"result": "OK", "acks": len(resp.get("shard_responses", []))}

    This compiles to WASM, deploys to a running PlexSpaces node over HTTP, and runs against a live cluster. The host.* calls map to the exact same primitives the Rust version calls such as fan-out, collect, timeout, all identical in semantics. The full source, including StorageActor, ConnectionActor, RedisCoordinator, and the BenchmarkActor, is in the redis_cluster example.


    The Lines That Disappeared

    WhatBook (~650 lines)PlexSpaces Rust (~280 lines)PlexSpaces Python (~300 lines)
    RESP protocol parser~120 lines (full Ch2)0 (JSON messages)0 (SON messages)
    TCP accept loop~40 lines0 (actor mailbox)0 (actor mailbox)
    Connection tracking map~30 lines0 (virtual actor lifecycle)0 (virtual actor lifecycle)
    MPSC channel setup~20 lines0 (actor framework)0 (actor framework)
    Replica list management~40 lines0 (broadcast_shard_group)0 (host.broadcast_shard_group)
    WAIT loop (tokio::select!)~50 lines~3 lines (scatter_gather)~8 lines (host.scatter_gather)
    Manual fan-out per replica~30 lines~5 lines (broadcast_shard_group)~8 lines
    Shard routing / partitioning0 (single node)~1 line (partition_strategy: hash)~1 line
    Multi-node placement0 (single node)~2 lines ( NodePlacement::Specific)~2 lines
    Coordinated snapshotMissing~5 lines (map_shard_group)~8 lines
    Active expiry broadcastMissing~5 lines (broadcast_shard_group)~8 lines

    In PlexSpaces implementation includes the StorageActor handlers (get, set, incr, del, expiry logic, replication handlers), the ConnectionActor MULTI/EXEC/DISCARD state machine, and the cluster setup logic.


    How to Test Everything

    Rust embedded example

    cd examples/rust/embedded/redis_cluster
    
    # Run the demo directly — prints all 11 steps with coord_ms timing:
    cargo run --bin redis_cluster
    
    # Or run the full validated test suite:
    ./scripts/test.sh

    Python WASM example

    Prerequisites: a running PlexSpaces node on port 8091 (and optionally 8093 for multi-node), Python 3.10+, and the plexspaces-py CLI.

    cd examples/python/apps/redis_cluster
    
    # Build actors to WASM:
    ./build.sh
    
    # Deploy, initialize, and test all 11 steps (single-node):
    ./test.sh
    # or explicitly: ./test.sh 8091
    
    # Multi-node (shards distributed across both nodes):
    ./test.sh 8091 8093

    Above example runs two nodes on ports 8091 and 803. The create_shard_group uses from_registry placement to spread shard actors across both nodes automatically. Every collective operation like scatter_gather, reduce, map, broadcast routes cross-node over gRPC, using the same primitives whether the shards are local or remote.

    Fixed ports, on purpose

    The Rust embedded example starts two in-process nodes on fixed ports, :8091 and :8093, rather than ephemeral ones:

    redis-master-node  ->  gRPC :8091
    redis-replica-node ->  gRPC :8093

    Both examples use the same ports, so you can point the Python test.sh at a cluster the Rust example already stood up, or vice versa. Run ./scripts/test.sh for the Rust side.

    What the test scripts actually check

    The Rust scripts/test.sh looks for: “Cluster ready”, “Basic operations”, “broadcast_shard_group”, “scatter_gather”, “reduce”, “map + concat”, “parallel map”, “Multi-Node”, “Throughput Benchmark”, “SET/sec”, “p50”, “Example Complete”.

    The Python test.sh checks: a setup response with "status".*"ok", GET/SET/INCR semantics, an expired key returning "found".*false, replication returning "acks", a snapshot returning "shard_count" and "shards", and the benchmark returning "tps" and "p50_ms".


    Learnings

    The purpose of the book was to teach how Redis works using Rust. I rebuilt it on PlexSpaces to show how abstractions can simplify building complex distributed applications. For example, the RESP parser, the connection lifecycle, the replica list, and the WAIT loop are not your job. The redis book teaches you to hand-roll patterns like an mpsc mailbox, a spawn-per-connection loop, a tokio::select! deadline race. ThePlexSpaces uses many of same primitives to define high level abstractions but it provides an actor runtime that hides all complexity. This allows you to build the actual business logic like the storage logic, the replication semantics, the transaction model. Actors just give you a place to put them that scales horizontally without rewriting any of the logic itself. The throughput numbers show rough cost of the actor primitives, e.g., cluster setup is expensive, so you pay it once; individual operations are fast. With p50s in the hundreds of microseconds; bulk operations get much cheaper per key as the coordination overhead amortizes over more keys. This is how distributed systems work in general so you need to measure performance overhead.


    GitHub: https://github.com/bhatti/PlexSpaces

    Related reading

    Example code and documentation:

    August 25, 2026

    Orchestrating Background AI Agents for Software Teams

    Filed under: Computing — admin @ 4:25 pm

    How to design agents as a graph, run them through a harness that keeps them reliable, sandbox them because you can’t just trust what they do. (This is a follow-up to AI Writes Code, You Own the Design and Declarative AI Coding Agents with an Orchestration System)


    I have been using various workflow systems for business process management, data pipelines and batch processing over twenty years. I built a declarative orchestration system over ten years, which was before GitHub Actions, CircleCI, and GitLab CI. The idea was simple: describe your work as a graph of tasks with pipes and filter patterns declaratively, give each task clear inputs and outputs, let a server run them on a schedule or in response to events, and report the results back. Automation has always been the core discipline of software engineering and we built tools for CI/CD, cron jobs, and data pipelines. With AI, your job is changing from the one who executes the steps to a conductor. Instead of manual coding, testing and other tasks, you now design the graph of nodes, define what each node is allowed to do, let agents work on the nodes while you monitor the work where it needs a human hand. You need a harness around each agent so the graph runs reliably, a sandbox for every agent because you fundamentally cannot trust what an LLM decides to do, and a skills layer so the same harness stays general-purpose.

    I have seen teams building tightly coupled monolithic agentic systems that includes poorly implemented orchestration, the integrations, the prompts, and the skills, all bolted together. This makes it harder to make changes or extend agent capabilities and skills independently. In this post I will walk through three open source projects: Formicary, the orchestration engine that turns workflows into a graph with the Slack interface; ai-dev-tools, the harness of small scripts that actually do the work; and you-got-skills, the library of skills for SDLC.


    Daily Routine

    Each morning, you typically have to scan Jira/Github board, check pull requests to review, triage any new blockers, scroll through Slack messages to respond before deciding what to actually work on. AI coding tools have automated the implementation but you still need to know what the team is doing, catching problems, following up on reviews, etc. The basic flaw in most “agentic” setups today is that they’re still fundamentally interactive. You close your laptop and the agent stops. You need a way to keep agents working in background and instead of building systems that make you go pull information, and build agents that push it to you instead. You need somewhere to define the graph of who-triggers-what in a sandbox environment. You need a harness around each step so a flaky script doesn’t take the whole pipeline down. And you need the agent’s behavior to live somewhere editable, not buried in a prompt string, so the system stays general enough to point at a new kind of problem without a rewrite.


    The Architecture: Three Tools, Four Layers

    This section describes three open source tools: Formicary is the orchestration engine with Slack integration; ai-dev-tools is the harness that actually runs each agent inside a sandbox; and you-got-skills is the knowledge layer that keeps the whole system general-purpose.

    Above diagram shows a graph where every box is a node with a job type, a task, a skill and every arrow is an edge that Formicary evaluates at runtime. Designing agents this way, instead of as one long prompt with a loop around it makes the system debuggable. The ai-dev-tools is the harness that runs inside each node and it turns “call an LLM” into “call an LLM inside a container, with a defined timeout, a defined exit-code contract in a sandbox environment. You should not trust an agent’s judgment about what’s safe to run but with sandbox you let the harness enforce the boundary.

    Why hand off state through files?

    Formicary provides declarative syntax to store artifacts or consume artifacts from a previous task. Every script starts by checking whether its own output already exists:

    # Every script starts with this pattern
    existing = read_json(config, issue_id, "plan_result.json")
    if existing and existing.get("status") == "DONE":
        print("Already done, skipping")
        sys.exit(0)
    

    If a task dies halfway through and gets re-run, it just picks up where it left off. Exit codes are the contract between a script and the orchestrator:

    Exit codeMeaningWhat Formicary does
    0SuccessMove on to the next task
    1Error (worth retrying)Retry with backoff
    2Blocked for a human inputPause the job indefinitely
    3Not finished yet / waiting on somethingPause, then resume on a trigger or after a delay

    The Skills Library

    Skills are the knowledge layer that allow general-purpose harness instead of hardcoded to whatever workflow you built. The graph and the sandbox don’t know anything about code review or standups as they just run nodes. A skill is what tells an agent how to do a specific kind of engineering work well. Add a new skill and you’ve extended the system to a new kind of task without touching Formicary or ai-dev-tools at all.

    you-got-skills/
    ??? skills/
        ??? ygs-standup/skill.md
        ??? ygs-risk-scan/skill.md
        ??? ygs-implement/skill.md
        ??? ygs-review-pr/skill.md
        ??? ygs-code-review/skill.md
        ??? ygs-security-review/skill.md
        ??? ygs-sre-review/skill.md
        ??? ygs-learn/skill.md
        ??? ygs-retro/skill.md
        ??? ygs-sprint-plan/skill.md
        ??? ygs-investigate/skill.md
        ??? ygs-qa/skill.md
        ??? ygs-wbs/skill.md
        ??? ygs-estimate/skill.md
        ??? ygs-ship/skill.md
        ??? ygs-triage/skill.md
        ??? ... 28 total
    

    Together they cover the whole lifecycle:

    PhaseSkills
    Requirementsygs-refine-prd, ygs-review-prd, ygs-refine-trd, ygs-review-trd
    Architectureygs-refine-architecture, ygs-review-architecture, ygs-spike
    Planningygs-sprint-plan, ygs-wbs, ygs-estimate, ygs-triage
    Executionygs-implement, ygs-qa, ygs-ship, ygs-uat
    Reviewygs-review-pr, ygs-code-review, ygs-security-review, ygs-sre-review, ygs-api-review, ygs-ui-review
    Team intelligenceygs-standup, ygs-risk-scan, ygs-pr-queue, ygs-sync
    Learningygs-learn, ygs-retro, ygs-investigate

    Skills produce and consume a consistent folder structure inside your project:

    your-project/
    ??? docs/
    ?   ??? prd/              # Product requirements (YYYY-MM-DD-slug.md)
    ?   ??? trd/               # Technical designs
    ?   ??? adr/               # Architecture decisions (NNN-slug.md)
    ?   ??? spikes/            # Spike findings
    ?   ??? learnings/         # Learnings pulled from PRs and incidents
    ??? tasks/
        ??? backlog/            # task-NNN.md
        ??? in-progress/        # moving a file here = starting it
        ??? done/               # moving a file here = finishing it
    

    Status is the folder. No ticket-state dropdown, no transition workflow to configure. mv tasks/backlog/task-042.md tasks/in-progress/ means the task has started. git blame tells you who moved it and when.

    Diverge, then converge

    The diverage and converge pattern allows running multiple independent LLM passes first (diverge), then merge and rank what came back (converge). A single reviewer may miss something but several independent reviewers will catch most of the issues. The ygs-review-pr skill runs four independent passes in parallel for correctness, security, API surface, and SRE concerns. A finalize step then merges everything and ranks it by severity.

    The standup workflow uses the same idea. ygs-standup gathers signals from the issue tracker and from Slack then cross-references the two.


    The Standup Workflow

    Every weekday at 8am, a cron job wakes up, queries your Jira sprint (or GitHub), reads the open PRs, pulls the last 26 hours of Slack messages from your standup channel, and turns all of it into a brief.

    # ai-standup-jira.yaml
    job_type: ai-standup-jira
    description: "Daily standup brief from Jira sprint + Bitbucket PRs + Slack signals"
    cron_trigger: "0 0 8 * * 1-5 *"
    max_concurrency: 1
    timeout: 900s
    
    tasks:
      - task_type: gather     # pulls Jira issues, Bitbucket PRs, Slack in parallel
      - task_type: synthesize # Claude + ygs-standup + ygs-risk-scan ? standup_brief.md
      - task_type: post       # renders HTML artifact, posts to Slack
    

    The synthesize task calls the ygs-standup skill, which follows a fairly strict protocol:

    **Alice:** Closed PROJ-42 (auth fix). Working on PROJ-51 (rate limiter)
    — PR open 28h, no review yet. [Slack: "waiting on infra cert renewal"]
    
    **Bob:** No tracker activity in the last 24h. Last Slack message Monday
    (3 days ago). [Slack: silent since Monday]
    
    **Carol:** PROJ-55 (data export) marked In Progress, no commits in 4 days.
    Blocked label present. [Tracker: blocked since Tuesday]
    

    Every claim traces back to a ticket, a PR, or a Slack message. ygs-risk-scan runs right after and appends a ranked list, using thresholds you can tune:

    ? HIGH  PROJ-55 blocked, blocks PROJ-60 (in progress, owned by different person)
    ? MED   PR #142 open 28h, single reviewer, no activity
    ? MED   Carol: no updates in 4 days, sprint ends Friday
    

    The thresholds live in a shared Markdown file:

    SignalDefault severityEscalates to HIGH if…
    Issue stale > 3 daysMEDIUMit blocks another issue
    Issue stale > 5 daysHIGH—
    PR open > 2 days, no reviewMEDIUMonly one reviewer assigned
    PR open > 4 daysHIGH—
    Person silent > 2 daysMEDIUMalso no tracker activity
    Blocked labelHIGH—
    Dependency chain: upstream is staleHIGH—
    Sprint ends in < 2 days, not startedHIGH—

    You can also trigger the standup on demand from Slack:

    @bot standup
    @bot risk

    Label an Issue, Get Back a Pull Request

    You can label any Jira or GitHub issue ai-ready. Every five minutes a cron job checks for newly labeled issues and kicks off a four-task pipeline. When it’s done, there’s an open PR with an implementation and tests, the label has flipped to ai-pr-open.

    # ai-gh-issue-picker.yaml
    job_type: ai-gh-issue-picker
    cron_trigger: "*/5 * * * *"
    tasks:
      - task_type: gather-issues
        script:
          - python -m scripts.gh.issue_picker
        on_exit_code:
          2: COMPLETED  # no issues = not an error
      - task_type: submit-jobs
        # Uses formicary template to fan out one ai-gh-implement job per issue
        script:
          - '{{SubmitJobsFromJSON "ai-gh-implement" .IssuesJSON}}'
    

    It has a built-in guard: if 10 or more implement jobs are already running or queued, it skips its turn. That stops someone from labeling 50 issues at once and blowing up the queue. The implementation pipeline itself:

    # ai-gh-implement.yaml
    job_type: ai-gh-implement
    max_concurrency: 5
    timeout: 86400s    # 24 hours — some implementations take a while
    
    tasks:
      - task_type: plan
        timeout: 15m
        script:
          - python -m scripts.gh.issue_picker --issue-id {{.IssueNumber}}
          - python -m scripts.gh.plan --issue-id {{.IssueNumber}}
        on_exit_code:
          2: PAUSE_JOB    # BLOCKED — needs a human before continuing
        on_completed: implement
    
      - task_type: implement
        timeout: 45m
        script:
          - python -m scripts.gh.implement --issue-id {{.IssueNumber}}
        on_completed: create-pr
    
      - task_type: self-review
        timeout: 10m
        script:
          - python -m scripts.review.run --mode self-review
              --issue-id {{.IssueNumber}} --base-branch main
        on_exit_code:
          2: PAUSE_JOB    # BLOCKED — critical finding, needs a human before the PR opens
        on_completed: create-pr
    
      - task_type: create-pr
        timeout: 10m
        script:
          - python -m scripts.gh.create_pr --issue-id {{.IssueNumber}}
        on_completed: poll-pr
    
      - task_type: poll-pr
        dependencies: [create-pr, poll-pr]  # self-referencing = loop
        script:
          - python -m scripts.gh.poll_pr --issue-id {{.IssueNumber}}
        on_exit_code:
          3: PAUSE_JOB    # PR still open — check back in `delay` seconds
        delay: "{{.PollInterval}}s"  # default 120s
    

    Here’s the full chain:

    Let’s walk through what actually happens for a real issue.

    Plan

    scripts/gh/plan.py calls Claude with up to 50 turns and a prompt that tells it to:

    1. Read CLAUDE.md, .cursorrules, or any repo-specific coding guidelines if they exist
    2. Discover .claude/skills/ in the repo — if a skill applies, plan to invoke it
    3. Before designing new abstractions, search utils/, shared/, common/ for existing utilities
    4. Check for monorepo structure
    5. Generate a concise plan covering:
       - Task breakdown with complexity estimates (S/M/H/XL)
       - Exact files to create/modify per task
       - Test strategy: write failing tests first, then implement
       - A "Failing Test Spec" section
       - Any risks or blockers
    6. Classify overall complexity: S/low (?3 files), M/medium (4-10), H/high (>10)
    7. Write the plan to PLANS/{slug}-{issue_id}-plan.md
    

    What it produces:

    /workspace/42/
    ??? issue.json         # issue title, body, labels, assignee
    ??? plan.md            # human-readable plan
    ??? plan_result.json   # {"status":"DONE","task_count":3,"total_complexity":"M"}
    ??? PLANS/
        ??? add-rate-limit-42-plan.md
    

    Implement

    scripts/gh/implement.py calls Claude with up to 200 turns. The key rules from the ygs-implement skill:

    • Check for existing utilities before writing new ones.
    • TDD: write the failing test first, then make it pass.
    • One commit per plan task, message format "task: <description>".
    • After all tasks are done, run the full test suite and iterate on failures up to twice.
    /workspace/42/
    ??? impl_result.json   # {"status":"DONE","files_changed":["src/..."],"commits":5,"tests_status":"passing"}
    ??? branch.txt         # "ai/42-add-rate-limiting"
    

    ygs-implement also uses “ceremony levels” so small tasks don’t get over-engineered:

    Light    (1-3 files, <300 lines):   proceed directly, skip checkpoints
    Standard (4-8 files, 300-800):      plan mode + checkpoint every 5 files
    Heavy    (8+ files, 800+ lines):    flag as oversized, ask user to split
    

    Picking the model by complexity

    The plan task classifies overall complexity and writes it to /workspace/plan_complexity.txt. The implement task reads that and picks the right model:

    # scripts/common/config.py
    COMPLEXITY_MODEL_MAP = {
        "low":    MODEL_BEDROCK_HAIKU,   # ?3 files, simple edits — fast and cheap
        "medium": MODEL_BEDROCK_SONNET,  # default — most issues
        "high":   MODEL_BEDROCK_OPUS,    # complex architecture changes
    }
    

    The plan prompt writes a single word (low, medium, or high) to that file, and the implement task reads it back:

    # In ai-gh-implement.yaml — implement task
    COMPLEXITY=$(cat /workspace/plan_complexity.txt 2>/dev/null || echo "medium")
    AI_MODEL="${ANTHROPIC_COMPLEXITY_${COMPLEXITY^^}_MODEL:-${ANTHROPIC_DEFAULT_SONNET_MODEL}}"
    

    AnthropicComplexityLowModel and AnthropicComplexityHighModel are set in your org config by deploy-ai-workflows.sh, and can be overridden per deployment in models.env.

    Polling the PR and responding to feedback

    Every 2 minutes, the poll task checks the PR for new comments. It only reacts to comments that start with ai-bot. The agent stays out of human-to-human review discussion, and it never responds to its own earlier comments. So when a reviewer writes:

    ai-bot please add a test for the rate limit exceeded case
    

    The poll task reads it, applies the feedback, marks the comment handled in processed_comments.json, and keeps polling. Once the PR merges, the learning step kicks off automatically.


    PR Review With a Human Gate

    Code review usesthe diverge-then-converge pattern to review the code with multiple perspectives like security, architecture, SRE, etc.

    # ai-gh-review.yaml
    job_type: ai-gh-review
    max_concurrency: 10
    
    tasks:
      - task_type: review        # Claude runs ygs-review-pr, writes findings.json
      - task_type: await-feedback # posts Block Kit to Slack, exits 3 ? PAUSE_JOB
      - task_type: finalize       # reads Decision, posts result to PR thread
    

    The review step runs ygs-review-pr with this instruction:

    1. Invoke the /ygs-review-pr skill to perform a full PR review
    2. After the skill completes, write findings to findings.json
    3. Output ONLY this JSON on the last line:
       {"status":"DONE","findings_count":N,"verdict":"APPROVE|REQUEST_CHANGES","summary":"..."}
    

    ygs-review-pr runs four passes in paralle:

    1. Correctness: logic errors, null handling, incomplete enum handling, partial failure, race conditions, off-by-one errors
    2. Security: injection vectors (SQL, command, XSS, SSRF, path traversal), auth/authorization, data exposure across tenant boundaries, etc.
    3. API surface: breaking changes, contract violations, backwards compatibility, versioning
    4. SRE: failure modes, blast radius, observability gaps, rollback safety, resource consumption, dependency risk

    Once all four are done, findings get merged and ranked: CRITICAL > HIGH > MEDIUM > LOW, and by confidence within each level. The verdict maps straight off the highest severity found:

    Any CRITICAL or HIGH finding  ? REQUEST_CHANGES
    Only MEDIUM/LOW findings      ? COMMENT
    No findings / only low conf.  ? APPROVE
    

    Deep review: seven domains in one pass

    Standard review runs four passes. Deep review adds three more: performance, testing quality, and architecture.

    @bot deep review https://github.com/org/repo/pull/42

    These all map to the same ai-gh-review (or ai-jira-review) job type, just with ReviewDepth=deep injected as a static job variable:

    # workflows.yml — deep-review entry
    - name: deep-review
      job_type: ai-gh-review
      triggers: ["deep review", "full review", "arch review"]
      target_kind: github
      extra_params:
        ReviewDepth: "deep"
      description: "Deep 7-domain GitHub PR review: standard + performance, testing quality, architecture"
    

    The workflow reads ReviewDepth and picks the right skill: ygs-review-deep when it’s set to “deep,” ygs-review-pr otherwise. ygs-review-deep is a superset of the standard four passes, plus:

    • Performance: algorithmic complexity, N+1 queries, cache misses, lock contention, unnecessary allocations
    • Testing quality: coverage gaps, brittle assertions, missing edge cases, test-code coupling
    • Architecture: single-responsibility violations, circular dependencies, premature abstractions, missing boundaries

    Self-review before the PR even opens

    The implement pipeline runs a self-review task right before create-pr. Before the PR exists, the agent reviews its own diff against the base branch:

      - task_type: self-review
        timeout: 10m
        script:
          - python -m scripts.review.run --mode self-review
              --issue-id {{.IssueNumber}} --base-branch {{.BaseBranch}}
        on_exit_code:
          2: PAUSE_JOB    # BLOCKED — critical finding, needs a human before the PR opens
        on_completed: create-pr
    

    This runs ygs-implement in review mode, compares the diff to the original plan, and writes self_review.json. The outcome maps directly to an exit code:

    self_review_statusExit codeWhat the pipeline does
    APPROVED0Proceed to create-pr
    NEEDS_FIX0Claude fixes it inline, then create-pr
    BLOCKED2PAUSE_JOB — a human needs to decide before the PR opens

    How Slack Messages Turn Into Workflows

    Socket Mode lives inside the Formicary queen itself. The queen opens an outbound WebSocket to Slack using your xapp- app-level token. When you mention the bot:

    @bot review https://github.com/myorg/myrepo/pull/142
    

    The queen’s SlackService does one thing, deterministically: it strips the mention, takes the first word, and looks it up against a route table in the queen’s config.

    # In k8s/formicary-leader.yaml ConfigMap, under slack.routes:
    slack:
      routes:
        - triggers: ["review", "pr"]
          job_type: ai-gh-review
          description: "PR review: correctness, security, API, SRE"
        - triggers: ["implement", "build"]
          job_type: ai-jira-implement
          description: "Full pipeline: plan ? implement ? PR"
        - triggers: ["standup", "status", "daily"]
          job_type: ai-standup-jira
          description: "Daily standup brief from Jira"
        - triggers: ["adhoc"]
          job_type: ai-adhoc
          description: "Ad-hoc Claude invocation with any skill"
    

    Everything after the trigger word gets passed through as the Prompt job parameter, verbatim. The queen’s whole job is mapping verb to job type and passing the text along. It does zero AI work of its own. All of that happens inside the ai-dev-tools container once the job actually starts. Once routing resolves, the queen submits the job and replies right in the thread:

    Started ai-gh-review (job req-7f3a2) — I'll post updates here.
    https://formicary.example.com/dashboard/jobs/requests/req-7f3a2
    

    Replying to the thread resumes a paused job. If a review job is paused waiting on a decision and you reply in that thread, the queen matches it by SlackThreadTs and resumes it with your reply text injected as Prompt.

    Registering as a developer

    Before Slack commands work for you, you DM the bot your Formicary API token, once:

    DM to @bot:
    setup eyJhbGc...  (your Formicary API token)
    

    The queen validates the token inline and from that point on, any @bot mention from you has a known Formicary identity behind it.

    How multi-tenant isolation actually works, end to end:

    When you type @bot review https://github.com/org/repo/pull/42, here’s what the queen does, entirely server-side:

    1. Reads your Slack user ID (U0A1HQL0C9J) off the Socket Mode event.
    2. Looks up slack_user_id = U0A1HQL0C9J in user_configs and finds your Formicary user record, including your UserID and OrganizationID.
    3. Calls SaveJobRequest(qc, req) and the server overwrites request.UserID and request.OrganizationID from that context.
    4. Schedules the job and the ant scheduler first looks for a worker registered under your org_id.

    Ant routing: when you connect your laptop as a worker with setup-ant-worker.sh --token <your-token>, the queen reads org_id out of your JWT at connect time and records it on that worker. Your jobs prefer your own worker.

    All the commands

    What you typeWhat runs
    @bot standupDaily brief: per-person status, risks, discussion questions (routes to Jira or GitHub via DEFAULT_TRACKER)
    @bot risk / @bot risksRanked sprint risks with a capacity check
    @bot prs / @bot open prs / @bot review queueOpen PRs grouped by reviewer status, sorted by age
    @bot pr comments <url>All inline feedback and open tasks for a PR
    @bot review <github-url>Standard PR review: correctness, security, API, SRE (4 domains)
    @bot review <bitbucket-url>Same, for Bitbucket PRs
    @bot deep review <url>Deep review: standard 4 domains + performance, testing, architecture (7 domains)
    @bot full review <url>Alias for deep review
    @bot arch review <url>Alias for deep review
    @bot security review <url>OWASP-focused security audit
    @bot sre review <url>Failure modes, observability, deploy safety
    @bot implement PROJ-123Full pipeline: plan ? implement ? self-review ? PR, for a Jira issue
    @bot implement 42Same, for a GitHub issue number (model picked by complexity: Haiku/Sonnet/Opus)
    @bot jira-query <term> / @bot qjira <term>Search open Jira issues by keyword, results as a Block Kit table
    @bot jira-analyze PROJ-1, PROJ-2Claude analyzes root cause + possible fixes for Jira issues
    @bot gh-query <term>Search open GitHub issues by keyword
    @bot gh-analyze #123, #456Claude analyzes root cause + possible fixes for GitHub issues
    @bot adhoc <free text>Run any Claude skill with a freeform prompt
    @bot helpList every command

    Ad-hoc Skill Execution

    ai-adhoc is a general-purpose runner: any you-got-skills skill can be invoked with a free-form prompt, and the result comes back into your Slack thread.

    # ai-adhoc.yaml
    job_type: ai-adhoc
    max_concurrency: 20
    timeout: 1800s
    
    variables:
      Skill:  { type: STRING, required: true }
      Prompt: { type: STRING, required: true }
    

    At runtime the script:

    1. Looks for the skill in a few candidate locations (/workspace/skills, ~/.claude/skills/you-got-skills/skills, ~/workplace/you-got-skills/skills).
    2. Writes .ygs/tracker.yml dynamically from environment variables (Jira or GitHub config, team members, sprint info).
    3. Invokes Claude with the skill content and your prompt.
    4. Strips Markdown formatting from the output so it renders cleanly in Slack.
    5. Posts up to 3000 characters back into the originating thread.

    Adding a new shorthand command is just one entry in the queen’s route table:

    # k8s/formicary-leader.yaml — slack.routes
    slack:
      routes:
        - triggers: ["prs", "open prs", "review queue"]
          job_type: ai-adhoc
          description: "Open PRs grouped by review status"
    

    @bot prs then submits ai-adhoc with Prompt="prs", and the container maps that to the ygs-pr-queue skill. No Python code changes needed.


    Querying and Analyzing Issues From Slack

    Two commands take you from a Slack message straight to structured Jira insight, no browser required.

    @bot query-jira: find issues by keyword

    @bot jira-query auth timeout
    

    This submits an ai-jira-query job. The script builds a JQL query scoped to your project and, optionally, your team’s custom field. Results come back as a structured Slack Block Kit table:

    Jira issues matching "flaky tests" (5 found)
    
    PROJ-1001  [Bug] Flaky test in auth service                    - clickable link
               Status: In Progress   Priority: High
               Assignee: alice        Date: 2026-07-28
    
    PROJ-995   [Story] Fix race condition in logger test
               Status: To Do         Priority: Medium
               Assignee: bob         Date: 2026-07-21
    ...
    

    @bot jira-analyze: root cause from issue keys

    @bot jira-analyze https://yourorg.atlassian.net/browse/PROJ-1001
    

    This routes to the same ai-jira-query job type. The result posts back to your thread.

    @bot gh-query / @bot gh-analyze the GitHub equivalents

    Same commands, gh- prefix, for teams on GitHub instead of Jira:

    @bot gh-query open authentication bugs
    @bot gh-analyze https://github.com/org/repo/issues/42
    

    Team filtering

    Both commands automatically filter to your configured team and sprint:

    • JIRA_SPACE (or BITBUCKET_WORKSPACE): the team/area filter value.
    • JIRA_TEAM_FIELD: the Jira custom field name (default EngScrumTeam, resolved to a field ID dynamically).
    • Set JIRA_TEAM_FIELD="" to turn the filter off entirely.

    The Learning Loop: Getting Better Over Time

    Most agent systems are stateless and every run starts from zero. In our workflow, after every PR merges, learn.py runs automatically as part of the poll-pr task. It reads the PR comments, the implementation artifacts, and the review findings, then invokes ygs-learn:

    ygs-learn protocol:
    1. Capture: What happened? Why does it matter? Category?
    2. Dedup: search docs/learnings/ for similar slug before creating new
    3. Write to docs/learnings/YYYY-MM-DD-slug.md
    

    A learning document looks like this:

    # Rate limiter key collision when user has multiple active sessions
    
    **Category:** Edge Case
    **Date:** 2025-08-03
    **Source:** PR review finding, PROJ-123
    
    ## Learning
    When a user has multiple active sessions, rate limiting by user_id counts across
    all sessions. A single slow client can exhaust the budget for all their tabs.
    
    ## Evidence
    Review comment on PR #142 flagged this. Reproduced locally with two
    concurrent sessions against the same account.
    
    ## Application
    When implementing per-user rate limits, check whether session isolation is
    intended. If counts should be per-session, key by session_id not user_id.
    

    Next time an implementation runs, that document is part of the context Claude reads.

    ygs-retro runs at sprint end. It reads tasks/done/, recent git history, and the sprint’s accumulated learnings, and asks pointed questions based on what it actually found.

    ygs-investigate enforces a debugging discipline: build a feedback loop first, e.g., a failing test, a log line, a REPL session before forming any hypotheses. Rank hypotheses 1 through 5. Instrument one variable at a time with tagged markers.


    Trust and Oversight

    Though, AI agents have solved most of coding and testing tasks but it still requires human review and feedback. We need a trusted autonomy, with minimal friction. You don’t trust the model to know its own limits; you build a harness that doesn’t need you to. That’s why every agent runs inside a container it can’t escape. You can’t ask an AI agent to self-police or prompt it to be careful. Every decision point is explicit. The agent never merges to main. It opens a branch, opens a PR, and stops there. A human approves and merges. For anything that needs a more formal sign-off, Formicary supports approval workflows with SLAs:

    - task_type: security-approval
      method: MANUAL
      approval_policy:
        min_approvals: 1
        sla_deadline: 4h
        timeout_action: ESCALATE
        escalation_recipients: "security-oncall@example.com,vp-eng@example.com"
        escalation_message: "Security approval SLA breached — deployment blocked"
      on_exit_code:
        APPROVED: deploy-prod
        REJECTED: notify-rejected
    

    Secrets live in Kubernetes Secrets, never in ConfigMaps or plain env files. The container runs as non-root (uid 1000). Every artifact is a plain file and every state transition is logged in Formicary.


    What Else You Can Build

    Everything above is running in production today, but the same architecture supports a much wider range of background agents.

    • Codebase quality agent. A weekly workflow runs ygs-code-review across everything changed in the last week, posts a ranked findings report to Slack, and files tasks in tasks/backlog/ for anything CRITICAL or HIGH.
    • Documentation drift detector. A webhook fires when a PR touching an API handler merges. The workflow checks whether the matching docs were updated.
    • Duplicate abstraction scanner. A periodic workflow compares utility functions across repos owned by different teams and posts “team B has something that looks like what you just built,”.
    • Security posture monitor. Nightly, ygs-security-review runs against everything merged in the last 24 hours that touches auth, authorization, or data access. Findings go to a security channel.
    • Sprint health check, Wednesday afternoons. A mid-sprint cron runs ygs-risk-scan and only posts if it finds something HIGH severity. Most weeks it says nothing.
    • Bug pattern finder. A workflow runs ygs-investigate against recent error logs, proposes the top three hypotheses for each recurring pattern.

    Getting Started

    Installing the skills locally

    git clone https://github.com/bhatti/you-got-skills.git
    cd you-got-skills && ./setup
    

    Now any skill runs right in your IDE:

    /ygs-standup
    /ygs-risk-scan
    /ygs-review-pr https://github.com/org/repo/pull/42
    /ygs-implement
    

    Loading extra skill repos at runtime

    Every job pod can pull in additional skill repos without a rebuild. Set EXTRA_SKILLS_REPOS before running the deploy script.

    # 1. Plain URL — sparse-clones only the skills directory (fast, default)
    EXTRA_SKILLS_REPOS=https://github.com/myorg/my-skills.git
    
    # 2. Comma-separated — multiple repos in one value, YAML-safe
    EXTRA_SKILLS_REPOS="https://github.com/bhatti/you-got-skills.git,skills-cli:nutlope/hallmark"
    
    # 3. JSON array — full control (use for org config; JSON breaks YAML template substitution)
    EXTRA_SKILLS_REPOS='[
      {"url": "https://github.com/bhatti/you-got-skills.git", "sparse": false},
      {"url": "nutlope/hallmark", "type": "skills-cli"}
    ]'
    

    Setup environment variables:

    GH_ORG=your-org
    GH_REPO=your-repo
    GH_TOKEN=ghp_your_token_here
    ANTHROPIC_API_KEY=sk-ant-your_key_here
    AI_MODEL=claude-sonnet-4-6
    
    # Controls which tracker bare commands like "standup" route to.
    # "jira" routes to ai-standup-jira; "github" routes to ai-standup-gh.
    DEFAULT_TRACKER=jira
    

    Start formicary server

    You can use kubernetes to get Formicary running.

    export COMMON_AUTH_JWT_SECRET="<stable-secret-never-rotate>"
    export COMMON_AUTH_GOOGLE_CLIENT_ID="<google-client-id>"
    export COMMON_AUTH_GOOGLE_CLIENT_SECRET="<google-client-secret>"
    export SLACK_BOT_TOKEN="xoxb-..."
    export SLACK_APP_TOKEN="xapp-..."
    
    ./scripts/deploy-formicary.sh --ec2-ip 10.X.X.X
    

    Connect your ant worker

    Jobs run on your own laptop’s local cluster.

    ./scripts/setup-ant-worker.sh \
      --queen formicary.example.com \
      --token "$FORMICARY_TOKEN"
    

    Deploying the Slack integration

    Slack is built into the Formicary queen, so there’s no separate router pod to deploy.

    Create a Slack app with Socket Mode

    In your Slack app settings:

    1. Enable Socket Mode: generate an xapp- app-level token (scope: connections:write).
    2. Add these bot token scopes under OAuth & Permissions:
    ScopePurpose
    app_mentions:readReceive @bot mentions
    channels:historyRead channel messages
    channels:readList channels
    chat:writePost messages and Block Kit
    groups:historyRead private channel messages
    groups:readList private channels
    im:historyRead DMs (for the setup registration flow)
    im:writeReply in DMs
    users:readResolve user display names
    1. Subscribe to bot events under Event Subscriptions:
      • app_mention — @bot mentions in channels
      • message.im — DMs, for the setup registration flow
    2. Install to your workspace (needs admin approval if app installs are restricted).

    Each developer registers

    Anyone who wants Slack commands to work DMs the bot once:

    DM to @bot:
    setup eyJhbGc...  (Formicary API token from dashboard ? API Tokens)
    

    In any channel the bot’s been invited to:

    @bot help          - lists all commands
    @bot standup       - runs standup, posts to thread
    

    Summary

    The core patterns in this post include declarative pipelines, cron triggers, file-based artifact handoff, exit-code contracts, approval gates. I have used these patterns to automate complex business processing, data pipelines, and CI/CD processes. I am now using it to automate AI backed tasks: a task can read code and form a judgment about it. When the graph, the harness, the sandbox, and the knowledge are four separate layers instead of one tangled system, each one evolves on its own. It allows you to update an environment with configuration, markdown files and configuration. Your job is now to design the graph, define what’s in each node, decide where the sandbox boundary sits, and write down what “good” looks like as a skill. Instead of manually gathering information, you use agents to do the low-level work. You then review the finished product, at the review verdict, at the escalation and decides what matters. That’s conducting, not playing every instrument, and it’s where your judgment actually belongs.


    Related Reading

    Code

    August 16, 2026

    Structured Concurrency in Modern Programming Languages Part V: The Coordination Models Behind It All (CSP, Actors, Linda, and async/await)

    Filed under: Computing — admin @ 4:37 pm

    This is a part of series on structured concurrency: Part I (the general problem and TypeScript), Part II (Erlang and Elixir), Part III (Go and Rust), and Part IV (Kotlin and Swift).

    In the earlier parts of this series I delved into how TypeScript, Erlang, Elixir, Go, Rust, Kotlin, and Swift each handle structured concurrency in practice such as spawning tasks, waiting for children to finish, propagating errors, and cancelling work cleanly. But I skipped over the the coordination models these languages are actually built on. For example, Go didn’t invent channels and Erlang didn’t invent actors. Both are engineering ideas that go back to the 1970s and 80s, and once you understand the original model, most of the “gotchas” you hit while using the language stop looking like bugs and start looking like predictable consequences of a design choice made decades ago.

    This post explains where each model came from, what it actually guarantees and how structured concurrency sits on top of all of them as a separate concern. It includes PlexSpaces, an actor-and-tuplespace framework I’ve been building in Rust that takes a pragmatic stance on this history, e.g., bounded mailboxes instead of Erlang’s unbounded ones, first-in-first-out matching instead of the classic tuple-space model’s unspecified ordering, and one small API instead of forcing you to learn several calculi at once.

    A short timeline

    It helps to see these ideas in the order they actually appeared:

    • 1973: Carl Hewitt proposes the actor model: small, isolated units of state that can only talk to each other by sending messages.
    • 1978: Tony Hoare publishes Communicating Sequential Processes (CSP), a mathematical notation (a “process algebra”) with a precise definition of what it means for two processes to synchronize.
    • 1985: David Gelernter publishes Linda, a coordination model built around a shared associative memory (the “tuple space”).
    • 1986: Gul Agha’s book extends the actor model with a fuller algebraic treatment.
    • 1986: Joe Armstrong and colleagues at Ericsson start building Erlang.
    • Early-to-mid 2000s: event-loop async/await goes mainstream: Node.js’s callback-then-promise evolution.
    • 2009: Go ships with goroutines and channels inspired by CSP.
    • 2018 onward structured concurrency (Trio in Python, Kotlin’s coroutines, Swift’s TaskGroup, Java’s StructuredTaskScope) formalizes a simple idea: a spawned task’s lifetime should never outlive the scope that spawned it.

    Notice that everything on that list except the last item is about how work talks to other work. Structured concurrency is about a completely different question, i.e., how work’s lifetime gets tracked.

    Concurrency and parallelism

    Concurrency is a property of how a program is structured: multiple logically independent activities are in progress, possibly interleaved on a single CPU core. Parallelism is a property of execution: things are genuinely happening at the same time, which requires more than one core. You can have concurrency without parallelism like Node.js’s single-threaded event loop juggling many pending requests on one core. You can also have parallelism without concurrency like a tight SIMD loop doing the same arithmetic on many numbers at once has no interleaved independent logic at all. Go’s own documentation defines it as: concurrency is about dealing with lots of things at once, parallelism is about doing lots of things at once. Goroutines give you concurrency; whether that concurrency turns into real parallelism depends on GOMAXPROCS and how many cores are actually available.

    This matter for CSP and actors because both are concurrency models and neither one is “more parallel” than the other. What actually differs between them is how they structure communication.

    Five models, one spectrum of coupling

    Every concurrency model is answering the same underlying question, i.e., how does one unit of work talk to another one.

    ModelOriginHow units talkCouplingFormal backing
    CSPHoare, 1978Synchronous rendezvous on a named channelTime-coupledFull algebra, checked by tools like FDR
    Go-style CSPGo, 2009Channel, synchronous or bufferedTime-decoupled if bufferedNone
    Actor modelHewitt 1973 / Erlang 1986Async message to a named addressIdentity-coupled, time-decoupledPartial (Clinger, Agha)
    async/awaitNode.js/C#/Python event loopsFuture/promise handleTime-decoupledNone
    Linda / tuple spaceGelernter, 1985Tuple matched by contentFully decoupled – no identity, no timingPartial (Klaim’s semantics)
    Structured concurrencyTrio/Kotlin/Swift, 2018+Whatever the underlying model usesLifetime-coupled to a scopeNone

    CSP: the algebra

    Before getting into Go’s implementation, let me explain algebra in CSP. I’ve written before about algebraic effects like resumable exceptions in OCaml 5 / Koka that let a function declare what it needs without saying who provides it. But algebra in CSP is a process algebra: a small set of operators like sequence, choice, parallel composition, hiding with equational laws. Because those laws exist, you can prove two CSP process descriptions behave identically. Tools like FDR (Failures-Divergences Refinement) do this mechanically, e.g., you describe your system as CSP processes, describe a specification as another CSP process, and FDR checks whether the implementation actually refines the spec, across every possible interleaving.

    Here’s what that looks like for a scatter-gather pattern, an orchestrator firing off requests to several workers and collecting exactly K responses:

    -- Specification: orchestrator collects exactly K results then stops
    SPEC = scatter -> (collect -> collect -> collect -> STOP)
    
    -- Implementation: N workers communicate via channels
    WORKER(i) = request.i -> response.i -> STOP
    SYSTEM = (||| i : {0..4} @ WORKER(i))
             [| {| response |} |]
             COLLECTOR(3)
    COLLECTOR(0) = STOP
    COLLECTOR(k) = response?i -> COLLECTOR(k-1)
    
    -- FDR checks: assert SYSTEM [T= SPEC (trace refinement)
    -- This PROVES: no deadlock, no livelock, exactly K responses collected
    

    A handful of operators do almost all the work here:

    CSP OperatorMeaningWhat FDR Proves
    P ? QExternal choice: the environment decidesDeadlock-freedom: at least one branch is always available
    P ? QInternal choice: the process decides nondeterministicallyLiveness: both branches are eventually reachable
    P ? QParallel composition, synchronized on shared eventsNo protocol deadlock between P and Q
    P ; QSequential composition: Q starts only after P terminatesTermination: P always reaches STOP
    P \ AHiding: internal events in set A become invisibleDivergence-freedom: no infinite internal loops

    In other words you can model your protocol in CSP and let FDR check every possible interleaving for you. For the scatter-gather pattern specifically, FDR would catch, automatically, before any code runs:

    • A worker that never responds (a deadlock)
    • A collector that waits for more responses than the workers can ever produce (also a deadlock)
    • A timeout path that accidentally creates an infinite retry loop (a livelock)

    Go and Rust can’t do any of this because the as soon as you add a buffer (make(chan int, 5)) or an async boundary, you’ve left the strictly synchronous world that FDR reasons about. Go’s race detector can find data races at runtime, after the fact. Rust’s borrow checker prevents a whole class of shared-state bugs at compile time. But, neither one can prove protocol-level, whole-system deadlock-freedom the way FDR can for pure CSP.

    Go’s channels

    Hoare’s CSP defines communication as synchronous by construction where a send and its matching receive aren’t two separate events that happen to line up in time but they’re the same event in the algebra. Go’s unbuffered channel matches that faithfully: ch <- x and <-ch really do rendezvous. A buffered channel doesn’t, and that one divergence from the original model explains most of the sharp edges Go developers run into. Also, real CSP processes have no persistent identity beyond the algebra describing them, and channels are closer to anonymous synchronization events than to objects you hold a reference to. Go’s channels, by contrast, are first-class values that you create one, pass it into ten different functions, and any of them can close it. Nothing in the language enforces “exactly one owner, exactly one closer” and this is the seed of several of the gotchas below.

    func worker(jobs <-chan int, results chan<- int) {
        for j := range jobs {
            results <- j * j
        }
    }
    
    func main() {
        jobs := make(chan int, 5)
        results := make(chan int, 5)
        go worker(jobs, results)
    
        for i := 1; i <= 5; i++ {
            jobs <- i
        }
        close(jobs) // safe: only the sender closes, and no sends follow
    
        for i := 0; i < 5; i++ {
            fmt.Println(<-results)
        }
    
        // jobs <- 6 // panics — sending on a closed channel always panics,
                     // whether or not anything is still listening
    }
    

    Three more gotchas that Go’s compiler won’t warn you about:

    • Receiving from a closed channel never panics. It returns the zero value and ok == false immediately instead of blocking.
    • A nil channel blocks forever, on both ends, with no panic. Occasionally this is useful on purpose but if it happens to an uninitialized struct field, it results in permanent hang with no error message pointing you at the cause.
    • Goroutines have no structure by default. go func(){}() creates nothing that ties that goroutine’s lifetime to the caller. A goroutine permanently blocked on a channel operation is invisible to the garbage collector and invisible to Go’s deadlock detector. It causes a partial leak where the program running fine, with one goroutine stuck forever in the background.

    One place Go actually stayed close to the algebra: select. Hoare’s algebra has external choice (?) as a first-class operator, and select‘s randomized tie-break among multiple ready cases matches CSP alegbra.

    select {
    case job := <-jobs:
        handle(job)
    case <-ctx.Done():
        return ctx.Err() // structured cancellation, Go-style
    default:
        // non-blocking probe — CSP has no built-in default arm,
        // but this is the standard way to build one
    }
    

    Best practices that have converged around Go’s channels

    • Confine, don’t share. Exactly one goroutine should own a channel’s write side and be the one to close it.
    • Thread context.Context through every long-running goroutine. A select that never watches ctx.Done() is a goroutine leak waiting to happen.
    • Size buffered channels as semaphores for bounding concurrency (worker pools, rate limiting) instead of letting goroutines fan out unbounded.
    • Use errgroup (or equivalent) for propagating the first error and coordinating cancellation across a group of goroutines, instead of hand-rolling error channels.
    • Treat structured concurrency as the governing principle anyway, even without language support: a goroutine’s lifetime should be scoped to, and never outlive, the function or request that spawned it.
    • Instrument the concurrency itself: race detector in CI, plus metrics and tracing on channel operations and goroutine counts in production because concurrency bugs are nondeterministic and hard to catch with a handful of unit tests.

    Actors: isolation you get structurally

    Actors give you a different, and in some ways weaker, guarantee than CSP but they give it to you structurally. An actor’s state is private, and it processes exactly one message at a time. There is no data race inside one actor, full stop. Instead of “processes synchronizing on named events,” Hewitt’s model says: everything is an actor. Each actor has a private mailbox (unbounded and asynchronous), private state and three things it’s allowed to do on receiving a message: send messages to other actors, create new actors, and decide how to handle its next message. There’s no synchronous handshake requirement anywhere.

    Here’s a worker pool in Erlang, matching the crawler pattern from Part II of this series:

    -module(worker_pool).
    -export([start_pool/1, dispatch/2, worker_loop/1]).
    
    start_pool(N) ->
        [spawn_link(fun() -> worker_loop(0) end) || _ <- lists:seq(1, N)].
    
    worker_loop(Count) ->
        receive
            {work, Job, From} ->
                From ! {result, do_work(Job)},
                worker_loop(Count + 1);
            {status, From} ->
                From ! {count, Count},
                worker_loop(Count)
            % No catch-all clause yet — see the gotcha below
        end.
    
    dispatch(Pid, Job) ->
        Pid ! {work, Job, self()},
        receive
            {result, R} -> R
        after 5000 ->
            {error, timeout}
        end.
    

    Two Erlang-specific gotchas worth knowing before you ship anything like this:

    • Selective receive skips a non-matching message instead of discarding it. receive scans the mailbox in arrival order against your clauses. Anything that matches none of them just sits there, and the next receive call starts scanning from the front all over again. Left unchecked, this is O(n²) behavior over time as junk quietly accumulates. The fix is a catch-all clause:
    worker_loop(Count) ->
        receive
            {work, Job, From} -> ...;
            {status, From} -> ...;
            Other ->
                logger:warning("unexpected message: ~p", [Other]),
                worker_loop(Count)  % drop it, don't let it pile up
        end.
    
    • Mailboxes have no bound by default. ! never blocks in Erlang and there’s no rendezvous. If a producer outpaces a slower worker, the worker’s mailbox just keeps growing until memory runs out. In practice, you either switch to a blocking gen_server:call for anything where backpressure actually matters, or you monitor process_info(Pid, message_queue_len) yourself and shed load manually.

    Supervision is the actor model’s answer to fault structure where a supervisor’s children are linked to it, and a crash triggers a restart strategy instead of taking the whole system down with it. But it is structured fault handling, not structured lifetime tracking. A supervisor doesn’t block waiting for its children to finish instead supervision answers “what happens when a child crashes.” Erlang also provides a location transparency, e.g., an Erlang Pid looks identical whether it points to a local process or one on another node ( Pid ! Msg). But that transparency is syntactic, not operational. A remote send can fail with nodedown or badrpc, latency is never zero so you cannot skip handling the failures that only show up once the mailbox is across a network.

    Where actors are simpler than channels

    A few structural reasons actors tend to feel simpler in practice than channel-based code:

    1. Ownership is enforced by the design itself. There’s no equivalent of “who’s allowed to write to this channel,”, every interaction is a message dropped into a mailbox that only the receiving actor ever drains.
    2. There’s no close semantics to get wrong. Actors don’t have anything like Go’s send-on-closed-panics / double-close-race. An actor’s lifecycle like start, running, terminated is a small, well-understood state machine, and you can monitor/link actor for detecting unexpected crash.
    3. Failure handling is first-class. Supervision trees and let-it-crash give you a systematic answer to “a worker just crashed, now what?” Go’s answer is recover() scattered wherever someone remembered to put it or manual errgroup/context-cancellation wiring to propagate failure to siblings.
    4. Location transparency. With channels, you need to build remoting capability yourself. With actors, it’s often just a deployment decision.

    Where actors are not automatically simpler: mailbox-based concurrency can hide backpressure problems, e.g., an actor with an unbounded mailbox can happily accept messages faster than it processes them and quietly balloon memory. Reasoning about message ordering across several independent actors’ mailboxes is also harder than reasoning about a single shared channel’s FIFO order. CSP’s synchronous rendezvous gives you stronger backpressure for free where an unbuffered send blocks until the receiver is ready (some of modern actor runtimes like Akka support mailbox bounding).

    Best practices for actor systems

    • Bound mailboxes and monitor mailbox depth as a first-class metric, e.g., an unbounded mailbox is the actor world’s version of an unbuffered-channel leak.
    • Design supervision hierarchies deliberately like one-for-one, one-for-all, rest-for-one.
    • Keep actor state small and serializable if you ever want migration or persistence.
    • Use location transparency deliberately, not accidentally.

    CSP/channels fit use cases when you have a fixed, well-understood pipeline topology like stream-processing stages, worker pools with a known fan-out shape. Actors suite when your system’s topology is dynamic like agents spawning agents and where failure isolation matters. I have built PlexSpaces, an actor-based framework, with facets for durability, supervision, and virtual-actor placement based on these lessons. For example, it provides location transparency, failure isolation, and the backpressure/mailbox-bounding. Here is how an actor lifecycle is managed in PlexSpaces:

    async/await

    Async/await never got a formal algebra or expressiveness proof. It’s syntactic sugar over futures and promises, running on a single-threaded event loop or a thread-pool-backed task scheduler. Within one event loop, there’s no preemption between await points, which quietly eliminates a lot of classic race conditions but it introduces its own flavor of the “who’s tracking this” problem:

    async function processOrder(order) {
      sendConfirmationEmail(order); // fire-and-forget — no await!
      return { status: "accepted" };
    }
    

    sendConfirmationEmail here returns a promise nobody is holding onto. If it rejects, nothing catches it and in most runtimes that becomes an unhandled-rejection warning nobody reads. If the process exits before it resolves, it just silently never finishes. Structurally, this is the exact same failure as an unstructured Go goroutine, a unit of work whose lifetime nothing owns. Promise.all and asyncio.gather fix this for the cases you remember to wrap explicitly.

    This is also where the “function coloring” problem lives, which I covered in the ADTs and algebraic effects post: once a single function is async, every caller up the chain has to become async too. Algebraic effects unrelated to CSP’s process algebra are one proposed fix: separate what a function needs from who provides it.

    Linda Memory Model

    Linda coordinates through a shared associative memory called a tuple space, with four operations:

    • out(t): write a tuple, don’t block
    • in(t): block until a tuple matches your template, then atomically remove it
    • rd(t): block until a tuple matches, but leave it there for others
    • eval(t): spawn a computation; its eventual result becomes an ordinary tuple once it finishes

    Neither side of a Linda interaction needs to know who the other one is. A producer can out() a tuple long before any consumer even exists. This is “generative communication”, data just floats in the shared space until something matching comes looking for it:

    // Pseudocode — classic Linda fan-out/fan-in
    for i in 0..n:
        eval(("result", i, compute(i)))    // spawn n concurrent computations
    
    count := 0
    while count < n:
        in(("result", ?i, ?r))              // blocking, destructive, matched by content
        collect(r)
        count += 1
    

    Notice there’s no worker identity anywhere in that collector loop at all. This is suitable for use cases like master/worker fan-out, blackboard-style coordination, barrier synchronization by counting tuples as they arrive. Linda has two gotchas of its own:

    • Which matching tuple you get is unspecified. If two tuples both match your template, the classic Linda spec never says which one in() hands you.
    • eval() returns nothing. No handle, no future, no promise object of any kind. The only way to know a spawned computation ever finished is to already know the shape of its result tuple and read it.

    These issues prevented Linda from going mainstream but its associative memory primitives are natural for coordination related use cases.

    Structured concurrency

    Kotlin’s coroutineScope, Swift’s TaskGroup, and Java 21’s StructuredTaskScope bind a spawned task’s lifetime to the lexical scope that spawned it. The scope literally cannot exit until every child has finished whether error or not. None of the five communication models above give you that by default:

    ModelWhat tracks a spawned unit’s completion
    CSP / GoNothing: go func(){}() has no parent link at all
    Actors / ErlangSupervision restarts a crashed child, but nothing blocks waiting for a healthy one to finish
    async/awaitNothing, an un-awaited promise just runs, or silently fails
    LindaNothing, eval() doesn’t even return a handle to check

    This is exactly why structured concurrency reads as an add-on layer rather than another communication model. It’s a lifetime discipline you can apply on top of channels, actors, promises, or tuples, e.g., Trio applies it to async/await, Kotlin applies it to coroutines that might be built on channels.

    How PlexSpaces answers these gotchas

    Most frameworks inherit one of above models’ specific historical rough edges along with its strengths. PlexSpaces is a actor-and-tuplespace framework I’ve been building, and wrote about in more depth here. It deliberately combines actors and Linda rather than picking one, but it doesn’t reproduce either one’s original sharp edges just for the sake of purity. Here’s the mapping from “gotcha described above” to “the specific fix PlexSpaces makes”:

    Gotcha, as described abovePlexSpaces’ pragmatic answer
    Erlang mailboxes have no bound, so a fast producer can grow one until memory runs outBounded mailboxes. An actor’s inbox has a real, configurable limit, a producer that outpaces its consumer gets backpressure instead of an unbounded memory leak.
    Classic Linda leaves the order of matching tuples unspecified, so which one you get is nondeterministic by designFIFO tuple matching. When more than one tuple matches a template, PlexSpaces returns them in the order they were written, not an arbitrary one, removing nondeterminism-by-specification.
    Full CSP requires learning a process algebra; full Linda requires learning a four-primitive calculus bolted onto a host language; Erlang requires learning OTP’s supervision idiomsOne small API surface. Actors expose a handful of primitives like send, ask, and the tuple-space operations.
    eval() in classic Linda returns no handle, so a spawned computation’s completion is untracked by defaultBecause the “worker” side of a fan-out/fan-in in PlexSpaces is an ordinary supervised actor rather than a bare eval(), its lifetime is owned by a supervisor even though its result is collected the Linda way.
    A crash mid-task loses whatever work was in flight, a concern none of CSP, actors, or Linda’s formulations really addressDurability journaling underneath everything. Messages are journaled at the actor-framework level, below application code, so a crash doesn’t silently lose in-flight work, replay picks the actor back up where it left off.

    A fan-out/fan-in worker pool shows the combination directly:

    // Coordinator spawns N supervised, bounded-mailbox workers.
    for i in 0..n {
        spawn_with_facets(
            &ctx, service_locator.clone(),
            "worker", "default",
            Worker::new(i), vec![],
        ).await?;
    }
    
    // Coordinator collects results Linda-style — associative, FIFO,
    // no ActorRef needed for any individual worker.
    let mut collected = 0;
    while collected < n {
        let tuple = ctx.tuple_space()
            .in_(template!["result", Wildcard, Wildcard])
            .await?;
        collect(tuple);
        collected += 1;
    }
    

    The workers are ordinary supervised actors, restartable, journaled, isolated, with bounded mailboxes so a slow coordinator can’t be flooded. The result collection is Linda-style associative matching, but FIFO instead of unspecified, so results come back in the order the workers actually produced them rather than in an arbitrary one.

    Example: scatter-gather with a timeout, across three models

    Comparisons are easier to trust when they’re concrete rather than abstract, so I implemented the same pattern, scatter-gather with a timeout across all three approaches. The problem: fan requests out to N services, collect the first K responses within a deadline, then cancel everything else, guaranteed. This is the pattern underneath every hedged-request system, every parallel-search aggregator, and every timeout-bounded fan-out you’ve seen in production.

    Go CSP: the naive version leaks goroutines

    The code below shows the mistake almost everyone makes on the first pass like spawning goroutines with no cancellation path at all:

    func ScatterGatherNaive(services []time.Duration, firstK int) []ServiceResponse {
        ch := make(chan ServiceResponse, len(services))
    
        for i, latency := range services {
            go func(id int, lat time.Duration) {
                time.Sleep(lat)
                ch <- ServiceResponse{ServiceID: id, Data: fmt.Sprintf("response-%d", id)}
            }(i, latency)
        }
    
        results := make([]ServiceResponse, 0, firstK)
        for range firstK {
            results = append(results, <-ch)
        }
        // BUG: N-K goroutines still running in background with no cancellation path
        return results
    }
    

    The fix combines context.WithTimeout with errgroup to get a real structured lifetime:

    func ScatterGatherStructured(services []time.Duration, firstK int, timeout time.Duration) []ServiceResponse {
        ctx, cancel := context.WithTimeout(context.Background(), timeout)
        defer cancel()
    
        var mu sync.Mutex
        results := make([]ServiceResponse, 0, firstK)
    
        g, ctx := errgroup.WithContext(ctx)
        for i, latency := range services {
            g.Go(func() error {
                resp, err := simulateService(ctx, i, latency)
                if err != nil { return nil }
                mu.Lock()
                defer mu.Unlock()
                if len(results) < firstK {
                    results = append(results, resp)
                    if len(results) >= firstK { cancel() }
                }
                return nil
            })
        }
        _ = g.Wait() // All goroutines done — structured lifetime guarantee
        return results
    }
    

    errgroup.Wait() guarantees every goroutine finishes before the function returns, and context.WithTimeout propagates cancellation down to the slow workers, so nothing leaks. The catch: you still have to write that plumbing by hand every time.

    Go CSP gotchas, side by side with the CSP algebra:

    GotchaWhat happensCSP algebra equivalent
    Goroutine leakBlocked goroutines run forever, invisible to GC and deadlock detectorN/A
    Nil channel recvBlocks forever – no panic, no warningN/A
    Nil channel sendAlso blocks forever silentlyN/A
    Send on closedRuntime panic – unrecoverable crashN/A
    Recv from closedReturns zero value + ok=false N/A
    Buffered vs unbufferedBreaks rendezvous = breaks the formal reasoning about synchronizationBuffered channels aren’t part of original CSP
    Select non-determinismA random ready case is chosen when several are readyCSP’s external choice (?) is nondeterministic by design

    Runnable, with tests: native/gotchas_test.go — 8 tests demonstrating every row above as failing-then-fixed code.

    // Gotcha: Goroutine leak — spawned workers block forever on an unread channel
    ch := make(chan int)
    for i := 0; i < 10; i++ {
        go func(id int) { ch <- id }(i)  // blocks forever — no reader
    }
    // 10 goroutines leaked: invisible to GC, invisible to deadlock detector
    
    // Gotcha: Nil channel blocks forever (both directions, no panic)
    var ch chan int  // nil
    <-ch            // blocks forever on recv
    ch <- 42        // blocks forever on send
    
    // Gotcha: Send on closed panics, recv from closed returns zero (asymmetric!)
    ch := make(chan int, 1); close(ch)
    ch <- 1         // PANIC: send on closed channel
    v, ok := <-ch   // v=0, ok=false — no panic, just zero value
    
    // Gotcha: Buffered channel breaks rendezvous
    buffered := make(chan int, 5)
    buffered <- 1   // sender proceeds without receiver — not CSP anymore
    

    Pure Rust: a Nursery, and select! as a guarded command

    Rust gives you ownership-enforced isolation without shared state by default but it doesn’t give you structured lifetime by default either. The Nursery type below binds a group of spawned tasks to a scope explicitly:

    pub struct Nursery<T: Send + 'static> {
        join_set: JoinSet<T>,
    }
    
    impl<T: Send + 'static> Nursery<T> {
        pub fn spawn<F>(&mut self, future: F)
        where F: Future<Output = T> + Send + 'static {
            self.join_set.spawn(future);
        }
    
        /// Wait for first K tasks OR timeout — then cancel everything else.
        pub async fn wait_first_k_or_timeout(mut self, k: usize, timeout: Duration) -> Vec<T> {
            let mut results = Vec::with_capacity(k);
            let deadline = tokio::time::Instant::now() + timeout;
            loop {
                if results.len() >= k { break; }
                tokio::select! {
                    maybe = self.join_set.join_next() => {
                        match maybe {
                            Some(Ok(value)) => results.push(value),
                            Some(Err(_)) => continue,
                            None => break,
                        }
                    }
                    _ = tokio::time::sleep_until(deadline) => { break; }
                }
            }
            self.join_set.abort_all(); // Structured cleanup — cancel stragglers
            results
        }
    }
    

    tokio::select! maps almost directly onto CSP’s guarded command / external choice operator, whichever branch is ready first wins, and the biased; modifier gives you deterministic priority ordering (unlike Go’s deliberately random tie-break):

    let results = scatter_gather_csp(&services, 3, Duration::from_millis(300)).await;
    // Nursery guarantees: all children cancelled before scope exits
    

    Ownership prevents shared-state bugs entirely, at compile time. JoinSet::abort_all() gives you a real, guaranteed cancellation. The nursery pattern gets you structured lifetime with essentially no runtime overhead. The catch: there’s no nursery built into the standard library and there’s still no formal algebra backing any of it (no FDR-style prover checking).

    Rust CSP gotcha, demonstrated in csp_channels/src/scatter_gather.rs:

    // Gotcha: Naive tokio::spawn has no parent link — tasks leak
    let mut handles = vec![];
    for svc in &services {
        handles.push(tokio::spawn(simulate_service(svc)));
    }
    // If we return early, spawned tasks run forever — no cancellation
    
    // Fix: JoinSet provides structured lifetime
    let mut join_set = JoinSet::new();
    for svc in &services { join_set.spawn(simulate_service(svc)); }
    // On drop or abort_all(), all tasks are cancelled — guaranteed
    

    PlexSpaces: supervised lifetime plus decoupled collection

    This is where actor supervision (structured fault handling) and Linda-style tuple-space coordination (decoupled result collection) get combined. Start with thin Linda-style wrappers over the tuple-space host functions:

    // Linda-style thin wrappers over tuplespace host functions
    fn linda_out(fields: &[Value]) -> Result<(), String> {
        let request = WriteRequest { tuples: vec![json_to_tuple(fields)?], .. };
        ts_write(&request.encode_to_vec()).map(|_| ())
    }
    
    fn linda_in(pattern: &[Value]) -> Result<Option<Vec<Value>>, String> {
        let request = ReadRequest { template: Some(to_pattern(pattern)?), take: true, .. };
        let bytes = ts_take(&request.encode_to_vec())?;
        Ok(decode_response(&bytes)?.first().map(to_json_array))
    }
    

    The orchestrator scatters by spawning supervised workers, then sets its own timeout with a self-message:

    // Scatter: spawn N workers under supervisor
    for i in 0..num_services {
        spawn("actor-csp-wasm", &format!("worker-{i}"), "", &init_json)?;
        send(&worker_id, "cast", &work_payload)?;
    }
    
    // Set timeout — send_after fires a collection message to self
    send_after(timeout_ms, "cast", &collect_msg)?;
    

    Each worker writes its result to the shared tuple space, with zero knowledge of who’s collecting it:

    // Worker: Linda OUT — write result tuple to shared tuplespace
    linda_out(&[
        Value::String("result".into()),
        Value::String(request_id.into()),
        Value::Number(service_id.into()),
        Value::String(result_data),
    ])?;
    

    And gathering reads back whatever arrived in time, then explicitly tells the supervisor to stop the rest:

    // Gather: Linda RD-ALL — collect whatever arrived before timeout
    let results = linda_rd_all(&["result", request_id, *, *])?;
    // Structured cleanup: stop remaining workers via supervisor
    for wid in &worker_ids { stop(wid)?; }
    

    Workers never need to know who’s collecting their results, that’s the Linda decoupling doing its job. The supervisor guarantees the worker lifecycle end to end, e.g., a crashed worker restarts automatically under a OneForOne strategy. Bounded mailboxes keep a slow coordinator from getting flooded. FIFO tuple matching makes the collection step deterministic instead of an open question.

    The three approaches, side by side

    PropertyGo CSP (errgroup)Rust (Nursery/select!)PlexSpaces (actors + Linda)
    Cancellationcontext.Cancel() propagatedJoinSet::abort_all()stop() via supervisor
    Structured lifetimeg.Wait() blocksNursery scope exitSupervisor manages lifecycle
    BackpressureBuffered channel capacityChannel capacityBounded mailbox
    Failure handlingerrgroup collects first errorJoinError on abortSupervisor restarts crashed worker
    CouplingWorkers know the result channelWorkers know the result typeWorkers only know the tuple shape (Linda)
    Formal backingNoneNoneNone (but FIFO + bounded mailboxes removes two classes of nondeterminism)
    DistributionSingle process onlySingle process onlyMulti-node, transparently

    Full runnable examples with tests live at: examples/rust/embedded/csp_channels, examples/go/apps/csp_structured/native, and examples/rust/apps/actor_csp.

    One more actor-model gotcha, demonstrated via supervisor behavior

    // Gotcha: Unbounded mailbox — fast producer OOMs the consumer
    // Fix: PlexSpaces uses bounded mailboxes with a configurable limit
    
    // Gotcha: No structured lifetime — actors are async, no scope to wait on
    // Fix: Supervisor + explicit stop() for child actors after collection
    
    // Gotcha: Orphaned actors — a spawned actor runs forever if nobody stops it
    // Fix: OneForOne supervisor manages worker lifecycle; orchestrator calls stop()
    

    The tldr;

    • CSP has real algebra and real tooling (FDR) behind it. Go borrows the vocabulary but drops the proof the when you add a buffer.
    • Actors have partial formal treatment and isolation by construction, but unbounded mailboxes and selective-receive skip are real, sharp edges the model doesn’t protect you from on its own.
    • async/await never had formal backing at all, and its fire-and-forget promise is the exact same “who’s tracking this” bug as an unstructured goroutine.
    • Linda is the most decoupled model on this list and the least adopted. The original spec leaves match order unspecified and spawned work untracked.
    • Structured concurrency isn’t a another concurrency model, instead it’s a lifetime discipline layered.
    • PlexSpaces’ bet is that you don’t have to inherit every historical rough edge along with the good ideas like bounded mailboxes, FIFO matching, and one small unified API let it combine actor supervision with Linda’s decoupled coordination.

    The rest of the series

    1. Part I — the general problem, concurrency constructs, and TypeScript
    2. Part II — Erlang and Elixir
    3. Part III — Go and Rust
    4. Part IV — Kotlin and Swift
    5. Building a Durable Actor Framework for Polyglot Serverless Apps
    6. 20+ Production Patterns for Distributed AI Agents Using Actors and TupleSpaces
    7. Building an Agent Harness and Eval Pipeline with Durable Actors
    8. Building a Self-Improving AI Agent with Durable Actors: MiniHermes
    9. Building Mini OpenClaw: Secure AI Agents with Actors, WASM, and Supervision
    10. Making Bad State Impossible: A Practical Guide to ADTs and Algebraic Effects
    11. Building PlexSpaces: Decades of Distributed Systems Distilled Into One Framework

    Code for everything above: github.com/bhatti/PlexSpaces

    July 22, 2026

    Migrating Off Cloudflare Durable Objects: Build Your Own Portable FAAS

    Filed under: Computing — admin @ 8:33 pm

    Introduction

    Every major cloud now supports Serverless FAAS capabilities like AWS Lambda/Step Functions, Azure Durable Functions, GCP Cloud Functions and Cloudflare Durable Objects where you write a function or a small stateful actor. This allows you to scale it, pay only for what runs but there is a catch, you build on a proprietary runtime, and the runtime’s storage model, invocation model, and IAM rules become part of your application whether you meant them to or not. You cannot easily rewrite it or run it somewhere else. I saw a recent post (Why we’re moving Wire off Cloudflare Durable Objects) from Wire, which ran every container on Cloudflare Durable Objects since day one. They wrote why they rebuilt their own data plane instead of staying, which included extra network hops on the hottest path, drift of state, separation of compute from data, rigid placement policies and lack of self-hosting. None of these are reliability complaints, instead they’re architectural ceilings baked into a runtime you don’t own. And the pattern generalizes past Cloudflare:

    • AWS Lambda: SAM and LocalStack approximate the runtime locally, but diverge on execution environment, IAM, and VPC behavior.
    • Azure Durable Functions: Azurite emulates the storage layer, but the replay-based orchestration engine behaves differently under real concurrent load than it does in the emulator.
    • GCP Cloud Functions: the Functions Framework runs locally, but Eventarc, Pub/Sub push, and Cloud Run triggers all need live GCP resources.
    • Cloudflare Workers/DO/Agents: wrangler dev simulates KV with SQLite and alarms with in-process timers, but never replicates the distributed routing that decides which data center actually holds your object.

    Every one of these runtimes gives you a great abstraction and takes your operational sovereignty in exchange. This post shows how to keep the abstraction like stateful actors, durable storage, alarms, WASM sandboxing, LLM calls, observability, an event bus, webhooks while running it on infrastructure you control with an open-source framework called PlexSpaces.


    PlexSpaces

    PlexSpaces is an open-source, polyglot actor framework that gives every actor durable KV storage and alarms, routes messages between actors on one node or across a gRPC mesh, and exposes a host API to Go, TypeScript, Python, and Rust. You write an actor once; the same WASM module deploys to your laptop, a Docker container, Kubernetes, bare metal, or several clouds at once, unchanged. The mental model sits close enough to Cloudflare Durable Objects that migrating existing DO code is mostly mechanical. The key difference: there’s no simulated version to diverge from production, because the production runtime is the development runtime.


    Part I: Durable Objects

    Cloudflare publishes a short list of rules for writing correct Durable Objects, which are easily mapped to a PlexSpaces constraint, and for most of them the mapping is tighter:

    Cloudflare RuleHow PlexSpaces handles it
    Don’t coordinate between objects from inside an object: use async messaginghost.send(actorId, op, payload) is the only cross-actor primitive. There is no shared-memory path on the same node.
    Don’t assume a single instance: globally unique, but workers can race to create oneActor IDs are content-addressed: {name}//{type}::{ns}@nodeId. The node that owns the ID wins.
    Store state before returning: in-memory state is lost on evictiongetState()/setState() checkpoints on every handler return. Durable KV is secondary store. Survive eviction and restart.
    Keep objects small: large objects cause cold-start latencyWASM heap is the actor’s private address space, isolated from the host. State serializes only on checkpoint.
    Use blockConcurrencyWhile() for initonInit() in TypeScript / @init_handler in Python / Init() in Go runs before the first message and restore persisted state.
    Alarm fires at-most-onceReminderFacet persists the alarm timestamp in durable storage and re-queues after restart.

    The one meaningful difference: Cloudflare guarantees globally unique placement. PlexSpaces virtual actors are unique per node or per cluster when pinned with @* (any node) or @nodeId (specific node). PlexSpaces provides an object-registry for managing cross-cluster deduplication.


    Durable KV Storage

    Cloudflare’s ctx.storage gives you a transactionally consistent store scoped to one object:

    // Cloudflare DO — JavaScript
    const count = await this.ctx.storage.get("count") ?? 0;
    await this.ctx.storage.put("count", count + 1);
    
    // Batch operations (storage.get([keys]) / storage.put(map))
    const [history, meta] = await this.ctx.storage.get(["room:history", "room:meta"]);
    await this.ctx.storage.put(new Map([["room:history", data], ["room:meta", index]]));

    PlexSpaces gives you the same guarantee through host.kv. Because each actor processes one message at a time, a plain read-modify-write is safe without extra locking. The migrating_cloudflare_workers TypeScript example restores room history on onInit() using batch KV similar to blockConcurrencyWhile:

    // PlexSpaces TypeScript — examples/typescript/apps/migrating_cloudflare_workers/guild_chat_actor.ts
    protected override onInit(config: Record<string, unknown>): void {
        this.state.room_id = String(config.actor_id ?? "");
    
        // blockConcurrencyWhile() equivalent — batch-fetch persisted state before first message
        const keys = [
            "room:" + this.state.room_id + ":history",
            "room:" + this.state.room_id + ":meta",
        ];
        const values = host.kv.multiGet(keys);   // like DO storage.get(["k1","k2"])
        const [historyRaw] = values;
        if (historyRaw) {
            const msgs = JSON.parse(historyRaw);
            if (Array.isArray(msgs)) {
                this.state.messages = msgs;
                this.state.message_seq = msgs[msgs.length - 1]?.seq ?? 0;
            }
        }
    }
    
    private persistHistory(): void {
        // Batch write — like DO storage.put(new Map([["k1",v1],["k2",v2]]))
        host.kv.multiPut({
            ["room:" + this.state.room_id + ":history"]: JSON.stringify(this.state.messages),
            ["room:" + this.state.room_id + ":meta"]: JSON.stringify({
                message_seq: this.state.message_seq,
                last_updated: host.nowMs(),
            }),
        });
    }
    # PlexSpaces Python — examples/python/apps/migrating_cloudflare_workers/guild_chat.py
    def _load_history(self) -> None:
        room_id = self._room_id()
        index_raw = host.kv.get(f"room:{room_id}:seq_index")
        if index_raw:
            seqs = json.loads(index_raw)
            keys = [f"room:{room_id}:msg:{seq}" for seq in seqs]
            values = host.kv.multi_get(keys)   # DO storage.get([k1,k2,...]) equivalent
            self.messages = [json.loads(v) for v in values if v]
            if self.messages:
                self.msg_seq = self.messages[-1]["seq"]
    
    def _persist_history(self) -> None:
        room_id = self._room_id()
        entries = {}
        seqs = []
        for msg in self.messages:
            entries[f"room:{room_id}:msg:{msg['seq']}"] = json.dumps(msg)
            seqs.append(msg["seq"])
        entries[f"room:{room_id}:seq_index"] = json.dumps(seqs)
        host.kv.multi_put(entries)   # DO storage.put({k:v,...}) equivalent

    PlexSpaces also ships atomic operations that Cloudflare leaves you to build yourself with a coordinator object:

    // PlexSpaces TypeScript — from RateLimiterActor in guild_chat_actor.ts
    // Atomic distributed counter — equivalent to Cloudflare KV atomic increment
    const distributedCount = host.kv.increment(windowKey, 1);
    
    // Compare-and-swap — equivalent to DO transactional storage read-modify-write
    const applied = await host.kvCas("lock_key", expectedValue, newValue);
    
    // KV with TTL — equivalent to DO storage.put with metadata expiration
    await host.kvPutWithTtl("session_token", token, 3600);

    The Python RateLimiterActor shows both in context:

    # PlexSpaces Python — examples/python/apps/migrating_cloudflare_workers/guild_chat.py
    @handler("check")
    def check(self, user_id: str = "") -> dict:
        # Atomic increment — survives actor restarts, equivalent to Cloudflare KV atomic
        host.kv.increment(f"rate:{user_id}:total", 1)
    
        # CAS — idempotent token slot reservation, like DO transactional put
        cas_key = f"rate:{user_id}:window"
        current_val = host.kv.get(cas_key) or ""
        host.kv.cas(cas_key, current_val, str(host.now_ms()))
    
        # Token bucket logic operates on in-memory state (safe: single-threaded actor)
        allowed = self.buckets[user_id]["tokens"] > 0
        if allowed:
            self.buckets[user_id]["tokens"] -= 1
        return {"allowed": allowed, "remaining": self.buckets[user_id]["tokens"]}

    Durable Alarms

    A DO schedules one future callback that survives node restarts:

    // Cloudflare DO
    await this.ctx.storage.setAlarm(Date.now() + 10_000);
    
    async alarm() {
        const count = await this.ctx.storage.get("count");
        await this.ctx.storage.delete("count");
        console.log(`Processing ${count} batched requests`);
    }

    PlexSpaces maps this directly through ReminderFacet, which persists the alarm to durable storage and re-queues it automatically after a restart. The AlarmDemoActor in the guild-chat example demonstrates the full lifecycle like set, query, fire, cancel:

    // PlexSpaces TypeScript — examples/typescript/apps/migrating_cloudflare_workers/guild_chat_actor.ts
    onEnqueue(_payload: Record<string, unknown>): Record<string, unknown> {
        this.state.queued++;
        if (this.state.queued === 1) {
            // First item — equivalent to: await this.state.storage.setAlarm(Date.now() + 30_000)
            host.alarm.set(host.nowMs() + 30_000);
        }
        return { status: "ok", queued: this.state.queued };
    }
    
    // Fires when the scheduled timestamp is reached
    // Equivalent to Cloudflare DO: async alarm() { ... }
    on__alarm__(_payload: Record<string, unknown>): Record<string, unknown> {
        const processed = this.state.queued;
        this.state.processed += processed;
        this.state.queued = 0;
        this.state.total_alarm_fires++;
        return { status: "ok", processed };
    }
    
    onStatus(_payload: Record<string, unknown>): Record<string, unknown> {
        const alarmAt = host.alarm.get();   // DO: await this.state.storage.getAlarm()
        return { ...this.state, alarm_at: alarmAt, alarm_set: alarmAt > 0 };
    }
    # PlexSpaces Python — examples/python/apps/migrating_cloudflare_workers/guild_chat.py
    @handler("start")
    def start(self, delay_ms: int = 30000) -> dict:
        # Equivalent to: this.state.storage.setAlarm(Date.now() + delay_ms)
        host.alarm.set(host.now_ms() + delay_ms)
        return {"status": "ok", "fire_at_ms": host.now_ms() + delay_ms}
    
    @handler("__alarm__")
    def on_alarm(self) -> dict:
        # Equivalent to Cloudflare DO: async alarm() { ... }
        processed = len(self.pending_requests)
        self.total_processed += processed
        self.pending_requests = []
        return {"status": "ok", "processed": processed}
    
    @handler("cancel")
    def cancel(self) -> dict:
        host.alarm.delete()   # DO: this.state.storage.deleteAlarm()
        return {"status": "ok", "action": "alarm_cancelled"}

    Get-or-Create (Virtual Actors)

    Cloudflare’s core abstraction is the globally unique object that appears on first access, at the cost of a binding declared in wrangler.toml:

    // Cloudflare DO — needs a binding in wrangler.toml
    const id = env.CHAT_ROOM.idFromName(roomId);
    const room = env.CHAT_ROOM.get(id);
    await room.fetch("/send", { method: "POST", body: JSON.stringify(msg) });

    PlexSpaces virtual actors give you the same behavior with no binding file. The actor ID ({name}//{actorType}::{namespace}@*) encodes both the shard key and the actor class, and the @* suffix lets the runtime place it on whichever node is best; pin it with @node-id for data locality:

    // PlexSpaces TypeScript
    import { getActorRef } from "@plexspaces/sdk";
    const room = getActorRef("ChatRoomActor", roomId, "default");
    const reply = await room.ask("send_message", { user_id: userId, content: msg });
    # PlexSpaces Python
    from plexspaces.host import get_actor_ref
    room = get_actor_ref("ChatRoomActor", room_id, "default")
    reply = room.ask("send_message", {"user_id": user_id, "content": msg}, timeout_ms=5000)

    The actor spins up on its first message, whether or not it existed before.


    Listing Actors by Namespace

    Cloudflare exposes a Durable Objects namespace list API for management and observability. PlexSpaces has a direct equivalent defined in its proto-first design. The ListActors RPC in actor_runtime.proto accepts namespace, actor_type, state, and node_id filters and returns paginated results:

    // proto/plexspaces/v1/actors/actor_runtime.proto
    message ListActorsRequest {
        string actor_type = 3;
        ActorState state  = 4;
        string node_id    = 5;
        // Namespace for tenant isolation — only actors in this namespace are returned
        string namespace  = 6;
        PageRequest page_request = 2;
    }
    
    message ListActorsResponse {
        repeated Actor actors        = 2;
        PageResponse page_response   = 3;
    }

    From a PlexSpaces client, listing all active ChatRoomActor instances in the default namespace:

    # HTTP API (equivalent to Cloudflare's list-objects endpoint)
    curl "http://localhost:8080/api/v1/actors/default/ChatRoomActor?state=active"
    // Go SDK
    actors, err := client.ListActors(ctx, &ListActorsRequest{
        ActorType: "ChatRoomActor",
        Namespace: "default",
        State:     ActorStateActive,
    })

    Cloudflare limits listing to metadata (ID, location, storage size). PlexSpaces returns full Actor records including state, facets, node assignment, resource usage, and tenant/namespace tags.


    WebSocket Handling

    Cloudflare’s Model

    Cloudflare gives the Durable Object two WebSocket modes. In the standard model, the object holds the socket directly:

    // Cloudflare DO — standard WebSocket
    async fetch(request) {
        const [client, server] = Object.values(new WebSocketPair());
        this.ctx.acceptWebSocket(server);
        return new Response(null, { status: 101, webSocket: client });
    }
    
    async webSocketMessage(ws, message) {
        for (const peer of this.ctx.getWebSockets()) {
            peer.send(`broadcast: ${message}`);
        }
    }
    
    async webSocketClose(ws, code) {
        ws.close(code, "connection closed");
    }

    The WebSocket Hibernation API (state.acceptWebSocket / getWebSockets()) is Cloudflare’s optimization for objects that hold many sockets but are mostly idle: the object is evicted when no message is being processed, and WebSocket state is restored from durable storage on the next message.

    PlexSpaces: A Cleaner Split

    PlexSpaces takes a different approach: room state and connection state are in separate actors. A ChatRoomActor holds only membership and message history; each browser connection is a thin-node client registered under its own actor ID. When the room fans out, host.send(actorId, "chat_message", event) routes each delivery through the WsActorTransportClient to the right WebSocket session like the room never holds a socket handle.

    This is the same split Discord uses internally in its Elixir stack, where session processes are separate from guild processes. Here’s the real onSend handler from examples/typescript/apps/ws_chat_room/:

    // PlexSpaces TypeScript — examples/typescript/apps/ws_chat_room/chat_server_actor.ts
    onSend(payload: SendPayload): unknown {
        const senderUsername = this.state.members[payload.sender_actor_id] ?? payload.sender_actor_id;
        const ts = host.nowMs();
        this.state.history.push({
            senderActorId: payload.sender_actor_id,
            sender: senderUsername,
            text: payload.text,
            ts,
        });
        if (this.state.history.length > MAX_HISTORY) {
            this.state.history = this.state.history.slice(-MAX_HISTORY);
        }
    
        const event = {
            sender: payload.sender_actor_id,
            sender_username: senderUsername,
            text: payload.text,
            room_id: this.state.roomId,
            ts,
        };
        // Fan out to all members including sender (delivery confirmation)
        // host.send() routes each tell through WsActorTransportClient ? WsRegistry ? thin-node WS session
        const memberIds = Object.keys(this.state.members);
        for (const actorId of memberIds) {
            host.send(actorId, "chat_message", event);
        }
        return { success: true, members_notified: memberIds.length };
    }

    A companion PresenceActor in the same file tracks online/offline state using a durable reminder (host.sendAfter) to mark a user offline after 55 seconds of silence without external cron job involved:

    // PlexSpaces TypeScript — examples/typescript/apps/ws_chat_room/chat_server_actor.ts
    onOnline(_payload: OnlinePayload): unknown {
        this.state.online = true;
        this.state.last_seen = host.nowMs();
        host.kv.putJson(`presence:${this.state.userId}`, { online: true, last_seen: this.state.last_seen });
        host.sendAfter(60_000, "timeout_check", {});   // durable reminder, survives restart
        return { success: true, online: true };
    }
    
    onTimeout_check(): unknown {
        const idleSince = host.nowMs() - this.state.last_seen;
        if (idleSince > 55_000) {
            this.state.online = false;
            host.kv.putJson(`presence:${this.state.userId}`, {
                online: false, last_seen: this.state.last_seen,
            });
        }
        return { checked: true, idle_ms: idleSince };
    }

    Gap Analysis: WebSocket Hibernation

    FeatureCloudflare DOPlexSpaces
    Accept WebSocket in objectctx.acceptWebSocket(ws)Thin-node client; actor never holds socket
    Per-socket tagsctx.acceptWebSocket(ws, tags)Actor ID is the tag – lookup by actor ID
    Get sockets by tagctx.getWebSockets(tag)host.send(actorId, ...) routes directly
    Evict during idleHibernation API (cost optimization)Actor checkpoints and can be evicted; reconnect restores via getState()
    Socket-level error handlingwebSocketError(ws, err)Session actor handles disconnect
    Outgoing WebSocket from DOnew WebSocket(url) in fetchhost.httpClient("link").fetch(...) or service link for outbound calls

    The PlexSpaces model costs more code upfront (session actor + room actor) and pays you back with independent scaling: fan-out scales with the number of receivers, not with the room actor’s memory; a crashed session doesn’t lock the room; reconnect logic is client-side only.


    Fan-Out: DO’s Missing host.send() Primitive

    In Cloudflare, cross-object fan-out means individual fetch() calls or Queues, which are not lean as a fire-and-forget tell. PlexSpaces’s host.send() is a fire-and-forget message routed through the actor mesh, no HTTP overhead, with in-process delivery for co-located actors. The Python guild-chat send_message handler shows this clearly:

    # PlexSpaces Python — examples/python/apps/migrating_cloudflare_workers/guild_chat.py
    @handler("send_message")
    def send_message(self, user_id: str = "", content: str = "") -> dict:
        msg = self._add_message(user_id, content, host.now_ms())
    
        # Fan-out: fire-and-forget to each member actor
        # Mirrors Discord's Manifold pattern for distributed fan-out
        # In Cloudflare DO, this would be individual fetch() calls — expensive
        fan_out_count = 0
        for member_id in list(self.members.keys()):
            if member_id != user_id:
                host.send(member_id, "receive_message", {
                    "room_id": self._room_id(),
                    "seq": msg["seq"],
                    "from": user_id,
                    "content": content,
                })
                fan_out_count += 1
    
        self._persist_history()   # batch multiPut — one KV call for the whole room
        return {"status": "ok", "seq": msg["seq"], "fan_out": fan_out_count}

    Part II: Cloudflare Agents SDK

    Cloudflare’s Agents SDK builds stateful AI agents on top of Durable Objects. PlexSpaces covers the same patterns; some are direct translations and a few require a different shape.

    Conversation State and Memory

    Cloudflare stores conversation history in the DO’s storage. PlexSpaces does the same through host.kv, with the same durability guarantee: the ChatAgentActor in examples/python/apps/chat_agent/ stores history under a well-known key and restores it across activations:

    # PlexSpaces Python — examples/python/apps/chat_agent/chat_agent.py
    @handler("chat")
    def chat(self, message: str = "") -> dict:
        # Load history — equivalent to: await this.storage.get('history')
        history = host.kv.get_json("history") or []
    
        history.append({"role": "user", "content": message, "timestamp": host.now_ms()})
    
        assistant_reply = self._call_llm(history)
    
        history.append({"role": "assistant", "content": assistant_reply, "timestamp": host.now_ms()})
    
        # Persist — equivalent to: await this.storage.put('history', history)
        host.kv.put_json("history", history)
        self.total_messages += 1
    
        # Schedule summarization alarm once history is long enough
        if len(history) > _ALARM_THRESHOLD and host.alarm.get() == 0:
            host.alarm.set(host.now_ms() + _ALARM_DELAY_MS)
    
        return {"status": "ok", "reply": assistant_reply, "history_length": len(history)}
    
    @handler("__alarm__")
    def on_alarm(self) -> dict:
        # Durable alarm callback — equivalent to Cloudflare Agents SDK onAlarm()
        history = host.kv.get_json("history") or []
        summary = self._call_llm([{
            "role": "user",
            "content": f"Summarize this conversation (2-3 sentences): {json.dumps(history)}"
        }])
        host.kv.put("summary", summary)
        host.kv.delete("history")   # clear after summarizing
        return {"status": "ok", "action": "summarized", "messages_summarized": len(history)}

    For long-term cross-session memory, store summaries under a user-scoped key (memory:{user_id}) and inject them into the next conversation’s system prompt similar to Cloudflare’s getMemory/setMemory implementation.


    Calling LLMs

    Cloudflare’s Agents SDK routes every call through env.AI, Cloudflare’s own inference gateway, locked to providers they support:

    // Cloudflare Agents SDK — locked to Cloudflare's AI gateway
    const response = await this.env.AI.run("@cf/meta/llama-3-8b-instruct", { messages });

    PlexSpaces actors call any provider through a named HTTP service link resolved at deploy time, so the actor code never mentions a specific vendor:

    # PlexSpaces Python — examples/python/apps/chat_agent/chat_agent.py
    def _call_llm(self, messages):
        http = ServiceHttpClient("llm-link")
        body = {
            "model": "claude-3-5-haiku-20241022",
            "max_tokens": 1024,
            "messages": [{"role": m["role"], "content": m["content"]} for m in messages],
        }
        resp = http.post("/v1/messages", body)
        # Parse Anthropic response
        if isinstance(resp, dict):
            content = resp.get("content", [])
            if content and isinstance(content, list):
                return content[0].get("text", "")
        return "[LLM unavailable]"

    app-config.toml points the link at Ollama locally, Anthropic or OpenAI in production, or an internal AI gateway in a regulated environment:

    # Local development — Ollama
    [[service_links]]
    name = "llm-link"
    url  = "http://localhost:11434"
    
    # Production — swap without touching actor code
    [[service_links]]
    name = "llm-link"
    url  = "https://api.anthropic.com"
    headers = { "x-api-key" = "${ANTHROPIC_API_KEY}" }

    TypeScript actors use the same pattern:

    // PlexSpaces TypeScript — examples/typescript/apps/chat_agent/
    const resp = host.httpClient("llm-link").post("/v1/messages", {
        model: "claude-3-5-haiku-20241022",
        max_tokens: 1024,
        messages,
    });

    Durable Workflows

    Cloudflare Agents SDK ships workflow primitives that checkpoint multi-step sequences. PlexSpaces WorkflowActor gives you the same like Run, Signal, and Query RPCs map to start, inject external events, and inspect state. For example, the payment workflow in examples/go/apps/migrating_cadence/payment_workflow.go shows the shape:

    // PlexSpaces Go — examples/go/apps/migrating_cadence/payment_workflow.go
    // PaymentWorkflow implements WorkflowActor for idempotent payment processing.
    // Steps: validate ? authorize (with retry) ? capture ? settle.
    // Signals: refund, cancel.  Queries: status, payment_id.
    type PaymentWorkflow struct {
        plexspaces.BaseActor
        PaymentID       string        `json:"payment_id"`
        Status          string        `json:"status"` // pending ? validated ? authorized ? captured ? settled
        Steps           []PaymentStep `json:"steps"`
        RefundRequested bool          `json:"refund_requested"`
    }
    
    func (p *PaymentWorkflow) Run(payloadJSON string) string {
        // Each step checkpoints via getState/setState before proceeding.
        // If the node crashes mid-run, the workflow resumes from the last checkpoint.
        p.Status = "validating"
        if err := p.validatePayment(); err != nil {
            p.Status = "failed"
            return marshal(map[string]any{"error": err.Error()})
        }
        p.addStep("validate")
    
        // Authorize with retries (idempotency key prevents double-charge)
        p.Status = "authorizing"
        for attempt := 0; attempt < 3; attempt++ {
            if err := p.authorizePayment(); err == nil {
                break
            }
        }
        p.addStep("authorize")
        // ... capture, settle
        return marshal(map[string]any{"status": p.Status, "payment_id": p.PaymentID})
    }
    
    func (p *PaymentWorkflow) Signal(name, _ string) {
        switch name {
        case "refund":
            p.RefundRequested = true
        case "cancel":
            p.Status = "cancelled"
        }
    }
    
    func (p *PaymentWorkflow) Query(name, _ string) string {
        return marshal(map[string]any{"status": p.Status, "payment_id": p.PaymentID})
    }

    For multi-agent orchestration, OrchestratorActor in examples/go/apps/miniclaw/ decomposes a task into sub-tasks, delegates each to a worker agent discovered via process group, and aggregates results through TupleSpace:

    // PlexSpaces Go — examples/go/apps/miniclaw/orchestrator.go
    func (o *OrchestratorActor) Run(payloadJSON string) string {
        task := stringVal(parsePayload(payloadJSON), "task", "")
        taskID := fmt.Sprintf("orch-%d", host.NowMs())
    
        o.Status = "running"
        o.TaskID = taskID
    
        // Discover available agents via process group membership
        agentID, err := pgFirst("svc:agent")
        if err != nil {
            return marshal(map[string]any{"error": "no agents in svc:agent process group"})
        }
    
        // Decompose and delegate sub-tasks
        subTasks := decompose(task)
        for i, subTask := range subTasks {
            o.Progress = (i + 1) * 100 / len(subTasks)
            result, err := host.Ask(agentID, "chat", map[string]any{
                "message":    subTask,
                "session_id": fmt.Sprintf("orch-%s-%d", taskID, i),
            }, 30000)
            if err != nil {
                return marshal(map[string]any{"error": "sub-task failed: " + err.Error()})
            }
            // Store result in TupleSpace for aggregation
            host.TS().Write([]any{"orch_result", taskID, i, result})
        }
    
        o.Status = "completed"
        return marshal(map[string]any{"task_id": taskID, "status": "completed", "sub_tasks": len(subTasks)})
    }

    Human-in-the-Loop

    Cloudflare’s human-in-the-loop pattern pauses an agent on a high-stakes action and waits for external approval before resuming. PlexSpaces implements this natively through a GenFSM actor, e.g., ApprovalGateActor in examples/python/apps/minipi/approval_gate.py:

    # PlexSpaces Python — examples/python/apps/minipi/approval_gate.py
    @fsm_actor(states=["idle", "awaiting_approval", "approved", "rejected"], initial="idle")
    class ApprovalGateActor:
        """
        FSM states: idle ? awaiting_approval ? approved / rejected ? idle
    
        Key insight: the agent can wait for days. DurabilityFacet preserves all state
        durably — no polling, no timeouts burning tokens.
        """
        fsm_state: str = state(default="idle")
        pending_request: dict = state(default_factory=dict)
        pending_agent_id: str = state(default="")
    
        @handler("request_approval")
        def request_approval(self, agent_id: str = "", action: str = "", context: dict = None) -> dict:
            """An agent requests human approval for a high-stakes action."""
            if self.fsm_state != "idle":
                return {"status": "busy", "current_agent": self.pending_agent_id}
    
            self.fsm_state = "awaiting_approval"
            self.pending_agent_id = agent_id
            self.pending_request = {"action": action, "context": context or {}, "requested_at_ms": host.now_ms()}
    
            # Store request for external review (dashboard, Slack notification, etc.)
            host.kv.put(f"approval_request:{self.actor_id}", json.dumps(self.pending_request))
    
            return {"status": "pending", "gate_id": self.actor_id}
    
        @handler("approve")
        def approve(self, approver: str = "", comment: str = "") -> dict:
            """Human approves — signals the suspended agent to resume."""
            agent_id = self.pending_agent_id
            self.fsm_state = "approved"
            self.decision_history.append({
                "action": self.pending_request.get("action"),
                "decision": "approved",
                "approver": approver,
                "decided_at_ms": host.now_ms(),
            })
    
            # Signal the waiting agent to resume with the decision
            host.send(agent_id, "workflow_signal:resume", {
                "decision": "approved",
                "approver": approver,
                "comment": comment,
            })
    
            self.fsm_state = "idle"
            self.pending_agent_id = ""
            return {"status": "approved", "agent_id": agent_id}
    
        @handler("reject")
        def reject(self, approver: str = "", reason: str = "") -> dict:
            """Human rejects — signals the agent with the rejection."""
            agent_id = self.pending_agent_id
            host.send(agent_id, "workflow_signal:resume", {
                "decision": "rejected",
                "approver": approver,
                "reason": reason,
            })
            self.fsm_state = "idle"
            return {"status": "rejected", "agent_id": agent_id}

    The agent on the other side calls host.ask("approval_gate", "request_approval", {...}) then processes the workflow_signal:resume message when it arrives. Because state is checkpointed durably, the agent can wait hours or days with no polling loop and no timeout burning tokens.


    Long-Running Agents

    Cloudflare’s long-running agent pattern uses alarms to wake a dormant agent on a schedule. PlexSpaces handles this identically, e.g., any actor with the ReminderFacet can schedule work arbitrarily far in the future. The ChatAgentActor summarization alarm is one example; for a true long-running loop:

    @actor
    class ChatAgentActor:
        """Minimal chat agent: conversation in KV, LLM via service link, alarm for summarization."""
    
        actor_id: str = state(default="")
        total_messages: int = state(default=0)
        total_summarizations: int = state(default=0)
    
        @init_handler
        def on_init(self, config: dict) -> None:
            self.actor_id = config.get("actor_id", "")
            host.info(f"ChatAgentActor init actor_id={self.actor_id}")
    
    
        @handler("__alarm__")
        def on_alarm(self) -> dict:
            """Durable alarm callback — equivalent to Cloudflare Agents SDK onAlarm().
    
            Summarizes conversation history and stores a summary KV key,
            then clears history.
            """
            host.info("ChatAgentActor: alarm fired — summarizing history")
    
            history = host.kv.get_json("history") or []
            if not history:
                return {"status": "ok", "action": "no_history_to_summarize"}
    
            # Summarize via LLM
            summary_prompt = (
                f"Summarize this conversation concisely (2-3 sentences): "
                f"{json.dumps([{'role': m['role'], 'content': m['content']} for m in history])}"
            )
            summary = self._call_llm([{"role": "user", "content": summary_prompt}])
    
            # Persist summary, clear history — equivalent to: storage.put('summary', s); storage.delete('history')
            host.kv.put("summary", summary)
            host.kv.delete("history")
    
            self.total_summarizations += 1
    
            host.info(f"ChatAgentActor: summarized {len(history)} messages")
            return {
                "status": "ok",
                "action": "summarized",
                "messages_summarized": len(history),
            }

    What PlexSpaces adds beyond Cloudflare’s model: the actor can also be signalled externally at any time via host.send(actorId, "wake_early", {...}) — you’re not limited to the alarm cadence.


    Agents Feature Map

    Cloudflare Agents SDKPlexSpaces
    this.storage.get/put (conversation history)host.kv.get_json / host.kv.put_json
    env.AI.run(model, messages)ServiceHttpClient("llm-link").post(...)
    storage.setAlarm / onAlarm()host.alarm.set() / @handler("__alarm__")
    connection.send(msg)host.send(actorId, op, payload)
    Workflow checkpointingWorkflowActor Run/Signal/Query + durable state
    Human-in-the-loop / approval gates@fsm_actor + workflow_signal:resume
    Long-running scheduled agentsReminderFacet + alarm reschedule
    Multi-agent orchestrationOrchestratorActor + process groups + TupleSpace
    env.AI binding in wrangler.toml[service_links] in app-config.toml
    Cloudflare edge onlyLocal, Docker, K8s, on-prem, multi-cloud

    Part III: Distributed Computation

    Cloudflare Workers and Lambda optimize for millisecond, latency-sensitive request handling. PlexSpaces handles a second problem class: large-scale computation across resources that come and go.

    ShardGroups: Scatter-Gather Without the Plumbing

    For data-parallel and ML-style workloads, PlexSpaces exposes MPI-style collectives directly through the host API:

    // Create a pool of 20 workers, hash-partitioned
    let pool_id = client.create_worker_pool(
        "worker-pool-1", "worker", 20,
        PartitionStrategy::Hash, HashMap::new(),
    ).await?;
    
    // Bulk update: 10,000 messages routed to the right shard by key
    client.parallel_update(&pool_id, updates, ConsistencyLevel::Eventual, false).await?;
    
    // Parallel map: query every shard simultaneously
    let results = client.parallel_map(&pool_id, json!({ "action": "get_total_count" })).await?;
    
    // Parallel reduce: aggregate stats across all shards
    let stats = client.parallel_reduce(
        &pool_id, json!({ "action": "stats" }),
        ShardGroupAggregationStrategy::Concat, 20,
    ).await?;

    Idle Browsers as Compute Nodes

    The mersenne_prime TypeScript example runs a GIMPS-style distributed primality search: browser tabs connect as thin WebSocket clients, receive worker JavaScript from a CodeServerActor, and run Lucas-Lehmer tests inside Web Workers. A CoordinatorActor running as WASM on the server assigns exponents, tracks per-worker CPU cores, and dispatches the next candidate immediately on each result:

    // PlexSpaces TypeScript — examples/typescript/apps/mersenne_prime/mersenne_actor.ts
    onResult(payload: ResultPayload): unknown {
        const item = this.state.work[String(payload?.p)];
        item.status = 'done';
        item.is_prime = Boolean(payload.is_prime);
        item.duration_ms = payload.duration_ms ?? 0;
    
        // Record every completed candidate in TupleSpace and bump a Prometheus counter
        host.ts.write(['result', String(payload.p), item.is_prime ? 'true' : 'false',
            String(item.duration_ms), payload.actor_id ?? 'unknown']);
        host.incrCounter('ts-mersenne-prime', item.is_prime ? 'primes_found' : 'composites_found');
    
        // Immediately hand the same worker the next pending candidate
        const next = this._nextPending(this.state.workers[payload.actor_id!]?.cpu_cores ?? 1);
        if (next) {
            next.status = 'assigned';
            host.send(payload.actor_id!, 'assign_work', { p: next.p, done: false });
        }
        return { ok: true };
    }

    Open the same URL in ten browser tabs and you have ten compute shards, coordinated by one WASM actor, with zero additional infrastructure.


    Comprehensive Feature Map

    FeatureCloudflare DO / AgentsPlexSpaces
    Durable KVctx.storage.get/puthost.kv.get/put
    Batch KV writestorage.put(new Map)host.kv.multiPut(entries)
    Batch KV readstorage.get([keys])host.kv.multiGet(keys)
    Atomic CASManual retryhost.kv.cas(key, expected, new)
    Atomic counterManual CAShost.kv.increment(key, delta)
    KV TTLputWithMetadatahost.kv.putWithTtl(key, val, secs)
    Durable alarmstorage.setAlarm(ts)host.alarm.set(ts)
    Alarm querystorage.getAlarm()host.alarm.get()
    Alarm cancelstorage.deleteAlarm()host.alarm.delete()
    Alarm callbackasync alarm()on__alarm__() / @handler("__alarm__")
    Get-or-createenv.BINDING.get(id)getActorRef(type, name, ns)
    List actors by namespaceCloudflare REST APIListActors gRPC / HTTP REST
    Init lifecycleblockConcurrencyWhileonInit() / @init_handler / Init()
    WebSocket (standard)DO holds socketThin-node session actor per connection
    WebSocket hibernationstate.acceptWebSocketRoom actor eviction + getState() restore
    LLM callsenv.AI.run(model, msgs)ServiceHttpClient("llm-link").post(...)
    Conversation statethis.storage.get('history')host.kv.get_json("history")
    Durable workflowsCloudflare WorkflowsWorkflowActor Run/Signal/Query
    Human-in-the-loopManual pause / external call@fsm_actor + workflow_signal:resume
    Long-running agentsscheduleAlarm()ReminderFacet + alarm reschedule
    Multi-agent orchestrationMultiple DO fetchesProcess groups + TupleSpace coordination
    MCP tool integrationMcpAgent classHTTP service link (client-side)
    Routing configwrangler.toml [[bindings]]app-config.toml [[children]]
    Cross-actor fan-outIndividual fetch() callshost.send() (in-process or mesh)
    Fire-and-forget delayExternal queue / DO alarmhost.sendAfter(delayMs, op, payload)
    Multi-languageTypeScript onlyGo, TypeScript, Python, Rust
    Local devwrangler dev (simulated)Same binary, full parity
    On-prem / self-hostNoYes
    Multi-cloudNoYes, via gRPC mesh
    ObservabilityCloudflare AnalyticsPrometheus + OTLP, self-hosted
    WebhooksWorker fetch()[[http_routes]] in config
    Process groupsNot supportedhost.pg.broadcast(group, op, payload)
    Data-parallel computeNot supportedShardGroups (scatter-gather, allreduce)

    A couple of things need real rework, not a find-and-replace:

    • WebSocket architecture. If your DO holds sockets directly today, plan for a thin node that hosts actors.
    • Bindings vs children. Cloudflare bindings live in wrangler.toml and show up as env properties. PlexSpaces declares the same relationships as supervision children in app-config.toml.

    Running Everywhere

    The whole point is that development and production run the identical binary:

    # Local development — exact production behavior
    plexspaces-node start --config app-config.toml
    
    # Docker — same binary, same config
    docker run -v $(pwd):/app plexspaces/node start --config /app/app-config.toml
    
    # Kubernetes — same config via Helm
    helm install my-app plexspaces/app-chart \
      --set config.path=app-config.toml \
      --set persistence.storage=postgres
    
    # Multi-cloud — nodes in GCP + AWS joined over a gRPC mesh; actors route transparently

    An alarm that fires in production runs through the exact code path you tested on your laptop. There’s no gap to debug between wrangler dev and prod, because there’s only one runtime.


    Working Examples

    Every example referenced above is in the repository (github.com/bhatti/PlexSpaces):

    • Guild chat (DO migration pattern): examples/{go,typescript,python,rust}/apps/migrating_cloudflare_workers/: ChatRoomActor with member fan-out, a token-bucket RateLimiterActor backed by host.kv.increment and host.kv.cas, and an AlarmDemoActor mirroring DO’s full alarm lifecycle
    • WebSocket chat room: examples/{typescript,python,go,rust}/apps/ws_chat_room/: ChatRoomActor plus a PresenceActor that uses a durable reminder to detect idle disconnects
    • AI chat agent (Cloudflare Agents SDK pattern): examples/{go,typescript,python,rust}/apps/chat_agent/: conversation history in KV, LLM calls through service link, durable summarization alarm
    • Human-in-the-loop approval gate: examples/python/apps/minipi/approval_gate.py: FSM-based approval workflow, workflow_signal:resume handoff, state durable across multi-day waits
    • Multi-agent orchestration: examples/go/apps/miniclaw/: OrchestratorActor decomposing tasks, delegating via process groups, aggregating through TupleSpace
    • Durable payment workflow: examples/go/apps/migrating_cadence/payment_workflow.go: WorkflowActor with Run/Signal/Query, idempotent retry, refund and cancel signals
    • Mersenne prime search (browser compute): examples/typescript/apps/mersenne_prime/: browser tabs as Lucas-Lehmer worker shards, coordinated by a WASM CoordinatorActor
    • wasmCloud migration: examples/python/apps/migrating_wasmcloud/session_store.py: capability-based session store showing host.kv.list, inter-actor ask, and timer-driven cleanup
    • Every abstraction in one place: examples/{go,typescript,python,rust}/apps/abstractions/: durable virtual-actor reactivation, workflow run/signal/query, process-group event delivery, timers, reminders, KV, tuple space, and blob storage

    Conclusion

    The Cloudflare model is the right model: stateful actors, durable storage, alarms, WebSocket session management, LLM calls baked in. Wire’s post makes the same point from the other direction as they’re not leaving because the model is wrong, they’re leaving because specific architectural ceilings came due at their scale. PlexSpaces keeps the model and removes the ceiling: you write the same actors, get the same alarms and durable storage and observability, and the binary that runs on your laptop is the exact binary that runs in production, on any cloud, on-prem, or across all of them at once. The migration from Cloudflare DO code is mostly mechanical. What you get back is the ability to run anywhere, own your infrastructure and never again debug a divergence between wrangler dev and prod.


    GitHub: github.com/bhatti/PlexSpaces

    Previous posts in this series:

    July 16, 2026

    The Fallback Trap: How Defensive Programming Silently Destroys Distributed Systems

    Filed under: Computing — admin @ 9:08 pm

    I have seen some systems never crash, they start cleanly, swallow every error, and keep running no matter what goes wrong. In my experience, they are also the hardest systems to debug, the most dangerous to operate, and the most expensive to maintain. I worked on a similar legacy system for distributed data platform that routed events between hundreds of thousands of nodes. It had zero unhandled exceptions in production. It also had silent authentication failures, invisible data loss, and configuration divergence that took days to diagnose. This is a follow-up to my earlier posts on building an observability platform in Rust, why DRY becomes a liability, and making bad state impossible with ADTs. Those posts were about the type system but this one is about a habit of mind that no type system fixes on its own: the instinct to catch every error with some fallback or default behavior.

    The culprit in that system was never a lack of error handling. It was error handling, applied in the wrong places, for the wrong reasons. I call it defensive programming as a religion: every function protects itself against every possible invalid input by inventing a fallback. Missing config? Use a default. Secret unavailable? Generate a random one. Database write failed? Log it and move on. The result is a system that looks healthy on every dashboard while quietly corrupting its own state underneath. In this post, I will walk through the patterns I found, why each one causes more damage than the crash, and what an alternative looks like instead. The whole argument rests on one idea: Every fallback creates a new source of truth, and two sources of truth always drift apart.


    I. Why a Fallback Is Worse Than It Looks

    It’s tempting to think of a fallback as just “hiding an error.” It’s worse than that. When a function invents a value because the real one is missing, that invented value doesn’t stay hidden, instead it becomes a fact in the system. From that moment on, the system is carrying two truths: the one that should exist but doesn’t, and the one that got made up and does. These two truths never stay in sync and when they drift apart, the failure almost never shows up where the fallback happened. Instead, it shows up somewhere else entirely, in a component with no obvious connection to the code that invented the value. For example, you’ll spend hours in the authentication layer before realizing the signing key was randomly generated at startup by a config migration function three layers away. To be clear about what I mean by “fallback,” I am not talking about:

    • Validating input at system boundaries: checking what a user typed, sanitizing data from outside. Correct and necessary.
    • Graceful degradation with an explicit signal: returning a typed Degraded state that the caller can see and react to.
    • Retrying transient I/O failures with backoff. Standard practice.

    I’m talking about code that silently invents state when the real state is missing, and then carries on as if nothing happened. Three things make a fallback harmful:

    1. It hides the root cause. The missing value was the bug. The fallback makes it disappear.
    2. It persists the invented value. Once it’s written to disk or sent over the wire, every future operation has to succeed against a value that was never correct in the first place.
    3. It relocates the symptom. The failure surfaces hours later, in a different component, in a different log file, with no visible thread connecting it back to the startup code that invented the wrong value.

    I am not advocating “crash on every error.” It just means: a function that requires X must fail when X is absent and it must never invent X.


    II. Postel’s Law

    There’s a reason smart engineers build these fallback-heavy systems: they’re following a respected principle: “be conservative in what you send, be liberal in what you accept.” that Jon Postel wrote as guidance for TCP implementations. That advice made a lot of sense in its original context. But this principle leaked out of the protocol layer and turned into a general design philosophy. Engineers started applying “be liberal in what you accept” to function signatures, config loading, and communication between services inside a system they fully control. A function that takes string | undefined and silently substitutes a random value gets called “being robust.” A startup sequence that swallows errors and keeps going gets called “being tolerant.”

    Inside your own system, that assumption doesn’t hold. For example, you can fix the sender because you own the caller. But when you own the migration script that’s supposed to write the auth token and you “liberally accept” a missing auth token by inventing a random one instead, you’re not enabling interoperability with an outside party instead you’re hiding a bug in code you wrote. This misapplication creates a ratchet effect. Every “liberal” receiver makes it harder to notice problems at the source. If every function tolerates missing input, the function that’s supposed to supply that input has no pressure to get it right. People also tend to forget that Postel’s Law has a second half: “be conservative in what you send.”

    The corrected version for internal systems is this: be strict with components you control, and liberal only at the boundaries where you genuinely can’t fix the sender like external APIs, user input, third-party integrations. Inside your own codebase, strictness isn’t fragility. A function that rejects invalid input tells you exactly where the bug lives.


    III. Inventing Values: The Most Dangerous Pattern

    This is the category that caused the most damage and I have seen countless bugs due to this anti-pattern. For example, the code generates a random value, assigns it to a security-critical field, and moves on as if that field were properly populated. The invented value becomes a durable fact somewhere and it’s always wrong.

    The Archetype: Inventing a Secret

    A function runs at startup and writes a signing secret into every worker group’s configuration. Workers and the coordinator use this secret to authenticate each other. If two groups end up with different values, every authentication attempt between them fails silently, showing up as delivery failures instead of auth failures.

    // The bug: if authToken is absent, invent a UUID and persist it
    const authToken = settings.distributed?.master?.authToken;
    const plaintext = authToken != null && authToken.length > 0
      ? authToken
      : randomUUID();  // <-- this line breaks authentication across the cluster
    

    When authToken is missing from the merged settings, which is a perfectly valid state on a fresh install but this function generates a randomUUID() and writes it as the signing secret for every group it touches. Each call produces a different UUID. Meanwhile the token store writes yet another value through a completely different code path. The function reports success for every group it writes to. The symptom 401 errors between nodes shows up minutes or hours later, in worker logs, pointing investigators toward the wrong layer entirely. This is the archetype of the whole problem: if the required input is missing, invent something plausible-looking and keep going. The fix is three lines:

    const authToken = settings.distributed?.master?.authToken;
    if (authToken == null || authToken.length === 0) {
      logger.warn('authToken absent from settings; cannot write signing secrets');
      return;  // do not proceed; do not invent
    }
    

    The “Disabled” Sentinel That Looks Just Like a Real Key

    A secret provider has a three-source fallback chain: encrypted store, config file, random bytes. That last fallback is supposed to act as a “poison key” that never validates:

    async getSecret(): Promise<string> {
      // Source 1: encrypted store
      const fromStore = await secretsMgr.get(KEY_ID).catch(() => null);
      if (fromStore) return fromStore;
    
      // Source 2: config file (which may itself contain a well-known default!)
      const fromFile = settings.distributed?.master?.authToken;
      if (fromFile) return fromFile;
    
      // Source 3: "disable" by returning random bytes — looks perfectly valid to the caller
      return random(16);
    }
    

    The caller gets back a plain string in all three cases. It has no way to tell “a real secret from the store” apart from “a well-known default from the config file” apart from “random bytes that will never work.” It signs a token with whatever string it received and sends it off. This is a type-system failure because the return type Promise<string> squashes three semantically different outcomes into one shape.

    The explicit contract fixes this at the type level:

    type SecretResult =
      | { kind: 'valid'; value: string; source: 'store' | 'file' }
      | { kind: 'unavailable'; reason: string };
    
    async getSecret(): Promise<SecretResult> {
      const fromStore = await secretsMgr.get(KEY_ID).catch((err) => {
        logger.error('Secret store unavailable', { error: err });
        return null;
      });
      if (fromStore) return { kind: 'valid', value: fromStore, source: 'store' };
    
      const fromFile = settings.distributed?.master?.authToken;
      if (fromFile) return { kind: 'valid', value: fromFile, source: 'file' };
    
      return { kind: 'unavailable', reason: 'no secret in store or config' };
    }
    

    Now a caller that receives unavailable marks itself as degraded. It doesn’t sign tokens and surfaces a health check failure instead.

    The Token Renewal That Signs With Random Bytes

    The token authenticator calls getSecret(), and if that fails, it falls back to random bytes anyway:

    let keyStr = await this.secretProvider.getSecret().catch(() => undefined);
    if (keyStr == null) {
      this.logger?.debug('Unable to generate a valid token. Disabling.');
      keyStr = random(16);  // this "signed" token will never be accepted by any peer
    }
    const token = jwt.sign(payload, keyStr);
    this.cachedToken = token;  // cached and reused for every future request
    

    The log message says “Disabling,” but nothing gets disabled. The code signs a token with a random key, caches it, and hands it out to every caller for the rest of the process’s life. Workers receive the token, fail to verify it, and log a 401, with nothing to suggest the root cause of the issue.

    The explicit contract: if the secret is unavailable, don’t sign anything. Set this.cachedToken = null. Let callers check for null and surface a real health degradation.

    The User ID That Regenerates on Every Call

    const userId = existing?.username ?? crypto.randomUUID();
    

    When existing is unexpectedly missing, every call generates a brand-new UUID. The user becomes a ghost: every request creates a fresh identity, invisible to deduplication, audit trails, and rate limiters.

    The Config Placeholder That Silently Materializes

    if (authToken?.token === 'REPLACE_ME') {
      authToken.token = uuidv4();  // silent replacement, no log of the value generated
    }
    

    This runs at startup and inside a database migration. If the startup write fails, the generated token is gone.


    IV. Swallowed Errors

    In one legacy codebase I worked had 656 instances of .catch(NOOP), an empty function attached to a promise rejection that turns any error into undefined and lets execution continue. Many of these were on cleanup paths, which is harmless. But a large number sat on critical data paths like durability, metrics transport, authentication state.

    The Durability Guarantee That Wasn’t

    A persistent queue exists for exactly one reason: to guarantee that events survive a destination outage. That’s the system’s durability promise.

    // PQ flush — the durability guarantee
    await this.flushBuffer().catch(NOOP);  // silently swallows disk I/O errors
    await this.commit().catch(NOOP);        // silently swallows commit failures
    

    If the flush fails like disk full, I/O error, permission denied, the buffered events are silently gone. The one component whose entire job is preventing data loss is itself a source of silent data loss.

    The explicit contract: treat PQ flush errors as fatal to the ingest path. For example, if the flush fails, apply backpressure and pause ingest. Log at error level with the event count at risk and emit a pq_flush_failure metric. Now the operator gets to choose: fix the disk, add capacity, or consciously accept the loss.

    Metrics That Vanish

    void saasMetrics.sendPacket(packet).catch(NOOP);
    

    Every metrics-send failure is silently swallowed.Dashboards go blank, and nobody knows why.

    The Config Load That Treats Corruption as “Empty”

    const groups = await conf.loadSystem('internal-groups').catch(() => ({}));
    

    One line here quietly conflates two very different situations:

    • “The file doesn’t exist yet”, which is normal on first boot –> return {}
    • “The file is corrupt, or a parse error, or permission was denied”, which is a real configuration bug –> also return {}

    Either way, the loop over groups never runs. The operator has no way to tell “healthy, no groups configured yet” apart from “broken, groups exist but couldn’t be read.”

    Startup That Succeeds Despite Total Failure

    export async function syncGroupSecrets(conf: Configuration): Promise<void> {
      try {
        // ... write secrets to all groups ...
      } catch (err) {
        logger.error('failed to sync secrets', { reason: err });
        // swallowed — startup continues, caller receives no signal
      }
    }
    

    The function catches everything at the top and resolves successfully no matter what. The caller awaits it and gets no signal that anything went wrong. If the sync fails for every group, the coordinator still starts, workers still connect, and authentication fails across the board but startup “succeeded.”


    V. Defensive Defaults

    This next category is more subtle. Instead of inventing a random value, the system injects a well-known one like a default credential, a default address, which makes it impossible for downstream code to tell that anything is missing at all.

    The Well-Known Default Credential

    The shared authentication token has a configuration setting, and by default, the settings loader injects a well-known string whenever nothing is configured:

    const { injectDefaultAuthToken = true } = options ?? {};
    if (injectDefaultAuthToken && conf.distributed?.master?.authToken == null) {
      conf.distributed.master.authToken = DEFAULT_AUTH_TOKEN;  // 'default-secret'
    }
    

    Any caller of getSettings() receives a truthy string for authToken. Code that checks if (authToken) proceeds as if a real token exists. The check passes. Authentication moves forward, using a credential that’s sitting right there in the source cod and every deployment that forgot to override it. The default here is opt-out, not opt-in. Every new call site has to remember to disable the injection.

    Workers That Connect to localhost on Config Failure

    const server = this.conf.distributed.master
      || { host: 'localhost', port: 5555, authToken: DEFAULT_AUTH_TOKEN };
    

    When a worker fails to load its coordinator address, it silently falls back to localhost:5555 with the default token. On a single-node dev box, this might accidentally work. On a production multi-node deployment, the worker ends up connecting to itself. The symptom looks like a connection timeout, not a config-load failure.

    The explicit contract: if conf.distributed.master is missing, throw. A worker cannot function without a coordinator address, and there is no valid default for it. Failing here with a clear message like “coordinator address not configured; set MASTER_URL or add distributed.master to instance.yml“.


    VI. Multiple Sources of Truth

    This is the pattern that ties everything else together. Every fallback chain creates more than one source of truth and any source of truth that isn’t explicitly designated as the source will eventually drift from the others. In a distributed system, that drift shows up as the hardest class of bug like intermittent or state-dependent failures.

    Two Functions, One Secret, Two Different Sources

    The bug that inspired this post exists because two functions write the same signing secret to group configs, but read the plaintext from two different physical sources:

    FunctionReads fromRuns when
    syncAllGroupSecretsConfig file (instance.yml)Startup for every group
    syncNewGroupSecretEncrypted token storeGroup creation for one group

    A migration populates the token store before syncAllGroupSecrets runs at boot, so the store is meant to be authoritative. But syncAllGroupSecrets predates the store’s existence and still reads straight from the config file. These two sources drift apart after a token rotation or some race condition.

    The explicit contract: one source of truth. Both functions read from the store. If the store is unavailable, both fail instead of silent fallback to a secondary source.

    The Three-Source Chain Is Three Sources of Truth

    Source 1: Encrypted store (authoritative)
      ? (unavailable ? silently falls through)
    Source 2: Config file (may hold a stale or default value)
      ? (absent ? silently falls through)
    Source 3: Random bytes (structurally valid, semantically useless)
    

    Each fallback quietly downgrades the security posture, and the caller gets back a plain string with no idea which source it came from.

    Three branches go into the same sign() call. Only one of them should ever be allowed to reach it, which is the whole argument for making unavailable its own explicit type instead of letting all three collapse into a plain string.

    The Cache Nobody Fully Trusts

    The settings system keeps a cache for high-availability mode. Some callers pass skipCache: true to bypass it; others don’t, and there’s no documented rule for which is which. This gives the system two sources of truth for the same data: the cache and the disk. If the config file changes between startup and an API call, the API may serve stale data. The skipCache escape hatch is a symptom, not a fix. It means someone stopped trusting the cache’s invalidation and punched a hole through it instead of repairing the underlying mechanism.

    The explicit contract: the cache invalidates on every write. Callers never need to know or care whether they’re reading from cache or disk. Remove skipCache as a public option entirely.


    VII. Redundant Guards

    When a function doesn’t trust its caller’s preconditions, it adds its own guard on top:

    // Caller (server.ts):
    if (isLeader && featureFlags.check('AUTH_TOKEN_MGMT')) {
      await syncGroupSecrets(conf);
    }
    
    // Callee (syncGroupSecrets):
    export async function syncGroupSecrets(conf: Configuration): Promise<void> {
      if (!isFreeTier() && !isRunningInSaaS()) return;  // guard 1
      if (!Product.isLeader(settings.distributed?.mode)) return;  // guard 2 (redundant!)
      if (!featureFlags.check('AUTH_TOKEN_MGMT')) return;  // guard 3 (redundant!)
      // ... actual work ...
    }
    

    This function lives in a directory named leader/ and is only ever called from the leader startup path, yet it re-checks isLeader internally anyway. The feature flag gets checked at the call site and again inside the function. Three layers of defense against calling this function in the wrong context.

    The checks themselves aren’t the problem. It’s what happens when they trip: nothing. The function returns quietly, and the caller gets no signal either way.

    The explicit contract: a function either does its job or throws. Preconditions get asserted, not silently absorbed. If the caller already guarantees the precondition, drop the internal check. If the function really can be called from multiple contexts and some of them are invalid, throw on the invalid ones instead of quietly returning.


    VIII. “Best-Effort” Writes to State That Isn’t Optional

    The most seductive justification for swallowing an error is: “the primary operation already succeeded so we don’t want a secondary failure to undo it.” That reasoning is correct in isolation and catastrophic in aggregate.

    The Store Upsert That Swallows Its Own Failure

    export async function mirrorTokenToStore(plaintext: string): Promise<void> {
      try {
        await store.upsert({ id: LEGACY_TOKEN_ID, token: encrypt(plaintext) });
      } catch (err) {
        logger.error('failed to mirror token to store', { reason: err });
        // swallowed — the caller sees success
      }
    }
    

    The comment above this function explains the intent: it’s “best-effort” so that a store-side failure doesn’t roll back a config file write that already succeeded. That reasoning holds up on its own but config file is now permanently out of sync.

    The explicit contract: add reconciliation. On startup, compare the store’s token to the config file’s. If they differ, update the store. Emit a token_store_divergence counter and surface it in a health check similar to reconciliation loops in Kubernetes.


    IX. Partial Operations Without a Way Back

    The Batch Write That Discards on Failure

    const batchTxn = db.transaction((ops) => { ops.forEach(fn => fn()); });
    try {
      batchTxn(this.mutationCache.splice(0, maxSize));  // splice removes BEFORE success
    } catch (err) {
      logger.error('Batch write failed', err);
      // operations already removed from cache — permanently lost
    }
    

    The comment in this code says: operations get removed from the cache regardless of whether the transaction succeeds. But the fix for “one bad operation might stall the queue” ends up being “discard the entire batch, including the good operations.” That isn’t a tradeoff instead it’s silent data loss.

    The explicit contract: only remove items from the cache after a confirmed commit. Isolate the failing operation into a dead-letter queue and retry the rest of the batch without it.

    The Config Reload Nobody Acknowledges

    await conf.triggerReload().catch(NOOP);  // worker continues with old config
    

    After receiving a new config bundle from the coordinator, a worker triggers a reload. If that reload fails, the worker just keeps running the old config. The two sides now disagree about what the worker is actually running.

    The explicit contract: report a reload failure back to the coordinator on the next heartbeat. The coordinator marks that worker as “stale config” and can retry or alert. This is how Kubernetes rolling updates work, the controller notices and either retries or halts the rollout.

    The Package Install That Saves Despite Partial Failure

    for (const op of ops) {
      try {
        switch (op.type) {
          case 'install': await installPackage(op); break;
          case 'uninstall': await uninstallPackage(op); break;
        }
        status.applied.push(op);
      } catch (error) {
        errors.push(error);
      }
    }
    await this.save(packageManifest);  // saves regardless of how many failed
    

    If three out of five packages install and two fail, the manifest still gets saved with those three. Next startup tries again from this partial state but the failed packages may have left behind lock files or half-written artifacts causing conflicts.

    The explicit contract: validate that every operation can succeed before running any of them (a dry-run pass). Execute atomically (ACID transactions).


    X. Silent Truncation With No Backpressure

    The Metrics Buffer That Silently Drops

    Workers piggyback metrics onto heartbeat messages sent to the coordinator, with a hard cap of 100,000 packets. Past that cap, excess metrics are silently dropped without any counter, logs or metrics. This is a nasty failure mode specifically because absence of data is itself meaningful data and silent truncation destroys that signal’s reliability.

    The explicit contract: when the buffer nears capacity, reduce granularity instead of dropping outright. When truncation does happen, include metrics_truncated: N in the heartbeat so the coordinator knows its picture is incomplete. Better, instead of piggyback metrics on heartbeats at all, give them their own transport.

    The TCP Sender That Zeroes Buffers on Disconnect

    // On disconnect: all in-transit events silently lost
    this.inTransitBufs = [];
    this.bufOffset = 0;
    this.bufferEventCount = 0;
    this.dropBytes += len;  // only evidence: a counter increment
    

    On a TCP disconnect, every in-transit event gets zeroed out. The only trace left behind is a dropBytes counter buried in internal metrics. Compare that to Kafka’s producer, where unacknowledged messages stay in the producer’s buffer and get retried on reconnect.

    The Unbounded Queue That Becomes an OOM

    protected queueBatch(): void {
      this.queuedBatches.push({ eventCount, eventsSize, events: this.currentBatch });
      // NO CHECK on length, size, or memory pressure
    }
    

    When the output destination is unreachable, failed batches get re-queued, and the queue grows without any bound. Memory climbs until the OOM killer steps in and terminates the process. An unbounded in-memory queue isn’t really a data structure, instead it’s a deferred OOM crash.

    The explicit contract: bound the queue. Once it’s full, either apply backpressure to ingest, spill overflow to the persistent queue, or trip a circuit breaker that rejects new events with a typed error. Let the pipeline decide from there: drop, buffer to disk, or pause the source.


    XI. Non-Atomic Writes

    The Lease File That Can Split-Brain

    await writeFile(this.leaseFile, stringify(content));
    

    The failover lease file, the mechanism that’s supposed to guarantee only one coordinator is ever active is written directly with writeFile. On NFS, which is where this system runs in HA mode, writes aren’t atomic. A power failure mid-write leaves behind a truncated or corrupt file. The standby coordinator reads that corrupt lease, fails to parse it, and ends up in an undefined state.

    The explicit contract: write to a temp file, fsync, then atomically rename it into place, with a checksum so readers can detect corruption. Every real database does this like SQLite’s WAL.

    The Multi-File Config Deploy Without a Journal

    Config deployment writes several YAML files in sequence like inputs.yml, outputs.yml, pipelines.yml,, etc. A crash midway through leaves the worker with a partial config and a worker that restarts after a partial deploy loads that inconsistent config.

    The explicit contract: write every file to a staging directory first, verify that everything references correctly, then swap atomically like rename the directory, or flip a symlink. This is the same idea behind Docker image layers, Kubernetes ConfigMaps, and Nix store paths.


    XII. Six Principles That Cover All of This

    Every pattern above breaks one of six well-established principles. They’re standard practice in any system that prioritizes correctness over the appearance of uptime.

    1. Fail fast at trust boundaries (Erlang’s “let it crash.”): When a precondition is violated, fail immediately and loudly. Erlang runs telecom infrastructure at 99.9999999% uptime on a philosophy of letting individual processes crash and having a supervisor restart them into known-good state.

    2. Make invalid states unrepresentable: If getSecret() can return random(16) as a plain string then every caller is stuck defensively guessing whether it’s “real.” If it returns Secret | Disabled as a discriminated union instead then the type system forces every caller to handle both cases at compile time. I wrote about this pattern in “Making Bad State Impossible: A Practical Guide to ADTs and Algebraic Effects.”

    3. Classify errors (transient vs. fatal): Without classification, every catch block faces an impossible choice: rethrow and break “resilience,” or swallow and hide a real bug. For example, gRPC solves this with status codes like UNAVAILABLE means retry, INVALID_ARGUMENT means don’t, INTERNAL means there’s a bug.

    4. Define delivery semantics for critical state: “Fire-and-forget” is fine for debug logs. It’s not fine for persistent queue flushes, token store upserts, or config reloads. If an operation mutates state that downstream code assumes succeeded, it needs at-least-once semantics.

    5. One source of truth without fallback chains for critical state: For any given piece of state, there should be exactly one authoritative source. A fallback chain isn’t graceful degradation instead it’s an implicit decision that secondary sources are acceptable substitutes for the truth. If that’s genuinely acceptable, make it explicit with TTLs, version vectors, or consistency levels.

    6. Atomic state transitions: State changes should be all-or-nothing: temp file, fsync, atomic rename for files; transactions with rollback for databases; staging plus swap for multi-file deployments.


    XIII. Cognitive Load

    All this conditional logic and fallback behavior creates the cognitive load, e.g., when any function might silently invent a value, you can’t trust a function’s output without reading its implementation. When errors are swallowed, a successful await no longer means the operation actually succeeded. When defaults get injected into config reads, a non-null value stops meaning “configured.”

    Debugging a production incident in a system like this means reading every function in the call chain, understanding every fallback along the way, and reconstructing which code path actually ran. crash tells you exactly where and when an invariant broke. New engineers ask why workers sometimes fail to authenticate after a token rotation, and the honest answer involves reading six functions across four files, understanding a three-source fallback chain. Compare that to: “the store upsert threw on failure, the rotation API returned a 500, the operator re-ran it, it worked.


    XIV. Conclusion

    Defensive programming isn’t inherently wrong, and neither is Postel’s Law like at the boundary it was designed for. Validating input at system boundaries, handling I/O errors gracefully, protecting against malformed external data are correct applications of defensive thinking. The problem shows up when the same techniques get applied to internal code, inside a system where you control both ends of every interface. The alternative, in short:

    • Throw Error when a required state that’s missing: The function refuses to proceed and the caller finds out immediately.
    • Explicit return when an optional state that’s missing: null, undefined, Option<T>, a discriminated union.
    • Transient failures = retry with backoff: Never .catch(NOOP).
    • One source of truth per piece of state: Not a fallback chain that quietly degrades. Not a cache with no real invalidation.
    • Bounded queues with backpressure: Not unbounded buffers waiting to OOM.
    • Atomic state transitions: Not multi-step operations that can half-finish.
    • Reconciliation loops for distributed state: Not one-shot “best-effort” writes that quietly drift apart.

    This mud didn’t accumulate overnight, and it won’t disappear overnight either. But every fallback you remove, every error you refuse to swallow will makes the next incident roughly ten times faster to diagnose.

    Related Blogs

    1. From Big Ball of Mud to Functional Pipeline
    2. The Reusability Trap: When DRY Becomes a Liability
    3. Making Bad State Impossible: A Practical Guide to ADTs and Algebraic Effects

    July 9, 2026

    Building an Agent Harness and Eval Pipeline with Durable Actors

    Filed under: Agentic AI,Computing — admin @ 1:09 pm

    You may start with a simple agent for demo that calls a tool, gets an answer, prints it. But building a production ready agentic system requires a full agent harness so that it doesn’t crash halfway through a task. For example, you might have an AI coding assistant that debugs production incidents by searching logs, isolating a root cause, creating a pull request for the fix. It might take several minutes with dozens of tool calls. The AI agent may crash mid-run, fail to call a tool reliably or the test suite for eval fails. These are not model problems that you can solve with a better model or a better prompt. Instead, you need a reliable infrastructure that an agent harness provides. This post shows how to build an agent harness and eval pipeline using PlexSpaces, a polyglot actor framework that treats agent infrastructure as a first-class problem instead of an afterthought.


    Agent = model + harness

    You can think of an agent as model + harness. The harness is everything that isn’t the model.

    The harness is the loop that decides when to stop, the tool calling that connects the model to the world, the state that survives a crash and resumes where it left off, the coordination that lets multiple agents share work, and the eval plumbing that tells you whether any of it actually worked. Most teams spend their time on the model like a different temperature here, a different prompt there, a bigger model if budget allows. I have seen teams build a prototype agentic system and then ship it to an entire organization without proper harness resulting in unexpected failures. The harness stays invisible until it breaks, and when it breaks, it looks exactly like a model problem.

    There are three levers that move agent quality. Model changes are the most expensive like fine-tuning, RL, moving to a bigger model. Harness changes are nearly free like loop logic, tool schemas, retry policies, agent topology. Memory changes are the cheapest like context window management, retrieval strategy. Teams reach for the model first. They should usually start with the harness.


    What the harness actually has to do

    Eight responsibilities show up in every serious agent deployment, in every framework, in every language. The only question is whether you build them on purpose or accumulate them by accident after the third production incident.

    Harness propertyWhat it doesPlexSpaces primitive
    Loop controlIteration limits, token budget, stop conditionsAgentLoop (max_iterations, token_budget)
    Tool callingDispatch, schema validation, error captureToolRegistryActor + SchemaValidationFacet
    State managementSurvives crashes, resumes from checkpointDurabilityFacet (journal replay)
    MemoryPrior context per agent, per runKV store (host.kv_get / host.kv_put)
    Multi-agent coordinationFan out work, collect results without tight couplingTupleSpace (write tuple, match pattern)
    SupervisionA subagent crash doesn’t take down the orchestratorSupervision tree (one_for_one)
    ObservabilityEvery step captured and queryable mid-runExecutionTraceFacet
    Eval plumbingTrajectories –> scores –> regression detectionEvalRunnerActor, ScorerActor

    Every one of these is a solved problem in actor frameworks generally. PlexSpaces just wires them together for agent workloads specifically.


    Why the actor model fits this problem

    The actor model was designed for systems that keep running when individual components fail, which happens to be exactly the property a multi-agent pipeline needs. Each actor is an isolated unit of state and behavior. Actors talk only through messages without sharing memory. There’s no global state, and no way for one actor to corrupt another actor’s state directly.

    This lines up with what distributed systems theory already tells us. The FLP theorem says that in a distributed system where even one failure is possible, you cannot guarantee both safety and liveness without explicit coordination. The actor model handles this by making failure a first-class citizen: actors crash, supervisors restart them, the system keeps running. It’s the same design that let Erlang run telecom systems for five nines of uptime.

    For agent systems, three consequences follow directly:

    • Crash isolation. When an AgentActor fails mid-eval, only that actor restarts. EvalRunnerActor and every other running agent keep going. In thread-per-agent or future-based systems, a crash in one agent typically propagates up and takes the rest down with it.
    • No shared-state corruption. Agents talk through messages and TupleSpace, not shared memory, so they can’t overwrite each other’s context. A hallucinating agent writing garbage stays contained to itself.
    • Journal replay without application code. The durability journal lives at the framework level, below your actor’s code. You don’t implement checkpointing yourself and the framework journals every message before the actor runs it.

    The building blocks

    Following PlexSpaces primitives do all the work in this harness.

    • Actors are the basic unit. Each one owns its state and handles messages one at a time. Actors talk by sending messages and you never reach into another actor’s state directly.
    • GenServer is a request-reply actor: send it a message, it processes and replies. LLMGatewayActor, ScorerActor, and DashboardActor are all GenServers.
    • WorkflowActor is a durable workflow. It checkpoints its state before each step, and if the process crashes, it replays from the last checkpoint on restart without application code. EvalRunnerActor and BenchmarkActor are WorkflowActors.
    • GenFSM is a state machine actor: you define states and transitions, and the state persists across crashes. ApprovalGateActor (human-in-the-loop) is a GenFSM.
    • Facets are cross-cutting behaviors you attach to any actor without touching its code. You can think them as middleware, but declared in app-config.toml instead of written in application code. Three facets carry the harness:
      • SchemaValidationFacet validates tool call arguments against JSON Schema before the actor ever sees the message.
      • DurabilityFacet journals every message before your actor’s code runs, so a crash-and-restart wakes up the actor with exactly the state it had.
      • ExecutionTraceFacet records every step in order and exports the full trace to KV storage when a workflow completes, which is what feeds eval.
    • Supervision trees enforce fault isolation. You declare a tree of actors and a restart strategy; one_for_one means one crash restarts only that actor, while the orchestrator and every sibling agent keep running. In LangGraph or CrewAI, a crash typically takes down the whole graph.
    • TupleSpace is a shared blackboard for multi-agent coordination, built on the Linda coordination model. Actors write tuples (["trajectory", run_id, data]) and read them by pattern (["trajectory", run_id, nil]). Producer and consumer stay decoupled without polling or sharing state.

    Two loops, two owners

    There’s a useful mental model for agent systems: two concentric loops, with different owners.

    The inner loop is the agent trying to accomplish the task: investigate, implement, test, report. The outer loop is the engineer deciding whether the agent’s output deserves trust: decide, verify, approve, own. The harness sits at the boundary. It’s where agent output turns into evidence like diffs, test results, trajectories, scores that the engineer can actually inspect before deciding to approve, redirect, or block.

    This framing matters for eval specifically because eval tooling that scores only the final answer misses most of what’s happening. An agent that stumbles onto the right answer through a wrong path scores fine on outcome-only eval, then fails the moment the task shifts slightly. What you actually want to evaluate is the trajectory or the path, not just the destination.


    Why pass@k beats pass/fail

    Agents are non-deterministic. For example, the same task, same model, same harness will succeed sometimes and fail other times. A single binary pass/fail on one run gives you noise, not signal. The right metric is pass@k: run the same scenario k times and count how many succeed. A score of 0.9 means the agent solved it 9 out of 10 tries. This is well established in code-generation benchmarks like SWE-bench, HumanEval, and MBPP. The same logic carries over to agent harnesses: you need pass@k across your own task distribution, not a one-shot score you happened to get lucky on. Comparing pass@k between two harness configurations gives you real evidence about which one is more reliable, without touching the model at all.

    ScorerActor produces the 0–1 signal pass@k needs, using rubric-based scoring:

    Step 9: ScorerActor — score trajectory
      score task_completion
      Score: 0.85  (rubric: task_completion)
      score tool_use
      Score: 0.80  (rubric: tool_use)

    Two rubrics run against the same trajectory. task_completion asks whether the agent reached the goal. tool_use asks whether it used the right tools in the right order. These can diverge, e.g., an agent might complete a task through a lucky shortcut that would fail on a harder variant. The trajectory rubric catches that divergence; the outcome-only rubric never sees it.


    In-runtime eval versus post-hoc eval tools

    The popular eval tools like LangSmith, DeepEval, Braintrust, Phoenix/Arize work the same way: run the agent, export traces or logs, evaluate afterward. That model has a structural flaw: eval doesn’t run in the same environment as production. Agent configuration, tool schemas, retry logic, and context management can all drift between the eval harness and the production deploy. When eval passes and production fails, there’s no way to tell whether the failure is in the model or just in the mismatch between the two setups.

    MiniPi’s eval runs inside the same PlexSpaces node, against the same actors, with the same tool schemas, under the same supervision tree as production. EvalRunnerActor isn’t a separate process logging to an external service, instead it’s an actor in the same supervision tree as the AgentActor it’s testing. If you change the schema in production, and eval will pick it up automatically, because it’s the same schema.

    That also makes eval a first-class durable workflow instead of a batch job. EvalRunnerActor is a WorkflowActor with DurabilityFacet attached. Kill it mid-suite, restart it, and it resumes from the last checkpoint, skipping every scenario already scored. Long eval suites survive node restarts, which is a property that no external eval tool offers.


    MiniPi: the example

    MiniPi is a 12-actor agent eval pipeline, ported five ways: Go, Python, TypeScript, Rust WASM, and Rust embedded. They all produce the same output against the same PlexSpaces node:

    examples/go/apps/minipi/          # 1.5M WASM
    examples/python/apps/minipi/      # 47M WASM
    examples/typescript/apps/minipi/  # 13M WASM
    examples/rust/apps/minipi/        # 6.3M WASM
    examples/rust/embedded/minipi/    # native Rust, no WASM

    The 12 actors cover the full harness stack:

    ActorTypeWhat it does
    LLMGatewayActorGenServerOllama integration with KV response cache and mock fallback
    ToolRegistryActorGenServer4 built-in tools with JSON Schema validation
    AgentActorWorkflowActorOODA loop (Observe, Orient, Decide, Act)
    EvalRunnerActorWorkflowActorRuns scenarios in parallel, collects trajectories
    ScenarioStoreActorGenServer10 built-in test scenarios
    ScorerActorGenServerScores trajectories against rubrics
    TrajectoryStoreActorGenServerPersists agent trajectories for eval
    RegressionDetectorActorGenServerCompares scores across runs, flags drops over 5%
    BenchmarkActorWorkflowActorRuns the same scenarios against different harness configs
    AdvisorActorGenServerTwo-tier LLM: cheap model plus expensive advisor on demand
    ApprovalGateActorGenFSMHuman-in-the-loop: idle –> awaiting_approval –> idle
    DashboardActorGenServerAggregate view across all eval runs

    All 12 are declared in app-config.toml under a one_for_one supervision strategy. The framework starts them, watches them, and restarts individual actors on crash. Here’s how they connect. Notice there’s no separate “eval environment” bolted on the side and EvalRunnerActor calls the exact same AgentActor that production traffic calls:

    You can swap the debugging assistant for a support-ticket triager, a claims processor, or a code-review bot, and the same 12 actors still apply, only ScenarioStoreActor‘s scenarios and the tool schemas change.


    The OODA loop

    The agent itself is an AgentActor, a WorkflowActor running an OODA loop (Observe, Orient, Decide, Act).

    Here’s the core loop, from the Go implementation:

    // agent.go — the OODA loop
    // DurabilityFacet (priority 90) journals every message before this code runs.
    // Kill the process mid-loop. Restart. It picks up from the last checkpoint.
    
    for !loop.IterationLimitReached() {
        if loop.BudgetExceeded() {
            // Over token budget — finalize trajectory and return cleanly
            traj := loop.FinalizeTrajectory("budget_exceeded", iterations)
            a.exportTrajectory(traj)
            return result("budget_exceeded", traj)
        }
    
        // OBSERVE: load prior context from KV memory
        observations := a.doObserve(loop, task)
    
        // ORIENT: ask LLM gateway what to do next
        // LLM gateway tries Ollama first, falls back to mock, caches in KV
        plan := a.doOrient(loop, observations)
    
        // DECIDE: pick action. Does it need human approval?
        action := a.doDecide(loop, plan)
        if needsApproval(action) {
            loop.Suspend("action_needs_approval")
            return result("suspended", nil)
        }
    
        // ACT: run the tool through ToolRegistryActor
        // SchemaValidationFacet (priority 95) validates the call before the tool runs
        a.doAct(loop, action)
        loop.IncrementIteration()
    }
    
    traj := loop.FinalizeTrajectory("completed", iterations)
    a.exportTrajectory(traj) // writes to KV + posts TupleSpace tuple for eval collection

    Four things happen here that you’d otherwise have to build by hand:

    • Crash recovery is automatic. DurabilityFacet journals each message before the actor runs. Kill the node at iteration 7, restart it, and the loop resumes at iteration 8 without re-burning tokens.
    • Budget enforcement lives in AgentLoop. It counts tokens across every LLM call and stops the loop before you overspend.
    • Trajectory capture happens in exportTrajectory. Every Observe/Orient/Decide/Act step gets recorded with timing and token counts, written to KV storage, and posted as a TupleSpace tuple so EvalRunnerActor can find it.
    • Human approval is a durable suspend, not a poll. The agent serializes its full state and returns. ApprovalGateActor holds the request. When a human approves, the signal resumes the agent exactly where it paused even in the middle of a multi-hour run.

    Real output from test.sh against a running node:

    Step 5: AgentActor — OODA loop run
      workflow_run
      Status: completed  Outcome: completed
      Steps: 27  Trajectory: traj-01K... (27 steps in KV + TupleSpace index)

    Crash recovery: replay at the system level

    The durability property is worth slowing down on, because it’s categorically different from checkpointing you write yourself. When DurabilityFacet journals a message, it does so at the actor framework level, below your code. On restart, the journal replays those messages and your actor’s state comes back exactly as it was without application code to handle “resume from crash”.

    There are a few differences compared to Temporal, which also relies on replay. First, Temporal requires you to write workflow code as a deterministic function that can be safely replayed; PlexSpaces lets the actor’s message handling look like ordinary code, because the framework journals at the message boundary instead. Second, Temporal has no supervision tree so a crashing activity gets retried by the workflow, but nothing independently restarts just that piece while the rest keeps running. In PlexSpaces, one_for_one means a crashed AgentActor on scenario 3 restarts in isolation while EvalRunnerActor, ScorerActor, and every other scenario agent keep going.

    [supervisor]
    strategy = "one_for_one"
    max_restarts = 10
    max_restart_window_seconds = 60

    In LangGraph or CrewAI, one agent crashing typically kills the whole graph. Temporal can be made to handle this, but it takes explicit error-handling code. In PlexSpaces, independent crash isolation is just the default. For example, you might have have 20 tool calls into a host that gets OOM-killed. With DurabilityFacet, the node restarts, the journal replays, and the agent picks up at tool call 21.


    Validating tool calls without touching agent code

    Agents call tools with malformed arguments, which is unavoidable because models make mistakes. The real question is where you catch it. Catching it inside the tool handler is too late; execution has already started, and now you’re cleaning up a half-run call.

    SchemaValidationFacet catches it before the actor sees the message at all. An empty web_search query never reaches the tool registry because the facet returns a structured error, the agent corrects the call, and retries. The schema itself lives in app-config.toml, not in code:

    [[supervisor.children]]
    name = "tool_registry"
    actor_type = "minipi_wasm"
    behavior_kind = "GenServer"
    
    [[supervisor.children.facets]]
    type = "schema_validation"
    priority = 95
    
    [supervisor.children.facets.config]
    validation_mode = "strict"
    
    [supervisor.children.facets.config.method_schemas]
    web_search  = '{"type":"object","required":["query"],"properties":{"query":{"type":"string","minLength":1}}}'
    calculator  = '{"type":"object","required":["expression"],"properties":{"expression":{"type":"string"}}}'
    kv_read     = '{"type":"object","required":["key"],"properties":{"key":{"type":"string"}}}'
    kv_write    = '{"type":"object","required":["key","value"]}'

    Nothing about the tool actor’s code changes. The guardrail lives entirely in configuration.

    Step 6: SchemaValidationFacet — reject invalid method input
      reject empty query
      Schema validation: REJECTED (before actor sees it)
      Error contains: validation or schema
      valid call still works
      Valid call accepted: "web_search" executed successfully

    For example, you might have a billing support agent with a refund_customer tool. The model occasionally hallucinates a negative amount, a missing currency code, or an order ID that isn’t a string. Without a facet catching this, that call reaches your payments system and either throws an ugly stack trace or, worse, silently coerces bad input. With the schema in app-config.toml, the malformed call never leaves the tool registry. Instead, it bounces back to the agent as a structured error it can correct on the next turn, and your payments code never has to defend against it.


    Eval, running in the same runtime as production

    EvalRunnerActor is a WorkflowActor. It fans out one AgentActor per scenario, collects trajectories through TupleSpace, and scores them.

    // eval_runner.go — fan out and collect
    for i, scenario := range scenarios {
        // Spawn a fresh AgentActor for each scenario
        agentID := fmt.Sprintf("eval-agent-%s-%d", evalRunID, i)
        spawnedID, _ := host.Spawn("minipi_wasm", agentID, "agent_runner", map[string]string{
            "eval_run_id": evalRunID,
            "scenario_id": scenario.ID,
        })
    
        // Run the agent — same OODA loop as production
        resp, _ := host.Ask(spawnedID, "workflow_run", map[string]any{
            "task":        scenario.Input,
            "eval_run_id": evalRunID,
        }, 60000)
    
        // Collect trajectory directly from response
        if traj, ok := resp["trajectory"]; ok {
            trajectories = append(trajectories, traj)
        }
    }
    
    // Score each trajectory against the scenario's rubric
    for _, traj := range trajectories {
        score, _ := host.Ask("scorer", "score", map[string]any{
            "trajectory": traj,
            "rubric":     scenario.Rubric,
        }, 10000)
        scores = append(scores, score)
    }

    Because EvalRunnerActor is a WorkflowActor, killing it mid-eval is safe. Restart it and it skips scenarios that already finished. Long eval suites survive node restarts without losing progress. Real output from a 5-scenario run (Go):

    Step 10: EvalRunnerActor — 5-scenario standard suite
      eval smoke suite
      Pass rate: 0.833  Avg score: 0.775  Completed: 5 / 5
      Harness metrics: total_ms=125  coord_overhead=92.8%  speedup=5x  scenarios/sec=48
    
      sc-math-01     [XXXXXXXX--] 0.85
      sc-search-01   [XXXX------] 0.40
      sc-calc-01     [XXXXXXXX--] 0.85
      sc-reason-01   [XXXXXXXX--] 0.85
      sc-budget-01   [XXXXXXXX--] 0.85

    The TypeScript port tracks real token cost per scenario:

    Step 10: EvalRunnerActor — 5-scenario standard suite
      Pass rate: 0.4  Avg score: 0.818  Completed: 5 / 5
      Tokens: 311 in / 223 out  (est. cost: $0.00018)
    
      sc-math-01     0.92  (53 in / 41 out)
      sc-search-01   0.92  (63 in / 45 out)
      sc-calc-01     0.76  (60 in / 44 out)
      sc-reason-01   0.79  (62 in / 44 out)
      sc-budget-01   0.70  (73 in / 49 out)

    coord_overhead is the harness’s own overhead like spawning agents, collecting via TupleSpace, scoring. Across a 5-agent parallel run, it stays flat while compute scales, which is exactly the property you want: harness cost shouldn’t grow with the number of agents. For example, a legal team may need to run a contract review nightly across 200 incoming documents. Each document gets its own AgentActor, spawned by EvalRunnerActor the same way scenarios are spawned here. Because coordination happens through TupleSpace instead of a shared in-memory queue, one document’s agent hanging on a malformed PDF doesn’t block the other 199 and the batch survives a restart if the node needs to redeploy halfway through the night.


    Regression detection

    A single eval score doesn’t tell you much on its own. What matters is whether it’s better or worse than last time. RegressionDetectorActor stores baseline scores and flags any scenario that drops more than 5%:

    # regression_detector.py
    @handler("compare")
    def compare(self, eval_run_id: str = "", baseline_run_id: str = "") -> dict:
        baseline = self._load_scores(baseline_run_id)
        current  = self._load_scores(eval_run_id)
    
        regressions  = []
        improvements = []
    
        for scenario_id, base_score in baseline.items():
            curr_score = current.get(scenario_id, base_score)
            delta = curr_score - base_score
            if delta < -0.05:   # more than 5% drop
                regressions.append({"scenario_id": scenario_id, "delta": delta})
            elif delta > 0.05:
                improvements.append({"scenario_id": scenario_id, "delta": delta})
    
        return {
            "regressions":  len(regressions),
            "improvements": len(improvements),
            "details":      regressions + improvements,
        }
    Step 11: RegressionDetectorActor
      set_baseline
      Baseline set from eval-smoke-001 actual scores
      compare
      Regressions: 1  (sc-search-01 degraded by 0.20)
      Improvements: 1  (sc-reason-01 improved by 0.05)
      Regression detector caught degradation in search scenario

    The search scenario dropped by 20-point regression that would be invisible if the only thing you watched was the aggregate pass rate.


    Benchmarking harness configs, not just models

    The most underused insight in agent engineering: harness changes are often cheaper and more impactful than model changes. BenchmarkActor runs the same scenarios against multiple harness configurations. Here’s the Python implementation:

    # benchmark.py
    @handler("run_benchmark")
    def run_benchmark(self, scenario_suite: str = "smoke", configs: list = None) -> dict:
        results = []
        for config in (configs or self._default_configs()):
            # Run a full eval with this config
            eval_result = host.ask("eval_runner", {
                "action":       "run_suite",
                "suite":        scenario_suite,
                "eval_run_id":  f"bench-{config['name']}",
                "agent_config": config,
            }, timeout_ms=120000)
    
            results.append({
                "config":    config["name"],
                "score":     eval_result.get("avg_score", 0),
                "pass_rate": eval_result.get("pass_rate", 0),
                "tokens":    config.get("token_budget", 0),
                "max_iter":  config.get("max_iterations", 0),
            })
    
        winner = max(results, key=lambda r: r["score"])
        return {"configs": results, "winner": winner["config"]}

    Output from the Rust WASM port, on harder multi-step scenarios:

    Step 12: BenchmarkActor — 3-config comparison
      Configs tested: 3  Winner: aggressive  Best score: 0.83  Worst: 0.73
    
      aggressive    [XXXXXXXX--] score=0.830  budget=8192tok  max_iter=20
      balanced      [XXXXXXXX--] score=0.800  budget=4096tok  max_iter=10
      conservative  [XXXXXXX---] score=0.730  budget=1024tok  max_iter=3

    Same model, same scenarios, same prompt but the harness config alone moves quality by 14%. That’s the case for measuring this before reaching for a bigger model. The Python port, on simpler scenarios, tells a different story:

    Step 12: BenchmarkActor — 3-config comparison
      Configs tested: 3  Winner: conservative  Best score: 0.7
    
      conservative  [XXXXXXX---] score=0.700  budget=1024tok  max_iter=3
      balanced      [XXXXXXX---] score=0.700  budget=4096tok  max_iter=10
      aggressive    [XXXXXXX---] score=0.700  budget=8192tok  max_iter=20
      (on simple arithmetic tasks, all configs tie — benchmark your actual tasks)

    On simple arithmetic, budget doesn’t matter, the task fits in 3 iterations no matter what you give it. On multi-step research tasks, the loop limit starts to bite. Benchmark your own workload rather than trusting either result blindly. For example, a content-moderation team may need to decide how much iteration budget to give a policy-review agent. A conservative config (low budget, few iterations) is cheap but might miss nuance in a borderline post. An aggressive config catches more edge cases but costs more per review. BenchmarkActor runs last month’s flagged-content scenarios against both configs and reports the actual quality delta.


    Two-tier LLM: the advisor pattern

    Not every turn of the OODA loop needs the expensive model. Most turns are routine; only a handful demand deep reasoning. AdvisorActor implements a two-tier pattern: a fast, cheap model handles everything by default, and escalates to the expensive model only when its own confidence drops below a threshold.

    # advisor.py
    @handler("advise")
    def advise(self, prompt: str = "", context: dict = None) -> dict:
        # Always try the cheap model first
        fast_result = self._call_executor(prompt, context)
        self.total_requests += 1
    
        if fast_result.get("confidence", 1.0) >= self.confidence_threshold:
            # Confident enough — return fast result
            return fast_result
    
        # Low confidence — escalate to expensive advisor
        self.escalated += 1
        self.advisor_tokens += fast_result.get("tokens", 0)
        advisor_result = self._call_advisor(prompt, context, fast_result)
        self.advisor_tokens += advisor_result.get("tokens", 0)
        return advisor_result

    Rust, with harder prompts and a 60% escalation rate:

    Step 14: AdvisorActor — two-tier LLM
      Confidence threshold: 0.8
      Total requests: 5  Escalated: 3 / 5
      Escalation rate: 60.0%
      Advisor token share: 57.3%

    Python, with simpler prompts and a 40% escalation rate:

    Step 14: AdvisorActor — two-tier LLM
      Escalation rate: 40.0%
      Advisor token share: 33.6%
      Two-tier routing working — advisor escalated high-complexity prompts

    The two metrics that matter are escalation rate and advisor token share. Run your eval suite at threshold 0.9, then 0.7, then 0.5, and feed each result into BenchmarkActor. That’s how you find where the quality/cost tradeoff actually sits for your own tasks, instead of guessing. For example, a support-ticket classifier handling 10,000 tickets a day. Most are routine like “reset my password,” “where’s my order” and a cheap model nails them at near-100% confidence. The 5–10% that involve conflicting account details or ambiguous intent escalate to the expensive advisor. Routing every ticket through the expensive model would be needlessly costly; routing none of them through it would tank accuracy on the hard cases. The advisor pattern gets you both.


    Human-in-the-loop, without polling

    High-stakes actions need approval before they execute. The naive approach polls a status endpoint, which holds resources open, doesn’t survive a restart, and forces the agent to stay running the whole time. ApprovalGateActor is a GenFSM: idle –> awaiting_approval –> idle. The state is durable, e.g., kill the node while a request is pending, restart it, and the request is still there because the FSM state was journaled before the crash ever happened.

    # approval_gate.py
    class ApprovalGateActor:
        state: str = "idle"
        pending_request: dict = None
    
        @handler("request_approval")
        def request_approval(self, request: dict) -> dict:
            if self.state != "idle":
                return {"status": "busy", "current_state": self.state}
            self.state = "awaiting_approval"
            self.pending_request = request
            return {"status": "pending", "state": self.state}
    
        @handler("approve")
        def approve(self, approver_id: str = "", notes: str = "") -> dict:
            if self.state != "awaiting_approval":
                return {"error": "not_awaiting_approval"}
            decision = {"decision": "approved", "approver_id": approver_id, "notes": notes}
            self.decision_history.append(decision)
            self.state = "idle"
            self.pending_request = None
            return {"status": "approved", "state": self.state}
    Step 13: ApprovalGateActor — human-in-the-loop
      get_status idle
      FSM state: idle
      request_approval
      FSM state: awaiting_approval
      approve
      Approved by: alice@example.com
      FSM state: idle
      Decisions in history: 1

    The agent never polls. It calls loop.Suspend(), serializes its state, and returns. When approval comes through, the PlexSpaces runtime sends a resume signal, and the agent wakes up exactly where it paused without loss of state or rerun. For example, your debugging assistant runs 30 tool calls, identifies the root cause, and proposes a deploy. At the deploy_to_production call, it suspends. An on-call engineer reviews the trajectory in the dashboard and clicks approve. The agent resumes, and only the deployment step runs.


    The aggregate view

    After several eval runs, DashboardActor rolls everything up — scores, pass rates, trends over time:

    Step 15: DashboardActor — aggregate results
      Total evals: 4  Avg score: 0.767
    
      bench-001          [XXXXXX--] score=0.700  pass=70%
      eval-smoke-001     [XXXXX---] score=0.760  pass=40%
      eval-smoke-002     [XXXXX---] score=0.730  pass=40%
      test-999           [XXXXXXX-] score=0.880  pass=90%

    Four runs in this session: two smoke evals, one benchmark, one direct test. test-999 at 0.880 is a single high-confidence call. The smoke runs sit at 0.730–0.760, dragged down by the harder search scenario that consistently scores 0.40.


    How PlexSpaces compares

    FeatureLangGraphAutoGenCrewAIRestateTemporalPlexSpaces
    Crash recoveryNoNoNoJournal replayJournal replayJournal replay
    Supervision treesNoNoNoNoNoYes (one_for_one, one_for_all, rest_for_one)
    Eval in same runtimeNo (LangSmith)NoNoNoNoYes (same actors, same facets)
    Tool schema validationApp codeApp codeApp codeApp codeApp codeSchemaValidationFacet (config only)
    Human-in-the-loopInterruptNo native supportNo native supportSignalSignalGenFSM (durable state)
    PolyglotPythonPythonPythonTS/Java/Python/Go/RustTS/Java/Python/GoGo/Python/TS/Rust via WASM
    Multi-agent coordinationGraph edgesShared memoryRole handoffKeyed stateWorkflow stepsTupleSpace (Linda model)
    WASM sandboxingNoNoNoNoNoYes

    What MiniPi tests covers

    Each language port runs the same 15-step integration test. Every step exercises a production pattern, not a mock shortcut:

    StepWhat it testsKey metric
    1ScenarioStore: seed 10 built-in scenariosscenarios_stored=10
    2LLMGateway: Ollama with mock fallback and KV cacheprovider=ollama or mock
    3ToolRegistry: 4 tools with JSON Schema registeredtools_registered=4
    4SchemaValidation: empty query rejected before actor runsrejected_before_actor=true
    5AgentActor: full OODA loop, 10 iterations, budget enforcedoutcome=completed, steps=27–40
    6TrajectoryStore: persist and retrieve by IDtrajectory_id=traj-…
    7ScorerActor: two rubrics on the same trajectorytask_completion=0.85, tool_use=0.80
    8EvalRunnerActor: 5-scenario smoke suite, parallelpass_rate=0.40–0.83, avg=0.76–0.82
    9RegressionDetector: compare against baselineregressions=1, improvements=1
    10BenchmarkActor: 3 harness configs, same scenarioswinner by score
    11ApprovalGateActor: durable FSM wait and resumeidle –> awaiting –> idle
    12Second eval suite: drift detectionpass_rate compared to step 8
    13DashboardActor: first aggregate viewtotal_evals=2
    14AdvisorActor: two-tier routing, token splitescalation_rate=40%–60%
    15DashboardActor: final aggregate across all runstotal_evals=2–4, avg_score=0.767–0.81

    Running it

    You need a PlexSpaces node on port 8091, plus the toolchain for whichever port you want to run. Ollama with llama3.2 pulled is optional and test.sh falls back to a deterministic mock if Ollama isn’t running.

    # Using Docker (recommended)
    docker run -p 8000:8000 plexspaces/node:latest
    
    # Or build from source
    git clone https://github.com/plexobject/plexspaces.git
    cd plexspaces && make build
    # Go — 1.5M WASM, 5x parallel speedup, fastest eval
    cd examples/go/apps/minipi
    ./build.sh && ./test.sh 8091
    
    # Python — 47M WASM, most readable actor code, full object registry
    cd examples/python/apps/minipi
    ./build.sh && ./test.sh 8091
    
    # TypeScript — 13M WASM, token cost tracking per scenario
    cd examples/typescript/apps/minipi
    ./build.sh && ./test.sh 8091
    
    # Rust WASM — 6.3M WASM, most complete benchmark scoring
    cd examples/rust/apps/minipi
    ./build.sh && ./test.sh 8091
    
    # Rust embedded — no WASM, in-process node, fastest startup
    cd examples/rust/embedded/minipi
    ./test.sh

    Summary

    Agents fail for harness reasons more often than model reasons. For example, the loop exits too early; a malformed tool call crashes the agent; an eval suite runs against mocks, passes, and hands you false confidence. These are infrastructure problems, and infrastructure problems have infrastructure solutions. For example, in a debugging assistant mentioned above, the harness is what lets you restart a 40-tool-call investigation from step 37 instead of step 1. It’s what lets you gate a deployment on human approval without holding a thread open for ten minutes. It’s what lets you run ten scenarios in parallel and get a pass@k score you can actually trust before you ship.

    PlexSpaces brings together four decades of actor-model thinking like supervision trees from Erlang, TupleSpace coordination from Linda, durable workflows in the spirit of Temporal and wires them together specifically for agent workloads. The same runtime runs on a laptop and in production. The same actors used for eval are the actors that run in production. There’s no mismatch between the two. The harness is half the agent so build it like infrastructure.

    GitHub: https://github.com/bhatti/PlexSpaces


    Further reading

    Previous posts in this series:

    Example code and documentation:

    Older Posts »

    Powered by WordPress