svn commit: r1934076 - spamassassin/trunk/lib/Mail/SpamAssassin/Plugin
[email protected] Mon, 11 May 2026 07:25:15 -0000
| Newsgroups | gmane.mail.spam.spamassassin.cvs |
|---|---|
| Message-ID | <177848431572.1282311.10558968041162884484@svn03-he-fi> |
Author: gbechis
Date: Mon May 11 07:25:15 2026
New Revision: 1934076
Log:
improve replay algorithm, improve locking
Modified:
spamassassin/trunk/lib/Mail/SpamAssassin/Plugin/NeuralNetwork.pm
Modified: spamassassin/trunk/lib/Mail/SpamAssassin/Plugin/NeuralNetwork.pm
==============================================================================
--- spamassassin/trunk/lib/Mail/SpamAssassin/Plugin/NeuralNetwork.pm Mon May 11 07:20:16 2026 (r1934075)
+++ spamassassin/trunk/lib/Mail/SpamAssassin/Plugin/NeuralNetwork.pm Mon May 11 07:25:15 2026 (r1934076)
@@ -44,7 +44,7 @@ use strict;
use warnings;
use re 'taint';
-my $VERSION = 0.8.4;
+my $VERSION = 0.9.0;
use AI::FANN qw(:all);
use Storable qw(store retrieve);
@@ -519,6 +519,7 @@ sub _text_to_features {
# When training, build per-document term sets to update doc counts
my $local_doc_increment = 0;
+ my %term_deltas;
if ($train == 1) {
foreach my $tok_ref (@token_lists) {
next unless ref($tok_ref) eq 'ARRAY' && @$tok_ref;
@@ -529,15 +530,19 @@ sub _text_to_features {
my %seen;
foreach my $t (@tokens) {
$vocabulary{terms}{$t}{total} = ($vocabulary{terms}{$t}{total} || 0) + 1;
+ $term_deltas{$t}{total} = ($term_deltas{$t}{total} || 0) + 1;
$seen{$t} = 1;
}
foreach my $t (keys %seen) {
$vocabulary{terms}{$t}{docs} = ($vocabulary{terms}{$t}{docs} || 0) + 1;
+ $term_deltas{$t}{docs} = ($term_deltas{$t}{docs} || 0) + 1;
# Track spam/ham sources
if (defined $label && $label == 1) {
$vocabulary{terms}{$t}{spam} = ($vocabulary{terms}{$t}{spam} || 0) + 1;
+ $term_deltas{$t}{spam} = ($term_deltas{$t}{spam} || 0) + 1;
} elsif (defined $label && $label == 0) {
$vocabulary{terms}{$t}{ham} = ($vocabulary{terms}{$t}{ham} || 0) + 1;
+ $term_deltas{$t}{ham} = ($term_deltas{$t}{ham} || 0) + 1;
}
}
}
@@ -556,7 +561,7 @@ sub _text_to_features {
my $vocab_path;
if (defined $conf->{neuralnetwork_dsn} && $self && $self->{dbh}) {
- $self->_save_vocabulary_to_sql(\%vocabulary, $self->{main}->{username});
+ $self->_save_vocabulary_to_sql(\%term_deltas, $self->{main}->{username});
} else {
$vocab_path = File::Spec->catfile($nn_data_dir, 'vocabulary-' . lc($self->{main}->{username}) . '.data');
$vocab_path = Mail::SpamAssassin::Util::untaint_file_path($vocab_path);
@@ -682,7 +687,7 @@ sub learn_message {
dbg("Not enough text, skipping neural network processing");
return;
}
- my $tokens_ref = $self->_extract_features_from_message($conf, $msg);
+ my $tokens_ref = $self->_extract_features_from_message($pms, $conf, $msg);
push(@training_data, { label => $isspam, tokens => $tokens_ref } );
}
@@ -690,7 +695,7 @@ sub learn_message {
$dataset_path = Mail::SpamAssassin::Util::untaint_file_path($dataset_path);
my $locker = $self->{main}->{locker};
- unless ($locker->safe_lock($dataset_path, $conf->{neuralnetwork_lock_timeout})) {
+unless ($locker->safe_lock($dataset_path, $conf->{neuralnetwork_lock_timeout})) {
dbg("Cannot acquire lock on '$dataset_path', skipping learning");
return;
}
@@ -725,14 +730,12 @@ sub learn_message {
my ($feature_vectors, $vocab_size, $vocab_keys_ref) = _text_to_features($self, $self->{main}->{conf}, $nn_data_dir, $update_vocab, $isspam, undef, @email_token_lists);
unless ($feature_vectors && @$feature_vectors) {
- $locker->safe_unlock($dataset_path);
return;
}
my $num_input = scalar(@{$feature_vectors->[0]{vec}});
if ($num_input == 0) {
dbg("No valid features found in message, skipping learning");
- $locker->safe_unlock($dataset_path);
return;
}
# Defer model creation until minimum training corpus is built
@@ -744,7 +747,6 @@ sub learn_message {
my $hc = $vocab->{_ham_count} || 0;
if ($sc < $min_spam || $hc < $min_ham) {
dbg("Deferring model creation: spam=$sc/$min_spam ham=$hc/$min_ham (vocabulary updated)");
- $locker->safe_unlock($dataset_path);
# Record message as learned to prevent re-learning.
if (defined $msg && defined $msgid && length($msgid) > 0) {
$self->_save_msgid_to_neural_seen($msgid, $isspam);
@@ -753,6 +755,10 @@ sub learn_message {
}
}
+ # Snapshot the on-disk model mtime so the save section can detect that
+ # another writer modified the model file while we were training.
+ my $lock1_mtime = (stat($dataset_path))[9];
+
# Two-layer hidden sizing: layer1 ~10% of inputs (clamped 32..512),
# layer2 half of layer1 (min 16).
my $num_hidden1 = int($num_input / 10);
@@ -788,14 +794,30 @@ sub learn_message {
$num_input = $model_size;
$network = $existing_network;
} else {
- dbg("Model vocabulary mismatch, rebuilding model with current vocabulary ($num_input terms)");
+ # Model and vocab files are inconsistent, discard the stale network and fall through to
+ # a fresh one.
+ my $stored_size = defined $stored_vocab_ref ? scalar(@$stored_vocab_ref) : 0;
+ info("Model/vocab mismatch (stored=$stored_size vs model=$model_size inputs); " .
+ "discarding stale model, starting fresh network");
+ undef $existing_network;
}
}
unless (defined $network) {
- # No existing model or inconsistent model, create a baseline from vocabulary
- $network = $self->_retrain_from_vocabulary($self->{main}->{conf}, $nn_data_dir, $num_input);
+ # No existing model, or a consistency error cleared it above.
+ # If the vocabulary already has enough balanced data, build a baseline
+ # from it; otherwise start with a randomly-initialised network that will
+ # be shaped by subsequent incremental training.
+ my $vocab = $self->{_last_train_vocab} // {};
+ my $sc = $vocab->{_spam_count} || 0;
+ my $hc = $vocab->{_ham_count} || 0;
+ my $min_spam = $self->{main}->{conf}->{neuralnetwork_min_spam_count};
+ my $min_ham = $self->{main}->{conf}->{neuralnetwork_min_ham_count};
+ if ($sc >= $min_spam && $hc >= $min_ham) {
+ $network = $self->_retrain_from_vocabulary($self->{main}->{conf}, $nn_data_dir, $num_input);
+ }
if (!defined $network) {
- # No vocabulary stats available yet, create a fresh network
+ # Vocabulary not yet balanced enough, or retrain returned undef.
+ # Create a blank network; it will learn through incremental training.
$network = AI::FANN->new_standard($num_input, $num_hidden1, $num_hidden2, $num_output_neurons);
my $act_fn = ($train_algorithm == FANN_TRAIN_RPROP) ? FANN_SIGMOID_STEPWISE : FANN_SIGMOID;
$network->hidden_activation_function($act_fn);
@@ -812,8 +834,10 @@ sub learn_message {
# reuse the cached vocabulary for class-balance accounting
my %vocab_for_balance = %{ delete $self->{_last_train_vocab} || {} };
- my $spam_docs = $vocab_for_balance{_spam_count} || 1;
- my $ham_docs = $vocab_for_balance{_ham_count} || 1;
+ my $raw_spam = $vocab_for_balance{_spam_count} || 0;
+ my $raw_ham = $vocab_for_balance{_ham_count} || 0;
+ my $spam_docs = $raw_spam || 1;
+ my $ham_docs = $raw_ham || 1;
# class_weight > 1 means this message belongs to the minority class and
# should be trained harder; < 1 means it belongs to the majority class.
@@ -841,18 +865,15 @@ sub learn_message {
($train_algorithm == FANN_TRAIN_RPROP)
&& defined $vocab_keys_ref
&& scalar(@$vocab_keys_ref) == $num_input
- && $spam_docs > 1 && $ham_docs > 1;
+ && $raw_spam >= 1 && $raw_ham >= 1;
if ($replay_eligible) {
($svec, $hvec) = _build_class_tfidf_vectors(\%vocab_for_balance, $vocab_keys_ref);
} elsif ($train_algorithm == FANN_TRAIN_RPROP && !defined $vocab_keys_ref) {
dbg("Skipping RPROP replay: vocab_keys unavailable");
}
- # drop the lock and re-acquire it later
my $locked_num_input = $num_input;
my $locked_vocab_keys_ref = $vocab_keys_ref;
- my $lock1_mtime = (stat($dataset_path))[9];
- $locker->safe_unlock($dataset_path);
for my $e (1 .. $weighted_epochs) {
for my $i (0 .. $#$feature_vectors) {
@@ -895,23 +916,10 @@ sub learn_message {
dbg("Prediction after learning: " . (defined $pred_after ? $pred_after : 'undef'));
}
- # re-acquire the lock briefly to save the model
- unless ($locker->safe_lock($dataset_path, $conf->{neuralnetwork_lock_timeout})) {
- info("Cannot re-acquire lock on '$dataset_path' for model save; skipping persistence (vocabulary already updated)");
- return 0;
- }
-
- # Log if we are going to overwrite the model
- my $current_mtime = (stat($dataset_path))[9] // 0;
- if (defined $lock1_mtime && $current_mtime > $lock1_mtime) {
- info("on-disk model changed during training; overwriting");
- }
-
- # Save the model atomically
- my $model_saved = 0;
+ # Pre-stage the FANN model to a temp file BEFORE acquiring the lock.
my $tmp_path;
my $file_mode = 0666 & ~umask();
- eval {
+ my $prestage_ok = eval {
my ($vol, $dir, undef) = File::Spec->splitpath($dataset_path);
my $tmp_dir = File::Spec->catpath($vol, $dir, '');
(undef, $tmp_path) = File::Temp::tempfile(
@@ -923,11 +931,46 @@ sub learn_message {
$tmp_path = Mail::SpamAssassin::Util::untaint_file_path($tmp_path);
chmod($file_mode, $tmp_path) or info("chmod $file_mode on '$tmp_path' failed: $!");
$network->save($tmp_path) or die "model save to temp '$tmp_path' failed";
+ 1;
+ };
+ if (!$prestage_ok) {
+ my $err = $@ || 'unknown';
+ info("Cannot pre-stage model save to '$dataset_path' ($err)");
+ if (defined $tmp_path && -f $tmp_path) {
+ unlink($tmp_path)
+ or info("Cannot remove temp model file '$tmp_path': $!");
+ }
+ return 0;
+ }
+ # Re-acquire the lock briefly to commit.
+ unless ($locker->safe_lock($dataset_path, $conf->{neuralnetwork_lock_timeout})) {
+ info("Cannot re-acquire lock on '$dataset_path' for model save; skipping persistence (vocabulary already updated)");
+ unlink($tmp_path)
+ or info("Cannot remove temp model file '$tmp_path': $!") if -f $tmp_path;
+ return 0;
+ }
+
+ my $current_mtime = (stat($dataset_path))[9] // 0;
+ if (defined $lock1_mtime && $current_mtime > $lock1_mtime) {
+ info("on-disk model changed during training; overwriting");
+ }
+
+ my $model_saved = 0;
+ eval {
if (defined $self->{main}->{conf}->{neuralnetwork_dsn} && $self->{dbh}) {
- $self->_save_model_vocab_to_sql($locked_vocab_keys_ref);
+ $self->_save_model_vocab_to_sql($locked_vocab_keys_ref)
+ or info("WARNING: model saved but vocab SQL write failed; " .
+ "model/vocab are now inconsistent and a full rebuild will " .
+ "occur on the next learn call");
} else {
$self->_save_model_vocab($locked_vocab_keys_ref, $nn_data_dir);
+ my $vocab_path = $self->_model_vocab_path($nn_data_dir);
+ unless (-f $vocab_path) {
+ info("WARNING: model saved but vocab file '$vocab_path' is missing; " .
+ "model/vocab are now inconsistent and a full rebuild will " .
+ "occur on the next learn call");
+ }
}
delete $self->{_model_vocab_cache};
delete $self->{_model_vocab_cache_t};
@@ -980,7 +1023,7 @@ sub forget_message {
}
# Decrement vocabulary counts for tokens in this message
- my $tokens_ref = $self->_extract_features_from_message($conf, $msg);
+ my $tokens_ref = $self->_extract_features_from_message(undef, $conf, $msg);
my $deleted_count = 0;
if (ref $tokens_ref eq 'ARRAY' && @$tokens_ref) {
@@ -1065,43 +1108,56 @@ sub forget_message {
my $full_terms = ref($full_vocab_ref) eq 'HASH' ? ($full_vocab_ref->{terms} || {}) : {};
my $full_vocab_size = scalar keys %$full_terms;
if ($full_vocab_size > 0) {
- my $locker = $self->{main}->{locker};
- my $got_lock = 0;
- eval { $got_lock = $locker->safe_lock($dataset_path, $conf->{neuralnetwork_lock_timeout}); 1 };
- my $rebuilt = eval { $self->_retrain_from_vocabulary($conf, $nn_data_dir, $full_vocab_size) };
- if ($rebuilt) {
+ my $rebuilt = eval { $self->_retrain_from_vocabulary($conf, $nn_data_dir, $full_vocab_size) };
+ if (!$rebuilt) {
+ dbg("NeuralNetwork: Retrain after forget failed: " . ($@ || 'undef'));
+ } else {
+ # Pre-stage the FANN model to a temp file BEFORE acquiring the lock.
+ my $tmp_path;
my $file_mode = 0666 & ~umask();
- eval {
+ my $prestage_ok = eval {
my ($vol, $dir) = File::Spec->splitpath($dataset_path);
my $tmp_dir = File::Spec->catpath($vol, $dir, '');
- my (undef, $tmp_path) = File::Temp::tempfile(
+ (undef, $tmp_path) = File::Temp::tempfile(
'fann-XXXXXX', DIR => $tmp_dir, SUFFIX => '.tmp', UNLINK => 0);
$tmp_path = Mail::SpamAssassin::Util::untaint_file_path($tmp_path);
chmod($file_mode, $tmp_path) or info("chmod $file_mode on '$tmp_path' failed: $!");
- if ($rebuilt->save($tmp_path)) {
- rename($tmp_path, $dataset_path) or die "rename failed: $!";
- } else {
- unlink $tmp_path;
- die "save failed";
- }
+ $rebuilt->save($tmp_path) or die "save failed";
1;
- } or do {
- info("NeuralNetwork: Could not persist retrained model after forget: " . ($@ || 'unknown'));
};
- my $new_vocab_keys = [ sort keys %$full_terms ];
- if (defined $conf->{neuralnetwork_dsn} && $self->{dbh}) {
- $self->_save_model_vocab_to_sql($new_vocab_keys);
+ if (!$prestage_ok) {
+ info("NeuralNetwork: Could not pre-stage retrained model after forget: " . ($@ || 'unknown'));
+ unlink($tmp_path) if defined $tmp_path && -f $tmp_path;
} else {
- $self->_save_model_vocab($new_vocab_keys, $nn_data_dir);
+ my $locker = $self->{main}->{locker};
+ my $got_lock = 0;
+ eval { $got_lock = $locker->safe_lock($dataset_path, $conf->{neuralnetwork_lock_timeout}); 1 };
+ if (!$got_lock) {
+ info("NeuralNetwork: Cannot acquire lock for save after forget; in-memory model not persisted");
+ unlink($tmp_path) if defined $tmp_path && -f $tmp_path;
+ } else {
+ eval {
+ rename($tmp_path, $dataset_path) or die "rename failed: $!";
+ $tmp_path = undef;
+ 1;
+ } or do {
+ info("NeuralNetwork: Could not persist retrained model after forget: " . ($@ || 'unknown'));
+ unlink($tmp_path) if defined $tmp_path && -f $tmp_path;
+ };
+ my $new_vocab_keys = [ sort keys %$full_terms ];
+ if (defined $conf->{neuralnetwork_dsn} && $self->{dbh}) {
+ $self->_save_model_vocab_to_sql($new_vocab_keys);
+ } else {
+ $self->_save_model_vocab($new_vocab_keys, $nn_data_dir);
+ }
+ $self->{neural_model} = $rebuilt;
+ $self->{_neural_model_load_time} = time();
+ delete $self->{_model_vocab_cache};
+ info("NeuralNetwork: Retrained model with $full_vocab_size vocabulary terms after forget");
+ $locker->safe_unlock($dataset_path);
+ }
}
- $self->{neural_model} = $rebuilt;
- $self->{_neural_model_load_time} = time();
- delete $self->{_model_vocab_cache};
- info("NeuralNetwork: Retrained model with $full_vocab_size vocabulary terms after forget");
- } else {
- dbg("NeuralNetwork: Retrain after forget failed");
}
- $locker->safe_unlock($dataset_path) if $got_lock;
}
}
}
@@ -1234,8 +1290,8 @@ sub _build_class_tfidf_vectors {
my $idf = log(($N + 1) / (($td->{docs} || 0) + 1)) + 1;
my $spam_rate = ($td->{spam} || 0) / $spam_docs;
my $ham_rate = ($td->{ham} || 0) / $ham_docs;
- $spam_vec[$i] = ($spam_rate > $ham_rate) ? ($spam_rate - $ham_rate) * $idf : 0;
- $ham_vec[$i] = ($ham_rate > $spam_rate) ? ($ham_rate - $spam_rate) * $idf : 0;
+ $spam_vec[$i] = $spam_rate * $idf;
+ $ham_vec[$i] = $ham_rate * $idf;
}
for my $vec (\@spam_vec, \@ham_vec) {
@@ -1359,7 +1415,7 @@ sub _check_neuralnetwork {
dbg("Too short email text");
return;
}
- my $tokens_ref = $self->_extract_features_from_message($conf, $msg);
+ my $tokens_ref = $self->_extract_features_from_message($pms, $conf, $msg);
my $nn_data_dir = $self->{main}->{conf}->{neuralnetwork_data_dir};
$nn_data_dir = Mail::SpamAssassin::Util::untaint_file_path($nn_data_dir);
@@ -1393,64 +1449,69 @@ sub _check_neuralnetwork {
dbg("Model file changed on disk, reloading");
}
- my $got_lock = 0;
- eval {
- $got_lock = $locker->safe_lock($dataset_path,
- $self->{main}->{conf}->{neuralnetwork_lock_timeout});
- 1;
- };
-
my $file_mode = 0666 & ~umask();
- eval {
- my $loaded = AI::FANN->new_from_file($dataset_path);
- $self->{neural_model} = $loaded;
+
+ my $loaded = eval { AI::FANN->new_from_file($dataset_path) };
+ if ($loaded) {
+ $self->{neural_model} = $loaded;
$self->{_neural_model_load_time} = time();
- 1;
- } or do {
+ $network = $loaded;
+ } else {
my $err = $@ || 'unknown';
my @stat = stat($dataset_path);
my $fsize = @stat ? $stat[7] : 'N/A';
dbg("Failed to load model for prediction: $err "
. "(path=$dataset_path, size=${fsize}B), attempting vocabulary rebuild");
- # rebuild an in-memory model from vocabulary statistics
undef $self->{neural_model};
- my $rebuild_vocab_ref = $self->_load_model_vocab($nn_data_dir);
- my $rebuild_vocab_size = (defined $rebuild_vocab_ref && @$rebuild_vocab_ref)
- ? scalar(@$rebuild_vocab_ref) : 0;
- my $rebuilt = eval {
- $self->_retrain_from_vocabulary($conf, $nn_data_dir, $rebuild_vocab_size);
- };
- if ($rebuilt) {
- dbg("Vocabulary rebuild succeeded");
- $self->{neural_model} = $rebuilt;
- $self->{_neural_model_load_time} = time();
- # Persist the rebuilt model
- eval {
- my ($vol, $dir) = File::Spec->splitpath($dataset_path);
- my $tmp_dir = File::Spec->catpath($vol, $dir, '');
- my (undef, $tmp_path) = File::Temp::tempfile(
- 'fann-XXXXXX', DIR => $tmp_dir, SUFFIX => '.tmp', UNLINK => 0);
- $tmp_path = Mail::SpamAssassin::Util::untaint_file_path($tmp_path);
- chmod($file_mode, $tmp_path) or info("chmod $file_mode on '$tmp_path' failed: $!");
- if ($rebuilt->save($tmp_path)) {
- rename($tmp_path, $dataset_path)
- or die "rename failed: $!";
- dbg("Persisted rebuilt model to '$dataset_path'");
- } else {
- unlink $tmp_path;
- }
- 1;
- } or dbg("Could not persist rebuilt model: " . ($@ || 'unknown'));
- } else {
+
+ # skip prediction if we cannot acquire the lock fast enough
+ my $got_lock = 0;
+ eval { $got_lock = $locker->safe_lock($dataset_path, 1); 1; };
+ if (!$got_lock) {
+ dbg("another worker is rebuilding the model; skipping prediction");
+ return;
+ }
+
+ my $rebuilt;
+ eval {
+ my $vref = $self->_load_model_vocab($nn_data_dir);
+ my $vsz = (defined $vref && @$vref) ? scalar(@$vref) : 0;
+ $rebuilt = $self->_retrain_from_vocabulary($conf, $nn_data_dir, $vsz);
+ 1;
+ } or dbg("Vocabulary rebuild failed: " . ($@ || 'unknown'));
+
+ # Release the rebuild lock before the FANN save.
+ $locker->safe_unlock($dataset_path);
+
+ if (!$rebuilt) {
dbg("Vocabulary rebuild failed");
- $locker->safe_unlock($dataset_path) if $got_lock;
return;
}
- };
- $network = $self->{neural_model};
- $locker->safe_unlock($dataset_path) if $got_lock;
+ dbg("Vocabulary rebuild succeeded");
+ $self->{neural_model} = $rebuilt;
+ $self->{_neural_model_load_time} = time();
+
+ eval {
+ my ($vol, $dir) = File::Spec->splitpath($dataset_path);
+ my $tmp_dir = File::Spec->catpath($vol, $dir, '');
+ my (undef, $tmp_path) = File::Temp::tempfile(
+ 'fann-XXXXXX', DIR => $tmp_dir, SUFFIX => '.tmp', UNLINK => 0);
+ $tmp_path = Mail::SpamAssassin::Util::untaint_file_path($tmp_path);
+ chmod($file_mode, $tmp_path) or info("chmod $file_mode on '$tmp_path' failed: $!");
+ if ($rebuilt->save($tmp_path)) {
+ rename($tmp_path, $dataset_path)
+ or die "rename failed: $!";
+ dbg("Persisted rebuilt model to '$dataset_path'");
+ } else {
+ unlink $tmp_path;
+ }
+ 1;
+ } or dbg("Could not persist rebuilt model: " . ($@ || 'unknown'));
+
+ $network = $rebuilt;
+ }
} else {
$network = $self->{neural_model};
}
@@ -1670,16 +1731,14 @@ sub _is_msgid_in_neural_seen {
}
sub _save_vocabulary_to_sql {
- my ($self, $vocabulary, $username) = @_;
- return unless $self->{dbh} && defined $vocabulary && ref($vocabulary) eq 'HASH';
+ my ($self, $term_deltas, $username) = @_;
+ return unless $self->{dbh} && defined $term_deltas && ref($term_deltas) eq 'HASH';
$username //= $self->{main}->{username};
eval {
- my $terms = $vocabulary->{terms} || {};
- return unless scalar keys %{$terms};
+ return unless scalar keys %{$term_deltas};
- # Use ON DUPLICATE KEY UPDATE for MySQL or ON CONFLICT for other databases
my $upsert_sql;
if ($self->{main}->{conf}->{neuralnetwork_dsn} =~ /^dbi:(?:mysql|MariaDB)/i) {
@@ -1687,20 +1746,20 @@ sub _save_vocabulary_to_sql {
INSERT INTO neural_vocabulary (username, keyword, total_count, docs_count, spam_count, ham_count)
VALUES (?, ?, ?, ?, ?, ?)
ON DUPLICATE KEY UPDATE
- total_count = VALUES(total_count),
- docs_count = VALUES(docs_count),
- spam_count = VALUES(spam_count),
- ham_count = VALUES(ham_count)
+ total_count = total_count + VALUES(total_count),
+ docs_count = docs_count + VALUES(docs_count),
+ spam_count = spam_count + VALUES(spam_count),
+ ham_count = ham_count + VALUES(ham_count)
";
} else {
$upsert_sql = "
INSERT INTO neural_vocabulary (username, keyword, total_count, docs_count, spam_count, ham_count)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT (username, keyword) DO UPDATE SET
- total_count = excluded.total_count,
- docs_count = excluded.docs_count,
- spam_count = excluded.spam_count,
- ham_count = excluded.ham_count
+ total_count = neural_vocabulary.total_count + excluded.total_count,
+ docs_count = neural_vocabulary.docs_count + excluded.docs_count,
+ spam_count = neural_vocabulary.spam_count + excluded.spam_count,
+ ham_count = neural_vocabulary.ham_count + excluded.ham_count
";
}
@@ -1708,21 +1767,21 @@ sub _save_vocabulary_to_sql {
my $count = 0;
$self->{dbh}->begin_work();
- foreach my $keyword (sort keys %{$terms}) {
- my $term_data = $terms->{$keyword};
+ foreach my $keyword (sort keys %{$term_deltas}) {
+ my $delta = $term_deltas->{$keyword};
$sth_upsert->execute(
lc($username),
$keyword,
- $term_data->{total} || 0,
- $term_data->{docs} || 0,
- $term_data->{spam} || 0,
- $term_data->{ham} || 0
+ $delta->{total} || 0,
+ $delta->{docs} || 0,
+ $delta->{spam} || 0,
+ $delta->{ham} || 0
);
$count++;
}
$self->{dbh}->commit();
- dbg("Saved $count vocabulary terms to SQL for user: $username");
+ dbg("Applied $count vocabulary term deltas to SQL for user: $username");
# Invalidate cache for this user
if (defined $self->{_vocab_cache}) {
@@ -1792,16 +1851,15 @@ sub _tokenize_uri {
}
sub _extract_uris_from_msg {
- my ($msg) = @_;
- return () unless defined $msg;
- my $body = eval { $msg->get_pristine_body() };
- return () unless defined $body;
- $body = join("\n", @$body) if ref $body eq 'ARRAY';
+ my ($pms) = @_;
+
+ return () unless defined $pms;
+
my @uris;
- while ($body =~ m{((?:https?|ftp)://[^\s<>"'()\[\]\\]+)}gi) {
- my $u = $1;
- $u =~ s/[\.,;:!?]+$//;
- push @uris, $u if length $u;
+ my $uris = $pms->get_uri_detail_list();
+ my %huris = %{$uris};
+ foreach my $uri (keys %huris) {
+ push(@uris, $uri);
}
return @uris;
}
@@ -1810,7 +1868,7 @@ sub _extract_uris_from_msg {
# invisible (HTML-hidden) body, a small whitelist of headers, the URIs in
# the message, and the MIME-part digests.
sub _extract_features_from_message {
- my ($self, $conf, $msg) = @_;
+ my ($self, $pms, $conf, $msg) = @_;
my @tokens;
return \@tokens unless defined $msg;
@@ -1837,7 +1895,7 @@ sub _extract_features_from_message {
push @tokens, _tokenize_header_value('H*rcv:', $lines[-1]) if @lines;
}
- for my $u (_extract_uris_from_msg($msg)) {
+ for my $u (_extract_uris_from_msg($pms)) {
push @tokens, _tokenize_uri($u);
}
@@ -1943,6 +2001,38 @@ sub _load_vocabulary_from_sql {
);
eval {
+ my $meta_sth = $self->{dbh}->prepare(
+ "SELECT COUNT(*),
+ SUM(CASE WHEN flag = 'S' THEN 1 ELSE 0 END),
+ SUM(CASE WHEN flag = 'H' THEN 1 ELSE 0 END)
+ FROM neural_seen WHERE username = ?"
+ );
+ $meta_sth->execute($username);
+ my $meta = $meta_sth->fetchrow_arrayref();
+ if ($meta) {
+ $vocabulary{_doc_count} = $meta->[0] || 0;
+ $vocabulary{_spam_count} = $meta->[1] || 0;
+ $vocabulary{_ham_count} = $meta->[2] || 0;
+ }
+ 1;
+ } or do {
+ dbg("Pre-check query failed: " . ($@ || 'unknown'));
+ };
+
+ my $conf = $self->{main}->{conf};
+ my $min_spam = $conf->{neuralnetwork_min_spam_count};
+ my $min_ham = $conf->{neuralnetwork_min_ham_count};
+
+ if ($vocabulary{_spam_count} < $min_spam || $vocabulary{_ham_count} < $min_ham) {
+ dbg("Pre-check: insufficient training data " .
+ "(spam=$vocabulary{_spam_count}/$min_spam, ham=$vocabulary{_ham_count}/$min_ham)" .
+ ", skipping full vocabulary load");
+ $self->{_vocab_cache}{$username} = \%vocabulary;
+ $self->{_vocab_cache_time}{$username} = time();
+ return \%vocabulary;
+ }
+
+ eval {
my $sth = $self->{dbh}->prepare("
SELECT keyword, total_count, docs_count, spam_count, ham_count
FROM neural_vocabulary
@@ -1964,21 +2054,6 @@ sub _load_vocabulary_from_sql {
$count++;
}
- my $meta_sth = $self->{dbh}->prepare("
- SELECT COUNT(*),
- SUM(CASE WHEN flag = 'S' THEN 1 ELSE 0 END),
- SUM(CASE WHEN flag = 'H' THEN 1 ELSE 0 END)
- FROM neural_seen
- WHERE username = ?
- ");
- $meta_sth->execute($username);
- my $meta = $meta_sth->fetchrow_arrayref();
- if ($meta) {
- $vocabulary{_doc_count} = $meta->[0] || 0;
- $vocabulary{_spam_count} = $meta->[1] || 0;
- $vocabulary{_ham_count} = $meta->[2] || 0;
- }
-
dbg("Loaded $count vocabulary terms from SQL for user: $username " .
"(spam_docs=$vocabulary{_spam_count}, ham_docs=$vocabulary{_ham_count})");
1;