Перейти к содержимому

Go SDK (`pkg/client`)

Это содержимое пока не доступно на вашем языке.

The reference Go client talks to Chronacta over gRPC. It uses the same protobuf contract as the CLI (api/proto/stream/v1/stream.proto) and never opens the data directory directly.

  • The SDK follows the compatibility policy: same minor release = wire-compatible with the server; patch releases are forward compatible for clients.
  • Breaking API changes require a documented minor/major bump and CHANGELOG entry.
  • CI smoke check: make sdk-smoke (runs pkg/client integration tests — auth, append/read, reconnect/idempotency, export/import).
  • Transport alternatives (REST/WebSocket) are documented in transports architecture; TypeScript and Python REST v2 clients live under sdk/typescript and sdk/python.

In this repository the client is imported as:

import "aignatov.com/chronacta/pkg/client"

External applications should depend on the same module path or a published module tag when available.

Default server address: 127.0.0.1:2113.

ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
cli, err := client.Dial(ctx, client.Config{Address: "127.0.0.1:2113"})
if err != nil {
log.Fatal(err)
}
defer cli.Close()

Environment variable used by examples: CHRONACTA_SERVER (default 127.0.0.1:2113).

Pass a TLSConfig when the server runs with TLS enabled:

cli, err := client.Dial(ctx, client.Config{
Address: "127.0.0.1:2113",
TLS: &client.TLSConfig{
CAFile: "/path/to/ca.pem",
CertFile: "/path/to/client.pem", // optional mTLS
KeyFile: "/path/to/client-key.pem",
ServerName: "chronacta.local",
},
})

CLI equivalent: -tls -tls-ca FILE [-tls-cert FILE -tls-key FILE -tls-server-name NAME].

When auth is enabled on the server:

cli, err := client.Dial(ctx, client.Config{Address: addr})
if err != nil {
log.Fatal(err)
}
defer cli.Close()
token, err := cli.Login(ctx, "admin", "secret")
if err != nil {
log.Fatal(err)
}
// token is stored on the client; all subsequent RPCs attach Bearer metadata
_ = token
// or set a token obtained elsewhere:
cli.SetToken(token)

Bootstrap (first admin, empty auth store):

_, err := cli.BootstrapAdmin(ctx, "admin", "secret")

Permissions follow server RBAC (stream.read, stream.append, backup, etc.). See gRPC API.

import (
"aignatov.com/chronacta/pkg/event"
"aignatov.com/chronacta/pkg/storage"
)
result, err := cli.AppendToStream(ctx, "orders-1", storage.ExpectedVersionNoStream, []*event.Event{{
EventType: "OrderCreated",
Data: []byte(`{"order_id":"1"}`),
Metadata: map[string]string{"source": "checkout"},
}})

Expected version constants (same as server):

Value Meaning
-2 (ExpectedVersionNoStream) Stream must not exist
-1 (ExpectedVersionAny) Any current version
>= 0 Exact current version

Schema-aware append: set SchemaName and SchemaVersion on the event.

result, err := cli.AppendToStream(ctx, streamID, storage.ExpectedVersionNoStream, events,
client.AppendOptions{IdempotencyKey: "checkout-req-42"},
)
if result.IdempotentReplay {
// server returned the original commit without duplicating events
}

Use the same key and payload after timeouts or reconnect. See duplicate policy.

stream, err := cli.ReadStream(ctx, "orders-1", 1, 100, storage.DirectionForward)
all, err := cli.ReadAll(ctx, 1, 100)
sub, err := cli.Subscribe(ctx, "orders-1", 0, func(ev *event.Event) error {
fmt.Printf("event %s@%d\n", ev.EventType, ev.StreamVersion)
return nil
})
defer sub.Close()
import "aignatov.com/chronacta/pkg/transfer"
count, err := cli.ExportStream(ctx, "orders-1", 1, func(stored *event.Event) error {
return transfer.WriteJSONLRecord(file, stored)
})
summary, err := cli.ImportEvents(ctx, transfer.ImportOptions{
DuplicatePolicy: transfer.DuplicateRejectExistingStream,
}, records)

Binary CLI files use transfer.ParseRecords with transfer.FormatAuto. See export/import.

When the server runs with CHRONACTA_NATS_URL:

status, err := cli.GetNATSBridgeStatus(ctx)
err = cli.ReplayNATSBridge(ctx, 0) // optional manual replay

See NATS bridge.

Use pkg/client/resilience for production code — one Connect() at startup, automatic reconnect and selective retries:

import (
"aignatov.com/chronacta/pkg/client"
"aignatov.com/chronacta/pkg/client/resilience"
)
rc, err := resilience.Connect(ctx, resilience.Config{
Config: client.Config{Address: "127.0.0.1:2113", Token: token},
// MaxReconnects: -1 (default, infinite)
// ReconnectWait: time.Second
})
if err != nil {
log.Fatal(err)
}
defer rc.Close()
result, err := rc.AppendToStream(ctx, streamID, expected, events,
client.AppendOptions{IdempotencyKey: "checkout-42"},
)
  • AppendToStream: auto-retries on transport errors only when IdempotencyKey is set.
  • ReadStream / ReadAll: safe to retry automatically.
  • SubscribeLive: auto-resubscribe transient streams from last position.
  • RunDurableWorker: pull/ack loop with reconnect pause (checkpoint on server).

See Resilient client ADR and examples/resilience-test.

Low-level client.Dial remains for tests. Manual dialWithRetry loops are legacy — prefer resilience.Connect.

Area Client methods
Health Health
Admin VerifyStorage, GetStreamInfo, ListStreams, ListTenants, GetTenantUsage, CreateBackup, UploadBackup, ListRemoteBackups, …
Schema RegisterSchema, GetSchema, ListSchemas, ValidateEvent
Durable subscriptions CreateSubscription, PullSubscriptionEvents, AckSubscription, …
Projections CreateProjection, RebuildProjection, GetProjectionResult, …
Auth admin CreateUser, CreateRole, GrantPermission, …

Full RPC definitions: gRPC API and api/proto/stream/v1/stream.proto.

For HTTP-only integrations, use the REST adapters documented in REST API. Clients mirror the OpenAPI map (pkg/openapi/grpc-http-map.yaml).

import { RestClient } from '@chronacta/client';
const client = new RestClient({ baseUrl: 'http://127.0.0.1:8081', token });
await client.append('orders', 'OrderCreated', { order_id: '1' });
const info = await client.streamInfo('orders');

Smoke: cd sdk/typescript && npm test.

from chronacta import RestClient
client = RestClient("http://127.0.0.1:8081", token=token)
client.append("orders", "OrderCreated", {"order_id": "1"})
info = client.stream_info("orders")

Smoke: cd sdk/python && PYTHONPATH=. python3 -m unittest discover -s tests -v.

CI aggregates Go + TS + Python via make sdk-smoke.

Smoke binaries (require a running server):

Terminal window
./bin/chronacta-example-basic
./bin/chronacta-example-projection
./bin/chronacta-example-subscription
./bin/chronacta-example-export-import
./bin/chronacta-example-sdk # auth, reconnect, idempotency

Optional env for chronacta-example-sdk:

Variable Purpose
CHRONACTA_SERVER gRPC address
CHRONACTA_USERNAME / CHRONACTA_PASSWORD Login when auth enabled
CHRONACTA_TLS_CA, CHRONACTA_TLS_CERT, CHRONACTA_TLS_KEY, CHRONACTA_TLS_SERVER_NAME TLS