Skip to content

Commit e042d73

Browse files
authored
feat: Expose error classifier (#2539)
This PR exposes a way for plugins to capture and classify errors at a plugin level rather than having to handle it on a per resolver basis. For example this can be utilized in our AWS plugin to handle 404s or AWS specific errors without having to replace or wrap the error handling in every single table resolver Example PR that uses this functionality: cloudquery/cloudquery-private#13197
1 parent aee8e3f commit e042d73

9 files changed

Lines changed: 373 additions & 19 deletions

scheduler/queue/scheduler.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ type Scheduler struct {
3333
metrics *metrics.Metrics
3434
invocationID string
3535
seed int64
36+
errorClassifier schema.ErrorClassifier
3637
}
3738

3839
type Option func(*Scheduler)
@@ -61,6 +62,12 @@ func WithInvocationID(invocationID string) Option {
6162
}
6263
}
6364

65+
func WithErrorClassifier(classifier schema.ErrorClassifier) Option {
66+
return func(d *Scheduler) {
67+
d.errorClassifier = classifier
68+
}
69+
}
70+
6471
func NewShuffleQueueScheduler(logger zerolog.Logger, m *metrics.Metrics, seed int64, opts ...Option) *Scheduler {
6572
scheduler := &Scheduler{
6673
logger: logger,
@@ -104,6 +111,7 @@ func (d *Scheduler) Sync(ctx context.Context, tableClients []WorkUnit, resolvedR
104111
d.deterministicCQID,
105112
d.metrics,
106113
msgChan,
114+
d.errorClassifier,
107115
).work(ctx, activeWorkSignal)
108116
return nil
109117
})

scheduler/queue/worker.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,8 @@ type worker struct {
3232
deterministicCQID bool
3333
metrics *metrics.Metrics
3434
// message channel for sending SyncError messages
35-
msgChan chan<- message.SyncMessage
35+
msgChan chan<- message.SyncMessage
36+
errorClassifier schema.ErrorClassifier
3637
}
3738

