DATABASE_HOST_TO_CONNECTOR_MAP = [
    'qa.db.qaorch.com': 'qa-kafka-connect-debezium-ar'
]

DATABASE_HOST_TO_TOPICS_MAP = [
    'qa.db.qaorch.com': '_kafka_connect_debezium_mysql_source_art_relations_streams-offsets'
]

pipeline {
    agent any
    options { timestamps () }
    triggers {
     parameterizedCron('''
30 01 * * * %DB_HOST=qa.db.qaorch.com
00 03 * * * %DB_HOST=data.db.devorch.com
00 04 * * * %DB_HOST=phys.db.devorch.com
00 05 * * * %DB_HOST=build.db.devorch.com
00 07 * * * %DB_HOST=distro.db.devorch.com
        ''')
    }
    parameters {
        choice(
            name: 'DB_HOST',
            choices: [
                'accounting.db.devorch.com',
                'build.db.devorch.com',
                'concepts.db.devorch.com',
                'content.db.devorch.com',
                'data.db.devorch.com',
                'distro.db.devorch.com',
                'distro.dddb.devorch.com',
                'infosec.db.devorch.com',
                'phys.db.devorch.com',
                'qa.db.qaorch.com',
            ],
            description: 'Database host to refresh'
        )
        string(
            name: 'SCALA_VERSION',
            defaultValue: '2.13',
            description: 'Version of scala for kafka cli tools'
        )
        string(
            name: 'KAFKA_VERSION',
            defaultValue: '3.5.2',
            description: 'Version of kafka for kafka cli tools'
        )
    }
    environment {
        AWS_DEFAULT_REGION = 'us-east-1'
        CLUSTER_NAME = DATABASE_HOST_TO_CONNECTOR_MAP.get(DB_HOST)
        TOPIC_NAMES = DATABASE_HOST_TO_TOPICS_MAP.get(DB_HOST)
        BOOTSTRAP_SERVERS = "b-3.qa-managed-kafka-cdc-d.2kgc64.c6.kafka.us-east-1.amazonaws.com:9094,b-2.qa-managed-kafka-cdc-d.2kgc64.c6.kafka.us-east-1.amazonaws.com:9094,b-1.qa-managed-kafka-cdc-d.2kgc64.c6.kafka.us-east-1.amazonaws.com:9094"
    }
    stages {
        stage('Pause Replication') {
            when {
                expression {
                    return DATABASE_HOST_TO_CONNECTOR_MAP.get(env.DB_HOST);
                }
            }
            steps {
                node('aws') {
                    echo "Assuming AWS Role..."
                    script {
                        AWS_PARAMS =  sh(
                            script: '''
                            aws sts assume-role \
                                --role-arn arn:aws:iam::437795906767:role/qa-db-refresh-pipeline \
                                --role-session-name lambda | \
                                    jq -r '.Credentials | \
                                    "\\(.AccessKeyId)\n\\(.SecretAccessKey)\n\\(.SessionToken)"'
                            ''',
                            returnStdout: true
                        ).trim().tokenize("\n")
                        env.AWS_ACCESS_KEY_ID = AWS_PARAMS[0]
                        env.AWS_SECRET_ACCESS_KEY = AWS_PARAMS[1]
                        env.AWS_SESSION_TOKEN = AWS_PARAMS[2]
                        env.SERVICE_NAME = sh(
                            script: '''
                            aws ecs list-services \
                                --cluster "$CLUSTER_NAME" \
                                --query "serviceArns[0]" | tr -d '"'
                            ''',
                            returnStdout: true
                        ).trim()
                        env.ORIGINAL_DESIRED_COUNT = sh(
                            script: '''
                                aws ecs describe-services \
                                    --cluster "$CLUSTER_NAME" \
                                    --services "$SERVICE_NAME" \
                                    --query "services[0].desiredCount"
                            ''',
                            returnStdout: true
                        ).trim()
                    }
                    echo "Pausing Replication..."
                    timeout(unit: 'SECONDS', time: 600) {
                        sh('''
                            if [ ${ORIGINAL_DESIRED_COUNT} = 0 ]; then
                                exit 0
                            fi
                            # Update service to 0 desired instances
                            aws ecs update-service \
                                --cluster "$CLUSTER_NAME" \
                                --service "$SERVICE_NAME" \
                                --desired-count 0 > /dev/null
                            RUNNING_COUNT="$ORIGINAL_DESIRED_COUNT"
                            while [ "${RUNNING_COUNT}" -gt 0 ]; do
                                RUNNING_COUNT=$(
                                    aws ecs describe-services \
                                        --cluster "$CLUSTER_NAME" \
                                        --services "$SERVICE_NAME" \
                                        --query "services[0].runningCount"
                                    )
                                    if [ "${RUNNING_COUNT}" -gt 0 ]; then
                                        sleep 10
                                    fi
                            done
                        ''')
                    }
                }
                node('aws') {
                    echo "Resetting Kafka Topic..."
                    timeout(unit: 'SECONDS', time: 1800) {
                        sh('''#!/bin/bash -e
                            KAFKA_TOOLS_VERSION="kafka_$SCALA_VERSION-$KAFKA_VERSION"
                            if [ ! -d "${KAFKA_TOOLS_VERSION}" ]; then
                                echo "Downloading kafka cli tools."
                                wget -O "$KAFKA_TOOLS_VERSION.tgz" "https://dlcdn.apache.org/kafka/$KAFKA_VERSION/$KAFKA_TOOLS_VERSION.tgz"
                                tar -zxf "$KAFKA_TOOLS_VERSION.tgz"
                                echo "security.protocol=SSL" > "$KAFKA_TOOLS_VERSION/fake.config"
                            fi

                            cd "$KAFKA_TOOLS_VERSION" || exit 1
                            IFS=, read -ra TOPIC_NAME_LIST <<< "$TOPIC_NAMES"
                            for TOPIC_NAME in "${TOPIC_NAME_LIST[@]}"; do
                                echo "Deleting $TOPIC_NAME"
                                ./bin/kafka-topics.sh \
                                    --command-config fake.config \
                                    --bootstrap-server "$BOOTSTRAP_SERVERS" \
                                    --topic "$TOPIC_NAME" \
                                    --delete || true
                            done
                            sleep 15
                            for TOPIC_NAME in "${TOPIC_NAME_LIST[@]}"; do
                                echo "Waiting for $TOPIC_NAME to be deleted"
                                TOPIC_FOUND="exists"
                                while [ ! -z "${TOPIC_FOUND}" ]; do
                                    TOPIC_FOUND=$(
                                        ./bin/kafka-topics.sh \
                                            --command-config "fake.config" \
                                            --bootstrap-server "$BOOTSTRAP_SERVERS" \
                                            --list | (grep "^$TOPIC_NAME" || echo "")
                                    )
                                    if [ ! -z "${TOPIC_FOUND}" ]; then
                                        ./bin/kafka-topics.sh \
                                            --command-config fake.config \
                                            --bootstrap-server "$BOOTSTRAP_SERVERS" \
                                            --topic "$TOPIC_NAME" \
                                            --delete || true
                                        sleep 15
                                    fi
                                done
                            done
                        ''')
                    }
                }
            }
        }
        stage('Refresh Database') {
            steps {
                script{
                    currentBuild.description = env.DB_HOST
                }
                node('ny-deploy') {
                    sshagent (credentials: ['d24cef52-81a6-4ed2-9b46-07fa885a17a3']) {
                        sh '''
                            SSH_CONN="root@${DB_HOST}"
                            REFRESH_COMMAND=$(ssh -o UserKnownHostsFile=/dev/null "$SSH_CONN" "find /root -type f -name refresh.sh")
                            ssh -o UserKnownHostsFile=/dev/null -o ServerAliveInterval=30 "$SSH_CONN" "$REFRESH_COMMAND"
                        '''
                    }
                }
            }
        }
        stage('Resume Replication') {
            when {
                expression {
                    return DATABASE_HOST_TO_CONNECTOR_MAP.get(env.DB_HOST);
                }
            }
            steps {
                echo "Resuming Replication..."
                node('aws') {
                    timeout(unit: 'SECONDS', time: 600) {
                        sh('''
                            if [ ${ORIGINAL_DESIRED_COUNT} = 0 ]; then
                                ORIGINAL_DESIRED_COUNT=2
                            fi
                            aws ecs update-service \
                                --cluster "$CLUSTER_NAME" \
                                --service "$SERVICE_NAME" \
                                --desired-count $ORIGINAL_DESIRED_COUNT > /dev/null
                            RUNNING_COUNT=0
                            while [ ${RUNNING_COUNT} -lt 1 ]; do
                                RUNNING_COUNT=$(
                                    aws ecs describe-services \
                                        --cluster "$CLUSTER_NAME" \
                                        --services "$SERVICE_NAME" \
                                        --query "services[0].runningCount"
                                )
                                echo "Found $RUNNING_COUNT running instances. Waiting for 1 or more."
                                if [ "${RUNNING_COUNT}" -lt 1 ]; then
                                    sleep 10
                                fi
                            done
                        '''
                        )
                    }
                }
            }
        }
        stage('Post-refresh Scripts') {
            when {
                expression {
                    env.DB_HOST == 'qa.db.qaorch.com'
                }
            }
            steps {
                build job: 'content-review-queue-reindex', parameters: [
                    [$class: 'StringParameterValue', name: 'Environment', value: 'qa'],
                    [$class: 'BooleanParameterValue', name: 'DB_SYNC', value: true],
                    [$class: 'BooleanParameterValue', name: 'CLEAR_INDEX', value: false]
                ], wait: false
            }
        }
    }
}
