Algorithm-Classifier-IsolationForest

 view release on metacpan or  search on metacpan

lib/Algorithm/Classifier/IsolationForest/App/Command/streamd.pm  view on Meta::CPAN

		[ 'subsample=f', 'Per-tree stream subsampling probability, in (0, 1] (new models only).' ],
		[ 's=i',         'Seed int (new models only).' ],
		[
			'c=f',
			'Contamination. Expected fraction of anomalies, in (0, 0.5]; the decision threshold is '
				. 'relearned from the window before every save (new models only).'
		],
		[
			't=s@',
			'Feature name tag. Pass once per feature; enables the tagged (JSON object) row form '
				. '(new models only).'
		],
		[
			'mungers=s',
			'JSON file of Algorithm::ToNumberMunger specs, keyed by feature tag (new models only; requires -t).',
			{ 'completion' => 'files' }
		],
		[
			'prototype=s',
			'JSON prototype file to create the model from (new models only). May not be combined '
				. 'with -t or --mungers. See PROTOTYPES in the module POD.',
			{ 'completion' => 'files' }
		],
	);
} ## end sub opt_spec

sub abstract { 'Run an Online Isolation Forest scoring daemon on a Unix socket, speaking JSON lines' }

sub description {
	'Runs a prequential scoring daemon around an Online Isolation Forest
model (Algorithm::Classifier::IsolationForest::Online): clients connect
to the Unix domain socket and exchange one JSON document per line.

At startup the daemon resumes from <model-dir>/latest.json when it
exists; otherwise it creates a new model from the creation knobs (-n,
--window, --eta, --growth, --subsample, -s, -c, -t, --mungers,
--prototype -- the same set `iforest stream` takes). The model is saved
to a timestamped file in --model-dir every --save-interval seconds
(only when something was learned), on SIGUSR1, on the save command, and
at shutdown; the symlink latest.json is atomically repointed at every
save, so a restart resumes the stream losing at most one interval.

Requests are JSON objects carrying exactly one of "row", "rows", or
"cmd", an optional "mode", and an optional "tag" (any JSON value,
echoed back verbatim in the reply -- a correlation tag, not to be
confused with feature tags):

  {"row": [0.1, 0.7]}                        -> {"score": 0.41, "label": 0}
  {"row": {"cpu": 0.1, "mem": 0.7}}          -> {"score": 0.41, "label": 0}
  {"rows": [[...], {...}], "tag": "b7"}      -> {"scores": [[0.41,0], ...], "tag": "b7"}
  {"rows": [[...]], "mode": "learn"}         -> {"ok": {"learned": 1}}
  {"cmd": "mode", "mode": "score"}           -> {"ok": {"mode": "score"}}
  {"cmd": "ping"}                            -> {"ok": "pong"}
  {"cmd": "stats"}                           -> {"ok": {"seen": ..., ...}}
  {"cmd": "save"}                            -> {"ok": {"saved": "oiforest-....json"}}
  {"cmd": "relearn-threshold"}               -> {"ok": {"threshold": 0.61}}
  anything invalid                           -> {"error": "...", "tag": ...}

The array row form is positional (scalar mungers applied, like stream
CSV input); the object form is a tagged row and runs the full munger
plan, including expanding and combining mungers -- and, being JSON, the
raw values may safely contain commas, newlines, or any unicode.

A worked tagged example.  Create the daemon around raw HTTP request
data, with mungers turning the raw values into numbers (mungers.json
here; a --prototype carrying the same schema works identically):

  { "method":       { "munger": "http_method_enum", "default": -1 },
    "path_len":     { "munger": "length",  "from": "path" },
    "host_entropy": { "munger": "entropy", "from": "host" } }

  iforest streamd --set web -t method -t path_len -t host_entropy \
      --mungers mungers.json -c 0.05

Clients then send the raw values themselves -- note the input fields
are the munger SOURCES (method, path, host), not the feature tags,
because the plan derives path_len and host_entropy from them:

  -> {"row": {"method": "GET", "path": "/index.html",
      "host": "www.example.com"}, "tag": "r-1"}
  <- {"score": 0.31, "label": 0, "tag": "r-1"}
  -> {"row": {"method": "BREW", "path": "/aa,a\"a.php",
      "host": "kq3xv9z2.biz"}, "tag": "r-2"}
  <- {"score": 0.74, "label": 1, "tag": "r-2"}

The same rows work from the shell via
`iforest streamc --set web --jsonl -i rows.jsonl`.

Modes are prequential (score each row against the model as it stood, then
learn it -- the default), learn (learn only), and score (score only);
"mode" on a row/rows message overrides the connection default set by
the mode command for that message.  A bad row gets an {"error": ...}
reply on that message only; the connection and the daemon live on (for
a "rows" batch, rows before the failing one were already processed).

Multiple concurrent connections are supported; rows are applied to the
one shared model in the order their lines arrive, which defines the
stream order.

--set NAME runs a named instance: the set name is appended to
--model-dir (so its saves, latest.json, and default log live under
their own subdirectory) and the socket/pid become <set>.sock /
<set>.pid under the run dir -- with --set, --socket and --pid name the
base run dir instead of the files.  Several sets run side by side, each
with its own model, resume state, and double-start protection:

  iforest streamd --set web
  iforest streamd --set dns --prototype dns-proto.json -c 0.02

Set names must match /\A[A-Za-z0-9+\-@_]+\z/; since the class has no
"." or "/", a set name can only ever create one new path segment.

Everything under --model-dir and the socket/pid directories is created
at startup when missing; when that fails (e.g. running unprivileged
with the /var defaults) the daemon dies immediately, before forking,
naming the directory and the flag to override.
';
} ## end sub description

