Back to Repositories

Testing Event Queue Resolution Implementation in Conductor-OSS

The EventQueueResolutionTest suite validates the event queue resolution functionality in Conductor’s task execution system, focusing on queue name computation and queue retrieval operations. This test suite ensures proper handling of dynamic queue names and different queue provider implementations.

Test Coverage Overview

The test suite provides comprehensive coverage of event queue resolution mechanisms:
  • Queue name computation for different sink parameters
  • Dynamic queue name resolution with variable substitution
  • Integration with multiple queue providers (SQS, Conductor)
  • Validation of queue properties and configurations
  • Error handling and task status management

Implementation Analysis

The testing approach utilizes Spring’s test context framework with JUnit4, implementing mock queue providers for isolated testing. The suite employs parameterized testing patterns to verify queue resolution across different scenarios, with particular attention to dynamic value substitution and provider-specific queue naming conventions.

Technical Details

Testing infrastructure includes:
  • Spring Test Context framework
  • JUnit 4 test runner
  • Mock queue providers for SQS and Conductor
  • ObjectMapper for JSON handling
  • ParametersUtils for dynamic parameter resolution

Best Practices Demonstrated

The test suite exemplifies several testing best practices:
  • Proper test isolation through mock providers
  • Comprehensive setup and teardown management
  • Clear test method naming and organization
  • Thorough assertion coverage for both positive and negative cases
  • Effective use of Spring’s dependency injection in tests

conductor-oss/conductor

core/src/test/java/com/netflix/conductor/core/execution/tasks/EventQueueResolutionTest.java

            
/*
 * Copyright 2022 Conductor Authors.
 * <p>
 * Licensed 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
 * <p>
 * http://www.apache.org/licenses/LICENSE-2.0
 * <p>
 * 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 com.netflix.conductor.core.execution.tasks;

import java.util.HashMap;
import java.util.Map;

import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringRunner;

import com.netflix.conductor.common.config.TestObjectMapperConfiguration;
import com.netflix.conductor.common.metadata.tasks.TaskType;
import com.netflix.conductor.common.metadata.workflow.WorkflowDef;
import com.netflix.conductor.core.events.EventQueueProvider;
import com.netflix.conductor.core.events.EventQueues;
import com.netflix.conductor.core.events.MockQueueProvider;
import com.netflix.conductor.core.events.queue.ObservableQueue;
import com.netflix.conductor.core.utils.ParametersUtils;
import com.netflix.conductor.model.TaskModel;
import com.netflix.conductor.model.WorkflowModel;

import com.fasterxml.jackson.databind.ObjectMapper;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;

/**
 * Tests the {@link Event#computeQueueName(WorkflowModel, TaskModel)} and {@link
 * Event#getQueue(String, String)} methods with a real {@link ParametersUtils} object.
 */
@ContextConfiguration(classes = {TestObjectMapperConfiguration.class})
@RunWith(SpringRunner.class)
public class EventQueueResolutionTest {

    private WorkflowDef testWorkflowDefinition;
    private EventQueues eventQueues;
    private ParametersUtils parametersUtils;

    @Autowired private ObjectMapper objectMapper;

    @Before
    public void setup() {
        Map<String, EventQueueProvider> providers = new HashMap<>();
        providers.put("sqs", new MockQueueProvider("sqs"));
        providers.put("conductor", new MockQueueProvider("conductor"));

        parametersUtils = new ParametersUtils(objectMapper);
        eventQueues = new EventQueues(providers, parametersUtils);

        testWorkflowDefinition = new WorkflowDef();
        testWorkflowDefinition.setName("testWorkflow");
        testWorkflowDefinition.setVersion(2);
    }

