Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion examples/simple_plugin/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ require (
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/cloudquery/cloudquery-api-go v1.14.13 // indirect
github.com/cloudquery/codegen v0.4.1 // indirect
github.com/cloudquery/plugin-pb-go v1.27.23 // indirect
github.com/cloudquery/plugin-pb-go v1.27.24-0.20261002091754-b8dcd1cbaee9 // indirect
github.com/cloudquery/plugin-sdk/v2 v2.7.0 // indirect
github.com/getsentry/sentry-go v0.49.0 // indirect
github.com/ghodss/yaml v1.0.0 // indirect
Expand Down
4 changes: 2 additions & 2 deletions examples/simple_plugin/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,8 @@ github.com/cloudquery/cloudquery-api-go v1.14.13 h1:+lu1mLKqVSwrc02eqf9vfqGj9c0R
github.com/cloudquery/cloudquery-api-go v1.14.13/go.mod h1:u4uzOBEss9hJZco/tULIZuObqiBD748kejlxZ2AqOxs=
github.com/cloudquery/codegen v0.4.1 h1:c9D18N925tUvnDeGHIl3JWKj37TyII9daHufkf8hU+Y=
github.com/cloudquery/codegen v0.4.1/go.mod h1:QWIOD6R1aCa+YM+th+9Qt9lZw+ztdJR9JDEMLWyazwM=
github.com/cloudquery/plugin-pb-go v1.27.23 h1:X08b+1rB1PKw4s+XFVlUUitTO0NrQ2oe9UHjJ/afSiA=
github.com/cloudquery/plugin-pb-go v1.27.23/go.mod h1:ZUJgTp6qCuurc/c8qnDe/U8eqQpWfDxTBsd/IK8DUeY=
github.com/cloudquery/plugin-pb-go v1.27.24-0.20261002091754-b8dcd1cbaee9 h1:mFYZDvEri1MMEWEqt5a9/FTXE8n/XXEOvs6HqzDOPo0=
github.com/cloudquery/plugin-pb-go v1.27.24-0.20261002091754-b8dcd1cbaee9/go.mod h1:06VOdOlk5Y64CahAxQnmGy6l7PXJYm3SAwNxpEWPwn4=
github.com/cloudquery/plugin-sdk/v2 v2.7.0 h1:hRXsdEiaOxJtsn/wZMFQC9/jPfU1MeMK3KF+gPGqm7U=
github.com/cloudquery/plugin-sdk/v2 v2.7.0/go.mod h1:pAX6ojIW99b/Vg4CkhnsGkRIzNaVEceYMR+Bdit73ug=
github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g=
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ require (
github.com/bradleyjkemp/cupaloy/v2 v2.8.0
github.com/cloudquery/cloudquery-api-go v1.14.13
github.com/cloudquery/codegen v0.4.1
github.com/cloudquery/plugin-pb-go v1.27.23
github.com/cloudquery/plugin-pb-go v1.27.24-0.20261002091754-b8dcd1cbaee9
github.com/cloudquery/plugin-sdk/v2 v2.7.0
github.com/getsentry/sentry-go v0.49.0
github.com/goccy/go-json v0.10.6
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,8 @@ github.com/cloudquery/codegen v0.4.1 h1:c9D18N925tUvnDeGHIl3JWKj37TyII9daHufkf8h
github.com/cloudquery/codegen v0.4.1/go.mod h1:QWIOD6R1aCa+YM+th+9Qt9lZw+ztdJR9JDEMLWyazwM=
github.com/cloudquery/jsonschema v0.0.0-20260703174721-45e7e20e0ed8 h1:s7B+c57yTtVL8Zhcebae5poFInJwTftuMakj625CRbw=
github.com/cloudquery/jsonschema v0.0.0-20260703174721-45e7e20e0ed8/go.mod h1:KMcD1TlufeD5r5DmbYvYWs+cyULzIVuKkd7jpoKWVjk=
github.com/cloudquery/plugin-pb-go v1.27.23 h1:X08b+1rB1PKw4s+XFVlUUitTO0NrQ2oe9UHjJ/afSiA=
github.com/cloudquery/plugin-pb-go v1.27.23/go.mod h1:ZUJgTp6qCuurc/c8qnDe/U8eqQpWfDxTBsd/IK8DUeY=
github.com/cloudquery/plugin-pb-go v1.27.24-0.20261002091754-b8dcd1cbaee9 h1:mFYZDvEri1MMEWEqt5a9/FTXE8n/XXEOvs6HqzDOPo0=
github.com/cloudquery/plugin-pb-go v1.27.24-0.20261002091754-b8dcd1cbaee9/go.mod h1:06VOdOlk5Y64CahAxQnmGy6l7PXJYm3SAwNxpEWPwn4=
github.com/cloudquery/plugin-sdk/v2 v2.7.0 h1:hRXsdEiaOxJtsn/wZMFQC9/jPfU1MeMK3KF+gPGqm7U=
github.com/cloudquery/plugin-sdk/v2 v2.7.0/go.mod h1:pAX6ojIW99b/Vg4CkhnsGkRIzNaVEceYMR+Bdit73ug=
github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g=
Expand Down
85 changes: 85 additions & 0 deletions internal/servers/plugin/v3/assess.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
package plugin

import (
"context"
"fmt"

pb "github.com/cloudquery/plugin-pb-go/pb/plugin/v3"
"github.com/cloudquery/plugin-sdk/v4/plugin"
"github.com/cloudquery/plugin-sdk/v4/schema"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)

func (s *Server) AssessTables(ctx context.Context, req *pb.AssessTables_Request) (*pb.AssessTables_Response, error) {
tables := make([]plugin.TablePair, len(req.Tables))
for i, pair := range req.Tables {
if len(pair.OldTable) == 0 && len(pair.NewTable) == 0 {
return nil, status.Errorf(codes.InvalidArgument, "table pair %d has neither an old nor a new table", i)
}
var err error
if tables[i].Old, err = tableFromBytes(pair.OldTable); err != nil {
return nil, status.Errorf(codes.InvalidArgument, "failed to decode old table: %v", err)
}
if tables[i].New, err = tableFromBytes(pair.NewTable); err != nil {
return nil, status.Errorf(codes.InvalidArgument, "failed to decode new table: %v", err)
}
}
findings, err := s.Plugin.AssessTables(ctx, tables, plugin.AssessOptions{MigrateForce: req.MigrateForce})
if err != nil {
return nil, status.Errorf(codes.Internal, "failed to assess tables: %v", err)
}
resp := &pb.AssessTables_Response{Tables: make([]*pb.AssessTables_TableFinding, len(findings))}
for i, f := range findings {
resp.Tables[i] = tableFindingToPB(f)
}
return resp, nil
}

func tableFromBytes(b []byte) (*schema.Table, error) {
if len(b) == 0 {
return nil, nil
}
sc, err := pb.NewSchemaFromBytes(b)
if err != nil {
return nil, err
}
table, err := schema.NewTableFromArrowSchema(sc)
if err != nil {
return nil, fmt.Errorf("failed to create table from schema: %w", err)
}
return table, nil
}

func tableFindingToPB(f plugin.TableFinding) *pb.AssessTables_TableFinding {
columns := make([]*pb.AssessTables_ColumnFinding, len(f.Columns))
for i, c := range f.Columns {
columns[i] = &pb.AssessTables_ColumnFinding{
ColumnName: c.ColumnName,
Category: pb.AssessTables_Category(c.Category),
OldType: c.OldType,
NewType: c.NewType,
SafeModeBehavior: c.SafeModeBehavior,
ForcedModeBehavior: c.ForcedModeBehavior,
Evidence: evidenceToPB(c.Evidence),
}
}
return &pb.AssessTables_TableFinding{
TableName: f.TableName,
Category: pb.AssessTables_Category(f.Category),
SafeModeBehavior: f.SafeModeBehavior,
ForcedModeBehavior: f.ForcedModeBehavior,
Columns: columns,
Evidence: evidenceToPB(f.Evidence),
CoverageIncomplete: f.CoverageIncomplete,
CoverageIncompleteReason: f.CoverageIncompleteReason,
}
}

func evidenceToPB(evidence []plugin.Evidence) []*pb.AssessTables_Evidence {
res := make([]*pb.AssessTables_Evidence, len(evidence))
for i, e := range evidence {
res[i] = &pb.AssessTables_Evidence{SyntheticValue: e.SyntheticValue, Before: e.Before, After: e.After}
}
return res
}
156 changes: 156 additions & 0 deletions internal/servers/plugin/v3/assess_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,156 @@
package plugin

import (
"context"
"net"
"testing"

"github.com/apache/arrow-go/v18/arrow"
pb "github.com/cloudquery/plugin-pb-go/pb/plugin/v3"
"github.com/cloudquery/plugin-sdk/v4/plugin"
"github.com/cloudquery/plugin-sdk/v4/schema"
"github.com/google/go-cmp/cmp"
"github.com/rs/zerolog"
"github.com/stretchr/testify/require"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/status"
"google.golang.org/grpc/test/bufconn"
"google.golang.org/protobuf/testing/protocmp"
)

type assessorClient struct {
mockSourceColumnAdderPluginClient
gotTables []plugin.TablePair
gotOptions plugin.AssessOptions
}

func (c *assessorClient) AssessTables(_ context.Context, tables []plugin.TablePair, options plugin.AssessOptions) ([]plugin.TableFinding, error) {
c.gotTables, c.gotOptions = tables, options
return []plugin.TableFinding{{
TableName: "test_table",
Category: plugin.AssessCategoryManualMigrationRequired,
SafeModeBehavior: "rejects this change",
ForcedModeBehavior: "drops and recreates the table",
Columns: []plugin.ColumnFinding{{
ColumnName: "tags",
Category: plugin.AssessCategoryManualMigrationRequired,
OldType: "text[]",
NewType: "jsonb",
SafeModeBehavior: "rejects this change",
ForcedModeBehavior: "drops and recreates the table",
Evidence: []plugin.Evidence{{SyntheticValue: `["env:prod"]`, Before: `{"tags":["env:prod"]}`, After: `{"tags":["env:prod"]}`}},
}},
Evidence: []plugin.Evidence{{SyntheticValue: "header", Before: "tags", After: "tags"}},
CoverageIncomplete: true,
CoverageIncompleteReason: "nested values not compared",
}}, nil
}

func newAssessClient(t *testing.T, newClient plugin.NewClientFunc) pb.PluginClient {
t.Helper()
lis := bufconn.Listen(1024 * 1024)
srv := grpc.NewServer()
pb.RegisterPluginServer(srv, &Server{Plugin: plugin.NewPlugin("test", "development", newClient), Logger: zerolog.Nop()})
go func() { _ = srv.Serve(lis) }()
t.Cleanup(srv.Stop)

conn, err := grpc.NewClient("passthrough:///bufnet",
grpc.WithContextDialer(func(context.Context, string) (net.Conn, error) { return lis.Dial() }),
grpc.WithTransportCredentials(insecure.NewCredentials()))
require.NoError(t, err)
t.Cleanup(func() { _ = conn.Close() })

client := pb.NewPluginClient(conn)
_, err = client.Init(context.Background(), &pb.Init_Request{NoConnection: true})
require.NoError(t, err)
return client
}

func tableBytes(t *testing.T, table *schema.Table) []byte {
t.Helper()
b, err := pb.SchemaToBytes(table.ToArrowSchema())
require.NoError(t, err)
return b
}

func TestAssessTablesRoundTrip(t *testing.T) {
assessor := &assessorClient{}
client := newAssessClient(t, func(context.Context, zerolog.Logger, []byte, plugin.NewClientOptions) (plugin.Client, error) {
return assessor, nil
})
oldTable := &schema.Table{Name: "test_table", Columns: []schema.Column{{Name: "tags", Type: arrow.ListOf(arrow.BinaryTypes.String)}}}
newTable := &schema.Table{Name: "test_table", Columns: []schema.Column{{Name: "tags", Type: arrow.BinaryTypes.String}}}

resp, err := client.AssessTables(context.Background(), &pb.AssessTables_Request{
Tables: []*pb.AssessTables_TablePair{{OldTable: tableBytes(t, oldTable), NewTable: tableBytes(t, newTable)}, {NewTable: tableBytes(t, newTable)}},
MigrateForce: true,
})
require.NoError(t, err)

require.True(t, assessor.gotOptions.MigrateForce)
require.Len(t, assessor.gotTables, 2)
require.Equal(t, "test_table", assessor.gotTables[0].Old.Name)
require.True(t, arrow.TypeEqual(arrow.ListOf(arrow.BinaryTypes.String), assessor.gotTables[0].Old.Columns[0].Type))
require.True(t, arrow.TypeEqual(arrow.BinaryTypes.String, assessor.gotTables[0].New.Columns[0].Type))
require.Nil(t, assessor.gotTables[1].Old)

want := &pb.AssessTables_Response{Tables: []*pb.AssessTables_TableFinding{{
TableName: "test_table",
Category: pb.AssessTables_CATEGORY_MANUAL_MIGRATION_REQUIRED,
SafeModeBehavior: "rejects this change",
ForcedModeBehavior: "drops and recreates the table",
Columns: []*pb.AssessTables_ColumnFinding{{
ColumnName: "tags",
Category: pb.AssessTables_CATEGORY_MANUAL_MIGRATION_REQUIRED,
OldType: "text[]",
NewType: "jsonb",
SafeModeBehavior: "rejects this change",
ForcedModeBehavior: "drops and recreates the table",
Evidence: []*pb.AssessTables_Evidence{{SyntheticValue: `["env:prod"]`, Before: `{"tags":["env:prod"]}`, After: `{"tags":["env:prod"]}`}},
}},
Evidence: []*pb.AssessTables_Evidence{{SyntheticValue: "header", Before: "tags", After: "tags"}},
CoverageIncomplete: true,
CoverageIncompleteReason: "nested values not compared",
}}}
require.Empty(t, cmp.Diff(want, resp, protocmp.Transform()))
}

func TestAssessTablesWithoutAssessorReturnsUnknown(t *testing.T) {
client := newAssessClient(t, getColumnAdderPlugin())
table := &schema.Table{Name: "test_table", Columns: []schema.Column{{Name: "id", Type: arrow.PrimitiveTypes.Int64}}}

resp, err := client.AssessTables(context.Background(), &pb.AssessTables_Request{
Tables: []*pb.AssessTables_TablePair{{OldTable: tableBytes(t, table)}},
})
require.NoError(t, err)

want := &pb.AssessTables_Response{Tables: []*pb.AssessTables_TableFinding{{
TableName: "test_table",
Category: pb.AssessTables_CATEGORY_UNKNOWN,
CoverageIncomplete: true,
CoverageIncompleteReason: plugin.AssessNotSupportedReason,
}}}
require.Empty(t, cmp.Diff(want, resp, protocmp.Transform()))
}

func TestAssessTablesRejectsEmptyPair(t *testing.T) {
client := newAssessClient(t, getColumnAdderPlugin())
_, err := client.AssessTables(context.Background(), &pb.AssessTables_Request{Tables: []*pb.AssessTables_TablePair{{}}})
require.Equal(t, codes.InvalidArgument, status.Code(err))
}

func TestAssessCategoryMatchesProto(t *testing.T) {
for category, want := range map[plugin.AssessCategory]pb.AssessTables_Category{
plugin.AssessCategoryUnknown: pb.AssessTables_CATEGORY_UNKNOWN,
plugin.AssessCategoryNoChange: pb.AssessTables_CATEGORY_NO_CHANGE,
plugin.AssessCategoryAutomaticallyMigratable: pb.AssessTables_CATEGORY_AUTOMATICALLY_MIGRATABLE,
plugin.AssessCategoryManualMigrationRequired: pb.AssessTables_CATEGORY_MANUAL_MIGRATION_REQUIRED,
plugin.AssessCategoryTableRemoved: pb.AssessTables_CATEGORY_TABLE_REMOVED,
plugin.AssessCategoryFileSchemaChanged: pb.AssessTables_CATEGORY_FILE_SCHEMA_CHANGED,
} {
require.Equal(t, want, tableFindingToPB(plugin.TableFinding{Category: category}).Category)
}
require.Len(t, pb.AssessTables_Category_name, 6)
}
94 changes: 94 additions & 0 deletions plugin/plugin_assess.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
package plugin

import (
"context"
"errors"

"github.com/cloudquery/plugin-sdk/v4/schema"
)

type AssessCategory int

const (
AssessCategoryUnknown AssessCategory = iota
AssessCategoryNoChange
AssessCategoryAutomaticallyMigratable
AssessCategoryManualMigrationRequired
AssessCategoryTableRemoved
AssessCategoryFileSchemaChanged
)

const AssessNotSupportedReason = "destination does not support assessment"

// TablePair holds a table before and after a schema change. Old is nil for an added table, New is nil for a removed one.
type TablePair struct {
Old *schema.Table
New *schema.Table
}

func (p TablePair) TableName() string {
if p.New != nil {
return p.New.Name
}
if p.Old != nil {
return p.Old.Name
}
return ""
}

type AssessOptions struct {
MigrateForce bool
}

type Evidence struct {
SyntheticValue string
Before string
After string
}

type ColumnFinding struct {
ColumnName string
Category AssessCategory
OldType string
NewType string
SafeModeBehavior string
ForcedModeBehavior string
Evidence []Evidence
}

type TableFinding struct {
TableName string
Category AssessCategory
SafeModeBehavior string
ForcedModeBehavior string
Columns []ColumnFinding
Evidence []Evidence
CoverageIncomplete bool
CoverageIncompleteReason string
}

// Assessor is an optional DestinationClient interface that reports how schema changes would be applied, without writing anything.
// It is called after Init with NoConnection set, so implementations must not open database or cloud connections.
type Assessor interface {
AssessTables(ctx context.Context, tables []TablePair, options AssessOptions) ([]TableFinding, error)
}

// AssessTables returns one Unknown finding per table when the client does not implement Assessor.
func (p *Plugin) AssessTables(ctx context.Context, tables []TablePair, options AssessOptions) ([]TableFinding, error) {
if p.client == nil {
return nil, errors.New("plugin not initialized. call Init() first")
}
if assessor, ok := p.client.(Assessor); ok {
return assessor.AssessTables(ctx, tables, options)
}
findings := make([]TableFinding, len(tables))
for i, table := range tables {
findings[i] = TableFinding{
TableName: table.TableName(),
Category: AssessCategoryUnknown,
CoverageIncomplete: true,
CoverageIncompleteReason: AssessNotSupportedReason,
}
}
return findings, nil
}
Loading
Loading