Calinora Pilot API (0.21.0)

Download OpenAPI specification:Download

Calinora Pilot API for comprehensive Kafka cluster management and monitoring

Calinora Pilot is an all-in-one, container-native tool that live-samples Kafka traffic, scores partition activity, exposes management APIs (including intelligent partition redistribution) and displays everything in a modern React dashboard.

Features

  • Real-time partition activity monitoring and sampling
  • Intelligent partition redistribution with multiple strategies
  • Broker maintenance mode management
  • Topic configuration management with bulk operations
  • Cluster health monitoring
  • Reassignment monitoring and cancellation

Authentication

When authentication is enabled, all API endpoints (except health, metrics, and auth login flows) require a valid session or token.

Session-based (UI): Users authenticate via OAuth2/OIDC and receive a session cookie.

Bearer token (API/MCP): API clients and MCP tools authenticate via Authorization: Bearer <token> header using a Personal Access Token (PAT). PATs are long-lived, revocable tokens created from the Pilot UI.

MCP endpoint: The Model Context Protocol server is available at /mcp for AI agent integration.

Agents

Pilot Agent management, deployment, enrollment, and broker lifecycle operations

List connected Pilot Agents

Returns the state of all agents known to Pilot. When agent management is disabled returns an empty list with enabled: false.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Download the agent binary

Returns the Pilot Agent binary for the requested architecture as an application/octet-stream. The binary is fetched from the configured upstream download URL and cached locally by Pilot.

query Parameters
arch
string
Default: "amd64"
Enum: "amd64" "arm64"

CPU architecture

Responses

Response samples

Content type
application/json
{
  • "success": false,
  • "error": "string",
  • "details": "string"
}

Get the agent binary SHA256 checksum

Returns the full SHA256 hex digest of the cached agent binary for the requested architecture.

query Parameters
arch
string
Default: "amd64"
Enum: "amd64" "arm64"

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get agent binary configuration

Returns the current upstream download URL, the agent version, and a map of architectures for which a binary is cached.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update agent binary download URL

Updates the upstream URL from which Pilot fetches the agent binary. The change is not persisted and reverts on restart. Requires a valid license (management mode).

Request Body schema: application/json
downloadUrl
required
string

Responses

Request samples

