Repository navigation
fix(warehouse): write Snowflake reconnects to the log #1427
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,49 @@ | ||
| /** | ||
| * A driver replacing a connection the warehouse closed (today: Snowflake after an idle | ||
| * timeout, a network drop or laptop sleep). The reopen is silent to the caller, which is | ||
| * the point, but without these events it was also missing from the log: a session that | ||
| * kept dropping looked exactly like one that never did. | ||
| */ | ||
| export interface ReconnectEvent { | ||
| warehouse: string | ||
| /** Identifies which account the connection belongs to, e.g. the Snowflake account locator. */ | ||
| account?: string | ||
| phase: "started" | "reconnected" | "failed" | ||
| /** How the closed session showed: the SDK reported it down before a statement, or a statement failed because | ||
| * the session had closed. The second does not say the statement did not run: a write may already have. */ | ||
| reason: "connection-down" | "closed-during-statement" | ||
| /** From the start of the reconnect. On `reconnected` and `failed`. */ | ||
| durationMs?: number | ||
| /** Session settings (USE, SET, ALTER SESSION…) restored on the new session. On `reconnected`. */ | ||
| settingsRestored?: number | ||
| /** Temporary objects or an open transaction went with the old session, or a statement that could create them | ||
| * was still running on it, so whether they did is not known yet. On `reconnected`. */ | ||
| sessionStateLost?: boolean | ||
| /** On `failed`. Never carries the session's statements (a replayed setting can hold anything the user ran). */ | ||
| error?: string | ||
| } | ||
|
|
||
| type Listener = (event: ReconnectEvent) => void | Promise<void> | ||
|
|
||
| // Process-global, for the same reason as the sign-in notices: the driver and its | ||
| // subscriber can be loaded through different module graphs. | ||
| const LISTENERS_KEY = Symbol.for("altimate.drivers.reconnectListeners") | ||
| const listeners: Set<Listener> = ((globalThis as Record<symbol, unknown>)[LISTENERS_KEY] as Set<Listener>) ?? | ||
| ((globalThis as Record<symbol, unknown>)[LISTENERS_KEY] = new Set<Listener>()) | ||
|
|
||
| /** Subscribe to reconnect events from every driver. Returns the unsubscribe. */ | ||
| export function onReconnect(listener: Listener): () => void { | ||
| listeners.add(listener) | ||
| return () => listeners.delete(listener) | ||
| } | ||
|
|
||
| export function emitReconnect(event: ReconnectEvent): void { | ||
| for (const listener of listeners) { | ||
| // A broken listener, sync or async, must not fail the reconnect or leave an unhandled rejection. | ||
| try { | ||
| void Promise.resolve(listener(event)).catch(() => {}) | ||
| } catch { | ||
| // as above | ||
| } | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -6,6 +6,7 @@ import * as fs from "fs" | |
| import type { ConnectionConfig, Connector, ConnectorResult, ExecuteOptions, SchemaColumn } from "./types" | ||
| import { loadOptionalDriver } from "./resolve" | ||
| import { emitBrowserSignIn, openInBrowser } from "./sign-in" | ||
| import { emitReconnect, type ReconnectEvent } from "./reconnect-events" | ||
|
|
||
| /** | ||
| * Run `fn` with stdout/stderr writes swallowed for the (synchronous) duration of | ||
|
|
@@ -214,6 +215,9 @@ export async function connect( | |
| /** The generation at which that state was lost: statements that started before it are not run, because they | ||
| * may depend on it; statements issued after it run on the new session. Never cleared by a late error. */ | ||
| let stateLostAt: number | undefined | ||
| /** Statements running now that may create temporary objects or open a transaction: one that finishes after a | ||
| * reconnect is reported to its caller then, so the reconnect cannot yet say the session lost nothing. */ | ||
| let stateStatementsRunning = 0 | ||
|
|
||
| function openConnection(): Promise<any> { | ||
| const account = connectOptions?.account | ||
|
|
@@ -290,9 +294,13 @@ export async function connect( | |
| } | ||
|
|
||
| /** Replace a connection Snowflake has closed. Concurrent callers share one reconnect. */ | ||
| function reconnect(): Promise<void> { | ||
| function reconnect(reason: ReconnectEvent["reason"]): Promise<void> { | ||
| if (!reconnecting) { | ||
| const previous = connection | ||
| const account = String(connectOptions?.account ?? "") | ||
| const startedAt = Date.now() | ||
| let settingsRestored = 0 | ||
| emitReconnect({ warehouse: "snowflake", account, phase: "started", reason }) | ||
| reconnecting = openConnection() | ||
| .then(async (conn) => { | ||
| const discard = () => { | ||
|
|
@@ -314,13 +322,22 @@ export async function connect( | |
| await runQuery(conn, setting) | ||
| } catch (err) { | ||
| discard() | ||
| throw new Error( | ||
| `Snowflake closed the session and its settings could not be restored on a new one (${setting}): ${(err as Error)?.message ?? err}`, | ||
| // The caller ran this setting and is told which; the reconnect event, which any subscriber in the | ||
| // process receives, gets the message without it. | ||
| const cause = (err as Error)?.message ?? err | ||
| throw Object.assign( | ||
| new Error( | ||
| `Snowflake closed the session and its settings could not be restored on a new one (${setting}): ${cause}`, | ||
| ), | ||
| { eventMessage: `Snowflake closed the session and its settings could not be restored on a new one: ${cause}` }, | ||
| ) | ||
| } | ||
| } | ||
| const unchanged = snapshot.length === sessionSettings.length && snapshot.every((v, i) => v === sessionSettings[i]) | ||
| if (unchanged) break | ||
| if (unchanged) { | ||
| settingsRestored = snapshot.length | ||
| break | ||
| } | ||
| if (pass + 1 >= MAX_REPLAY_PASSES) { | ||
| // Settings keep arriving: a connection installed now would run later queries with some of them missing. | ||
| discard() | ||
|
|
@@ -346,6 +363,26 @@ export async function connect( | |
| } catch { | ||
| // already terminated — nothing to release | ||
| } | ||
| emitReconnect({ | ||
| warehouse: "snowflake", | ||
| account, | ||
| phase: "reconnected", | ||
| reason, | ||
| durationMs: Date.now() - startedAt, | ||
| settingsRestored, | ||
| sessionStateLost: stateLostAt === generation || stateStatementsRunning > 0, | ||
| }) | ||
| }) | ||
| .catch((err) => { | ||
| emitReconnect({ | ||
| warehouse: "snowflake", | ||
| account, | ||
| phase: "failed", | ||
| reason, | ||
| durationMs: Date.now() - startedAt, | ||
| error: String((err as { eventMessage?: string })?.eventMessage ?? (err as Error)?.message ?? err), | ||
| }) | ||
| throw err | ||
| }) | ||
| .finally(() => { | ||
| reconnecting = undefined | ||
|
|
@@ -372,7 +409,7 @@ export async function connect( | |
| async function ensureLive(): Promise<void> { | ||
| if (reconnecting) return reconnecting | ||
| if (connection && connectOptions && typeof connection.isUp === "function" && !connection.isUp()) { | ||
| await reconnect() | ||
| await reconnect("connection-down") | ||
| } | ||
| } | ||
|
|
||
|
|
@@ -421,6 +458,17 @@ export async function connect( | |
| * it may already have run (see `isRetrySafe`). | ||
| */ | ||
| async function executeQuery(sql: string, binds?: any[]): Promise<{ columns: string[]; rows: any[][] }> { | ||
| const mayHoldState = opensTransaction(sql) || holdsSessionState(sql) || (looksLikeSessionChange(sql) && !isSessionSetting(sql)) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. WARNING: Bind-bearing session settings are omitted from in-flight state tracking
Reply with |
||
| if (!mayHoldState) return executeTracked(sql, binds) | ||
| stateStatementsRunning++ | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win Count only statements running on the old session. If a caller issues 🤖 Prompt for AI AgentsThere was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. P2: This increments before Prompt for AI agents |
||
| try { | ||
| return await executeTracked(sql, binds) | ||
| } finally { | ||
| stateStatementsRunning-- | ||
| } | ||
| } | ||
|
|
||
| async function executeTracked(sql: string, binds?: any[]): Promise<{ columns: string[]; rows: any[][] }> { | ||
| // Captured before any wait: a statement queued behind a reconnect started on the old session's assumptions. | ||
| const started = generation | ||
| const startedInEpoch = epoch | ||
|
|
@@ -436,7 +484,7 @@ export async function connect( | |
| if (!connectOptions || !isClosedConnectionError(err)) throw err | ||
| // Reopen only if the failed connection is still the current one: a late error from a connection | ||
| // another statement already replaced must not tear down its replacement. | ||
| if (connection === used) await reconnect() | ||
| if (connection === used) await reconnect("closed-during-statement") | ||
| else if (reconnecting) await reconnecting | ||
| const cause = String((err as Error)?.message ?? err) | ||
| if (lostFor(started)) throw sessionLostError(cause) | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,76 @@ | ||
| // altimate_change - new file | ||
| /** | ||
| * Writes a driver's reconnects to `opencode.log`. When Snowflake closes a session (idle | ||
| * timeout, VPN drop, laptop sleep) the driver opens a new one and the statement goes on, | ||
| * so the user sees nothing; without these lines the log did not show it either, and a | ||
| * session that kept dropping looked the same as one that never did. | ||
| */ | ||
| import { onReconnect, type ReconnectEvent } from "@altimateai/drivers" | ||
| import { fileLog } from "@/altimate/util/file-log" | ||
| import { Telemetry } from "@/altimate/telemetry" | ||
|
|
||
| type Write = typeof fileLog | ||
|
|
||
| /** Connection names by warehouse account: the driver knows the account, the log lines name the connection. */ | ||
| const names = new Map<string, Set<string>>() | ||
|
|
||
| /** Called before each connect, so a later reconnect on that account can name its connection. */ | ||
| export function remember(account: string, name: string): void { | ||
| if (!account) return | ||
| const set = names.get(account) ?? new Set<string>() | ||
|
cubic-dev-ai[bot] marked this conversation as resolved.
|
||
| set.add(name) | ||
| names.set(account, set) | ||
| } | ||
|
|
||
| /** Called when a connection is removed, so a later reconnect on its account is not put down to it. */ | ||
| export function forget(name: string): void { | ||
| for (const [account, set] of names) { | ||
| set.delete(name) | ||
| if (set.size === 0) names.delete(account) | ||
| } | ||
| } | ||
|
|
||
| /** One name only when the account maps to exactly one connection; two connections on one account are both possible. */ | ||
| function nameFor(account: string | undefined): string | undefined { | ||
| const set = account ? names.get(account) : undefined | ||
| return set && set.size === 1 ? [...set][0] : undefined | ||
| } | ||
|
|
||
| export function handle(e: ReconnectEvent, write: Write = fileLog): void { | ||
| const base = { | ||
| ...(nameFor(e.account) ? { name: nameFor(e.account) } : {}), | ||
| type: e.warehouse, | ||
| ...(e.account ? { account: e.account } : {}), | ||
| reason: e.reason, | ||
| } | ||
| if (e.phase === "started") { | ||
| write("INFO", "warehouse-connect", "reconnecting", base) | ||
| } else if (e.phase === "reconnected") { | ||
| write("INFO", "warehouse-connect", "reconnected", { | ||
| ...base, | ||
| duration_ms: e.durationMs, | ||
| settings_restored: e.settingsRestored, | ||
| session_state_lost: e.sessionStateLost, | ||
| }) | ||
| } else { | ||
| write("WARN", "warehouse-connect", "reconnect failed", { | ||
| ...base, | ||
| duration_ms: e.durationMs, | ||
| error: Telemetry.maskString(String(e.error ?? "")).slice(0, 500), | ||
| }) | ||
| } | ||
| } | ||
|
|
||
| let unsubscribe: (() => void) | undefined | ||
|
|
||
| /** Idempotent; called before every warehouse connect. */ | ||
| export function install(write: Write = fileLog): void { | ||
| if (unsubscribe) return | ||
| unsubscribe = onReconnect((e) => handle(e, write)) | ||
| } | ||
|
|
||
| export function resetForTests(): void { | ||
| unsubscribe?.() | ||
| unsubscribe = undefined | ||
| names.clear() | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
WARNING: SDK replay errors can still disclose SQL to every reconnect listener
The new
eventMessageomitssettingbut appends the SDK'scauseunchanged. The checked-in failure fixture already shows a cause containing the schema named byUSE SCHEMA ANALYTICS; if the SDK echoes a setting or one of its literals, that text is delivered raw to every process-globalonReconnectlistener via thefailedevent. Masking inreconnect-log.tshappens only after delivery. Redact SQL-derived text in the event error, keeping the detailed SDK error on the caller's exception.Reply with
@kilocode-bot fix itto have Kilo Code address this issue.