This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new c4f7ca958cc branch-4.1: [test](fe) Cover internal query audit failure
in FE unit test #65696 (#65763)
c4f7ca958cc is described below
commit c4f7ca958cc4762f533d5e1882811516ab812437
Author: yujun <[email protected]>
AuthorDate: Mon Jul 20 12:23:23 2026 +0800
branch-4.1: [test](fe) Cover internal query audit failure in FE unit test
#65696 (#65763)
cherry-pick: #65696
---
.../doris/qe/StmtExecutorInternalQueryTest.java | 45 ++++++++
.../test_audit_log_internal_query_failure.groovy | 119 ---------------------
2 files changed, 45 insertions(+), 119 deletions(-)
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorInternalQueryTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorInternalQueryTest.java
index 528d37895a8..139711bf8d6 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorInternalQueryTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/qe/StmtExecutorInternalQueryTest.java
@@ -18,13 +18,19 @@
package org.apache.doris.qe;
import org.apache.doris.analysis.StatementBase;
+import org.apache.doris.catalog.Env;
import org.apache.doris.common.ErrorCode;
+import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.nereids.NereidsPlanner;
+import org.apache.doris.plugin.AuditEvent;
+import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr;
import mockit.Mock;
import mockit.MockUp;
import org.junit.Assert;
import org.junit.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
public class StmtExecutorInternalQueryTest {
@Test
@@ -73,4 +79,43 @@ public class StmtExecutorInternalQueryTest {
ConnectContext.remove();
}
}
+
+ @Test
+ public void testExecuteInternalQuerySubmitsErrorAuditEventOnFailure() {
+ ConnectContext ctx = new ConnectContext();
+ ctx.setThreadLocalInfo();
+ StmtExecutor executor = new StmtExecutor(ctx, "select * from table1");
+ Env env = Env.getCurrentEnv();
+ WorkloadRuntimeStatusMgr originalWorkloadRuntimeStatusMgr =
env.getWorkloadRuntimeStatusMgr();
+ WorkloadRuntimeStatusMgr workloadRuntimeStatusMgr =
Mockito.mock(WorkloadRuntimeStatusMgr.class);
+ ArgumentCaptor<AuditEvent> auditEventCaptor =
ArgumentCaptor.forClass(AuditEvent.class);
+
+ Deencapsulation.setField(env, "workloadRuntimeStatusMgr",
workloadRuntimeStatusMgr);
+ new MockUp<NereidsPlanner>() {
+ @Mock
+ public void plan(StatementBase queryStmt,
org.apache.doris.thrift.TQueryOptions queryOptions) {
+ throw new RuntimeException("mock plan failure");
+ }
+ };
+ try {
+ Assert.assertThrows(RuntimeException.class,
executor::executeInternalQuery);
+
Mockito.verify(workloadRuntimeStatusMgr).submitFinishQueryToAudit(auditEventCaptor.capture());
+ } finally {
+ Deencapsulation.setField(env, "workloadRuntimeStatusMgr",
originalWorkloadRuntimeStatusMgr);
+ ConnectContext.remove();
+ }
+
+ AuditEvent auditEvent = auditEventCaptor.getValue();
+ Assert.assertEquals(AuditEvent.EventType.AFTER_QUERY, auditEvent.type);
+ Assert.assertEquals("ERR", auditEvent.state);
+ Assert.assertEquals(ErrorCode.ERR_INTERNAL_ERROR.getCode(),
auditEvent.errorCode);
+ Assert.assertNotNull(auditEvent.errorMessage);
+ Assert.assertTrue("error message should mention root cause, got: " +
auditEvent.errorMessage,
+ auditEvent.errorMessage.contains("mock plan failure"));
+ Assert.assertTrue("audit event should be marked as internal",
auditEvent.isInternal);
+ Assert.assertTrue("audit event should be marked as query",
auditEvent.isQuery);
+ Assert.assertTrue("audit event should be marked as nereids",
auditEvent.isNereids);
+ Assert.assertEquals("select * from table1", auditEvent.stmt);
+ Assert.assertEquals("a8ec30e5ad0820f8c5bd16a82a4491ca",
auditEvent.sqlHash);
+ }
}
diff --git
a/regression-test/suites/fault_injection_p0/test_audit_log_internal_query_failure.groovy
b/regression-test/suites/fault_injection_p0/test_audit_log_internal_query_failure.groovy
deleted file mode 100644
index f6d44b19227..00000000000
---
a/regression-test/suites/fault_injection_p0/test_audit_log_internal_query_failure.groovy
+++ /dev/null
@@ -1,119 +0,0 @@
-// Licensed to the Apache Software Foundation (ASF) under one
-// or more contributor license agreements. See the NOTICE file
-// distributed with this work for additional information
-// regarding copyright ownership. The ASF licenses this file
-// to you under the Apache License, Version 2.0 (the
-// "License"); you may not use this file except in compliance
-// with the License. You may obtain a copy of the License at
-//
-// http://www.apache.org/licenses/LICENSE-2.0
-//
-// Unless required by applicable law or agreed to in writing,
-// software distributed under the License is distributed on an
-// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-// KIND, either express or implied. See the License for the
-// specific language governing permissions and limitations
-// under the License.
-
-// Regression test for CIR-20019: when an internal query fails (for example
-// the column-statistics gathering SQL that ANALYZE issues against a user
-// table), the audit log entry must record state=ERR with a descriptive
-// error_message instead of the previous misleading state=OK / return_rows=0.
-suite('test_audit_log_internal_query_failure', 'nonConcurrent') {
- def tbl = 'test_audit_log_internal_query_failure_t1'
-
- setGlobalVarTemporary([enable_audit_plugin: true], {
- def sqlCacheOrigValue = null
- try {
- sqlCacheOrigValue = (sql "select @@enable_sql_cache")[0][0]
- // The audit_log table is loaded asynchronously. Repeated polling
- // can otherwise cache an early empty result before the target
audit
- // batch becomes visible.
- sql "set enable_sql_cache=false"
-
- sql "drop table if exists ${tbl}"
- sql """
- create table ${tbl} (k int, v int)
- duplicate key(k)
- distributed by hash(k) buckets 1
- properties('replication_num'='1')
- """
- sql "insert into ${tbl} values(1,10),(2,20),(3,30)"
-
- // Limit the injected IO error to this table's tablet so concurrent
- // reads on system tables (such as audit_log itself) are
unaffected.
- def tabletRows = sql_return_maparray "show tablets from ${tbl}"
- assertFalse(tabletRows.isEmpty(), "expected at least one tablet
for ${tbl}")
- def tabletId = tabletRows[0].TabletId
-
- // Capture the FE-side current timestamp to filter audit rows so
- // stale entries from previous runs do not satisfy the assertion.
- def startTime = (sql_return_maparray "select now() as
ts")[0].ts.toString()
-
- GetDebugPoint().clearDebugPointsForAllBEs()
- try {
- GetDebugPoint().enableDebugPointForAllBEs(
- "LocalFileReader::read_at_impl.io_error",
- [ sub_path: "/${tabletId}/" ])
-
- test {
- sql "analyze table ${tbl} with sync"
- exception "IO_ERROR"
- }
- } finally {
- GetDebugPoint().clearDebugPointsForAllBEs()
- }
-
- // Force a flush so the failed internal query is queryable from
- // __internal_schema.audit_log.
- // The failed gather SQL reads from our user table and runs as an
- // internal query; it must show up with state=ERR. Filter by start
- // time to avoid matching stale entries from previous runs.
- // Match the injected IO error instead of queried_tables_and_views:
- // failed internal queries do not always populate that field before
- // the error is audited, but the propagated error message contains
- // the tablet-scoped debug point evidence.
- def query = """select state, error_code, error_message
- from __internal_schema.audit_log
- where is_internal = 1
- and state = 'ERR'
- and error_message like '%IO_ERROR%'
- and error_message like '%${tabletId}%'
- and `time` >= '${startTime}'
- order by `time` desc limit 1"""
- def diagnosticQuery = """select `time`, state, error_code,
error_message, queried_tables_and_views
- from __internal_schema.audit_log
- where is_internal = 1
- and `time` >= '${startTime}'
- order by `time` desc limit 5"""
- def res = []
- int retry = 90
- while (res.isEmpty() && retry-- > 0) {
- sql "call flush_audit_log()"
- sleep(2000)
- res = sql_return_maparray "${query}"
- }
- if (res.isEmpty()) {
- logger.info("recent internal audit rows after injected IO
error: ${sql_return_maparray(diagnosticQuery)}")
- }
- assertFalse(res.isEmpty(),
- "expected an audit_log entry with state=ERR for the
injected IO_ERROR")
- assertEquals('ERR', res[0].state.toString())
- assertNotEquals('0', res[0].error_code.toString())
- assertNotNull(res[0].error_message)
- assertTrue(!res[0].error_message.toString().isEmpty(),
- "audit_log error_message should not be empty, got:
${res[0].error_message}")
- } finally {
- try {
- if (sqlCacheOrigValue != null) {
- sql "set enable_sql_cache=${sqlCacheOrigValue}"
- }
- } catch (Throwable ignored) {
- }
- try {
- sql "drop table if exists ${tbl}"
- } catch (Throwable ignored) {
- }
- }
- })
-}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]