#------------------------------------------------------------
# Copyright (C) 2006 RoyaltyShare, Inc.   All Rights Reserved
#------------------------------------------------------------
package Distribution::Job::Matcher;
use strict;
use warnings;

use lib '/app/tools/common/lib';
use Common::RSApp;
use Common::Util;
use Common::Assert;
use Common::WriteXML;
use Common::RSApp;
use Common::Client;

use lib '/app/tools/rps/lib';
use RPS::Product::Product;
use RPS::Catalog::Album;
use RPS::Catalog::DistributionCache;
use RPS::Distribution::ParameterFactory;
use RPS::Product::DistributionService;

use lib '/app/tools/distribution/lib';
use Distribution::DB::Item::DistributionService;
use Distribution::DB::Item::ProductDistribution;
use Distribution::DB::Item::DistributionProcess;
use Distribution::DistributionProcess;

use constant kDistributionBatchSize     => 25;
use constant kDistributionMaxPacketSize => 700000;
use constant kDistributionMinSizeCheck  => 18;
use constant kInitialJobPriority        => 2000;
use constant kJobPriorityStep           => 5;

use base 'RPS::CronScript';

sub _init {
    my $self = shift;
    my %args = @_;

    $self->{_albumID}    = $args{album_id};
    $self->{_force}      = $args{force};
    $self->{_bundleSize} = $args{bundle_size};

    return $self->SUPER::_init(@_);
}

sub run_process {
    my $self = shift;
    my @list = $self->_get_client_list();

    $self->debug("\nStarting job processing");

    $self->debug("  Client list: @list");

    foreach (@list) {

        # Skip RSDEMO RSTEST AND CM unless a specific album is specified
        if ( ( $_ == 112 || $_ == 60 || $_ == 259 || $_ == 202 || $_ == 217 ) && !( $self->{_albumID} || $self->{_force} ) ) {
            $self->debug("Skipping client id: $_");
        } else {
            $self->_process_client($_);
        }
    }
}

sub _process_client {
    my $self      = shift;
    my $client_id = shift;

    assert($self);
    assert($client_id);

    $self->resetPriority();

    $self->log("\n\n**** Starting process for client id: $client_id\n");

    my $singleton = new Common::RSApp( clientID => $client_id );

    if ( !$self->_isDistributionClient() ) {
        $self->log(" ++ not a distribution client, ignoring");
        return;
    }

    die("Invalid client id '$client_id'")
      unless ( $singleton && Common::RSDB::ClientIDToDBName($client_id) );

    eval { $self->_process($client_id); };
    if ($@) {
        print STDERR "Failed to process client!!! Client ID: $client_id\n";
        print STDERR "$@\n";
    }

    $singleton = undef;
}

sub _process() {
    my $self      = shift;
    my $client_id = shift;
    my $found;

    assert($client_id);

    if ( $self->{_albumID} ) {
        my $album = new RPS::Catalog::Album( albumID => $self->{_albumID} );

        unless ( $album && $album->AlbumID ) {
            $self->print("Album not found ID: $self->{_albumID}");
            return;
        }

        unless ( $album->validateForDistribution() ) {
            $self->print("Album ID: $self->{_albumID} - Not distribution ready");
            return;
        }
    }

    my @services = $self->_getClientServices($client_id);

    unless (@services) {
        $self->log(" - No services found for this client") unless (@services);
        return;
    }

    unless ( $self->_hasDistributionAlbums() ) {
        $self->log(" - No albums found to deliver for this client");
        return;
    }

    # Ensure the validation cache is up to date.
    $self->_refreshCache();

    foreach my $service (@services) {
        $self->_processService( $client_id, $service );
    }
}

sub _getClientServices {
    my $self      = shift;
    my $client_id = Common::Client::Current()->ClientID;

    my $services = Distribution::DB::Item::DistributionService->GetAllClientServices( $client_id, status => 'active' );
    my @services;

    while ( $services->hasNext ) {
        my $service = $services->next();
        my $serviceObj = new RPS::Product::DistributionService( _dbItem => $service );

        push( @services, $service )
          if ( $serviceObj->Configured && $serviceObj->{ClientAllowed} );
    }

    return @services;
}

