Skip to content
Merged
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
94 changes: 76 additions & 18 deletions assess.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,23 +10,43 @@ import (
"github.com/apache/arrow-go/v18/arrow"
"github.com/apache/arrow-go/v18/arrow/array"
"github.com/apache/arrow-go/v18/arrow/memory"
csvfile "github.com/cloudquery/filetypes/v4/csv"
"github.com/cloudquery/filetypes/v4/parquet"
"github.com/cloudquery/filetypes/v4/types"
"github.com/cloudquery/plugin-sdk/v4/plugin"
"github.com/cloudquery/plugin-sdk/v4/schema"
)

const (
behaviorNewColumns = "new files add the new columns, existing files are not changed"
behaviorNewTable = "new files add the new table"
)

// AssessTable reports how the table change affects the files written with the configured format.
// Parquet compares generated schemas. JSON and CSV serialize equivalent synthetic values under both schemas.
func (cl *Client) AssessTable(pair plugin.TablePair) (plugin.TableFinding, error) {
if pq, ok := cl.filetype.(*parquet.Client); ok {
return pq.AssessTable(pair)
}
finding := plugin.TableFinding{TableName: pair.TableName(), Category: plugin.AssessCategoryNoChange}
if pair.New == nil {
finding.Category = plugin.AssessCategoryTableRemoved
return finding, nil
return plugin.TableFinding{TableName: pair.TableName(), Category: plugin.AssessCategoryTableRemoved}, nil
}
finding, outputChanged, err := cl.compareTables(pair)
if err != nil {
return plugin.TableFinding{}, err
}
for i, column := range finding.Columns {
if column.OldType == "" && cl.isAdditiveColumn(pair, column.ColumnName) {
finding.Columns[i].Category = plugin.AssessCategoryAutomaticallyMigratable
}
}
classifyTable(&finding, pair, outputChanged)
return finding, nil
}

func (cl *Client) compareTables(pair plugin.TablePair) (plugin.TableFinding, bool, error) {
if pq, ok := cl.filetype.(*parquet.Client); ok {
finding, err := pq.AssessTable(pair)
return finding, false, err
}
finding := plugin.TableFinding{TableName: pair.TableName()}
oldTable := pair.Old
if oldTable == nil {
oldTable = &schema.Table{Name: pair.New.Name}
Expand All @@ -43,7 +63,7 @@ func (cl *Client) AssessTable(pair plugin.TablePair) (plugin.TableFinding, error
}
column, err := cl.assessColumn(oldTable.Name, oldColumn, *newColumn)
if err != nil {
return plugin.TableFinding{}, err
return plugin.TableFinding{}, false, err
}
finding.Columns = append(finding.Columns, column)
}
Expand All @@ -52,33 +72,71 @@ func (cl *Client) AssessTable(pair plugin.TablePair) (plugin.TableFinding, error
finding.Columns = append(finding.Columns, outputChange(newColumn.Name, "", newColumn.Type.String()))
}
}
if !slices.Equal(oldTable.Columns.Names(), pair.New.Columns.Names()) {
evidence, err := cl.rowEvidence(oldTable, pair.New)
if err != nil {
return plugin.TableFinding{}, err
}
finding.Evidence = append(finding.Evidence, evidence)
if evidence.Before != evidence.After {
finding.Category = plugin.AssessCategoryFileSchemaChanged
}
if pair.Old == nil || slices.Equal(oldTable.Columns.Names(), pair.New.Columns.Names()) {
return finding, false, nil
}
evidence, err := cl.rowEvidence(oldTable, pair.New)
if err != nil {
return plugin.TableFinding{}, false, err
}
finding.Evidence = append(finding.Evidence, evidence)
existingColumns, err := cl.rowEvidence(oldTable, retainedColumns(oldTable, pair.New))
if err != nil {
return plugin.TableFinding{}, false, err
}
return finding, existingColumns.Before != existingColumns.After, nil
}

func retainedColumns(oldTable, newTable *schema.Table) *schema.Table {
retained := *newTable
retained.Columns = slices.DeleteFunc(slices.Clone(newTable.Columns), func(column schema.Column) bool {
return oldTable.Columns.Get(column.Name) == nil
})
return &retained
}

func (cl *Client) isAdditiveColumn(pair plugin.TablePair, name string) bool {
if pair.Old == nil {
return true
}
if column := pair.New.Columns.Get(name); column == nil || column.NotNull {
return false
}
csvClient, ok := cl.filetype.(*csvfile.Client)
if !ok {
return true
}
oldNames, newNames := pair.Old.Columns.Names(), pair.New.Columns.Names()
return csvClient.IncludeHeaders && len(newNames) >= len(oldNames) && slices.Equal(newNames[:len(oldNames)], oldNames)
}

