diff --git a/Build.PL b/Build.PL index 3da0f2c3..238ae4af 100644 --- a/Build.PL +++ b/Build.PL @@ -43,6 +43,7 @@ my $builder = Module::Build->new( #'Gearman::Task' => 0, 'Digest::MD5' => 0, 'Data::UUID' => 0, + 'UUID' => 0, 'Net::Address::IP::Local' => 0, 'Date::Manip' => 0, 'Date::Calc' => 0, diff --git a/debian/control b/debian/control index 9013771f..9cca15e3 100644 --- a/debian/control +++ b/debian/control @@ -38,6 +38,7 @@ Depends: libgearman-client-perl, libdigest-md5-perl, libdata-uuid-perl, + libuuid-perl, libnet-address-ip-local-perl, libdate-manip-perl, libdate-calc-perl, diff --git a/lib/NGCP/BulkProcessor/Dao/Trunk/billing/contacts.pm b/lib/NGCP/BulkProcessor/Dao/Trunk/billing/contacts.pm new file mode 100644 index 00000000..f2cd3a05 --- /dev/null +++ b/lib/NGCP/BulkProcessor/Dao/Trunk/billing/contacts.pm @@ -0,0 +1,132 @@ +package NGCP::BulkProcessor::Dao::Trunk::billing::contacts; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_billing_db + +); + +use NGCP::BulkProcessor::SqlProcessor qw( + checktableinfo + insert_record + copy_row +); +use NGCP::BulkProcessor::SqlRecord qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord); +our @EXPORT_OK = qw( + gettablename + check_table + insert_row +); + +#my $logger = getlogger(__PACKAGE__); + +my $tablename = 'contacts'; +my $get_db = \&get_billing_db; + +my $expected_fieldnames = [ + 'id', + 'reseller_id', + 'gender', + 'firstname', + 'lastname', + 'comregnum', + 'company', + 'street', + 'postcode', + 'city', + 'country', + 'phonenumber', + 'mobilenumber', + 'email', + 'newsletter', + 'modify_timestamp', + 'create_timestamp', + 'faxnumber', + 'iban', + 'bic', + 'vatnum', + 'bankname', + 'gpp0', + 'gpp1', + 'gpp2', + 'gpp3', + 'gpp4', + 'gpp5', + 'gpp6', + 'gpp7', + 'gpp8', + 'gpp9', +]; + +my $indexes = {}; +# 'balance_interval' => [ 'contract_id','start','end' ], +# 'invoice_idx' => [ 'invoice_id' ], +#}; + +my $insert_unique_fields = []; #[ 'contract_id','start','end' ]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::SqlRecord->new($get_db, + $tablename, + $expected_fieldnames,$indexes); + + bless($self,$class); + + copy_row($self,shift,$expected_fieldnames); + + return $self; + +} + +sub insert_row { + + my ($data,$insert_ignore) = @_; + check_table(); + #return insert_record($get_db,$tablename,$data,$insert_ignore,$unique_fields) = @_; + +} + +sub buildrecords_fromrows { + + my ($rows,$load_recursive) = @_; + + my @records = (); + my $record; + + if (defined $rows and ref $rows eq 'ARRAY') { + foreach my $row (@$rows) { + $record = __PACKAGE__->new($row); + + # transformations go here ... + + push @records,$record; + } + } + + return \@records; + +} + +sub gettablename { + + return $tablename; + +} + +sub check_table { + + return checktableinfo($get_db, + $tablename, + $expected_fieldnames, + $indexes); + +} + +1; diff --git a/lib/NGCP/BulkProcessor/FakeTime.pm b/lib/NGCP/BulkProcessor/FakeTime.pm index 1bfc15d4..420b719e 100644 --- a/lib/NGCP/BulkProcessor/FakeTime.pm +++ b/lib/NGCP/BulkProcessor/FakeTime.pm @@ -17,25 +17,24 @@ use NGCP::BulkProcessor::Logging qw( use NGCP::BulkProcessor::LogError qw( faketimeerror - restwarn); +); require Exporter; our @ISA = qw(Exporter); our @EXPORT_OK = qw( - set_time - get_now - current_unix + set_fake_time + get_fake_now + get_fake_now_string + fake_current_unix infinite_future is_infinite_future datetime_to_string datetime_from_string - ); - -#my $logger = getlogger(__PACKAGE__); +); my $is_fake_time = 0; -sub set_time { +sub set_fake_time { my ($o) = @_; if (defined $o) { _set_fake_time($o); @@ -48,15 +47,15 @@ sub set_time { } } -sub _get_fake_clienttime_now { +sub get_fake_now_string { return datetime_to_string(_current_local()); } -sub get_now { +sub get_fake_now { return _current_local(); } -sub current_unix { +sub fake_current_unix { if ($is_fake_time) { return Time::Warp::time; } else { diff --git a/lib/NGCP/BulkProcessor/FileProcessor.pm b/lib/NGCP/BulkProcessor/FileProcessor.pm index 684cd76a..8e039932 100644 --- a/lib/NGCP/BulkProcessor/FileProcessor.pm +++ b/lib/NGCP/BulkProcessor/FileProcessor.pm @@ -176,7 +176,7 @@ sub process { if (defined $init_reader_context_code) { &$init_reader_context_code($self,$context); } - if ('CODE' eq ref $init_process_context_code) { + if (defined $init_process_context_code and 'CODE' eq ref $init_process_context_code) { &$init_process_context_code($context); } my $extractlines_code = (ref $self)->can('extractlines'); @@ -260,7 +260,7 @@ sub process { } eval { - if ('CODE' eq ref $uninit_process_context_code) { + if (defined $uninit_process_context_code and 'CODE' eq ref $uninit_process_context_code) { &$uninit_process_context_code($context); } }; @@ -422,7 +422,7 @@ sub _process { my $blockcount = 0; eval { - if ('CODE' eq ref $context->{init_process_context_code}) { + if (defined $context->{init_process_context_code} and 'CODE' eq ref $context->{init_process_context_code}) { &{$context->{init_process_context_code}}($context); } while (not _get_stop_consumer_thread($context,$tid)) { @@ -454,7 +454,7 @@ sub _process { my $err = $@; filethreadingdebug($err ? '[' . $tid . '] processor thread error: ' . $err : '[' . $tid . '] processor thread finished (' . $blockcount . ' blocks)',getlogger(__PACKAGE__)); eval { - if ('CODE' eq ref $context->{uninit_process_context_code}) { + if (defined $context->{uninit_process_context_code} and 'CODE' eq ref $context->{uninit_process_context_code}) { &{$context->{uninit_process_context_code}}($context); } }; diff --git a/lib/NGCP/BulkProcessor/LoadConfig.pm b/lib/NGCP/BulkProcessor/LoadConfig.pm index 0b4f1b77..ed3a0ccf 100644 --- a/lib/NGCP/BulkProcessor/LoadConfig.pm +++ b/lib/NGCP/BulkProcessor/LoadConfig.pm @@ -102,12 +102,12 @@ sub load_config { ); my ($result,$loadconfig_args,$postprocesscode) = update_masterconfig(%context); _splashinfo($configfile); - if ('ARRAY' eq ref $loadconfig_args) { + if (defined $loadconfig_args and 'ARRAY' eq ref $loadconfig_args) { foreach my $loadconfig_arg (@$loadconfig_args) { $result &= load_config(@$loadconfig_arg); } } - if ('CODE' eq ref $postprocesscode) { + if (defined $postprocesscode and 'CODE' eq ref $postprocesscode) { $result &= &$postprocesscode(%context); } return $result; diff --git a/lib/NGCP/BulkProcessor/LogError.pm b/lib/NGCP/BulkProcessor/LogError.pm index a7f1c4a5..d680b8d1 100644 --- a/lib/NGCP/BulkProcessor/LogError.pm +++ b/lib/NGCP/BulkProcessor/LogError.pm @@ -43,6 +43,7 @@ require Exporter; our @ISA = qw(Exporter); our @EXPORT_OK = qw( notimplementederror + faketimeerror dberror dbwarn fieldnamesdiffer @@ -269,6 +270,17 @@ sub notimplementederror { } +sub faketimeerror { + + my ($message, $logger) = @_; + if (defined $logger) { + $logger->error($message); + } + + terminate($message, $logger); + +} + sub dberror { my ($db, $message, $logger) = @_; diff --git a/lib/NGCP/BulkProcessor/Logging.pm b/lib/NGCP/BulkProcessor/Logging.pm index d1087539..fa8f6adc 100644 --- a/lib/NGCP/BulkProcessor/Logging.pm +++ b/lib/NGCP/BulkProcessor/Logging.pm @@ -82,6 +82,9 @@ our @EXPORT_OK = qw( processing_info + faketimeinfo + faketimedebug + restthreadingdebug restprocessingstarted restprocessingdone @@ -706,6 +709,23 @@ sub configurationinfo { } +sub faketimeinfo { + + my ($message, $logger) = @_; + if (defined $logger) { + $logger->info($message); + } + +} + +sub faketimedebug { + + my ($message, $logger) = @_; + if (defined $logger) { + $logger->debug($message); + } + +} sub scriptinfo { diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Batch.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Batch.pm index 7bd33d62..a82e9d72 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Batch.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Batch.pm @@ -224,7 +224,6 @@ sub getinsertstatement { sub getupsertstatement { - my ($exists_delta,$new_delta) = @_; check_table(); my $db = &$get_db(); my $table = $db->tableidentifier($tablename); @@ -233,9 +232,9 @@ sub getupsertstatement { my @values = (); foreach my $fieldname (@$expected_fieldnames) { if ('delta' eq $fieldname) { - my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' . + my $stmt = 'SELECT \'' . $updated_delta . '\' FROM ' . $table . ' WHERE ' . $db->columnidentifier('number') . ' = ?'; - push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')'); + push(@values,'COALESCE((' . $stmt . '), \'' . $added_delta . '\')'); } else { push(@values,'?'); } diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOption.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOption.pm index 5f17d5b5..07c95fe5 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOption.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOption.pm @@ -236,7 +236,6 @@ sub getinsertstatement { sub getupsertstatement { - my ($exists_delta,$new_delta) = @_; check_table(); my $db = &$get_db(); my $table = $db->tableidentifier($tablename); @@ -245,10 +244,10 @@ sub getupsertstatement { my @values = (); foreach my $fieldname (@$expected_fieldnames) { if ('delta' eq $fieldname) { - my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' . + my $stmt = 'SELECT \'' . $updated_delta . '\' FROM ' . $table . ' WHERE ' . $db->columnidentifier('subscribernumber') . ' = ?' . ' AND ' . $db->columnidentifier('option') . ' = ?'; - push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')'); + push(@values,'COALESCE((' . $stmt . '), \'' . $added_delta . '\')'); } else { push(@values,'?'); } diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOptionSetItem.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOptionSetItem.pm index 11dc4643..30215797 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOptionSetItem.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/FeatureOptionSetItem.pm @@ -243,7 +243,6 @@ sub getinsertstatement { sub getupsertstatement { - my ($exists_delta,$new_delta) = @_; check_table(); my $db = &$get_db(); my $table = $db->tableidentifier($tablename); @@ -252,11 +251,11 @@ sub getupsertstatement { my @values = (); foreach my $fieldname (@$expected_fieldnames) { if ('delta' eq $fieldname) { - my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' . + my $stmt = 'SELECT \'' . $updated_delta . '\' FROM ' . $table . ' WHERE ' . $db->columnidentifier('subscribernumber') . ' = ?' . ' AND ' . $db->columnidentifier('option') . ' = ?' . ' AND ' . $db->columnidentifier('optionsetitem') . ' = ?'; - push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')'); + push(@values,'COALESCE((' . $stmt . '), \'' . $added_delta . '\')'); } else { push(@values,'?'); } diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Lnp.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Lnp.pm index 136c91d6..9b5b2b28 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Lnp.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Lnp.pm @@ -242,7 +242,6 @@ sub getinsertstatement { sub getupsertstatement { - my ($exists_delta,$new_delta) = @_; check_table(); my $db = &$get_db(); my $table = $db->tableidentifier($tablename); @@ -251,10 +250,10 @@ sub getupsertstatement { my @values = (); foreach my $fieldname (@$expected_fieldnames) { if ('delta' eq $fieldname) { - my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' . + my $stmt = 'SELECT \'' . $updated_delta . '\' FROM ' . $table . ' WHERE ' . $db->columnidentifier('lrn_code') . ' = ?' . ' AND ' . $db->columnidentifier('ported_number') . ' = ?'; - push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')'); + push(@values,'COALESCE((' . $stmt . '), \'' . $added_delta . '\')'); } else { push(@values,'?'); } diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Subscriber.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Subscriber.pm index 04cea222..b360365a 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Subscriber.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/Subscriber.pm @@ -287,7 +287,6 @@ sub getinsertstatement { sub getupsertstatement { - my ($exists_delta,$new_delta) = @_; check_table(); my $db = &$get_db(); my $table = $db->tableidentifier($tablename); @@ -296,11 +295,11 @@ sub getupsertstatement { my @values = (); foreach my $fieldname (@$expected_fieldnames) { if ('delta' eq $fieldname) { - my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' . + my $stmt = 'SELECT \'' . $updated_delta . '\' FROM ' . $table . ' WHERE ' . $db->columnidentifier('country_code') . ' = ?' . ' AND ' . $db->columnidentifier('area_code') . ' = ?' . ' AND ' . $db->columnidentifier('dial_number') . ' = ?'; - push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')'); + push(@values,'COALESCE((' . $stmt . '), \'' . $added_delta . '\')'); } else { push(@values,'?'); } diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/UsernamePassword.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/UsernamePassword.pm index b656f1d2..e8158790 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/UsernamePassword.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Dao/import/UsernamePassword.pm @@ -36,6 +36,7 @@ our @EXPORT_OK = qw( countby_delta $deleted_delta + $unchanged_delta $updated_delta $added_delta @@ -50,16 +51,20 @@ my $expected_fieldnames = [ 'fqdn', 'username', 'password', - 'auth_level', + #'auth_level', 'delta', ]; # table creation: my $primarykey_fieldnames = [ 'fqdn' ]; -my $indexes = { $tablename . '_delta' => [ 'delta(7)' ]}; +my $indexes = { + $tablename . '_username_password' => [ 'username(11)', 'password(32)' ], + $tablename . '_delta' => [ 'delta(7)' ], +}; #my $fixtable_statements = []; our $deleted_delta = 'DELETED'; +our $unchanged_delta = 'UNCHANGED'; our $updated_delta = 'UPDATED'; our $added_delta = 'ADDED'; @@ -219,7 +224,6 @@ sub getinsertstatement { sub getupsertstatement { - my ($exists_delta,$new_delta) = @_; check_table(); my $db = &$get_db(); my $table = $db->tableidentifier($tablename); @@ -228,9 +232,11 @@ sub getupsertstatement { my @values = (); foreach my $fieldname (@$expected_fieldnames) { if ('delta' eq $fieldname) { - my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' . + my $unchanged_stmt = 'SELECT \'' . $unchanged_delta . '\' FROM ' . $table . ' WHERE ' . + $db->columnidentifier('fqdn') . ' = ? AND ' . $db->columnidentifier('password') . ' = ?'; + my $updated_stmt = 'SELECT \'' . $updated_delta . '\' FROM ' . $table . ' WHERE ' . $db->columnidentifier('fqdn') . ' = ?'; - push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')'); + push(@values,'COALESCE((' . $unchanged_stmt . '), (' . $updated_stmt . '), \'' . $added_delta . '\')'); } else { push(@values,'?'); } diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/FileProcessors/UserPasswordFile.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/FileProcessors/UserPasswordFile.pm index 1739a7d5..8843d3f5 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/FileProcessors/UserPasswordFile.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/FileProcessors/UserPasswordFile.pm @@ -17,8 +17,8 @@ require Exporter; our @ISA = qw(Exporter NGCP::BulkProcessor::FileProcessor); our @EXPORT_OK = qw(); -my $lineseparator = '\\r\\n'; -my $fieldseparator = " +"; +my $lineseparator = '\\n'; #'\\r\\n'; +my $fieldseparator = ','; #" +"; my $encoding = 'UTF-8'; my $buffersize = 100 * 1024; diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Import.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Import.pm index b7349e68..b3a16c21 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Import.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Import.pm @@ -18,9 +18,11 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::Settings qw( $ignore_lnp_unique $user_password_import_numofthreads $ignore_user_password_unique + $username_prefix $batch_import_numofthreads $ignore_batch_unique - $dry + $subscribernumber_pattern + $skip_errors ); use NGCP::BulkProcessor::Logging qw ( getlogger @@ -75,7 +77,7 @@ sub import_features_define { # prepare parse: my $importer = NGCP::BulkProcessor::Projects::Migration::IPGallery::FileProcessors::FeaturesDefineFile->new($features_define_import_numofthreads); - $importer->stoponparseerrors(!$dry); + $importer->stoponparseerrors(!$skip_errors); my $upsert = _import_features_define_reset_delta(); @@ -105,8 +107,9 @@ sub import_features_define { } next unless defined $row; foreach my $subscriber_number (keys %$row) { + next unless _check_subscribernumber($context,$subscriber_number,$rownum); foreach my $option (@{$row->{$subscriber_number}}) { - if ('HASH' eq ref $option) { + if (defined $option and 'HASH' eq ref $option) { foreach my $setoption (keys %$option) { foreach my $setoptionitem (@{$skip_duplicate_setoptionitems ? removeduplicates($option->{$setoption}) : $option->{$setoption}}) { if ($context->{upsert}) { @@ -139,14 +142,14 @@ sub import_features_define { } if ((scalar @featureoption_rows) > 0) { - if ($dry) { + if ($skip_errors) { eval { _insert_featureoption_rows($context,\@featureoption_rows); }; } else { _insert_featureoption_rows($context,\@featureoption_rows); } } if ((scalar @featureoptionsetitem_rows) > 0) { - if ($dry) { + if ($skip_errors) { eval { _insert_featureoptionsetitem_rows($context,\@featureoptionsetitem_rows); }; } else { _insert_featureoptionsetitem_rows($context,\@featureoptionsetitem_rows); @@ -200,10 +203,8 @@ sub _insert_featureoption_rows { my ($context,$featureoption_rows) = @_; $context->{db}->db_do_begin( ($context->{upsert} ? - NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::getupsertstatement( - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::updated_delta, - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::added_delta - ) : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::getinsertstatement($ignore_options_unique)), + NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::getupsertstatement() + : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::getinsertstatement($ignore_options_unique)), #NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::gettablename(), #lock - $import_multithreading ); @@ -215,10 +216,8 @@ sub _insert_featureoptionsetitem_rows { my ($context,$featureoptionsetitem_rows) = @_; $context->{db}->db_do_begin( ($context->{upsert} ? - NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::getupsertstatement( - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::updated_delta, - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::added_delta - ) : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::getinsertstatement($ignore_setoptionitems_unique)), + NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::getupsertstatement() + : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::getinsertstatement($ignore_setoptionitems_unique)), #NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::gettablename(), #lock ); @@ -252,12 +251,14 @@ sub import_subscriber_define { if ($record->{dial_number} =~ $subscribernumer_exclude_pattern) { if ($record->{dial_number} =~ $subscribernumer_exclude_exception_pattern) { processing_info($context->{tid},'record ' . $rownum . ' - exclude exception pattern match: ' . $record->{dial_number},getlogger(__PACKAGE__)); + next unless _check_subscribernumber($context,$record->{dial_number},$rownum); next unless _import_subscriber_define_referential_checks($context,$record,$rownum); } else { processing_info($context->{tid},'record ' . $rownum . ' - skipped, exclude pattern match: ' . $record->{dial_number},getlogger(__PACKAGE__)); next; } } else { + next unless _check_subscribernumber($context,$record->{dial_number},$rownum); next unless _import_subscriber_define_referential_checks($context,$record,$rownum); } my @subscriber_row = @$row; @@ -270,7 +271,7 @@ sub import_subscriber_define { } if ((scalar @subscriber_rows) > 0) { - if ($dry) { + if ($skip_errors) { eval { _insert_subscriber_rows($context,\@subscriber_rows); }; } else { _insert_subscriber_rows($context,\@subscriber_rows); @@ -310,10 +311,10 @@ sub _import_subscriber_define_referential_checks { } } } else { - $result &= 0; - if ($dry) { + if ($skip_errors) { fileprocessingwarn($context->{filename},'record ' . $rownum . ' - no features records for subscriber found: ' . $record->{dial_number},getlogger(__PACKAGE__)); } else { + $result &= 0; fileprocessingerror($context->{filename},'record ' . $rownum . ' - no features records for subscriber found: ' . $record->{dial_number},getlogger(__PACKAGE__)); } } @@ -321,10 +322,10 @@ sub _import_subscriber_define_referential_checks { if (NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::countby_fqdn($record->subscribernumber()) > 0) { } else { - $result &= 0; - if ($dry) { + if ($skip_errors) { fileprocessingwarn($context->{filename},'record ' . $rownum . ' - no username password record for subscriber found: ' . $record->{dial_number},getlogger(__PACKAGE__)); } else { + $result &= 0; fileprocessingerror($context->{filename},'record ' . $rownum . ' - no username password record for subscriber found: ' . $record->{dial_number},getlogger(__PACKAGE__)); } } @@ -342,7 +343,7 @@ sub _import_subscriber_define_checks { }; if ($@ or $optioncount == 0) { fileprocessingerror($file,'please import subscriber features first',getlogger(__PACKAGE__)); - $result = 0; #even in dry mode.. + $result = 0; #even in skip-error mode.. } my $userpasswordcount = 0; eval { @@ -350,7 +351,7 @@ sub _import_subscriber_define_checks { }; if ($@ or $userpasswordcount == 0) { fileprocessingerror($file,'please import user passwords first',getlogger(__PACKAGE__)); - $result = 0; #even in dry mode.. + $result = 0; #even in skip-error mode.. } return $result; } @@ -371,10 +372,7 @@ sub _insert_subscriber_rows { my ($context,$subscriber_rows) = @_; $context->{db}->db_do_begin( ($context->{upsert} ? - NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::getupsertstatement( - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::updated_delta, - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::added_delta - ) + NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::getupsertstatement() : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::getinsertstatement($ignore_subscriber_unique)), #NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::gettablename(), #lock @@ -417,7 +415,7 @@ sub import_lnp_define { } if ((scalar @lnp_rows) > 0) { - if ($dry) { + if ($skip_errors) { eval { _insert_lnp_rows($context,\@lnp_rows); }; } else { _insert_lnp_rows($context,\@lnp_rows); @@ -457,10 +455,7 @@ sub _insert_lnp_rows { my ($context,$lnp_rows) = @_; $context->{db}->db_do_begin( ($context->{upsert} ? - NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::getupsertstatement( - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::updated_delta, - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::added_delta - ) + NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::getupsertstatement() : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::getinsertstatement($ignore_lnp_unique)), #NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::gettablename(), #lock @@ -495,9 +490,12 @@ sub import_user_password { foreach my $row (@$rows) { $rownum++; my @usernamepassword_row = @$row; + shift @usernamepassword_row; + unshift(@usernamepassword_row,($username_prefix // '') . $usernamepassword_row[0]); my $record = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword->new(\@usernamepassword_row); + next unless _check_subscribernumber($context,$record->{fqdn},$rownum); if ($context->{upsert}) { - push(@usernamepassword_row,$record->{fqdn}); + push(@usernamepassword_row,$record->{fqdn},$record->{password},$record->{fqdn}); } else { push(@usernamepassword_row,$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::added_delta); } @@ -505,7 +503,7 @@ sub import_user_password { } if ((scalar @usernamepassword_rows) > 0) { - if ($dry) { + if ($skip_errors) { eval { _insert_usernamepassword_rows($context,\@usernamepassword_rows); }; } else { _insert_usernamepassword_rows($context,\@usernamepassword_rows); @@ -545,10 +543,8 @@ sub _insert_usernamepassword_rows { my ($context,$usernamepassword_rows) = @_; $context->{db}->db_do_begin( ($context->{upsert} ? - NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::getupsertstatement( - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::updated_delta, - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::added_delta - ) : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::getinsertstatement($ignore_user_password_unique)), + NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::getupsertstatement() + : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::getinsertstatement($ignore_user_password_unique)), #NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::gettablename(), #lock - $import_multithreading ); @@ -578,6 +574,7 @@ sub import_batch { foreach my $row (@$rows) { $rownum++; my $record = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Batch->new($row); + next unless _check_subscribernumber($context,$record->{number},$rownum); next unless _import_batch_referential_checks($context,$record,$rownum); my @batch_row = @$row; if ($context->{upsert}) { @@ -589,7 +586,7 @@ sub import_batch { } if ((scalar @batch_rows) > 0) { - if ($dry) { + if ($skip_errors) { eval { _insert_batch_rows($context,\@batch_rows); }; } else { _insert_batch_rows($context,\@batch_rows); @@ -619,10 +616,10 @@ sub _import_batch_referential_checks { if (NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_subscribernumber($record->{number}) > 0) { } else { - $result &= 0; - if ($dry) { + if ($skip_errors) { fileprocessingwarn($context->{filename},'record ' . $rownum . ' - no subscriber record for batch number found: ' . $record->{number},getlogger(__PACKAGE__)); } else { + $result &= 0; fileprocessingerror($context->{filename},'record ' . $rownum . ' - no subscriber record for batch number found: ' . $record->{number},getlogger(__PACKAGE__)); } } @@ -640,7 +637,7 @@ sub _import_batch_checks { }; if ($@ or $subscribercount == 0) { fileprocessingerror($file,'please import subscribers first',getlogger(__PACKAGE__)); - $result = 0; #even in dry mode.. + $result = 0; #even in skip-error mode.. } return $result; } @@ -661,10 +658,7 @@ sub _insert_batch_rows { my ($context,$batch_rows) = @_; $context->{db}->db_do_begin( ($context->{upsert} ? - NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Batch::getupsertstatement( - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Batch::updated_delta, - $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Batch::added_delta - ) + NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Batch::getupsertstatement() : NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Batch::getinsertstatement($ignore_batch_unique)), #NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::gettablename(), #lock @@ -673,4 +667,20 @@ sub _insert_batch_rows { $context->{db}->db_finish(); } +sub _check_subscribernumber { + my ($context,$subscribernumber,$rownum) = @_; + my $result = 1; + if (defined $subscribernumber_pattern) { + if ($subscribernumber !~ $subscribernumber_pattern) { + if ($skip_errors) { + fileprocessingwarn($context->{filename},'record ' . $rownum . ' - invalid subscriber number found: ' . $subscribernumber,getlogger(__PACKAGE__)); + } else { + $result &= 0; + fileprocessingerror($context->{filename},'record ' . $rownum . ' - no features records for subscriber found: ' . $subscribernumber,getlogger(__PACKAGE__)); + } + } + } + return $result; +} + 1; diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Settings.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Settings.pm index 4b648a0f..aa2d621d 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Settings.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Settings.pm @@ -44,6 +44,7 @@ our @EXPORT_OK = qw( $import_multithreading $run_id $dry + $skip_errors $force $import_db_file @@ -67,11 +68,14 @@ our @EXPORT_OK = qw( $user_password_filename $user_password_import_numofthreads $ignore_user_password_unique + $username_prefix $batch_filename $batch_import_numofthreads $ignore_batch_unique + $subscribernumber_pattern + ); our $defaultconfig = 'config.cfg'; @@ -83,6 +87,7 @@ our $rollback_path = $working_path . 'rollback/'; our $force = 0; our $dry = 0; +our $skip_errors = 0; our $run_id = ''; our $import_db_file = _get_import_db_file($run_id,'import'); our $import_multithreading = $enablemultithreading; @@ -107,11 +112,14 @@ our $ignore_lnp_unique = 1; our $user_password_filename = undef; our $user_password_import_numofthreads = $cpucount; our $ignore_user_password_unique = 0; +our $username_prefix = undef; our $batch_filename = undef; our $batch_import_numofthreads = $cpucount; our $ignore_batch_unique = 0; +our $subscribernumber_pattern = undef; + sub update_settings { my ($data,$configfile) = @_; @@ -125,6 +133,7 @@ sub update_settings { $result &= _prepare_working_paths(1); $dry = $data->{dry} if exists $data->{dry}; + $skip_errors = $data->{skip_errors} if exists $data->{skip_errors}; $import_db_file = _get_import_db_file($run_id,'import'); $import_multithreading = $data->{import_multithreading} if exists $data->{import_multithreading}; @@ -141,12 +150,18 @@ sub update_settings { (my $regexp_result,$subscribernumer_exclude_exception_pattern) = parse_regexp($subscribernumer_exclude_exception_pattern,$configfile); $result &= $regexp_result; + $subscribernumber_pattern = $data->{subscribernumber_pattern} if exists $data->{subscribernumber_pattern}; + (my $regexp_result,$subscribernumber_pattern) = parse_regexp($subscribernumber_pattern,$configfile); + $result &= $regexp_result; + $lnp_define_filename = _get_import_filename($lnp_define_filename,$data,'lnp_define_filename'); $lnp_define_import_numofthreads = _get_import_numofthreads($cpucount,$data,'lnp_define_import_numofthreads'); $user_password_filename = _get_import_filename($user_password_filename,$data,'user_password_filename'); $user_password_import_numofthreads = _get_import_numofthreads($cpucount,$data,'user_password_import_numofthreads'); + $username_prefix = $data->{username_prefix} if exists $data->{username_prefix}; + $batch_filename = _get_import_filename($batch_filename,$data,'batch_filename'); $batch_import_numofthreads = _get_import_numofthreads($cpucount,$data,'batch_import_numofthreads'); diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Provisioning.pm b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/SubscriberProvisioning.pm similarity index 58% rename from lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Provisioning.pm rename to lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/SubscriberProvisioning.pm index 1ee46d20..97546081 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/Provisioning.pm +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/SubscriberProvisioning.pm @@ -1,4 +1,4 @@ -package NGCP::BulkProcessor::Projects::Migration::IPGallery::Provisioning; +package NGCP::BulkProcessor::Projects::Migration::IPGallery::SubscriberProvisioning; use strict; ## no critic @@ -21,30 +21,6 @@ sub test { my $result = 1; - return $result && NGCP::BulkProcessor::RestRequests::Trunk::BillingProfiles::process_items( - process_code => sub { - my ($context,$records,$row_offset) = @_; - my $rownum = $row_offset; - print "!!!!!!$row_offset!!!!!!!\n"; - return 1; - }, - init_process_context_code => sub { - my ($context)= @_; - }, - uninit_process_context_code => sub { - my ($context)= @_; - destroy_all_dbs(); - }, - load_recursive => 1, - multithreading => 1, - numofthreads => 4, - ); -} - -sub test1 { - - my $result = 1; - return $result && NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::process_records( process_code => sub { my ($context,$records,$row_offset) = @_; diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/process.pl b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/process.pl index 41c2b4d9..dc1860c3 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/process.pl +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/process.pl @@ -18,6 +18,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::Settings qw( $defaultsettings $defaultconfig $dry + $skip_errors $force $run_id $features_define_filename @@ -78,7 +79,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::Import qw( import_batch ); -use NGCP::BulkProcessor::Projects::Migration::IPGallery::Provisioning qw( +use NGCP::BulkProcessor::Projects::Migration::IPGallery::SubscriberProvisioning qw( test ); @@ -143,6 +144,7 @@ sub init { "task=s" => $tasks, "run=s" => \$run_id, "dry" => \$dry, + "skip-errors" => \$skip_errors, "force" => \$force, ); # or scripterror('error in command line arguments',getlogger(getscriptpath())); @@ -163,7 +165,8 @@ sub main() { my $result = 1; my $completion = 0; - if ('ARRAY' eq ref $tasks and (scalar @$tasks) > 0) { + if (defined $tasks and 'ARRAY' eq ref $tasks and (scalar @$tasks) > 0) { + scriptinfo('skip-errors: processing won\'t stop upon errors',getlogger(__PACKAGE__)) if $skip_errors; foreach my $task (@$tasks) { if (lc($check_task_opt) eq lc($task)) { @@ -489,10 +492,14 @@ sub import_user_password_task { $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::added_delta ); $stats .= "\n new: $added_count rows"; - my $existing_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::countby_delta( + my $unchanged_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::countby_delta( + $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::unchanged_delta + ); + $stats .= "\n unchanged: $unchanged_count rows"; + my $updated_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::countby_delta( $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::updated_delta ); - $stats .= "\n existing: $existing_count rows"; + $stats .= "\n updated: $updated_count rows"; my $deleted_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::countby_delta( $NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::UsernamePassword::deleted_delta ); diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/settings.cfg b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/settings.cfg index 6b77a5ce..853faa62 100644 --- a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/settings.cfg +++ b/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/settings.cfg @@ -1,7 +1,10 @@ #dry=0 +#skip_errors=0 import_multithreading = 0 +subscribernumber_pattern = ^356\d{8} + features_define_filename = /home/rkrenn/test/Features_Define.cfg features_define_import_numofthreads = 2 @@ -13,8 +16,9 @@ subscribernumer_exclude_exception_pattern = ^35627702770$ lnp_define_filename = /home/rkrenn/test/LNP_Define.cfg lnp_define_import_numofthreads = 2 -user_password_filename = /home/rkrenn/test/username_passwords.txt +user_password_filename = /home/rkrenn/test/FSISDN and password v2.csv user_password_import_numofthreads = 2 +username_prefix = 356 batch_filename = /home/rkrenn/test/blah.txt batch_import_numofthreads = 2 diff --git a/lib/NGCP/BulkProcessor/Projects/t/ApiTest.pm b/lib/NGCP/BulkProcessor/Projects/t/ApiTest.pm new file mode 100644 index 00000000..a4b6369a --- /dev/null +++ b/lib/NGCP/BulkProcessor/Projects/t/ApiTest.pm @@ -0,0 +1,38 @@ +package NGCP::BulkProcessor::Projects::t::ApiTest; +use strict; + +## no critic + +use NGCP::BulkProcessor::RestRequests::Trunk::BillingProfiles qw(); + +require Exporter; +our @ISA = qw(Exporter); +our @EXPORT_OK = qw( + test +); + +sub test { + + my $result = 1; + + return $result && NGCP::BulkProcessor::RestRequests::Trunk::BillingProfiles::process_items( + process_code => sub { + my ($context,$records,$row_offset) = @_; + my $rownum = $row_offset; + print "!!!!!!$row_offset!!!!!!!\n"; + return 1; + }, + init_process_context_code => sub { + my ($context)= @_; + }, + uninit_process_context_code => sub { + my ($context)= @_; + destroy_all_dbs(); + }, + load_recursive => 1, + multithreading => 1, + numofthreads => 4, + ); +} + +1; diff --git a/lib/NGCP/BulkProcessor/Projects/t/config.cfg b/lib/NGCP/BulkProcessor/Projects/t/config.cfg new file mode 100644 index 00000000..5c385acb --- /dev/null +++ b/lib/NGCP/BulkProcessor/Projects/t/config.cfg @@ -0,0 +1,41 @@ +##general settings: +#working_path = /var/sipwise/Migration/IPGallery +cpucount = 4 +enablemultithreading = 1 + +##gearman/service listener config: +jobservers = 127.0.0.1:4730 + +#provisioning_conf = /etc/ngcp-panel/provisioning.conf + +##NGCP MySQL connectivity - "accounting" db: +accounting_host = 127.0.0.1 +accounting_port = 3306 +accounting_databasename = accounting +accounting_username = root +accounting_password = + +##NGCP MySQL connectivity - "billing" db: +billing_host = 192.168.0.77 +billing_port = 3306 +billing_databasename = billing +billing_username = root +billing_password = + +##NGCP REST-API connectivity: +ngcprestapi_uri = https://127.0.0.1:1443 +ngcprestapi_username = administrator +ngcprestapi_password = administrator +ngcprestapi_realm = api_admin_http + +##sending email: +emailenable = 0 +erroremailrecipient = +warnemailrecipient = +completionemailrecipient = rkrenn@sipwise.com +doneemailrecipient = + +##logging: +fileloglevel = OFF +screenloglevel = INFO +emailloglevel = OFF diff --git a/lib/NGCP/BulkProcessor/Projects/t/test_api.pl b/lib/NGCP/BulkProcessor/Projects/t/test_api.pl new file mode 100644 index 00000000..55f46221 --- /dev/null +++ b/lib/NGCP/BulkProcessor/Projects/t/test_api.pl @@ -0,0 +1,112 @@ +#!/usr/bin/perl +use warnings; +use strict; + +use File::Basename; +use Cwd; +use lib Cwd::abs_path(File::Basename::dirname(__FILE__) . '/../../../../'); + +use NGCP::BulkProcessor::LoadConfig qw( + load_config +); + +use NGCP::BulkProcessor::RestRequests::Trunk::SystemContacts qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::Contracts qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::Resellers qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::BillingProfiles qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::BillingZones qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::BillingFees qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::Domains qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::CustomerContacts qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::Customers qw(); +use NGCP::BulkProcessor::RestRequests::Trunk::Subscribers qw(); + +load_config('config.cfg'); + +{ + my $t = time(); + my $systemcontact_id = NGCP::BulkProcessor::RestRequests::Trunk::SystemContacts::create_item({ + email => "test$t\@system.com", + }); + my $contract_id = NGCP::BulkProcessor::RestRequests::Trunk::Contracts::create_item({ + contact_id => $systemcontact_id, + billing_profile_id => 1, #default profile id + type => 'reseller', + status => 'active', + }); + my $reseller_id = NGCP::BulkProcessor::RestRequests::Trunk::Resellers::create_item({ + contract_id => $contract_id, + name => "reseller$t", + status => 'active', + }); + + my $billing_profile_id = NGCP::BulkProcessor::RestRequests::Trunk::BillingProfiles::create_item({ + handle => "profile$t", + name => "profile$t", + reseller_id => $reseller_id, + }); + + my $billing_zone_id = NGCP::BulkProcessor::RestRequests::Trunk::BillingZones::create_item({ + billing_profile_id => $billing_profile_id, + zone => "zone$t", + }); + + my $billing_fee_id = NGCP::BulkProcessor::RestRequests::Trunk::BillingFees::create_item({ + billing_profile_id => $billing_profile_id, + billing_zone_id => $billing_zone_id, + direction => 'out', + destination => '.*', + onpeak_init_rate => 10, + onpeak_init_interval => 1, + onpeak_follow_rate => 5, + onpeak_follow_interval => 2, + offpeak_init_rate => 8, + offpeak_init_interval => 1, + offpeak_follow_rate => 4, + offpeak_follow_interval => 2, + }); + + my $success = NGCP::BulkProcessor::RestRequests::Trunk::BillingFees::delete_item($billing_fee_id); + + my $contract = NGCP::BulkProcessor::RestRequests::Trunk::Contracts::set_item($contract_id,{ + contact_id => $systemcontact_id, + billing_profile_id => $billing_profile_id, + type => 'reseller', + status => 'active', + }); + + $contract = NGCP::BulkProcessor::RestRequests::Trunk::Contracts::update_item($contract_id,{ + contact_id => $systemcontact_id, + billing_profile_id => $billing_profile_id, + #type => 'reseller', + status => 'active', + }); + + my $domain_id = NGCP::BulkProcessor::RestRequests::Trunk::Domains::create_item({ + domain => "test$t.com", + reseller_id => $reseller_id, + }); + + my $customercontact_id = NGCP::BulkProcessor::RestRequests::Trunk::CustomerContacts::create_item({ + email => "test$t\@customer.com", + reseller_id => $reseller_id, + }); + + my $customer_id = NGCP::BulkProcessor::RestRequests::Trunk::Customers::create_item({ + contact_id => $customercontact_id, + billing_profile_id => $billing_profile_id, + status => "active", + type => "sipaccount", + }); + + my $subscriber_id = NGCP::BulkProcessor::RestRequests::Trunk::Subscribers::create_item({ + customer_id => $customer_id, + domain_id => $domain_id, + username => "subscriber$t", + password => "subscriber$t", + }); + + print "blah"; +} + +exit; diff --git a/lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/test_dsl.pl b/lib/NGCP/BulkProcessor/Projects/t/test_dsl.pl similarity index 100% rename from lib/NGCP/BulkProcessor/Projects/Migration/IPGallery/test_dsl.pl rename to lib/NGCP/BulkProcessor/Projects/t/test_dsl.pl diff --git a/lib/NGCP/BulkProcessor/RestConnector.pm b/lib/NGCP/BulkProcessor/RestConnector.pm index dca8cd80..ffaf4df3 100644 --- a/lib/NGCP/BulkProcessor/RestConnector.pm +++ b/lib/NGCP/BulkProcessor/RestConnector.pm @@ -23,7 +23,9 @@ use NGCP::BulkProcessor::Utils qw(threadid); require Exporter; our @ISA = qw(Exporter); -our @EXPORT_OK = qw(); +our @EXPORT_OK = qw( + _add_headers +); #my $logger = getlogger(__PACKAGE__); @@ -95,6 +97,7 @@ sub _create_ua { resterror($self,'base URL not set',getlogger(__PACKAGE__)); } my $ua = LWP::UserAgent->new(); + restdebug($self,"ua created",getlogger(__PACKAGE__)); $self->_setup_ua($ua,$self->{netloc}); return $ua; diff --git a/lib/NGCP/BulkProcessor/RestConnectors/NGCPRestApi.pm b/lib/NGCP/BulkProcessor/RestConnectors/NGCPRestApi.pm index cf276a74..a0ad70ae 100644 --- a/lib/NGCP/BulkProcessor/RestConnectors/NGCPRestApi.pm +++ b/lib/NGCP/BulkProcessor/RestConnectors/NGCPRestApi.pm @@ -22,7 +22,9 @@ use NGCP::BulkProcessor::LogError qw( restrequesterror restresponseerror); -use NGCP::BulkProcessor::RestConnector; +use NGCP::BulkProcessor::RestConnector qw(_add_headers); + +use NGCP::BulkProcessor::FakeTime qw(get_fake_now_string); require Exporter; our @ISA = qw(Exporter NGCP::BulkProcessor::RestConnector); @@ -41,6 +43,9 @@ my $first_page_num = 1; my $contenttype = 'application/json'; my $patchcontenttype = 'application/json-patch+json'; +my $defaultfaketime = 0; +my $faketime_header = 'X-Fake-Clienttime'; + our $ITEM_REL_PARAM = 'item_rel'; #my $logger = getlogger(__PACKAGE__); @@ -63,11 +68,12 @@ sub new { sub setup { my $self = shift; - my ($baseuri,$username,$password,$realm) = @_; + my ($baseuri,$username,$password,$realm,$faketime) = @_; $self->baseuri($baseuri // $defaulturi); $self->{username} = $username // $defaultusername; $self->{password} = $password // $defaultpassword; $self->{realm} = $realm // $defaultrealm; + $self->{faketime} = $faketime // $defaultfaketime; } @@ -93,6 +99,7 @@ sub _setup_ua { if ($self->{username}) { $ua->credentials($netloc, $self->{realm}, $self->{username}, $self->{password}); } + restdebug($self,"ua configured",getlogger(__PACKAGE__)); } @@ -105,7 +112,7 @@ sub _encode_request_content { sub _decode_response_content { my $self = shift; my ($data) = @_; - return JSON::from_json($data); + return ($data ? JSON::from_json($data) : undef); } sub _add_post_headers { @@ -113,6 +120,7 @@ sub _add_post_headers { my ($req,$headers) = @_; _add_headers($req,{ 'Content-Type' => $contenttype, + ($self->{faketime} ? ($faketime_header => get_fake_now_string()) : ()), }); # allow providing custom headers to post(), # e.g { 'X-Fake-Clienttime' => ... } @@ -122,6 +130,9 @@ sub _add_post_headers { sub _add_get_headers { my $self = shift; my ($req,$headers) = @_; + _add_headers($req,{ + ($self->{faketime} ? ($faketime_header => get_fake_now_string()) : ()), + }); $self->SUPER::_add_get_headers($req,$headers); } @@ -131,6 +142,7 @@ sub _add_patch_headers { _add_headers($req,{ 'Prefer' => 'return=representation', 'Content-Type' => $patchcontenttype, + ($self->{faketime} ? ($faketime_header => get_fake_now_string()) : ()), }); $self->SUPER::_add_patch_headers($req,$headers); } @@ -149,6 +161,7 @@ sub _add_put_headers { _add_headers($req,{ 'Prefer' => 'return=representation', 'Content-Type' => $contenttype, + ($self->{faketime} ? ($faketime_header => get_fake_now_string()) : ()), }); $self->SUPER::_add_put_headers($req,$headers); } @@ -156,6 +169,9 @@ sub _add_put_headers { sub _add_delete_headers { my $self = shift; my ($req,$headers) = @_; + _add_headers($req,{ + ($self->{faketime} ? ($faketime_header => get_fake_now_string()) : ()), + }); $self->SUPER::_add_delete_headers($req,$headers); } @@ -181,9 +197,10 @@ sub extract_collection_items { my $self = shift; my ($data,$page_size,$page_num,$params) = @_; my $result = undef; - if ('HASH' eq ref $data - and 'HASH' eq ref $data->{'_embedded'}) { + if (defined $data and 'HASH' eq ref $data + and defined $data->{'_embedded'} and 'HASH' eq ref $data->{'_embedded'}) { $result = $data->{'_embedded'}->{$params->{$ITEM_REL_PARAM}}; + undef $result unless ref $result; } $result //= []; return shared_clone($result); @@ -194,13 +211,97 @@ sub get_defaultcollectionpagesize { return $default_collection_page_size; } +sub _request_error { + my $self = shift; + my $msg = undef; + if (defined $self->responsedata() + and 'HASH' eq ref $self->responsedata()) { + $msg = $self->responsedata()->{'message'}; + } + resterror($self,$self->response->code . ' ' . $self->response->message . + (defined $msg && length($msg) > 0 ? ': ' . $msg : ''),getlogger(__PACKAGE__)); +} + +sub _extract_ids_from_response_location { + + my $self = shift; + my $location = $self->response()->header('Location'); + my @ids = (); + foreach my $segment (split('/',$location)) { + push(@ids,$segment) if $segment =~ /^\d+$/; + } + return @ids; + +} + sub get { my $self = shift; if ($self->_get(@_)->code() != HTTP_OK) { - resterror($self,$self->response->code . ' ' . $self->response->message,getlogger(__PACKAGE__)); + $self->_request_error(); + return undef; } else { return $self->responsedata(); } } +sub post { + my $self = shift; + if ($self->_post(@_)->code() != HTTP_CREATED) { + $self->_request_error(); + return (); + } else { + return $self->_extract_ids_from_response_location(); + } +} + +sub post_get { + my $self = shift; + my ($path_query,$post_headers,$get_headers) = @_; + if ($self->_post($path_query,$post_headers)->code() != HTTP_CREATED) { + $self->_request_error(); + return undef; + } else { + return $self->get($self->response()->header('Location'),$get_headers); + } +} + +sub put { + my $self = shift; + if ($self->_put(@_)->code() != HTTP_OK) { + $self->_request_error(); + return undef; + } else { + return $self->responsedata(); + } +} + +sub patch { + my $self = shift; + if ($self->_patch(@_)->code() != HTTP_OK) { + $self->_request_error(); + return undef; + } else { + return $self->responsedata(); + } +} + +sub delete { + my $self = shift; + if ($self->_delete(@_)->code() != HTTP_NO_CONTENT) { + $self->_request_error(); + return 0; + } else { + return 1; + } +} + +sub faketime { + my $self = shift; + if (@_) { + $self->{faketime} = shift; + restdebug($self,"fake time " . ($self->{faketime} ? 'enabled' : 'disabled'),getlogger(__PACKAGE__)); + } + return $self->{faketime}; +} + 1; diff --git a/lib/NGCP/BulkProcessor/RestProcessor.pm b/lib/NGCP/BulkProcessor/RestProcessor.pm index 0514368e..694cf220 100644 --- a/lib/NGCP/BulkProcessor/RestProcessor.pm +++ b/lib/NGCP/BulkProcessor/RestProcessor.pm @@ -188,7 +188,7 @@ sub process_collection { my $rowblock_result = 1; my $blockcount = 0; eval { - if ('CODE' eq ref $init_process_context_code) { + if (defined $init_process_context_code and 'CODE' eq ref $init_process_context_code) { &$init_process_context_code($context); } @@ -223,7 +223,7 @@ sub process_collection { } eval { - if ('CODE' eq ref $uninit_process_context_code) { + if (defined $uninit_process_context_code and 'CODE' eq ref $uninit_process_context_code) { &$uninit_process_context_code($context); } }; @@ -331,7 +331,7 @@ sub _process { my $blockcount = 0; eval { - if ('CODE' eq ref $context->{init_process_context_code}) { + if (defined $context->{init_process_context_code} and 'CODE' eq ref $context->{init_process_context_code}) { &{$context->{init_process_context_code}}($context); } while (not _get_stop_consumer_thread($context,$tid)) { @@ -363,7 +363,7 @@ sub _process { my $err = $@; restthreadingdebug($err ? '[' . $tid . '] processor thread error: ' . $err : '[' . $tid . '] processor thread finished (' . $blockcount . ' blocks)',getlogger(__PACKAGE__)); eval { - if ('CODE' eq ref $context->{uninit_process_context_code}) { + if (defined $context->{uninit_process_context_code} and 'CODE' eq ref $context->{uninit_process_context_code}) { &{$context->{uninit_process_context_code}}($context); } }; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingFees.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingFees.pm new file mode 100644 index 00000000..21f6ef03 --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingFees.pm @@ -0,0 +1,119 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::BillingFees; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item + delete_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'billingfees'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'billing_zone_id', + 'destination', + 'direction', + 'offpeak_follow_interval', + 'offpeak_follow_rate', + 'offpeak_init_interval', + 'offpeak_init_rate', + 'onpeak_follow_interval', + 'onpeak_follow_rate', + 'onpeak_init_interval', + 'onpeak_init_rate', + 'purge_existing', + 'source', + 'use_free_time', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub delete_item { + + my ($id,$headers) = @_; + my $restapi = &$get_restapi(); + ($id) = $restapi->delete(&$get_item_path_query($id),$headers); + return $id; + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingProfiles.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingProfiles.pm index 08de8b3c..e7779969 100644 --- a/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingProfiles.pm +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingProfiles.pm @@ -19,26 +19,38 @@ use NGCP::BulkProcessor::RestItem qw(); require Exporter; our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); our @EXPORT_OK = qw( + get_item + create_item process_items ); - +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'billingprofiles'; +my $item_relation = 'ngcp:' . $resource; my $get_item_path_query = sub { - my ($billing_profile_id) = @_; - return 'api/billingprofile/' . $billing_profile_id; + my ($id) = @_; + return 'api/' . $resource . '/' . $id; }; -my $collection_path_query = 'api/billingprofiles/'; -my $item_relation = 'ngcp:billingprofiles'; -my $get_restapi = \&get_ngcp_restapi; +my $collection_path_query = 'api/' . $resource . '/'; my $fieldnames = [ - 'id', - 'start_date', - 'end_date', - 'billing_profile_id', - 'contract_id', - 'product_id', - 'network_id', + 'currency', + 'fraud_daily_limit', + 'fraud_daily_lock', + 'fraud_daily_notify', + 'fraud_interval_limit', + 'fraud_interval_lock', + 'fraud_interval_notify', + 'fraud_use_reseller_rates', + 'handle', + 'interval_charge', + 'interval_free_cash', + 'interval_free_time', + 'name', + 'peaktime_special', + 'peaktime_weekdays', + 'prepaid', + 'reseller_id', ]; sub new { @@ -54,16 +66,35 @@ sub new { } +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} sub builditems_fromrows { my ($rows,$load_recursive) = @_; - my @items = (); my $item; if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); foreach my $row (@$rows) { $item = __PACKAGE__->new($row); @@ -71,9 +102,12 @@ sub builditems_fromrows { push @items,$item; } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; } - - return \@items; + return undef; } diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingZones.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingZones.pm new file mode 100644 index 00000000..5523ad7e --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/BillingZones.pm @@ -0,0 +1,98 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::BillingZones; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'billingzones'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'billing_profile_id', + 'detail', + 'zone', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/Contracts.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Contracts.pm new file mode 100644 index 00000000..c3c429f3 --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Contracts.pm @@ -0,0 +1,120 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::Contracts; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item + set_item + update_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'contracts'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'billing_profile_definition', + 'billing_profile_id', + 'billing_profiles', + 'contact_id', + 'external_id', + 'status', + 'type', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub set_item { + + my ($id,$data,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->put(&$get_item_path_query($id),$data,$headers),$load_recursive); + +} + +sub update_item { + + my ($id,$data,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->patch(&$get_item_path_query($id),$data,$headers),$load_recursive); + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/CustomerContacts.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/CustomerContacts.pm new file mode 100644 index 00000000..c3d69131 --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/CustomerContacts.pm @@ -0,0 +1,122 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::CustomerContacts; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'customercontacts'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'bankname', + 'bic', + 'city', + 'company', + 'comregnum', + 'country', + 'email', + 'faxnumber', + 'firstname', + 'gpp0', + 'gpp1', + 'gpp2', + 'gpp3', + 'gpp4', + 'gpp5', + 'gpp6', + 'gpp7', + 'gpp8', + 'gpp9', + 'iban', + 'lastname', + 'mobilenumber', + 'phonenumber', + 'postcode', + 'reseller_id', + 'street', + 'vatnum', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/Customers.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Customers.pm new file mode 100644 index 00000000..9c895cfb --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Customers.pm @@ -0,0 +1,128 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::Customers; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item + set_item + update_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'customers'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'add_vat', + 'billing_profile_definition', + 'billing_profile_id', + 'billing_profiles', + 'contact_id', + 'external_id', + 'invoice_email_template', + 'invoice_template', + 'max_subscribers', + 'passreset_email_template', + 'profile_package_id', + 'status', + 'subscriber_email_template', + 'type', + 'vat_rate', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub set_item { + + my ($id,$data,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->put(&$get_item_path_query($id),$data,$headers),$load_recursive); + +} + +sub update_item { + + my ($id,$data,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->patch(&$get_item_path_query($id),$data,$headers),$load_recursive); + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/Domains.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Domains.pm new file mode 100644 index 00000000..01be47f4 --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Domains.pm @@ -0,0 +1,98 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::Domains; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + process_collection + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'domains'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'domain', + 'reseller_id', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/Resellers.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Resellers.pm new file mode 100644 index 00000000..647811ed --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Resellers.pm @@ -0,0 +1,100 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::Resellers; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'resellers'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'contract_id', + 'enable_rtc', + 'name', + 'rtc_networks', + 'status', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/Subscribers.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Subscribers.pm new file mode 100644 index 00000000..8cf522b1 --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/Subscribers.pm @@ -0,0 +1,145 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::Subscribers; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item + set_item + update_item + delete_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'subscribers'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'administrative', + 'alias_numbers', + 'customer_id', + 'display_name', + 'domain', + 'domain_id', + 'email', + 'external_id', + 'is_pbx_group', + 'is_pbx_pilot', + 'lock', + 'password', + 'pbx_extension', + 'pbx_group_ids', + 'pbx_groupmember_ids_id', + 'primary_number', + 'profile_id', + 'profile_set_id', + 'status', + 'username', + 'webpassword', + 'webusername', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub set_item { + + my ($id,$data,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->put(&$get_item_path_query($id),$data,$headers),$load_recursive); + +} + +sub update_item { + + my ($id,$data,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->patch(&$get_item_path_query($id),$data,$headers),$load_recursive); + +} + +sub delete_item { + + my ($id,$headers) = @_; + my $restapi = &$get_restapi(); + ($id) = $restapi->delete(&$get_item_path_query($id),$headers); + return $id; + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/RestRequests/Trunk/SystemContacts.pm b/lib/NGCP/BulkProcessor/RestRequests/Trunk/SystemContacts.pm new file mode 100644 index 00000000..9929a14d --- /dev/null +++ b/lib/NGCP/BulkProcessor/RestRequests/Trunk/SystemContacts.pm @@ -0,0 +1,121 @@ +package NGCP::BulkProcessor::RestRequests::Trunk::SystemContacts; +use strict; + +## no critic + +use NGCP::BulkProcessor::ConnectorPool qw( + get_ngcp_restapi + +); + +use NGCP::BulkProcessor::RestProcessor qw( + copy_row +); + +use NGCP::BulkProcessor::RestConnectors::NGCPRestApi qw(); +use NGCP::BulkProcessor::RestItem qw(); + +require Exporter; +our @ISA = qw(Exporter NGCP::BulkProcessor::RestItem); +our @EXPORT_OK = qw( + get_item + create_item +); + +my $get_restapi = \&get_ngcp_restapi; +my $resource = 'systemcontacts'; +my $item_relation = 'ngcp:' . $resource; +my $get_item_path_query = sub { + my ($id) = @_; + return 'api/' . $resource . '/' . $id; +}; +my $collection_path_query = 'api/' . $resource . '/'; + +my $fieldnames = [ + 'bankname', + 'bic', + 'city', + 'company', + 'comregnum', + 'country', + 'email', + 'faxnumber', + 'firstname', + 'gpp0', + 'gpp1', + 'gpp2', + 'gpp3', + 'gpp4', + 'gpp5', + 'gpp6', + 'gpp7', + 'gpp8', + 'gpp9', + 'iban', + 'lastname', + 'mobilenumber', + 'phonenumber', + 'postcode', + 'street', + 'vatnum', +]; + +sub new { + + my $class = shift; + my $self = NGCP::BulkProcessor::RestItem->new($fieldnames); + + bless($self,$class); + + copy_row($self,shift,$fieldnames); + + return $self; + +} + +sub get_item { + + my ($id,$load_recursive,$headers) = @_; + my $restapi = &$get_restapi(); + return builditems_fromrows($restapi->get(&$get_item_path_query($id),$headers),$load_recursive); + +} + +sub create_item { + + my ($data,$load,$load_recursive,$post_headers,$get_headers) = @_; + my $restapi = &$get_restapi(); + if ($load) { + return builditems_fromrows($restapi->post_get($collection_path_query,$data,$post_headers,$get_headers),$load_recursive); + } else { + my ($id) = $restapi->post($collection_path_query,$data,$post_headers); + return $id; + } + +} + +sub builditems_fromrows { + + my ($rows,$load_recursive) = @_; + + my $item; + + if (defined $rows and ref $rows eq 'ARRAY') { + my @items = (); + foreach my $row (@$rows) { + $item = __PACKAGE__->new($row); + + # transformations go here ... + + push @items,$item; + } + return \@items; + } elsif (defined $rows and ref $rows eq 'HASH') { + $item = __PACKAGE__->new($rows); + return $item; + } + return undef; + +} + +1; diff --git a/lib/NGCP/BulkProcessor/SqlProcessor.pm b/lib/NGCP/BulkProcessor/SqlProcessor.pm index 7003fe0b..4f92dc3a 100644 --- a/lib/NGCP/BulkProcessor/SqlProcessor.pm +++ b/lib/NGCP/BulkProcessor/SqlProcessor.pm @@ -487,7 +487,7 @@ sub insert_record { my @unique_fields = (); my @unique_vals = (); - if ('ARRAY' eq ref $unique_count_fields) { + if (defined $unique_count_fields and 'ARRAY' eq ref $unique_count_fields) { foreach my $fieldname (@$unique_count_fields) { push(@unique_fields,$fieldname); if (exists $row->{$fieldname}) { @@ -1217,7 +1217,7 @@ sub process_table { my $context = { tid => $tid }; my $rowblock_result = 1; eval { - if ('CODE' eq ref $init_process_context_code) { + if (defined $init_process_context_code and 'CODE' eq ref $init_process_context_code) { &$init_process_context_code($context); } @@ -1256,7 +1256,7 @@ sub process_table { } eval { - if ('CODE' eq ref $uninit_process_context_code) { + if (defined $uninit_process_context_code and 'CODE' eq ref $uninit_process_context_code) { &$uninit_process_context_code($context); } }; @@ -1440,7 +1440,7 @@ sub _reader { # if thread cleanup has a problem... $reader_db->db_disconnect(); } - if ('CODE' eq ref $context->{destroy_dbs_code}) { + if (defined $context->{destroy_dbs_code} and 'CODE' eq ref $context->{destroy_dbs_code}) { &{$context->{destroy_dbs_code}}(); } lock $context->{errorstates}; @@ -1495,7 +1495,7 @@ sub _writer { # if thread cleanup has a problem... $writer_db->db_disconnect(); } - if ('CODE' eq ref $context->{destroy_dbs_code}) { + if (defined $context->{destroy_dbs_code} and 'CODE' eq ref $context->{destroy_dbs_code}) { &{$context->{destroy_dbs_code}}(); } lock $context->{errorstates}; @@ -1524,7 +1524,7 @@ sub _process { my $blockcount = 0; eval { - if ('CODE' eq ref $context->{init_process_context_code}) { + if (defined $context->{init_process_context_code} and 'CODE' eq ref $context->{init_process_context_code}) { &{$context->{init_process_context_code}}($context); } #$writer_db = &{$context->{get_target_db}}($writer_connection_name); @@ -1568,7 +1568,7 @@ sub _process { my $err = $@; tablethreadingdebug($err ? '[' . $tid . '] processor thread error: ' . $err : '[' . $tid . '] processor thread finished (' . $blockcount . ' blocks)',getlogger(__PACKAGE__)); eval { - if ('CODE' eq ref $context->{uninit_process_context_code}) { + if (defined $context->{uninit_process_context_code} and 'CODE' eq ref $context->{uninit_process_context_code}) { &{$context->{uninit_process_context_code}}($context); } }; diff --git a/lib/NGCP/BulkProcessor/Utils.pm b/lib/NGCP/BulkProcessor/Utils.pm index d25c4e73..a278f30a 100644 --- a/lib/NGCP/BulkProcessor/Utils.pm +++ b/lib/NGCP/BulkProcessor/Utils.pm @@ -9,14 +9,15 @@ use threads; use POSIX qw(strtod locale_h); setlocale(LC_NUMERIC, 'C'); -use Data::UUID; +use Data::UUID qw(); +use UUID qw(); -use Net::Address::IP::Local; +use Net::Address::IP::Local qw(); #use FindBin qw($Bin); #use File::Spec::Functions qw(splitdir catdir); use Net::Domain qw(hostname hostfqdn hostdomain); -use Cwd 'abs_path'; +use Cwd qw(abs_path); #use File::Basename qw(fileparse); use Date::Manip qw(Date_Init ParseDate UnixDate); @@ -26,9 +27,9 @@ Date_Init('DateFormat=US'); use Date::Calc qw(Normalize_DHMS Add_Delta_DHMS); -use Text::Wrap; +use Text::Wrap qw(); #use FindBin qw($Bin); -use Digest::MD5; #qw(md5 md5_hex md5_base64); +use Digest::MD5 qw(); #qw(md5 md5_hex md5_base64); use File::Temp qw(tempfile tempdir); use File::Path qw(remove_tree make_path); @@ -62,6 +63,7 @@ our @EXPORT_OK = qw( cat_file wrap_text create_guid + create_uuid urlencode urldecode timestamp @@ -368,6 +370,13 @@ sub create_guid { } +sub create_uuid { + my ($bin, $str); + UUID::generate($bin); + UUID::unparse($bin, $str); + return $str; +} + sub urlencode { my ($urltoencode) = @_; $urltoencode =~ s/([^a-zA-Z0-9\/_\-.])/uc sprintf("%%%02x",ord($1))/eg;