MT#18663 row bulk processing framework WIP #10

+reject invalid subscriber numbers by pattern
+detect changed passwords in subscriber+username.xslx
+separate --skip-errors option from --dry option
+rest request classes for ngcp resources:
 +SystemContacts
 +Contracts
 +BillingProfiles
 +Resellers
 +BillingZones
 +BillingFees
 +Domains
 +CustomerContacts
 +Customers
 +Subscribers
+making post, get, put, patch, delete work for
 NGCP rest-api

Change-Id: Ic55caed7d5425716adfbba9e2bc6b9153334481c
changes/95/6995/6
Rene Krenn 10 years ago
parent 5e47afdfc9
commit 84575f90e6

@ -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,

1
debian/control vendored

@ -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,

@ -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;

@ -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 {

@ -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);
}
};

@ -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;

@ -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) = @_;

@ -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 {

@ -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,'?');
}

@ -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,'?');
}

@ -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,'?');
}

@ -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,'?');
}

@ -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,'?');
}

@ -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,'?');
}

@ -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;

@ -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;

@ -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');

@ -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) = @_;

@ -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
);

@ -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

@ -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;

@ -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

@ -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;

@ -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;

@ -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;

@ -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);
}
};

@ -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;

@ -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;
}

@ -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;

@ -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;

@ -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;

@ -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;

@ -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;

@ -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;

@ -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;

@ -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;

@ -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);
}
};

@ -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;

Loading…
Cancel
Save