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 95265876df7 branch-4.1: [improvement](fe) Improve audit logs for S3 
streaming insert jobs #67489 (#67756)
95265876df7 is described below

commit 95265876df795b248e893a3072eb0a8c9b1f55b3
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Thu Sep 10 21:28:16 2026 +0800

    branch-4.1: [improvement](fe) Improve audit logs for S3 streaming insert 
jobs #67489 (#67756)
    
    Cherry-picked from #67489
    
    ---------
    
    Co-authored-by: wudi <[email protected]>
---
 .../insert/streaming/StreamingInsertTask.java      |  51 +++-
 .../job/offset/s3/S3SourceOffsetProvider.java      |   2 +-
 .../streaming/StreamingInsertTaskAuditTest.java    | 270 +++++++++++++++++++++
 3 files changed, 321 insertions(+), 2 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
index df23f107240..da74592287c 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTask.java
@@ -20,7 +20,9 @@ package org.apache.doris.job.extensions.insert.streaming;
 import org.apache.doris.analysis.UserIdentity;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.common.Config;
+import org.apache.doris.common.ErrorCode;
 import org.apache.doris.common.FeConstants;
+import org.apache.doris.common.Pair;
 import org.apache.doris.common.Status;
 import org.apache.doris.common.util.Util;
 import org.apache.doris.job.base.Job;
@@ -30,16 +32,22 @@ import org.apache.doris.job.extensions.insert.InsertTask;
 import org.apache.doris.job.offset.SourceOffsetProvider;
 import org.apache.doris.load.loadv2.LoadJob;
 import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.analyzer.UnboundTVFRelation;
 import org.apache.doris.nereids.glue.LogicalPlanAdapter;
 import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.trees.plans.commands.info.BaseViewInfo;
 import 
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand;
+import org.apache.doris.nereids.util.SqlLiteralUtils;
+import org.apache.doris.qe.AuditLogHelper;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.QueryState;
 import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.tablefunction.S3TableValuedFunction;
 import org.apache.doris.thrift.TCell;
 import org.apache.doris.thrift.TRow;
 import org.apache.doris.thrift.TStatusCode;
 
+import com.google.common.base.Preconditions;
 import lombok.Getter;
 import lombok.extern.log4j.Log4j2;
 import org.apache.commons.lang3.StringUtils;
@@ -49,6 +57,8 @@ import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
+import java.util.TreeMap;
+import java.util.stream.Collectors;
 
 @Log4j2
 @Getter
