Skip to content
Open
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 CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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)
Expand Down
4 changes: 3 additions & 1 deletion internal/backend/thrift/backend.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
6 changes: 3 additions & 3 deletions internal/rows/arrowbased/arrowRows.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
6 changes: 3 additions & 3 deletions internal/rows/arrowbased/batchloader.go
Original file line number Diff line number Diff line change
Expand Up @@ -505,18 +505,18 @@ 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)
}
}

if lastStatus != 0 {
// 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 {
Expand Down
124 changes: 124 additions & 0 deletions internal/rows/arrowbased/category_test.go
Original file line number Diff line number Diff line change
@@ -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))
})
}
51 changes: 51 additions & 0 deletions internal/rows/category_test.go
Original file line number Diff line number Diff line change
@@ -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))
}
7 changes: 6 additions & 1 deletion internal/rows/rows.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
62 changes: 62 additions & 0 deletions internal/rows/rowscanner/category_test.go
Original file line number Diff line number Diff line change
@@ -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))
}
7 changes: 6 additions & 1 deletion internal/rows/rowscanner/resultPageIterator.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
Loading