diff --git a/.github/workflows/adapters.yml b/.github/workflows/adapters.yml new file mode 100644 index 0000000..74c9f3d --- /dev/null +++ b/.github/workflows/adapters.yml @@ -0,0 +1,56 @@ +name: Adapters + +# The adapter modules live in their own go.mod files so the core module graph +# stays stdlib-only. They are not covered by the root CI's ./... , so build and +# test each one independently here. + +on: + push: + branches: [main] + paths: + - 'adapters/**' + - 'runtime/**' + - 'store/**' + - 'go.mod' + - '.github/workflows/adapters.yml' + pull_request: + branches: [main] + paths: + - 'adapters/**' + - 'runtime/**' + - 'store/**' + - 'go.mod' + - '.github/workflows/adapters.yml' + +permissions: + contents: read + +jobs: + test: + name: ${{ matrix.module }} + runs-on: ubuntu-latest + strategy: + fail-fast: false + matrix: + module: [adapters/otel, adapters/sql] + defaults: + run: + working-directory: ${{ matrix.module }} + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-go@v5 + with: + go-version: '1.26' + check-latest: true + - name: Verify formatting + run: | + fmt_out="$(gofmt -l .)" + if [ -n "$fmt_out" ]; then + echo "Not gofmt-clean:"; echo "$fmt_out"; exit 1 + fi + - name: Build + run: go build ./... + - name: Vet + run: go vet ./... + - name: Test (race) + run: go test -race ./... diff --git a/adapters/otel/README.md b/adapters/otel/README.md new file mode 100644 index 0000000..26e2775 --- /dev/null +++ b/adapters/otel/README.md @@ -0,0 +1,37 @@ +# Isopace OpenTelemetry adapter + +An implementation of the Isopace `runtime.Observer` traces-and-metrics facade +over [OpenTelemetry](https://opentelemetry.io). + +This is a **separate module** so the Isopace core never imports a telemetry SDK +(the core defines a tiny, dependency-free `Observer` interface; this adapter +bridges it to OTel). Wire it in at the edge of your application: + +```go +import ( + oteladapter "github.com/teqpace-services/isopace/adapters/otel" + "github.com/teqpace-services/isopace/runtime" +) + +// From explicit providers: +obs := oteladapter.New(tracerProvider, meterProvider) + +// …or from the globally-registered OpenTelemetry providers: +obs := oteladapter.Default() + +host := runtime.NewHost(runtime.WithObserver(obs)) +``` + +## Mapping + +| Isopace | OpenTelemetry | +|---|---| +| `Observer.StartSpan` | `Tracer.Start` (attributes attached at start) | +| `Span.End` / `SetError` / `SetAttr` | `Span.End` / `RecordError`+`SetStatus(Error)` / `SetAttributes` | +| `Observer.Counter(name).Add` | `Int64Counter.Add` | +| `Observer.Histogram(name).Observe` | `Float64Histogram.Record` | +| `runtime.Attr` | `attribute.KeyValue` (string/bool/int/int64/float64; else stringified) | + +Note: the `Counter`/`Histogram` contract carries no `context.Context`, so the +background context is used for measurements — OTel metric aggregation does not +depend on request context. diff --git a/adapters/otel/go.mod b/adapters/otel/go.mod new file mode 100644 index 0000000..11032e9 --- /dev/null +++ b/adapters/otel/go.mod @@ -0,0 +1,25 @@ +// Isopace OpenTelemetry observer adapter — a separate module so the stdlib-only +// core never imports a telemetry SDK. require/replace mirror the other adapters; +// the OpenTelemetry requirements are filled in by `go mod tidy`. +module github.com/teqpace-services/isopace/adapters/otel + +go 1.26 + +require ( + github.com/teqpace-services/isopace v0.3.0 + go.opentelemetry.io/otel v1.44.0 + go.opentelemetry.io/otel/metric v1.44.0 + go.opentelemetry.io/otel/sdk v1.44.0 + go.opentelemetry.io/otel/trace v1.44.0 +) + +require ( + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-logr/stdr v1.2.2 // indirect + github.com/google/uuid v1.6.0 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + golang.org/x/sys v0.45.0 // indirect +) + +replace github.com/teqpace-services/isopace => ../.. diff --git a/adapters/otel/go.sum b/adapters/otel/go.sum new file mode 100644 index 0000000..280996e --- /dev/null +++ b/adapters/otel/go.sum @@ -0,0 +1,35 @@ +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= +go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= +go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= +go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= +go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= +go.opentelemetry.io/otel/sdk v1.44.0 h1:nHYwb9lK+fJPU/dnT6s7W7Z8itMWyqrnVfbheVYrZ58= +go.opentelemetry.io/otel/sdk v1.44.0/go.mod h1:Osuydd3Se74nqjAKxid74N5eC+jfEqfTegHRnq58oK0= +go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRks6si09iEfI= +go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA= +go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= +go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= +golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/adapters/otel/otel.go b/adapters/otel/otel.go new file mode 100644 index 0000000..6b43793 --- /dev/null +++ b/adapters/otel/otel.go @@ -0,0 +1,146 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// +// Copyright (C) 2026 Teqpace Services Ltd. +// +// This file is part of Isopace, a financial transaction framework. +// +// Isopace is dual-licensed: +// - under the GNU Affero General Public License v3.0 or later (see LICENSE); or +// - under a commercial license from Teqpace Services Ltd. (see COMMERCIAL-LICENSE.md). +// +// Authorship is recorded in the AUTHORS file. + +// Package otel implements the Isopace runtime.Observer traces-and-metrics facade +// over OpenTelemetry. It lives in a separate module so the Isopace core never +// imports a telemetry SDK; wire it in at the edge: +// +// host := runtime.NewHost(runtime.WithObserver(otel.Default())) +package otel + +import ( + "context" + "fmt" + + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/metric" + "go.opentelemetry.io/otel/trace" + + "github.com/teqpace-services/isopace/runtime" +) + +// scopeName is the instrumentation scope reported to OpenTelemetry. +const scopeName = "github.com/teqpace-services/isopace" + +// Observer implements runtime.Observer over OpenTelemetry providers. +type Observer struct { + tracer trace.Tracer + meter metric.Meter +} + +var _ runtime.Observer = (*Observer)(nil) + +// New builds an Observer from explicit OpenTelemetry providers. +func New(tp trace.TracerProvider, mp metric.MeterProvider) *Observer { + return &Observer{tracer: tp.Tracer(scopeName), meter: mp.Meter(scopeName)} +} + +// Default builds an Observer from the globally-registered OpenTelemetry +// providers (otel.GetTracerProvider / otel.GetMeterProvider). +func Default() *Observer { + return New(otel.GetTracerProvider(), otel.GetMeterProvider()) +} + +// StartSpan begins a span; the returned context carries it for propagation. +func (o *Observer) StartSpan(ctx context.Context, name string, attrs ...runtime.Attr) (context.Context, runtime.Span) { + ctx, sp := o.tracer.Start(ctx, name, trace.WithAttributes(kvs(attrs)...)) + return ctx, &span{sp: sp} +} + +// Counter returns a monotonic Int64 counter by name. +func (o *Observer) Counter(name string) runtime.Counter { + c, err := o.meter.Int64Counter(name) + if err != nil { + return noopCounter{} + } + return &counter{c: c} +} + +// Histogram returns a Float64 distribution instrument by name. +func (o *Observer) Histogram(name string) runtime.Histogram { + h, err := o.meter.Float64Histogram(name) + if err != nil { + return noopHistogram{} + } + return &histogram{h: h} +} + +type span struct{ sp trace.Span } + +func (s *span) End() { s.sp.End() } + +func (s *span) SetError(err error) { + if err == nil { + return + } + s.sp.RecordError(err) + s.sp.SetStatus(codes.Error, err.Error()) +} + +func (s *span) SetAttr(attrs ...runtime.Attr) { s.sp.SetAttributes(kvs(attrs)...) } + +type counter struct{ c metric.Int64Counter } + +// Add records an increment. The runtime.Counter contract carries no context, so +// the background context is used; OpenTelemetry metric instruments do not depend +// on request context for aggregation. +func (c *counter) Add(n int64, attrs ...runtime.Attr) { + c.c.Add(context.Background(), n, metric.WithAttributes(kvs(attrs)...)) +} + +type histogram struct{ h metric.Float64Histogram } + +func (h *histogram) Observe(v float64, attrs ...runtime.Attr) { + h.h.Record(context.Background(), v, metric.WithAttributes(kvs(attrs)...)) +} + +// noop instruments returned when the meter fails to construct one, so callers +// never need a nil check. +type noopCounter struct{} + +func (noopCounter) Add(int64, ...runtime.Attr) {} + +type noopHistogram struct{} + +func (noopHistogram) Observe(float64, ...runtime.Attr) {} + +// kvs converts Isopace attributes to OpenTelemetry key/values. +func kvs(attrs []runtime.Attr) []attribute.KeyValue { + if len(attrs) == 0 { + return nil + } + out := make([]attribute.KeyValue, 0, len(attrs)) + for _, a := range attrs { + out = append(out, kv(a)) + } + return out +} + +func kv(a runtime.Attr) attribute.KeyValue { + k := attribute.Key(a.Key) + switch v := a.Value.(type) { + case string: + return k.String(v) + case bool: + return k.Bool(v) + case int: + return k.Int(v) + case int64: + return k.Int64(v) + case float64: + return k.Float64(v) + default: + return k.String(fmt.Sprintf("%v", v)) + } +} diff --git a/adapters/otel/otel_test.go b/adapters/otel/otel_test.go new file mode 100644 index 0000000..be5df2f --- /dev/null +++ b/adapters/otel/otel_test.go @@ -0,0 +1,76 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// +// Copyright (C) 2026 Teqpace Services Ltd. +// +// This file is part of Isopace, a financial transaction framework. +// +// Isopace is dual-licensed: +// - under the GNU Affero General Public License v3.0 or later (see LICENSE); or +// - under a commercial license from Teqpace Services Ltd. (see COMMERCIAL-LICENSE.md). +// +// Authorship is recorded in the AUTHORS file. + +package otel_test + +import ( + "context" + "errors" + "testing" + + "go.opentelemetry.io/otel/codes" + noopmetric "go.opentelemetry.io/otel/metric/noop" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" + nooptrace "go.opentelemetry.io/otel/trace/noop" + + oteladapter "github.com/teqpace-services/isopace/adapters/otel" + "github.com/teqpace-services/isopace/runtime" +) + +var _ runtime.Observer = (*oteladapter.Observer)(nil) + +func TestStartSpanRecordsNameAttributesAndError(t *testing.T) { + sr := tracetest.NewSpanRecorder() + tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(sr)) + obs := oteladapter.New(tp, noopmetric.NewMeterProvider()) + + _, span := obs.StartSpan(context.Background(), "auth", runtime.A("mti", "0200")) + span.SetAttr(runtime.A("rc", "00")) + span.SetError(errors.New("declined")) + span.End() + + spans := sr.Ended() + if len(spans) != 1 { + t.Fatalf("recorded %d spans, want 1", len(spans)) + } + s := spans[0] + if s.Name() != "auth" { + t.Errorf("span name = %q want %q", s.Name(), "auth") + } + got := map[string]string{} + for _, a := range s.Attributes() { + got[string(a.Key)] = a.Value.AsString() + } + for k, want := range map[string]string{"mti": "0200", "rc": "00"} { + if got[k] != want { + t.Errorf("attribute %q = %q want %q (all=%v)", k, got[k], want, got) + } + } + if s.Status().Code != codes.Error { + t.Errorf("span status = %v want Error", s.Status().Code) + } +} + +func TestCounterAndHistogramDoNotPanic(t *testing.T) { + obs := oteladapter.New(nooptrace.NewTracerProvider(), noopmetric.NewMeterProvider()) + obs.Counter("isopace_test_total").Add(1, runtime.A("k", "v")) + obs.Histogram("isopace_test_latency_ms").Observe(1.5, runtime.A("k", "v")) + _, sp := obs.StartSpan(context.Background(), "noop") + sp.End() +} + +func TestDefaultUsesGlobalProviders(t *testing.T) { + if oteladapter.Default() == nil { + t.Fatal("Default() returned nil") + } +} diff --git a/adapters/sql/README.md b/adapters/sql/README.md new file mode 100644 index 0000000..5b051c1 --- /dev/null +++ b/adapters/sql/README.md @@ -0,0 +1,39 @@ +# Isopace SQL store adapter + +A `store.Store` implementation backed by the standard library `database/sql`. + +This is a **separate module** so the Isopace core stays stdlib-only: the adapter +imports only `database/sql`; **you** supply the driver, so no SQL driver enters +your module graph except the one you already chose. + +```go +import ( + "database/sql" + + _ "github.com/lib/pq" // your driver + sqlstore "github.com/teqpace-services/isopace/adapters/sql" +) + +db, _ := sql.Open("postgres", dsn) +st, _ := sqlstore.New(db, sqlstore.Postgres) +_ = st.EnsureSchema(ctx) // creates the (collection, k, v) table if absent +// st satisfies store.Store +``` + +## Dialects + +`New` takes a `Dialect` carrying the two things that differ across drivers — the +positional placeholder syntax and the opaque-value column type. Built-ins: +`SQLite`, `MySQL`, `Postgres`. Construct your own `Dialect` for others. + +Writes use a portable UPDATE-then-INSERT upsert (no dialect-specific +`ON CONFLICT` / `ON DUPLICATE KEY`). + +## Testing + +Unit tests (dialects, validation) run with no database. The round-trip test is +skipped unless you set `ISOPACE_SQL_DRIVER`, `ISOPACE_SQL_DSN` (and +`ISOPACE_SQL_DIALECT`) **and** import the driver into the test binary. + +> Follow-up: wire a Postgres service container into CI to run the round-trip on +> every build (see `ROADMAP-to-v1.md`, B2). diff --git a/adapters/sql/go.mod b/adapters/sql/go.mod new file mode 100644 index 0000000..2f08e43 --- /dev/null +++ b/adapters/sql/go.mod @@ -0,0 +1,13 @@ +// Isopace SQL store adapter — a separate module so the stdlib-only core never +// gains a database dependency. The adapter itself imports only database/sql; +// the integrator supplies the driver. +module github.com/teqpace-services/isopace/adapters/sql + +go 1.26 + +require github.com/teqpace-services/isopace v0.3.0 + +// In-repo builds (and CI building this module from a checkout) resolve the core +// from the working tree. External consumers ignore replace and use the require +// above. +replace github.com/teqpace-services/isopace => ../.. diff --git a/adapters/sql/sql.go b/adapters/sql/sql.go new file mode 100644 index 0000000..6d55008 --- /dev/null +++ b/adapters/sql/sql.go @@ -0,0 +1,168 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// +// Copyright (C) 2026 Teqpace Services Ltd. +// +// This file is part of Isopace, a financial transaction framework. +// +// Isopace is dual-licensed: +// - under the GNU Affero General Public License v3.0 or later (see LICENSE); or +// - under a commercial license from Teqpace Services Ltd. (see COMMERCIAL-LICENSE.md). +// +// Authorship is recorded in the AUTHORS file. + +// Package sql implements the Isopace store.Store interface over the standard +// library database/sql package. The caller supplies a configured *sql.DB (and +// therefore chooses the driver), so this adapter imports no SQL driver and adds +// no third-party dependency to a consumer's module graph beyond the driver they +// already chose. +package sql + +import ( + "context" + "database/sql" + "errors" + "fmt" + "regexp" + "strconv" + + "github.com/teqpace-services/isopace/store" +) + +// Store implements store.Store. +var _ store.Store = (*Store)(nil) + +// Dialect captures the few SQL details that differ across drivers: the +// positional placeholder syntax and the column type for opaque byte values. +type Dialect struct { + Name string + Placeholder func(n int) string // 1-based positional placeholder, e.g. "?" or "$1" + BlobType string // column type for opaque values +} + +// Built-in dialects. Add others by constructing a Dialect. +var ( + SQLite = Dialect{Name: "sqlite", Placeholder: qmark, BlobType: "BLOB"} + MySQL = Dialect{Name: "mysql", Placeholder: qmark, BlobType: "BLOB"} + Postgres = Dialect{Name: "postgres", Placeholder: dollar, BlobType: "BYTEA"} +) + +func qmark(int) string { return "?" } +func dollar(n int) string { return "$" + strconv.Itoa(n) } + +// validIdent guards the table name, which is interpolated into SQL (placeholders +// cannot parameterise identifiers). +var validIdent = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`) + +// Store is a collection/key/value store backed by a SQL table with columns +// (collection, k, v) and a composite primary key (collection, k). +type Store struct { + db *sql.DB + dialect Dialect + table string +} + +// Option configures a Store. +type Option func(*Store) + +// WithTable sets the backing table name (default "isopace_kv"). The name must be +// a plain SQL identifier. +func WithTable(name string) Option { return func(s *Store) { s.table = name } } + +// New returns a Store over db using the given Dialect. db must already be +// configured with a registered driver. Call EnsureSchema once to create the +// backing table if needed. +func New(db *sql.DB, d Dialect, opts ...Option) (*Store, error) { + s := &Store{db: db, dialect: d, table: "isopace_kv"} + for _, o := range opts { + o(s) + } + if d.Placeholder == nil { + return nil, errors.New("sqlstore: dialect has no Placeholder func") + } + if !validIdent.MatchString(s.table) { + return nil, fmt.Errorf("sqlstore: invalid table name %q", s.table) + } + return s, nil +} + +func (s *Store) ph(n int) string { return s.dialect.Placeholder(n) } + +// EnsureSchema creates the backing table if it does not already exist. +func (s *Store) EnsureSchema(ctx context.Context) error { + q := fmt.Sprintf( + `CREATE TABLE IF NOT EXISTS %s (collection TEXT NOT NULL, k TEXT NOT NULL, v %s, PRIMARY KEY (collection, k))`, + s.table, s.dialect.BlobType) + _, err := s.db.ExecContext(ctx, q) + return err +} + +// Get returns the value stored under (collection, key), or store.ErrNotFound. +func (s *Store) Get(ctx context.Context, collection, key string) ([]byte, error) { + q := fmt.Sprintf(`SELECT v FROM %s WHERE collection = %s AND k = %s`, s.table, s.ph(1), s.ph(2)) + var v []byte + switch err := s.db.QueryRowContext(ctx, q, collection, key).Scan(&v); { + case errors.Is(err, sql.ErrNoRows): + return nil, store.ErrNotFound + case err != nil: + return nil, err + } + return v, nil +} + +// Put stores value under (collection, key), inserting or updating as needed. +func (s *Store) Put(ctx context.Context, collection, key string, value []byte) error { + // Portable upsert: UPDATE first, INSERT only if nothing was updated. This + // avoids dialect-specific ON CONFLICT / ON DUPLICATE KEY syntax. + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer func() { _ = tx.Rollback() }() + + upd := fmt.Sprintf(`UPDATE %s SET v = %s WHERE collection = %s AND k = %s`, s.table, s.ph(1), s.ph(2), s.ph(3)) + res, err := tx.ExecContext(ctx, upd, value, collection, key) + if err != nil { + return err + } + n, err := res.RowsAffected() + if err != nil { + return err + } + if n == 0 { + ins := fmt.Sprintf(`INSERT INTO %s (collection, k, v) VALUES (%s, %s, %s)`, s.table, s.ph(1), s.ph(2), s.ph(3)) + if _, err := tx.ExecContext(ctx, ins, collection, key, value); err != nil { + return err + } + } + return tx.Commit() +} + +// Delete removes (collection, key). Deleting a missing key is not an error. +func (s *Store) Delete(ctx context.Context, collection, key string) error { + q := fmt.Sprintf(`DELETE FROM %s WHERE collection = %s AND k = %s`, s.table, s.ph(1), s.ph(2)) + _, err := s.db.ExecContext(ctx, q, collection, key) + return err +} + +// List returns the keys in a collection in ascending order. +func (s *Store) List(ctx context.Context, collection string) ([]string, error) { + q := fmt.Sprintf(`SELECT k FROM %s WHERE collection = %s ORDER BY k`, s.table, s.ph(1)) + rows, err := s.db.QueryContext(ctx, q, collection) + if err != nil { + return nil, err + } + defer func() { _ = rows.Close() }() + + var keys []string + for rows.Next() { + var k string + if err := rows.Scan(&k); err != nil { + return nil, err + } + keys = append(keys, k) + } + return keys, rows.Err() +} + +// Close closes the underlying *sql.DB. +func (s *Store) Close() error { return s.db.Close() } diff --git a/adapters/sql/sql_test.go b/adapters/sql/sql_test.go new file mode 100644 index 0000000..c336289 --- /dev/null +++ b/adapters/sql/sql_test.go @@ -0,0 +1,107 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later +// +// Copyright (C) 2026 Teqpace Services Ltd. +// +// This file is part of Isopace, a financial transaction framework. +// +// Isopace is dual-licensed: +// - under the GNU Affero General Public License v3.0 or later (see LICENSE); or +// - under a commercial license from Teqpace Services Ltd. (see COMMERCIAL-LICENSE.md). +// +// Authorship is recorded in the AUTHORS file. + +package sql_test + +import ( + "bytes" + "context" + "database/sql" + "errors" + "os" + "testing" + + sqlstore "github.com/teqpace-services/isopace/adapters/sql" + "github.com/teqpace-services/isopace/store" +) + +func TestDialectPlaceholders(t *testing.T) { + if got := sqlstore.Postgres.Placeholder(1); got != "$1" { + t.Errorf("Postgres placeholder = %q want $1", got) + } + if got := sqlstore.Postgres.Placeholder(3); got != "$3" { + t.Errorf("Postgres placeholder = %q want $3", got) + } + if got := sqlstore.SQLite.Placeholder(2); got != "?" { + t.Errorf("SQLite placeholder = %q want ?", got) + } +} + +func TestInvalidTableRejected(t *testing.T) { + if _, err := sqlstore.New(nil, sqlstore.SQLite, sqlstore.WithTable("kv; DROP TABLE x")); err == nil { + t.Fatal("expected an invalid table name to be rejected") + } +} + +func TestMissingDialectPlaceholder(t *testing.T) { + if _, err := sqlstore.New(nil, sqlstore.Dialect{Name: "x"}); err == nil { + t.Fatal("expected a dialect with no Placeholder to be rejected") + } +} + +// TestRoundTrip exercises the store against a real database. It is skipped unless +// ISOPACE_SQL_DRIVER and ISOPACE_SQL_DSN are set AND the chosen driver is +// imported into the test binary (add a blank import in your fork, e.g. +// _ "github.com/lib/pq"). CI wires this with a database service. +func TestRoundTrip(t *testing.T) { + driver, dsn := os.Getenv("ISOPACE_SQL_DRIVER"), os.Getenv("ISOPACE_SQL_DSN") + if driver == "" || dsn == "" { + t.Skip("set ISOPACE_SQL_DRIVER and ISOPACE_SQL_DSN (and import the driver) to run") + } + db, err := sql.Open(driver, dsn) + if err != nil { + t.Fatalf("open: %v", err) + } + dialect := sqlstore.SQLite + if d := os.Getenv("ISOPACE_SQL_DIALECT"); d == "postgres" { + dialect = sqlstore.Postgres + } else if d == "mysql" { + dialect = sqlstore.MySQL + } + st, err := sqlstore.New(db, dialect, sqlstore.WithTable("isopace_kv_test")) + if err != nil { + t.Fatalf("new: %v", err) + } + defer func() { _ = st.Close() }() + + ctx := context.Background() + if err := st.EnsureSchema(ctx); err != nil { + t.Fatalf("ensure schema: %v", err) + } + + if _, err := st.Get(ctx, "c", "missing"); !errors.Is(err, store.ErrNotFound) { + t.Errorf("Get(missing) err = %v want ErrNotFound", err) + } + if err := st.Put(ctx, "c", "k1", []byte("v1")); err != nil { + t.Fatalf("put: %v", err) + } + if err := st.Put(ctx, "c", "k1", []byte("v1b")); err != nil { // upsert + t.Fatalf("upsert: %v", err) + } + got, err := st.Get(ctx, "c", "k1") + if err != nil || !bytes.Equal(got, []byte("v1b")) { + t.Errorf("Get = %q, %v want v1b", got, err) + } + if err := st.Put(ctx, "c", "k2", []byte("v2")); err != nil { + t.Fatalf("put k2: %v", err) + } + keys, err := st.List(ctx, "c") + if err != nil || len(keys) != 2 { + t.Errorf("List = %v, %v want 2 keys", keys, err) + } + if err := st.Delete(ctx, "c", "k1"); err != nil { + t.Fatalf("delete: %v", err) + } + if _, err := st.Get(ctx, "c", "k1"); !errors.Is(err, store.ErrNotFound) { + t.Errorf("after delete, Get err = %v want ErrNotFound", err) + } +}