KubeMQ
Operate

kmq CLI

Drive KubeMQ from the terminal with the kmq command-line client — messaging, observability, contexts, roles, and the installable agent skill.

kmq is KubeMQ's native command-line client for macOS, Linux, and Windows. It manages onboarding, protected credentials, saved deployment plans, messaging, and operations. Messaging and server inspection use the authenticated management API; deployment commands invoke your explicitly selected Docker, Podman, Helm, and Kubernetes tools.

Overview

kmq exposes every messaging pattern — Queues, Events, Events Store, and RPC (Commands/Queries) — plus full observability (status, metrics, connections, audit, connectors, agents) and a set of meta/discovery commands. Four properties make it suited to automated and agent-driven use:

  • Token-efficient — a kmq queue send is a few hundred tokens versus a multi-round exchange over a protocol like MCP.
  • Deterministic — typed exit codes and machine-readable json/ndjson output make results easy to branch on in a script or an agent loop.
  • Self-describing — kmq schema emits the whole command tree, offline, with no network call.
  • Bounded by default — every stream/subscribe/replay command honors --count/--duration/--idle, so an agent can never loop forever waiting on it.

Scope

kmq does not embed the broker. To run a server, use Install with kmq or Install with Docker. A successful install command is followed by authentication, intended-license verification, and a real messaging check. kmq mcp inspects the server's registered tools; it does not make kmq itself a protocol server.

Architecture

Messaging and server-operation commands use the management API: one-shot actions dispatch via POST /api/request ({type, data}) or dedicated REST routes, and streaming commands use Server-Sent Events / WebSocket subscription endpoints. Meta commands (schema, cheat, skills) serve content compiled into the binary, so they work fully offline and never need the server.

Installation

Install kmq on macOS, Linux or Windows as shown in Try KubeMQ. kmq version --check tells you whether a newer release exists, and kmq update installs the latest release.

Contexts & configuration

A context is a named connection profile (API address, token, TLS, defaults). Contexts live under $XDG_CONFIG_HOME/kmq/ (~/.config/kmq/ by default):

$XDG_CONFIG_HOME/kmq/
├── contexts/
│   ├── default.json
│   └── prod.json
└── current-context      # pointer file (active context name)
kmq auth login --username admin --api-address https://kubemq.example.com:8080 --ca-file trusted-ca.pem --context-name prod
kmq context use prod
kmq context list
kmq context current
kmq context delete staging

The login command uses the supplied customer username and prompts for the password, creates a dedicated service credential, saves it through protected credential storage, and binds its reference to the endpoint. Use kmq auth setup when initial setup is required. Use kmq auth import --token-stdin to receive an existing service credential through protected input. Do not pass passwords or service tokens in command arguments.

For an installation record, add --installation INSTALLATION_ID. For Kubernetes bind one context per verified server endpoint. Public certificate files go in --ca-file; never disable verification for a non-loopback management endpoint. Local single-server loopback can use http://127.0.0.1:8080.

Native credential storage is the default. Where unavailable, explicitly select --credential-backend file and a private --credential-dir; protect that directory and include it in your credential backup policy. Ordinary context files store credential references, not plaintext service keys.

Configuration is resolved by precedence, highest to lowest:

