diff --git a/CHANGELOG.md b/CHANGELOG.md index 24e227ed..ca5d3ea2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,7 +1,7 @@ # Release History ## Unreleased -- Enrich telemetry error classification with source-declared error categories. No behavior change yet: the categories are defined and read by `classifyError` but not attached at any error source in this change (databricks/databricks-sql-go#414, #415) +- Enrich telemetry error classification with source-declared error categories, so CloudFetch download, Arrow schema parsing, and result-fetch failures now report a specific `error_name` (`chunk_download_error`, `arrow_schema_parsing_error`, `result_set_error`) instead of the generic `error` (databricks/databricks-sql-go#414, #415) ## v1.14.0 (2026-07-13) - **Minimum Go version is now 1.25.0** (previously 1.20): the `go` directive was raised to 1.25.0 while clearing OSV-Scanner findings and updating dependencies. Consumers building with an older toolchain will need to upgrade Go (databricks/databricks-sql-go#368) diff --git a/internal/backend/thrift/backend.go b/internal/backend/thrift/backend.go index 0f934ccf..51f7398c 100644 --- a/internal/backend/thrift/backend.go +++ b/internal/backend/thrift/backend.go @@ -409,7 +409,9 @@ func (b *Backend) pollOperation(ctx context.Context, opHandle *cli_service.TOper if err != nil { log.Err(err).Msg("error polling operation status") if status == sentinel.WatchTimeout { - err = dbsqlerrint.NewRequestError(ctx, dbsqlerr.ErrSentinelTimeout, err) + // Unreachable today (production Watch uses timeout=0); tagged so it + // classifies correctly if a nonzero poll timeout is ever enabled. + err = dbsqlerrint.NewRequestError(ctx, dbsqlerr.ErrSentinelTimeout, err).WithCategory(dbsqlerrint.CategoryStatementTimeout) } return nil, err } diff --git a/internal/rows/arrowbased/arrowRows.go b/internal/rows/arrowbased/arrowRows.go index 4b99eb9d..58970d85 100644 --- a/internal/rows/arrowbased/arrowRows.go +++ b/internal/rows/arrowbased/arrowRows.go @@ -560,20 +560,20 @@ func tGetResultSetMetadataRespToArrowSchema(resultSetMetadata *cli_service.TGetR arrowSchema, err = tTableSchemaToArrowSchema(resultSetMetadata.Schema, &arrowConfig) if err != nil { logger.Err(err).Msg(errArrowRowsConvertSchema) - return nil, nil, dbsqlerrint.NewDriverError(ctx, errArrowRowsConvertSchema, err) + return nil, nil, dbsqlerrint.NewDriverError(ctx, errArrowRowsConvertSchema, err).WithCategory(dbsqlerrint.CategoryArrowSchemaParsing) } // serialize the arrow schema schemaBytes, err = getArrowSchemaBytes(arrowSchema, ctx) if err != nil { logger.Err(err).Msg(errArrowRowsSerializeSchema) - return nil, nil, dbsqlerrint.NewDriverError(ctx, errArrowRowsSerializeSchema, err) + return nil, nil, dbsqlerrint.NewDriverError(ctx, errArrowRowsSerializeSchema, err).WithCategory(dbsqlerrint.CategoryArrowSchemaParsing) } } else { br := bytes.NewReader(schemaBytes) rdr, err := ipc.NewReader(br) if err != nil { - return nil, nil, dbsqlerrint.NewDriverError(ctx, errArrowRowsUnableToReadBatch, err) + return nil, nil, dbsqlerrint.NewDriverError(ctx, errArrowRowsUnableToReadBatch, err).WithCategory(dbsqlerrint.CategoryArrowSchemaParsing) } defer rdr.Release() diff --git a/internal/rows/arrowbased/batchloader.go b/internal/rows/arrowbased/batchloader.go index 0a8a248e..1b83dac1 100644 --- a/internal/rows/arrowbased/batchloader.go +++ b/internal/rows/arrowbased/batchloader.go @@ -505,7 +505,7 @@ func fetchBatchBytes( if !retry.IsRetryableStatus(res.StatusCode) { msg := fmt.Sprintf("%s: %s %d", errArrowRowsCloudFetchDownloadFailure, "HTTP error", res.StatusCode) - return nil, dbsqlerrint.NewDriverError(ctx, msg, nil) + return nil, dbsqlerrint.NewDriverError(ctx, msg, nil).WithCategory(dbsqlerrint.CategoryChunkDownload) } } @@ -513,10 +513,10 @@ func fetchBatchBytes( // lastErr is nil here by construction: the HTTP-status branch above // explicitly clears it on every iteration. The status code is captured // in msg, so there's no underlying error to wrap. - return nil, dbsqlerrint.NewDriverError(ctx, fmt.Sprintf("%s: %s %d (after %d retries)", errArrowRowsCloudFetchDownloadFailure, "HTTP error", lastStatus, retryMax), nil) + return nil, dbsqlerrint.NewDriverError(ctx, fmt.Sprintf("%s: %s %d (after %d retries)", errArrowRowsCloudFetchDownloadFailure, "HTTP error", lastStatus, retryMax), nil).WithCategory(dbsqlerrint.CategoryChunkDownload) } msg := fmt.Sprintf("%s: %v (after %d retries)", errArrowRowsCloudFetchDownloadFailure, lastErr, retryMax) - return nil, dbsqlerrint.NewDriverError(ctx, msg, lastErr) + return nil, dbsqlerrint.NewDriverError(ctx, msg, lastErr).WithCategory(dbsqlerrint.CategoryChunkDownload) } func getReader(r io.Reader, useLz4Compression bool) io.Reader { diff --git a/internal/rows/arrowbased/category_test.go b/internal/rows/arrowbased/category_test.go new file mode 100644 index 00000000..c8dbb599 --- /dev/null +++ b/internal/rows/arrowbased/category_test.go @@ -0,0 +1,124 @@ +package arrowbased + +import ( + "context" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/databricks/databricks-sql-go/internal/cli_service" + "github.com/databricks/databricks-sql-go/internal/config" + dbsqlerrint "github.com/databricks/databricks-sql-go/internal/errors" + "github.com/stretchr/testify/assert" +) + +// A CloudFetch download failure must surface from the batch iterator carrying +// CategoryChunkDownload, on the same object path the telemetry hook reads. +func TestCloudFetchDownloadErrorCarriesCategory(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusNotFound) + })) + defer server.Close() + + startRowOffset := int64(0) + links := []*cli_service.TSparkArrowResultLink{ + { + FileLink: server.URL, + ExpiryTime: time.Now().Add(10 * time.Minute).Unix(), + StartRowOffset: startRowOffset, + RowCount: 1, + }, + } + + cfg := config.WithDefaults() + cfg.UseLz4Compression = false + cfg.MaxDownloadThreads = 1 + + bi, err := NewCloudBatchIterator(context.Background(), links, startRowOffset, nil, cfg, nil) + assert.Nil(t, err) + + _, nextErr := bi.Next() + assert.NotNil(t, nextErr) + assert.Equal(t, dbsqlerrint.CategoryChunkDownload, dbsqlerrint.CategoryFromError(nextErr), + "download failure must carry chunk_download_error so telemetry classifies it correctly") +} + +// Retry-exhausted (persistent retryable status) and transport-error download +// failures also carry CategoryChunkDownload — the two branches after the retry +// loop, distinct from the non-retryable-status branch above. +func TestCloudFetchDownloadErrorCarriesCategory_RetryBranches(t *testing.T) { + link := func(url string) []*cli_service.TSparkArrowResultLink { + return []*cli_service.TSparkArrowResultLink{{ + FileLink: url, + ExpiryTime: time.Now().Add(10 * time.Minute).Unix(), + StartRowOffset: 0, + RowCount: 1, + }} + } + cfg := config.WithDefaults() + cfg.UseLz4Compression = false + cfg.MaxDownloadThreads = 1 + cfg.RetryMax = 1 + cfg.RetryWaitMin = 1 * time.Millisecond + cfg.RetryWaitMax = 5 * time.Millisecond + + t.Run("retries exhausted on persistent 503", func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusServiceUnavailable) + })) + defer server.Close() + + bi, err := NewCloudBatchIterator(context.Background(), link(server.URL), 0, nil, cfg, nil) + assert.Nil(t, err) + _, nextErr := bi.Next() + assert.NotNil(t, nextErr) + assert.Equal(t, dbsqlerrint.CategoryChunkDownload, dbsqlerrint.CategoryFromError(nextErr)) + }) + + t.Run("transport error (unreachable server)", func(t *testing.T) { + // Start then immediately close the server so the connection is refused, + // driving the transport-error branch (lastErr set, lastStatus 0). + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {})) + url := server.URL + server.Close() + + bi, err := NewCloudBatchIterator(context.Background(), link(url), 0, nil, cfg, nil) + assert.Nil(t, err) + _, nextErr := bi.Next() + assert.NotNil(t, nextErr) + assert.Equal(t, dbsqlerrint.CategoryChunkDownload, dbsqlerrint.CategoryFromError(nextErr)) + }) +} + +// A schema-parse failure while building the row scanner must carry +// CategoryArrowSchemaParsing. Covers the schema-convert branch (broken decimal +// type) and the IPC-read branch (invalid pre-supplied arrow schema bytes). +func TestArrowSchemaParseErrorCarriesCategory(t *testing.T) { + t.Run("schema convert failure", func(t *testing.T) { + rowSet := &cli_service.TRowSet{ArrowBatches: []*cli_service.TSparkArrowBatch{{RowCount: 2}}} + schema := getAllTypesSchema() + // Break the decimal column so tTableSchemaToArrowSchema fails. + schema.Columns[13].TypeDesc.Types[0].PrimitiveEntry.TypeQualifiers = nil + metadataResp := getMetadataResp(schema) + + cfg := config.Config{} + cfg.UseArrowBatches = true + cfg.ArrowConfig.UseArrowNativeDecimal = true + + _, err := NewArrowRowScanner(metadataResp, rowSet, &cfg, nil, context.Background(), nil) + assert.NotNil(t, err) + assert.Equal(t, dbsqlerrint.CategoryArrowSchemaParsing, dbsqlerrint.CategoryFromError(err)) + }) + + t.Run("invalid IPC arrow schema bytes", func(t *testing.T) { + rowSet := &cli_service.TRowSet{ArrowBatches: []*cli_service.TSparkArrowBatch{{RowCount: 2}}} + metadataResp := getMetadataResp(getAllTypesSchema()) + // A non-nil but invalid ArrowSchema forces the ipc.NewReader branch to fail. + metadataResp.ArrowSchema = []byte("not a valid arrow ipc stream") + + _, err := NewArrowRowScanner(metadataResp, rowSet, &config.Config{}, nil, context.Background(), nil) + assert.NotNil(t, err) + assert.Equal(t, dbsqlerrint.CategoryArrowSchemaParsing, dbsqlerrint.CategoryFromError(err)) + }) +} diff --git a/internal/rows/category_test.go b/internal/rows/category_test.go new file mode 100644 index 00000000..1b50f3d2 --- /dev/null +++ b/internal/rows/category_test.go @@ -0,0 +1,51 @@ +package rows + +import ( + "context" + "testing" + + "github.com/databricks/databricks-sql-go/internal/cli_service" + "github.com/databricks/databricks-sql-go/internal/client" + dbsqlerrint "github.com/databricks/databricks-sql-go/internal/errors" + "github.com/pkg/errors" + "github.com/stretchr/testify/assert" +) + +// A GetResultSetMetadata failure must carry CategoryResultSet, and (regression +// for the former err/err2 bug) must wrap the real underlying cause rather than +// dropping it. +func TestResultSetMetadataErrorCarriesCategoryAndCause(t *testing.T) { + cause := errors.New("metadata rpc failed") + failingClient := &client.TestClient{ + FnGetResultSetMetadata: func(ctx context.Context, req *cli_service.TGetResultSetMetadataReq) (*cli_service.TGetResultSetMetadataResp, error) { + return nil, cause + }, + } + + r := &rows{client: failingClient, ctx: context.Background()} + + _, err := r.getResultSetSchema() + assert.NotNil(t, err) + assert.Equal(t, dbsqlerrint.CategoryResultSet, dbsqlerrint.CategoryFromError(err)) + // The real cause must be preserved (previously the wrong variable was wrapped). + assert.True(t, errors.Is(err, cause)) +} + +// A metadata fetch that fails when the results context is already aborted must +// NOT be tagged result_set_error. +func TestResultSetMetadataCancelNotTaggedResultSet(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + failingClient := &client.TestClient{ + FnGetResultSetMetadata: func(ctx context.Context, req *cli_service.TGetResultSetMetadataReq) (*cli_service.TGetResultSetMetadataResp, error) { + return nil, errors.New("boom") + }, + } + + r := &rows{client: failingClient, ctx: ctx} + + _, err := r.getResultSetSchema() + assert.NotNil(t, err) + assert.Equal(t, dbsqlerrint.ErrorCategory(""), dbsqlerrint.CategoryFromError(err)) +} diff --git a/internal/rows/rows.go b/internal/rows/rows.go index a87e0ceb..1ba3bb7c 100644 --- a/internal/rows/rows.go +++ b/internal/rows/rows.go @@ -526,7 +526,12 @@ func (r *rows) getResultSetSchema() (*cli_service.TTableSchema, dbsqlerr.DBError resp, err2 := r.client.GetResultSetMetadata(r.ctx, &req) if err2 != nil { r.logger().Err(err2).Msg(err2.Error()) - return nil, dbsqlerr_int.NewRequestError(r.ctx, errRowsMetadataFetchFailed, err) + // A fetch aborted via the results context (e.g. Close) isn't a + // result-set failure; leave it untagged. Mirrors the CloudFetch path. + if r.ctx != nil && r.ctx.Err() != nil { + return nil, dbsqlerr_int.NewRequestError(r.ctx, errRowsMetadataFetchFailed, err2) + } + return nil, dbsqlerr_int.NewRequestError(r.ctx, errRowsMetadataFetchFailed, err2).WithCategory(dbsqlerr_int.CategoryResultSet) } r.resultSetMetadata = resp diff --git a/internal/rows/rowscanner/category_test.go b/internal/rows/rowscanner/category_test.go new file mode 100644 index 00000000..d5292da3 --- /dev/null +++ b/internal/rows/rowscanner/category_test.go @@ -0,0 +1,62 @@ +package rowscanner + +import ( + "context" + "testing" + + "github.com/databricks/databricks-sql-go/internal/cli_service" + "github.com/databricks/databricks-sql-go/internal/client" + dbsqlerrint "github.com/databricks/databricks-sql-go/internal/errors" + dbsqllog "github.com/databricks/databricks-sql-go/logger" + "github.com/pkg/errors" + "github.com/stretchr/testify/assert" +) + +// A FetchResults failure surfaced from the result page iterator must carry +// CategoryResultSet so telemetry reports "result_set_error". +func TestResultPageFetchErrorCarriesCategory(t *testing.T) { + failingClient := &client.TestClient{ + FnFetchResults: func(ctx context.Context, req *cli_service.TFetchResultsReq) (*cli_service.TFetchResultsResp, error) { + return nil, errors.New("boom") + }, + } + + rpf := &resultPageIterator{ + Delimiter: NewDelimiter(0, 0), + client: failingClient, + ctx: context.Background(), // live ctx, like production, so the guard's real branch runs + logger: dbsqllog.WithContext("connId", "correlationId", ""), + connectionId: "connId", + correlationId: "correlationId", + } + + _, err := rpf.Next() + assert.NotNil(t, err) + assert.Equal(t, dbsqlerrint.CategoryResultSet, dbsqlerrint.CategoryFromError(err)) +} + +// A page fetch that fails when the results context is already aborted must NOT +// be tagged result_set_error. +func TestResultPageFetchCancelNotTaggedResultSet(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + failingClient := &client.TestClient{ + FnFetchResults: func(ctx context.Context, req *cli_service.TFetchResultsReq) (*cli_service.TFetchResultsResp, error) { + return nil, errors.New("boom") + }, + } + + rpf := &resultPageIterator{ + Delimiter: NewDelimiter(0, 0), + client: failingClient, + ctx: ctx, + logger: dbsqllog.WithContext("connId", "correlationId", ""), + connectionId: "connId", + correlationId: "correlationId", + } + + _, err := rpf.Next() + assert.NotNil(t, err) + assert.Equal(t, dbsqlerrint.ErrorCategory(""), dbsqlerrint.CategoryFromError(err)) +} diff --git a/internal/rows/rowscanner/resultPageIterator.go b/internal/rows/rowscanner/resultPageIterator.go index 0743443a..22f2af25 100644 --- a/internal/rows/rowscanner/resultPageIterator.go +++ b/internal/rows/rowscanner/resultPageIterator.go @@ -194,7 +194,12 @@ func (rpf *resultPageIterator) getNextPage() (*cli_service.TFetchResultsResp, er fetchResult, err = rpf.client.FetchResults(rpf.ctx, &req) if err != nil { rpf.logger.Err(err).Msg("databricks: Rows instance failed to retrieve results") - return nil, dbsqlerrint.NewRequestError(rpf.ctx, errRowsResultFetchFailed, err) + // A fetch aborted via the results context (e.g. Close) isn't a + // result-set failure; leave it untagged. Mirrors the CloudFetch path. + if rpf.ctx != nil && rpf.ctx.Err() != nil { + return nil, dbsqlerrint.NewRequestError(rpf.ctx, errRowsResultFetchFailed, err) + } + return nil, dbsqlerrint.NewRequestError(rpf.ctx, errRowsResultFetchFailed, err).WithCategory(dbsqlerrint.CategoryResultSet) } rpf.Delimiter = NewDelimiter(fetchResult.Results.StartRowOffset, CountRows(fetchResult.Results))