From cf69486027dba957eca57cd79b2aac99e1b1bc65 Mon Sep 17 00:00:00 2001 From: Prathamesh Baviskar Date: Mon, 20 Jul 2026 11:15:37 +0000 Subject: [PATCH 1/2] Tag CloudFetch, Arrow schema, and result-fetch errors for telemetry Third step of PECOBLR-3537. Attaches source-declared error categories at the error sites on the result-materialization paths, so telemetry reports a precise error_name instead of the generic "error" fallback: - chunk_download_error: CloudFetch batch download failures (internal/rows/arrowbased/batchloader.go). - arrow_schema_parsing_error: Arrow schema convert/serialize/read failures (internal/rows/arrowbased/arrowRows.go). - result_set_error: result-page fetch failure (internal/rows/rowscanner/resultPageIterator.go) and result-set metadata fetch failure (internal/rows/rows.go). A caller cancellation/deadline is left untagged at these sites so it still classifies as cancelled/timeout. - statement_execution_timeout: sentinel WatchTimeout in the poll loop (internal/backend/thrift/backend.go). This branch is unreachable in production today (the sole Watch call uses timeout=0), so it emits no telemetry; it is tagged and commented so it classifies correctly if a nonzero poll timeout is ever enabled. All tags use the existing WithCategory chaining, which returns the same concrete pointer, so errors.Is/As, the dbsqlerr.DBError assertion in the Arrow scan path, Error() strings, and sentinel identity are unchanged. Each live tag reaches classifyError via the row-iteration path (Next -> iterationErr -> AfterExecute), verified by the added tests. Separately, this fixes a latent bug in rows.getResultSetSchema where a GetResultSetMetadata failure wrapped the wrong (nil) variable instead of the real cause, dropping it from the error chain. This does not change the telemetry category (the tag wins regardless) but restores Cause()/Unwrap(). unsupported_operation (staging default case) is deferred to a follow-up: it is on a separate path and depends on more fragile invariants. Signed-off-by: Prathamesh Baviskar --- CHANGELOG.md | 2 +- internal/backend/thrift/backend.go | 4 +- internal/rows/arrowbased/arrowRows.go | 6 +- internal/rows/arrowbased/batchloader.go | 6 +- internal/rows/arrowbased/category_test.go | 124 ++++++++++++++++++ internal/rows/category_test.go | 51 +++++++ internal/rows/rows.go | 7 +- internal/rows/rowscanner/category_test.go | 62 +++++++++ .../rows/rowscanner/resultPageIterator.go | 7 +- 9 files changed, 259 insertions(+), 10 deletions(-) create mode 100644 internal/rows/arrowbased/category_test.go create mode 100644 internal/rows/category_test.go create mode 100644 internal/rows/rowscanner/category_test.go 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)) From 4f499c256891e6088fcc2123394d5978fb5bfaeb Mon Sep 17 00:00:00 2001 From: Prathamesh Baviskar Date: Mon, 27 Jul 2026 12:07:14 +0000 Subject: [PATCH 2/2] Simplify Unreleased CHANGELOG entry to a high-level summary Addresses review nit: the changelog line enumerated internal driver category names, which is too low-level for an external changelog. Co-authored-by: Isaac Signed-off-by: Prathamesh Baviskar --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ca5d3ea2..a5f2c72a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,7 +1,7 @@ # Release History ## Unreleased -- 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) +- Improve telemetry error reporting: driver failures are now categorized by cause instead of reported as a generic error (databricks/databricks-sql-go#414, #415, #417) ## 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)