PrecedenceSource
1 (highest)Persistent flags (--context, --api-address)
2Environment variables
3Active-context file
4 (lowest)Built-in defaults (http://127.0.0.1:8080)
Env varPurpose
KMQ_TOKENService-account Bearer key (kmq_<keyid>_<secret>) — preferred in CI
KMQ_CONTEXTActive context name (overrides the pointer file)
KMQ_API_ADDRESSManagement API URL

Authentication & roles

The onboarding recipes enable server authentication. See Security for the account model and roles. Check whether it is on with kmq doctor -o json | jq .auth (on/off). When auth is on, supply a service-account key through protected credential storage or a securely injected KMQ_TOKEN in automation. Avoid shell history and process arguments. Service-account roles gate what the CLI can do:

RoleGrants
read_onlylist/inspect, metrics, status, overview, schema, doctor, license --details
read_writethe above + send/receive/stream/subscribe/purge, channel create/delete
adminthe above + audit, account management

Auth-exempt commands

A handful of commands and routes succeed with no token, even when server auth is on: kmq doctor, kmq metrics scrape (routes /ready, /health, /metrics, /api/v1/auth/status), and kmq license, which shows the license state of a Docker server or a Kubernetes cluster without signing in. Don't read the role table above as universal gating — these are the exceptions. kmq license --details returns the full license record and needs a signed-in context.

config set/revert and account management require the admin role. Service accounts never carry admin, so those operations are deliberate non-goals of the CLI.

Global flags & output discipline

Persistent flags are inherited by every subcommand:

FlagDefaultPurpose
-o, --outputjsonOutput format: json | ndjson | yaml | table
--context—Use a specific context (overrides current-context)
--api-address—Target :8080 endpoint (overrides context)
--no-colorfalseDisable color in table output
--verbosefalseRequest timing to stderr (token redacted)
--yesfalseConfirm destructive operations without prompting
--dry-runfalseRender the action without executing
--fields—Project output to these camelCase wire fields
--detailsummarysummary | full verbosity, for commands that support it

Output discipline:

  • Data → stdout, warnings/errors/diagnostics → stderr — safe to pipe.
  • One-shot commands default to compact json; streaming commands (queue stream, *subscribe, *replay, conn watch, command/query receive) default to ndjson (one record per line, flushed per record).
  • metrics scrape emits raw Prometheus text; cheat and skills get emit raw markdown — for those, -o is ignored.
kmq queue receive orders --count 10 -o ndjson | jq .body      # stream, per-line
kmq estore replay telemetry --from-first --count 100 -o json # buffer then array
kmq status --fields is_healthy,channels,clients

Exit codes

kmq returns typed exit codes so scripts and agents can branch deterministically. On error it also writes a JSON envelope to stderr: {"error":{"code":"...","message":"...","retryable":bool}}.

CodeNameMeaningRetry?
0OKSuccess—
1GenericUnclassified errorNo
2UsageBad flags / usageNo
3NotFoundResource not foundNo
4AuthAuth required / failed / forbiddenNo
5ConnServer unreachableYes (server down?)
6TimeoutRequest timed outYes
7PartialPartial successCase-by-case
8RetryableServer initializing / too many attemptsYes — kmq auto-retries

Server wire codes map onto this table as follows: auth_required, auth_failed, forbidden, must_change, and tls_required all map to exit code 4; auth_initializing and too_many_attempts map to exit code 8 and are auto-retried; seed_read_only, duplicate, and service_admin_forbidden map to exit code 2; any empty or unknown wire code maps to exit code 1.

Command reference

Every send command reads its body from (in priority) a positional argument → --body-base64 → -f/--file → stdin, and accepts --message-id/--client-id/ --metadata/-M/--tag k:v (repeatable). Every stream/subscribe/replay/receive command accepts the bounding flags --count N, --duration 5m, --idle 30s — whichever fires first stops it; Ctrl-C also exits cleanly with code 0.

Onboarding and deployment

CommandPurpose
trial request --platform docker|kubernetes, verify, claim, resumeRequest a trial key for one Docker server or three Kubernetes servers, prove your email, and save the key as a protected credential
license import, list, show, export, recoverManage protected credentials. license alone shows the license state of the running server (Docker and Kubernetes); see Check license status
license fingerprint --kube-context CONTEXT --namespace NSRead the Kubernetes fingerprint through the server ServiceAccount after independent operator installation
deploy prepareStage pinned artifacts and missing Kubernetes prerequisites
deploy plan --out PRIVATE_PLANInspect the target and save its identities and artifact references
deploy apply --plan PRIVATE_PLANApply or resume the same saved operation
onboardPrepare a plan and apply a simple installation
deploy status, restart, remove --installation IDOperate on the exact saved owned installation
deploy update --installation ID --credential REFReplace a license while preserving data and installation identity
deploy forward --installation ID --server-index NHold a bounded, identity-checked forward to one Kubernetes server
deploy diagnose --installation IDReport bounded reviewed runtime failure categories
deploy connectRecord an existing target for connection without lifecycle authority
auth setup, login, import, useBootstrap or bind protected management credentials
deploy verify --installation IDCheck intended license acceptance and authenticated messaging

Step-by-step use: Install with kmq and Install on Kubernetes. Non-interactive workflows return an actionable state when email proof, browser verification, credentials, or user input is still needed. They must not guess missing secrets or turn a partial operation into success. Reuse saved operation IDs after a lost response; do not create duplicate trials or runtimes.

Messaging

PatternSendReceive / SubscribeNotes
Queuekmq queue send <ch> <body>kmq queue receive <ch> [--count N]persistent, at-least-once
Queue (peek)—kmq queue peek <ch> [--count N]non-consuming
Queue (interactive)—kmq queue stream <ch> [--visibility 60] [--wait 5] [--auto-ack]WS poll/ack/reject session
Queue (drain)—kmq queue purge <ch> --yesdestructive
Eventskmq events send <ch> <body>kmq events subscribe <ch> [--group g]fire-and-forget pub/sub
Events Storekmq estore send <ch> <body>kmq estore subscribe <ch>persistent + replayable
Events Store (replay)—kmq estore replay <ch> <offset-mode>historical replay — one offset mode: --new-only (default), --from-first, --from-last, --from-sequence N, --from-time <RFC3339>, --since-seconds N, plus --group
Command (RPC)kmq command send <ch> <body> --timeout 30kmq command receive <ch> [--respond-body '{}' | --command 'sh']request/ack
Query (RPC)kmq query send <ch> <body> --timeout 30kmq query receive <ch> [--respond-body '{}' | --command 'sh']request/data

Queue send extras: --max-receive-count, --dead-letter <ch>, --expiration-seconds, --delay-seconds.

RPC responders: receive with --respond-body '<json>' echoes a static reply; with --command '<shell>' the inbound body is piped to the command's stdin and its stdout becomes the reply; with neither, requests are printed for manual handling. --respond-body and --command are mutually exclusive.

Channels

kmq channel create <name> --type queues|events|events_store|commands|queries
kmq channel delete <name> --type queues --yes
kmq channel list [--type <t>]
kmq channel inspect <type> <name>          # full detail: clients, rates, totals
# per-family shortcuts also exist: kmq queue list / kmq queue inspect <name>

Cluster

kmq cluster info [--node]      # (alias: snapshot) cluster-merged, or --node for local
kmq cluster health            # /ready — leadership role, ready/healthy
kmq cluster nodes             # topology: node, type, unavailable nodes

Observability

kmq status                    # composite digest: /ready + snapshot, one-line health
kmq overview [--detail full]  # per-pattern rollups (queues / pubsub / request-reply / totals)
kmq metrics scrape            # raw Prometheus /metrics (auth-exempt)
kmq metrics history [--metric message_rate|volume_rate|error_rate|messages|bytes]
kmq conn list                 # active connections (first SSE snapshot)
kmq conn inspect <id>
kmq conn watch                # live SSE stream of connection events
kmq audit query  [--from ..] [--to ..] [--event-type queue.send] [--limit N] ...
kmq audit stats  [--group-by event_type|client_id|category|channel|transport]
kmq connector list            # aws / gcp / amqp / amqp10 / stomp / mqtt / ce / kafka
kmq connector inspect <name> [--operation <service/op>]
kmq agent list [--limit N] [--offset N]     # A2A agents
kmq agent inspect <id>
kmq mcp list                  # server's registered MCP tools
kmq mcp inspect <tool>

Kafka

kmq kafka probe --bootstrap host:port[,host:port]   # dial a source cluster, read-only
kmq kafka probe -b host:port --dial-timeout 30s   # override the 15s default
kmq kafka probe -b host:port \
  --aws-access-key <key> --aws-secret-key <secret> --aws-session-token <token>  # MSK IAM

kmq kafka probe dials a source Apache Kafka / MSK / Confluent cluster over the Kafka wire protocol and reports back the brokers plus their ApiVersions. It is read-only — it never produces, commits, or auto-creates topics — so it's safe to point at a production cluster. Flags: -b/--bootstrap host:port[,host:port...], --dial-timeout (default 15s), and AWS MSK-IAM auth via --aws-access-key/--aws-secret-key/--aws-session-token. Run it as a first connectivity check ahead of the Migration family below, or see Migrate from Kafka for the full workflow.

kmq kafka share-groups list                        # every share group the coordinator has sampled
kmq kafka share-groups describe <group>            # members + per-partition start offset, lag, in-flight
kmq kafka share-groups reset-offsets <group> --topic <t> --to-earliest            # preview only
kmq kafka share-groups reset-offsets <group> --topic <t> --to-offset 42 --execute # write it

kmq kafka share-groups works with Kafka share groups on the KubeMQ server named by the current context. list and describe read the same view as GET /api/kafka/share-groups, at most about 10 seconds old. reset-offsets writes over the Kafka wire (AlterShareGroupOffsets): pass exactly one of --to-earliest, --to-latest or --to-offset N (clamped to the partition's log), one or more --topic, and optionally --partition. Without --execute it only previews each partition's current and target start offset. Where the group has no start offset yet it is created, so a group's first consumer starts from it — the way to set share-group start positions after a migration. The bootstrap address defaults to the context's host on port 9092; pass --bootstrap when that does not apply. Exit codes: 2 for a bad flag or a group id KubeMQ does not accept (containing /, * or >), 3 for an unknown topic or a group id already in use as a classic group, 8 while a consumer is attached to the group — kmq does not retry that one itself; stop the consumer and run it again. If --execute fails on only some partitions, the others were written: retry just the failed ones with --partition.

Migration

CommandFlagsNotes
kmq assess kafka--bootstrap, --tls, --sasl-mechanismRead-only fit assessment of an external Kafka cluster → per-topic READY/CAVEAT/UNKNOWN/BLOCKED + a T1–T4 verdict. Never writes to the source.
kmq migrate assess|replicate|translate|cutover--bootstrap (all four); --target-bootstrap, --state (replicate, translate, cutover); --target-api (required once per target node) and --target-context (replicate, cutover); --dry-run (cutover); --check (assess)Four-phase migration (assess → replicate → translate → cutover) from a Kafka cluster to KubeMQ. Assess reads only the source. Replicate and a real cutover refuse a target unless every node reports strict acknowledgment. cutover --force is deprecated and refused.
kmq migrate plan--input, --output, --source-profile, --target-profile, --target-api, --target-context, --dry-runAuthor the immutable migration plan from a JSON declaration; the output file is never overwritten.
kmq migrate prepare--input, --output, --dry-run, --apply PLAN_FILE, --approve-hash HASHPreview, then apply, bounded topic preparation on the target (replicate needs every target topic created first); --apply takes the preview's plan file and the exact hash it printed.
kmq migrate job submit--plan, --operation, --state-pvc, --state-user-id, --profiles-configmap, --service-account, --cutover-service-account, --credential-secret, --target-apiRun one migration operation as a serialized in-cluster Kubernetes Job, for a target whose Kafka listener is reachable only inside the cluster. The ConfigMap, ServiceAccounts, state volume and any credential Secret must already exist. Cutover runs as separate staged operations.
kmq migrate job status--planRead the observed state of the plan's migration Job.

For the full narrative (per-source auth, MirrorMaker 2 hybrid, rollback, staged dry-run) see Migrate from Kafka.

Meta & discovery

kmq version                   # binary version
kmq version --check           # compare with the latest published release
kmq update                    # update this binary to the latest release
kmq whoami                    # identity + role (auth on), or {"authenticated":false,"auth":"disabled"} (auth off)
kmq doctor                    # connectivity + auth check (no auth required)
kmq config get [--fields ..]  # server config (read-only, server-redacted)
kmq license                   # license state of the server or cluster (no sign-in; --details signs in)
kmq schema -o json            # full command tree + roles + exit codes (offline)
kmq cheat [topic]             # embedded recipes (offline)
kmq docs [topic] [--open]     # signpost to the online docs + LLM corpus
kmq skills ...                # serve/install the agent skill (see next section)

kmq cheat topics (embedded, offline): queue, events, estore, rpc, drain, subscribe, health, auth, context, output.

kmq schema -o json is the machine-readable contract — the full command tree (path, description, required role, flags), the exit-code table, the server wire error-code catalogue, and the CLI/server versions — with no network and no auth. It is the recommended way for an agent to introspect the CLI.

Agent-skill distribution

The discovery skill is optional. An agent with terminal access can run kmq skills get core directly. That guide is embedded in the installed binary, so it matches the kmq version the agent is using. Installing the discovery skill helps agents find the guide when a KubeMQ task begins.

kmq skills                    # list skills (alias: kmq skills list)   → 'core'
kmq skills get core           # print the core skill (raw markdown)
kmq skills get core --full    # + the full command reference
kmq skills get --all          # every skill
kmq skills path [name]        # skills source dir, or '(embedded)'
kmq skills install [--global] [--force]   # Claude Code skill + AGENTS.md pointer (no Node.js)
Install pathCommandScope
Multi-agent installernpx skills add kubemq-io/kmqSelect a supported agent and project or user scope. Requires Node.js.
Claude Code pluginclaude plugin marketplace add kubemq-io/kmqClaude Code
Built-in installerkmq skills installClaude Code skill plus an AGENTS.md pointer in the current directory; --global writes to the user directory.

The multi-agent installer lists supported agents, including Claude Code (claude-code), Codex (codex), Cursor (cursor), GitHub Copilot (github-copilot), Gemini CLI (gemini-cli), and Windsurf (windsurf). To target one explicitly, run npx skills add kubemq-io/kmq --agent codex. Installation places the skill in the agent's expected directory; verify that the agent actually loads it before claiming runtime support for a particular version or configuration.

Install kmq first on macOS or Linux or on Windows x86-64. On Windows, download install.ps1 and run it in PowerShell. After installation, run kmq version and kmq skills get core in a new terminal. The skill installer does not install the kmq executable.

Server communication

kmq speaks to the management API over HTTP. One-shot actions dispatch as POST /api/request with {type, data} (or dedicated REST routes); the server replies HTTP 200 with an envelope:

{ "error": false, "error_string": "", "code": "", "data": {} }

error: false decodes data into the requested output; error: true maps the code field to an exit code via the wire-code table above. A handful of routes are auth-exempt (no Bearer needed) — the same ones called out in the Auth-exempt commands note above.

Retryable server states (auth_initializing, too_many_attempts) are auto-retried with exponential backoff, cancellable by SIGINT/SIGTERM. Streaming uses SSE/WebSocket subscription endpoints and stops on the first of --count/--duration/--idle or a signal. SIGINT/SIGTERM cancel the root context — streaming/waiting commands close gracefully and exit 0 on clean cancellation.

Was this page helpful?

On this page