Документация Протокольные клиенты SDK (gRPC, Kafka, RabbitMQ, SOAP, GraphQL, SSE, WS)

Протокольные клиенты SDK

SDK Mockarty предоставляют тестовые клиенты для каждого протокола, который поддерживает платформа — не только чтобы настроить mock, а чтобы прогнать систему под тестом из CI-скриптов и отправить тайм-лайн вызовов обратно в прогон TCM.

Одинаковая поверхность в Go, Python, Java:

Протокол Go package Python module Java package
gRPC protocols/grpc (пока нет) ru.mockarty.protocols.grpc
Kafka protocols/kafka (пока нет) ru.mockarty.protocols.kafka
RabbitMQ protocols/rabbitmq (пока нет) ru.mockarty.protocols.rabbitmq
SOAP (пока нет) mockarty.protocols.soap ru.mockarty.protocols.soap
GraphQL (пока нет) mockarty.protocols.graphql ru.mockarty.protocols.graphql
SSE (пока нет) mockarty.protocols.sse ru.mockarty.protocols.sse
WebSocket (пока нет) mockarty.protocols.websocket ru.mockarty.protocols.websocket

Об URL в примерах: все примеры используют localhost:5770 как адрес Mockarty по умолчанию. Если ваш инстанс работает на удалённом сервере — замените localhost:5770 на его адрес. Подробнее: Полезные функции и советы.

Связанные страницы: Руководство по SDK · SDK — быстрый старт · Тест-планы в CI/CD

Зачем нужен протокольный клиент, а не просто HTTP?

Типовой CI-тест по gRPC занимает 30 строк boilerplate’а: сгенерированные stubs, JSON↔protobuf, маппинг status-кодов, тайминг и отдельный загрузчик отчётов. Протокольные клиенты в SDK сводят всё это к:

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

Та же строка записывает шаг TCM (start/end/duration/status/payload preview) и отправляет его через externalruns, так что отчёт о прогоне показывает по-RPC тайм-лайн после завершения теста. Никаких сгенерированных stub’ов, никакого protoc-шага в тест-репо, никакого кастомного reporter’а.

Паттерн Step Recorder

Каждый протокольный клиент принимает опциональный StepRecorder. По умолчанию — NopRecorder, который выбрасывает каждый шаг. Скрипты, которым TCM-тайм-лайн не нужен, работают «из коробки».

Подключите реальный recorder, когда тайм-лайн нужен:

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()

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 } }")

# В конце теста — заливаем тайм-лайн в 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 } }");
client.externalRuns().report("qa",
    new ExternalRunRequest().caseName("my case").status(ExternalRunRequest.STATUS_PASSED));

Cross-language инварианты

Три SDK используют одинаковые формы шагов на стороне сервера, чтобы смешанный тест-сьют корректно сворачивался по (namespace, run, step_key):

  • Таксономия статуса: "passed" / "failed" / "broken" / "skipped". Пустое → "passed".
  • Step key: <name>#<seq>, где <seq> — монотонный счётчик per-client. Ретраи сворачиваются на том же ключе.
  • Имя шага:
    • gRPC: голое полное имя метода (acme.UserService/GetUser).
    • Kafka: produce:<topic> / consume:<topic>.
    • RabbitMQ: publish:<exchange>/<routingKey> / consume:<queue> / declare-queue:<name>.
    • SOAP / GraphQL / SSE / WebSocket: с префиксом протокола (soap:GetUser, graphql:GetUser, sse:collect, ws:send).
  • Маркер обрезки: …(truncated <N>B) (многоточие U+2026). По умолчанию 1024 байта; 0 отключает; отрицательное обнуляется.
  • Классификация ошибок:
    • Cancel/deadline контекста → "broken" (среда, не assertion).
    • Типовая ошибка протокола (gRPC status, kafka.Error, amqp.Error, SOAP fault, GraphQL errors[], HTTP 4xx/5xx) → "failed".
    • Untyped transport (сеть, TLS, DNS) → "broken".

Эргономика тестов (wrap / eventually / parallel)

Три сквозных хелпера делают многошаговые сюиты читаемыми и устойчивыми. Работают
с любым протокольным фасетом и доступны идентично во всех трёх SDK.

  • wrap(name, body) — сгруппировать цепочки, запущенные внутри body, под
    одним именованным родительским шагом, чтобы отчёт рендерился деревом, а не плоско.
  • eventually(within, interval, attempt) — повторять attempt, пока он не
    пройдёт или не истечёт бюджет, со сном interval между попытками. В отчёте
    остаются шаги только успешной (или, при таймауте, последней) попытки —
    промежуточные провалы откатываются. Для eventual consistency (async-эффекты,
    read-after-write лаг). Возвращает, прошло ли.
  • parallel(branch...) — запустить каждую ветку параллельно на своём
    branch-local тестере (общий транспорт / base URL / заголовки + снапшот
    переменных) и слить шаги обратно в порядке веток. Записи переменных внутри
    ветки изолированы.

