Docs SDK Protocol Clients (gRPC, Kafka, RabbitMQ, SOAP, GraphQL, SSE, WS)

SDK Protocol Clients

The Mockarty SDKs ship test clients for every protocol the platform supports — not just for configuring a mock, but for driving the system under test from CI scripts and emitting the per-call timeline back into a TCM run.

Same surface in Go, Python, Java:

Protocol Go package Python module Java package
gRPC protocols/grpc (not yet) ru.mockarty.protocols.grpc
Kafka protocols/kafka (not yet) ru.mockarty.protocols.kafka
RabbitMQ protocols/rabbitmq (not yet) ru.mockarty.protocols.rabbitmq
SOAP (not yet) mockarty.protocols.soap ru.mockarty.protocols.soap
GraphQL (not yet) mockarty.protocols.graphql ru.mockarty.protocols.graphql
SSE (not yet) mockarty.protocols.sse ru.mockarty.protocols.sse
WebSocket (not yet) mockarty.protocols.websocket ru.mockarty.protocols.websocket

About URLs in examples: all examples use localhost:5770 as the default Mockarty address. If your instance runs on a remote server, replace localhost:5770 with its actual address. See Tips & Useful Features for details.

Related pages: SDK Guide · SDK Quick Start · Test Plans in CI/CD

Why a protocol client and not just HTTP?

A typical CI test that exercises a gRPC service ends up writing 30 lines of boilerplate: generated stubs, JSON-to-protobuf, status-code mapping, time measurement, and a separate report uploader. The protocol clients in the SDK collapse all of that into:

err := conn.InvokeJSON(ctx, "acme.UserService/GetUser",
    map[string]any{"id": "u-42"}, &resp)

The same line records a Mockarty TCM step (start/end/duration/status/payload preview) and ships it through externalruns so the run report shows a per-RPC timeline once the test finishes. No generated stubs, no protoc step in the test repo, no custom reporter.

The Step Recorder pattern

Every protocol client takes an optional StepRecorder. The default is NopRecorder — it drops every step on the floor, so scripts that don’t care about the TCM timeline still work out of the box.

Wire a real recorder when you want the timeline:

Go

import (
    "github.com/mockarty/mockarty-go/externalruns"
    "github.com/mockarty/mockarty-go/protocols/grpc"
    "github.com/mockarty/mockarty-go/protocols/telemetry"
)

runs, _ := externalruns.NewClient("http://localhost:5770", "sandbox", apiToken)
run, _ := runs.CreateRun(ctx, externalruns.CreateRunRequest{
    Name: "grpc smoke", Framework: "go-test",
})
defer runs.FinishRun(ctx, run.ID, externalruns.FinishRunRequest{})

recorder := telemetry.NewExternalRunsRecorder(runs, run.ID)
defer recorder.Close()

conn, _ := grpc.Dial(ctx, "localhost:50051", grpc.WithRecorder(recorder))
defer conn.Close()
// ... your invokes ...

Python

from mockarty import MockartyClient
from mockarty.protocols import AccumulatingRecorder
from mockarty.protocols.graphql import GraphQLClient

client = MockartyClient(base_url="http://localhost:5770", api_key="...")
recorder = AccumulatingRecorder()

gql = GraphQLClient("http://app/graphql", recorder=recorder)
gql.execute("query GetUser { user { id } }")

# At the end of the test, push the captured timeline to TCM:
client.external_runs().report(
    case_name="my case",
    status="passed",
    steps=recorder.steps(),
)

Java

import ru.mockarty.protocols.telemetry.AccumulatingRecorder;
import ru.mockarty.protocols.graphql.GraphQLClient;

AccumulatingRecorder rec = new AccumulatingRecorder();
GraphQLClient gql = new GraphQLClient("http://app/graphql", opts -> opts.recorder(rec));
gql.execute("query GetUser { user { id } }");
// At test finish:
client.externalRuns().report("qa",
    new ExternalRunRequest().caseName("my case").status(ExternalRunRequest.STATUS_PASSED));

Cross-language invariants

