From 94cea7e3df157b68d1063ee14a1b3c372391ebdd Mon Sep 17 00:00:00 2001 From: James Robinson Date: Tue, 4 Aug 2026 14:19:29 -0400 Subject: [PATCH 1/4] Ability to sumbit Flink statements in snapshot mode. --- .claude/rules/testing/unit-tests.md | 29 ++++++ .vscode/launch.json | 4 +- src/codelens/flinkSqlProvider.test.ts | 105 ++++++++++++++++++--- src/codelens/flinkSqlProvider.ts | 32 ++++++- src/commands/documents.test.ts | 63 +++++++++++++ src/commands/documents.ts | 26 +++++ src/commands/flinkStatements.ts | 9 +- src/flinkSql/statementUtils.test.ts | 56 ++++++++++- src/flinkSql/statementUtils.ts | 26 ++++- src/models/flinkStatement.test.ts | 19 ++++ src/models/flinkStatement.ts | 22 +++++ src/storage/constants.ts | 2 + tests/unit/testResources/flinkStatement.ts | 5 +- 13 files changed, 376 insertions(+), 22 deletions(-) diff --git a/.claude/rules/testing/unit-tests.md b/.claude/rules/testing/unit-tests.md index b09388e2b4..6cc291409e 100644 --- a/.claude/rules/testing/unit-tests.md +++ b/.claude/rules/testing/unit-tests.md @@ -51,3 +51,32 @@ assignment instead: `obj["methodName"] = sandbox.stub()`. - Unit test fixtures in `tests/fixtures/` - Shared stubs in `tests/stubs/` + +## macOS: Local Electron Test Host Won't Launch (e.g. VS Code 1.131.0) + +`npx gulp test` downloads a real VS Code build into `.vscode-test/vscode-darwin-arm64-/` +(gitignored cache) and spawns it via `@vscode/test-electron@2.3.9` as the Extension Development +Host. Two independent problems show up together on newer VS Code releases (first seen with 1.131.0): + +**1. `spawn .../Electron ENOENT`** — `@vscode/test-electron@2.3.9` hardcodes the launcher binary +name as `Electron`, but this VS Code build ships it renamed to `Code`. Fix locally (never commit — +this only touches the gitignored cache): + +```bash +cd ".vscode-test/vscode-darwin-arm64-/Visual Studio Code.app/Contents/MacOS/" +ln -s Code Electron +``` + +**2. `"Visual Studio Code" is damaged and can't be opened`** — adding that symlink modifies the +signed `.app` bundle's contents, which invalidates its code signature. macOS then refuses to launch +it and reports it as "damaged" (this is a broken-signature error, not a Gatekeeper quarantine issue +— no `com.apple.quarantine` xattr needs to be present for it to happen). Fix by re-signing ad hoc +after adding the symlink: + +```bash +codesign --force --deep --sign - "Visual Studio Code.app" +``` + +Re-run `npx gulp test` after both steps. If it still fails with `SIGKILL` or an Electron "bad +option" error, that's a genuine headless/no-GUI-session environment (e.g. a sandboxed agent session +without display access) — not a code or signing issue, and not fixable with the above. diff --git a/.vscode/launch.json b/.vscode/launch.json index de0df51373..51c097ff6c 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -10,7 +10,9 @@ "type": "extensionHost", "request": "launch", "args": [ - "--extensionDevelopmentPath=${workspaceFolder}/out" + "--extensionDevelopmentPath=${workspaceFolder}/out", + "--disable-extension", + "confluentinc.vscode-confluent" ], "outFiles": [ "${workspaceFolder}/out/**/*.js" diff --git a/src/codelens/flinkSqlProvider.test.ts b/src/codelens/flinkSqlProvider.test.ts index 0a8aac4bf3..ace1e98b6f 100644 --- a/src/codelens/flinkSqlProvider.test.ts +++ b/src/codelens/flinkSqlProvider.test.ts @@ -14,6 +14,7 @@ import { FLINK_CONFIG_COMPUTE_POOL, FLINK_CONFIG_DATABASE } from "../extensionSe import type { CCloudResourceLoader } from "../loaders"; import { CCloudEnvironment } from "../models/environment"; import { CCloudFlinkComputePool } from "../models/flinkComputePool"; +import { FlinkSnapshotMode } from "../models/flinkStatement"; import { CCloudKafkaCluster } from "../models/kafkaCluster"; import type { EnvironmentId } from "../models/resource"; import * as ccloud from "../sidecar/connections/ccloud"; @@ -27,6 +28,7 @@ import { getCatalogDatabaseFromMetadata, getComputePoolFromMetadata, getDefaultCatalogDatabase, + getSnapshotModeFromMetadata, } from "./flinkSqlProvider"; const testUri = Uri.parse("file:///test/file.sql"); @@ -202,11 +204,22 @@ describe("codelens/flinkSqlProvider.ts", () => { const codeLenses: CodeLens[] = await provider.provideCodeLenses(fakeDocument); - assert.strictEqual(codeLenses.length, 3); + assert.strictEqual(codeLenses.length, 4); - const poolLens = codeLenses[0]; - const dbLens = codeLenses[1]; - const resetLens = codeLenses[2]; + const snapshotModeLens = codeLenses[0]; + const poolLens = codeLenses[1]; + const dbLens = codeLenses[2]; + const resetLens = codeLenses[3]; + + assert.strictEqual( + snapshotModeLens.command?.command, + "confluent.document.flinksql.toggleSnapshotMode", + ); + assert.strictEqual(snapshotModeLens.command?.title, "Mode: Streaming"); + assert.deepStrictEqual(snapshotModeLens.command?.arguments, [ + fakeDocument.uri, + FlinkSnapshotMode.STREAMING, + ]); assert.strictEqual( dbLens.command?.command, @@ -244,11 +257,18 @@ describe("codelens/flinkSqlProvider.ts", () => { const codeLenses: CodeLens[] = await provider.provideCodeLenses(fakeDocument); - assert.strictEqual(codeLenses.length, 3); + assert.strictEqual(codeLenses.length, 4); - const poolLens = codeLenses[0]; - const dbLens = codeLenses[1]; - const resetLens = codeLenses[2]; + const snapshotModeLens = codeLenses[0]; + const poolLens = codeLenses[1]; + const dbLens = codeLenses[2]; + const resetLens = codeLenses[3]; + + assert.strictEqual( + snapshotModeLens.command?.command, + "confluent.document.flinksql.toggleSnapshotMode", + ); + assert.strictEqual(snapshotModeLens.command?.title, "Mode: Streaming"); assert.strictEqual( dbLens.command?.command, @@ -291,17 +311,28 @@ describe("codelens/flinkSqlProvider.ts", () => { const codeLenses: CodeLens[] = await provider.provideCodeLenses(fakeDocument); - assert.strictEqual(codeLenses.length, 4); + assert.strictEqual(codeLenses.length, 5); const submitLens = codeLenses[0]; - const poolLens = codeLenses[1]; - const dbLens = codeLenses[2]; - const resetLens = codeLenses[3]; + const snapshotModeLens = codeLenses[1]; + const poolLens = codeLenses[2]; + const dbLens = codeLenses[3]; + const resetLens = codeLenses[4]; assert.strictEqual(submitLens.command?.command, "confluent.statements.create"); assert.strictEqual(submitLens.command?.title, "▶️ Submit Statement"); assert.deepStrictEqual(submitLens.command?.arguments, [fakeDocument.uri, pool, database]); + assert.strictEqual( + snapshotModeLens.command?.command, + "confluent.document.flinksql.toggleSnapshotMode", + ); + assert.strictEqual(snapshotModeLens.command?.title, "Mode: Streaming"); + assert.deepStrictEqual(snapshotModeLens.command?.arguments, [ + fakeDocument.uri, + FlinkSnapshotMode.STREAMING, + ]); + assert.strictEqual(dbLens.command?.command, "confluent.document.flinksql.setCCloudDatabase"); assert.strictEqual( dbLens.command?.title, @@ -323,6 +354,56 @@ describe("codelens/flinkSqlProvider.ts", () => { assert.strictEqual(resetLens.command?.title, "Clear Settings"); assert.deepStrictEqual(resetLens.command?.arguments, [fakeDocument.uri]); }); + + it("should provide 'Mode: Snapshot' codelens when snapshot mode metadata is set to BATCH", async () => { + const pool: CCloudFlinkComputePool = TEST_CCLOUD_FLINK_COMPUTE_POOL; + const database: CCloudKafkaCluster = TEST_CCLOUD_KAFKA_CLUSTER; + resourceManagerStub.getUriMetadata.resolves({ + [UriMetadataKeys.FLINK_COMPUTE_POOL_ID]: pool.id, + [UriMetadataKeys.FLINK_CATALOG_ID]: TEST_CCLOUD_ENVIRONMENT.id, + [UriMetadataKeys.FLINK_CATALOG_NAME]: TEST_CCLOUD_ENVIRONMENT.name, + [UriMetadataKeys.FLINK_DATABASE_ID]: database.id, + [UriMetadataKeys.FLINK_DATABASE_NAME]: database.name, + [UriMetadataKeys.FLINK_SNAPSHOT_MODE]: FlinkSnapshotMode.BATCH, + }); + ccloudLoaderStub.getEnvironments.resolves([testEnvWithPoolAndCluster]); + + const codeLenses: CodeLens[] = await provider.provideCodeLenses(fakeDocument); + + const snapshotModeLens = codeLenses[1]; + assert.strictEqual( + snapshotModeLens.command?.command, + "confluent.document.flinksql.toggleSnapshotMode", + ); + assert.strictEqual(snapshotModeLens.command?.title, "Mode: Snapshot"); + assert.deepStrictEqual(snapshotModeLens.command?.arguments, [ + fakeDocument.uri, + FlinkSnapshotMode.BATCH, + ]); + }); + }); + + describe("getSnapshotModeFromMetadata()", () => { + it("should return STREAMING when metadata is undefined", () => { + assert.strictEqual(getSnapshotModeFromMetadata(undefined), FlinkSnapshotMode.STREAMING); + }); + + it("should return STREAMING when snapshot mode metadata is unset", () => { + const metadata: UriMetadata = { [UriMetadataKeys.FLINK_COMPUTE_POOL_ID]: "pool-123" }; + assert.strictEqual(getSnapshotModeFromMetadata(metadata), FlinkSnapshotMode.STREAMING); + }); + + it("should return STREAMING when snapshot mode metadata is explicitly cleared (null)", () => { + const metadata: UriMetadata = { [UriMetadataKeys.FLINK_SNAPSHOT_MODE]: null }; + assert.strictEqual(getSnapshotModeFromMetadata(metadata), FlinkSnapshotMode.STREAMING); + }); + + it("should return BATCH when snapshot mode metadata is set to BATCH", () => { + const metadata: UriMetadata = { + [UriMetadataKeys.FLINK_SNAPSHOT_MODE]: FlinkSnapshotMode.BATCH, + }; + assert.strictEqual(getSnapshotModeFromMetadata(metadata), FlinkSnapshotMode.BATCH); + }); }); describe("getComputePoolFromMetadata()", () => { diff --git a/src/codelens/flinkSqlProvider.ts b/src/codelens/flinkSqlProvider.ts index 94d3f2cba1..d7f073303e 100644 --- a/src/codelens/flinkSqlProvider.ts +++ b/src/codelens/flinkSqlProvider.ts @@ -6,6 +6,7 @@ import { CCloudResourceLoader } from "../loaders"; import { Logger } from "../logging"; import type { CCloudEnvironment } from "../models/environment"; import type { CCloudFlinkComputePool } from "../models/flinkComputePool"; +import { FlinkSnapshotMode } from "../models/flinkStatement"; import type { CCloudKafkaCluster } from "../models/kafkaCluster"; import { hasCCloudAuthSession } from "../sidecar/connections/ccloud"; import { UriMetadataKeys } from "../storage/constants"; @@ -87,6 +88,20 @@ export class FlinkSqlCodelensProvider extends DisposableCollection implements Co const computePool: CCloudFlinkComputePool | undefined = await getComputePoolFromMetadata(uriMetadata); const { catalog, database } = await getCatalogDatabaseFromMetadata(uriMetadata, computePool); + const snapshotMode: FlinkSnapshotMode = getSnapshotModeFromMetadata(uriMetadata); + + // codelens for toggling between streaming (default, continuous) and snapshot (one-shot, + // bounded) statement execution mode + const isSnapshotMode = snapshotMode === FlinkSnapshotMode.BATCH; + const toggleSnapshotModeCommand: Command = { + title: isSnapshotMode ? "Mode: Snapshot" : "Mode: Streaming", + command: "confluent.document.flinksql.toggleSnapshotMode", + tooltip: isSnapshotMode + ? "Statement will run once against a snapshot of current data, then complete. Click to switch to streaming mode." + : "Statement will run continuously against new data as it arrives. Click to switch to one-time snapshot mode.", + arguments: [document.uri, snapshotMode], + }; + const snapshotModeLens = new CodeLens(range, toggleSnapshotModeCommand); // codelens for selecting a compute pool, which we'll use to derive the rest of the properties // needed for various Flink operations (env ID, provider/region, etc) @@ -129,11 +144,11 @@ export class FlinkSqlCodelensProvider extends DisposableCollection implements Co arguments: [document.uri, computePool, database], }; const submitLens = new CodeLens(range, submitCommand); - // show the "Submit Statement" | | codelenses - codeLenses.push(submitLens, computePoolLens, databaseLens, resetLens); + // show the "Submit Statement" | mode | | codelenses + codeLenses.push(submitLens, snapshotModeLens, computePoolLens, databaseLens, resetLens); } else { // don't show the submit codelens if we don't have a compute pool and database - codeLenses.push(computePoolLens, databaseLens, resetLens); + codeLenses.push(snapshotModeLens, computePoolLens, databaseLens, resetLens); } return codeLenses; @@ -162,6 +177,17 @@ export async function getComputePoolFromMetadata( return await loader.getFlinkComputePool(computePoolId); } +/** + * Get the snapshot ("batch") vs streaming execution mode from the metadata stored in the document. + * Defaults to {@link FlinkSnapshotMode.STREAMING} when unset or explicitly cleared. + * @param metadata The metadata stored in the document. + */ +export function getSnapshotModeFromMetadata(metadata: UriMetadata | undefined): FlinkSnapshotMode { + return metadata?.[UriMetadataKeys.FLINK_SNAPSHOT_MODE] === FlinkSnapshotMode.BATCH + ? FlinkSnapshotMode.BATCH + : FlinkSnapshotMode.STREAMING; +} + export interface CatalogDatabase { catalog?: CCloudEnvironment; database?: CCloudKafkaCluster; diff --git a/src/commands/documents.test.ts b/src/commands/documents.test.ts index 2a71c39139..538d483f08 100644 --- a/src/commands/documents.test.ts +++ b/src/commands/documents.test.ts @@ -13,6 +13,7 @@ import { } from "../extensionSettings/constants"; import * as statementUtils from "../flinkSql/statementUtils"; import type { CCloudResourceLoader } from "../loaders"; +import { FlinkSnapshotMode } from "../models/flinkStatement"; import * as flinkComputePoolsQuickPick from "../quickpicks/flinkComputePools"; import * as flinkDatabaseQuickpick from "../quickpicks/kafkaClusters"; import * as ccloudConnections from "../sidecar/connections/ccloud"; @@ -22,6 +23,7 @@ import { resetCCloudMetadataForUriCommand, setCCloudComputePoolForUriCommand, setCCloudDatabaseForUriCommand, + toggleSnapshotModeForUriCommand, } from "./documents"; const testUri = Uri.parse("file:///path/to/test.sql"); @@ -428,8 +430,69 @@ describe("commands/documents.ts resetCCloudMetadataForUriCommand()", () => { [UriMetadataKeys.FLINK_CATALOG_NAME]: null, [UriMetadataKeys.FLINK_DATABASE_ID]: null, [UriMetadataKeys.FLINK_DATABASE_NAME]: null, + [UriMetadataKeys.FLINK_SNAPSHOT_MODE]: null, }); sinon.assert.calledOnce(uriMetadataSetFireStub); sinon.assert.calledOnceWithExactly(uriMetadataSetFireStub, testUri); }); }); + +describe("commands/documents.ts toggleSnapshotModeForUriCommand()", () => { + let sandbox: sinon.SinonSandbox; + + let setFlinkDocumentMetadataStub: sinon.SinonStub; + let hasCCloudAuthSessionStub: sinon.SinonStub; + + beforeEach(() => { + sandbox = sinon.createSandbox(); + + setFlinkDocumentMetadataStub = sandbox + .stub(statementUtils, "setFlinkDocumentMetadata") + .resolves(); + hasCCloudAuthSessionStub = sandbox + .stub(ccloudConnections, "hasCCloudAuthSession") + .returns(true); + }); + + afterEach(() => { + sandbox.restore(); + }); + + it("should do nothing when no Uri is provided", async () => { + await toggleSnapshotModeForUriCommand(); + + sinon.assert.notCalled(setFlinkDocumentMetadataStub); + }); + + it("should do nothing when no CCloud auth session is available", async () => { + hasCCloudAuthSessionStub.returns(false); + + await toggleSnapshotModeForUriCommand(testUri); + + sinon.assert.notCalled(setFlinkDocumentMetadataStub); + }); + + it("should set mode to BATCH when current mode is undefined", async () => { + await toggleSnapshotModeForUriCommand(testUri, undefined); + + sinon.assert.calledOnceWithExactly(setFlinkDocumentMetadataStub, testUri, { + snapshotMode: FlinkSnapshotMode.BATCH, + }); + }); + + it("should set mode to BATCH when current mode is STREAMING", async () => { + await toggleSnapshotModeForUriCommand(testUri, FlinkSnapshotMode.STREAMING); + + sinon.assert.calledOnceWithExactly(setFlinkDocumentMetadataStub, testUri, { + snapshotMode: FlinkSnapshotMode.BATCH, + }); + }); + + it("should set mode to STREAMING when current mode is BATCH", async () => { + await toggleSnapshotModeForUriCommand(testUri, FlinkSnapshotMode.BATCH); + + sinon.assert.calledOnceWithExactly(setFlinkDocumentMetadataStub, testUri, { + snapshotMode: FlinkSnapshotMode.STREAMING, + }); + }); +}); diff --git a/src/commands/documents.ts b/src/commands/documents.ts index e745bd7df4..a4b8c3c46a 100644 --- a/src/commands/documents.ts +++ b/src/commands/documents.ts @@ -14,6 +14,7 @@ import { CCloudResourceLoader } from "../loaders"; import { Logger } from "../logging"; import type { CCloudEnvironment } from "../models/environment"; import { CCloudFlinkComputePool } from "../models/flinkComputePool"; +import { FlinkSnapshotMode } from "../models/flinkStatement"; import type { CCloudKafkaCluster } from "../models/kafkaCluster"; import { showInfoNotificationWithButtons } from "../notifications"; import { flinkComputePoolQuickPick } from "../quickpicks/flinkComputePools"; @@ -142,6 +143,26 @@ export async function setCCloudDatabaseForUriCommand(uri?: Uri, pool?: CCloudFli } } +export async function toggleSnapshotModeForUriCommand( + uri?: Uri, + currentMode?: FlinkSnapshotMode, +): Promise { + if (!(uri instanceof Uri)) { + return; + } + if (!hasCCloudAuthSession()) { + // shouldn't happen since callers shouldn't be able to call this command without a valid CCloud + // connection, but just in case + logger.warn("not toggling snapshot mode for URI: no CCloud auth session"); + return; + } + + const newMode: FlinkSnapshotMode = + currentMode === FlinkSnapshotMode.BATCH ? FlinkSnapshotMode.STREAMING : FlinkSnapshotMode.BATCH; + + await setFlinkDocumentMetadata(uri, { snapshotMode: newMode }); +} + export async function resetCCloudMetadataForUriCommand(uri?: Uri) { if (!(uri instanceof Uri)) { return; @@ -163,6 +184,7 @@ export async function resetCCloudMetadataForUriCommand(uri?: Uri) { [UriMetadataKeys.FLINK_CATALOG_NAME]: null, [UriMetadataKeys.FLINK_DATABASE_ID]: null, [UriMetadataKeys.FLINK_DATABASE_NAME]: null, + [UriMetadataKeys.FLINK_SNAPSHOT_MODE]: null, }); uriMetadataSet.fire(uri); } @@ -181,5 +203,9 @@ export function registerDocumentCommands(): Disposable[] { "confluent.document.flinksql.resetCCloudMetadata", resetCCloudMetadataForUriCommand, ), + registerCommandWithLogging( + "confluent.document.flinksql.toggleSnapshotMode", + toggleSnapshotModeForUriCommand, + ), ]; } diff --git a/src/commands/flinkStatements.ts b/src/commands/flinkStatements.ts index f2e8931fa3..8bf426dd15 100644 --- a/src/commands/flinkStatements.ts +++ b/src/commands/flinkStatements.ts @@ -1,6 +1,9 @@ import * as vscode from "vscode"; import { registerCommandWithLogging } from "."; -import { getCatalogDatabaseFromMetadata } from "../codelens/flinkSqlProvider"; +import { + getCatalogDatabaseFromMetadata, + getSnapshotModeFromMetadata, +} from "../codelens/flinkSqlProvider"; import { FLINKSTATEMENT_URI_SCHEME, FlinkStatementDocumentProvider, @@ -128,6 +131,7 @@ export async function submitFlinkStatementCommand( const document = await getEditorOrFileContents(statementBodyUri); const statement = document.content; const uriMetadata = await ResourceManager.getInstance().getUriMetadata(statementBodyUri); + const snapshotMode = getSnapshotModeFromMetadata(uriMetadata); // 2. Choose the statement name const statementName = await determineFlinkStatementName(); @@ -184,6 +188,7 @@ export async function submitFlinkStatementCommand( organizationId: organization.id, hidden: false, // Do not create a hidden statement, the user authored it. properties: currentDatabaseKafkaCluster.toFlinkSpecProperties(), + snapshotMode, }; const newStatement = await submitFlinkStatement(submission); @@ -197,6 +202,7 @@ export async function submitFlinkStatementCommand( compute_pool_id: computePool.id, failure_reason: newStatement.status.detail, from_flink_workspace: isFromFlinkWorkspace(uriMetadata), + snapshot_mode: snapshotMode, }); // limit the error message content so the notification isn't hidden automatically @@ -211,6 +217,7 @@ export async function submitFlinkStatementCommand( sql_kind: newStatement.sqlKind, compute_pool_id: computePool.id, from_flink_workspace: isFromFlinkWorkspace(uriMetadata), + snapshot_mode: snapshotMode, }); // Refresh the statements view onto the compute pool in question, diff --git a/src/flinkSql/statementUtils.test.ts b/src/flinkSql/statementUtils.test.ts index 04ee74d85f..7f08f9411b 100644 --- a/src/flinkSql/statementUtils.test.ts +++ b/src/flinkSql/statementUtils.test.ts @@ -16,6 +16,7 @@ import { import { TEST_CCLOUD_ORGANIZATION_ID } from "../../tests/unit/testResources/organization"; import { createResponseError } from "../../tests/unit/testUtils"; import type { + CreateSqlv1StatementOperationRequest, GetSqlv1Statement200Response, GetSqlv1StatementResult200Response, } from "../clients/flinkSql"; @@ -24,7 +25,7 @@ import { uriMetadataSet } from "../emitters"; import { FLINK_CONFIG_STATEMENT_PREFIX } from "../extensionSettings/constants"; import type { CCloudResourceLoader } from "../loaders"; import * as flinkStatementModels from "../models/flinkStatement"; -import { FlinkSpecProperties, FlinkStatement } from "../models/flinkStatement"; +import { FlinkSnapshotMode, FlinkSpecProperties, FlinkStatement } from "../models/flinkStatement"; import type { CCloudFlinkDbKafkaCluster } from "../models/kafkaCluster"; import type { EnvironmentId } from "../models/resource"; import type * as sidecar from "../sidecar"; @@ -252,6 +253,47 @@ describe("flinkSql/statementUtils.ts", function () { ); }); } + + const snapshotModeCases: Array<{ + snapshotMode: FlinkSnapshotMode | undefined; + expectedProperty: string | undefined; + }> = [ + { snapshotMode: undefined, expectedProperty: undefined }, + { snapshotMode: FlinkSnapshotMode.STREAMING, expectedProperty: undefined }, + { snapshotMode: FlinkSnapshotMode.BATCH, expectedProperty: "now" }, + ]; + + for (const { snapshotMode, expectedProperty } of snapshotModeCases) { + it(`sets "sql.snapshot.mode" property correctly for snapshotMode=${snapshotMode}`, async function () { + const params: IFlinkStatementSubmitParameters = { + statement: "SELECT * FROM my_table", + statementName: "test-statement", + computePool: TEST_CCLOUD_FLINK_COMPUTE_POOL, + organizationId: TEST_CCLOUD_ORGANIZATION_ID, + hidden: false, + properties: new FlinkSpecProperties({}), + snapshotMode, + }; + + const createSqlv1StatementStub = sandbox.stub().resolves(TEST_CCLOUD_FLINK_STATEMENT); + sandbox + .stub(flinkStatementModels, "restFlinkStatementToModel") + .returns(TEST_CCLOUD_FLINK_STATEMENT); + mockSidecar.getFlinkSqlStatementsApi.returns({ + createSqlv1Statement: createSqlv1StatementStub, + } as any); + + await submitFlinkStatement(params); + + sinon.assert.calledOnce(createSqlv1StatementStub); + const request = createSqlv1StatementStub.firstCall + .args[0] as CreateSqlv1StatementOperationRequest; + const spec = request.CreateSqlv1StatementRequest?.spec as { + properties?: Record; + }; + assert.strictEqual(spec.properties?.["sql.snapshot.mode"], expectedProperty); + }); + } }); describe("waitForStatement* tests", () => { @@ -488,6 +530,18 @@ describe("flinkSql/statementUtils.ts", function () { sinon.assert.calledWith(uriMetadataSetFireStub, uri); }); + + it("should set the snapshot mode metadata when provided", async () => { + await setFlinkDocumentMetadata(uri, { + snapshotMode: FlinkSnapshotMode.BATCH, + }); + + sinon.assert.calledWith(rmSetUriMetadataStub, uri, { + [UriMetadataKeys.FLINK_SNAPSHOT_MODE]: FlinkSnapshotMode.BATCH, + }); + + sinon.assert.calledWith(uriMetadataSetFireStub, uri); + }); }); describe("isFromFlinkWorkspace()", function () { diff --git a/src/flinkSql/statementUtils.ts b/src/flinkSql/statementUtils.ts index dd04870ec9..4689f90ded 100644 --- a/src/flinkSql/statementUtils.ts +++ b/src/flinkSql/statementUtils.ts @@ -15,7 +15,11 @@ import { Logger } from "../logging"; import type { CCloudEnvironment } from "../models/environment"; import type { CCloudFlinkComputePool } from "../models/flinkComputePool"; import type { FlinkSpecProperties, FlinkStatement } from "../models/flinkStatement"; -import { restFlinkStatementToModel, TERMINAL_PHASES } from "../models/flinkStatement"; +import { + FlinkSnapshotMode, + restFlinkStatementToModel, + TERMINAL_PHASES, +} from "../models/flinkStatement"; import type { CCloudFlinkDbKafkaCluster } from "../models/kafkaCluster"; import { getSidecar } from "../sidecar"; import { UriMetadataKeys } from "../storage/constants"; @@ -46,6 +50,12 @@ export interface IFlinkStatementSubmitParameters { /** Metadata hints for the statement execution */ properties: FlinkSpecProperties; + /** + * Batch ("snapshot") vs streaming execution mode to submit the statement with. + * Defaults to {@link FlinkSnapshotMode.STREAMING} when omitted. + */ + snapshotMode?: FlinkSnapshotMode; + /** * False if user directly gestured / wrote this statement, true if it was created by the extension * (system catalog queries to support our view providers, ...). @@ -61,6 +71,11 @@ export async function submitFlinkStatement( ): Promise { const handle = await getSidecar(); + const properties: Record = params.properties.toProperties(); + if (params.snapshotMode === FlinkSnapshotMode.BATCH) { + properties["sql.snapshot.mode"] = "now"; + } + const requestInner: CreateSqlv1StatementRequest = { api_version: CreateSqlv1StatementRequestApiVersionEnum.SqlV1, kind: CreateSqlv1StatementRequestKindEnum.Statement, @@ -70,7 +85,7 @@ export async function submitFlinkStatement( spec: { statement: params.statement, compute_pool_id: params.computePool.id, - properties: params.properties.toProperties(), + properties, }, }; @@ -429,11 +444,12 @@ export async function setFlinkDocumentMetadata( database?: CCloudFlinkDbKafkaCluster; computePool?: CCloudFlinkComputePool; fromWorkspace?: boolean; + snapshotMode?: FlinkSnapshotMode; }, ): Promise { const metadata: UriMetadata = {}; - const { catalog: environment, database, computePool, fromWorkspace } = opts; + const { catalog: environment, database, computePool, fromWorkspace, snapshotMode } = opts; if (environment) { metadata[UriMetadataKeys.FLINK_CATALOG_ID] = environment.id; @@ -453,6 +469,10 @@ export async function setFlinkDocumentMetadata( metadata[UriMetadataKeys.FLINK_FROM_WORKSPACE] = true; } + if (snapshotMode) { + metadata[UriMetadataKeys.FLINK_SNAPSHOT_MODE] = snapshotMode; + } + logger.debug(`setting Flink catalog / database / compute pool metadata for URI`, { uri: documentUri.toString(), metadata, diff --git a/src/models/flinkStatement.test.ts b/src/models/flinkStatement.test.ts index dbd882221a..db9b2b26e8 100644 --- a/src/models/flinkStatement.test.ts +++ b/src/models/flinkStatement.test.ts @@ -9,6 +9,7 @@ import type { SqlV1StatementStatus } from "../clients/flinkSql"; import { CCLOUD_BASE_PATH } from "../constants"; import { IconNames } from "../icons"; import { + FlinkSnapshotMode, FlinkStatement, FlinkStatementTreeItem, Phase, @@ -296,6 +297,24 @@ describe("FlinkStatement", () => { }); }); + describe("mode", () => { + it("returns STREAMING when sql.snapshot.mode property is absent", () => { + const statement = createFlinkStatement({}); + assert.strictEqual(statement.mode, FlinkSnapshotMode.STREAMING); + }); + + it("returns BATCH when sql.snapshot.mode property is 'now'", () => { + const statement = createFlinkStatement({ mode: FlinkSnapshotMode.BATCH }); + assert.strictEqual(statement.mode, FlinkSnapshotMode.BATCH); + }); + + it("returns STREAMING when sql.snapshot.mode property is some other unrecognized value", () => { + const statement = createFlinkStatement({}); + statement.spec.properties = { ...statement.spec.properties, "sql.snapshot.mode": "later" }; + assert.strictEqual(statement.mode, FlinkSnapshotMode.STREAMING); + }); + }); + describe("possiblyViewable", () => { const ONE_DAY_MS = 24 * 60 * 60 * 1000; const now = new Date(); diff --git a/src/models/flinkStatement.ts b/src/models/flinkStatement.ts index e72d18b071..2091131c55 100644 --- a/src/models/flinkStatement.ts +++ b/src/models/flinkStatement.ts @@ -254,6 +254,17 @@ export class FlinkStatement implements IResourceBase, IdItem, ISearchable, IEnvP return this.spec.properties?.["sql.current-database"]; } + /** + * Whether the statement was (or will be) evaluated in bounded "batch" (snapshot) mode or + * continuous "streaming" mode. + * @see https://docs.confluent.io/cloud/current/flink/concepts/snapshot-queries.html#snapshot-mode + */ + get mode(): FlinkSnapshotMode { + return this.spec.properties?.["sql.snapshot.mode"] === "now" + ? FlinkSnapshotMode.BATCH + : FlinkSnapshotMode.STREAMING; + } + /** Returns true if the statement is in a failed or failing phase. */ get failed(): boolean { return FAILED_PHASES.includes(this.phase); @@ -441,6 +452,17 @@ export function createFlinkStatementTooltip(resource: FlinkStatement) { return tooltip; } +/** + * Statement execution mode, derived from (or set via) the `sql.snapshot.mode` statement property. + * @see https://docs.confluent.io/cloud/current/flink/concepts/snapshot-queries.html#snapshot-mode + */ +export enum FlinkSnapshotMode { + /** Bounded, one-shot evaluation against a snapshot of the data as of submission time. */ + BATCH = "batch", + /** Continuous, unbounded evaluation against data as it arrives. This is the CCloud default. */ + STREAMING = "streaming", +} + export class FlinkSpecProperties { currentCatalog: string | undefined = undefined; currentDatabase: string | undefined = undefined; diff --git a/src/storage/constants.ts b/src/storage/constants.ts index 574953bd96..8bcff65e25 100644 --- a/src/storage/constants.ts +++ b/src/storage/constants.ts @@ -45,6 +45,8 @@ export enum UriMetadataKeys { FLINK_DATABASE_NAME = "flinkDatabaseName", /** True when the document was opened from a Flink workspace deep link. */ FLINK_FROM_WORKSPACE = "flinkFromWorkspace", + /** Batch ("snapshot") vs streaming execution mode to submit the statement with. */ + FLINK_SNAPSHOT_MODE = "flinkSnapshotMode", } export enum SecretStorageKeys { diff --git a/tests/unit/testResources/flinkStatement.ts b/tests/unit/testResources/flinkStatement.ts index ae763e8293..360966a58a 100644 --- a/tests/unit/testResources/flinkStatement.ts +++ b/tests/unit/testResources/flinkStatement.ts @@ -1,7 +1,7 @@ import type { ColumnDetails, SqlV1StatementStatus } from "../../../src/clients/flinkSql"; import { ConnectionType } from "../../../src/clients/sidecar"; import { CCLOUD_CONNECTION_ID } from "../../../src/constants"; -import { FlinkStatement, Phase } from "../../../src/models/flinkStatement"; +import { FlinkSnapshotMode, FlinkStatement, Phase } from "../../../src/models/flinkStatement"; import type { EnvironmentId, OrganizationId } from "../../../src/models/resource"; import { TEST_CCLOUD_ENVIRONMENT, TEST_CCLOUD_PROVIDER, TEST_CCLOUD_REGION } from "./environments"; import { TEST_CCLOUD_FLINK_COMPUTE_POOL_ID } from "./flinkComputePool"; @@ -27,6 +27,8 @@ export interface CreateFlinkStatementArgs { appendOnly?: boolean; upsertColumns?: number[]; schemaColumns?: ColumnDetails[]; + + mode?: FlinkSnapshotMode; } export function createFlinkStatement(overrides: CreateFlinkStatementArgs = {}): FlinkStatement { @@ -65,6 +67,7 @@ export function createFlinkStatement(overrides: CreateFlinkStatementArgs = {}): "sql.current-catalog": TEST_CCLOUD_ENVIRONMENT.name, "sql.current-database": TEST_CCLOUD_KAFKA_CLUSTER.name, "sql.local-time-zone": "GMT-04:00", + ...(overrides.mode === FlinkSnapshotMode.BATCH ? { "sql.snapshot.mode": "now" } : {}), }, stopped: false, }, From 22aecd7578b68583aa2d12e795937f5bc0bb54aa Mon Sep 17 00:00:00 2001 From: James Robinson Date: Tue, 4 Aug 2026 14:29:06 -0400 Subject: [PATCH 2/4] Avoid any. --- src/flinkSql/statementUtils.test.ts | 16 +++++++++------- 1 file changed, 9 insertions(+), 7 deletions(-) diff --git a/src/flinkSql/statementUtils.test.ts b/src/flinkSql/statementUtils.test.ts index 7f08f9411b..d264f8d3e6 100644 --- a/src/flinkSql/statementUtils.test.ts +++ b/src/flinkSql/statementUtils.test.ts @@ -16,6 +16,7 @@ import { import { TEST_CCLOUD_ORGANIZATION_ID } from "../../tests/unit/testResources/organization"; import { createResponseError } from "../../tests/unit/testUtils"; import type { + CreateSqlv1Statement201Response, CreateSqlv1StatementOperationRequest, GetSqlv1Statement200Response, GetSqlv1StatementResult200Response, @@ -275,19 +276,20 @@ describe("flinkSql/statementUtils.ts", function () { snapshotMode, }; - const createSqlv1StatementStub = sandbox.stub().resolves(TEST_CCLOUD_FLINK_STATEMENT); + const stubbedStatementsApi = sandbox.createStubInstance(StatementsSqlV1Api); + stubbedStatementsApi.createSqlv1Statement.resolves( + TEST_CCLOUD_FLINK_STATEMENT as unknown as CreateSqlv1Statement201Response, + ); sandbox .stub(flinkStatementModels, "restFlinkStatementToModel") .returns(TEST_CCLOUD_FLINK_STATEMENT); - mockSidecar.getFlinkSqlStatementsApi.returns({ - createSqlv1Statement: createSqlv1StatementStub, - } as any); + mockSidecar.getFlinkSqlStatementsApi.returns(stubbedStatementsApi); await submitFlinkStatement(params); - sinon.assert.calledOnce(createSqlv1StatementStub); - const request = createSqlv1StatementStub.firstCall - .args[0] as CreateSqlv1StatementOperationRequest; + sinon.assert.calledOnce(stubbedStatementsApi.createSqlv1Statement); + const request: CreateSqlv1StatementOperationRequest = + stubbedStatementsApi.createSqlv1Statement.firstCall.args[0]; const spec = request.CreateSqlv1StatementRequest?.spec as { properties?: Record; }; From 7487fe46b8e29516163b6517cc248c64f47f0422 Mon Sep 17 00:00:00 2001 From: James Robinson Date: Tue, 4 Aug 2026 14:35:18 -0400 Subject: [PATCH 3/4] Changeloggery. --- CHANGELOG.md | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 873016a541..68591a654b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,12 @@ All notable changes to this extension will be documented in this file. ## Unreleased +### Added + +- Support for submitting Flink statements in + "[snapshot](https://docs.confluent.io/cloud/current/flink/concepts/snapshot-queries.html)" mode. + New codelens control for toggling streaming vs snapshot submission. + ## 2.3.1 ### Fixed From bf0a2b53aafb5cefe3d6cb1f0f9cbf2d4c53b227 Mon Sep 17 00:00:00 2001 From: James Robinson Date: Tue, 4 Aug 2026 15:19:01 -0400 Subject: [PATCH 4/4] Include new command in manifest. --- package.json | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/package.json b/package.json index 05352cf8e8..4aab552c6e 100644 --- a/package.json +++ b/package.json @@ -420,6 +420,12 @@ "title": "Set CCloud Flink Database for Flink Statement", "category": "Confluent: Flink SQL" }, + { + "command": "confluent.document.flinksql.toggleSnapshotMode", + "icon": "$(confluent-logo)", + "title": "Toggle Snapshot Mode for Flink Statement", + "category": "Confluent: Flink SQL" + }, { "command": "confluent.flink.configureFlinkDefaults", "icon": "$(settings-gear)", @@ -1447,6 +1453,10 @@ "command": "confluent.document.flinksql.setCCloudDatabase", "when": "false" }, + { + "command": "confluent.document.flinksql.toggleSnapshotMode", + "when": "false" + }, { "command": "confluent.flink.configureFlinkDefaults", "when": "true"