sub _hasDistributionAlbums {
    my $self = shift;

    my @services = $self->_getClientServices();
    my @service_ids;

    foreach my $service (@services) {
        push( @service_ids, $service->service_id() ) if ( $service->{configured} );
    }

    my %albums = ( album_ids => [ $self->{_albumID} ] ) if ( $self->{_albumID} );
    my $p_dists = Distribution::DB::Item::ProductDistribution->GetAllToDistribute( %albums, service_ids => \@service_ids );

    $p_dists->hasNext();
}

sub _refreshCache {
    my $self = shift;
    $self->_validateDistributionReady();
    $self->_updateDirtyDistributionStatus();
}

sub _updateDirtyDistributionStatus() {
    my $self  = shift;
    my $cache = new RPS::Catalog::DistributionCache;
    my ( $dist, $album );
    my @services = $self->_getClientServices();
    my @service_ids;

    foreach my $service (@services) {
        push( @service_ids, $service->service_id() ) if ( $service->{configured} );
    }

    # Itterate through all "dirty" albums status and revalidate if we are trying
    # to deliver that album.
    $self->log(" - Validate albums with updated metadata");

    my $p_dists = RPS::DB::Item::AlbumDistributionReady->GetAllInvalidCache();

    my $size    = $p_dists->totalSize();
    my $current = 1;

    while ( $p_dists->hasNext() ) {
        $dist = $p_dists->next;

        if (
            Distribution::DB::Item::ProductDistribution->GetStagedDistributionCount(
                album_id    => $dist->{album_id},
                service_ids => \@service_ids
            )
          ) {
            $self->log( sprintf( "    + Validating album: %-7d ( %d of %d )", $dist->{album_id}, $current++, $size ) );
            $album = new RPS::Catalog::Album( albumID => $dist->{album_id} );
            $cache->refreshAlbumStatus( album => $album );
        } else {
            $self->log( sprintf( "    + Validating album: %-7d ( %d of %d ) ... %s", $dist->{album_id}, $current++, $size, "Skipped" ) );
        }
    }
}

sub _validateDistributionReady {
    my $self = shift;

    my $cache = new RPS::Catalog::DistributionCache;
    my ( $dist, $album );
    my @services = $self->_getClientServices();
    my @service_ids;

    foreach my $service (@services) {
        push( @service_ids, $service->service_id() ) if ( $service->{configured} );
    }

    # First we want to iterate through all albums we think are valid and force
    # a validation to ensure that the cache isn't spoiled. (We don't trust the
    # record unless we validated it ourself).
    my %albums = ( album_ids => [ $self->{_albumID} ] ) if ( $self->{_albumID} );
    my $p_dists = Distribution::DB::Item::ProductDistribution->GetAllToDistribute( %albums, service_ids => \@service_ids );

    my $size    = $p_dists->totalSize();
    my $current = 1;

    $self->log(" - Validate distribution ready albums");

    while ( $p_dists->hasNext() ) {
        $dist = $p_dists->next;

        if (
            Distribution::DB::Item::ProductDistribution->GetStagedDistributionCount(
                album_id    => $dist->{album_id},
                service_ids => \@service_ids
            )
          ) {
            $self->log( sprintf( "    + Validating album: %-7d ( %d of %d )", $dist->{album_id}, $current++, $size ) );
            $album = new RPS::Catalog::Album( albumID => $dist->{album_id} );
            $cache->refreshAlbumStatus( album => $album );
        } else {
            $self->log( sprintf( "    + Validating album: %-7d ( %d of %d ) ... %s", $dist->{album_id}, $current++, $size, "Skipped" ) );
            die;
        }
    }

}