The three SDKs use identical step shapes server-side so a mixed-language test suite collapses correctly on (namespace, run, step_key):

  • Status taxonomy: "passed" / "failed" / "broken" / "skipped". Empty defaults to "passed".
  • Step key: <name>#<seq> where <seq> is a per-client monotonic counter. Retries collapse on the same key.
  • Step name:
    • gRPC: the bare fully-qualified method (acme.UserService/GetUser).
    • Kafka: produce:<topic> / consume:<topic>.
    • RabbitMQ: publish:<exchange>/<routingKey> / consume:<queue> / declare-queue:<name>.
    • SOAP / GraphQL / SSE / WebSocket: protocol-prefixed (soap:GetUser, graphql:GetUser, sse:collect, ws:send).
  • Truncation marker: …(truncated <N>B) (U+2026 ellipsis). Default cap 1024 bytes; 0 disables; negative clamps to 0.
  • Error classification:
    • Context cancel / deadline → "broken" (env, not a failed assertion).
    • Typed protocol error (gRPC status, kafka.Error, amqp.Error, SOAP fault, GraphQL errors[], HTTP 4xx/5xx) → "failed".
    • Untyped transport (network, TLS, DNS) → "broken".

Test ergonomics (wrap / eventually / parallel)

Three cross-cutting helpers make multi-step suites readable and resilient. They
work with every protocol facet and are available identically in all three SDKs.

  • wrap(name, body) — group the chains fired inside body under one named
    parent step so the report renders as a tree, not a flat list.
  • eventually(within, interval, attempt) — retry attempt until it passes
    or the budget elapses, sleeping interval between tries. Only the successful
    (or, on timeout, the final) attempt’s steps stay in the report — intermediate
    failures are rolled back. Use it for eventual consistency (async side effects,
    read-after-write lag). Returns whether it passed.
  • parallel(branch...) — run each branch concurrently on its own
    branch-local tester (shared transport / base URL / headers + a snapshot of the
    vars) and merge the steps back in branch order. Variable writes inside a branch
    stay isolated.

The success signal is idiomatic per language — Go returns an error (nil = pass),
Python returns None or an Exception, Java returns a boolean — but the
semantics match exactly.

Go

t.Wrap("checkout flow", func() {
    t.HTTP().POST("/cart").JSON(map[string]any{"sku": "abc"}).ExpectStatus(201)
    t.HTTP().POST("/checkout").ExpectStatus(200)
})

ok := t.Eventually(5*time.Second, 200*time.Millisecond, func() error {
    t.HTTP().GET("/orders/42").ExpectStatus(200)
    if !t.OK() {
        return errors.New("order not ready")
    }
    return nil
})

t.Parallel(
    func(t *tester.Tester) { t.HTTP().GET("/a").ExpectStatus(200) },
    func(t *tester.Tester) { t.HTTP().GET("/b").ExpectStatus(200) },
)

Python

from mockarty.tester import wrap, eventually, parallel

def checkout():
    t.http().post("/cart").json({"sku": "abc"}).expect_status(201)
    t.http().post("/checkout").expect_status(200)

wrap(t, "checkout flow", checkout)

def order_ready():
    t.http().get("/orders/42").expect_status(200)
    return None if t.ok() else RuntimeError("order not ready")

ok = eventually(t, within=5.0, interval=0.2, fn=order_ready)

parallel(t,
    lambda b: b.http().get("/a").expect_status(200),
    lambda b: b.http().get("/b").expect_status(200))

Java

t.wrap("checkout flow", () -> {
    t.http().post("/cart").json(Map.of("sku", "abc")).expectStatus(201);
    t.http().post("/checkout").expectStatus(200);
});

boolean ok = t.eventually(Duration.ofSeconds(5), Duration.ofMillis(200), () -> {
    t.http().get("/orders/42").expectStatus(200);
    return t.ok();
});

t.parallel(
    b -> b.http().get("/a").expectStatus(200),
    b -> b.http().get("/b").expectStatus(200));

gRPC

Reflection-driven JSON-shaped client. Air-gapped installs that don’t expose grpc.reflection can pre-compile descriptor sets (Java) or pass .proto files (Go).

Go

conn, _ := mgrpc.Dial(ctx, "localhost:50051",
    mgrpc.WithRecorder(recorder),
    mgrpc.WithProtoFile("acme/user.proto"),   // optional
    mgrpc.WithImportDir("./protos"),          // optional
    mgrpc.WithReflection(true),               // default
)
defer conn.Close()

var resp map[string]any
_ = conn.InvokeJSON(ctx, "acme.UserService/GetUser",
    map[string]any{"id": "u-42"}, &resp)

