import groovy.json.JsonOutput

RDS_REFRESH_STATE_MACHINE_ARN = 'arn:aws:states:us-east-1:086679231553:stateMachine:shared-rds-refresh-state-machine'
SHARED_ACCOUNT_ID = '086679231553'
STEP_FUNCTION_RUNNER_ROLE = 'shared-rds-refresh-state-machine-runner-role'

SOURCE_ENVIRONMENTS = [
    prod: [accountId: '437795906767'],
    'prod-songwhip': [accountId: '926734670777']
]

TARGET_ENVIRONMENTS = [
    'qa': [
        restoreRole: 'qa-rds-refresh-restore-role',
        dataSanitisationFunction: 'qa-lambda-sanitise-rds-data',
        kafkaBootstrapServers: [
            "b-1.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-3.qa-managed-kafka-cdc-d.2kgc64.c6.kafka.us-east-1.amazonaws.com:9094"
        ]
    ],
    'dev': [
        restoreRole: 'dev-rds-refresh-restore-role',
        dataSanitisationFunction: 'dev-lambda-sanitise-rds-data',
        kafkaBootstrapServers: [
            "b-1.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094",
            "b-2.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094",
            "b-3.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094"
        ]
    ],
    'uat': [
        restoreRole: 'uat-rds-refresh-restore-role',
        dataSanitisationFunction: 'uat-lambda-sanitise-rds-data',
        kafkaBootstrapServers: []
    ]
]

SOURCE_ENVIRONMENT_NAMES = SOURCE_ENVIRONMENTS.keySet() as List
TARGET_ENVIRONMENT_NAMES = TARGET_ENVIRONMENTS.keySet() as List

SERVICES = [
    'art-relations': [
        targetEnvironments: [
            dev: [
                accountId: '103233932089'
            ],
            qa: [
                accountId: '437795906767'
            ]
        ],
        debeziumConnectors: [
            qa: [name: 'qa-kafka-connect-debezium-ar', kafkaConnectGroupId: 'art_relations_streams', connectorName: 'debezium_mysql_source']
        ]
    ],
    'content-review': [
        debeziumConnectors: [
            qa: [name: 'qa-kafka-connect-debezium-cr', kafkaConnectGroupId: 'debezium_content_review', connectorName: 'debezium_mysql_source']
        ]
    ],
    'direct-delivery-cluster': [:],
    'ddex-ingester': [:],
    'integrations': [:],
    'neighbouring-rights-delivery': [:],
    'neighbouring-rights-ownership-delivery': [:],
    'ows-asset-transcoder': [:],
    'ows-assets': [
        debeziumConnectors: [
            dev: [name: 'dev-kafka-connect-debezium-ats', kafkaConnectGroupId: 'orchard', connectorName: 'debezium_mysql_source_ows_assets'],
            qa : [name: 'qa-kafka-connect-debezium-ats', kafkaConnectGroupId: 'ows_assets_ows_assets', connectorName: 'debezium_mysql_source_ows_assets']
        ]
    ],
//    'ows-bulk-encoder': [:], // Postgresql not supported yet
    'ows-collaborator': [
        debeziumConnectors: [
            qa: [name: 'qa-kafka-connect-debezium-cp', kafkaConnectGroupId: 'debezium_ows_collaborator', connectorName: 'debezium_mysql_source']
        ],
        fivetranSync: [
            qa: [enabled: true, connector_id: 'armful_good', historical_sync: 'true']
        ],
    ],
    'ows-content-review': [:],
    'ows-dmp': [:],
//    'ows-grid-generator': [:], // Postgresql not supported yet
    'ows-podcast': [:],
    'ows-pricing': [:],
    'ows-product-store-mapping': [:],
    'ows-store-availability': [:],
    'ows-track': [
        debeziumConnectors: [
            qa: [name: 'qa-kafka-connect-debezium-ot', kafkaConnectGroupId: 'debezium_ows_track', connectorName: 'debezium_mysql_source']
        ]
    ],
    'ows-transcoding': [:],
//    'ows-vector-job-rules': [:], // Not currently supported due to QA and prod being different types. Prod must be converted to a cluster first.
    'ows-video': [:],
    'publishing-compositions': [:],
    'prs-toolkit': [:],
    'royalty-accounting': [
        targetEnvironments: [
            qa: [
                accountId: '989790945997'
            ],
            uat: [
                accountId: '989790945997'
            ]
        ],
        snowflakeRefresh: [
            dev: [enabled: true],
            qa: [enabled: true],
            uat: [enabled: true]
        ],
        debeziumConnectors: [
            qa: [
                name: 'qa-kafka-connect-debezium-ra',
                accountId: '437795906767',
                kafkaConnectGroupId: 'royalty_accounting',
                connectorName: 'debezium_mysql_source'
            ]
        ],
        slack: [
            dev: [notify: false, channel: '#accounting-tech'],
            qa: [notify: true, channel: '#accounting-tech'],
            uat: [notify: true, channel: '#accounting-tech'],
        ],
        redeployServices: [
            'ows-abacus-account': [qa: [account: 'accounting-qa'], uat: [account: 'accounting-qa']],
            'ows-abacus-contract': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-abacus-event': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-abacus-legacy-sync': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-abacus-schedule': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-abacus-state': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-abacus-worksheet': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-ledger': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-payee': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-payment': [qa: [account: 'prod'], uat: [account: 'accounting-qa']],
            'ows-royalties': [qa: [account: 'prod'], uat: [account: 'accounting-qa']]
        ],
        fivetranSync: [
            qa: [enabled: true, connector_id: 'elected_sixtyfold'],
            uat: [enabled: true, connector_id: 'were_ethereal']
        ]
    ],
    'songwhip-api': [
        targetEnvironments: [
            qa: [
                accountId: '619719722105'
            ]
        ],
        fivetranSync: [
            qa: [enabled: true, connector_id: 'spies_flooded']
        ]
    ]
]