sub _processService {
    my $self      = shift;
    my $client_id = shift;
    my $service   = shift;
    my $dist;
    my $product;
    my $album;
    my $size;

    my $batchSize = $self->{_bundleSize} || kDistributionBatchSize;

    my $dist_process;
    my $count = 0;

    assert( $client_id && $service );

    $self->log( " - Processing service " . $service->service_name );

    unless ( $service->{configured} ) {
        $self->log("   + service not configured");
        return;
    }

    my %albums = ( album_ids => [ $self->{_albumID} ] ) if ( $self->{_albumID} );
    my $p_dists = Distribution::DB::Item::ProductDistribution->GetAllToDistribute( service_id => $service->service_id, %albums );

    $self->log("   + No products found for this service") unless ( $p_dists->size() );

    while ( $p_dists->hasNext() ) {
        $dist = $p_dists->next();

        $product = new RPS::Product::Product( productID => $dist->product_id );
        next if ( $self->{_albumID} && $self->{_albumID} != $product->AssetID() );
        $album = new RPS::Catalog::Album( albumID => $product->AssetID );

        $self->log( "    + processing album: " . $product->AssetID . " product: " . $product->ProductCode );

        # Rerun validation.  This is to ensure the status wasn't dirtied between
        # start time and now.
        unless ( $album->validateForDistribution() ) {
            $self->log( "    + Failed album validation for " . $product->ProductCode );
            next;
        }

        #
        # The album is validate so lets add it to a job queue.
        #

        # Only get the size if we have reached our minimum size check
        if ( $count < kDistributionMinSizeCheck ) {
            $size = 0;
        } else {
            my $paramObj = RPS::Distribution::ParameterFactory::CreateDistributionParameters( $dist_process->distribution_process_id );
            my $xml = Common::WriteXML::GetXMLString( $paramObj, undef, skip_array_indexing => 1 );
            $size = length($xml);
            $self->log( "    + size check: " . length($xml) );
        }

        # First see if we need to create a dist job, either we hit our max albums
        # or max xml size.
        if ( $count >= $batchSize || $size > kDistributionMaxPacketSize ) {
            $self->finializeDistributionProcess( $dist_process, $service->service_id );
            $dist_process = undef;
            $count        = 0;
        }

        unless ($dist_process) {
            $dist_process = Distribution::DB::Item::DistributionProcess->Create( created => 1 );
            $dist_process->save();
            assert($dist_process);
        }

        $dist->distribution_process_id( $dist_process->distribution_process_id );
        $dist->save();
        $dist->setQueuedDate( $dist->product_distribution_id );

        $count++;
    }

    $self->finializeDistributionProcess( $dist_process, $service->service_id ) if ($dist_process);
}

sub finializeDistributionProcess {
    my $self         = shift;
    my $dist_process = shift;
    my $service_id   = shift;

    assert($self);
    assert($dist_process);
    assert($service_id);

    my $priority = $self->getNextPriority();

    my $dist_process_id = $dist_process->distribution_process_id;

    $self->log("    + Finalizing Distribution Process.  ID: $dist_process_id");

    assert( !$dist_process->xml_data, "This process already contains data.  I'm not going to overwrite" );

    # Generate the XML
    my $paramObj = RPS::Distribution::ParameterFactory::CreateDistributionParameters($dist_process_id);
    my $xmlString = Common::WriteXML::GetXMLString( $paramObj, undef, skip_array_indexing => 1 );

    assert($xmlString);

    $dist_process->xml_data($xmlString);
    $dist_process->save();

    my $process = Distribution::DistributionProcess->new( distributionProcessID => $dist_process_id );
    $process->enqueuePackager( priority => $priority );
}

sub getNextPriority {
    my ($self) = @_;

    if ( !$self->{_priority} || $self->{_priority} <= 0 ) {
        $self->{_priority} = kInitialJobPriority;
    }

    my $returnPriority = $self->{_priority};

    # Lower the job priority for the next job
    $self->{_priority} -= kJobPriorityStep;

    return $returnPriority;
}

sub resetPriority {
    my ($self) = @_;
    $self->{_priority} = undef;
}

sub _isDistributionClient {
    my $self = shift;

    return Common::Client::Current()->isClientType(Common::Client::kDistribution);
}

1;