3839
func (w *worker) work(ctx context.Context, activeWorkSignal *activeWorkSignal) {
@@ -57,6 +58,7 @@ func newWorker(
5758
deterministicCQID bool,
5859
m *metrics.Metrics,
5960
msgChan chan<- message.SyncMessage,
61+
errorClassifier schema.ErrorClassifier,
6062
) *worker {
6163
return &worker{
6264
jobs: jobs,
@@ -68,6 +70,7 @@ func newWorker(
6870
invocationID: invocationID,
6971
metrics: m,
7072
msgChan: msgChan,
73+
errorClassifier: errorClassifier,
7174
}
7275
}
7376

@@ -113,6 +116,11 @@ func (w *worker) resolveTable(ctx context.Context, table *schema.Table, client s
113116
close(res)
114117
}()
115118
if err := table.Resolver(ctx, client, parent, res); err != nil {
119+
event := schema.ErrorEvent{Table: table, Client: client, Phase: schema.ErrorPhaseTableResolver}
120+
if w.errorClassifier.Suppress(ctx, err, event) {
121+
logger.Debug().Err(err).Msg("table resolver finished with suppressed error")
122+
return
123+
}
116124
logger.Error().Err(err).Msg("table resolver finished with error")
117125
w.metrics.AddErrors(ctx, 1, selector)
118126
// Send SyncError message
@@ -155,7 +163,7 @@ func (w *worker) resolveResource(ctx context.Context, table *schema.Table, clien
155163
wg.Add(1)
156164
go func() {
157165
defer wg.Done()
158-
resolvedResources := resolvers.ResolveResourcesChunk(ctx, w.logger, w.metrics, table, client, parent, chunks[i], w.caser)
166+
resolvedResources := resolvers.ResolveResourcesChunkWithClassifier(ctx, w.logger, w.metrics, table, client, parent, chunks[i], w.caser, w.errorClassifier)
159167
for _, resolvedResource := range resolvedResources {
160168
if err := resolvedResource.CalculateCQID(w.deterministicCQID); err != nil {
161169
w.logger.Error().Err(err).Str("table", table.Name).Str("client", client.ID()).Msg("resource resolver finished with primary key calculation error")
Lines changed: 125 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,125 @@
1+
package resolvers
2+
3+
import (
4+
"context"
5+
"errors"
6+
"testing"
7+
8+
"github.com/apache/arrow-go/v18/arrow"
9+
"github.com/cloudquery/plugin-sdk/v4/caser"
10+
"github.com/cloudquery/plugin-sdk/v4/scheduler/metrics"
11+
"github.com/cloudquery/plugin-sdk/v4/schema"
12+
"github.com/rs/zerolog"
13+
"github.com/stretchr/testify/require"
14+
)
15+
16+
var errResolverBoom = errors.New("resolver boom")
17+
18+
// TestResolveResourcesChunk_ErrorClassifier verifies that the error classifier controls
19+
// whether resolver errors at each phase are counted as errors, and that it receives an
20+
// ErrorEvent describing the phase that failed.
21+
func TestResolveResourcesChunk_ErrorClassifier(t *testing.T) {
22+
for _, tc := range []struct {
23+
name string
24+
table func() *schema.Table
25+
wantPhase schema.ErrorPhase
26+
wantColumn string
27+
}{
28+
{
29+
name: "column resolver",
30+
table: func() *schema.Table {
31+
return &schema.Table{
32+
Name: "test_table",
33+
Columns: []schema.Column{
34+
{
35+
Name: "test_column",
36+
Type: arrow.PrimitiveTypes.Int64,
37+
Resolver: func(_ context.Context, _ schema.ClientMeta, _ *schema.Resource, _ schema.Column) error {
38+
return errResolverBoom
39+
},
40+
},
41+
},
42+
}
43+
},
44+
wantPhase: schema.ErrorPhaseColumnResolver,
45+
wantColumn: "test_column",
46+
},
47+
{
48+
name: "pre resource chunk resolver",
49+
table: func() *schema.Table {
50+
return &schema.Table{
51+
Name: "test_table",
52+
PreResourceChunkResolver: &schema.RowsChunkResolver{
53+
ChunkSize: 10,
54+
RowsResolver: func(_ context.Context, _ schema.ClientMeta, _ []*schema.Resource) error {
55+
return errResolverBoom
56+
},
57+
},
58+
Columns: []schema.Column{{Name: "test_column", Type: arrow.PrimitiveTypes.Int64}},
59+
}
60+
},
61+
wantPhase: schema.ErrorPhasePreResourceChunkResolver,
62+
},
63+
{
64+
name: "pre resource resolver",
65+
table: func() *schema.Table {
66+
return &schema.Table{
67+
Name: "test_table",
68+
PreResourceResolver: func(_ context.Context, _ schema.ClientMeta, _ *schema.Resource) error {
69+
return errResolverBoom
70+
},
71+
Columns: []schema.Column{{Name: "test_column", Type: arrow.PrimitiveTypes.Int64}},
72+
}
73+
},
74+
wantPhase: schema.ErrorPhasePreResourceResolver,
75+
},
76+
{
77+
name: "post resource resolver",
78+
table: func() *schema.Table {
79+
return &schema.Table{
80+
Name: "test_table",
81+
PostResourceResolver: func(_ context.Context, _ schema.ClientMeta, _ *schema.Resource) error {
82+
return errResolverBoom
83+
},
84+
Columns: []schema.Column{{Name: "test_column", Type: arrow.PrimitiveTypes.Int64}},
85+
}
86+
},
87+
wantPhase: schema.ErrorPhasePostResourceResolver,
88+
},
89+
} {
90+
t.Run(tc.name, func(t *testing.T) {
91+
client := testClient{}
92+
chunk := []any{0}
93+
logger := zerolog.New(zerolog.NewTestWriter(t))
94+
95+
t.Run("raised without classifier", func(t *testing.T) {
96+
table := tc.table()
97+
m := metrics.NewMetrics()
98+
m.InitWithClients(table, []schema.ClientMeta{client})
99+
ResolveResourcesChunkWithClassifier(context.Background(), logger, m, table, client, nil, chunk, caser.New(), nil)
100+
require.Equal(t, uint64(1), m.GetErrors(m.NewSelector(client.ID(), table.Name)))
101+
})
102+
103+
t.Run("suppressed by classifier", func(t *testing.T) {
104+
table := tc.table()
105+
m := metrics.NewMetrics()
106+
m.InitWithClients(table, []schema.ClientMeta{client})
107+
var gotEvent schema.ErrorEvent
108+
classifier := func(_ context.Context, _ error, event schema.ErrorEvent) bool {
109+
gotEvent = event
110+
return true
111+
}
112+
ResolveResourcesChunkWithClassifier(context.Background(), logger, m, table, client, nil, chunk, caser.New(), classifier)
113+
require.Equal(t, uint64(0), m.GetErrors(m.NewSelector(client.ID(), table.Name)), "suppressed error should not be counted")
114+
require.Equal(t, tc.wantPhase, gotEvent.Phase)
115+
require.Equal(t, table, gotEvent.Table)
116+
if tc.wantColumn != "" {
117+
require.NotNil(t, gotEvent.Column)
118+
require.Equal(t, tc.wantColumn, gotEvent.Column.Name)
119+
} else {
120+
require.Nil(t, gotEvent.Column)
121+
}
122+
})
123+
})
124+
}
125+
}

scheduler/resolvers/resolvers.go

Lines changed: 49 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ import (
1414
"github.com/thoas/go-funk"
1515
)
1616

17-
func resolveColumn(ctx context.Context, logger zerolog.Logger, m *metrics.Metrics, selector metrics.Selector, client schema.ClientMeta, resource *schema.Resource, column schema.Column, c *caser.Caser) {
17+
func resolveColumn(ctx context.Context, logger zerolog.Logger, m *metrics.Metrics, selector metrics.Selector, client schema.ClientMeta, resource *schema.Resource, column schema.Column, c *caser.Caser, classifier schema.ErrorClassifier) {
1818
columnStartTime := time.Now()
1919
defer func() {
2020
if err := recover(); err != nil {
@@ -29,25 +29,38 @@ func resolveColumn(ctx context.Context, logger zerolog.Logger, m *metrics.Metric
2929
}
3030
}()
3131

32+
handleErr := func(err error) {
33+
event := schema.ErrorEvent{Table: resource.Table, Client: client, Phase: schema.ErrorPhaseColumnResolver, Column: &column}
34+
if classifier.Suppress(ctx, err, event) {
35+
logger.Debug().Str("column", column.Name).Err(err).Msg("column resolver finished with suppressed error")
36+
return
37+
}
38+
logger.Error().Err(err).Msg("column resolver finished with error")
39+
m.AddErrors(ctx, 1, selector)
40+
}
41+
3242
if column.Resolver != nil {
3343
if err := column.Resolver(ctx, client, resource, column); err != nil {
34-
logger.Error().Err(err).Msg("column resolver finished with error")
35-
m.AddErrors(ctx, 1, selector)
44+
handleErr(err)
3645
}
3746
} else {
3847
// base use case: try to get column with CamelCase name
3948
v := funk.Get(resource.GetItem(), c.ToPascal(column.Name), funk.WithAllowZero())
4049
if v != nil {
41-
err := resource.Set(column.Name, v)
42-
if err != nil {
43-
logger.Error().Err(err).Msg("column resolver finished with error")
44-
m.AddErrors(ctx, 1, selector)
50+
if err := resource.Set(column.Name, v); err != nil {
51+
handleErr(err)
4552
}
4653
}
4754
}
4855
}
4956

57+
// Deprecated: use ResolveResourcesChunkWithClassifier. This retains the original
58+
// signature and resolves with a nil classifier, so every error is raised.
5059
func ResolveResourcesChunk(ctx context.Context, logger zerolog.Logger, m *metrics.Metrics, table *schema.Table, client schema.ClientMeta, parent *schema.Resource, chunk []any, c *caser.Caser) []*schema.Resource {
60+
return ResolveResourcesChunkWithClassifier(ctx, logger, m, table, client, parent, chunk, c, nil)
61+
}
62+
63+
func ResolveResourcesChunkWithClassifier(ctx context.Context, logger zerolog.Logger, m *metrics.Metrics, table *schema.Table, client schema.ClientMeta, parent *schema.Resource, chunk []any, c *caser.Caser, classifier schema.ErrorClassifier) []*schema.Resource {
5164
ctx, cancel := context.WithTimeout(ctx, 10*time.Minute)
5265
defer cancel()
5366

@@ -72,8 +85,13 @@ func ResolveResourcesChunk(ctx context.Context, logger zerolog.Logger, m *metric
7285

7386
if table.PreResourceChunkResolver != nil {
7487
if err := table.PreResourceChunkResolver.RowsResolver(ctx, client, resources); err != nil {
75-
tableLogger.Error().Stack().Err(err).Msg("pre resource chunk resolver finished with error")
76-
m.AddErrors(ctx, 1, selector)
88+
event := schema.ErrorEvent{Table: table, Client: client, Phase: schema.ErrorPhasePreResourceChunkResolver}
89+
if classifier.Suppress(ctx, err, event) {
90+
tableLogger.Debug().Err(err).Msg("pre resource chunk resolver finished with suppressed error")
91+
} else {
92+
tableLogger.Error().Stack().Err(err).Msg("pre resource chunk resolver finished with error")
93+
m.AddErrors(ctx, 1, selector)
94+
}
7795
return nil
7896
}
7997
}
@@ -82,30 +100,45 @@ func ResolveResourcesChunk(ctx context.Context, logger zerolog.Logger, m *metric
82100
filtered := resources[:0]
83101
for _, resource := range resources {
84102
if err := table.PreResourceResolver(ctx, client, resource); err != nil {
85-
if ctx.Err() != nil {
103+
event := schema.ErrorEvent{Table: table, Client: client, Phase: schema.ErrorPhasePreResourceResolver}
104+
suppress := classifier.Suppress(ctx, err, event)
105+
switch {
106+
case suppress && ctx.Err() != nil:
107+
tableLogger.Debug().Err(err).Msg("pre resource resolver failed, context cancelled (suppressed)")
108+
return nil
109+
case suppress:
110+
tableLogger.Debug().Err(err).Msg("pre resource resolver failed (suppressed)")
111+
continue
112+
case ctx.Err() != nil:
86113
tableLogger.Error().Err(err).Msg("pre resource resolver failed, context cancelled")
87114
m.AddErrors(ctx, 1, selector)
88115
return nil
116+
default:
117+
tableLogger.Error().Err(err).Msg("pre resource resolver failed")
118+
m.AddErrors(ctx, 1, selector)
119+
continue
89120
}
90-
tableLogger.Error().Err(err).Msg("pre resource resolver failed")
91-
m.AddErrors(ctx, 1, selector)
92-
continue
93121
}
94122
filtered = append(filtered, resource)
95123
}
96124
resources = filtered
97125
}
98126
for _, resource := range resources {
99127
for _, column := range table.Columns {
100-
resolveColumn(ctx, tableLogger, m, selector, client, resource, column, c)
128+
resolveColumn(ctx, tableLogger, m, selector, client, resource, column, c, classifier)
101129
}
102130
}
103131

104132
if table.PostResourceResolver != nil {
105133
for _, resource := range resources {
106134
if err := table.PostResourceResolver(ctx, client, resource); err != nil {
107-
tableLogger.Error().Stack().Err(err).Msg("post resource resolver finished with error")
108-
m.AddErrors(ctx, 1, selector)
135+
event := schema.ErrorEvent{Table: table, Client: client, Phase: schema.ErrorPhasePostResourceResolver}
136+
if classifier.Suppress(ctx, err, event) {
137+
tableLogger.Debug().Err(err).Msg("post resource resolver finished with suppressed error")
138+
} else {
139+
tableLogger.Error().Stack().Err(err).Msg("post resource resolver finished with error")
140+
m.AddErrors(ctx, 1, selector)
141+
}
109142
}
110143
}
111144
}

scheduler/scheduler.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,21 @@ const (
3737
StrategyShuffleQueue
3838
)
3939

40+
// Re-exported from the schema package for use with WithErrorClassifier.
41+
type (
42+
ErrorClassifier = schema.ErrorClassifier
43+
ErrorEvent = schema.ErrorEvent
44+
ErrorPhase = schema.ErrorPhase
45+
)
46+
47+
const (
48+
ErrorPhaseTableResolver = schema.ErrorPhaseTableResolver
49+
ErrorPhasePreResourceChunkResolver = schema.ErrorPhasePreResourceChunkResolver
50+
ErrorPhasePreResourceResolver = schema.ErrorPhasePreResourceResolver
51+
ErrorPhaseColumnResolver = schema.ErrorPhaseColumnResolver
52+
ErrorPhasePostResourceResolver = schema.ErrorPhasePostResourceResolver
53+
)
54+
4055
type Option func(*Scheduler)
4156

4257
func WithLogger(logger zerolog.Logger) Option {
@@ -45,6 +60,14 @@ func WithLogger(logger zerolog.Logger) Option {
4560
}
4661
}
4762

63+
// WithErrorClassifier sets the classifier consulted for every resolver error. See
64+
// schema.ErrorClassifier.
65+
func WithErrorClassifier(classifier ErrorClassifier) Option {
66+
return func(s *Scheduler) {
67+
s.errorClassifier = classifier
68+
}
69+
}
70+
4871
func WithConcurrency(concurrency int) Option {
4972
return func(s *Scheduler) {
5073
s.concurrency = concurrency
@@ -122,6 +145,8 @@ type Scheduler struct {
122145
batchSettings *BatchSettings
123146

124147
invocationID string
148+
149+
errorClassifier schema.ErrorClassifier
125150
}
126151

127152
type shard struct {

scheduler/scheduler_dfs.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -128,6 +128,11 @@ func (s *syncClient) resolveTableDfs(ctx context.Context, table *schema.Table, c
128128
close(res)
129129
}()
130130
if err := table.Resolver(ctx, client, parent, res); err != nil {
131+
event := schema.ErrorEvent{Table: table, Client: client, Phase: schema.ErrorPhaseTableResolver}
132+
if s.scheduler.errorClassifier.Suppress(ctx, err, event) {
133+
logger.Debug().Err(err).Msg("table resolver finished with suppressed error")
134+
return
135+
}
131136
logger.Error().Err(err).Msg("table resolver finished with error")
132137
s.metrics.AddErrors(ctx, 1, selector)
133138
// Send SyncError message
@@ -192,7 +197,7 @@ func (s *syncClient) resolveResourcesDfs(ctx context.Context, table *schema.Tabl
192197
defer resourceSem.Release(1)
193198
defer s.scheduler.resourceSem.Release(1)
194199
defer wg.Done()
195-
resolvedResources := resolvers.ResolveResourcesChunk(ctx, s.logger, s.metrics, table, client, parent, chunks[i], s.scheduler.caser)
200+
resolvedResources := resolvers.ResolveResourcesChunkWithClassifier(ctx, s.logger, s.metrics, table, client, parent, chunks[i], s.scheduler.caser, s.scheduler.errorClassifier)
196201
if len(resolvedResources) == 0 {
197202
return
198203
}

0 commit comments

Comments
 (0)