sub validate {
	my ( $self, $opt, $args ) = @_;

lib/Algorithm/Classifier/IsolationForest/App/Command/streamd.pm  view on Meta::CPAN

		my $mode = $msg->{mode};
		if ( !defined $mode || ref $mode || $mode !~ /\A(?:prequential|learn|score)\z/ ) {
			return _reply( $c, { error => q{mode must be "prequential", "learn", or "score"}, @$tag } );
		}
		$c->{mode} = $mode;
		return _reply( $c, { ok => { mode => $mode }, @$tag } );
	}
	if ( $cmd eq 'stats' ) {
		return _reply(
			$c,
			{
				ok => {
					seen        => 0 + $OIF->seen,
					window      => 0 + $OIF->window_count,
					n_features  => ( defined $OIF->{n_features} ? 0 + $OIF->{n_features} : undef ),
					threshold   => 0 + _threshold(),
					connections => 0 + scalar( keys %CONN ),
					dirty       => ( $DIRTY ? 1 : 0 ),
					set         => $OPT{'set'},
				},
				@$tag
			}
		);
	} ## end if ( $cmd eq 'stats' )
	if ( $cmd eq 'save' ) {
		my $name = eval { _save_model('command') };
		if ($@) {
			( my $err = $@ ) =~ s/ at \S+ line \d+\.?\s*\z//s;
			return _reply( $c, { error => $err, @$tag } );
		}
		return _reply( $c, { ok => { saved => $name }, @$tag } );
	}
	if ( $cmd eq 'relearn-threshold' ) {
		my $ok = eval { $OIF->relearn_threshold; 1 };
		if ( !$ok ) {
			( my $err = $@ ) =~ s/ at \S+ line \d+\.?\s*\z//s;
			chomp $err;
			return _reply( $c, { error => $err, @$tag } );
		}
		return _reply( $c, { ok => { threshold => 0 + $OIF->decision_threshold }, @$tag } );
	}
	return _reply( $c, { error => 'unknown cmd "' . $cmd . '"', @$tag } );
} ## end sub _handle_cmd

sub _reply {
	my ( $c, $reply ) = @_;
	$c->{outbuf} .= $JSON->encode($reply) . "\n";
	return 1;
}

# The effective label cutoff, resolved per message so relearns and the
# contamination refresh at save time take effect immediately.
sub _threshold {
	return
		  defined $OPT{'threshold'}        ? $OPT{'threshold'}
		: defined $OIF->decision_threshold ? $OIF->decision_threshold
		:                                    0.5;
}

# One row through the model.  A JSON object is a tagged row (full munger
# plan, expanding/combining mungers included); a JSON array is positional
# (scalar mungers, like stream CSV input).  Either way the final vector
# is validated numeric before it touches the model -- JSON delivers
# typed values, so anything non-numeric left after munging is a caller
# bug worth an explicit error rather than Perl's silent string-to-0.
# Returns the score, or undef in learn mode.  Croaks on any problem.
sub _apply_row {
	my ( $row, $mode ) = @_;

	my $vec;
	if ( ref $row eq 'HASH' ) {
		$vec = $OIF->tagged_row_to_array( $row, 'streamd' );
	} elsif ( ref $row eq 'ARRAY' ) {
		$vec = $row;
		if ( ref $OIF->{mungers} eq 'HASH' && %{ $OIF->{mungers} } ) {
			$vec = $OIF->munge_rows( [$row] )->[0];
		}
	} else {
		die 'row must be a JSON array (positional) or object (tagged)' . "\n";
	}

	for my $col ( 0 .. $#$vec ) {
		next if !defined $vec->[$col];    # undef defers to the model's missing policy
		die 'column ' . ( $col + 1 ) . ' is not a number after munging' . "\n"
			unless looks_like_number( $vec->[$col] );
	}

	if ( $mode eq 'learn' ) {
		$OIF->learn( [$vec] );
		$DIRTY = 1;
		return undef;
	}
	if ( $mode eq 'score' ) {
		return $OIF->score_samples( [$vec] )->[0];
	}
	my $score = $OIF->score_learn( [$vec] )->[0];
	$DIRTY = 1;
	return $score;
} ## end sub _apply_row

return 1;



( run in 0.387 second using v1.01-cache-2.11-cpan-a9496e3eb41 )