Сигнал успеха идиоматичен для каждого языка — Go возвращает error (nil = ок),
Python — None или Exception, Java — boolean — но семантика совпадает.

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

JSON-shaped клиент на reflection. Air-gapped инстанс без grpc.reflection может работать с пред-собранными descriptor set (Java) или .proto файлами (Go).

Go

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

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

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);

Стриминг: в v1 поддержаны только unary-вызовы gRPC (InvokeJSON). Server-streaming, client-streaming и bidi в v1 не предоставляются — они удваивают API, а CI-кейсы, которым они нужны, редки.

Kafka

Pure-Go (segmentio/kafka-go) и pure-Java (kafka-clients) — без librdkafka и 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));

Один шаг на каждый consume (не на каждое сообщение); параметр count показывает реально полученное количество. Decode-ошибка downgrade’ит статус с "passed" на "failed" до отправки.

RabbitMQ

Pure-Go (rabbitmq/amqp091-go) и 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)

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);

ContentType по умолчанию — application/json, потребитель может полагаться без per-call boilerplate.

SOAP

JDK DOM (Java) и stdlib xml.etree (Python). Без lxml. XXE-защита включена.

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

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());
}

Параметры шага: operation, http_status, request, response, плюс fault_<key> если в ответе SOAP Fault. SOAP 1.1 (text/xml) и SOAP 1.2 (application/soap+xml) — оба через опцию version.

GraphQL

Стандартный POST envelope {"query": "...", "variables": {...}, "operationName": "..."}. Непустой errors[] → "failed" даже на 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)

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());
    }
}

Имя операции авто-извлекается из тела запроса, чтобы метка шага была graphql:GetUser, а не graphql:anonymous.

SSE

WHATWG-спек парсер для data: / event: / id: / retry: / комментариев. Два режима: collect (cap + deadline) и stream (lazy-генератор).

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)

Java

try (SseClient sse = new SseClient("http://app/events", opts -> opts.recorder(rec))) {
    List<SseEvent> events = sse.collect(5, Duration.ofSeconds(10));
}

collect Python’а классифицирует deadline с нулём событий как "broken" (поток молчит — это env-failure). Cap-hit и deadline + ≥1 событие остаются "passed". Параметр reason (cap | deadline) показывает почему collect вернулся.

WebSocket

JDK 11+ java.net.http.WebSocket (Java) и websockets package (Python, опционально через mockarty[protocols]).

Python

pip install "mockarty[protocols]"   # ставит 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)

Java

try (WebSocketClient ws = new WebSocketClient("ws://app/socket", opts -> opts.recorder(rec))) {
    ws.send("{\"type\":\"subscribe\",\"channel\":\"orders\"}"); // send() принимает String
    String msg = ws.recv(Duration.ofSeconds(5));
}

Шаги — ws:send и ws:recv. Долгие подписки лучше обрабатывать через SseClient, где это применимо.

Air-gapped

Каждый протокольный клиент — pure-language (нет native-зависимостей с отдельной установкой):

  • Go: segmentio/kafka-go, rabbitmq/amqp091-go, jhump/protoreflect — все pure-Go.
  • Python: SOAP / GraphQL / SSE на httpx (core dep); WebSocket — websockets через protocols extra.
  • Java: gRPC через grpc-netty-shaded (Netty встроена); Kafka / RabbitMQ — без native-deps; SOAP / GraphQL / SSE / WebSocket на JDK java.net.http и javax.xml.parsers.

Те же бинарники работают в distroless-контейнерах против Mockarty-mock’ов без внешних инструментов.

Траблшутинг

  • «streaming method not supported in v1» (gRPC) — метод client-streaming или bidi. Используйте язык-специфичный stub — streaming-методов в SDK нет.
  • Пустые шаги в TCM-тайм-лайне — не подключён recorder. По умолчанию NopRecorder выбрасывает всё; передайте WithRecorder(...) (Go/Java) или recorder=... (Python).
  • Шаги захвачены, но отсутствуют в TCM-прогоне — recorder буферизованный. В Python тест должен в конце вызвать client.external_runs().report(steps=recorder.steps()). В Go / Java ExternalRunsRecorder сбрасывает данные асинхронно — вызовите recorder.Close(), чтобы сбросить их перед завершением процесса.
  • Replacement char (�) в payload preview — версия SDK до 2026-05-19 (UTF-8 boundary fix). Обновите.
  • recorder=AccumulatingRecorder() записал ноль шагов в Python — версия SDK до 2026-05-19 (falsy-recorder fix; пустой recorder был bool(rec) == False, и recorder or NopRecorder() его отбрасывал).