From 6f74ae97b3e88739cc72c9013ee19582713e1bfc Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Fri, 27 Mar 2026 22:48:04 -0700 Subject: [PATCH 1/9] Initial imple --- connectors/couchdb/conn.go | 630 +++++++++++++++++++++++++ connectors/couchdb/conn_test.go | 223 +++++++++ connectors/couchdb/integration_test.go | 268 +++++++++++ go.mod | 1 + go.sum | 10 + internal/app/options/connectorflags.go | 53 +++ 6 files changed, 1185 insertions(+) create mode 100644 connectors/couchdb/conn.go create mode 100644 connectors/couchdb/conn_test.go create mode 100644 connectors/couchdb/integration_test.go diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go new file mode 100644 index 00000000..63b4a9f9 --- /dev/null +++ b/connectors/couchdb/conn.go @@ -0,0 +1,630 @@ +/* + * Copyright (C) 2024 Adiom, Inc. + * + * SPDX-License-Identifier: AGPL-3.0-or-later + */ +package couchdb + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "net/url" + "slices" + "strconv" + "strings" + "time" + + "connectrpc.com/connect" + adiomv1 "github.com/adiom-data/dsync/gen/adiom/v1" + "github.com/adiom-data/dsync/gen/adiom/v1/adiomv1connect" + "github.com/go-kivik/kivik/v4" + _ "github.com/go-kivik/kivik/v4/couchdb" + "github.com/mitchellh/hashstructure" +) + +const ( + connectorDBType string = "CouchDB" + connectorSpec string = "Apache CouchDB / IBM Cloudant" +) + +var ExcludedSystemDatabases = []string{"_users", "_replicator", "_global_changes"} + +type ConnectorSettings struct { + Uri string + ServerConnectTimeout time.Duration + PingTimeout time.Duration + WriterMaxBatchSize int + TargetDocCountPerPartition int64 + MaxPageSize int + IncludeSystemDbs bool +} + +func setDefault[T comparable](field *T, defaultValue T) { + if *field == *new(T) { + *field = defaultValue + } +} + +type conn struct { + adiomv1connect.UnimplementedConnectorServiceHandler + client *kivik.Client + settings ConnectorSettings + dsn string +} + +func generateConnectorID(connectionString string) string { + id, err := hashstructure.Hash(connectionString, nil) + if err != nil { + panic(fmt.Sprintf("Failed to hash the connection string: %v", err)) + } + return strconv.FormatUint(id, 16) +} + +func convertUriToDSN(uri string) (string, error) { + lower := strings.ToLower(uri) + if strings.HasPrefix(lower, "couchdb://") { + return "http://" + strings.TrimPrefix(uri, "couchdb://"), nil + } + if strings.HasPrefix(lower, "couchdbs://") { + return "https://" + strings.TrimPrefix(uri, "couchdbs://"), nil + } + if strings.HasPrefix(lower, "cloudant://") { + return "https://" + strings.TrimPrefix(uri, "cloudant://"), nil + } + if strings.HasPrefix(lower, "http://") || strings.HasPrefix(lower, "https://") { + return uri, nil + } + return "", fmt.Errorf("unsupported URI scheme: %s", uri) +} + +func getAllDatabases(ctx context.Context, client *kivik.Client, includeSystem bool) ([]string, error) { + dbs, err := client.AllDBs(ctx) + if err != nil { + return nil, fmt.Errorf("failed to list databases: %w", err) + } + + if includeSystem { + return dbs, nil + } + + filtered := slices.DeleteFunc(dbs, func(db string) bool { + return strings.HasPrefix(db, "_") || slices.Contains(ExcludedSystemDatabases, db) + }) + + return filtered, nil +} + +func (c *conn) GetInfo(ctx context.Context, r *connect.Request[adiomv1.GetInfoRequest]) (*connect.Response[adiomv1.GetInfoResponse], error) { + version, err := c.client.Version(ctx) + if err != nil { + return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to get server version: %w", err)) + } + + dbType := connectorDBType + if strings.Contains(strings.ToLower(version.Vendor), "cloudant") { + dbType = "Cloudant" + } + + return connect.NewResponse(&adiomv1.GetInfoResponse{ + Id: generateConnectorID(c.settings.Uri), + DbType: dbType, + Version: version.Version, + Spec: connectorSpec, + Capabilities: &adiomv1.Capabilities{ + Source: &adiomv1.Capabilities_Source{ + SupportedDataTypes: []adiomv1.DataType{adiomv1.DataType_DATA_TYPE_JSON_ID}, + LsnStream: true, + MultiNamespacePlan: true, + DefaultPlan: true, + }, + Sink: &adiomv1.Capabilities_Sink{ + SupportedDataTypes: []adiomv1.DataType{adiomv1.DataType_DATA_TYPE_JSON_ID}, + }, + }, + }), nil +} + +func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.GeneratePlanRequest]) (*connect.Response[adiomv1.GeneratePlanResponse], error) { + var namespaces []string + + if len(r.Msg.GetNamespaces()) > 0 { + namespaces = r.Msg.GetNamespaces() + } else { + dbs, err := getAllDatabases(ctx, c.client, c.settings.IncludeSystemDbs) + if err != nil { + return nil, connect.NewError(connect.CodeInternal, err) + } + namespaces = dbs + } + + var partitions []*adiomv1.Partition + var updatesPartitions []*adiomv1.UpdatesPartition + + for _, ns := range namespaces { + db := c.client.DB(ns) + + stats, err := db.Stats(ctx) + if err != nil { + slog.Warn(fmt.Sprintf("Failed to get stats for database %s: %v", ns, err)) + continue + } + + if r.Msg.GetInitialSync() { + count := int64(stats.DocCount) + if count <= c.settings.TargetDocCountPerPartition { + partitions = append(partitions, &adiomv1.Partition{ + Namespace: ns, + EstimatedCount: uint64(count), + }) + } else { + numPartitions := (count + c.settings.TargetDocCountPerPartition - 1) / c.settings.TargetDocCountPerPartition + docsPerPartition := count / numPartitions + + rows := db.AllDocs(ctx, kivik.Params(map[string]interface{}{ + "limit": numPartitions - 1, + "skip": docsPerPartition, + })) + defer rows.Close() + + var boundaries []string + for rows.Next() { + id, err := rows.ID() + if err != nil { + continue + } + boundaries = append(boundaries, id) + if len(boundaries) >= int(numPartitions)-1 { + break + } + } + + var prevKey string + for i, boundary := range boundaries { + cursor := encodeCursor(prevKey, boundary) + partitions = append(partitions, &adiomv1.Partition{ + Namespace: ns, + Cursor: cursor, + EstimatedCount: uint64(docsPerPartition), + }) + prevKey = boundary + if i == len(boundaries)-1 { + partitions = append(partitions, &adiomv1.Partition{ + Namespace: ns, + Cursor: encodeCursor(boundary, ""), + EstimatedCount: uint64(docsPerPartition), + }) + } + } + + if len(boundaries) == 0 { + partitions = append(partitions, &adiomv1.Partition{ + Namespace: ns, + EstimatedCount: uint64(count), + }) + } + } + } else { + partitions = append(partitions, &adiomv1.Partition{ + Namespace: ns, + }) + } + + if r.Msg.GetUpdates() { + changes := db.Changes(ctx, kivik.Params(map[string]interface{}{ + "descending": true, + "limit": 1, + })) + var lastSeq string + for changes.Next() { + lastSeq = changes.Seq() + } + if err := changes.Err(); err != nil { + slog.Warn(fmt.Sprintf("Failed to get changes for database %s: %v", ns, err)) + } + changes.Close() + + if lastSeq == "" { + info := db.Changes(ctx, kivik.Params(map[string]interface{}{ + "since": "now", + "limit": 0, + })) + for info.Next() { + } + meta, err := info.Metadata() + if err == nil && meta != nil { + lastSeq = meta.LastSeq + } + info.Close() + } + + updatesPartitions = append(updatesPartitions, &adiomv1.UpdatesPartition{ + Namespaces: []string{ns}, + Cursor: []byte(lastSeq), + }) + } + } + + return connect.NewResponse(&adiomv1.GeneratePlanResponse{ + Partitions: partitions, + UpdatesPartitions: updatesPartitions, + }), nil +} + +func encodeCursor(startKey, endKey string) []byte { + if startKey == "" && endKey == "" { + return nil + } + data, _ := json.Marshal(map[string]string{"s": startKey, "e": endKey}) + return data +} + +func decodeCursor(cursor []byte) (startKey, endKey string) { + if len(cursor) == 0 { + return "", "" + } + var m map[string]string + if err := json.Unmarshal(cursor, &m); err != nil { + return "", "" + } + return m["s"], m["e"] +} + +func (c *conn) GetNamespaceMetadata(ctx context.Context, r *connect.Request[adiomv1.GetNamespaceMetadataRequest]) (*connect.Response[adiomv1.GetNamespaceMetadataResponse], error) { + db := c.client.DB(r.Msg.GetNamespace()) + stats, err := db.Stats(ctx) + if err != nil { + return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to get database stats: %w", err)) + } + + return connect.NewResponse(&adiomv1.GetNamespaceMetadataResponse{ + Count: uint64(stats.DocCount), + }), nil +} + +func (c *conn) ListData(ctx context.Context, r *connect.Request[adiomv1.ListDataRequest]) (*connect.Response[adiomv1.ListDataResponse], error) { + partition := r.Msg.GetPartition() + db := c.client.DB(partition.GetNamespace()) + + pageSize := c.settings.MaxPageSize + if pageSize == 0 { + pageSize = 1000 + } + + params := map[string]interface{}{ + "include_docs": true, + "limit": pageSize + 1, // fetch one extra to detect if there's a next page + } + + pageCursor := r.Msg.GetCursor() + if len(pageCursor) > 0 { + // Page cursor from previous call - start after this key + params["startkey"] = string(pageCursor) + params["skip"] = 1 // skip the document we already returned + } else { + // Initial call - use partition cursor if present + startKey, endKey := decodeCursor(partition.GetCursor()) + if startKey != "" { + params["startkey"] = startKey + params["skip"] = 1 + } + if endKey != "" { + params["endkey"] = endKey + } + } + + rows := db.AllDocs(ctx, kivik.Params(params)) + defer rows.Close() + + var data [][]byte + var lastID string + hasMore := false + + for rows.Next() { + if len(data) >= pageSize { + // We got one extra - there's more data + hasMore = true + break + } + + var doc map[string]interface{} + if err := rows.ScanDoc(&doc); err != nil { + slog.Warn(fmt.Sprintf("Failed to scan document: %v", err)) + continue + } + + delete(doc, "_rev") + + jsonBytes, err := json.Marshal(doc) + if err != nil { + slog.Warn(fmt.Sprintf("Failed to marshal document: %v", err)) + continue + } + + if id, err := rows.ID(); err == nil { + lastID = id + } + data = append(data, jsonBytes) + } + + if err := rows.Err(); err != nil { + if !errors.Is(err, context.Canceled) { + slog.Error(fmt.Sprintf("Error iterating documents: %v", err)) + } + return nil, connect.NewError(connect.CodeInternal, err) + } + + var nextCursor []byte + if hasMore && lastID != "" { + nextCursor = []byte(lastID) + } + + return connect.NewResponse(&adiomv1.ListDataResponse{ + Data: data, + NextCursor: nextCursor, + }), nil +} + +func (c *conn) WriteData(ctx context.Context, r *connect.Request[adiomv1.WriteDataRequest]) (*connect.Response[adiomv1.WriteDataResponse], error) { + db := c.client.DB(r.Msg.GetNamespace()) + + var docs []interface{} + for _, data := range r.Msg.GetData() { + var doc map[string]interface{} + if err := json.Unmarshal(data, &doc); err != nil { + return nil, connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("failed to unmarshal document: %w", err)) + } + delete(doc, "_rev") + docs = append(docs, doc) + + if c.settings.WriterMaxBatchSize > 0 && len(docs) >= c.settings.WriterMaxBatchSize { + if _, err := db.BulkDocs(ctx, docs); err != nil { + return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk insert documents: %w", err)) + } + docs = nil + } + } + + if len(docs) > 0 { + if _, err := db.BulkDocs(ctx, docs); err != nil { + return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk insert documents: %w", err)) + } + } + + return connect.NewResponse(&adiomv1.WriteDataResponse{}), nil +} + +func (c *conn) WriteUpdates(ctx context.Context, r *connect.Request[adiomv1.WriteUpdatesRequest]) (*connect.Response[adiomv1.WriteUpdatesResponse], error) { + db := c.client.DB(r.Msg.GetNamespace()) + + for _, update := range r.Msg.GetUpdates() { + if len(update.GetId()) == 0 { + continue + } + + docID := string(update.GetId()[0].GetData()) + + switch update.GetType() { + case adiomv1.UpdateType_UPDATE_TYPE_INSERT, adiomv1.UpdateType_UPDATE_TYPE_UPDATE: + var doc map[string]interface{} + if err := json.Unmarshal(update.GetData(), &doc); err != nil { + slog.Error(fmt.Sprintf("Failed to unmarshal update data: %v", err)) + continue + } + + var existingRev string + result := db.Get(ctx, docID) + var existing map[string]interface{} + if err := result.ScanDoc(&existing); err == nil { + if rev, ok := existing["_rev"].(string); ok { + existingRev = rev + } + } + + if existingRev != "" { + doc["_rev"] = existingRev + } else { + delete(doc, "_rev") + } + + if _, err := db.Put(ctx, docID, doc); err != nil { + slog.Error(fmt.Sprintf("Failed to upsert document %s: %v", docID, err)) + } + + case adiomv1.UpdateType_UPDATE_TYPE_DELETE: + result := db.Get(ctx, docID) + var existing map[string]interface{} + if err := result.ScanDoc(&existing); err == nil { + if rev, ok := existing["_rev"].(string); ok { + if _, err := db.Delete(ctx, docID, rev); err != nil { + slog.Error(fmt.Sprintf("Failed to delete document %s: %v", docID, err)) + } + } + } + } + } + + return connect.NewResponse(&adiomv1.WriteUpdatesResponse{}), nil +} + +func (c *conn) StreamUpdates(ctx context.Context, r *connect.Request[adiomv1.StreamUpdatesRequest], s *connect.ServerStream[adiomv1.StreamUpdatesResponse]) error { + namespaces := r.Msg.GetNamespaces() + since := string(r.Msg.GetCursor()) + + if len(namespaces) != 1 { + return connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("expected exactly 1 namespace for streaming")) + } + + ns := namespaces[0] + db := c.client.DB(ns) + + params := map[string]interface{}{ + "feed": "continuous", + "include_docs": true, + "heartbeat": 1000, // 1 second heartbeat + } + + // Only set 'since' if we have a valid cursor + if since != "" { + params["since"] = since + } + + changes := db.Changes(ctx, kivik.Params(params)) + defer changes.Close() + + for changes.Next() { + seq := changes.Seq() + id := changes.ID() + idBytes := []byte(id) + + var update *adiomv1.Update + + if changes.Deleted() { + update = &adiomv1.Update{ + Id: []*adiomv1.BsonValue{{ + Data: idBytes, + Type: 2, // string type + Name: "_id", + }}, + Type: adiomv1.UpdateType_UPDATE_TYPE_DELETE, + } + } else { + var doc map[string]interface{} + if err := changes.ScanDoc(&doc); err != nil { + slog.Error(fmt.Sprintf("Failed to scan change doc: %v", err)) + continue + } + + delete(doc, "_rev") + + jsonBytes, err := json.Marshal(doc) + if err != nil { + slog.Error(fmt.Sprintf("Failed to marshal change doc: %v", err)) + continue + } + + update = &adiomv1.Update{ + Id: []*adiomv1.BsonValue{{ + Data: idBytes, + Type: 2, + Name: "_id", + }}, + Type: adiomv1.UpdateType_UPDATE_TYPE_UPDATE, + Data: jsonBytes, + } + } + + // Send each update immediately for continuous feed + if err := s.Send(&adiomv1.StreamUpdatesResponse{ + Updates: []*adiomv1.Update{update}, + Namespace: ns, + NextCursor: []byte(seq), + }); err != nil { + if errors.Is(err, context.Canceled) { + return nil + } + return connect.NewError(connect.CodeInternal, err) + } + } + + if err := changes.Err(); err != nil { + if !errors.Is(err, context.Canceled) { + return connect.NewError(connect.CodeInternal, err) + } + } + + return nil +} + +func (c *conn) StreamLSN(ctx context.Context, r *connect.Request[adiomv1.StreamLSNRequest], s *connect.ServerStream[adiomv1.StreamLSNResponse]) error { + namespaces := r.Msg.GetNamespaces() + since := string(r.Msg.GetCursor()) + + if len(namespaces) != 1 { + return connect.NewError(connect.CodeInvalidArgument, fmt.Errorf("expected exactly 1 namespace for LSN streaming")) + } + + ns := namespaces[0] + db := c.client.DB(ns) + + params := map[string]interface{}{ + "feed": "continuous", + "heartbeat": 10000, + "timeout": 60000, + } + + if since != "" { + params["since"] = since + } + + changes := db.Changes(ctx, kivik.Params(params)) + defer changes.Close() + + var lsn uint64 + for changes.Next() { + lsn++ + seq := changes.Seq() + + if err := s.Send(&adiomv1.StreamLSNResponse{ + Lsn: lsn, + NextCursor: []byte(seq), + }); err != nil { + if errors.Is(err, context.Canceled) { + return nil + } + return connect.NewError(connect.CodeInternal, err) + } + } + + if err := changes.Err(); err != nil { + if !errors.Is(err, context.Canceled) { + return connect.NewError(connect.CodeInternal, err) + } + } + + return nil +} + +func (c *conn) Teardown() { + if c.client != nil { + _ = c.client.Close() + } +} + +func NewConn(settings ConnectorSettings) (adiomv1connect.ConnectorServiceHandler, error) { + setDefault(&settings.ServerConnectTimeout, 10*time.Second) + setDefault(&settings.PingTimeout, 5*time.Second) + setDefault(&settings.TargetDocCountPerPartition, 50000) + setDefault(&settings.MaxPageSize, 1000) + + dsn, err := convertUriToDSN(settings.Uri) + if err != nil { + return nil, fmt.Errorf("failed to convert URI to DSN: %w", err) + } + + // In Kivik v4, authentication is handled via the DSN URL itself + // (username:password in the URL) - no separate Authenticate call needed + client, err := kivik.New("couch", dsn) + if err != nil { + return nil, fmt.Errorf("failed to create CouchDB client: %w", err) + } + + pingCtx, pingCancel := context.WithTimeout(context.Background(), settings.PingTimeout) + defer pingCancel() + + if _, err := client.Ping(pingCtx); err != nil { + return nil, fmt.Errorf("failed to ping CouchDB server: %w", err) + } + + parsedURL, _ := url.Parse(dsn) + baseURL := fmt.Sprintf("%s://%s", parsedURL.Scheme, parsedURL.Host) + slog.Info(fmt.Sprintf("Connected to CouchDB at %s", baseURL)) + + return &conn{ + client: client, + settings: settings, + dsn: dsn, + }, nil +} diff --git a/connectors/couchdb/conn_test.go b/connectors/couchdb/conn_test.go new file mode 100644 index 00000000..0fcdc3e4 --- /dev/null +++ b/connectors/couchdb/conn_test.go @@ -0,0 +1,223 @@ +/* + * Copyright (C) 2024 Adiom, Inc. + * + * SPDX-License-Identifier: AGPL-3.0-or-later + */ +package couchdb + +import ( + "testing" +) + +// Unit tests - always run + +func TestConvertUriToDSN(t *testing.T) { + tests := []struct { + name string + uri string + want string + wantErr bool + }{ + { + name: "couchdb scheme", + uri: "couchdb://localhost:5984", + want: "http://localhost:5984", + wantErr: false, + }, + { + name: "couchdbs scheme (TLS)", + uri: "couchdbs://localhost:5984", + want: "https://localhost:5984", + wantErr: false, + }, + { + name: "cloudant scheme", + uri: "cloudant://user:pass@account.cloudant.com", + want: "https://user:pass@account.cloudant.com", + wantErr: false, + }, + { + name: "http passthrough", + uri: "http://localhost:5984", + want: "http://localhost:5984", + wantErr: false, + }, + { + name: "https passthrough", + uri: "https://localhost:5984", + want: "https://localhost:5984", + wantErr: false, + }, + { + name: "unsupported scheme", + uri: "mongodb://localhost:27017", + want: "", + wantErr: true, + }, + { + name: "couchdb with auth", + uri: "couchdb://admin:secret@localhost:5984", + want: "http://admin:secret@localhost:5984", + wantErr: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := convertUriToDSN(tt.uri) + if (err != nil) != tt.wantErr { + t.Errorf("convertUriToDSN() error = %v, wantErr %v", err, tt.wantErr) + return + } + if got != tt.want { + t.Errorf("convertUriToDSN() = %v, want %v", got, tt.want) + } + }) + } +} + +func TestEncodeCursor(t *testing.T) { + tests := []struct { + name string + startKey string + endKey string + wantNil bool + }{ + { + name: "both empty", + startKey: "", + endKey: "", + wantNil: true, + }, + { + name: "only startKey", + startKey: "doc1", + endKey: "", + wantNil: false, + }, + { + name: "only endKey", + startKey: "", + endKey: "doc2", + wantNil: false, + }, + { + name: "both keys", + startKey: "doc1", + endKey: "doc2", + wantNil: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := encodeCursor(tt.startKey, tt.endKey) + if tt.wantNil && got != nil { + t.Errorf("encodeCursor() = %v, want nil", got) + } + if !tt.wantNil && got == nil { + t.Errorf("encodeCursor() = nil, want non-nil") + } + }) + } +} + +func TestDecodeCursor(t *testing.T) { + tests := []struct { + name string + cursor []byte + wantStartKey string + wantEndKey string + }{ + { + name: "nil cursor", + cursor: nil, + wantStartKey: "", + wantEndKey: "", + }, + { + name: "empty cursor", + cursor: []byte{}, + wantStartKey: "", + wantEndKey: "", + }, + { + name: "valid cursor with both keys", + cursor: []byte(`{"s":"doc1","e":"doc2"}`), + wantStartKey: "doc1", + wantEndKey: "doc2", + }, + { + name: "valid cursor with only startKey", + cursor: []byte(`{"s":"doc1","e":""}`), + wantStartKey: "doc1", + wantEndKey: "", + }, + { + name: "invalid JSON", + cursor: []byte(`invalid`), + wantStartKey: "", + wantEndKey: "", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + gotStart, gotEnd := decodeCursor(tt.cursor) + if gotStart != tt.wantStartKey { + t.Errorf("decodeCursor() startKey = %v, want %v", gotStart, tt.wantStartKey) + } + if gotEnd != tt.wantEndKey { + t.Errorf("decodeCursor() endKey = %v, want %v", gotEnd, tt.wantEndKey) + } + }) + } +} + +func TestEncodeDecode_Roundtrip(t *testing.T) { + tests := []struct { + startKey string + endKey string + }{ + {"", ""}, + {"doc1", ""}, + {"", "doc2"}, + {"doc1", "doc2"}, + {"special-chars!@#", "unicode-日本語"}, + } + + for _, tt := range tests { + cursor := encodeCursor(tt.startKey, tt.endKey) + gotStart, gotEnd := decodeCursor(cursor) + + if tt.startKey == "" && tt.endKey == "" { + if cursor != nil { + t.Errorf("expected nil cursor for empty keys") + } + continue + } + + if gotStart != tt.startKey { + t.Errorf("roundtrip startKey = %v, want %v", gotStart, tt.startKey) + } + if gotEnd != tt.endKey { + t.Errorf("roundtrip endKey = %v, want %v", gotEnd, tt.endKey) + } + } +} + +func TestGenerateConnectorID(t *testing.T) { + id1 := generateConnectorID("couchdb://localhost:5984") + id2 := generateConnectorID("couchdb://localhost:5984") + id3 := generateConnectorID("couchdb://otherhost:5984") + + if id1 != id2 { + t.Errorf("same connection string should produce same ID: %v != %v", id1, id2) + } + if id1 == id3 { + t.Errorf("different connection strings should produce different IDs") + } + if id1 == "" { + t.Errorf("connector ID should not be empty") + } +} diff --git a/connectors/couchdb/integration_test.go b/connectors/couchdb/integration_test.go new file mode 100644 index 00000000..9b35c0ee --- /dev/null +++ b/connectors/couchdb/integration_test.go @@ -0,0 +1,268 @@ +//go:build external +// +build external + +/* + * Copyright (C) 2024 Adiom, Inc. + * + * SPDX-License-Identifier: AGPL-3.0-or-later + */ +package couchdb + +import ( + "context" + "encoding/json" + "fmt" + "os" + "testing" + "time" + + "connectrpc.com/connect" + adiomv1 "github.com/adiom-data/dsync/gen/adiom/v1" + "github.com/adiom-data/dsync/gen/adiom/v1/adiomv1connect" + pkgtest "github.com/adiom-data/dsync/pkg/test" + "github.com/go-kivik/kivik/v4" + _ "github.com/go-kivik/kivik/v4/couchdb" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/suite" +) + +const ( + CouchDBEnvironmentVariable = "COUCHDB_TEST" +) + +var TestCouchDBConnectionString = os.Getenv(CouchDBEnvironmentVariable) + +func DBString() string { + if r := os.Getenv("COUCHDB_TEST_DB"); r != "" { + return r + } + return "dsync_test" +} + +func TestCouchDBConnectorSuite(t *testing.T) { + if TestCouchDBConnectionString == "" { + t.Skip("COUCHDB_TEST environment variable not set") + } + + dsn, err := convertUriToDSN(TestCouchDBConnectionString) + if err != nil { + t.Fatalf("Failed to convert URI: %v", err) + } + + client, err := kivik.New("couch", dsn) + if err != nil { + t.Fatalf("Failed to create CouchDB client: %v", err) + } + + dbName := DBString() + db := client.DB(dbName) + + tSuite := pkgtest.NewConnectorTestSuite(dbName, func() adiomv1connect.ConnectorServiceClient { + conn, err := NewConn(ConnectorSettings{ + Uri: TestCouchDBConnectionString, + MaxPageSize: 2, + }) + if err != nil { + t.FailNow() + } + return pkgtest.ClientFromHandler(conn) + }, func(ctx context.Context) error { + // Bootstrap: create database and insert test data + _ = client.DestroyDB(ctx, dbName) + if err := client.CreateDB(ctx, dbName); err != nil { + return fmt.Errorf("failed to create database: %w", err) + } + + db = client.DB(dbName) + + // Insert 3 test documents + docs := []map[string]interface{}{ + {"_id": "doc1", "data": "hi"}, + {"_id": "doc2", "data": "hi2"}, + {"_id": "doc3", "data": "hi3"}, + } + for _, doc := range docs { + if _, err := db.Put(ctx, doc["_id"].(string), doc); err != nil { + return fmt.Errorf("failed to insert document: %w", err) + } + } + + return nil + }, func(ctx context.Context) error { + // Insert update for CDC test + doc := map[string]interface{}{"_id": "update_doc", "data": "update"} + if _, err := db.Put(ctx, "update_doc", doc); err != nil { + return fmt.Errorf("failed to insert update document: %w", err) + } + // Small delay to ensure change is propagated + time.Sleep(100 * time.Millisecond) + return nil + }, 2, 3) // 2 pages (with pageSize=2, 3 docs = 2 pages), 3 items total + + tSuite.AssertExists = func(ctx context.Context, a *assert.Assertions, id []*adiomv1.BsonValue, exists bool) error { + docID := string(id[0].GetData()) + row := db.Get(ctx, docID) + var doc map[string]interface{} + err := row.ScanDoc(&doc) + if exists { + a.NoError(err, "Document should exist") + } else { + a.Error(err, "Document should not exist") + } + return nil + } + + // CouchDB uses JSON, not BSON - skip the BSON-specific write tests + tSuite.SkipWriteUpdatesTest = true + + suite.Run(t, tSuite) +} + +func TestCouchDBWriteDataJSON(t *testing.T) { + if TestCouchDBConnectionString == "" { + t.Skip("COUCHDB_TEST environment variable not set") + } + + ctx := context.Background() + + dsn, err := convertUriToDSN(TestCouchDBConnectionString) + assert.NoError(t, err) + + client, err := kivik.New("couch", dsn) + assert.NoError(t, err) + + dbName := "dsync_write_test" + _ = client.DestroyDB(ctx, dbName) + err = client.CreateDB(ctx, dbName) + assert.NoError(t, err) + defer client.DestroyDB(ctx, dbName) + + conn, err := NewConn(ConnectorSettings{Uri: TestCouchDBConnectionString}) + assert.NoError(t, err) + + c := pkgtest.ClientFromHandler(conn) + + // Test WriteData with JSON + doc := map[string]interface{}{ + "_id": "json_test_doc", + "name": "test", + "value": 123, + } + jsonBytes, _ := json.Marshal(doc) + + _, err = c.WriteData(ctx, connect.NewRequest(&adiomv1.WriteDataRequest{ + Namespace: dbName, + Data: [][]byte{jsonBytes}, + Type: adiomv1.DataType_DATA_TYPE_JSON_ID, + })) + assert.NoError(t, err) + + // Verify document was written + db := client.DB(dbName) + row := db.Get(ctx, "json_test_doc") + var result map[string]interface{} + err = row.ScanDoc(&result) + assert.NoError(t, err) + assert.Equal(t, "test", result["name"]) + assert.Equal(t, float64(123), result["value"]) +} + +func TestCouchDBWriteUpdatesJSON(t *testing.T) { + if TestCouchDBConnectionString == "" { + t.Skip("COUCHDB_TEST environment variable not set") + } + + ctx := context.Background() + + dsn, err := convertUriToDSN(TestCouchDBConnectionString) + assert.NoError(t, err) + + client, err := kivik.New("couch", dsn) + assert.NoError(t, err) + + dbName := "dsync_updates_test" + _ = client.DestroyDB(ctx, dbName) + err = client.CreateDB(ctx, dbName) + assert.NoError(t, err) + defer client.DestroyDB(ctx, dbName) + + conn, err := NewConn(ConnectorSettings{Uri: TestCouchDBConnectionString}) + assert.NoError(t, err) + + c := pkgtest.ClientFromHandler(conn) + + // Test INSERT + doc := map[string]interface{}{ + "_id": "update_test_doc", + "name": "initial", + } + jsonBytes, _ := json.Marshal(doc) + + _, err = c.WriteUpdates(ctx, connect.NewRequest(&adiomv1.WriteUpdatesRequest{ + Namespace: dbName, + Updates: []*adiomv1.Update{{ + Id: []*adiomv1.BsonValue{{ + Data: []byte("update_test_doc"), + Type: 2, // string + Name: "_id", + }}, + Type: adiomv1.UpdateType_UPDATE_TYPE_INSERT, + Data: jsonBytes, + }}, + Type: adiomv1.DataType_DATA_TYPE_JSON_ID, + })) + assert.NoError(t, err) + + // Verify insert + db := client.DB(dbName) + row := db.Get(ctx, "update_test_doc") + var result map[string]interface{} + err = row.ScanDoc(&result) + assert.NoError(t, err) + assert.Equal(t, "initial", result["name"]) + + // Test UPDATE + doc["name"] = "updated" + jsonBytes, _ = json.Marshal(doc) + + _, err = c.WriteUpdates(ctx, connect.NewRequest(&adiomv1.WriteUpdatesRequest{ + Namespace: dbName, + Updates: []*adiomv1.Update{{ + Id: []*adiomv1.BsonValue{{ + Data: []byte("update_test_doc"), + Type: 2, + Name: "_id", + }}, + Type: adiomv1.UpdateType_UPDATE_TYPE_UPDATE, + Data: jsonBytes, + }}, + Type: adiomv1.DataType_DATA_TYPE_JSON_ID, + })) + assert.NoError(t, err) + + // Verify update + row = db.Get(ctx, "update_test_doc") + err = row.ScanDoc(&result) + assert.NoError(t, err) + assert.Equal(t, "updated", result["name"]) + + // Test DELETE + _, err = c.WriteUpdates(ctx, connect.NewRequest(&adiomv1.WriteUpdatesRequest{ + Namespace: dbName, + Updates: []*adiomv1.Update{{ + Id: []*adiomv1.BsonValue{{ + Data: []byte("update_test_doc"), + Type: 2, + Name: "_id", + }}, + Type: adiomv1.UpdateType_UPDATE_TYPE_DELETE, + }}, + Type: adiomv1.DataType_DATA_TYPE_JSON_ID, + })) + assert.NoError(t, err) + + // Verify delete + row = db.Get(ctx, "update_test_doc") + err = row.ScanDoc(&result) + assert.Error(t, err, "Document should be deleted") +} diff --git a/go.mod b/go.mod index 91b45aba..532a271a 100644 --- a/go.mod +++ b/go.mod @@ -21,6 +21,7 @@ require ( github.com/cenkalti/backoff/v4 v4.3.0 github.com/cespare/xxhash v1.1.0 github.com/gdamore/tcell/v2 v2.8.1 + github.com/go-kivik/kivik/v4 v4.5.2 github.com/jackc/pglogrepl v0.0.0-20250509230407-a9884f6bd75a github.com/jackc/pgx-shopspring-decimal v0.0.0-20220624020537-1d36b5a1853e github.com/jackc/pgx/v5 v5.7.5 diff --git a/go.sum b/go.sum index 2be3d933..8abdf076 100644 --- a/go.sum +++ b/go.sum @@ -21,6 +21,8 @@ github.com/BurntSushi/toml v1.5.0 h1:W5quZX/G/csjUnuI8SUYlsHs9M38FC7znL0lIO+DvMg github.com/BurntSushi/toml v1.5.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= github.com/IBM/sarama v1.46.3 h1:njRsX6jNlnR+ClJ8XmkO+CM4unbrNr/2vB5KK6UA+IE= github.com/IBM/sarama v1.46.3/go.mod h1:GTUYiF9DMOZVe3FwyGT+dtSPceGFIgA+sPc5u6CBwko= +github.com/Masterminds/semver/v3 v3.2.1 h1:RN9w6+7QoMeJVGyfmbcgs28Br8cvmnucEXnY0rYXWg0= +github.com/Masterminds/semver/v3 v3.2.1/go.mod h1:qvl/7zhW3nngYb5+80sSMF+FG2BjYrf8m9wsX0PNOMQ= github.com/OneOfOne/xxhash v1.2.2 h1:KMrpdQIwFcEqXDklaen+P1axHaj9BSKzvpUUfnHldSE= github.com/OneOfOne/xxhash v1.2.2/go.mod h1:HSdplMjZKSmBqAxg5vPj2TmRDmfkzw+cTzAElWljhcU= github.com/PuerkitoBio/purell v1.1.1/go.mod h1:c11w/QuzBsJSee3cPx9rAFu61PvFxuPbtSwDGJws/X0= @@ -105,6 +107,8 @@ github.com/gdamore/encoding v1.0.1 h1:YzKZckdBL6jVt2Gc+5p82qhrGiqMdG/eNs6Wy0u3Uh github.com/gdamore/encoding v1.0.1/go.mod h1:0Z0cMFinngz9kS1QfMjCP8TY7em3bZYeeklsSDPivEo= github.com/gdamore/tcell/v2 v2.8.1 h1:KPNxyqclpWpWQlPLx6Xui1pMk8S+7+R37h3g07997NU= github.com/gdamore/tcell/v2 v2.8.1/go.mod h1:bj8ori1BG3OYMjmb3IklZVWfZUJ1UBQt9JXrOCOhGWw= +github.com/go-kivik/kivik/v4 v4.5.2 h1:Xi6QyjscrWrSRQEEW/25vaWyh20ZCg3LhTm6AqFRZCc= +github.com/go-kivik/kivik/v4 v4.5.2/go.mod h1:5YlQJZim4qvaJ3T0fCAS6U4oaN4hzXK6CVY9nvN4Phg= github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= @@ -195,11 +199,15 @@ github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX github.com/google/uuid v1.1.1/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= 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/gopherjs/gopherjs v1.20.1 h1:22uLWFvVcxhJ+j3dJ99NNfwGyHynxCmjhYsrcwqbY60= +github.com/gopherjs/gopherjs v1.20.1/go.mod h1:h+FTmmLgbXMmmtuZFp9bUqXciN429Wx0sJEJuMnpyfM= github.com/gorilla/securecookie v1.1.1/go.mod h1:ra0sb63/xPlUeL+yeDciTfxMRAA+MP+HVt/4epWDjd4= github.com/gorilla/sessions v1.2.1/go.mod h1:dk2InVEVJ0sfLlnXv9EAgkf6ecYs/i80K/zI+bUmuGM= github.com/hashicorp/go-uuid v1.0.2/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro= github.com/hashicorp/go-uuid v1.0.3 h1:2gKiV6YVmrJ1i2CKKa9obLvRieoRGviZFL26PcT/Co8= github.com/hashicorp/go-uuid v1.0.3/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro= +github.com/icza/dyno v0.0.0-20230330125955-09f820a8d9c0 h1:nHoRIX8iXob3Y2kdt9KsjyIb7iApSvb3vgsd93xb5Ow= +github.com/icza/dyno v0.0.0-20230330125955-09f820a8d9c0/go.mod h1:c1tRKs5Tx7E2+uHGSyyncziFjvGpgv4H2HrqXeUQ/Uk= github.com/inconshreveable/mousetrap v1.0.0/go.mod h1:PxqpIevigyE2G7u3NXJIT2ANytuPF1OarO4DADm73n8= github.com/jackc/pgio v1.0.0 h1:g12B9UwVnzGhueNavwioyEEpAmqMe1E/BN9ES+8ovkE= github.com/jackc/pgio v1.0.0/go.mod h1:oP+2QK2wFfUWgr+gxjoBH9KGBb31Eio69xUb0w5bYf8= @@ -391,6 +399,8 @@ github.com/youmark/pkcs8 v0.0.0-20181117223130-1be2e3e5546d/go.mod h1:rHwXgn7Jul github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 h1:ilQV1hzziu+LLM3zUTJ0trRztfwgjqKnBWNtSRkbmwM= github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78/go.mod h1:aL8wCCfTfSfmXjznFBSZNN13rSJjlIOI1fUNAtF7rmI= github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= +gitlab.com/flimzy/testy v0.15.0 h1:69TL12IpxqGUyL8NuRV3Z5OhIDszXLNqLtfBDhOV3ys= +gitlab.com/flimzy/testy v0.15.0/go.mod h1:KbAJWCwB++0hEFzeeQRbC7vdZYP/yEha94s4X1wVFrw= go.akshayshah.org/attest v1.0.2 h1:qOv9PXCG2mwnph3g0I3yZj0rLAwLyUITs8nhxP+wS44= go.akshayshah.org/attest v1.0.2/go.mod h1:PnWzcW5j9dkyGwTlBmUsYpPnHG0AUPrs1RQ+HrldWO0= go.akshayshah.org/memhttp v0.1.0 h1:Enf7JeZnm+A8iRur0FYvs4ZjWa1VVMc2gG4EirG+aNE= diff --git a/internal/app/options/connectorflags.go b/internal/app/options/connectorflags.go index 3ffd5198..ee8d1857 100644 --- a/internal/app/options/connectorflags.go +++ b/internal/app/options/connectorflags.go @@ -13,6 +13,7 @@ import ( "github.com/IBM/sarama" "github.com/adiom-data/dsync/connectors/airbyte" "github.com/adiom-data/dsync/connectors/cosmos" + "github.com/adiom-data/dsync/connectors/couchdb" "github.com/adiom-data/dsync/connectors/dynamodb" fileconnector "github.com/adiom-data/dsync/connectors/file" "github.com/adiom-data/dsync/connectors/kafka" @@ -518,6 +519,21 @@ func GetRegisteredConnectors() []RegisteredConnector { })(args, as) }, }, + { + Name: "CouchDB", + IsConnector: func(s string) bool { + lower := strings.ToLower(s) + return strings.HasPrefix(lower, "couchdb://") || + strings.HasPrefix(lower, "couchdbs://") || + strings.HasPrefix(lower, "cloudant://") + }, + Create: func(args []string, as AdditionalSettings) (adiomv1connect.ConnectorServiceHandler, []string, error) { + settings := couchdb.ConnectorSettings{Uri: args[0]} + return CreateHelper("CouchDB", "couchdb://user:pass@host:port OR cloudant://user:pass@account.cloudant.com [options]", CouchDBFlags(&settings), func(_ *cli.Context, _ []string, _ AdditionalSettings) (adiomv1connect.ConnectorServiceHandler, error) { + return couchdb.NewConn(settings) + })(args, as) + }, + }, { Name: "weaviate", IsConnector: func(s string) bool { @@ -1062,6 +1078,43 @@ var postgresSettingsDefault = postgres.PostgresSettings{ EnableReplicaMode: true, } +func CouchDBFlags(settings *couchdb.ConnectorSettings) []cli.Flag { + return []cli.Flag{ + altsrc.NewDurationFlag(&cli.DurationFlag{ + Name: "server-timeout", + Usage: "Connection timeout for the CouchDB server", + Destination: &settings.ServerConnectTimeout, + }), + altsrc.NewDurationFlag(&cli.DurationFlag{ + Name: "ping-timeout", + Usage: "Ping timeout for the CouchDB server", + Destination: &settings.PingTimeout, + }), + altsrc.NewIntFlag(&cli.IntFlag{ + Name: "writer-batch-size", + Usage: "Maximum batch size for bulk writes", + Destination: &settings.WriterMaxBatchSize, + }), + altsrc.NewInt64Flag(&cli.Int64Flag{ + Name: "doc-partition", + Usage: "Target number of documents per partition", + Value: 50000, + Destination: &settings.TargetDocCountPerPartition, + }), + altsrc.NewIntFlag(&cli.IntFlag{ + Name: "max-page-size", + Usage: "Maximum page size for listing data", + Value: 1000, + Destination: &settings.MaxPageSize, + }), + altsrc.NewBoolFlag(&cli.BoolFlag{ + Name: "include-system-dbs", + Usage: "Include system databases (_users, _replicator, etc.)", + Destination: &settings.IncludeSystemDbs, + }), + } +} + func PostgresFlags(settings *postgres.PostgresSettings) []cli.Flag { return []cli.Flag{ altsrc.NewIntFlag(&cli.IntFlag{ From 708393959e8f687f22a9769ac14b7455f99099a4 Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Sat, 28 Mar 2026 08:14:50 -0700 Subject: [PATCH 2/9] fix --- connectors/couchdb/conn.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index 63b4a9f9..a18878f4 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -466,9 +466,11 @@ func (c *conn) StreamUpdates(ctx context.Context, r *connect.Request[adiomv1.Str "heartbeat": 1000, // 1 second heartbeat } - // Only set 'since' if we have a valid cursor + // Only set 'since' if we have a valid cursor, otherwise start from now if since != "" { params["since"] = since + } else { + params["since"] = "now" } changes := db.Changes(ctx, kivik.Params(params)) @@ -557,6 +559,8 @@ func (c *conn) StreamLSN(ctx context.Context, r *connect.Request[adiomv1.StreamL if since != "" { params["since"] = since + } else { + params["since"] = "now" } changes := db.Changes(ctx, kivik.Params(params)) From d0e27b5ac7b6569554df0ee8aff8eee6931924af Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Sat, 28 Mar 2026 08:15:10 -0700 Subject: [PATCH 3/9] don't treat everything as bson when logging json --- connectors/null/connector.go | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/connectors/null/connector.go b/connectors/null/connector.go index f6cb27e8..e4f71239 100644 --- a/connectors/null/connector.go +++ b/connectors/null/connector.go @@ -117,8 +117,12 @@ func (c *conn) WriteUpdates(ctx context.Context, r *connect.Request[adiomv1.Writ var idOutput []any for _, id := range updates.GetId() { var v any - if err := bson.UnmarshalValue(bsontype.Type(id.GetType()), id.GetData(), &v); err != nil { - return nil, connect.NewError(connect.CodeInternal, err) + if r.Msg.GetType() == adiomv1.DataType_DATA_TYPE_JSON_ID { + v = string(id.GetData()) + } else { + if err := bson.UnmarshalValue(bsontype.Type(id.GetType()), id.GetData(), &v); err != nil { + return nil, connect.NewError(connect.CodeInternal, err) + } } idOutput = append(idOutput, v) } From 7889e24c0a2d43439f272fe1ad090be803ff2c96 Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Sun, 29 Mar 2026 09:06:07 -0700 Subject: [PATCH 4/9] fixes --- connectors/couchdb/conn.go | 113 +++++++++++++++++++++++++++---------- 1 file changed, 84 insertions(+), 29 deletions(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index a18878f4..2c3efbc8 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -164,10 +164,10 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene docsPerPartition := count / numPartitions rows := db.AllDocs(ctx, kivik.Params(map[string]interface{}{ - "limit": numPartitions - 1, - "skip": docsPerPartition, + "include_docs": false, + "limit": numPartitions - 1, + "skip": docsPerPartition, })) - defer rows.Close() var boundaries []string for rows.Next() { @@ -180,6 +180,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene break } } + rows.Close() var prevKey string for i, boundary := range boundaries { @@ -237,7 +238,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene if err == nil && meta != nil { lastSeq = meta.LastSeq } - info.Close() + _ = info.Close() } updatesPartitions = append(updatesPartitions, &adiomv1.UpdatesPartition{ @@ -380,17 +381,29 @@ func (c *conn) WriteData(ctx context.Context, r *connect.Request[adiomv1.WriteDa docs = append(docs, doc) if c.settings.WriterMaxBatchSize > 0 && len(docs) >= c.settings.WriterMaxBatchSize { - if _, err := db.BulkDocs(ctx, docs); err != nil { + results, err := db.BulkDocs(ctx, docs) + if err != nil { return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk insert documents: %w", err)) } + for _, result := range results { + if result.Error != nil { + slog.Error(fmt.Sprintf("Failed to insert document %s: %v", result.ID, result.Error)) + } + } docs = nil } } if len(docs) > 0 { - if _, err := db.BulkDocs(ctx, docs); err != nil { + results, err := db.BulkDocs(ctx, docs) + if err != nil { return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk insert documents: %w", err)) } + for _, result := range results { + if result.Error != nil { + slog.Error(fmt.Sprintf("Failed to insert document %s: %v", result.ID, result.Error)) + } + } } return connect.NewResponse(&adiomv1.WriteDataResponse{}), nil @@ -398,8 +411,50 @@ func (c *conn) WriteData(ctx context.Context, r *connect.Request[adiomv1.WriteDa func (c *conn) WriteUpdates(ctx context.Context, r *connect.Request[adiomv1.WriteUpdatesRequest]) (*connect.Response[adiomv1.WriteUpdatesResponse], error) { db := c.client.DB(r.Msg.GetNamespace()) + updates := r.Msg.GetUpdates() + + if len(updates) == 0 { + return connect.NewResponse(&adiomv1.WriteUpdatesResponse{}), nil + } - for _, update := range r.Msg.GetUpdates() { + // Collect all document IDs to fetch revisions in batch + var docIDs []string + docIDSet := make(map[string]struct{}) + for _, update := range updates { + if len(update.GetId()) == 0 { + continue + } + docID := string(update.GetId()[0].GetData()) + if _, exists := docIDSet[docID]; !exists { + docIDs = append(docIDs, docID) + docIDSet[docID] = struct{}{} + } + } + + // Batch fetch all existing revisions using _all_docs with keys + revMap := make(map[string]string) + if len(docIDs) > 0 { + rows := db.AllDocs(ctx, kivik.Params(map[string]interface{}{ + "keys": docIDs, + })) + for rows.Next() { + id, err := rows.ID() + if err != nil { + continue + } + var value struct { + Rev string `json:"rev"` + } + if err := rows.ScanValue(&value); err == nil && value.Rev != "" { + revMap[id] = value.Rev + } + } + rows.Close() + } + + // Prepare batch operations + var upsertDocs []interface{} + for _, update := range updates { if len(update.GetId()) == 0 { continue } @@ -414,34 +469,34 @@ func (c *conn) WriteUpdates(ctx context.Context, r *connect.Request[adiomv1.Writ continue } - var existingRev string - result := db.Get(ctx, docID) - var existing map[string]interface{} - if err := result.ScanDoc(&existing); err == nil { - if rev, ok := existing["_rev"].(string); ok { - existingRev = rev - } - } - - if existingRev != "" { - doc["_rev"] = existingRev + if rev, exists := revMap[docID]; exists { + doc["_rev"] = rev } else { delete(doc, "_rev") } + upsertDocs = append(upsertDocs, doc) - if _, err := db.Put(ctx, docID, doc); err != nil { - slog.Error(fmt.Sprintf("Failed to upsert document %s: %v", docID, err)) + case adiomv1.UpdateType_UPDATE_TYPE_DELETE: + if rev, exists := revMap[docID]; exists { + upsertDocs = append(upsertDocs, map[string]interface{}{ + "_id": docID, + "_rev": rev, + "_deleted": true, + }) } + } + } - case adiomv1.UpdateType_UPDATE_TYPE_DELETE: - result := db.Get(ctx, docID) - var existing map[string]interface{} - if err := result.ScanDoc(&existing); err == nil { - if rev, ok := existing["_rev"].(string); ok { - if _, err := db.Delete(ctx, docID, rev); err != nil { - slog.Error(fmt.Sprintf("Failed to delete document %s: %v", docID, err)) - } - } + // Execute bulk operation + if len(upsertDocs) > 0 { + results, err := db.BulkDocs(ctx, upsertDocs) + if err != nil { + return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk write documents: %w", err)) + } + // Log any individual document errors + for _, result := range results { + if result.Error != nil { + slog.Error(fmt.Sprintf("Failed to write document %s: %v", result.ID, result.Error)) } } } From ce82a4428be28450f44da33af3c127d66086fdf0 Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Sun, 29 Mar 2026 09:15:19 -0700 Subject: [PATCH 5/9] writedata optimization for conflicts --- connectors/couchdb/conn.go | 103 ++++++++++++++++++++++++++++++------- 1 file changed, 84 insertions(+), 19 deletions(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index 2c3efbc8..0dd53e88 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -368,10 +368,89 @@ func (c *conn) ListData(ctx context.Context, r *connect.Request[adiomv1.ListData }), nil } +func (c *conn) writeDataBatch(ctx context.Context, db *kivik.DB, docs []map[string]interface{}) error { + if len(docs) == 0 { + return nil + } + + // Convert to interface slice for BulkDocs + bulkDocs := make([]interface{}, len(docs)) + for i, doc := range docs { + bulkDocs[i] = doc + } + + // First attempt: optimistic insert + results, err := db.BulkDocs(ctx, bulkDocs) + if err != nil { + return fmt.Errorf("failed to bulk insert documents: %w", err) + } + + // Collect conflicts (409 errors) + var conflictIDs []string + conflictDocs := make(map[string]map[string]interface{}) + for i, result := range results { + if result.Error != nil { + if kivik.HTTPStatus(result.Error) == 409 { + docID, ok := docs[i]["_id"].(string) + if ok { + conflictIDs = append(conflictIDs, docID) + conflictDocs[docID] = docs[i] + } + } else { + slog.Error(fmt.Sprintf("Failed to insert document %s: %v", result.ID, result.Error)) + } + } + } + + // Retry conflicts with revisions + if len(conflictIDs) > 0 { + rows := db.AllDocs(ctx, kivik.Params(map[string]interface{}{ + "keys": conflictIDs, + })) + for rows.Next() { + id, err := rows.ID() + if err != nil { + continue + } + var value struct { + Rev string `json:"rev"` + } + if err := rows.ScanValue(&value); err == nil && value.Rev != "" { + if doc, exists := conflictDocs[id]; exists { + doc["_rev"] = value.Rev + } + } + } + rows.Close() + + // Retry with revisions + var retryDocs []interface{} + for _, doc := range conflictDocs { + if _, hasRev := doc["_rev"]; hasRev { + retryDocs = append(retryDocs, doc) + } + } + + if len(retryDocs) > 0 { + retryResults, err := db.BulkDocs(ctx, retryDocs) + if err != nil { + return fmt.Errorf("failed to bulk upsert conflicting documents: %w", err) + } + for _, result := range retryResults { + if result.Error != nil { + slog.Error(fmt.Sprintf("Failed to upsert document %s: %v", result.ID, result.Error)) + } + } + } + } + + return nil +} + func (c *conn) WriteData(ctx context.Context, r *connect.Request[adiomv1.WriteDataRequest]) (*connect.Response[adiomv1.WriteDataResponse], error) { db := c.client.DB(r.Msg.GetNamespace()) - var docs []interface{} + var docs []map[string]interface{} for _, data := range r.Msg.GetData() { var doc map[string]interface{} if err := json.Unmarshal(data, &doc); err != nil { @@ -381,29 +460,15 @@ func (c *conn) WriteData(ctx context.Context, r *connect.Request[adiomv1.WriteDa docs = append(docs, doc) if c.settings.WriterMaxBatchSize > 0 && len(docs) >= c.settings.WriterMaxBatchSize { - results, err := db.BulkDocs(ctx, docs) - if err != nil { - return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk insert documents: %w", err)) - } - for _, result := range results { - if result.Error != nil { - slog.Error(fmt.Sprintf("Failed to insert document %s: %v", result.ID, result.Error)) - } + if err := c.writeDataBatch(ctx, db, docs); err != nil { + return nil, connect.NewError(connect.CodeInternal, err) } docs = nil } } - if len(docs) > 0 { - results, err := db.BulkDocs(ctx, docs) - if err != nil { - return nil, connect.NewError(connect.CodeInternal, fmt.Errorf("failed to bulk insert documents: %w", err)) - } - for _, result := range results { - if result.Error != nil { - slog.Error(fmt.Sprintf("Failed to insert document %s: %v", result.ID, result.Error)) - } - } + if err := c.writeDataBatch(ctx, db, docs); err != nil { + return nil, connect.NewError(connect.CodeInternal, err) } return connect.NewResponse(&adiomv1.WriteDataResponse{}), nil From 8a3d87998b65c40a5626cabcb68765b587429919 Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Sun, 29 Mar 2026 09:45:32 -0700 Subject: [PATCH 6/9] fix gosec issues --- connectors/couchdb/conn.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index 0dd53e88..f7655d22 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -180,7 +180,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene break } } - rows.Close() + _ = rows.Close() var prevKey string for i, boundary := range boundaries { @@ -225,7 +225,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene if err := changes.Err(); err != nil { slog.Warn(fmt.Sprintf("Failed to get changes for database %s: %v", ns, err)) } - changes.Close() + _ = changes.Close() if lastSeq == "" { info := db.Changes(ctx, kivik.Params(map[string]interface{}{ @@ -421,7 +421,7 @@ func (c *conn) writeDataBatch(ctx context.Context, db *kivik.DB, docs []map[stri } } } - rows.Close() + _ = rows.Close() // Retry with revisions var retryDocs []interface{} @@ -514,7 +514,7 @@ func (c *conn) WriteUpdates(ctx context.Context, r *connect.Request[adiomv1.Writ revMap[id] = value.Rev } } - rows.Close() + _ = rows.Close() } // Prepare batch operations From 5edc1735b2a8cd7c3686b4a6310290163d468c22 Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Mon, 30 Mar 2026 10:54:49 -0700 Subject: [PATCH 7/9] fix planning --- connectors/couchdb/conn.go | 39 ++++++++++++++++++++++++-------------- 1 file changed, 25 insertions(+), 14 deletions(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index f7655d22..d47f3b90 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -163,24 +163,35 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene numPartitions := (count + c.settings.TargetDocCountPerPartition - 1) / c.settings.TargetDocCountPerPartition docsPerPartition := count / numPartitions - rows := db.AllDocs(ctx, kivik.Params(map[string]interface{}{ - "include_docs": false, - "limit": numPartitions - 1, - "skip": docsPerPartition, - })) - + // Find partition boundaries using chained startkey queries + // Each query skips docsPerPartition from the previous boundary var boundaries []string - for rows.Next() { - id, err := rows.ID() - if err != nil { - continue + var startKey string + + for i := 1; i < int(numPartitions); i++ { + params := map[string]interface{}{ + "include_docs": false, + "limit": 1, + "skip": docsPerPartition, + } + if startKey != "" { + params["startkey"] = startKey } - boundaries = append(boundaries, id) - if len(boundaries) >= int(numPartitions)-1 { - break + + rows := db.AllDocs(ctx, kivik.Params(params)) + if rows.Next() { + id, err := rows.ID() + if err == nil { + boundaries = append(boundaries, id) + startKey = id + } + } + rows.Close() + + if len(boundaries) < i { + break // No more docs available } } - _ = rows.Close() var prevKey string for i, boundary := range boundaries { From 4e54c72540b31f033c11e80ddcef539227b2f9fe Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Mon, 30 Mar 2026 11:15:25 -0700 Subject: [PATCH 8/9] fix listdata --- connectors/couchdb/conn.go | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index d47f3b90..9bb1b3ce 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -310,21 +310,23 @@ func (c *conn) ListData(ctx context.Context, r *connect.Request[adiomv1.ListData "limit": pageSize + 1, // fetch one extra to detect if there's a next page } + // Always apply partition end bound + startKey, endKey := decodeCursor(partition.GetCursor()) + if endKey != "" { + params["endkey"] = endKey + } + pageCursor := r.Msg.GetCursor() if len(pageCursor) > 0 { // Page cursor from previous call - start after this key params["startkey"] = string(pageCursor) params["skip"] = 1 // skip the document we already returned } else { - // Initial call - use partition cursor if present - startKey, endKey := decodeCursor(partition.GetCursor()) + // Initial call - use partition start key if present if startKey != "" { params["startkey"] = startKey params["skip"] = 1 } - if endKey != "" { - params["endkey"] = endKey - } } rows := db.AllDocs(ctx, kivik.Params(params)) From e93cf56de48828554fbeb34ccbdd8a136b67d77e Mon Sep 17 00:00:00 2001 From: Alexander Komyagin Date: Mon, 30 Mar 2026 11:20:08 -0700 Subject: [PATCH 9/9] gosec fix --- connectors/couchdb/conn.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connectors/couchdb/conn.go b/connectors/couchdb/conn.go index 9bb1b3ce..343cb4e2 100644 --- a/connectors/couchdb/conn.go +++ b/connectors/couchdb/conn.go @@ -186,7 +186,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene startKey = id } } - rows.Close() + _ = rows.Close() if len(boundaries) < i { break // No more docs available