SERVICE_NAMES = SERVICES.keySet() as List // Cast to a list because choice parameter definition requires a list

pipeline {
    agent {
        label 'aws'
    }

    options {
        lock "${params.TARGET_ENVIRONMENT}-${params.SERVICE_NAME}"
        timestamps()
    }

    parameters {
        choice(name: 'SERVICE_NAME', choices: SERVICE_NAMES, description: 'The service to refresh the database for.')
        choice(name: 'SOURCE_ENVIRONMENT', choices: SOURCE_ENVIRONMENT_NAMES, description: 'The environment to refresh from.')
        choice(name: 'TARGET_ENVIRONMENT', choices: TARGET_ENVIRONMENT_NAMES, description: 'The environment to refresh to.')
        booleanParam(name: 'FORCE_SNAPSHOT', defaultValue: false, description: 'DEV CLUSTERS ONLY: Force snapshot-based refresh.')
        booleanParam(name: 'SNOWFLAKE_REFRESH', defaultValue: true, description: 'SNOWFLAKE-ENABLED CLUSTERS ONLY: Trigger Snowflake-refresh, has no effect on clusters without Snowflake scripts.')
        string(name: 'RESTORE_SNAPSHOT_ID', defaultValue: '', description: 'Restore from a pre-existing snapshot in the source (prod) account. It is shared with the target account automatically (and, for standalone databases, copied into it). The creating team owns the snapshot lifecycle. Leave blank for a normal refresh.')
        string(name: 'SHARED_LIBRARIES_VERSION', defaultValue: 'master', description: 'The version of the Jenkins shared libraries to use. Can be a branch, tag or Git revision.')
    }

    triggers {
        parameterizedCron('''
            # Add configuration here to perform a refresh on a schedule
            # Each schedule is defined on a separate line.
            30 23 * * * %SERVICE_NAME=art-relations;SOURCE_ENVIRONMENT=prod;TARGET_ENVIRONMENT=qa
            30 22 * * * %SERVICE_NAME=art-relations;SOURCE_ENVIRONMENT=prod;TARGET_ENVIRONMENT=dev;FORCE_SNAPSHOT=true
            30 23 * * * %SERVICE_NAME=direct-delivery-cluster;SOURCE_ENVIRONMENT=prod;TARGET_ENVIRONMENT=qa
            45 23 * * * %SERVICE_NAME=ows-track;SOURCE_ENVIRONMENT=prod;TARGET_ENVIRONMENT=qa
            30  1 * * * %SERVICE_NAME=ows-collaborator;SOURCE_ENVIRONMENT=prod;TARGET_ENVIRONMENT=qa
        ''')
    }

    environment {
        AWS_DEFAULT_REGION = 'us-east-1'
        SOURCE_DB_NAME = "${params.SOURCE_ENVIRONMENT.split('-')[0]}-${params.SERVICE_NAME}"
        TARGET_DB_NAME = "${params.TARGET_ENVIRONMENT.split('-')[0]}-${params.SERVICE_NAME}"
    }

    stages {
        stage('Load Shared Libraries') {
            steps {
                library "jenkins-global-libraries@${params.SHARED_LIBRARIES_VERSION}"
            }
        }
        stage('Notify Start') {
            when {
                expression {
                    def slackConfig = SERVICES[params.SERVICE_NAME]?.slack?.get(params.TARGET_ENVIRONMENT)
                    slackConfig?.notify
                }
            }
            steps {
                script {
                  def slackConfig = SERVICES[params.SERVICE_NAME]?.slack?.get(params.TARGET_ENVIRONMENT)
                  slackSend channel: slackConfig.channel, color: 'good', message: """${env.JOB_NAME} Started (<${env.BUILD_URL}|Open>)
*Service:* ${env.SERVICE_NAME}
*Source Environment:* ${env.SOURCE_ENVIRONMENT}
*Target Environment:* ${env.TARGET_ENVIRONMENT}"""
                }
            }
        }
        stage('Verify Sanitise Scripts') {
            when {
                not { triggeredBy 'GitHubPushCause' }
            }
            steps {
                script {
                    currentBuild.description = "${params.SERVICE_NAME}: ${params.SOURCE_ENVIRONMENT} -> ${params.TARGET_ENVIRONMENT}"
                    String scriptsDir = 'lambda/sanitise_rds_data/scripts'
                    if (!fileExists("${scriptsDir}/${TARGET_DB_NAME}")) {
                        error """No data sanitisation scripts found for target database ${TARGET_DB_NAME} in scripts directory ${scriptsDir}. \
See https://github.com/theorchard/python-rds-utils/tree/master/${scriptsDir} for more information.
"""
                    }
                }
            }
        }
        stage('Refresh') {
            when {
                not { triggeredBy 'GitHubPushCause' }
            }
            environment {
                STATE_MACHINE_ARN = "${RDS_REFRESH_STATE_MACHINE_ARN}"
                STATE_MACHINE_EXECUTION_NAME = "${env.BUILD_TAG}"
                STATE_MACHINE_INPUT = "${getStateMachineInput()}"
                STATE_MACHINE_POLLING_INTERVAL = "10"
            }
            steps {
                withEcr {
                    withAWS(
                        duration: 10800,
                        role: STEP_FUNCTION_RUNNER_ROLE,
                        roleAccount: SHARED_ACCOUNT_ID,
                        roleSessionName: env.BUILD_TAG,
                        useNode: true
                    ) {
                        sh '''
                            docker run --rm \
                                --pull always \
                                -e AWS_ACCESS_KEY_ID \
                                -e AWS_SECRET_ACCESS_KEY \
                                -e AWS_SESSION_TOKEN \
                                -e STATE_MACHINE_ARN \
                                -e STATE_MACHINE_EXECUTION_NAME \
                                -e STATE_MACHINE_INPUT \
                                -e STATE_MACHINE_POLLING_INTERVAL \
                                086679231553.dkr.ecr.us-east-1.amazonaws.com/state-machine-runner:latest
                        '''
                    }
                }
            }
        }
        stage('Post-refresh Scripts') {
            when {
                allOf {
                    environment name: 'TARGET_DB_NAME', value: 'qa-art-relations'
                    not { triggeredBy 'GitHubPushCause' }
                }
            }
            steps {
                build job: 'content-review-queue-reindex', parameters: [
                    string(name: 'Environment', value: 'qa'),
                    booleanParam(name: 'DB_SYNC', value: true),
                    booleanParam(name: 'CLEAR_INDEX', value: false)
                ], wait: false
            }
        }
        stage('Redeploy Services') {
            when {
                expression { SERVICES[params.SERVICE_NAME]?.redeployServices }
            }
            steps {
                script {
                    def targetEnv = params.TARGET_ENVIRONMENT
                    parallel(
                        SERVICES[params.SERVICE_NAME]?.redeployServices
                            .findAll { service, cfg -> cfg.containsKey(targetEnv) }
                            .collectEntries { service, cfg ->
                                def redeployCfg = cfg[targetEnv]
                                [
                                    (service): {
                                        echo "Redeploying service ${service}"
                                        build job: 'fargate-redeploy', parameters: [
                                            string(name: 'ENV', value: targetEnv),
                                            string(name: 'SERVICE_NAME', value: service),
                                            string(name: 'AWS_ACCOUNT', value: redeployCfg.account),
                                        ], wait: false
                                    }
                                ]
                            })
                }
            }
        }
    }

    post {
        success {
            script {
                def slackConfig = SERVICES[params.SERVICE_NAME]?.slack?.get(params.TARGET_ENVIRONMENT)
                if (slackConfig?.notify) {
                    def channel = slackConfig.channel ?: '#devops'
                    slackSend channel: channel, color: 'good', message: """${env.JOB_NAME} Finished (<${env.BUILD_URL}|Open>)
*Service:* ${env.SERVICE_NAME}
*Source Environment:* ${env.SOURCE_ENVIRONMENT}
*Target Environment:* ${env.TARGET_ENVIRONMENT}"""
                }
            }
        }
        failure {
            script {
                def slackConfig = SERVICES[params.SERVICE_NAME]?.slack?.get(params.TARGET_ENVIRONMENT)
                def channel = slackConfig?.channel ?: '#devops'
                slackSend channel: channel, color: 'danger', message: """${env.JOB_NAME} Failed (<${env.BUILD_URL}|Open>)
*Service:* ${env.SERVICE_NAME}
*Source Environment:* ${env.SOURCE_ENVIRONMENT}
*Target Environment:* ${env.TARGET_ENVIRONMENT}"""
            }
        }
        cleanup {
            cleanWs()
        }
    }
}