Content type
application/json
{
  • "downloadUrl": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get agent manual install script

Returns a shell script that downloads the agent binary and provides systemd unit template instructions. Intended to be piped into bash on the target host.

Responses

Sign an agent certificate from a CSR

Called by an agent during enrollment. Accepts either a single-use enrollment token or the persistent master bootstrap token. Single-use tokens are consumed on successful enrollment. Not license-gated because it is part of the agent bootstrap flow. Bootstrap requests over plain HTTP are logged with a security warning.

Request Body schema: application/json
agentId
required
string
token
required
string
csrPEM
required
string

PEM-encoded certificate signing request

Responses

Request samples

Content type
application/json
{
  • "agentId": "string",
  • "token": "string",
  • "csrPEM": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Create a single-use agent enrollment token

Creates a new single-use bootstrap token that an agent can exchange for a signed certificate. Requires a valid license (management mode).

Request Body schema: application/json
label
string
ttlMinutes
integer

Token TTL in minutes. 0 or omitted uses the default TTL.

Responses

Request samples

Content type
application/json
{
  • "label": "string",
  • "ttlMinutes": 0
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

List active agent enrollment tokens

Returns active enrollment tokens with masked token values.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Revoke an agent enrollment token

Revokes a previously issued single-use enrollment token. Requires a valid license (management mode).

Request Body schema: application/json
token
required
string

Responses

Request samples

Content type
application/json
{
  • "token": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Deploy a Pilot Agent to a host via SSH

Deploys a Pilot Agent to a remote host using SSH (optionally via one or more jump hosts). Requires HTTPS and a valid license (management mode).

Request Body schema: application/json
required
object (SSHHop)
Array of objects (SSHHop)
brokerId
integer <int32>
pilotServerAddress
required
string

gRPC address used by the agent to connect back to Pilot (e.g. pilot:9190).

bootstrapUrl
string

HTTP bootstrap URL used by the agent for certificate renewal. Defaults to the incoming request URL.

arch
string
Default: "amd64"
Enum: "amd64" "arm64"
serviceUser
string
Default: "root"

OS user the agent systemd service runs as. Defaults to root. Set to pilot-agent to run the agent under a dedicated non-root user (the deploy flow will create the user and configure a narrow sudoers drop-in automatically).

kafkaUnit
string

Systemd unit name of the Kafka broker on the target host. Optional; auto-detected from the running Kafka JVM process when omitted. Comma-separated list for hosts running multiple brokers.

kafkaGroup
string

Group that owns Kafka log files. Optional; auto-detected from the running unit when omitted. The pilot-agent user is added to this group so it can tail Kafka logs.

sudoPassword
string

Password for sudo on the target host. Empty means passwordless sudo is attempted.

Responses

Request samples

Content type
application/json
{
  • "target": {
    },
  • "proxyHops": [
    ],
  • "brokerId": 0,
  • "pilotServerAddress": "string",
  • "bootstrapUrl": "string",
  • "arch": "amd64",
  • "serviceUser": "root",
  • "kafkaUnit": "kafka.service",
  • "kafkaGroup": "kafka",
  • "sudoPassword": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Deploy Pilot Agents to multiple hosts via SSH

Deploys Pilot Agents to up to 50 target hosts sequentially. Requires HTTPS and a valid license (management mode).

Request Body schema: application/json
required
Array of objects (BulkDeployTarget) <= 50 items
Array of objects (SSHHop)
pilotServerAddress
required
string

gRPC address used by the agent to connect back to Pilot (e.g. pilot:9190).

bootstrapUrl
string

HTTP bootstrap URL used by the agent for certificate renewal. Defaults to the incoming request URL.

arch
string

Shared architecture. Empty string auto-detects per host.

serviceUser
string
Default: "root"

OS user the agent systemd service runs as. Defaults to root. Set to pilot-agent to run the agent under a dedicated non-root user (the deploy flow will create the user and configure a narrow sudoers drop-in automatically). Applied to every host in the bulk batch.

kafkaUnit
string

Systemd unit name of the Kafka broker on each target host. Optional; auto-detected from the running Kafka JVM process when omitted. Comma-separated list for hosts running multiple brokers. Applied to every host in the bulk batch.

kafkaGroup
string

Group that owns Kafka log files. Optional; auto-detected from the running unit when omitted. The pilot-agent user is added to this group so it can tail Kafka logs. Applied to every host in the bulk batch.

sudoPassword
string

Password for sudo on the target host. Empty means passwordless sudo is attempted.

Responses

Request samples

Content type
application/json
{
  • "targets": [
    ],
  • "proxyHops": [
    ],
  • "pilotServerAddress": "string",
  • "bootstrapUrl": "string",
  • "arch": "string",
  • "serviceUser": "root",
  • "kafkaUnit": "kafka.service",
  • "kafkaGroup": "kafka",
  • "sudoPassword": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Test SSH connectivity for agent deploy

Runs the SSH connection, auth, and discovery steps of a deploy without installing anything. Returns host key material for trust-on-first-use. Requires HTTPS and a valid license (management mode).

Request Body schema: application/json
required
object (SSHHop)
Array of objects (SSHHop)
brokerId
integer <int32>
pilotServerAddress
required
string

gRPC address used by the agent to connect back to Pilot (e.g. pilot:9190).

bootstrapUrl
string

HTTP bootstrap URL used by the agent for certificate renewal. Defaults to the incoming request URL.

arch
string
Default: "amd64"
Enum: "amd64" "arm64"
serviceUser
string
Default: "root"

OS user the agent systemd service runs as. Defaults to root. Set to pilot-agent to run the agent under a dedicated non-root user (the deploy flow will create the user and configure a narrow sudoers drop-in automatically).

kafkaUnit
string

Systemd unit name of the Kafka broker on the target host. Optional; auto-detected from the running Kafka JVM process when omitted. Comma-separated list for hosts running multiple brokers.

kafkaGroup
string

Group that owns Kafka log files. Optional; auto-detected from the running unit when omitted. The pilot-agent user is added to this group so it can tail Kafka logs.

sudoPassword
string

Password for sudo on the target host. Empty means passwordless sudo is attempted.

Responses

Request samples

Content type
application/json
{
  • "target": {
    },
  • "proxyHops": [
    ],
  • "brokerId": 0,
  • "pilotServerAddress": "string",
  • "bootstrapUrl": "string",
  • "arch": "amd64",
  • "serviceUser": "root",
  • "kafkaUnit": "kafka.service",
  • "kafkaGroup": "kafka",
  • "sudoPassword": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Deploy a Pilot Agent and stream progress via SSE

Same as /agents/deploy but emits Server-Sent Events with per-step progress. Events: step for each deploy step, done with the final result, or error on failure. Requires HTTPS and a valid license (management mode).

Request Body schema: application/json
required
object (SSHHop)
Array of objects (SSHHop)
brokerId
integer <int32>
pilotServerAddress
required
string

gRPC address used by the agent to connect back to Pilot (e.g. pilot:9190).

bootstrapUrl
string

HTTP bootstrap URL used by the agent for certificate renewal. Defaults to the incoming request URL.

arch
string
Default: "amd64"
Enum: "amd64" "arm64"
serviceUser
string
Default: "root"

OS user the agent systemd service runs as. Defaults to root. Set to pilot-agent to run the agent under a dedicated non-root user (the deploy flow will create the user and configure a narrow sudoers drop-in automatically).

kafkaUnit
string

Systemd unit name of the Kafka broker on the target host. Optional; auto-detected from the running Kafka JVM process when omitted. Comma-separated list for hosts running multiple brokers.

kafkaGroup
string

Group that owns Kafka log files. Optional; auto-detected from the running unit when omitted. The pilot-agent user is added to this group so it can tail Kafka logs.

sudoPassword
string

Password for sudo on the target host. Empty means passwordless sudo is attempted.

Responses

Request samples

Content type
application/json
{
  • "target": {
    },
  • "proxyHops": [
    ],
  • "brokerId": 0,
  • "pilotServerAddress": "string",
  • "bootstrapUrl": "string",
  • "arch": "amd64",
  • "serviceUser": "root",
  • "kafkaUnit": "kafka.service",
  • "kafkaGroup": "kafka",
  • "sudoPassword": "string"
}

Response samples

Content type
application/json
{
  • "success": false,
  • "error": "string",
  • "details": "string"
}

Bulk deploy Pilot Agents and stream progress via SSE

Same as /agents/deploy/bulk but emits Server-Sent Events with per-host progress. Events: step for each host step, done with the final summary, or error on failure. Requires HTTPS and a valid license (management mode).

Request Body schema: application/json
required
Array of objects (BulkDeployTarget) <= 50 items
Array of objects (SSHHop)
pilotServerAddress
required
string

gRPC address used by the agent to connect back to Pilot (e.g. pilot:9190).

bootstrapUrl
string

HTTP bootstrap URL used by the agent for certificate renewal. Defaults to the incoming request URL.

arch
string

Shared architecture. Empty string auto-detects per host.

serviceUser
string
Default: "root"

OS user the agent systemd service runs as. Defaults to root. Set to pilot-agent to run the agent under a dedicated non-root user (the deploy flow will create the user and configure a narrow sudoers drop-in automatically). Applied to every host in the bulk batch.

kafkaUnit
string

Systemd unit name of the Kafka broker on each target host. Optional; auto-detected from the running Kafka JVM process when omitted. Comma-separated list for hosts running multiple brokers. Applied to every host in the bulk batch.

kafkaGroup
string

Group that owns Kafka log files. Optional; auto-detected from the running unit when omitted. The pilot-agent user is added to this group so it can tail Kafka logs. Applied to every host in the bulk batch.

sudoPassword
string

Password for sudo on the target host. Empty means passwordless sudo is attempted.

Responses

Request samples

Content type
application/json
{
  • "targets": [
    ],
  • "proxyHops": [
    ],
  • "pilotServerAddress": "string",
  • "bootstrapUrl": "string",
  • "arch": "string",
  • "serviceUser": "root",
  • "kafkaUnit": "kafka.service",
  • "kafkaGroup": "kafka",
  • "sudoPassword": "string"
}

Response samples

Content type
application/json
{
  • "success": false,
  • "error": "string",
  • "details": "string"
}

Get agent state for a single broker

path Parameters
brokerId
required
integer <int32>

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Upgrade an agent via SSH

Copies a new agent binary to the target host over SSH and restarts the agent service. Requires HTTPS and a valid license (management mode).

path Parameters
brokerId
required
integer <int32>
Request Body schema: application/json
required
object (SSHHop)
Array of objects (SSHHop)
brokerId
integer <int32>
arch
string
Default: "amd64"
Enum: "amd64" "arm64"
sudoPassword
string

Password for sudo on the target host. Empty means passwordless sudo is attempted.

Responses

Request samples

Content type
application/json
{
  • "target": {
    },
  • "proxyHops": [
    ],
  • "brokerId": 0,
  • "arch": "amd64",
  • "sudoPassword": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Upgrade an agent via gRPC self-upgrade

Sends an upgrade command over the existing gRPC stream; the agent downloads the new binary from Pilot and restarts itself. No SSH required. Requires a valid license (management mode).

path Parameters
brokerId
required
integer <int32>

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Restart a broker via its agent

Sends a restart command to the broker's Pilot Agent. Requires a valid license (management mode).

path Parameters
brokerId
required
integer <int32>
Request Body schema: application/json
graceful
boolean
Default: true
timeoutSeconds
integer [ 1 .. 600 ]
Default: 300

Responses

Request samples

Content type
application/json
{
  • "graceful": true,
  • "timeoutSeconds": 300
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Stop a broker via its agent

Requires a valid license (management mode).

path Parameters
brokerId
required
integer <int32>
Request Body schema: application/json
graceful
boolean
Default: true
timeoutSeconds
integer [ 1 .. 600 ]
Default: 300

Responses

Request samples

Content type
application/json
{
  • "graceful": true,
  • "timeoutSeconds": 300
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Start a broker via its agent

Requires a valid license (management mode).

path Parameters
brokerId
required
integer <int32>
Request Body schema: application/json
graceful
boolean
Default: true
timeoutSeconds
integer [ 1 .. 600 ]
Default: 300

Responses

Request samples

Content type
application/json
{
  • "graceful": true,
  • "timeoutSeconds": 300
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Stream broker logs via SSE

Subscribes to the broker's log stream via Server-Sent Events. Filters can be applied by level or substring pattern. Recent buffered entries are backfilled on connect. Each event is event: log with a JSON payload.

path Parameters
brokerId
required
integer <int32>
query Parameters
level
string
Enum: "TRACE" "DEBUG" "INFO" "WARN" "ERROR"

Minimum log level

pattern
string

Substring filter applied to each log line

tail
integer <= 1000
Default: 200

Number of recent entries to replay on connect

Responses

Response samples

Content type
application/json
{
  • "success": false,
  • "error": "string",
  • "details": "string"
}

Search buffered broker logs

Searches the in-memory log buffer for a broker and returns matching entries.

path Parameters
brokerId
required
integer <int32>
query Parameters
q
string

Substring search term

level
string
Enum: "TRACE" "DEBUG" "INFO" "WARN" "ERROR"
from
integer <int64>

Lower bound timestamp (unix millis)

to
integer <int64>

Upper bound timestamp (unix millis)

limit
integer <= 1000
Default: 500

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Rolling Restart

Rolling broker restart orchestration and preflight safety gates

Start a rolling restart

Starts a rolling restart across all brokers using the connected Pilot Agents. Runs preflight safety checks (ISR health, URP counts, connected agents, pending reassignments) before starting. Returns 409 if preflight fails or a restart is already running. Requires a valid license (management mode).

Request Body schema: application/json
gracefulShutdown
boolean
Default: true
stallTimeoutSeconds
integer [ 30 .. 600 ]
Default: 120

Universal stall timeout. Phases abort after this many seconds of zero progress. Clamped to [30, 600].

stabilizationSeconds
integer [ 0 .. 120 ]
Default: 0

Dwell time between brokers after a successful restart. Clamped to [0, 120].

readinessSoakSeconds
integer [ 1 .. 60 ]
Default: 10

Soak period after initial readiness before proceeding. Clamped to [1, 60].

Responses

Request samples

Content type
application/json
{
  • "gracefulShutdown": true,
  • "stallTimeoutSeconds": 120,
  • "stabilizationSeconds": 0,
  • "readinessSoakSeconds": 10
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get rolling restart status

Returns the current rolling restart state. When agent management is disabled or no restart has been run, returns {status: "IDLE"}.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Cancel an active rolling restart

Cancels the currently running rolling restart. Remaining brokers are marked as skipped. If a broker was stopped when cancelled, a best-effort start command is sent. Requires a valid license (management mode).

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Run rolling restart preflight checks

Runs all preflight safety gates (ISR, URPs, connected agents, pending reassignments, controller location) without starting a restart. Returns the individual gate results, the computed broker order, and whether the restart can proceed.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Health

Health check operations

Liveness probe

Lightweight liveness probe for Kubernetes. Returns 200 OK if the service is running. Does not check external dependencies. Use /ready for readiness checks.

Responses

Response samples

Content type
application/json
{
  • "status": "ok"
}

Readiness probe

Readiness probe for Kubernetes. Returns 200 OK if the service is ready to accept traffic. Checks that Kafka brokers are available.

Responses

Response samples

Content type
application/json
{
  • "status": "ready",
  • "reason": "no available Kafka brokers"
}

Get application version

Returns the version and build information for the Pilot service

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "status": "healthy",
  • "version": "0.10.0",
  • "service": "pilot",
  • "go_version": "go1.23.0",
  • "build_time": "2025-01-04T12:00:00Z",
  • "git_commit": "abc123def"
}

Get self-healing status

Returns the current state of all three self-healing loops (activity, critical, RF), including scheduling, execution status, and configuration settings.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    },
  • "timestamp": "2019-08-24T14:15:22Z"
}

Cluster

Cluster information and management

Get cluster information

Returns basic cluster information including brokers and topics

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get comprehensive cluster information

Returns comprehensive cluster information including sampling status and detailed partition data

query Parameters
summary
boolean
Default: false

When true, excludes per-partition details and returns only topic names and partition counts (60-80% smaller response)

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get cluster log directories

Returns log directory information for all brokers in the cluster

query Parameters
summary
boolean
Default: false

When true, returns only aggregated broker-level data without per-partition details (90%+ smaller response)

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": { }
}

Get cluster health

Returns cached cluster health information including broker availability and partition health

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get follower lag metrics

Returns aggregated follower lag (offset lag) at cluster, broker, topic, and partition levels. Data is derived from the DescribeLogDirs API OffsetLag field, refreshed every ~10 seconds. Also classifies under-replicated partitions by cause (offline broker vs follower lag).

query Parameters
topic
string

Filter to a specific topic (includes per-partition detail)

broker
integer

Filter to a specific broker

summary
boolean
Default: false

When true, omits per-partition details from topics (smaller response)

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Refresh cluster metadata

Manually triggers a metadata refresh for the cluster and sampler

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "timestamp": "2019-08-24T14:15:22Z",
  • "message": "string",
  • "successful_components": [
    ],
  • "failed_components": [
    ]
}

Partitions

Partition activity and monitoring

Get partition activity

Returns real-time partition activity data including sampling statistics and rates

query Parameters
minimal
boolean
Default: false

When true, excludes detailed fields (messageCount, avgMessageSize, samplingStatus, healthStatus, watermarks etc.) returning only core partition metadata and rates (70-80% smaller response)

fields
string
Example: fields=topic,partition,leader,replicas,isr,messageRate,byteRate,minISR

Comma-separated list of fields to include per partition (e.g. "topic,partition,leader,replicas,isr,messageRate,byteRate"). Only specified fields are returned. Applies after minimal/full mode field population.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Redistribution

Partition redistribution operations

Redistribute topic partitions

Redistributes all partitions of a topic using the specified strategy

path Parameters
topic
required
string

Name of the topic to redistribute

Request Body schema: application/json
strategy
required
string
Enum: "ROUND_ROBIN" "RANDOM" "BALANCED" "RACK_AWARE"

Redistribution strategy to use

excludedBrokers
Array of integers <int32> [ items <int32 > ]

Broker IDs to exclude from redistribution

Responses

Request samples

Content type
application/json
{
  • "strategy": "BALANCED",
  • "excludedBrokers": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string",
  • "reassignmentId": "string"
}

Redistribute single partition

Redistributes a single partition to different brokers

path Parameters
topic
required
string

Name of the topic

partition
required
integer <int32>

Partition number to redistribute

Request Body schema: application/json
targetBrokers
required
Array of integers <int32> [ items <int32 > ]

Target broker IDs for the partition

Responses

Request samples

Content type
application/json
{
  • "targetBrokers": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string",
  • "reassignmentId": "string"
}

Bulk partition reassignment

Reassigns multiple partitions in a single operation

path Parameters
topic
required
string

Name of the topic

Request Body schema: application/json
description
string

Optional description for the bulk operation

required
Array of objects (RedistributePartitionRequest) non-empty

List of partition redistribution requests

Responses

Request samples

Content type
application/json
{
  • "description": "Rebalancing user-events topic for better distribution",
  • "partitions": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string",
  • "reassignmentId": "string"
}

Brokers

Broker management operations

Get broker racks

Returns the rack layout for rack-aware algorithms

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Set broker maintenance mode

Marks a broker for maintenance mode and redistributes its partitions to other brokers

path Parameters
brokerId
required
integer <int32>

ID of the broker to put in maintenance mode

Request Body schema: application/json
strategy
required
string
Enum: "ROUND_ROBIN" "RANDOM" "BALANCED" "RACK_AWARE"

Redistribution strategy to use for evacuating partitions

preserveLeaders
boolean
Default: true

Whether to preserve current leaders when possible

preview
boolean
Default: false

If true, only return a preview of changes without executing them

force
boolean
Default: false

Allow maintenance even if it causes under-replicated partitions

targetBrokers
Array of integers

Optional list of specific brokers to move partitions to. If not provided, partitions will be distributed across all available brokers.

Responses

Request samples

Content type
application/json
{
  • "strategy": "RACK_AWARE",
  • "preserveLeaders": true,
  • "preview": false,
  • "force": false,
  • "targetBrokers": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Remove broker from maintenance mode

Removes a broker from maintenance mode

path Parameters
brokerId
required
integer <int32>

ID of the broker to remove from maintenance mode

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Move broker log directories

Moves replicas stored in source log directories to target log directories on the specified broker

path Parameters
brokerId
required
integer

ID of the broker to move log directories on

Request Body schema: application/json
strategy
string
Default: "BALANCE_ALL"
Enum: "BALANCE_ALL" "BALANCE_SELECTED"

Move strategy. Both strategies use the same two-pass balancing algorithm (disk size first, then partition count). BALANCE_ALL operates on all log dirs. BALANCE_SELECTED operates on the specified source + target dirs only.

sourceLogDirs
Array of strings

Log directory paths to move replicas from (required for BALANCE_SELECTED)

targetLogDirs
Array of strings

Log directory paths to move replicas to (required for BALANCE_SELECTED)

Responses

Request samples

Content type
application/json
{
  • "strategy": "BALANCE_ALL",
  • "sourceLogDirs": [
    ],
  • "targetLogDirs": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "brokerId": 0,
  • "totalPartitions": 0,
  • "movedPartitions": 0,
  • "failedPartitions": 0,
  • "errors": [
    ]
}

Topic Configuration

Topic configuration management

Get topic configuration

Returns the configuration for a specific topic

path Parameters
topic
required
string

Name of the topic

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "configs": {
    }
}

Update topic configuration

Updates the configuration for a specific topic

path Parameters
topic
required
string

Name of the topic

Request Body schema: application/json
required
object

Configuration key-value pairs to update

Responses

Request samples

Content type
application/json
{
  • "configs": {
    }
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Get explicit topic configuration

Returns only explicitly set (non-default) configurations for a topic

path Parameters
topic
required
string

Name of the topic

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "configs": {
    }
}

Get default configurations

Returns cluster default topic configurations

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "defaults": {
    },
  • "message": "string"
}

Bulk update topic configurations

Updates configurations for multiple topics in a single operation

Request Body schema: application/json
required
Array of objects

Responses

Request samples

Content type
application/json
{
  • "topics": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Search topics by configuration

Searches topics based on configuration criteria

Request Body schema: application/json
required
Array of objects

Responses

Request samples

Content type
application/json
{
  • "criteria": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "topics": [
    ],
  • "totalMatches": 0
}

Client Quotas

Client quota management

List client quotas

List configured client quota entities, including defaults.

query Parameters
entityType
string
Enum: "default" "default-user" "default-client-id" "default-user-client-id" "user" "user-default-client-id" "client-id" "user-client-id"

Filter by entity scope

user
string

User principal to filter by

clientId
string

Client ID to filter by

includeEffective
boolean

Include effective quotas for the provided user/clientId pair

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Resolve effective client quotas

Resolve the effective quotas for a user and client-id pair.

Request Body schema: application/json
user
string
clientId
string

Responses

Request samples

Content type
application/json
{
  • "user": "string",
  • "clientId": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Analyze quota bottlenecks

Correlates configured client quotas with observed throughput to detect potential bottlenecks.

Consumer analysis (high confidence): Compares each consumer group member's actual consumption byte rate against their resolved consumer_byte_rate quota. Flags as bottleneck when utilization >= 80%.

Producer analysis (medium/low confidence): Heuristically compares topic production byte rates against configured producer_byte_rate quotas. Since Pilot uses metadata-only monitoring, it cannot determine which specific producer writes to which topic. Named client-id entities receive "medium" confidence; default entities receive "low" confidence.

Results are cached for 30 seconds and automatically refreshed after each consumer group collection cycle.

query Parameters
minUtilization
number
Default: 0

Only return entries with utilization percentage at or above this value

direction
string
Enum: "consume" "produce"

Filter by quota direction

bottlenecksOnly
boolean
Default: false

Only return entries flagged as bottlenecks or at-risk

search
string

Case-insensitive substring search across groupId, clientId, topics (consumers) and user, clientId, entityType (producers)

sortBy
string
Default: "utilization"

Field to sort results by. Consumer table: utilization (default), groupId, clientId, quotaByteRate, consumeByteRate. Producer table: riskPercent (default), user, clientId, quotaByteRate.

sortDir
string
Default: "desc"
Enum: "asc" "desc"

Sort direction

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    },
  • "meta": {
    }
}

Get default client quotas

Returns configured default quotas (empty if none configured).

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update default client quotas

Create or update default quota keys.

Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete default client quotas

Removes all configured default quotas.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get default user quotas

Returns quotas for the default user entity.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update default user quotas

Create or update default user quota keys.

Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete default user quotas

Removes all configured default user quotas.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get default client-id quotas

Returns quotas for the default client-id entity.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update default client-id quotas

Create or update default client-id quota keys.

Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete default client-id quotas

Removes all configured default client-id quotas.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get default user + client-id quotas

Returns configured quotas for the default user and specific client-id.

path Parameters
clientId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update default user + client-id quotas

Create or update quotas for the default user and client-id.

path Parameters
clientId
required
string
Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete default user + client-id quotas

Removes all quotas for the default user and client-id.

path Parameters
clientId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get user + default client-id quotas

Returns configured quotas for a user across all client-ids.

path Parameters
user
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update user + default client-id quotas

Create or update quotas for a user across all client-ids.

path Parameters
user
required
string
Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete user + default client-id quotas

Removes all quotas for a user across all client-ids.

path Parameters
user
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get user quotas

Returns configured quotas for a user.

path Parameters
user
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update user quotas

Create or update user quota keys.

path Parameters
user
required
string
Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete user quotas

Removes all quotas for a user.

path Parameters
user
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get client-id quotas

Returns configured quotas for a client-id.

path Parameters
clientId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update client-id quotas

Create or update client-id quota keys.

path Parameters
clientId
required
string
Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete client-id quotas

Removes all quotas for a client-id.

path Parameters
clientId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get combined user and client-id quotas

Returns configured quotas for a user and client-id pair.

path Parameters
user
required
string
clientId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update combined quotas

Create or update quotas for a user and client-id pair.

path Parameters
user
required
string
clientId
required
string
Request Body schema: application/json
object
delete
Array of strings

Responses

Request samples

Content type
application/json
{
  • "set": {
    },
  • "delete": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete combined quotas

Removes all quotas for a user and client-id pair.

path Parameters
user
required
string
clientId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Proposals

Redistribution proposal operations

Retrieve current redistribution proposal

Returns the current multi-objective, rack-aware redistribution proposal. Proposals are generated in the background on a schedule. If a proposal is not yet available or generation is in progress, the endpoint returns HTTP 202 with success: true and metadata describing the next generation time.

Request Body schema: application/json
immediate
boolean
Default: false

If true, triggers synchronous generation and returns fresh proposal

Responses

Request samples

Content type
application/json
{
  • "immediate": false
}

Response samples

Content type
application/json
{
  • "success": true,
  • "proposal": {
    },
  • "advisories": { },
  • "metadata": {
    }
}

List proposals

Returns proposals from the in-memory store (newest first)

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Get proposal by ID

Returns a stored proposal by ID. Special ID "current" returns the currently generated proposal if available.

path Parameters
proposalId
required
string

Responses

Response samples

Content type
application/json
Example
{
  • "success": true,
  • "data": {
    }
}

Delete proposal by ID

Deletes a stored proposal by ID. Use "current" to clear the in-memory current proposal cache.

path Parameters
proposalId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Apply proposal

Applies the specified proposal. Special ID "current" applies the current proposal.

path Parameters
proposalId
required
string

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string",
  • "data": {
    }
}

Get current proposal

Returns the currently generated proposal if available.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Optimization

Partition optimization recommendations

Get partition optimization

Returns partition optimization recommendations

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": { }
}

Leadership

Leader election operations

Cluster-wide preferred leader election

Triggers preferred leader election for the entire cluster or specified topics

Request Body schema: application/json

Optional. When omitted or empty, triggers election for all non-internal topics cluster-wide.

required
Array of objects

List of topics and their partitions for preferred leader election

Responses

Request samples

Content type
application/json
{
  • "topics": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Topic-scoped preferred leader election

Triggers preferred leader election for a specific topic

path Parameters
topic
required
string

Name of the topic

Request Body schema: application/json
required
Array of objects

List of topics and their partitions for preferred leader election

Responses

Request samples

Content type
application/json
{
  • "topics": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Reassignments

Reassignment monitoring and management

Get all reassignments

Returns all partition reassignments

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Get active reassignments

Returns currently running partition reassignments

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Get reassignment history

Returns completed partition reassignments

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Get reassignment tuning

Returns current throttle and concurrency settings for active reassignments

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update reassignment tuning

Updates throttle rate and/or max concurrent moves per broker for active reassignments

Request Body schema: application/json
throttleRateMBps
number <double>

Throttle rate in MB/s (leader/follower)

maxConcurrentMovesPerBroker
integer

Maximum concurrent moves per broker

Responses

Request samples

Content type
application/json
{
  • "throttleRateMBps": 150,
  • "maxConcurrentMovesPerBroker": 20
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Resync reassignment monitor

Forces an immediate reassessment of active reassignments to clear stale states

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Get reassignment by ID

Returns detailed status of a specific reassignment

path Parameters
id
required
string

ID of the reassignment

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Cancel reassignment

Cancels a running partition reassignment

path Parameters
id
required
string

ID of the reassignment to cancel

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

License

License status and management

Get license status

Returns current license state and limits

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get third-party license notices

Returns the full text of third-party software license notices for all open-source dependencies bundled into the Pilot binary.

Responses

Response samples

Content type
application/json
{
  • "success": false,
  • "error": "string",
  • "details": "string"
}

Consumer Groups

Consumer group monitoring, lag tracking, and management

List consumer groups

Returns all consumer groups with state, member count, total lag, and velocity. Uses cached data from the background collector when available.

query Parameters
state
string

Filter by group state (e.g. Stable, Rebalancing, Dead, Empty)

search
string

Filter by group name substring (case-insensitive)

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Get consumer group summary

Returns aggregated consumer group statistics including total groups, lag, groups by state, and consumption rates by topic and broker. Requires the background collector to be running.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get consumer group detail

Returns detailed information about a single consumer group including members, assignments, and coordinator.

path Parameters
group
required
string

Consumer group ID

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete a consumer group

Deletes a consumer group from the cluster. The group must be empty (no active members) before it can be deleted. Requires a valid license (management mode).

path Parameters
group
required
string

Consumer group ID

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get consumer group lag

Returns per-partition lag for a consumer group, including current offsets, log end offsets, lag, velocity, and consumer ID assignments. Uses cached data from the background collector when available.

path Parameters
group
required
string

Consumer group ID

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Reset consumer group offsets

Resets the committed offsets for a consumer group. The group must be empty (no active members) before offsets can be reset. Supports multiple strategies: earliest, latest, to a specific timestamp, or to specific offsets per partition. Use dryRun to preview the offset changes without applying them. Requires a valid license (management mode).

path Parameters
group
required
string

Consumer group ID

Request Body schema: application/json
strategy
required
string
Enum: "earliest" "latest" "timestamp" "specific"

The offset reset strategy to use

topics
Array of strings

Optional list of topics to reset offsets for. If omitted, all topics consumed by the group are reset.

timestamp
integer <int64>

Unix timestamp in milliseconds. Required when strategy is "timestamp".

object

Map of topic-partition to target offset. Required when strategy is "specific". Keys are in the format "topic:partition".

dryRun
boolean
Default: false

If true, return the offset changes without actually applying them

Responses

Request samples

Content type
application/json
{
  • "strategy": "earliest",
  • "topics": [
    ],
  • "timestamp": 1708099200000,
  • "offsets": {
    },
  • "dryRun": false
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Auth

Authentication operations

List configured OAuth providers

Responses

Start OAuth login flow

path Parameters
provider
required
string

Responses

OAuth callback endpoint

path Parameters
provider
required
string

Responses

Get current authenticated user

Responses

Logout current session

Responses

Refresh access token using refresh token

Responses

Access Tokens

Personal access token management for API and MCP authentication

List personal access tokens

Returns all active (non-revoked) PATs for the authenticated user. Token hashes are never exposed.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Create a personal access token

Creates a new PAT for the authenticated user. The plaintext token is returned once in the response and cannot be retrieved again. Requires an OAuth session (PATs cannot create other PATs).

Request Body schema: application/json
name
required
string

Display name for the token

scope
string
Default: "read"
Enum: "read" "write"

Token scope. read allows GET operations only. write allows all operations.

expiresInDays
integer

Number of days until expiry. 0 or omitted means no expiry.

Responses

Request samples

Content type
application/json
{
  • "name": "My MCP token",
  • "scope": "read",
  • "expiresInDays": 90
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Revoke a personal access token

Revokes a PAT by ID. Only the owning user can revoke their own tokens. MCP clients using this token will lose access immediately. Requires an OAuth session.

path Parameters
tokenId
required
string

The token ID to revoke

Responses

MCP

Model Context Protocol server for AI agent integration

Model Context Protocol endpoint

MCP (Model Context Protocol) server endpoint for AI agent integration. Accepts MCP JSON-RPC messages and returns tool results. Supports Streamable HTTP transport.

When authentication is enabled, requires a valid Authorization: Bearer <pat_token> header.

The server exposes 67+ tools for Kafka cluster management across three categories:

  • Read tools (30+): Cluster overview, topic configs, consumer groups, ACLs, quotas, proposals, message browsing
  • Mutate tools (20+): Reassignments, ACL/quota management, leader elections, maintenance mode (require approval)
  • Meta tools (3): Approve, reject, and list pending mutation approvals

Client configuration example (VS Code / GitHub Copilot):

{
  "servers": {
    "pilot": {
      "type": "http",
      "url": "https://pilot.example.com/api/v1/mcp",
      "headers": {
        "Authorization": "Bearer pat_your_token_here"
      }
    }
  }
}
Request Body schema: application/json
object

MCP JSON-RPC message

Responses

Request samples

Content type
application/json
{ }

Audit

Audit log viewing and filtering

List audit log events

Returns audit log events from the in-memory store with optional server-side filtering. Requires authentication. Only available when audit logging is enabled.

query Parameters
action
string

Filter by exact action type (e.g., topics.config.update, broker.maintenance.set)

userId
string

Filter by user ID (case-insensitive substring match)

after
string <date-time>

Only include events after this RFC3339 timestamp (inclusive)

before
string <date-time>

Only include events before this RFC3339 timestamp (inclusive)

search
string

Free-text search across userId, path, action, and topic fields

limit
integer
Default: 500

Maximum number of events to return

offset
integer
Default: 0

Pagination offset

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Revert an audited change

Reverses a previously audited mutation by applying the inverse operation. Accepts the audit event data and restores the state captured in the before field. The revert itself is audit-logged with the action suffixed by .revert.

Request Body schema: application/json
action
required
string

The audit action to revert. Supported actions: topics.config.update, acls.create, acls.update, acls.delete, quotas.update, quotas.delete, reassignments.apply, consumer-groups.reset-offsets

before
object or null

State before the original change (used to restore)

after
object or null

State after the original change (used to identify what to undo)

topic
string

Topic name (required for topic config and reassignment reverts)

path
string

Original request path (required for quota reverts to identify the entity)

Responses

Request samples

Content type
application/json
{
  • "action": "string",
  • "before": { },
  • "after": { },
  • "topic": "string",
  • "path": "string"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Message Browser

Browse and stream messages from Kafka topics

Browse topic messages

Fetches messages from a Kafka topic with pagination support. Supports seeking by offset, timestamp, or latest/earliest. Returns messages with auto-detected encoding (JSON, text, or binary/base64). Values are truncated by default - use the single message endpoint for full values. Use keySerde/valueSerde to force a specific deserialization mode.

path Parameters
topic
required
string

Topic name

query Parameters
partition
integer
Default: -1

Partition to read from (-1 for all partitions)

offset
integer <int64>

Start reading from this offset (mutually exclusive with timestamp and seek)

timestamp
integer <int64>

Start reading from messages at/after this Unix millisecond timestamp

seek
string
Default: "latest"
Enum: "earliest" "latest"

Seek to earliest or latest offset (used when neither offset nor timestamp is provided)

limit
integer <= 500
Default: 50

Maximum number of messages to return

direction
string
Default: "backward"
Enum: "forward" "backward"

Read direction - backward returns newest first, forward returns oldest first

keyFilter
string

Substring filter on message key

valueFilter
string

Substring filter on decoded message value

maxValueSize
integer
Default: 10240

Max bytes of value to return per message (truncated beyond this)

keySerde
string
Default: "auto"
Enum: "auto" "string" "json" "hex" "base64" "int32" "int64" "float64"

Deserialization mode for message keys

valueSerde
string
Default: "auto"
Enum: "auto" "string" "json" "hex" "base64" "int32" "int64" "float64"

Deserialization mode for message values

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get a single message

Retrieves a single message at an exact partition and offset. Returns the full untruncated value (up to 1MB). Use this to lazy-load truncated values from the browse endpoint.

path Parameters
topic
required
string

Topic name

partition
required
integer <int32>

Partition number

offset
required
integer <int64>

Message offset

query Parameters
maxValueSize
integer
Default: 1048576

Max bytes of value to return (default 1MB)

keySerde
string
Default: "auto"
Enum: "auto" "string" "json" "hex" "base64" "int32" "int64" "float64"

Deserialization mode for message key

valueSerde
string
Default: "auto"
Enum: "auto" "string" "json" "hex" "base64" "int32" "int64" "float64"

Deserialization mode for message value

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Get topic watermarks

Returns the low and high watermark offsets for each partition of a topic. Useful for the UI to show seek controls, progress bars, and partition selection.

path Parameters
topic
required
string

Topic name

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Live tail topic messages (SSE)

Streams new messages from a topic in real-time using Server-Sent Events (SSE). The connection stays open and sends message events as new messages arrive. A keepalive event is sent every 15 seconds to prevent timeouts. Close the connection to stop tailing.

path Parameters
topic
required
string

Topic name

query Parameters
partition
integer
Default: -1

Partition to tail (-1 for all)

maxValueSize
integer
Default: 10240

Max bytes of value to return per message

keySerde
string
Default: "auto"
Enum: "auto" "string" "json" "hex" "base64" "int32" "int64" "float64"

Deserialization mode for message keys

valueSerde
string
Default: "auto"
Enum: "auto" "string" "json" "hex" "base64" "int32" "int64" "float64"

Deserialization mode for message values

Responses

Simulator

What-if simulation and blast radius analysis

Run what-if simulation

Simulates hypothetical cluster mutations (broker failures, rack failures, traffic spikes, broker additions) and returns the projected cluster state compared to the current state. This is a read-only operation that never affects the live cluster.

Request Body schema: application/json
brokerFailures
Array of integers

List of broker IDs to mark as offline

rackFailures
Array of strings

List of rack names to mark as offline (all brokers in rack)

Array of objects

Hypothetical new brokers to add

object

Topic name to new replication factor

object

Topic name to traffic multiplier (e.g. 2.0 = double traffic)

Responses

Request samples

Content type
application/json
{
  • "brokerFailures": [
    ],
  • "rackFailures": [
    ],
  • "addBrokers": [
    ],
  • "topicRFChanges": {
    },
  • "trafficMultipliers": {
    }
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Analyze blast radius of an entity

Computes the full dependency graph and failure impact for a broker, rack, or topic. Shows affected topics, consumer groups, and a simulated failure scenario. This is a read-only operation.

Request Body schema: application/json
entityType
string
Enum: "broker" "rack" "topic"

Type of entity to analyze (single-entity mode)

entityId
string

Entity identifier - broker ID, rack name, or topic name (single-entity mode)

brokerFailures
Array of integers <int32> [ items <int32 > ]

Broker IDs to simulate as failed (multi-failure mode)

rackFailures
Array of strings

Rack names to simulate as failed (multi-failure mode)

Responses

Request samples

Content type
application/json
{
  • "entityType": "broker",
  • "entityId": "string",
  • "brokerFailures": [
    ],
  • "rackFailures": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Chat

AI chat assistant for natural-language Kafka cluster management. Mutations require explicit user approval.

Get chat feature status

Returns whether the AI chat assistant feature is enabled and, if so, which model is configured.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Create a new chat conversation

Creates a new conversation for the authenticated user. Each user can have up to PILOT_CHAT_MAX_CONVERSATIONS active conversations (default 10). Oldest conversations are evicted when the limit is reached.

Request Body schema: application/json
title
string

Optional conversation title

Responses

Request samples

Content type
application/json
{
  • "title": "Topic configuration review"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

List user's chat conversations

Returns all active conversations for the authenticated user, ordered by most recently updated.

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": [
    ]
}

Get a conversation with messages

Returns a single conversation including its full message history.

path Parameters
id
required
string

Conversation ID

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete a conversation

Permanently deletes a conversation and all its messages.

path Parameters
id
required
string

Conversation ID

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Send a message and stream the response (SSE)

Sends a user message to the AI assistant and streams the response back using Server-Sent Events. The assistant can use tools to query and manage the Kafka cluster. Read-only tools execute automatically; mutation tools require explicit approval via the approve/reject endpoints.

SSE Event Types:

  • text_delta - Incremental text content from the assistant
  • tool_use - A tool is being called (includes tool name and arguments)
  • tool_result - Tool execution completed (includes result and status)
  • approval_required - A mutation tool needs user approval before execution
  • done - Response stream complete
  • error - An error occurred during processing

The connection stays open until the assistant finishes responding or the client disconnects. Rate limited to PILOT_CHAT_RATE_LIMIT messages per user per minute (default 20).

path Parameters
id
required
string

Conversation ID

Request Body schema: application/json
content
required
string <= 10000 characters

The user's message text (max 10,000 characters)

Responses

Request samples

Content type
application/json
{
  • "content": "Which topics have min.insync.replicas below 2?"
}

Approve a pending mutation tool call

Approves a pending mutation operation that was requested by the AI assistant during a chat conversation. The tool will be executed and the result sent back to the assistant to continue the conversation. Approval tokens expire after PILOT_CHAT_APPROVAL_TIMEOUT (default 5 minutes).

path Parameters
id
required
string

Conversation ID

toolCallId
required
string

Tool call ID from the approval_required SSE event

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Reject a pending mutation tool call

Rejects a pending mutation operation. The assistant will be informed of the rejection (with optional reason) and can suggest alternatives.

path Parameters
id
required
string

Conversation ID

toolCallId
required
string

Tool call ID from the approval_required SSE event

Request Body schema: application/json
reason
string

Optional reason for rejection

Responses

Request samples

Content type
application/json
{
  • "reason": "Not during peak hours"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "message": "string"
}

Topics

Move a topic's replicas on a broker to a target log directory

Moves all partition replicas of the topic that currently live on the specified broker to the target log directory. Replicas already on the target are skipped.

path Parameters
topic
required
string

Topic name

Request Body schema: application/json
brokerId
required
integer <int32>

Broker hosting the replicas to move

targetLogDir
required
string

Destination log directory on the broker

Responses

Request samples

Content type
application/json
{
  • "brokerId": 0,
  • "targetLogDir": "/var/lib/kafka/data-2"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "brokerId": 0,
  • "totalPartitions": 0,
  • "movedPartitions": 0,
  • "failedPartitions": 0,
  • "errors": [
    ]
}

Move a single partition replica on a broker to a target log directory

Moves the replica of (topic, partition) on the specified broker to the target log directory.

path Parameters
topic
required
string

Topic name

partition
required
integer <int32>

Partition index

Request Body schema: application/json
brokerId
required
integer <int32>

Broker hosting the replica to move

targetLogDir
required
string

Destination log directory on the broker

Responses

Request samples

Content type
application/json
{
  • "brokerId": 0,
  • "targetLogDir": "/var/lib/kafka/data-2"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "brokerId": 0,
  • "totalPartitions": 0,
  • "movedPartitions": 0,
  • "failedPartitions": 0,
  • "errors": [
    ]
}

ACLs

List ACL bindings

Returns all ACL bindings in the cluster. Supports optional query parameters to filter results by resource type, resource name, pattern type, principal, host, operation, or permission type.

query Parameters
resourceType
string
Enum: "Topic" "Group" "Cluster" "TransactionalID" "DelegationToken"

Filter by resource type

resourceName
string

Filter by resource name

patternType
string
Enum: "Literal" "Prefixed" "Any" "Match"

Filter by pattern type

principal
string

Filter by principal (e.g. "User:alice")

host
string

Filter by host

operation
string
Enum: "All" "Read" "Write" "Create" "Delete" "Alter" "Describe" "ClusterAction" "DescribeConfigs" "AlterConfigs" "IdempotentWrite" "Any"

Filter by operation

permissionType
string
Enum: "Allow" "Deny" "Any"

Filter by permission type

Responses

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Create ACL bindings

Creates one or more ACL bindings. Requires a valid license. For Cluster resource types, the resourceName defaults to "kafka-cluster" if omitted.

Request Body schema: application/json
required
Array of objects (AclBinding) non-empty

List of ACL bindings to create

Responses

Request samples

Content type
application/json
{
  • "bindings": [
    ]
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Update an ACL binding

Replaces an existing ACL binding with a new one. The old binding is deleted and the new binding is created atomically. If creation of the new binding fails, the old binding is restored. Requires a valid license.

Request Body schema: application/json
required
object (AclBinding)
required
object (AclBinding)

Responses

Request samples

Content type
application/json
{
  • "old": {
    },
  • "new": {
    }
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}

Delete ACL bindings

Deletes ACL bindings matching the provided filter criteria. At least one filter field is required to prevent accidental deletion of all ACLs. Requires a valid license.

Request Body schema: application/json
resourceType
string (AclResourceType)
Enum: "Topic" "Group" "Cluster" "TransactionalID" "DelegationToken"

Kafka resource type for the ACL binding

resourceName
string

Filter by resource name

patternType
string (AclPatternType)
Enum: "Literal" "Prefixed" "Any" "Match"

Resource name pattern type

principal
string

Filter by principal

host
string

Filter by host

operation
string (AclOperation)
Enum: "All" "Read" "Write" "Create" "Delete" "Alter" "Describe" "ClusterAction" "DescribeConfigs" "AlterConfigs" "IdempotentWrite" "Any"

Kafka ACL operation

permission
string (AclPermissionType)
Enum: "Allow" "Deny" "Any"

ACL permission type

Responses

Request samples

Content type
application/json
{
  • "resourceType": "Topic",
  • "resourceName": "string",
  • "patternType": "Literal",
  • "principal": "string",
  • "host": "string",
  • "operation": "All",
  • "permission": "Allow"
}

Response samples

Content type
application/json
{
  • "success": true,
  • "data": {
    }
}