From 9679d48824da372ba7e8a1e8b360f481aeaa4653 Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Sun, 19 Jul 2026 12:47:46 +0800 Subject: [PATCH] fix(sqlserver): return informational server messages --- Cargo.lock | 2 + .../main/java/com/dbx/agent/JdbcExecutor.java | 58 +++- .../java/com/dbx/agent/JdbcExecutorTest.java | 117 ++++++++ crates/dbx-core/Cargo.toml | 2 + crates/dbx-core/src/db/sqlserver.rs | 279 ++++++++++++++---- .../tests/sqlserver_server_messages.rs | 75 +++++ 6 files changed, 468 insertions(+), 65 deletions(-) create mode 100644 crates/dbx-core/tests/sqlserver_server_messages.rs diff --git a/Cargo.lock b/Cargo.lock index 1b5e9e980..9037bab87 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1918,6 +1918,8 @@ dependencies = [ "tokio-postgres", "tokio-postgres-rustls", "tokio-util", + "tracing", + "tracing-subscriber", "uuid", "webpki-roots 0.26.11", "zip 4.6.1", diff --git a/agents/common/src/main/java/com/dbx/agent/JdbcExecutor.java b/agents/common/src/main/java/com/dbx/agent/JdbcExecutor.java index e21ea8680..3f937ff39 100644 --- a/agents/common/src/main/java/com/dbx/agent/JdbcExecutor.java +++ b/agents/common/src/main/java/com/dbx/agent/JdbcExecutor.java @@ -6,13 +6,16 @@ import java.sql.Clob; import java.sql.ResultSet; import java.sql.ResultSetMetaData; import java.sql.SQLException; +import java.sql.SQLWarning; import java.sql.SQLXML; import java.sql.Statement; import java.sql.Types; import java.util.ArrayList; import java.util.Collections; +import java.util.IdentityHashMap; import java.util.List; import java.util.Locale; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.function.Function; @@ -104,20 +107,22 @@ public final class JdbcExecutor { // Do not translate them to Connection.commit(), which requires autoCommit=false. boolean hasResultSet = stmt.execute(trimmedSql); long elapsed = System.currentTimeMillis() - start; + QueryResult result; if (hasResultSet) { try (ResultSet rs = stmt.getResultSet()) { - return readResultSet(rs, elapsed, effectiveMaxRows, valueReader); + result = readResultSet(rs, elapsed, effectiveMaxRows, valueReader); } + } else { + int updateCount = stmt.getUpdateCount(); + result = new QueryResult( + Collections.emptyList(), + Collections.emptyList(), + updateCount >= 0 ? updateCount : 0, + elapsed, + false + ); } - - int updateCount = stmt.getUpdateCount(); - return new QueryResult( - Collections.emptyList(), - Collections.emptyList(), - updateCount >= 0 ? updateCount : 0, - elapsed, - false - ); + return withStatementWarnings(result, stmt); } finally { activeStatements.remove(stmt); } @@ -653,6 +658,39 @@ public final class JdbcExecutor { } } + private static QueryResult withStatementWarnings(QueryResult result, Statement stmt) { + if (!result.getColumns().isEmpty() || !result.getRows().isEmpty()) { + return result; + } + + List> rows = new ArrayList<>(); + try { + Set seen = Collections.newSetFromMap(new IdentityHashMap<>()); + for (SQLWarning warning = stmt.getWarnings(); warning != null && seen.add(warning); warning = warning.getNextWarning()) { + String message = warning.getMessage(); + if (message != null && !message.trim().isEmpty()) { + rows.add(Collections.singletonList(message)); + } + } + stmt.clearWarnings(); + } catch (SQLException ignored) { + // Warning retrieval is advisory; a driver bug here must not turn a + // successfully executed statement into a query failure. + } + + if (rows.isEmpty()) { + return result; + } + return new QueryResult( + Collections.singletonList("Message"), + Collections.singletonList("nvarchar"), + rows, + result.getAffected_rows(), + result.getExecution_time_ms(), + result.getTruncated() + ); + } + private QueryResult emptyQueryResult(long start) { return new QueryResult( Collections.emptyList(), diff --git a/agents/common/src/test/java/com/dbx/agent/JdbcExecutorTest.java b/agents/common/src/test/java/com/dbx/agent/JdbcExecutorTest.java index 3210f0b9d..45955dc00 100644 --- a/agents/common/src/test/java/com/dbx/agent/JdbcExecutorTest.java +++ b/agents/common/src/test/java/com/dbx/agent/JdbcExecutorTest.java @@ -9,15 +9,18 @@ import java.sql.Connection; import java.sql.ResultSet; import java.sql.ResultSetMetaData; import java.sql.SQLException; +import java.sql.SQLWarning; import java.sql.Statement; import java.sql.Types; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import javax.sql.rowset.serial.SerialBlob; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; class JdbcExecutorTest { @Test @@ -67,6 +70,72 @@ class JdbcExecutorTest { assertEquals(2, fixture.getColumnTypeCalls()); } + @Test + void executeReturnsMultipleStatementWarningsForNoResultStatements() { + SQLWarning first = new SQLWarning("identity value is 443", "S0003", 7998); + first.setNextWarning(new SQLWarning("DBCC execution completed", "S0001", 2528)); + AtomicInteger clearWarningsCalls = new AtomicInteger(); + + QueryResult result = JdbcExecutor.INSTANCE.execute( + executionConnection(false, -1, first, clearWarningsCalls, null, null), + "DBCC CHECKIDENT ('dbo.tVillage', RESEED)", + "", + schema -> "" + ); + + assertEquals(Arrays.asList("Message"), result.getColumns()); + assertEquals(Arrays.asList("nvarchar"), result.getColumn_types()); + assertEquals( + Arrays.asList(Arrays.asList("identity value is 443"), Arrays.asList("DBCC execution completed")), + result.getRows() + ); + assertEquals(1, clearWarningsCalls.get()); + } + + @Test + void executeKeepsEmptyNoResultStatementsUnchangedWithoutWarnings() { + QueryResult result = JdbcExecutor.INSTANCE.execute( + executionConnection(false, 3, null, new AtomicInteger(), null, null), + "UPDATE people SET active = 1", + "", + schema -> "" + ); + + assertEquals(Collections.emptyList(), result.getColumns()); + assertEquals(Collections.emptyList(), result.getRows()); + assertEquals(3L, result.getAffected_rows()); + } + + @Test + void executeDoesNotReplaceOrdinaryResultSetsWithWarnings() { + CountingResultSetFixture fixture = countingResultSet(new Object[][]{{1, "Ada"}}); + QueryResult result = JdbcExecutor.INSTANCE.execute( + executionConnection(true, -1, new SQLWarning("informational"), new AtomicInteger(), fixture.resultSet(), null), + "SELECT id, name FROM people", + "", + schema -> "" + ); + + assertEquals(Arrays.asList("id", "name"), result.getColumns()); + assertEquals(Arrays.asList(Arrays.asList(1, "Ada")), result.getRows()); + } + + @Test + void executeStillPropagatesStatementErrors() { + SQLException failure = new SQLException("permission denied", "42000", 229); + RuntimeException thrown = assertThrows( + RuntimeException.class, + () -> JdbcExecutor.INSTANCE.execute( + executionConnection(false, -1, null, new AtomicInteger(), null, failure), + "DBCC CHECKIDENT ('dbo.tVillage', RESEED)", + "", + schema -> "" + ) + ); + + assertEquals(failure, thrown.getCause()); + } + @Test void schemaSwitcherPrefersDriverSpecificSql() throws Exception { List calls = new ArrayList<>(); @@ -214,6 +283,54 @@ class JdbcExecutorTest { ); } + private static Connection executionConnection( + boolean hasResultSet, + int updateCount, + SQLWarning warning, + AtomicInteger clearWarningsCalls, + ResultSet resultSet, + SQLException executeFailure + ) { + InvocationHandler statementHandler = (Object unused, Method method, Object[] args) -> { + switch (method.getName()) { + case "execute": + if (executeFailure != null) { + throw executeFailure; + } + return hasResultSet; + case "getResultSet": + return resultSet; + case "getUpdateCount": + return updateCount; + case "getWarnings": + return warning; + case "clearWarnings": + clearWarningsCalls.incrementAndGet(); + return null; + case "close": + return null; + default: + return defaultValue(method.getReturnType()); + } + }; + Statement statement = (Statement) Proxy.newProxyInstance( + Statement.class.getClassLoader(), + new Class[]{Statement.class}, + statementHandler + ); + InvocationHandler connectionHandler = (Object unused, Method method, Object[] args) -> { + if (method.getName().equals("createStatement")) { + return statement; + } + return defaultValue(method.getReturnType()); + }; + return (Connection) Proxy.newProxyInstance( + Connection.class.getClassLoader(), + new Class[]{Connection.class}, + connectionHandler + ); + } + private static Object defaultValue(Class type) { if (type == Boolean.TYPE) { return false; diff --git a/crates/dbx-core/Cargo.toml b/crates/dbx-core/Cargo.toml index e0226301e..f96fc29fa 100644 --- a/crates/dbx-core/Cargo.toml +++ b/crates/dbx-core/Cargo.toml @@ -46,6 +46,8 @@ rayon = "1" percent-encoding = "2" encoding_rs = "0.8" log = "0.4" +tracing = "0.1" +tracing-subscriber = "0.3" tokio = { version = "1", features = ["full"] } tokio-util = { version = "0.7", features = ["compat"] } chrono = { version = "0.4", features = ["serde"] } diff --git a/crates/dbx-core/src/db/sqlserver.rs b/crates/dbx-core/src/db/sqlserver.rs index fb7147c77..1d9e64fa7 100644 --- a/crates/dbx-core/src/db/sqlserver.rs +++ b/crates/dbx-core/src/db/sqlserver.rs @@ -7,11 +7,16 @@ use crate::types::{ use futures::{FutureExt, TryStreamExt}; use std::future::Future; use std::panic::AssertUnwindSafe; +use std::sync::{Arc as StdArc, Mutex as StdMutex}; use std::time::{Duration, Instant}; use tiberius::{AuthMethod, Client, ColumnData, Config, FromSql, QueryItem, QueryStream, SqlBrowser}; use tokio::net::TcpStream; use tokio_util::compat::{Compat, TokioAsyncWriteCompatExt}; use tokio_util::sync::CancellationToken; +use tracing::instrument::WithSubscriber; +use tracing::{Event, Level, Subscriber}; +use tracing_subscriber::layer::{Context as LayerContext, SubscriberExt}; +use tracing_subscriber::Layer; pub type SqlServerClient = Client>; pub const SQLSERVER_DRIVER_PANIC_ERROR_PREFIX: &str = "SQL Server driver panic:"; @@ -200,6 +205,98 @@ fn column_types_from_metadata(metadata: &tiberius::ResultMetadata) -> Vec>>, +} + +impl Layer for SqlServerMessageLayer +where + S: Subscriber, +{ + fn on_event(&self, event: &Event<'_>, _ctx: LayerContext<'_, S>) { + let metadata = event.metadata(); + if metadata.level() != &Level::INFO + || metadata.target() != "tiberius::tds::stream::token" + || metadata.line() != Some(TIBERIUS_INFO_TOKEN_EVENT_LINE) + { + return; + } + + let mut visitor = SqlServerMessageVisitor::default(); + event.record(&mut visitor); + if let Some(message) = visitor.message.filter(|message| !message.trim().is_empty()) { + if let Ok(mut messages) = self.messages.lock() { + messages.push(message); + } + } + } +} + +#[derive(Default)] +struct SqlServerMessageVisitor { + message: Option, +} + +impl tracing::field::Visit for SqlServerMessageVisitor { + fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) { + if field.name() == "message" { + self.message = Some(format!("{value:?}").trim_matches('"').to_string()); + } + } +} + +async fn capture_sqlserver_messages(future: F) -> (T, Vec) +where + F: Future, +{ + let layer = SqlServerMessageLayer::default(); + let messages = layer.messages.clone(); + // Tiberius consumes TDS INFO tokens internally and exposes them only as + // tracing events. Scope collection to this future to isolate connections. + let output = future.with_subscriber(tracing_subscriber::registry().with(layer)).await; + let messages = messages.lock().map(|messages| messages.clone()).unwrap_or_default(); + (output, messages) +} + +fn query_result_with_server_messages(mut result: QueryResult, messages: Vec) -> QueryResult { + if messages.is_empty() || !result.columns.is_empty() || !result.rows.is_empty() { + return result; + } + + result.columns = vec![SQLSERVER_MESSAGE_COLUMN.to_string()]; + result.column_types = vec!["nvarchar".to_string()]; + result.rows = messages.into_iter().map(|message| vec![serde_json::Value::String(message)]).collect(); + result +} + +fn server_messages_query_result(messages: Vec, start: Instant) -> Option { + if messages.is_empty() { + return None; + } + + Some(query_result_with_server_messages( + QueryResult { + columns: vec![], + column_types: vec![], + column_sortables: vec![], + rows: vec![], + affected_rows: 0, + execution_time_ms: start.elapsed().as_millis(), + truncated: false, + session_id: None, + has_more: false, + }, + messages, + )) +} + async fn collect_first_result_limited( mut stream: QueryStream<'_>, start: Instant, @@ -1688,37 +1785,52 @@ pub async fn execute_query_with_max_rows( Ok(Some(sql)) => sql, Ok(None) | Err(_) => sql.to_string(), }; - let stream = sqlserver_driver_result(client.query(query_sql.as_str(), &[])).await?; - let mut result = sqlserver_driver_result(collect_first_result_limited(stream, start, max_rows)).await?; + let (result, messages) = capture_sqlserver_messages(async { + let stream = sqlserver_driver_result(client.query(query_sql.as_str(), &[])).await?; + sqlserver_driver_result(collect_first_result_limited(stream, start, max_rows)).await + }) + .await; + let mut result = query_result_with_server_messages(result?, messages); strip_dbx_sqlserver_row_number_column(&mut result, sql); Ok(result) } else if requires_simple_query_batch(sql) || contains_transaction_control(sql) { - let stream = sqlserver_driver_result(client.simple_query(sql)).await?; - let _ = sqlserver_driver_result(collect_result_sets_limited(stream, start, max_rows)).await?; - Ok(QueryResult { - columns: vec![], - column_types: Vec::new(), - column_sortables: vec![], - rows: vec![], - affected_rows: 0, - execution_time_ms: start.elapsed().as_millis(), - truncated: false, - session_id: None, - has_more: false, + let (result, messages) = capture_sqlserver_messages(async { + let stream = sqlserver_driver_result(client.simple_query(sql)).await?; + sqlserver_driver_result(collect_result_sets_limited(stream, start, max_rows)).await }) + .await; + let _ = result?; + Ok(query_result_with_server_messages( + QueryResult { + columns: vec![], + column_types: Vec::new(), + column_sortables: vec![], + rows: vec![], + affected_rows: 0, + execution_time_ms: start.elapsed().as_millis(), + truncated: false, + session_id: None, + has_more: false, + }, + messages, + )) } else { - let result = sqlserver_driver_result(client.execute(sql, &[])).await?; - Ok(QueryResult { - columns: vec![], - column_types: Vec::new(), - column_sortables: vec![], - rows: vec![], - affected_rows: result.rows_affected().iter().sum::(), - execution_time_ms: start.elapsed().as_millis(), - truncated: false, - session_id: None, - has_more: false, - }) + let (result, messages) = capture_sqlserver_messages(sqlserver_driver_result(client.execute(sql, &[]))).await; + let result = result?; + Ok(query_result_with_server_messages( + QueryResult { + columns: vec![], + column_types: Vec::new(), + column_sortables: vec![], + rows: vec![], + affected_rows: result.rows_affected().iter().sum::(), + execution_time_ms: start.elapsed().as_millis(), + truncated: false, + session_id: None, + has_more: false, + }, + messages, + )) } } @@ -1733,29 +1845,36 @@ pub async fn execute_batch_with_max_rows( ) -> Result, String> { let start = Instant::now(); if sqlserver_batch_can_use_execute(sql) { - let result = sqlserver_driver_result(client.execute(sql, &[])).await?; - return Ok(vec![QueryResult { - columns: vec![], - column_types: Vec::new(), - column_sortables: vec![], - rows: vec![], - affected_rows: result.rows_affected().iter().sum::(), - execution_time_ms: start.elapsed().as_millis(), - truncated: false, - session_id: None, - has_more: false, - }]); + let (result, messages) = capture_sqlserver_messages(sqlserver_driver_result(client.execute(sql, &[]))).await; + let result = result?; + return Ok(vec![query_result_with_server_messages( + QueryResult { + columns: vec![], + column_types: Vec::new(), + column_sortables: vec![], + rows: vec![], + affected_rows: result.rows_affected().iter().sum::(), + execution_time_ms: start.elapsed().as_millis(), + truncated: false, + session_id: None, + has_more: false, + }, + messages, + )]); } if is_single_sqlserver_select(sql) { if let Ok(Some(query_sql)) = sqlserver_unsafe_type_query(client, sql).await { - let stream = sqlserver_driver_result(client.query(query_sql.as_str(), &[])).await?; - return sqlserver_driver_result(collect_first_result_limited(stream, start, max_rows)).await.map( - |mut result| { - strip_dbx_sqlserver_row_number_column(&mut result, sql); - vec![result] - }, - ); + let (result, messages) = capture_sqlserver_messages(async { + let stream = sqlserver_driver_result(client.query(query_sql.as_str(), &[])).await?; + sqlserver_driver_result(collect_first_result_limited(stream, start, max_rows)).await + }) + .await; + return result.map(|result| { + let mut result = query_result_with_server_messages(result, messages); + strip_dbx_sqlserver_row_number_column(&mut result, sql); + vec![result] + }); } } execute_simple_batch_with_max_rows(client, sql, max_rows).await @@ -1772,13 +1891,19 @@ pub async fn execute_simple_batch_with_max_rows( max_rows: Option, ) -> Result, String> { let start = Instant::now(); - let stream = sqlserver_driver_result(client.simple_query(sql)).await?; - let mut results = sqlserver_driver_result(collect_result_sets_limited(stream, start, max_rows)).await?; + let (results, messages) = capture_sqlserver_messages(async { + let stream = sqlserver_driver_result(client.simple_query(sql)).await?; + sqlserver_driver_result(collect_result_sets_limited(stream, start, max_rows)).await + }) + .await; + let mut results = results?; for result in &mut results { strip_dbx_sqlserver_row_number_column(result, sql); } - if results.is_empty() { + if let Some(message_result) = server_messages_query_result(messages, start) { + results.push(message_result); + } else if results.is_empty() { results.push(QueryResult { columns: vec![], column_types: Vec::new(), @@ -1958,13 +2083,13 @@ fn first_sql_tokens(sql: &str, limit: usize) -> Vec { #[cfg(test)] mod tests { use super::{ - build_sqlserver_unsafe_type_query, format_sqlserver_numeric, is_sqlserver_spatial_column, - is_sqlserver_variant_column, requires_simple_query_batch, sqlserver_batch_can_use_execute, - sqlserver_cell_to_json, sqlserver_columns_sql, sqlserver_completion_assistant_sql, - sqlserver_dml_output_returns_rows, sqlserver_hidden_schema_names, sqlserver_indexes_sql, - sqlserver_list_objects_sql, sqlserver_list_schemas_sql, sqlserver_list_tables_sql, sqlserver_table_comment_sql, - sqlserver_visible_object_predicate, strip_dbx_sqlserver_row_number_column, SqlServerDescribedColumn, - SqlServerResultSet, + build_sqlserver_unsafe_type_query, capture_sqlserver_messages, format_sqlserver_numeric, + is_sqlserver_spatial_column, is_sqlserver_variant_column, query_result_with_server_messages, + requires_simple_query_batch, sqlserver_batch_can_use_execute, sqlserver_cell_to_json, sqlserver_columns_sql, + sqlserver_completion_assistant_sql, sqlserver_dml_output_returns_rows, sqlserver_hidden_schema_names, + sqlserver_indexes_sql, sqlserver_list_objects_sql, sqlserver_list_schemas_sql, sqlserver_list_tables_sql, + sqlserver_table_comment_sql, sqlserver_visible_object_predicate, strip_dbx_sqlserver_row_number_column, + SqlServerDescribedColumn, SqlServerResultSet, }; use crate::types::{ CompletionAssistantMatchMode, CompletionAssistantObjectKind, CompletionAssistantRequest, QueryResult, @@ -1973,6 +2098,50 @@ mod tests { use std::{borrow::Cow, time::Instant}; use tiberius::{ColumnData, IntoSql}; + #[tokio::test] + async fn sqlserver_ignores_non_info_tiberius_events() { + let (_, messages) = capture_sqlserver_messages(async { + tracing::event!(target: "tiberius::tds::stream::token", tracing::Level::ERROR, "permission denied"); + tracing::event!(target: "dbx_core::db::sqlserver", tracing::Level::INFO, "not a TDS token"); + }) + .await; + + assert!(messages.is_empty()); + } + + #[test] + fn sqlserver_server_messages_fill_only_empty_results() { + let empty = QueryResult { + columns: vec![], + column_types: vec![], + column_sortables: vec![], + rows: vec![], + affected_rows: 0, + execution_time_ms: 1, + truncated: false, + session_id: None, + has_more: false, + }; + let result = query_result_with_server_messages(empty, vec!["DBCC execution completed".to_string()]); + assert_eq!(result.columns, vec!["Message"]); + assert_eq!(result.rows, vec![vec![serde_json::json!("DBCC execution completed")]]); + + let select = QueryResult { + columns: vec!["id".to_string()], + column_types: vec!["int".to_string()], + column_sortables: vec![], + rows: vec![vec![serde_json::json!(1)]], + affected_rows: 0, + execution_time_ms: 1, + truncated: false, + session_id: None, + has_more: false, + }; + let result = query_result_with_server_messages(select, vec!["informational".to_string()]); + assert_eq!(result.columns, vec!["id"]); + assert_eq!(result.rows, vec![vec![serde_json::json!(1)]]); + } + #[test] fn sqlserver_endpoint_splits_named_instance_hosts() { assert_eq!( diff --git a/crates/dbx-core/tests/sqlserver_server_messages.rs b/crates/dbx-core/tests/sqlserver_server_messages.rs new file mode 100644 index 000000000..c851f1f5b --- /dev/null +++ b/crates/dbx-core/tests/sqlserver_server_messages.rs @@ -0,0 +1,75 @@ +use dbx_core::db::sqlserver; +use std::time::Duration; + +#[tokio::test] +#[ignore = "requires DBX_TEST_SQLSERVER_HOST and DBX_TEST_SQLSERVER_PASSWORD"] +async fn sqlserver_dbcc_messages_do_not_replace_results_or_errors() { + let host = std::env::var("DBX_TEST_SQLSERVER_HOST").expect("DBX_TEST_SQLSERVER_HOST"); + let port = std::env::var("DBX_TEST_SQLSERVER_PORT").ok().and_then(|value| value.parse().ok()).unwrap_or(1433); + let user = std::env::var("DBX_TEST_SQLSERVER_USER").unwrap_or_else(|_| "sa".to_string()); + let password = std::env::var("DBX_TEST_SQLSERVER_PASSWORD").expect("DBX_TEST_SQLSERVER_PASSWORD"); + let mut client = sqlserver::connect_with_port_explicit( + &host, + port, + true, + &user, + &password, + Some("master"), + Duration::from_secs(15), + ) + .await + .expect("connect to SQL Server"); + + let table = "dbo.dbx_issue_3583_messages"; + let setup = format!( + "IF OBJECT_ID('{table}', 'U') IS NOT NULL DROP TABLE {table}; \ + CREATE TABLE {table} (id BIGINT IDENTITY(1,1) PRIMARY KEY, name NVARCHAR(20)); \ + INSERT INTO {table} (name) VALUES (N'test');" + ); + sqlserver::execute_batch(&mut client, &setup).await.expect("create DBCC fixture"); + + let dbcc = sqlserver::execute_query(&mut client, &format!("DBCC CHECKIDENT ('{table}', RESEED)",)) + .await + .expect("execute DBCC CHECKIDENT"); + assert_eq!(dbcc.columns, vec!["Message"]); + assert!(dbcc.rows.len() >= 2, "expected SQL Server identity and completion messages: {dbcc:?}"); + assert!(dbcc.rows.iter().any(|row| row[0].as_str().is_some_and(|message| message.contains("identity")))); + + let select = sqlserver::execute_query(&mut client, &format!("SELECT id, name FROM {table}")) + .await + .expect("ordinary SELECT remains available"); + assert_eq!(select.columns, vec!["id", "name"]); + assert_eq!(select.rows.len(), 1); + + let multi = sqlserver::execute_batch(&mut client, "SELECT 1 AS first; SELECT 2 AS second") + .await + .expect("multiple result sets remain available"); + assert_eq!(multi.len(), 2); + assert_eq!(multi[0].columns, vec!["first"]); + assert_eq!(multi[1].columns, vec!["second"]); + + let dml = sqlserver::execute_query(&mut client, &format!("UPDATE {table} SET name = N'updated'")) + .await + .expect("ordinary DML remains available"); + assert_eq!(dml.affected_rows, 1); + assert!(dml.columns.is_empty()); + + let use_database = + sqlserver::execute_query(&mut client, "USE master").await.expect("execute database context change"); + assert_eq!(use_database.columns, vec!["Message"]); + assert!( + use_database.rows.iter().all(|row| !row[0] + .as_str() + .is_some_and(|message| message.starts_with("Database change") + || message.starts_with("SQL collation") + || message.starts_with("Packet size change"))), + "internal TDS environment changes must not leak into server messages: {use_database:?}" + ); + + let error = sqlserver::execute_query(&mut client, "SELECT * FROM dbo.dbx_issue_3583_missing") + .await + .expect_err("real SQL Server errors must remain failures"); + assert!(!error.trim().is_empty()); + + sqlserver::execute_query(&mut client, &format!("DROP TABLE {table}")).await.expect("clean up DBCC fixture"); +}