    @Test
    public void testSinkParam() {
        String sink = "sqs:queue_name";

        WorkflowDef def = new WorkflowDef();
        def.setName("wf0");

        WorkflowModel workflow = new WorkflowModel();
        workflow.setWorkflowDefinition(def);

        TaskModel task1 = new TaskModel();
        task1.setReferenceTaskName("t1");
        task1.addOutput("q", "t1_queue");
        workflow.getTasks().add(task1);

        TaskModel task2 = new TaskModel();
        task2.setReferenceTaskName("t2");
        task2.addOutput("q", "task2_queue");
        workflow.getTasks().add(task2);

        TaskModel task = new TaskModel();
        task.setReferenceTaskName("event");
        task.getInputData().put("sink", sink);
        task.setTaskType(TaskType.EVENT.name());
        workflow.getTasks().add(task);

        Event event = new Event(eventQueues, parametersUtils, objectMapper);
        String queueName = event.computeQueueName(workflow, task);
        ObservableQueue queue = event.getQueue(queueName, task.getTaskId());
        assertNotNull(task.getReasonForIncompletion(), queue);
        assertEquals("queue_name", queue.getName());
        assertEquals("sqs", queue.getType());

        sink = "sqs:${t1.output.q}";
        task.getInputData().put("sink", sink);
        queueName = event.computeQueueName(workflow, task);
        queue = event.getQueue(queueName, task.getTaskId());
        assertNotNull(queue);
        assertEquals("t1_queue", queue.getName());
        assertEquals("sqs", queue.getType());

        sink = "sqs:${t2.output.q}";
        task.getInputData().put("sink", sink);
        queueName = event.computeQueueName(workflow, task);
        queue = event.getQueue(queueName, task.getTaskId());
        assertNotNull(queue);
        assertEquals("task2_queue", queue.getName());
        assertEquals("sqs", queue.getType());

        sink = "conductor";
        task.getInputData().put("sink", sink);
        queueName = event.computeQueueName(workflow, task);
        queue = event.getQueue(queueName, task.getTaskId());
        assertNotNull(queue);
        assertEquals(
                workflow.getWorkflowName() + ":" + task.getReferenceTaskName(), queue.getName());
        assertEquals("conductor", queue.getType());

        sink = "sqs:static_value";
        task.getInputData().put("sink", sink);
        queueName = event.computeQueueName(workflow, task);
        queue = event.getQueue(queueName, task.getTaskId());
        assertNotNull(queue);
        assertEquals("static_value", queue.getName());
        assertEquals("sqs", queue.getType());
    }

    @Test
    public void testDynamicSinks() {
        Event event = new Event(eventQueues, parametersUtils, objectMapper);
        WorkflowModel workflow = new WorkflowModel();
        workflow.setWorkflowDefinition(testWorkflowDefinition);

        TaskModel task = new TaskModel();
        task.setReferenceTaskName("task0");
        task.setTaskId("task_id_0");
        task.setStatus(TaskModel.Status.IN_PROGRESS);
        task.getInputData().put("sink", "conductor:some_arbitary_queue");

        String queueName = event.computeQueueName(workflow, task);
        ObservableQueue queue = event.getQueue(queueName, task.getTaskId());
        assertEquals(TaskModel.Status.IN_PROGRESS, task.getStatus());
        assertNotNull(queue);
        assertEquals("testWorkflow:some_arbitary_queue", queue.getName());
        assertEquals("testWorkflow:some_arbitary_queue", queue.getURI());
        assertEquals("conductor", queue.getType());

        task.getInputData().put("sink", "conductor");
        queueName = event.computeQueueName(workflow, task);
        queue = event.getQueue(queueName, task.getTaskId());
        assertEquals(
                "not in progress: " + task.getReasonForIncompletion(),
                TaskModel.Status.IN_PROGRESS,
                task.getStatus());
        assertNotNull(queue);
        assertEquals("testWorkflow:task0", queue.getName());

        task.getInputData().put("sink", "sqs:my_sqs_queue_name");
        queueName = event.computeQueueName(workflow, task);
        queue = event.getQueue(queueName, task.getTaskId());
        assertEquals(
                "not in progress: " + task.getReasonForIncompletion(),
                TaskModel.Status.IN_PROGRESS,
                task.getStatus());
        assertNotNull(queue);
        assertEquals("my_sqs_queue_name", queue.getName());
        assertEquals("sqs", queue.getType());
    }
}