---
title: "Architecture - Pilot Docs"
description: "How Pilot works under the hood"
url: "https://docs.calinora.io/features/architecture/"
---

# Architecture

Pilot is a single binary that connects to your Kafka cluster, continuously samples partition metrics, generates rebalancing proposals, and executes reassignments - all without an external database.

## High-Level Data Flow

```mermaid
graph TD
    K["Kafka Cluster<br/><small>Brokers · Topics · Partitions</small>"]

    K -->|"metadata + health"| CC["Kafka Client<br/><small>single connection pool</small>"]
    K -->|"watermarks + disk sizes"| MS["Metadata Sampler<br/><small>activity scoring</small>"]

    CC -->|"broker health<br/>partition metadata"| PE["Proposal Engine<br/><small>multi-phase pipeline</small>"]
    MS -->|"7 metrics per broker"| PE

    PE -->|"rebalancing proposal"| RE["Reassignment Executor<br/><small>streaming submission · throttle mgmt</small>"]
    RE -->|"alter_replica_assignment"| K
```

## Core Components

### Kafka Client

A single connection pool handles every Kafka interaction:

| Surface | Examples |
| - | - |
| **Metadata + health** | Cluster info, broker availability, partition watermarks, log directory sizes |
| **Reads** | Message browser, audit consumer, state topic consumers |
| **Mutations** | Reassignments, config updates, leader elections, ACLs, quotas |
| **State persistence** | `__pilot_broker_state`, `__pilot_audit_log`, `__pilot_pat_tokens` |

### Metadata Sampler

Pilot infers partition activity from Kafka metadata APIs - the sampling loop does not consume topic messages.

Every collection cycle (default 10s):

1. **Watermark deltas** - high watermark changes over time give message rates per partition
2. **Log directory sizes** - disk usage per partition per broker gives byte rates
3. **Follower lag aggregation** - the `OffsetLag` field from the `DescribeLogDirs` response is extracted and aggregated at cluster, broker, topic, and partition granularity, piggybacking on the same log directory call with no additional Kafka requests
4. **Follower calibration** - learns how much replication traffic followers actually generate relative to leader writes, smoothed over time per topic class

These are combined into seven metrics per broker: leader count, follower count, disk bytes, producer message rate, producer byte rate, consumer message rate, and consumer byte rate.

The **activity score** is a composite indicator that blends message rate, byte rate, disk size, and cumulative volume into a single 0-100 number per partition. Each component is log-normalized so large traffic differences don’t dominate the score. It is a dashboard “how busy is this broker?” summary only - the rebalancing optimizer works directly off the seven metrics above. The activity score is deliberately not one of them: it is a capped re-encoding of metrics the optimizer already balances, so including it would double-count producer throughput and disk in the objective.

### State Persistence

Pilot stores state in Kafka itself - no external database:

| Topic | Purpose | Policy | Created |
| - | - | - | - |
| `__pilot_broker_state` | Broker status (online/offline/maintenance) | Compacted | Always on startup |
| `__pilot_audit_log` | Audit trail of all mutations | Delete (30-day retention) | On startup when `AUDIT_ENABLED=true` (default) |
| `__pilot_pat_tokens` | Personal access token hashes | Compacted | When PATs are enabled |