services, _ := conn.ListServices(ctx)         // discovery
methods,  _ := conn.ListMethods(ctx, "acme.UserService")

Java

byte[] descSet = Files.readAllBytes(Path.of("build/user.desc"));   // protoc --descriptor_set_out

GrpcClient client = new GrpcClient("localhost:50051", opts -> opts
    .recorder(rec)
    .protoDescriptorSet(descSet)
    .reflection(true));

Map<String, Object> resp = client.invokeJson(
    "acme.UserService/GetUser",
    Map.of("id", "u-42"),
    Map.class);

Streaming: v1 supports unary gRPC calls only (InvokeJSON). Server-streaming, client-streaming and bidi are not exposed in v1 — they double the API surface and the CI test cases that need them are rare.

Kafka

Pure-Go (segmentio/kafka-go) and pure-Java (kafka-clients) — no librdkafka, no CGO.

Go

cli, _ := kafka.NewClient([]string{"kafka:9092"},
    kafka.WithRecorder(recorder),
    kafka.WithAutoTopicCreation(true))
defer cli.Close()

_ = cli.Produce(ctx, "orders", "order-42",
    map[string]any{"id": 42, "amount": 19.99},
    map[string]string{"x-tenant": "demo"})

msgs, _ := cli.Consume(ctx, kafka.ConsumeOptions{
    Topic: "orders", GroupID: "smoke", MaxMessages: 1,
})

Java

KafkaClient cli = new KafkaClient(List.of("kafka:9092"), opts -> opts
    .recorder(rec)
    .requiredAcks("all"));

cli.produce("orders", "order-42",
    Map.of("id", 42, "amount", 19.99),
    Map.of("x-tenant", "demo"));

List<ConsumedMessage> got = cli.consume(new ConsumeOptions("orders", "smoke", 1));

One step per consume call (not per message); count parameter records messages actually fetched. Decode failures downgrade the verdict from "passed" to "failed" before emission.

RabbitMQ

Pure-Go (rabbitmq/amqp091-go) and pure-Java (com.rabbitmq:amqp-client).

Go

cli, _ := rabbitmq.NewClient("amqp://guest:guest@rabbitmq:5672/",
    rabbitmq.WithRecorder(recorder))
defer cli.Close()

_ = cli.DeclareQueue(ctx, "orders", rabbitmq.DeclareQueueOptions{Durable: true})
_ = cli.Publish(ctx, "", "orders",
    map[string]any{"id": 42, "status": "ready"}, nil)

msgs, _ := cli.Consume(ctx, rabbitmq.ConsumeOptions{
    Queue: "orders", MaxMessages: 1,
})

Java

RabbitMqClient cli = new RabbitMqClient("amqp://guest:guest@rabbitmq:5672/",
    opts -> opts.recorder(rec));

cli.declareQueue("orders", new DeclareQueueOptions(true, false, false, Map.of()));
cli.publish("", "orders", Map.of("id", 42, "status", "ready"), null);

List<ConsumedMessage> got = cli.consume(new ConsumeOptions("orders", 1, false));

ContentType defaults to application/json so consumers can rely on it without per-call boilerplate.

SOAP

JDK DOM (Java) and stdlib xml.etree (Python). No lxml dependency. XXE protection enabled.

Python

from mockarty.protocols.soap import SoapClient

with SoapClient("http://app/svc", soap_action="GetUser") as cli:
    resp = cli.call("GetUser", "<GetUser><id>u-1</id></GetUser>")
    assert resp.status_code == 200
    assert resp.fault is None
    root = resp.root()  # ElementTree.Element

Java

try (SoapClient cli = new SoapClient("http://app/svc", opts -> opts
        .soapAction("GetUser")
        .recorder(rec))) {
    SoapResponse resp = cli.call("GetUser", "<GetUser><id>u-1</id></GetUser>");
    assertEquals(200, resp.getStatusCode());
    assertNull(resp.getFault());
}

Step parameters carry operation, http_status, request, response, plus fault_<key> entries when the response carries a SOAP Fault. SOAP 1.1 (text/xml) and SOAP 1.2 (application/soap+xml) are both supported via the version option.

GraphQL

Standard POST envelope {"query": "...", "variables": {...}, "operationName": "..."}. Treats non-empty errors[] as "failed" even on HTTP 200.

Python

from mockarty.protocols.graphql import GraphQLClient, GraphQLError