func classifyTable(finding *plugin.TableFinding, pair plugin.TablePair, outputChanged bool) {
finding.Category = plugin.AssessCategoryNoChange
var unknownColumns []string
for _, column := range finding.Columns {
switch column.Category {
case plugin.AssessCategoryFileSchemaChanged:
finding.Category = plugin.AssessCategoryFileSchemaChanged
outputChanged = true
case plugin.AssessCategoryAutomaticallyMigratable:
finding.Category = plugin.AssessCategoryAutomaticallyMigratable
case plugin.AssessCategoryUnknown:
unknownColumns = append(unknownColumns, column.ColumnName)
}
}
switch {
case outputChanged:
finding.Category = plugin.AssessCategoryFileSchemaChanged
case finding.Category == plugin.AssessCategoryAutomaticallyMigratable && pair.Old == nil:
finding.SafeModeBehavior, finding.ForcedModeBehavior = behaviorNewTable, behaviorNewTable
case finding.Category == plugin.AssessCategoryAutomaticallyMigratable:
finding.SafeModeBehavior, finding.ForcedModeBehavior = behaviorNewColumns, behaviorNewColumns
}
if len(unknownColumns) > 0 {
if finding.Category == plugin.AssessCategoryNoChange {
finding.Category = plugin.AssessCategoryUnknown
}
finding.IncompleteCoverageReason = "unable to compare: no equivalent values for columns " + strings.Join(unknownColumns, ", ")
}
return finding, nil
}

