This is an automated email from the ASF dual-hosted git repository.

dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 502268a67e [INLONG-8522][Manager] optimize log print and status 
judgement in preProcessLabelFileTasks (#8523)
502268a67e is described below

commit 502268a67e59dc4fc7017a94b25b8063f2edfa6b
Author: kipshi <[email protected]>
AuthorDate: Thu Jul 13 17:37:05 2023 +0800

    [INLONG-8522][Manager] optimize log print and status judgement in 
preProcessLabelFileTasks (#8523)
    
    Co-authored-by: kipshi <[email protected]>
---
 .../service/core/impl/AgentServiceImpl.java        | 12 ++++----
 .../service/core/impl/AgentServiceTest.java        | 34 +++++++++++++---------
 2 files changed, 27 insertions(+), 19 deletions(-)

diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java
index 174929741e..b20d76a8c0 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java
@@ -202,7 +202,7 @@ public class AgentServiceImpl implements AgentService {
 
         // Update task status, other tasks with status 20x will change to 30x 
in next request
         if (CollectionUtils.isEmpty(request.getCommandInfo())) {
-            LOGGER.info("task result was empty in request: {}, just return", 
request);
+            LOGGER.debug("task result was empty in request: {}, just return", 
request);
             return;
         }
         for (CommandEntity command : request.getCommandInfo()) {
@@ -339,7 +339,7 @@ public class AgentServiceImpl implements AgentService {
     private List<DataConfig> processQueuedTasks(TaskRequest request) {
         HashSet<SourceStatus> needAddStatusSet = 
Sets.newHashSet(SourceStatus.TOBE_ISSUED_SET);
         if (PullJobTypeEnum.NEVER == 
PullJobTypeEnum.getPullJobType(request.getPullJobType())) {
-            LOGGER.warn("agent pull job type is [NEVER], just pull to be 
active tasks");
+            LOGGER.debug("agent pull job type is [NEVER], just pull to be 
active tasks");
             needAddStatusSet.remove(SourceStatus.TO_BE_ISSUED_ADD);
         }
 
@@ -369,7 +369,7 @@ public class AgentServiceImpl implements AgentService {
     private void preProcessNonFileTasks(TaskRequest taskRequest) {
         List<Integer> needAddStatusList;
         if (PullJobTypeEnum.NEVER == 
PullJobTypeEnum.getPullJobType(taskRequest.getPullJobType())) {
-            LOGGER.warn("agent pull job type is [NEVER], just pull to be 
active tasks");
+            LOGGER.debug("agent pull job type is [NEVER], just pull to be 
active tasks");
             needAddStatusList = 
Collections.singletonList(SourceStatus.TO_BE_ISSUED_ACTIVE.getCode());
         } else {
             needAddStatusList = 
Arrays.asList(SourceStatus.TO_BE_ISSUED_ADD.getCode(),
@@ -482,12 +482,12 @@ public class AgentServiceImpl implements AgentService {
                     SourceStatus.SOURCE_NORMAL,
                     SourceStatus.TO_BE_ISSUED_ADD,
                     SourceStatus.TO_BE_ISSUED_ACTIVE);
-            Set<GroupStatus> exceptedMatchedGroupStatus = 
Sets.newHashSet(GroupStatus.SUSPENDED,
-                    GroupStatus.SUSPENDING);
+            Set<GroupStatus> matchedGroupStatus = Sets.newHashSet(
+                    GroupStatus.CONFIG_SUCCESSFUL, GroupStatus.RESTARTED);
             if (matchGroup(sourceEntity, clusterNodeEntity)
                     && groupEntity != null
                     && 
!exceptedMatchedSourceStatus.contains(SourceStatus.forCode(sourceEntity.getStatus()))
-                    && 
!exceptedMatchedGroupStatus.contains(GroupStatus.forCode(groupEntity.getStatus())))
 {
+                    && 
matchedGroupStatus.contains(GroupStatus.forCode(groupEntity.getStatus()))) {
                 LOGGER.info("Transform task({}) from {} to {} because tag 
rematch "
                         + "for agent({}) in cluster({})", 
sourceEntity.getAgentIp(),
                         sourceEntity.getStatus(), 
SourceStatus.TO_BE_ISSUED_ACTIVE.getCode(),
diff --git 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java
 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java
index d27c018653..6f0cbebaf3 100644
--- 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java
+++ 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java
@@ -48,7 +48,6 @@ import org.apache.commons.lang3.tuple.ImmutablePair;
 import org.apache.commons.lang3.tuple.Pair;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.springframework.beans.factory.annotation.Autowired;
@@ -161,18 +160,13 @@ class AgentServiceTest extends ServiceBaseTest {
         sourceService.logicDeleteAll(groupId, streamId, GLOBAL_OPERATOR);
     }
 
-    @BeforeAll
-    public static void setUp(
-            @Autowired AgentService agentService,
-            @Autowired HeartbeatService heartbeatService) {
-        agent = new MockAgent(agentService, heartbeatService, 2);
-        agent.sendHeartbeat();
-    }
-
     @BeforeEach
     public void setupEach() {
+        agent = new MockAgent(agentService, heartbeatService, 2);
+        agent.sendHeartbeat();
         groupStreamCache = new ArrayList<>();
         groupCache = new ArrayList<>();
+        super.login();
     }
 
     @AfterEach
@@ -251,13 +245,27 @@ class AgentServiceTest extends ServiceBaseTest {
         // bind group and rematch
         bindGroup(true, "group1");
         TaskResult t2 = agent.pullTask();
-        Assertions.assertEquals(1, t2.getDataConfigs().size());
-        Assertions.assertEquals(1, t2.getDataConfigs().stream()
+        Assertions.assertEquals(0, t2.getDataConfigs().size());
+        Assertions.assertEquals(0, t2.getDataConfigs().stream()
+                .filter(dataConfig -> Integer.valueOf(dataConfig.getOp()) == 
ManagerOpEnum.ACTIVE.getType())
+                .collect(Collectors.toSet())
+                .size());
+
+        // update group to config success
+        final String groupId = sourceService.listSource(groupStream.getLeft(), 
groupStream.getRight()).stream()
+                .filter(source -> source.getTemplateId() != null)
+                .findAny()
+                .get()
+                .getInlongGroupId();
+        groupMapper.updateStatus(groupId, 
GroupStatus.CONFIG_SUCCESSFUL.getCode(), GLOBAL_OPERATOR);
+        TaskResult t3 = agent.pullTask();
+        Assertions.assertEquals(1, t3.getDataConfigs().size());
+        Assertions.assertEquals(1, t3.getDataConfigs().stream()
                 .filter(dataConfig -> Integer.valueOf(dataConfig.getOp()) == 
ManagerOpEnum.ACTIVE.getType())
                 .collect(Collectors.toSet())
                 .size());
-        DataConfig d2 = t2.getDataConfigs().get(0);
-        Assertions.assertEquals(sourceId, d2.getTaskId());
+        DataConfig d3 = t3.getDataConfigs().get(0);
+        Assertions.assertEquals(sourceId, d3.getTaskId());
     }
 
     /**

Reply via email to