@@ -61,6 +71,8 @@ public class StreamingInsertTask extends 
AbstractStreamingTask {
     private StreamingJobProperties jobProperties;
     private Map<String, String> originTvfProps;
     private String cloudCluster;
+    private String auditSql;
+    private final boolean auditEnabled;
     SourceOffsetProvider offsetProvider;
 
     public StreamingInsertTask(long jobId,
@@ -79,10 +91,12 @@ public class StreamingInsertTask extends 
AbstractStreamingTask {
         this.jobProperties = jobProperties;
         this.originTvfProps = originTvfProps;
         this.cloudCluster = cloudCluster;
+        this.auditEnabled = 
S3TableValuedFunction.NAME.equalsIgnoreCase(offsetProvider.getSourceType());
     }
 
     @Override
     public void before() throws Exception {
+        auditSql = null;
         if (getIsCanceled().get()) {
             log.info("streaming insert task has been canceled, task id is {}", 
getTaskId());
             return;
@@ -111,6 +125,14 @@ public class StreamingInsertTask extends 
AbstractStreamingTask {
         this.taskCommand = offsetProvider.rewriteTvfParams(baseCommand, 
runningOffset, getTaskId());
         this.taskCommand.setLabelName(Optional.of(labelName));
         this.stmtExecutor = new StmtExecutor(ctx, new 
LogicalPlanAdapter(taskCommand, ctx.getStatementContext()));
+        if (auditEnabled) {
+            ctx.setExecutor(stmtExecutor);
+            try {
+                this.auditSql = buildAuditSql();
+            } catch (Exception e) {
+                log.warn("Failed to prepare audit SQL, label {}; skipping 
audit for this attempt", labelName, e);
+            }
+        }
     }
 
     @Override
@@ -121,6 +143,9 @@ public class StreamingInsertTask extends 
AbstractStreamingTask {
         }
         log.info("start to run streaming insert task, label {}, offset is {}", 
labelName, runningOffset.toString());
         String errMsg = null;
+        if (auditEnabled) {
+            ctx.setStartTime();
+        }
         try {
             taskCommand.run(ctx, stmtExecutor);
             if (ctx.getState().getStateType() == QueryState.MysqlStateType.OK) 
{
@@ -130,12 +155,36 @@ public class StreamingInsertTask extends 
AbstractStreamingTask {
             }
             throw new JobException(errMsg);
         } catch (Exception e) {
+            String errorMessage = Util.getRootCauseMessage(e);
+            if (auditEnabled && ctx.getState().getStateType() != 
QueryState.MysqlStateType.ERR) {
+                ctx.getState().setError(ErrorCode.ERR_INTERNAL_ERROR, 
errorMessage);
+            }
             log.warn("execute insert task error, label is {},offset is {}", 
taskCommand.getLabelName(),
                     runningOffset.toString(), e);
-            throw new JobException(Util.getRootCauseMessage(e));
+            throw new JobException(errorMessage);
+        } finally {
+            if (auditSql != null) {
+                AuditLogHelper.logAuditLog(ctx, auditSql, 
stmtExecutor.getParsedStmt(),
+                        stmtExecutor.getQueryStatisticsForAuditLog(), true);
+            }
         }
     }
 
+    private String buildAuditSql() {
+        TreeMap<Pair<Integer, Integer>, String> replacements = new 
TreeMap<>(new Pair.PairComparator<>());
+        new NereidsParser().parseForEncryption(sql, replacements);
+        List<UnboundTVFRelation> tvfRelations = 
taskCommand.getAllTVFRelation();
+        Preconditions.checkState(replacements.size() == 1 && 
tvfRelations.size() == 1,
+                "S3 streaming insert must contain exactly one TVF");
+        String rewrittenProperties = 
tvfRelations.get(0).getProperties().getMap().entrySet().stream()
+                .map(entry -> 
SqlLiteralUtils.quoteStringLiteral(entry.getKey()) + " = "
+                        + SqlLiteralUtils.quoteStringLiteral(entry.getValue()))
+                .collect(Collectors.joining(", "));
+        Pair<Integer, Integer> tvfPropertiesRange = replacements.firstKey();
+        replacements.replace(tvfPropertiesRange, rewrittenProperties);
+        return BaseViewInfo.rewriteSql(replacements, sql);
+    }
+
     @Override
     public List<Long> getScanBackendIds() {
         if (stmtExecutor != null && stmtExecutor.getCoord() != null) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
index 630b423ce9d..a3eed67b378 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/job/offset/s3/S3SourceOffsetProvider.java
@@ -135,7 +135,7 @@ public class S3SourceOffsetProvider implements 
SourceOffsetProvider {
     public InsertIntoTableCommand rewriteTvfParams(InsertIntoTableCommand 
originCommand,
             Offset runningOffset, long taskId) {
         S3Offset offset = (S3Offset) runningOffset;
-        Map<String, String> props = new HashMap<>();
+        Map<String, String> props = 
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
         // rewrite plan
         Plan rewritePlan = originCommand.getParsedPlan().get().rewriteUp(plan 
-> {
             if (plan instanceof UnboundTVFRelation) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
new file mode 100644
index 00000000000..e07caa6edd1
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertTaskAuditTest.java
@@ -0,0 +1,270 @@
+// 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.
+
+package org.apache.doris.job.extensions.insert.streaming;
+
+import org.apache.doris.analysis.StmtType;
+import org.apache.doris.analysis.UserIdentity;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.common.profile.SummaryProfile;
+import org.apache.doris.datasource.CatalogMgr;
+import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.job.exception.JobException;
+import org.apache.doris.job.extensions.insert.InsertTask;
+import org.apache.doris.job.offset.Offset;
+import org.apache.doris.job.offset.SourceOffsetProvider;
+import org.apache.doris.job.offset.jdbc.JdbcTvfSourceOffsetProvider;
+import org.apache.doris.job.offset.s3.S3Offset;
+import org.apache.doris.job.offset.s3.S3SourceOffsetProvider;
+import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.analyzer.UnboundTVFRelation;
+import org.apache.doris.nereids.glue.LogicalPlanAdapter;
+import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.trees.expressions.Properties;
+import org.apache.doris.nereids.trees.plans.RelationId;
+import 
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.plugin.AuditEvent;
+import org.apache.doris.qe.AuditLogHelper;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.QueryState;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedConstruction;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.atomic.AtomicLong;
+
+public class StreamingInsertTaskAuditTest {
+    private static final String ORIGIN_URI = "s3://bucket/input/*.csv";
+    private static final String RESOLVED_URI = 
"s3://bucket/input/{1.csv,2.csv}";
+    private static final String S3_SQL = "insert into target_table select * 
from s3("
+            + "\"uri\" = \"" + ORIGIN_URI + "\", "
+            + "\"s3.secret_key\" = \"private-value\", "
+            + "\"enclose\" = \"\\\"\")";
+
+    @Test
+    public void testS3RunSubmitsAuditEvent() throws Exception {
+        AuditEvent auditEvent = runS3Task(null);
+
+        Assertions.assertEquals(AuditEvent.EventType.AFTER_QUERY, 
auditEvent.type);
+        Assertions.assertEquals(StmtType.INSERT.name(), auditEvent.stmtType);
+        Assertions.assertFalse(auditEvent.stmt.contains(ORIGIN_URI));
+        Assertions.assertTrue(auditEvent.stmt.contains(RESOLVED_URI));
+        Assertions.assertFalse(auditEvent.stmt.contains("private-value"));
+        Assertions.assertEquals("OK", auditEvent.state);
+        Assertions.assertTrue(auditEvent.isInternal);
+    }
+
+    @Test
+    public void testFailedS3RunSubmitsErrorAuditEvent() throws Exception {
+        AuditEvent auditEvent = runS3Task(new RuntimeException("insert 
failed"));
+
+        Assertions.assertEquals(AuditEvent.EventType.AFTER_QUERY, 
auditEvent.type);
+        Assertions.assertEquals(StmtType.INSERT.name(), auditEvent.stmtType);
+        Assertions.assertEquals("ERR", auditEvent.state);
+        Assertions.assertTrue(auditEvent.errorMessage.contains("insert 
failed"));
+        Assertions.assertTrue(auditEvent.isInternal);
+    }
+
+    @Test
+    public void testCdcRunDoesNotSubmitAuditEvent() throws Exception {
+        String sql = "insert into target_table select * from cdc_stream("
+                + "\"type\" = \"mysql\", \"jdbc_url\" = 
\"jdbc:mysql://127.0.0.1:3306\", "
+                + "\"table\" = \"source_table\", \"offset\" = \"latest\")";
+        runTask(null, sql, Mockito.mock(JdbcTvfSourceOffsetProvider.class), 
Mockito.mock(Offset.class), false);
+    }
+
+    @Test
+    public void testAuditParsingFailureDoesNotFailInsert() throws Exception {
+        runWithAuditPreparationFailure(true);
+    }
+
+    @Test
+    public void testAuditRangeFailureDoesNotFailInsert() throws Exception {
+        runWithAuditPreparationFailure(false);
+    }
+
+    private void runWithAuditPreparationFailure(boolean parserFailure) throws 
Exception {
+        ConnectContext ctx = Mockito.mock(ConnectContext.class);
+        QueryState state = new QueryState();
+        Mockito.when(ctx.getState()).thenReturn(state);
+        StreamingJobProperties properties = 
Mockito.mock(StreamingJobProperties.class);
+        SourceOffsetProvider provider = 
Mockito.mock(SourceOffsetProvider.class);
+        Mockito.when(provider.getSourceType()).thenReturn("s3");
+        S3Offset offset = new S3Offset();
+        offset.setFileLists(RESOLVED_URI);
+        Mockito.when(provider.getNextOffset(Mockito.eq(properties), 
Mockito.anyMap())).thenReturn(offset);
+        InsertIntoTableCommand baseCommand = 
Mockito.mock(InsertIntoTableCommand.class);
+        
Mockito.when(baseCommand.getParsedPlan()).thenReturn(Optional.of(Mockito.mock(LogicalPlan.class)));
+        InsertIntoTableCommand taskCommand = 
Mockito.mock(InsertIntoTableCommand.class);
+        UnboundTVFRelation tvf = Mockito.mock(UnboundTVFRelation.class);
+        
Mockito.when(taskCommand.getAllTVFRelation()).thenReturn(Collections.singletonList(tvf));
+        Mockito.when(provider.rewriteTvfParams(Mockito.eq(baseCommand), 
Mockito.eq(offset), Mockito.anyLong()))
+                .thenReturn(taskCommand);
+        Mockito.doAnswer(invocation -> {
+            state.setOk();
+            return null;
+        }).when(taskCommand).run(Mockito.eq(ctx), 
Mockito.any(StmtExecutor.class));
+
+        try (MockedStatic<InsertTask> insertTask = 
Mockito.mockStatic(InsertTask.class);
+                MockedConstruction<StmtExecutor> executors = 
Mockito.mockConstruction(StmtExecutor.class);
+                MockedConstruction<NereidsParser> parsers = 
Mockito.mockConstruction(NereidsParser.class,
+                        (parser, construction) -> {
+                            
Mockito.when(parser.parseSingle(S3_SQL)).thenReturn(baseCommand);
+                            if (parserFailure) {
+                                
Mockito.when(parser.parseForEncryption(Mockito.eq(S3_SQL), Mockito.anyMap()))
+                                        .thenThrow(new 
IllegalStateException("audit parsing failed"));
+                            }
+                            // Otherwise leave replacements empty to exercise 
the range assertion.
+                        });
+                MockedStatic<AuditLogHelper> audit = 
Mockito.mockStatic(AuditLogHelper.class)) {
+            insertTask.when(() -> 
InsertTask.makeConnectContext(UserIdentity.ROOT, "test_db")).thenReturn(ctx);
+            StreamingInsertTask task = new StreamingInsertTask(1L, 2L, S3_SQL, 
provider, "test_db", properties,
+                    Collections.emptyMap(), UserIdentity.ROOT, null);
+            // A retry must not audit files from an earlier attempt after 
preparation fails.
+            Deencapsulation.setField(task, "auditSql", "stale audit SQL from a 
previous attempt");
+            task.before();
+            Assertions.assertNull(task.getAuditSql());
+            Assertions.assertEquals(2, parsers.constructed().size());
+            Mockito.verify(parsers.constructed().get(0)).parseSingle(S3_SQL);
+            
Mockito.verify(parsers.constructed().get(1)).parseForEncryption(Mockito.eq(S3_SQL),
 Mockito.anyMap());
+            task.run();
+            Assertions.assertEquals(QueryState.MysqlStateType.OK, 
state.getStateType());
+            Mockito.verify(taskCommand).run(ctx, 
executors.constructed().get(1));
+            audit.verifyNoInteractions();
+        }
+    }
+
+    @Test
+    public void testRewriteS3UriCaseInsensitive() {
+        Map<String, String> originProperties = new HashMap<>();
+        originProperties.put("URI", ORIGIN_URI);
+        UnboundTVFRelation originTvf = new UnboundTVFRelation(
+                new RelationId(1), "s3", new Properties(originProperties));
+        InsertIntoTableCommand originCommand = 
Mockito.mock(InsertIntoTableCommand.class);
+        
Mockito.when(originCommand.getParsedPlan()).thenReturn(Optional.of(originTvf));
+
+        S3Offset offset = new S3Offset();
+        offset.setFileLists(RESOLVED_URI);
+        Env env = Mockito.mock(Env.class);
+        Mockito.when(env.isMaster()).thenReturn(false);
+        try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+            mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+            InsertIntoTableCommand rewritten = new S3SourceOffsetProvider()
+                    .rewriteTvfParams(originCommand, offset, 1L);
+
+            Map<String, String> rewrittenProperties =
+                    
rewritten.getAllTVFRelation().get(0).getProperties().getMap();
+            Assertions.assertEquals(1, rewrittenProperties.size());
+            Assertions.assertEquals(RESOLVED_URI, 
rewrittenProperties.get("URI"));
+        }
+    }
+
+    private AuditEvent runS3Task(RuntimeException commandFailure) throws 
Exception {
+        S3Offset offset = new S3Offset();
+        offset.setFileLists(RESOLVED_URI);
+        return runTask(commandFailure, S3_SQL, new S3SourceOffsetProvider(), 
offset, true);
+    }
+
+    private AuditEvent runTask(RuntimeException commandFailure, String sql,
+            SourceOffsetProvider offsetProvider, Offset offset, boolean 
expectAudit) throws Exception {
+        Env env = Mockito.mock(Env.class);
+        CatalogMgr catalogMgr = Mockito.mock(CatalogMgr.class);
+        InternalCatalog catalog = Mockito.mock(InternalCatalog.class);
+        WorkloadRuntimeStatusMgr statusMgr = 
Mockito.mock(WorkloadRuntimeStatusMgr.class);
+        Mockito.when(env.getCatalogMgr()).thenReturn(catalogMgr);
+        Mockito.when(env.getInternalCatalog()).thenReturn(catalog);
+        
Mockito.when(catalogMgr.getCatalog(Mockito.anyString())).thenReturn(catalog);
+        Mockito.when(catalog.getName()).thenReturn("internal");
+        Mockito.when(env.getWorkloadRuntimeStatusMgr()).thenReturn(statusMgr);
+
+        try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+            mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+
+            ConnectContext ctx = 
InsertTask.makeConnectContext(UserIdentity.ROOT, "test_db");
+            ctx.getState().setOk();
+            InsertIntoTableCommand command = 
Mockito.mock(InsertIntoTableCommand.class);
+            if (expectAudit) {
+                Map<String, String> rewrittenTvfProps = new HashMap<>();
+                rewrittenTvfProps.put("uri", RESOLVED_URI);
+                rewrittenTvfProps.put("s3.secret_key", "private-value");
+                rewrittenTvfProps.put("enclose", "\"");
+                UnboundTVFRelation tvf = 
Mockito.mock(UnboundTVFRelation.class);
+                Mockito.when(tvf.getProperties()).thenReturn(new 
Properties(rewrittenTvfProps));
+                
Mockito.when(command.getAllTVFRelation()).thenReturn(Collections.singletonList(tvf));
+            }
+            AtomicLong commandStartTime = new AtomicLong();
+            Mockito.doAnswer(invocation -> {
+                commandStartTime.set(ctx.getStartTime());
+                if (commandFailure != null) {
+                    throw commandFailure;
+                }
+                return null;
+            }).when(command).run(Mockito.eq(ctx), 
Mockito.any(StmtExecutor.class));
+
+            LogicalPlanAdapter parsedStmt = new LogicalPlanAdapter(
+                    new NereidsParser().parseSingle(sql), new 
StatementContext());
+            StmtExecutor executor = Mockito.mock(StmtExecutor.class);
+            Mockito.when(executor.getParsedStmt()).thenReturn(parsedStmt);
+            
Mockito.when(executor.getSummaryProfile()).thenReturn(Mockito.mock(SummaryProfile.class));
+            ctx.setExecutor(executor);
+
+            StreamingInsertTask task = new StreamingInsertTask(
+                    1L, 2L, sql, offsetProvider, "test_db", null,
+                    Collections.emptyMap(), UserIdentity.ROOT, null);
+            Deencapsulation.setField(task, "ctx", ctx);
+            Deencapsulation.setField(task, "taskCommand", command);
+            Deencapsulation.setField(task, "stmtExecutor", executor);
+            Deencapsulation.setField(task, "runningOffset", offset);
+            if (expectAudit) {
+                Deencapsulation.setField(task, "auditSql",
+                        Deencapsulation.invoke(task, "buildAuditSql"));
+            }
+
+            if (commandFailure == null) {
+                task.run();
+            } else {
+                Assertions.assertThrows(JobException.class, task::run);
+            }
+
+            if (!expectAudit) {
+                Mockito.verify(statusMgr, Mockito.never())
+                        
.submitFinishQueryToAudit(Mockito.any(AuditEvent.class));
+                return null;
+            }
+            ArgumentCaptor<AuditEvent> auditEventCaptor = 
ArgumentCaptor.forClass(AuditEvent.class);
+            
Mockito.verify(statusMgr).submitFinishQueryToAudit(auditEventCaptor.capture());
+            AuditEvent auditEvent = auditEventCaptor.getValue();
+            Assertions.assertTrue(commandStartTime.get() > 0);
+            Assertions.assertEquals(commandStartTime.get(), 
auditEvent.timestamp);
+            return auditEvent;
+        } finally {
+            ConnectContext.remove();
+        }
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to