func outputChange(name, oldType, newType string) plugin.ColumnFinding {
Expand Down
124 changes: 120 additions & 4 deletions assess_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,11 +26,17 @@ var (
escapingValue = `"comma, \"quote\", back\\slash\nnew line\ttab é"`
jsonSpec = &filetypes.FileSpec{Format: filetypes.FormatTypeJSON}
csvSpec = &filetypes.FileSpec{Format: filetypes.FormatTypeCSV}
parquetSpec = &filetypes.FileSpec{Format: filetypes.FormatTypeParquet}
csvNoHeaderSpec = &filetypes.FileSpec{Format: filetypes.FormatTypeCSV, FormatSpec: map[string]any{"skip_header": true, "delimiter": ";"}}
datadogTagsPair = plugin.TablePair{Old: table(idColumn, tagsStringList), New: table(idColumn, tagsJSON)}
datadogTagsTypes = plugin.ColumnFinding{ColumnName: "tags", OldType: "list<item: utf8, nullable>", NewType: "json"}
)

func notNull(column schema.Column) schema.Column {
column.NotNull = true
return column
}

func columnFinding(base plugin.ColumnFinding, category plugin.AssessCategory, evidence ...plugin.Evidence) plugin.ColumnFinding {
base.Category = category
base.Evidence = evidence
Expand Down Expand Up @@ -155,7 +161,7 @@ func TestAssessTable(t *testing.T) {
Category: plugin.AssessCategoryFileSchemaChanged,
Columns: []plugin.ColumnFinding{
{ColumnName: "name", Category: plugin.AssessCategoryFileSchemaChanged, OldType: "utf8"},
{ColumnName: "tags", Category: plugin.AssessCategoryFileSchemaChanged, NewType: "json"},
{ColumnName: "tags", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "json"},
},
Evidence: []plugin.Evidence{{Before: `{"id":42,"name":"env:prod"}`, After: `{"id":42,"tags":{"env":"prod"}}`}},
},
Expand Down Expand Up @@ -212,16 +218,126 @@ func TestAssessTable(t *testing.T) {
name: "added table",
spec: csvSpec,
pair: plugin.TablePair{New: table(idColumn)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryAutomaticallyMigratable,
SafeModeBehavior: "new files add the new table",
ForcedModeBehavior: "new files add the new table",
Columns: []plugin.ColumnFinding{{ColumnName: "id", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "int64"}},
},
},
{
name: "json added nullable column extends output",
spec: jsonSpec,
pair: plugin.TablePair{Old: table(idColumn, nameString), New: table(tagsJSON, idColumn, nameString)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryAutomaticallyMigratable,
SafeModeBehavior: "new files add the new columns, existing files are not changed",
ForcedModeBehavior: "new files add the new columns, existing files are not changed",
Columns: []plugin.ColumnFinding{{ColumnName: "tags", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "json"}},
Evidence: []plugin.Evidence{{Before: `{"id":42,"name":"env:prod"}`, After: `{"id":42,"name":"env:prod","tags":{"env":"prod"}}`}},
},
},
{
name: "json added not null column changes output",
spec: jsonSpec,
pair: plugin.TablePair{Old: table(idColumn), New: table(idColumn, notNull(nameString))},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryFileSchemaChanged,
Columns: []plugin.ColumnFinding{{ColumnName: "name", Category: plugin.AssessCategoryFileSchemaChanged, NewType: "utf8"}},
Evidence: []plugin.Evidence{{Before: `{"id":42}`, After: `{"id":42,"name":"env:prod"}`}},
},
},
{
name: "csv column appended with header extends output",
spec: csvSpec,
pair: plugin.TablePair{Old: table(idColumn), New: table(idColumn, nameString)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryAutomaticallyMigratable,
SafeModeBehavior: "new files add the new columns, existing files are not changed",
ForcedModeBehavior: "new files add the new columns, existing files are not changed",
Columns: []plugin.ColumnFinding{{ColumnName: "name", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "utf8"}},
Evidence: []plugin.Evidence{{Before: "id\n42", After: "id,name\n42,env:prod"}},
},
},
{
name: "csv column appended without header changes output",
spec: csvNoHeaderSpec,
pair: plugin.TablePair{Old: table(idColumn), New: table(idColumn, nameString)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryFileSchemaChanged,
Columns: []plugin.ColumnFinding{{ColumnName: "id", Category: plugin.AssessCategoryFileSchemaChanged, NewType: "int64"}},
Evidence: []plugin.Evidence{{Before: "", After: "id\n42"}},
Columns: []plugin.ColumnFinding{{ColumnName: "name", Category: plugin.AssessCategoryFileSchemaChanged, NewType: "utf8"}},
Evidence: []plugin.Evidence{{Before: "42", After: "42;env:prod"}},
},
},
{
name: "csv column inserted before existing columns changes output",
spec: csvSpec,
pair: plugin.TablePair{Old: table(idColumn), New: table(nameString, idColumn)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryFileSchemaChanged,
Columns: []plugin.ColumnFinding{{ColumnName: "name", Category: plugin.AssessCategoryFileSchemaChanged, NewType: "utf8"}},
Evidence: []plugin.Evidence{{Before: "id\n42", After: "name,id\nenv:prod,42"}},
},
},
{
name: "parquet added nullable columns extend schema",
spec: parquetSpec,
pair: plugin.TablePair{Old: table(notNull(idColumn)), New: table(nameString, notNull(idColumn), tagsStringList)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryAutomaticallyMigratable,
SafeModeBehavior: "new files add the new columns, existing files are not changed",
ForcedModeBehavior: "new files add the new columns, existing files are not changed",
Columns: []plugin.ColumnFinding{
{ColumnName: "name", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "optional byte_array (String)"},
{ColumnName: "tags", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "optional group (List) {list: repeated group {element: optional byte_array (String)}}"},
},
},
},
{
name: "parquet added required column changes schema",
spec: parquetSpec,
pair: plugin.TablePair{Old: table(nameString), New: table(nameString, notNull(idColumn))},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryFileSchemaChanged,
Columns: []plugin.ColumnFinding{{ColumnName: "id", Category: plugin.AssessCategoryFileSchemaChanged, NewType: "required int64 (Int(bitWidth=64, isSigned=true))"}},
},
},
{
name: "parquet optional to required column changes schema",
spec: parquetSpec,
pair: plugin.TablePair{Old: table(idColumn), New: table(notNull(idColumn), nameString)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryFileSchemaChanged,
Columns: []plugin.ColumnFinding{
{ColumnName: "id", Category: plugin.AssessCategoryFileSchemaChanged, OldType: "optional int64 (Int(bitWidth=64, isSigned=true))", NewType: "required int64 (Int(bitWidth=64, isSigned=true))"},
{ColumnName: "name", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "optional byte_array (String)"},
},
},
},
{
name: "parquet added table",
spec: parquetSpec,
pair: plugin.TablePair{New: table(notNull(idColumn))},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Category: plugin.AssessCategoryAutomaticallyMigratable,
SafeModeBehavior: "new files add the new table",
ForcedModeBehavior: "new files add the new table",
Columns: []plugin.ColumnFinding{{ColumnName: "id", Category: plugin.AssessCategoryAutomaticallyMigratable, NewType: "required int64 (Int(bitWidth=64, isSigned=true))"}},
},
},
{
name: "parquet compares generated schemas",
spec: &filetypes.FileSpec{Format: filetypes.FormatTypeParquet},
spec: parquetSpec,
pair: plugin.TablePair{Old: table(idColumn, nameString), New: table(idColumn)},
want: plugin.TableFinding{
TableName: "datadog_monitors",
Expand Down
Loading