From f663f4156477306a50278ad44796a1284cb9765d Mon Sep 17 00:00:00 2001 From: "ben.hansen" Date: Thu, 1 Oct 2026 16:37:01 +0000 Subject: [PATCH 1/5] [air] Record workload submission and source staging telemetry --- cmd/air/runsubmit.go | 19 +-- cmd/air/runsubmit_test.go | 22 +++ cmd/air/snapshot.go | 16 ++- cmd/air/snapshot_upload.go | 66 +++++---- cmd/air/telemetry.go | 61 ++++++++ cmd/air/telemetry_test.go | 197 ++++++++++++++++++++++++++ libs/telemetry/protos/README.md | 23 +++ libs/telemetry/protos/air_run.go | 46 ++++++ libs/telemetry/protos/air_run_test.go | 34 +++++ libs/telemetry/protos/frontend_log.go | 1 + libs/testserver/fake_workspace.go | 3 + 11 files changed, 444 insertions(+), 44 deletions(-) create mode 100644 cmd/air/telemetry.go create mode 100644 cmd/air/telemetry_test.go create mode 100644 libs/telemetry/protos/air_run.go create mode 100644 libs/telemetry/protos/air_run_test.go diff --git a/cmd/air/runsubmit.go b/cmd/air/runsubmit.go index 0404d9009fe..dfc847fcb52 100644 --- a/cmd/air/runsubmit.go +++ b/cmd/air/runsubmit.go @@ -10,6 +10,7 @@ import ( "path" "strconv" "strings" + "time" "github.com/databricks/cli/libs/auth" "github.com/databricks/cli/libs/cmdio" @@ -321,17 +322,20 @@ func stageRunArtifacts(ctx context.Context, launchWriter fileWriter, items []upl }) } - if err := group.Wait(); err != nil { - return snapshotResult{}, err - } - return snap, nil + err := group.Wait() + // Preserve measurements from attempted snapshot phases even if staging fails. + return snap, err } // submitWorkload runs the submit happy path: ensure the experiment directory, // upload the launch artifacts, assemble the Jobs payload, and submit it. It // returns the new run_id and its dashboard URL. showProgress enables the stderr // staging spinner (text mode only). -func submitWorkload(ctx context.Context, w *databricks.WorkspaceClient, cfg *runConfig, configPath, idempotencyKey string, showProgress bool) (int64, string, error) { +func submitWorkload(ctx context.Context, w *databricks.WorkspaceClient, cfg *runConfig, configPath, idempotencyKey string, showProgress bool) (runID int64, dashboardURL string, err error) { + start := time.Now() + var snap snapshotResult + defer func() { logRunEvent(ctx, cfg, snap, runID, time.Since(start), err) }() + // Compute and validate the actual submission path before creating artifacts. base, funcDir, commandPath, err := prospectiveLaunchPaths(ctx, w, cfg) if err != nil { @@ -399,7 +403,6 @@ func submitWorkload(ctx context.Context, w *databricks.WorkspaceClient, cfg *run } } - var snap snapshotResult err = withSpinner(ctx, showProgress, "Staging run artifacts…", func() error { var stageErr error snap, stageErr = stageRunArtifacts(ctx, fc, items, stageSnapshot) @@ -426,12 +429,12 @@ func submitWorkload(ctx context.Context, w *databricks.WorkspaceClient, cfg *run // Submit returns as soon as the run is created; we don't wait for it to finish. // Permissions are granted by the caller, after the submit result is shown, so // the best-effort grant never delays the success line. - runID, err := submitRun(ctx, w, payload, poolID, priorityClass, cfg.unityCatalogImagePath(), containers) + runID, err = submitRun(ctx, w, payload, poolID, priorityClass, cfg.unityCatalogImagePath(), containers) if err != nil { return 0, "", err } - dashboardURL := strings.TrimRight(w.Config.Host, "/") + "/jobs/runs/" + strconv.FormatInt(runID, 10) + dashboardURL = strings.TrimRight(w.Config.Host, "/") + "/jobs/runs/" + strconv.FormatInt(runID, 10) return runID, dashboardURL, nil } diff --git a/cmd/air/runsubmit_test.go b/cmd/air/runsubmit_test.go index 8d80dc85c2e..8a0f607b09d 100644 --- a/cmd/air/runsubmit_test.go +++ b/cmd/air/runsubmit_test.go @@ -781,6 +781,16 @@ code_source: // zero new import-file calls (a real skip), not a re-upload to the same name. assert.Equal(t, first.CodeSourcePath, second.CodeSourcePath) assert.Equal(t, afterFirst, snapshotUploads, "unchanged plain_tar should skip the second upload") + require.NotNil(t, first.SizeBytes) + assert.Positive(t, *first.SizeBytes) + assert.Equal(t, first.SizeBytes, second.SizeBytes) + assert.Equal(t, new(false), second.UsesGit) + assert.Equal(t, new(modePlainTar), first.PackagingMode) + assert.Equal(t, first.PackagingMode, second.PackagingMode) + assert.NotNil(t, first.PackagingDurationMs) + assert.NotNil(t, first.UploadDurationMs) + assert.Equal(t, new(int64(0)), second.PackagingDurationMs) + assert.Equal(t, new(int64(0)), second.UploadDurationMs) } // A git_archive snapshot is content-addressed by (commit, include_paths): submitting @@ -837,6 +847,16 @@ code_source: // (the second submit is a cache hit and moves no bytes). assert.Equal(t, first.CodeSourcePath, second.CodeSourcePath) assert.Len(t, uploaded, 1, "git_archive cache hit should skip the second upload") + require.NotNil(t, first.SizeBytes) + assert.Positive(t, *first.SizeBytes) + assert.Equal(t, first.SizeBytes, second.SizeBytes) + assert.Equal(t, new(true), second.UsesGit) + assert.Equal(t, new(modeGitArchive), first.PackagingMode) + assert.Equal(t, first.PackagingMode, second.PackagingMode) + assert.NotNil(t, first.PackagingDurationMs) + assert.NotNil(t, first.UploadDurationMs) + assert.Equal(t, new(int64(0)), second.PackagingDurationMs) + assert.Equal(t, new(int64(0)), second.UploadDurationMs) } // When enabled, a code source uploads provenance sidecars (git_state.json and @@ -877,6 +897,8 @@ code_source: assert.Empty(t, snap.GitStatePath) assert.Empty(t, snap.GitDiffPath) + assert.Equal(t, new(true), snap.UsesGit) + assert.Equal(t, new(modePlainTar), snap.PackagingMode) _, err = sidecarStore.Read(ctx, gitStateName) assert.ErrorIs(t, err, fs.ErrNotExist) diff --git a/cmd/air/snapshot.go b/cmd/air/snapshot.go index 603b7e6c0f8..fcf439ebb24 100644 --- a/cmd/air/snapshot.go +++ b/cmd/air/snapshot.go @@ -10,13 +10,17 @@ import ( "github.com/databricks/cli/libs/env" ) -// snapshotResult holds the code_source_path wired into the submit payload (the -// uploaded code archive's remote path) plus the remote paths of the best-effort git -// provenance sidecars (empty when not a git repo or upload failed). +// snapshotResult holds artifact paths and snapshot measurements. Measurements +// survive staging failures; CodeSourcePath is set only after upload or a cache hit. type snapshotResult struct { - CodeSourcePath string - GitStatePath string - GitDiffPath string + CodeSourcePath string + GitStatePath string + GitDiffPath string + SizeBytes *int64 + UsesGit *bool + PackagingMode *snapshotMode + PackagingDurationMs *int64 + UploadDurationMs *int64 } // resolveRootPath resolves a code_source snapshot root_path: expand environment diff --git a/cmd/air/snapshot_upload.go b/cmd/air/snapshot_upload.go index 9c44a064d98..b36de94be35 100644 --- a/cmd/air/snapshot_upload.go +++ b/cmd/air/snapshot_upload.go @@ -49,7 +49,7 @@ func uploadSnapshot(ctx context.Context, w *databricks.WorkspaceClient, snap *sn } result, err := uploadSnapshotTarball(ctx, w, repoPath, plan, snapshotArtifactPath) if err != nil { - return snapshotResult{}, err + return result, err } // Upload git provenance sidecars (git_state.json / git_diff.patch) next to the @@ -158,49 +158,69 @@ const airArtifactInternalDir = ".internal" // for git_archive and by the working-tree metadata fingerprint for plain_tar — so if // the identical object already exists we skip packaging and upload entirely and reuse it. func uploadSnapshotTarball(ctx context.Context, w *databricks.WorkspaceClient, repoPath string, plan snapshotPlan, artifactPath string) (snapshotResult, error) { + result := snapshotResult{UsesGit: &plan.isGitRepo, PackagingMode: &plan.mode} f, uploadPath, err := snapshotUploadFiler(ctx, w, artifactPath) if err != nil { - return snapshotResult{}, err + return result, err } tarName, files, err := snapshotTarName(ctx, repoPath, plan) if err != nil { - return snapshotResult{}, err + return result, err } // code_source_path is content-addressed by tarName, so it is the same whether we // upload the bytes now or reuse an object already in the store. remote := path.Join(uploadPath, tarName) - exists, err := snapshotExists(ctx, f, tarName) - if err != nil { - return snapshotResult{}, err - } - if exists { + info, err := f.Stat(ctx, tarName) + if err == nil { log.Debugf(ctx, "snapshot upload skipped; reusing %s", remote) - return snapshotResult{CodeSourcePath: remote}, nil + size := info.Size() + result.CodeSourcePath = remote + result.SizeBytes = &size + result.PackagingDurationMs = new(int64(0)) + result.UploadDurationMs = new(int64(0)) + return result, nil + } + if !errors.Is(err, fs.ErrNotExist) { + return result, fmt.Errorf("failed to check snapshot cache: %w", err) } tmp, err := os.MkdirTemp("", "air-snapshot-*") if err != nil { - return snapshotResult{}, err + return result, err } defer os.RemoveAll(tmp) tarball := filepath.Join(tmp, tarName) - if err := packageSnapshot(ctx, repoPath, plan, files, tarball); err != nil { - return snapshotResult{}, err + packagingStart := time.Now() + err = packageSnapshot(ctx, repoPath, plan, files, tarball) + result.PackagingDurationMs = new(time.Since(packagingStart).Milliseconds()) + if err != nil { + return result, err } file, err := os.Open(tarball) if err != nil { - return snapshotResult{}, err + return result, err } defer file.Close() + // Size collection is best-effort and must not prevent a valid upload. + if info, err := file.Stat(); err == nil { + size := info.Size() + result.SizeBytes = &size + } else { + log.Debugf(ctx, "failed to measure code snapshot: %v", err) + } cmdio.LogProgress(ctx, fmt.Sprintf("Uploading %s...", tarName)) - if err := f.Write(ctx, tarName, file, filer.OverwriteIfExists, filer.CreateParentDirectories); err != nil { - return snapshotResult{}, fmt.Errorf("failed to upload snapshot %s: %w", tarName, err) + uploadStart := time.Now() + err = f.Write(ctx, tarName, file, filer.OverwriteIfExists, filer.CreateParentDirectories) + result.UploadDurationMs = new(time.Since(uploadStart).Milliseconds()) + if err != nil { + return result, fmt.Errorf("failed to upload snapshot %s: %w", tarName, err) } - return snapshotResult{CodeSourcePath: remote}, nil + result.CodeSourcePath = remote + return result, nil } // snapshotUploadFiler returns a filer rooted at /.internal plus that @@ -223,17 +243,3 @@ func snapshotUploadFiler(ctx context.Context, w *databricks.WorkspaceClient, art f, err := filer.NewWorkspaceFilesClient(w, uploadPath) return f, uploadPath, err } - -// snapshotExists reports whether name already exists in the artifact store, used to -// short-circuit a content-addressed upload (either mode). A not-found is a clean miss -// (false, nil); any other error is surfaced. -func snapshotExists(ctx context.Context, store filer.Filer, name string) (bool, error) { - _, err := store.Stat(ctx, name) - if err == nil { - return true, nil - } - if errors.Is(err, fs.ErrNotExist) { - return false, nil - } - return false, fmt.Errorf("failed to check snapshot cache: %w", err) -} diff --git a/cmd/air/telemetry.go b/cmd/air/telemetry.go new file mode 100644 index 00000000000..d56dbabd63e --- /dev/null +++ b/cmd/air/telemetry.go @@ -0,0 +1,61 @@ +package aircmd + +import ( + "context" + "strconv" + "time" + + "github.com/databricks/cli/libs/telemetry" + "github.com/databricks/cli/libs/telemetry/protos" +) + +// logRunEvent records the submission outcome independently of a later --watch +// outcome. cfg has already passed local validation before submission starts. +func logRunEvent(ctx context.Context, cfg *runConfig, snap snapshotResult, runID int64, elapsed time.Duration, err error) { + perNode, _ := gpusPerNode(gpuType(cfg.Compute.AcceleratorType)) + _, hasDependencies := cfg.inlineDependencies() + event := &protos.AirRunEvent{ + GPUType: telemetryGPUType(gpuType(cfg.Compute.AcceleratorType)), + NumGPUs: cfg.Compute.NumAccelerators, + NumNodes: cfg.Compute.NumAccelerators / perNode, + HasDockerImage: cfg.unityCatalogImagePath() != "", + HasCodeSnapshot: cfg.CodeSource != nil && cfg.CodeSource.Snapshot != nil, + HasRequirements: hasDependencies, + HasParameters: len(cfg.Parameters) > 0, + MaxRetries: cfg.maxRetries(), + HasTimeout: cfg.TimeoutMinutes != nil, + SubmittedSuccessfully: err == nil, + SubmitLatencyMs: elapsed.Milliseconds(), + CodeSourceUsesGit: snap.UsesGit, + CodeSourceSizeBytes: snap.SizeBytes, + CodeSourcePackagingDurationMs: snap.PackagingDurationMs, + CodeSourceUploadDurationMs: snap.UploadDurationMs, + } + if snap.PackagingMode != nil { + switch *snap.PackagingMode { + case modeGitArchive: + event.CodeSourcePackagingMode = protos.AirPackagingModeGitArchive + case modePlainTar: + event.CodeSourcePackagingMode = protos.AirPackagingModePlainTar + } + } + if err == nil { + event.JobRunID = strconv.FormatInt(runID, 10) + } + telemetry.Log(ctx, protos.DatabricksCliLog{AirRunEvent: event}) +} + +func telemetryGPUType(g gpuType) protos.AirGPUType { + switch g { + case gpuType1xA10: + return protos.AirGPUType1xA10 + case gpuType1xH100: + return protos.AirGPUType1xH100 + case gpuType8xH100: + return protos.AirGPUType8xH100 + case gpuType8xB300: + return protos.AirGPUType8xB300 + default: + return protos.AirGPUTypeUnspecified + } +} diff --git a/cmd/air/telemetry_test.go b/cmd/air/telemetry_test.go new file mode 100644 index 00000000000..00b313a2d17 --- /dev/null +++ b/cmd/air/telemetry_test.go @@ -0,0 +1,197 @@ +package aircmd + +import ( + "context" + "encoding/json" + "net/http" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/databricks/cli/libs/cmdctx" + "github.com/databricks/cli/libs/cmdio" + "github.com/databricks/cli/libs/telemetry" + "github.com/databricks/cli/libs/telemetry/protos" + "github.com/databricks/cli/libs/testserver" + "github.com/databricks/databricks-sdk-go" + "github.com/databricks/databricks-sdk-go/service/jobs" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func captureAirTelemetry(t *testing.T, server *testserver.Server, w *databricks.WorkspaceClient) (context.Context, *[]protos.FrontendLog) { + t.Helper() + var events []protos.FrontendLog + server.Handle("POST", "/telemetry-ext", func(req testserver.Request) any { + var body telemetry.RequestBody + require.NoError(t, json.Unmarshal(req.Body, &body)) + for _, raw := range body.ProtoLogs { + var event protos.FrontendLog + require.NoError(t, json.Unmarshal([]byte(raw), &event)) + events = append(events, event) + } + return telemetry.ResponseBody{NumProtoSuccess: int64(len(body.ProtoLogs))} + }) + ctx := telemetry.WithNewLogger(cmdio.MockDiscard(t.Context())) + return cmdctx.SetConfigUsed(ctx, w.Config), &events +} + +func TestAirRunTelemetry(t *testing.T) { + for _, tc := range []struct { + name string + snapshot bool + fail bool + failUpload bool + }{ + {name: "without snapshot"}, + {name: "uploaded snapshot", snapshot: true}, + {name: "failed submit", snapshot: true, fail: true}, + {name: "failed upload", snapshot: true, failUpload: true}, + } { + t.Run(tc.name, func(t *testing.T) { + server := testserver.New(t) + t.Cleanup(server.Close) + server.Handle("POST", "/api/2.2/jobs/runs/submit", func(req testserver.Request) any { + if tc.fail { + return testserver.Response{StatusCode: http.StatusBadRequest, Body: map[string]string{ + "error_code": "INVALID_PARAMETER_VALUE", "message": "submission rejected", + }} + } + return jobs.SubmitRunResponse{RunId: 555} + }) + var uploadedSize int64 + server.Handle("POST", "/api/2.0/workspace-files/import-file/{path...}", func(req testserver.Request) any { + p := req.Vars["path"] + if strings.HasSuffix(p, ".tar.gz") { + uploadedSize = int64(len(req.Body)) + time.Sleep(20 * time.Millisecond) + if tc.failUpload { + return testserver.Response{StatusCode: http.StatusForbidden, Body: map[string]string{ + "error_code": "PERMISSION_DENIED", "message": "upload rejected", + }} + } + } + return req.Workspace.WorkspaceFilesImportFile(p, req.Body, req.URL.Query().Get("overwrite") == "true") + }) + stubValidateConfig(server) + w, err := databricks.NewWorkspaceClient(&databricks.Config{Host: server.URL, Token: "token"}) + require.NoError(t, err) + ctx, events := captureAirTelemetry(t, server, w) + testserver.AddDefaultHandlers(server) + + configYAML := minimalConfig + if tc.snapshot { + repo := filepath.Join(t.TempDir(), "src") + writeRepoFile(t, repo, "train.py", "print('private source')") + configYAML += "\ncode_source:\n type: snapshot\n snapshot:\n root_path: " + filepath.ToSlash(repo) + "\n" + } + cfgPath := writeConfigFile(t, "run.yaml", configYAML) + cfg, err := loadRunConfig(cfgPath) + require.NoError(t, err) + _, _, err = submitWorkload(ctx, w, cfg, cfgPath, "idem", false) + if tc.fail || tc.failUpload { + require.Error(t, err) + } else { + require.NoError(t, err) + } + require.NoError(t, telemetry.Upload(ctx, protos.ExecutionContext{Command: "air_run"})) + require.Len(t, *events, 1) + event := (*events)[0].Entry.DatabricksCliLog.AirRunEvent + require.NotNil(t, event) + assert.Equal(t, !tc.fail && !tc.failUpload, event.SubmittedSuccessfully) + assert.Equal(t, tc.snapshot, event.HasCodeSnapshot) + assert.Equal(t, protos.AirGPUType1xH100, event.GPUType) + assert.Equal(t, 1, event.NumGPUs) + assert.Equal(t, 1, event.NumNodes) + assert.GreaterOrEqual(t, event.SubmitLatencyMs, int64(0)) + if tc.fail || tc.failUpload { + assert.Empty(t, event.JobRunID) + } else { + assert.Equal(t, "555", event.JobRunID) + } + if tc.snapshot { + require.NotNil(t, event.CodeSourceSizeBytes) + assert.Positive(t, uploadedSize) + assert.Equal(t, uploadedSize, *event.CodeSourceSizeBytes) + require.NotNil(t, event.CodeSourceUsesGit) + assert.False(t, *event.CodeSourceUsesGit) + assert.Equal(t, protos.AirPackagingModePlainTar, event.CodeSourcePackagingMode) + require.NotNil(t, event.CodeSourcePackagingDurationMs) + assert.GreaterOrEqual(t, *event.CodeSourcePackagingDurationMs, int64(0)) + require.NotNil(t, event.CodeSourceUploadDurationMs) + assert.GreaterOrEqual(t, *event.CodeSourceUploadDurationMs, int64(20)) + } else { + assert.Nil(t, event.CodeSourceSizeBytes) + assert.Nil(t, event.CodeSourceUsesGit) + assert.Empty(t, event.CodeSourcePackagingMode) + assert.Nil(t, event.CodeSourcePackagingDurationMs) + assert.Nil(t, event.CodeSourceUploadDurationMs) + } + raw, err := json.Marshal(event) + require.NoError(t, err) + assert.NotContains(t, string(raw), cfg.ExperimentName) + assert.NotContains(t, string(raw), "private source") + assert.NotContains(t, string(raw), "root_path") + }) + } +} + +func TestAirRunTelemetryConfig(t *testing.T) { + server := testserver.New(t) + t.Cleanup(server.Close) + w, err := databricks.NewWorkspaceClient(&databricks.Config{Host: server.URL, Token: "token"}) + require.NoError(t, err) + ctx, events := captureAirTelemetry(t, server, w) + cfgPath := writeConfigFile(t, "run.yaml", minimalConfig) + cfg, err := loadRunConfig(cfgPath) + require.NoError(t, err) + cfg.Compute.AcceleratorType = string(gpuType8xH100) + cfg.Compute.NumAccelerators = 16 + cfg.MaxRetries = new(0) + cfg.TimeoutMinutes = new(5) + cfg.Parameters = map[string]any{"private-parameter": "private-value"} + cfg.Environment = &environmentConfig{UnityCatalogImage: "private.schema.image:tag"} + logRunEvent(ctx, cfg, snapshotResult{ + PackagingMode: new(modeGitArchive), + PackagingDurationMs: new(int64(12)), + UploadDurationMs: new(int64(23)), + }, 123, 42*time.Millisecond, nil) + require.NoError(t, telemetry.Upload(ctx, protos.ExecutionContext{})) + require.Len(t, *events, 1) + event := (*events)[0].Entry.DatabricksCliLog.AirRunEvent + assert.Equal(t, 2, event.NumNodes) + assert.Equal(t, 0, event.MaxRetries) + assert.True(t, event.HasDockerImage) + assert.True(t, event.HasParameters) + assert.True(t, event.HasTimeout) + assert.Equal(t, int64(42), event.SubmitLatencyMs) + assert.Equal(t, protos.AirPackagingModeGitArchive, event.CodeSourcePackagingMode) + assert.Equal(t, new(int64(12)), event.CodeSourcePackagingDurationMs) + assert.Equal(t, new(int64(23)), event.CodeSourceUploadDurationMs) +} + +func TestSnapshotPackagingFailureMeasurements(t *testing.T) { + server := testserver.New(t) + t.Cleanup(server.Close) + testserver.AddDefaultHandlers(server) + w, err := databricks.NewWorkspaceClient(&databricks.Config{Host: server.URL, Token: "token"}) + require.NoError(t, err) + // The commit is no longer available when git archive attempts to package it. + result, err := uploadSnapshotTarball(t.Context(), w, newTestRepo(t), snapshotPlan{ + mode: modeGitArchive, commitSHA: "missing-commit", isGitRepo: true, + }, testSnapshotArtifactPath) + require.Error(t, err) + assert.Equal(t, new(modeGitArchive), result.PackagingMode) + require.NotNil(t, result.PackagingDurationMs) + assert.GreaterOrEqual(t, *result.PackagingDurationMs, int64(0)) + assert.Nil(t, result.UploadDurationMs) + assert.Nil(t, result.SizeBytes) + assert.Empty(t, result.CodeSourcePath) +} + +func TestTelemetryGPUTypeCoversSupportedTypes(t *testing.T) { + for _, g := range gpuTypes { + assert.NotEqual(t, protos.AirGPUTypeUnspecified, telemetryGPUType(g), "missing telemetry GPU type: %s", g) + } +} diff --git a/libs/telemetry/protos/README.md b/libs/telemetry/protos/README.md index 7dcc75e17cd..d760c256f9e 100644 --- a/libs/telemetry/protos/README.md +++ b/libs/telemetry/protos/README.md @@ -1,2 +1,25 @@ The types in this package are equivalent to the lumberjack protos defined in Universe. You can find all lumberjack protos for the Databricks CLI in the `proto/logs/frontend/databricks_cli` directory. + +AIR submission attempts use `databricks_cli_log.air_run_event`, mirrored by +`air_run.proto` in Universe. They record accelerator type/count, node count, +configuration-presence flags, retries, submission outcome and latency, Git usage, +and the successful Jobs run ID. `code_source_size_bytes` measures the compressed +tarball, using local file metadata for uploads and remote metadata for cache hits. +It is absent when no archive was produced or its size could not be measured. +Source contents, paths, dependency names, parameters, and image names are not sent. + +`code_source_packaging_mode` distinguishes `GIT_ARCHIVE` from `PLAIN_TAR`, including +on cache hits. `code_source_packaging_duration_ms` measures archive creation only, +excluding source discovery and cache lookup. `code_source_upload_duration_ms` +measures the tarball write, including directory creation and retries. Both are +wall-clock milliseconds and remain available when the attempted phase fails. +Cache hits explicitly report zero for both durations; phases never attempted +because of an earlier failure, and runs without code sources, omit the duration. + +The shared CLI logger uploads at command exit and respects +`DATABRICKS_CLI_DISABLE_TELEMETRY`. Dry runs do not emit a submission event. +Unlike Python AIR's `sgcli_log`, this event does not include `--via` attribution +or a separate first-log timing event, and watched submissions remain buffered +until the command exits. The Universe schema must land before ingestion can +retain the new event. diff --git a/libs/telemetry/protos/air_run.go b/libs/telemetry/protos/air_run.go new file mode 100644 index 00000000000..a3b50af7694 --- /dev/null +++ b/libs/telemetry/protos/air_run.go @@ -0,0 +1,46 @@ +package protos + +type AirGPUType string + +const ( + AirGPUTypeUnspecified AirGPUType = "GPU_TYPE_UNSPECIFIED" + AirGPUType1xA10 AirGPUType = "GPU_1X_A10" + AirGPUType1xH100 AirGPUType = "GPU_1X_H100" + AirGPUType8xH100 AirGPUType = "GPU_8X_H100" + AirGPUType8xB300 AirGPUType = "GPU_8X_B300" +) + +type AirPackagingMode string + +const ( + AirPackagingModeGitArchive AirPackagingMode = "GIT_ARCHIVE" + AirPackagingModePlainTar AirPackagingMode = "PLAIN_TAR" +) + +// AirRunEvent describes an AIR submission, matching the workload measurements +// collected by the Python AIR CLI. It contains no user-authored names or paths. +type AirRunEvent struct { + GPUType AirGPUType `json:"gpu_type"` + NumGPUs int `json:"num_gpus"` + NumNodes int `json:"num_nodes"` + HasDockerImage bool `json:"has_docker_image"` + HasCodeSnapshot bool `json:"has_code_snapshot"` + HasRequirements bool `json:"has_requirements"` + HasParameters bool `json:"has_parameters"` + MaxRetries int `json:"max_retries"` + HasTimeout bool `json:"has_timeout"` + SubmittedSuccessfully bool `json:"submitted_successfully"` + SubmitLatencyMs int64 `json:"submit_latency_ms"` + JobRunID string `json:"job_run_id,omitempty"` + CodeSourceUsesGit *bool `json:"code_source_uses_git,omitempty"` + + // Compressed tarball bytes, including when reusing a cached snapshot. + // Absent when no archive was measured; a measured zero remains explicit. + CodeSourceSizeBytes *int64 `json:"code_source_size_bytes,omitempty"` + + CodeSourcePackagingMode AirPackagingMode `json:"code_source_packaging_mode,omitempty"` + // Wall-clock milliseconds for archive creation and the tarball write, + // respectively. Zero on cache hits; absent if the phase was never attempted. + CodeSourcePackagingDurationMs *int64 `json:"code_source_packaging_duration_ms,omitempty"` + CodeSourceUploadDurationMs *int64 `json:"code_source_upload_duration_ms,omitempty"` +} diff --git a/libs/telemetry/protos/air_run_test.go b/libs/telemetry/protos/air_run_test.go new file mode 100644 index 00000000000..c809da479ec --- /dev/null +++ b/libs/telemetry/protos/air_run_test.go @@ -0,0 +1,34 @@ +package protos + +import ( + "encoding/json" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestAirRunEventOptionalMeasurements(t *testing.T) { + for _, size := range []*int64{nil, new(int64(0)), new(int64(123))} { + raw, err := json.Marshal(AirRunEvent{ + CodeSourceSizeBytes: size, + CodeSourcePackagingDurationMs: size, + CodeSourceUploadDurationMs: size, + }) + require.NoError(t, err) + var fields map[string]json.RawMessage + require.NoError(t, json.Unmarshal(raw, &fields)) + for _, name := range []string{"code_source_size_bytes", "code_source_packaging_duration_ms", "code_source_upload_duration_ms"} { + if size == nil { + assert.NotContains(t, fields, name) + } else { + var got int64 + require.NoError(t, json.Unmarshal(fields[name], &got)) + assert.Equal(t, *size, got) + } + } + assert.Equal(t, "0", string(fields["num_gpus"])) + assert.Equal(t, "0", string(fields["max_retries"])) + assert.Equal(t, "false", string(fields["submitted_successfully"])) + } +} diff --git a/libs/telemetry/protos/frontend_log.go b/libs/telemetry/protos/frontend_log.go index 50182765437..197856bab05 100644 --- a/libs/telemetry/protos/frontend_log.go +++ b/libs/telemetry/protos/frontend_log.go @@ -25,4 +25,5 @@ type DatabricksCliLog struct { BundleConfigRemoteSyncEvent *BundleConfigRemoteSyncEvent `json:"bundle_config_remote_sync_event,omitempty"` AitoolsInstallEvent *AitoolsInstallEvent `json:"aitools_install_event,omitempty"` SetupLocalEvent *SetupLocalEvent `json:"setup_local_event,omitempty"` + AirRunEvent *AirRunEvent `json:"air_run_event,omitempty"` } diff --git a/libs/testserver/fake_workspace.go b/libs/testserver/fake_workspace.go index b73803751b6..f709b3d3aea 100644 --- a/libs/testserver/fake_workspace.go +++ b/libs/testserver/fake_workspace.go @@ -660,6 +660,9 @@ func (s *FakeWorkspace) WorkspaceGetStatus(requestPath string, returnGitInfo boo info = dirInfo } else if entry, ok := s.files[cleaned]; ok { info = entry.Info + if info.ObjectType == workspace.ObjectTypeFile { + info.Size = int64(len(entry.Data)) + } } else { // Match the real Workspace API wording, which echoes the requested path. return Response{ From 097f6851d179cc9b2472ff1cd0f10bdc1c1aebe2 Mon Sep 17 00:00:00 2001 From: "ben.hansen" Date: Wed, 7 Oct 2026 21:25:43 +0000 Subject: [PATCH 2/5] [air] Align AIR run telemetry with the landed Universe schema The landed air_run.proto defines gpu_type as an optional string holding the API's canonical accelerator name, not an enum, so send GPU_1xH100-style values and omit the field when unresolved. Count rank-partitioned containers as using a Unity Catalog image, since each container must set one. Co-authored-by: Isaac --- cmd/air/telemetry.go | 19 ++-------------- cmd/air/telemetry_test.go | 20 ++++++++++++----- libs/telemetry/protos/README.md | 3 +-- libs/telemetry/protos/air_run.go | 37 ++++++++++++-------------------- 4 files changed, 32 insertions(+), 47 deletions(-) diff --git a/cmd/air/telemetry.go b/cmd/air/telemetry.go index d56dbabd63e..0b016c357eb 100644 --- a/cmd/air/telemetry.go +++ b/cmd/air/telemetry.go @@ -15,10 +15,10 @@ func logRunEvent(ctx context.Context, cfg *runConfig, snap snapshotResult, runID perNode, _ := gpusPerNode(gpuType(cfg.Compute.AcceleratorType)) _, hasDependencies := cfg.inlineDependencies() event := &protos.AirRunEvent{ - GPUType: telemetryGPUType(gpuType(cfg.Compute.AcceleratorType)), + GPUType: cfg.Compute.AcceleratorType, NumGPUs: cfg.Compute.NumAccelerators, NumNodes: cfg.Compute.NumAccelerators / perNode, - HasDockerImage: cfg.unityCatalogImagePath() != "", + HasDockerImage: cfg.unityCatalogImagePath() != "" || len(cfg.Containers) > 0, // every container sets its own image HasCodeSnapshot: cfg.CodeSource != nil && cfg.CodeSource.Snapshot != nil, HasRequirements: hasDependencies, HasParameters: len(cfg.Parameters) > 0, @@ -44,18 +44,3 @@ func logRunEvent(ctx context.Context, cfg *runConfig, snap snapshotResult, runID } telemetry.Log(ctx, protos.DatabricksCliLog{AirRunEvent: event}) } - -func telemetryGPUType(g gpuType) protos.AirGPUType { - switch g { - case gpuType1xA10: - return protos.AirGPUType1xA10 - case gpuType1xH100: - return protos.AirGPUType1xH100 - case gpuType8xH100: - return protos.AirGPUType8xH100 - case gpuType8xB300: - return protos.AirGPUType8xB300 - default: - return protos.AirGPUTypeUnspecified - } -} diff --git a/cmd/air/telemetry_test.go b/cmd/air/telemetry_test.go index 00b313a2d17..35a0cad984d 100644 --- a/cmd/air/telemetry_test.go +++ b/cmd/air/telemetry_test.go @@ -101,7 +101,7 @@ func TestAirRunTelemetry(t *testing.T) { require.NotNil(t, event) assert.Equal(t, !tc.fail && !tc.failUpload, event.SubmittedSuccessfully) assert.Equal(t, tc.snapshot, event.HasCodeSnapshot) - assert.Equal(t, protos.AirGPUType1xH100, event.GPUType) + assert.Equal(t, string(gpuType1xH100), event.GPUType) assert.Equal(t, 1, event.NumGPUs) assert.Equal(t, 1, event.NumNodes) assert.GreaterOrEqual(t, event.SubmitLatencyMs, int64(0)) @@ -190,8 +190,18 @@ func TestSnapshotPackagingFailureMeasurements(t *testing.T) { assert.Empty(t, result.CodeSourcePath) } -func TestTelemetryGPUTypeCoversSupportedTypes(t *testing.T) { - for _, g := range gpuTypes { - assert.NotEqual(t, protos.AirGPUTypeUnspecified, telemetryGPUType(g), "missing telemetry GPU type: %s", g) - } +func TestAirRunTelemetryContainersHaveDockerImage(t *testing.T) { + server := testserver.New(t) + t.Cleanup(server.Close) + w, err := databricks.NewWorkspaceClient(&databricks.Config{Host: server.URL, Token: "token"}) + require.NoError(t, err) + ctx, events := captureAirTelemetry(t, server, w) + cfg, err := loadRunConfig(writeConfigFile(t, "run.yaml", minimalConfig)) + require.NoError(t, err) + cfg.Command = nil + cfg.Containers = []containerConfig{{Name: "trainer", UnityCatalogImage: "private.schema.image:tag"}} + logRunEvent(ctx, cfg, snapshotResult{}, 123, time.Millisecond, nil) + require.NoError(t, telemetry.Upload(ctx, protos.ExecutionContext{})) + require.Len(t, *events, 1) + assert.True(t, (*events)[0].Entry.DatabricksCliLog.AirRunEvent.HasDockerImage) } diff --git a/libs/telemetry/protos/README.md b/libs/telemetry/protos/README.md index d760c256f9e..8e6c3055d6b 100644 --- a/libs/telemetry/protos/README.md +++ b/libs/telemetry/protos/README.md @@ -21,5 +21,4 @@ The shared CLI logger uploads at command exit and respects `DATABRICKS_CLI_DISABLE_TELEMETRY`. Dry runs do not emit a submission event. Unlike Python AIR's `sgcli_log`, this event does not include `--via` attribution or a separate first-log timing event, and watched submissions remain buffered -until the command exits. The Universe schema must land before ingestion can -retain the new event. +until the command exits. diff --git a/libs/telemetry/protos/air_run.go b/libs/telemetry/protos/air_run.go index a3b50af7694..e956a79eb51 100644 --- a/libs/telemetry/protos/air_run.go +++ b/libs/telemetry/protos/air_run.go @@ -1,15 +1,5 @@ package protos -type AirGPUType string - -const ( - AirGPUTypeUnspecified AirGPUType = "GPU_TYPE_UNSPECIFIED" - AirGPUType1xA10 AirGPUType = "GPU_1X_A10" - AirGPUType1xH100 AirGPUType = "GPU_1X_H100" - AirGPUType8xH100 AirGPUType = "GPU_8X_H100" - AirGPUType8xB300 AirGPUType = "GPU_8X_B300" -) - type AirPackagingMode string const ( @@ -20,19 +10,20 @@ const ( // AirRunEvent describes an AIR submission, matching the workload measurements // collected by the Python AIR CLI. It contains no user-authored names or paths. type AirRunEvent struct { - GPUType AirGPUType `json:"gpu_type"` - NumGPUs int `json:"num_gpus"` - NumNodes int `json:"num_nodes"` - HasDockerImage bool `json:"has_docker_image"` - HasCodeSnapshot bool `json:"has_code_snapshot"` - HasRequirements bool `json:"has_requirements"` - HasParameters bool `json:"has_parameters"` - MaxRetries int `json:"max_retries"` - HasTimeout bool `json:"has_timeout"` - SubmittedSuccessfully bool `json:"submitted_successfully"` - SubmitLatencyMs int64 `json:"submit_latency_ms"` - JobRunID string `json:"job_run_id,omitempty"` - CodeSourceUsesGit *bool `json:"code_source_uses_git,omitempty"` + // Canonical API accelerator type name, e.g. GPU_1xH100. + GPUType string `json:"gpu_type,omitempty"` + NumGPUs int `json:"num_gpus"` + NumNodes int `json:"num_nodes"` + HasDockerImage bool `json:"has_docker_image"` + HasCodeSnapshot bool `json:"has_code_snapshot"` + HasRequirements bool `json:"has_requirements"` + HasParameters bool `json:"has_parameters"` + MaxRetries int `json:"max_retries"` + HasTimeout bool `json:"has_timeout"` + SubmittedSuccessfully bool `json:"submitted_successfully"` + SubmitLatencyMs int64 `json:"submit_latency_ms"` + JobRunID string `json:"job_run_id,omitempty"` + CodeSourceUsesGit *bool `json:"code_source_uses_git,omitempty"` // Compressed tarball bytes, including when reusing a cached snapshot. // Absent when no archive was measured; a measured zero remains explicit. From 856fa6bec80a90c33166fd53f2e1a9672ef4d61d Mon Sep 17 00:00:00 2001 From: "ben.hansen" Date: Wed, 7 Oct 2026 22:59:10 +0000 Subject: [PATCH 3/5] [air] Shorten AIR telemetry README section Co-authored-by: Isaac --- libs/telemetry/protos/README.md | 25 ++++--------------------- 1 file changed, 4 insertions(+), 21 deletions(-) diff --git a/libs/telemetry/protos/README.md b/libs/telemetry/protos/README.md index 8e6c3055d6b..67964fc89d2 100644 --- a/libs/telemetry/protos/README.md +++ b/libs/telemetry/protos/README.md @@ -1,24 +1,7 @@ The types in this package are equivalent to the lumberjack protos defined in Universe. You can find all lumberjack protos for the Databricks CLI in the `proto/logs/frontend/databricks_cli` directory. -AIR submission attempts use `databricks_cli_log.air_run_event`, mirrored by -`air_run.proto` in Universe. They record accelerator type/count, node count, -configuration-presence flags, retries, submission outcome and latency, Git usage, -and the successful Jobs run ID. `code_source_size_bytes` measures the compressed -tarball, using local file metadata for uploads and remote metadata for cache hits. -It is absent when no archive was produced or its size could not be measured. -Source contents, paths, dependency names, parameters, and image names are not sent. - -`code_source_packaging_mode` distinguishes `GIT_ARCHIVE` from `PLAIN_TAR`, including -on cache hits. `code_source_packaging_duration_ms` measures archive creation only, -excluding source discovery and cache lookup. `code_source_upload_duration_ms` -measures the tarball write, including directory creation and retries. Both are -wall-clock milliseconds and remain available when the attempted phase fails. -Cache hits explicitly report zero for both durations; phases never attempted -because of an earlier failure, and runs without code sources, omit the duration. - -The shared CLI logger uploads at command exit and respects -`DATABRICKS_CLI_DISABLE_TELEMETRY`. Dry runs do not emit a submission event. -Unlike Python AIR's `sgcli_log`, this event does not include `--via` attribution -or a separate first-log timing event, and watched submissions remain buffered -until the command exits. +`databricks air run` logs one `air_run_event` per submission: requested compute, +which config options were set, the outcome and latency, and how long packaging and +uploading the code snapshot took. It never includes code, paths, names, parameters, +or image references. Set `DATABRICKS_CLI_DISABLE_TELEMETRY=1` to opt out. From 365d7014294026fbfd946e9f93773cc2d13f2b44 Mon Sep 17 00:00:00 2001 From: "ben.hansen" Date: Thu, 8 Oct 2026 23:01:18 +0000 Subject: [PATCH 4/5] [air] Omit snapshot size on cache hits A cache hit uploads nothing, so leave code_source_size_bytes unset instead of reading the remote file size. This drops the fake workspace change that only existed to support it, and the JSON marshalling test in protos. Co-authored-by: Isaac --- cmd/air/runsubmit_test.go | 4 ++-- cmd/air/snapshot_upload.go | 4 +--- libs/telemetry/protos/air_run.go | 4 ++-- libs/telemetry/protos/air_run_test.go | 34 --------------------------- libs/testserver/fake_workspace.go | 3 --- 5 files changed, 5 insertions(+), 44 deletions(-) delete mode 100644 libs/telemetry/protos/air_run_test.go diff --git a/cmd/air/runsubmit_test.go b/cmd/air/runsubmit_test.go index 8a0f607b09d..b24740317af 100644 --- a/cmd/air/runsubmit_test.go +++ b/cmd/air/runsubmit_test.go @@ -783,7 +783,7 @@ code_source: assert.Equal(t, afterFirst, snapshotUploads, "unchanged plain_tar should skip the second upload") require.NotNil(t, first.SizeBytes) assert.Positive(t, *first.SizeBytes) - assert.Equal(t, first.SizeBytes, second.SizeBytes) + assert.Nil(t, second.SizeBytes) assert.Equal(t, new(false), second.UsesGit) assert.Equal(t, new(modePlainTar), first.PackagingMode) assert.Equal(t, first.PackagingMode, second.PackagingMode) @@ -849,7 +849,7 @@ code_source: assert.Len(t, uploaded, 1, "git_archive cache hit should skip the second upload") require.NotNil(t, first.SizeBytes) assert.Positive(t, *first.SizeBytes) - assert.Equal(t, first.SizeBytes, second.SizeBytes) + assert.Nil(t, second.SizeBytes) assert.Equal(t, new(true), second.UsesGit) assert.Equal(t, new(modeGitArchive), first.PackagingMode) assert.Equal(t, first.PackagingMode, second.PackagingMode) diff --git a/cmd/air/snapshot_upload.go b/cmd/air/snapshot_upload.go index b36de94be35..53816cdd926 100644 --- a/cmd/air/snapshot_upload.go +++ b/cmd/air/snapshot_upload.go @@ -172,12 +172,10 @@ func uploadSnapshotTarball(ctx context.Context, w *databricks.WorkspaceClient, r // upload the bytes now or reuse an object already in the store. remote := path.Join(uploadPath, tarName) - info, err := f.Stat(ctx, tarName) + _, err = f.Stat(ctx, tarName) if err == nil { log.Debugf(ctx, "snapshot upload skipped; reusing %s", remote) - size := info.Size() result.CodeSourcePath = remote - result.SizeBytes = &size result.PackagingDurationMs = new(int64(0)) result.UploadDurationMs = new(int64(0)) return result, nil diff --git a/libs/telemetry/protos/air_run.go b/libs/telemetry/protos/air_run.go index e956a79eb51..b8f5d02a0fe 100644 --- a/libs/telemetry/protos/air_run.go +++ b/libs/telemetry/protos/air_run.go @@ -25,8 +25,8 @@ type AirRunEvent struct { JobRunID string `json:"job_run_id,omitempty"` CodeSourceUsesGit *bool `json:"code_source_uses_git,omitempty"` - // Compressed tarball bytes, including when reusing a cached snapshot. - // Absent when no archive was measured; a measured zero remains explicit. + // Compressed tarball bytes for a new upload. Absent on cache hits and + // when no archive was measured; a measured zero remains explicit. CodeSourceSizeBytes *int64 `json:"code_source_size_bytes,omitempty"` CodeSourcePackagingMode AirPackagingMode `json:"code_source_packaging_mode,omitempty"` diff --git a/libs/telemetry/protos/air_run_test.go b/libs/telemetry/protos/air_run_test.go deleted file mode 100644 index c809da479ec..00000000000 --- a/libs/telemetry/protos/air_run_test.go +++ /dev/null @@ -1,34 +0,0 @@ -package protos - -import ( - "encoding/json" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -func TestAirRunEventOptionalMeasurements(t *testing.T) { - for _, size := range []*int64{nil, new(int64(0)), new(int64(123))} { - raw, err := json.Marshal(AirRunEvent{ - CodeSourceSizeBytes: size, - CodeSourcePackagingDurationMs: size, - CodeSourceUploadDurationMs: size, - }) - require.NoError(t, err) - var fields map[string]json.RawMessage - require.NoError(t, json.Unmarshal(raw, &fields)) - for _, name := range []string{"code_source_size_bytes", "code_source_packaging_duration_ms", "code_source_upload_duration_ms"} { - if size == nil { - assert.NotContains(t, fields, name) - } else { - var got int64 - require.NoError(t, json.Unmarshal(fields[name], &got)) - assert.Equal(t, *size, got) - } - } - assert.Equal(t, "0", string(fields["num_gpus"])) - assert.Equal(t, "0", string(fields["max_retries"])) - assert.Equal(t, "false", string(fields["submitted_successfully"])) - } -} diff --git a/libs/testserver/fake_workspace.go b/libs/testserver/fake_workspace.go index f709b3d3aea..b73803751b6 100644 --- a/libs/testserver/fake_workspace.go +++ b/libs/testserver/fake_workspace.go @@ -660,9 +660,6 @@ func (s *FakeWorkspace) WorkspaceGetStatus(requestPath string, returnGitInfo boo info = dirInfo } else if entry, ok := s.files[cleaned]; ok { info = entry.Info - if info.ObjectType == workspace.ObjectTypeFile { - info.Size = int64(len(entry.Data)) - } } else { // Match the real Workspace API wording, which echoes the requested path. return Response{ From eda943aae053714159116e54d009177f14090d1e Mon Sep 17 00:00:00 2001 From: "ben.hansen" Date: Thu, 8 Oct 2026 23:10:14 +0000 Subject: [PATCH 5/5] [air] Restore snapshotExists The cache check was inlined only to read the remote size, which cache hits no longer report. Co-authored-by: Isaac --- cmd/air/snapshot_upload.go | 24 +++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/cmd/air/snapshot_upload.go b/cmd/air/snapshot_upload.go index 53816cdd926..1d4e9a533b8 100644 --- a/cmd/air/snapshot_upload.go +++ b/cmd/air/snapshot_upload.go @@ -172,17 +172,17 @@ func uploadSnapshotTarball(ctx context.Context, w *databricks.WorkspaceClient, r // upload the bytes now or reuse an object already in the store. remote := path.Join(uploadPath, tarName) - _, err = f.Stat(ctx, tarName) - if err == nil { + exists, err := snapshotExists(ctx, f, tarName) + if err != nil { + return result, err + } + if exists { log.Debugf(ctx, "snapshot upload skipped; reusing %s", remote) result.CodeSourcePath = remote result.PackagingDurationMs = new(int64(0)) result.UploadDurationMs = new(int64(0)) return result, nil } - if !errors.Is(err, fs.ErrNotExist) { - return result, fmt.Errorf("failed to check snapshot cache: %w", err) - } tmp, err := os.MkdirTemp("", "air-snapshot-*") if err != nil { @@ -241,3 +241,17 @@ func snapshotUploadFiler(ctx context.Context, w *databricks.WorkspaceClient, art f, err := filer.NewWorkspaceFilesClient(w, uploadPath) return f, uploadPath, err } + +// snapshotExists reports whether name already exists in the artifact store, used to +// short-circuit a content-addressed upload (either mode). A not-found is a clean miss +// (false, nil); any other error is surfaced. +func snapshotExists(ctx context.Context, store filer.Filer, name string) (bool, error) { + _, err := store.Stat(ctx, name) + if err == nil { + return true, nil + } + if errors.Is(err, fs.ErrNotExist) { + return false, nil + } + return false, fmt.Errorf("failed to check snapshot cache: %w", err) +}