with GraphQLClient("http://app/graphql", recorder=rec) as gql:
    resp = gql.execute(
        "query GetUser($id: ID!) { user(id: $id) { name } }",
        variables={"id": "u-1"},
    )
    if not resp.ok:
        print(resp.errors)
    # Or raise on errors[]:
    gql.execute("query Bad { x }", raise_for_errors=True)  # raises GraphQLError

Java

try (GraphQLClient gql = new GraphQLClient("http://app/graphql", opts -> opts.recorder(rec))) {
    GraphQLResponse resp = gql.execute(
        "query GetUser($id: ID!) { user(id: $id) { name } }",
        Map.of("id", "u-1"));
    if (!resp.isOk()) {
        log.warn("errors: {}", resp.getErrors());
    }
}

Operation name is auto-extracted from the query body so the step label looks like graphql:GetUser rather than graphql:anonymous.

SSE

WHATWG-spec parser handles data: / event: / id: / retry: / comments. Two read modes: collect (cap + deadline) and stream (lazy generator).

Python

from mockarty.protocols.sse import SseClient

with SseClient("http://app/events", recorder=rec) as sse:
    events = sse.collect(max_events=5, max_duration=10.0)
    for ev in events:
        print(ev.event, ev.data)

    # OR lazy:
    for ev in sse.stream():
        if ev.data == "done":
            break

Java

try (SseClient sse = new SseClient("http://app/events", opts -> opts.recorder(rec))) {
    List<SseEvent> events = sse.collect(5, Duration.ofSeconds(10));
    events.forEach(ev -> System.out.println(ev.getEvent() + " → " + ev.getData()));
}

The Python collect records "broken" when the deadline expires with zero events received — that’s an env-level failure (the stream is silent). Cap-hit and deadline + ≥1 events stay "passed". The reason parameter (cap | deadline) on the step shows why collect returned.

WebSocket

JDK 11+ java.net.http.WebSocket (Java) and websockets package (Python, optional via mockarty[protocols] extra).

Python

pip install "mockarty[protocols]"   # pulls websockets
from mockarty.protocols.websocket import WebSocketClient

with WebSocketClient("ws://app/socket", recorder=rec) as ws:
    ws.send({"type": "subscribe", "channel": "orders"})
    msg = ws.recv(timeout=5)
    print(msg)

Java

try (WebSocketClient ws = new WebSocketClient("ws://app/socket", opts -> opts.recorder(rec))) {
    ws.send("{\"type\":\"subscribe\",\"channel\":\"orders\"}"); // send() takes a String
    String msg = ws.recv(Duration.ofSeconds(5));
    System.out.println(msg);
}

Steps land as ws:send and ws:recv. Long-running subscriptions are better served by SseClient where applicable.

Air-gapped notes

Every protocol client is pure-language (no native deps that need a separate install). Specifically:

  • Go: segmentio/kafka-go, rabbitmq/amqp091-go, jhump/protoreflect are all pure-Go.
  • Python: SOAP / GraphQL / SSE ride on httpx (already a core dep); WebSocket pulls websockets via the protocols extra.
  • Java: gRPC uses grpc-netty-shaded which embeds Netty; Kafka / RabbitMQ have no native deps; SOAP / GraphQL / SSE / WebSocket use the JDK java.net.http and javax.xml.parsers directly.

The same binaries run in distroless containers against Mockarty-fronted mocks without any external tooling.

Troubleshooting

  • “streaming method not supported in v1” (gRPC) — the requested method is client-streaming or bidi. Use the language-specific stub — the SDK does not carry streaming methods.
  • Empty steps in the TCM timeline — no recorder was wired. The default NopRecorder drops everything; pass WithRecorder(...) (Go/Java) or recorder=... (Python) to enable capture.
  • Steps captured but absent in the TCM run — the recorder is buffered. In Python, the test must call client.external_runs().report(steps=recorder.steps()) at the end of the suite. In Go / Java, the ExternalRunsRecorder flushes asynchronously — call recorder.Close() to drain before the process exits.
  • Replacement char (�) in step payload previews — your SDK version predates the UTF-8 boundary fix (Go: pre-2026-05-19; Python: same; Java: same). Upgrade.
  • recorder=AccumulatingRecorder() recorded zero steps in Python — your SDK version predates the falsy-recorder fix (pre-2026-05-19; the empty recorder was bool(rec) == False, so recorder or NopRecorder() discarded it).