`__pilot_broker_state` and `__pilot_pat_tokens` are created with 1 partition and the cluster’s default replication factor (1 if that fails). The audit topic follows `AUDIT_PARTITIONS` and `AUDIT_REPLICATION_FACTOR`; see [Audit Logging](https://docs.calinora.io/configuration/audit-logging/#topic-configuration). The cleanup policy is enforced idempotently on every startup, so manual topic misconfiguration is self-correcting.

See [Audit Logging](https://docs.calinora.io/configuration/audit-logging/) for audit topic configuration and [Authentication](https://docs.calinora.io/configuration/authentication/#personal-access-tokens-pats) for PAT configuration.

## Proposal Pipeline

The proposal engine generates rebalancing plans through a multi-phase pipeline:

```mermaid
graph TD
    CS["Cluster State"] --> PREP

    subgraph PREP["Preparation"]
        direction TB
        A1["Fetch metadata & check safety"] --> A2["Filter exclusions & fix critical issues"]
    end

    PREP --> OPT

    subgraph OPT["Optimization"]
        direction TB
        B1["Multi-level search"] --> B2["Cost reduction"]
    end

    OPT --> FIN

    subgraph FIN["Finalization"]
        direction TB
        C1["Spread & validate"] --> C2["Prune moves & rebuild metrics"] --> C3["Check benefit & diagnose limits"]
    end

    FIN --> OUT["Proposal"]
```

**Preparation** collects metadata, checks cluster safety (in-flight reassignments, ISR shrinks), filters excluded or cooling-down partitions, and fixes critical issues like under-replicated partitions and rack violations. Brokers with active follower lag are marked so the optimizer avoids sending them new inter-broker data. If every broker is within 10% of its fair share of every load with traffic (see [No Traffic](https://docs.calinora.io/features/proposals/#no-traffic)) and no topic is [concentrated](https://docs.calinora.io/features/proposals/#concentrated-topics), the pipeline short-circuits here.

**Optimization** is the core of the engine. It searches for partition moves that reduce imbalance using focused, broad, and swap candidate strategies. Each candidate is scored by how much it improves balance versus how expensive the move is. A second pass then reduces proposal cost by pruning unnecessary moves and replacing expensive inter-broker transfers with cheaper leader-only shifts.

**Finalization** moves leadership off [concentrated topics](https://docs.calinora.io/features/proposals/#concentrated-topics), validates correctness, and evaluates each partition’s net transition. It removes no-ops and changes that do not help, then rebuilds projected broker metrics from the final placements. Optional plans must still be [worth applying](https://docs.calinora.io/features/proposals/#reducing-unnecessary-movement) after pruning and any partition cap. If correctness repairs are needed, or a broker is [down](https://docs.calinora.io/features/proposals/#unavailable-brokers) and not in maintenance, the final proposal contains only repairs and defers optional balancing. Loads [at their limit](https://docs.calinora.io/features/proposals/#at-its-limit) produce **advisories** that say what would help.

Completed transitions start a 30-minute cooldown, in which optional balancing leaves those partitions alone. See [Reducing Unnecessary Movement](https://docs.calinora.io/features/proposals/#reducing-unnecessary-movement) for the worth-applying check and [Recent Moves](https://docs.calinora.io/features/proposals/#recent-moves) for the cooldown.

Application admission checks measurement readiness and evaluates the exact final plan against three fresh, compatible observations from the bounded measurement history. Separate optimizer runs do not need to select matching moves. Manual and automatic submission recheck the plan against current inputs; changes to its starting placements or destination eligibility require recalculation.

### Scoring

Metrics are weighted by operational priority. Leader and follower counts have the highest weight, followed by disk usage, then byte rates, then message rates.

### Movement Cost Model

The engine prefers cheaper moves over expensive ones. Leader-only changes (no data transfer) are cheapest, followed by inter-broker moves (proportional to partition size), with cross-rack transfers being the most expensive due to network cost.

## Reassignment Execution

When a proposal is applied, the executor uses **streaming submission**:

```mermaid
graph TD
    P["Proposal<br/><small>N partition moves</small>"] --> Q["Pending Queue"]

    Q --> B1["Broker 1<br/><small>active: 3/5 → submit</small>"]
    Q --> B2["Broker 2<br/><small>active: 5/5 → wait</small>"]
    Q --> B3["Broker 3<br/><small>active: 1/5 → submit</small>"]

    B1 --> KA["Kafka<br/><small>alter_replica_assignment</small>"]
    B3 --> KA

    KA --> RM["Reassignment Monitor<br/><small>Track progress · Detect stuck · Retry</small>"]
    RM -->|"completed"| Q
```

Unlike batch/wave execution, streaming continuously submits new moves as soon as a broker has capacity. This maximizes cluster utilization while respecting per-broker concurrency limits (`PILOT_MOVE_MAX_PER_BROKER`).

Both the concurrency limit and throttle rate (`PILOT_THROTTLE_RATE_MB`) can be adjusted live while a reassignment is in progress - changes take effect immediately without restarting. Throttle configs are cleaned up automatically after reassignment completes.

## Self-Healing

Three independent background loops run at different intervals:

```mermaid
graph LR
    subgraph Healing Loops
        direction TB
        C["Critical Fixes<br/><small>every 5m · highest priority</small><br/><small>URP repairs, rack violations</small>"]
        R["RF Increases<br/><small>every 15m · medium priority</small><br/><small>Replication factor adjustments</small>"]
        A["Activity Balancing<br/><small>every 30m · lowest priority</small><br/><small>Load rebalancing across brokers</small>"]
    end

    C -->|"blocks"| R
    R -->|"blocks"| A

    C --> G{Safe to apply?}
    R --> G
    A --> G

    G -->|"yes"| E["Execute Proposal"]
    G -->|"no"| S["Skip<br/><small>cooldown / window / dry-run</small>"]
```

Each loop obtains a proposal, checks safety signals, and applies eligible changes. Critical fixes run first, then RF increases, then activity balancing. While one loop runs, the others skip their cycle. All loops respect time windows, cooldowns, and dry-run mode. Activity balancing shares the default manual application’s measurement readiness, sustained-benefit and freshness requirements. Its per-run partition limit is included in generation so the final plan’s benefit is checked before application.

## AI Integration

### Chat Assistant

The chat engine connects a large language model to Pilot’s full API surface. Supported providers are Anthropic and OpenAI. You can also point at self-hosted or third-party models via `PILOT_CHAT_BASE_URL`, as long as they implement the Anthropic or OpenAI chat API format.

```mermaid
graph LR
    U["User Message"] --> LLM["LLM<br/><small>Claude / GPT</small>"]
    LLM -->|"tool call"| MCP["MCP Tool Executor"]

    MCP --> RO{Read-only?}
    RO -->|"yes"| EX["Execute"]
    RO -->|"no"| AP["Approval Request"]
    AP -->|"approved"| EX
    AP -->|"denied"| DN["Return denial"]

    EX -->|"tool result"| LLM
    DN -->|"tool result"| LLM
```

Mutation tools (reassign partitions, alter configs, delete ACLs) require explicit user approval before execution. Read-only tools execute immediately.

### Model Context Protocol (MCP)

The MCP server exposes 56 tools that external AI agents can invoke over HTTP. Tools are categorized as read (auto-execute) or mutate (require approval). See [MCP](https://docs.calinora.io/features/mcp/) for the full tool catalog.