def getStateMachineInput() {
    // Pull out the composite key and the plain env part
    def targetEnv = params.TARGET_ENVIRONMENT

    def sourceEnvironment = SOURCE_ENVIRONMENTS[params.SOURCE_ENVIRONMENT]
    def targetEnvironment = TARGET_ENVIRONMENTS[targetEnv]
    def service           = SERVICES[params.SERVICE_NAME]
    def accountId         = service?.targetEnvironments?.get(targetEnv)?.accountId ?: '437795906767' // Default to QA account ID if not specified

    def stateMachineInput = [
        sanitise_data_function_name: targetEnvironment['dataSanitisationFunction'],
        source_account_id          : sourceEnvironment['accountId'],
        source_db_name             : SOURCE_DB_NAME,
        target_account_id          : accountId,
        target_account_role        : targetEnvironment['restoreRole'],
        target_db_name             : TARGET_DB_NAME,
        force_snapshot             : params.FORCE_SNAPSHOT,
        restore_snapshot_id        : params.RESTORE_SNAPSHOT_ID?.trim() ?: ''
    ]

    def debeziumConnector = service?.debeziumConnectors?.get(targetEnv)
    def debeziumConnectorAccountId = debeziumConnector?.get('accountId') ?: accountId
    if (debeziumConnector) {
        def kafka_topics = ['configs', 'offsets', 'status']
            .collect{ "_kafka_connect_debezium_mysql_source_${debeziumConnector['kafkaConnectGroupId']}-${it}" }

        stateMachineInput << [
            kafka_connector         : debeziumConnector['name'],
            kafka_bootstrap_servers : targetEnvironment['kafkaBootstrapServers'],
            kafka_topics            : kafka_topics,
            connector_name          : debeziumConnector['connectorName'],
            connector_account_id    : debeziumConnectorAccountId
        ]
    }

    def snowflakeRefresh = service?.snowflakeRefresh?.get(targetEnv)
    def fivetranSync = service?.fivetranSync?.get(targetEnv)

    if (snowflakeRefresh && params.SNOWFLAKE_REFRESH) {
        stateMachineInput << [
            requires_snowflake_refresh: snowflakeRefresh['enabled']
        ]
    }

    if (fivetranSync) {
        stateMachineInput << [
            requires_fivetran_sync: fivetranSync['enabled'],
            fivetran_connector_id: fivetranSync['connector_id'],
            fivetran_historical_sync: fivetranSync.get('historical_sync', 'false')
        ]
    }

    return JsonOutput.toJson(stateMachineInput)
}
