mirror of
https://github.com/zeek/zeek.git
synced 2025-10-10 02:28:21 +00:00
completely change interface again.
compiles, not really tested. basic test works 70% of the time, coredumps in the other 30 - but was not easy to debug on a first glance (most interestingly the crash happens in the logging framework - I wonder how that works). Other tests are not adjusted to the new interface yet.
This commit is contained in:
parent
b4e6971aab
commit
57ffe1be77
14 changed files with 403 additions and 1072 deletions
|
@ -31,14 +31,25 @@ declare(PDict, InputHash);
|
|||
|
||||
class Manager::Filter {
|
||||
public:
|
||||
EnumVal* id;
|
||||
string name;
|
||||
string source;
|
||||
|
||||
int mode;
|
||||
|
||||
FilterType filter_type; // to distinguish between event and table filters
|
||||
|
||||
EnumVal* type;
|
||||
ReaderFrontend* reader;
|
||||
|
||||
virtual ~Filter();
|
||||
};
|
||||
|
||||
Manager::Filter::~Filter() {
|
||||
Unref(type);
|
||||
|
||||
delete(reader);
|
||||
}
|
||||
|
||||
class Manager::TableFilter: public Manager::Filter {
|
||||
public:
|
||||
|
||||
|
@ -85,10 +96,6 @@ Manager::EventFilter::EventFilter() {
|
|||
filter_type = EVENT_FILTER;
|
||||
}
|
||||
|
||||
Manager::Filter::~Filter() {
|
||||
Unref(id);
|
||||
}
|
||||
|
||||
Manager::TableFilter::~TableFilter() {
|
||||
Unref(tab);
|
||||
Unref(itype);
|
||||
|
@ -99,41 +106,6 @@ Manager::TableFilter::~TableFilter() {
|
|||
delete lastDict;
|
||||
}
|
||||
|
||||
struct Manager::ReaderInfo {
|
||||
EnumVal* id;
|
||||
EnumVal* type;
|
||||
ReaderFrontend* reader;
|
||||
|
||||
map<int, Manager::Filter*> filters; // filters that can prevent our actions
|
||||
|
||||
bool HasFilter(int id);
|
||||
|
||||
~ReaderInfo();
|
||||
};
|
||||
|
||||
Manager::ReaderInfo::~ReaderInfo() {
|
||||
map<int, Manager::Filter*>::iterator it = filters.begin();
|
||||
|
||||
while ( it != filters.end() ) {
|
||||
delete (*it).second;
|
||||
++it;
|
||||
}
|
||||
|
||||
Unref(type);
|
||||
Unref(id);
|
||||
|
||||
delete(reader);
|
||||
}
|
||||
|
||||
bool Manager::ReaderInfo::HasFilter(int id) {
|
||||
map<int, Manager::Filter*>::iterator it = filters.find(id);
|
||||
if ( it == filters.end() ) {
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
struct ReaderDefinition {
|
||||
bro_int_t type; // the type
|
||||
const char *name; // descriptive name for error messages
|
||||
|
@ -203,37 +175,34 @@ ReaderBackend* Manager::CreateBackend(ReaderFrontend* frontend, bro_int_t type)
|
|||
}
|
||||
|
||||
// create a new input reader object to be used at whomevers leisure lateron.
|
||||
ReaderFrontend* Manager::CreateStream(EnumVal* id, RecordVal* description)
|
||||
bool Manager::CreateStream(Filter* info, RecordVal* description)
|
||||
{
|
||||
{
|
||||
ReaderInfo *i = FindReader(id);
|
||||
if ( i != 0 ) {
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
reporter->Error("Trying create already existing input stream %s", desc.Description());
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
ReaderDefinition* ir = input_readers;
|
||||
|
||||
RecordType* rtype = description->Type()->AsRecordType();
|
||||
if ( ! same_type(rtype, BifType::Record::Input::StreamDescription, 0) )
|
||||
if ( ! ( same_type(rtype, BifType::Record::Input::TableDescription, 0) || same_type(rtype, BifType::Record::Input::EventDescription, 0) ) )
|
||||
{
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
reporter->Error("Streamdescription argument not of right type for new input stream %s", desc.Description());
|
||||
return 0;
|
||||
reporter->Error("Streamdescription argument not of right type for new input stream");
|
||||
return false;
|
||||
}
|
||||
|
||||
Val* name_val = description->LookupWithDefault(rtype->FieldOffset("name"));
|
||||
string name = name_val->AsString()->CheckString();
|
||||
Unref(name_val);
|
||||
|
||||
{
|
||||
Filter *i = FindFilter(name);
|
||||
if ( i != 0 ) {
|
||||
reporter->Error("Trying create already existing input stream %s", name.c_str());
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
EnumVal* reader = description->LookupWithDefault(rtype->FieldOffset("reader"))->AsEnumVal();
|
||||
EnumVal* mode = description->LookupWithDefault(rtype->FieldOffset("mode"))->AsEnumVal();
|
||||
Val *autostart = description->LookupWithDefault(rtype->FieldOffset("autostart"));
|
||||
bool do_autostart = ( autostart->InternalInt() == 1 );
|
||||
Unref(autostart); // Ref'd by LookupWithDefault
|
||||
|
||||
ReaderFrontend* reader_obj = new ReaderFrontend(reader->InternalInt());
|
||||
assert(reader_obj);
|
||||
|
||||
ReaderFrontend* reader_obj = new ReaderFrontend(reader->InternalInt());
|
||||
assert(reader_obj);
|
||||
|
||||
// get the source...
|
||||
Val* sourceval = description->LookupWithDefault(rtype->FieldOffset("source"));
|
||||
|
@ -241,42 +210,45 @@ ReaderFrontend* Manager::CreateStream(EnumVal* id, RecordVal* description)
|
|||
const BroString* bsource = sourceval->AsString();
|
||||
string source((const char*) bsource->Bytes(), bsource->Len());
|
||||
Unref(sourceval);
|
||||
|
||||
EnumVal* mode = description->LookupWithDefault(rtype->FieldOffset("mode"))->AsEnumVal();
|
||||
info->mode = mode->InternalInt();
|
||||
Unref(mode);
|
||||
|
||||
ReaderInfo* info = new ReaderInfo;
|
||||
info->reader = reader_obj;
|
||||
info->type = reader->AsEnumVal(); // ref'd by lookupwithdefault
|
||||
info->id = id->Ref()->AsEnumVal();
|
||||
info->name = name;
|
||||
info->source = source;
|
||||
|
||||
readers.push_back(info);
|
||||
|
||||
reader_obj->Init(source, mode->InternalInt(), do_autostart);
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Successfully created new input stream %s",
|
||||
desc.Description());
|
||||
name.c_str());
|
||||
#endif
|
||||
|
||||
return reader_obj;
|
||||
return true;
|
||||
|
||||
}
|
||||
|
||||
bool Manager::AddEventFilter(EnumVal *id, RecordVal* fval) {
|
||||
ReaderInfo *i = FindReader(id);
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Stream not found");
|
||||
return false;
|
||||
}
|
||||
bool Manager::CreateEventStream(RecordVal* fval) {
|
||||
|
||||
RecordType* rtype = fval->Type()->AsRecordType();
|
||||
if ( ! same_type(rtype, BifType::Record::Input::EventFilter, 0) )
|
||||
if ( ! same_type(rtype, BifType::Record::Input::EventDescription, 0) )
|
||||
{
|
||||
reporter->Error("filter argument not of right type");
|
||||
return false;
|
||||
}
|
||||
|
||||
EventFilter* filter = new EventFilter();
|
||||
{
|
||||
bool res = CreateStream(filter, fval);
|
||||
if ( res == false ) {
|
||||
delete filter;
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Val* name = fval->LookupWithDefault(rtype->FieldOffset("name"));
|
||||
RecordType *fields = fval->LookupWithDefault(rtype->FieldOffset("fields"))->AsType()->AsTypeType()->Type()->AsRecordType();
|
||||
|
||||
Val *want_record = fval->LookupWithDefault(rtype->FieldOffset("want_record"));
|
||||
|
@ -352,42 +324,36 @@ bool Manager::AddEventFilter(EnumVal *id, RecordVal* fval) {
|
|||
}
|
||||
|
||||
Unref(fields); // ref'd by lookupwithdefault
|
||||
EventFilter* filter = new EventFilter();
|
||||
filter->name = name->AsString()->CheckString();
|
||||
Unref(name); // ref'd by lookupwithdefault
|
||||
filter->id = id->Ref()->AsEnumVal();
|
||||
filter->num_fields = fieldsV.size();
|
||||
filter->fields = fields->Ref()->AsRecordType();
|
||||
filter->event = event_registry->Lookup(event->GetID()->Name());
|
||||
filter->want_record = ( want_record->InternalInt() == 1 );
|
||||
Unref(want_record); // ref'd by lookupwithdefault
|
||||
|
||||
int filterid = 0;
|
||||
if ( i->filters.size() > 0 ) {
|
||||
filterid = i->filters.rbegin()->first + 1; // largest element is at beginning of map-> new id = old id + 1->
|
||||
}
|
||||
i->filters[filterid] = filter;
|
||||
i->reader->AddFilter( filterid, fieldsV.size(), logf );
|
||||
assert(filter->reader);
|
||||
filter->reader->Init(filter->source, filter->mode, filter->num_fields, logf );
|
||||
|
||||
readers[filter->reader] = filter;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool Manager::AddTableFilter(EnumVal *id, RecordVal* fval) {
|
||||
ReaderInfo *i = FindReader(id);
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Stream not found");
|
||||
return false;
|
||||
}
|
||||
|
||||
bool Manager::CreateTableStream(RecordVal* fval) {
|
||||
RecordType* rtype = fval->Type()->AsRecordType();
|
||||
if ( ! same_type(rtype, BifType::Record::Input::TableFilter, 0) )
|
||||
if ( ! same_type(rtype, BifType::Record::Input::TableDescription, 0) )
|
||||
{
|
||||
reporter->Error("filter argument not of right type");
|
||||
return false;
|
||||
}
|
||||
|
||||
TableFilter* filter = new TableFilter();
|
||||
{
|
||||
bool res = CreateStream(filter, fval);
|
||||
if ( res == false ) {
|
||||
delete filter;
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
Val* name = fval->LookupWithDefault(rtype->FieldOffset("name"));
|
||||
Val* pred = fval->LookupWithDefault(rtype->FieldOffset("pred"));
|
||||
|
||||
RecordType *idx = fval->LookupWithDefault(rtype->FieldOffset("idx"))->AsType()->AsTypeType()->Type()->AsRecordType();
|
||||
|
@ -493,9 +459,6 @@ bool Manager::AddTableFilter(EnumVal *id, RecordVal* fval) {
|
|||
fields[i] = fieldsV[i];
|
||||
}
|
||||
|
||||
TableFilter* filter = new TableFilter();
|
||||
filter->name = name->AsString()->CheckString();
|
||||
filter->id = id->Ref()->AsEnumVal();
|
||||
filter->pred = pred ? pred->AsFunc() : 0;
|
||||
filter->num_idx_fields = idxfields;
|
||||
filter->num_val_fields = valfields;
|
||||
|
@ -508,25 +471,22 @@ bool Manager::AddTableFilter(EnumVal *id, RecordVal* fval) {
|
|||
filter->want_record = ( want_record->InternalInt() == 1 );
|
||||
|
||||
Unref(want_record); // ref'd by lookupwithdefault
|
||||
Unref(name);
|
||||
Unref(pred);
|
||||
|
||||
if ( valfields > 1 ) {
|
||||
assert(filter->want_record);
|
||||
}
|
||||
|
||||
|
||||
assert(filter->reader);
|
||||
filter->reader->Init(filter->source, filter->mode, fieldsV.size(), fields );
|
||||
|
||||
readers[filter->reader] = filter;
|
||||
|
||||
int filterid = 0;
|
||||
if ( i->filters.size() > 0 ) {
|
||||
filterid = i->filters.rbegin()->first + 1; // largest element is at beginning of map-> new id = old id + 1->
|
||||
}
|
||||
i->filters[filterid] = filter;
|
||||
i->reader->AddFilter( filterid, fieldsV.size(), fields );
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Successfully created new table filter %s for stream",
|
||||
filter->name.c_str(), desc.Description());
|
||||
DBG_LOG(DBG_INPUT, "Successfully created table stream %s",
|
||||
filter->name.c_str());
|
||||
#endif
|
||||
|
||||
return true;
|
||||
|
@ -583,16 +543,8 @@ bool Manager::IsCompatibleType(BroType* t, bool atomic_only)
|
|||
}
|
||||
|
||||
|
||||
bool Manager::RemoveStream(const EnumVal* id) {
|
||||
ReaderInfo *i = 0;
|
||||
for ( vector<ReaderInfo *>::iterator s = readers.begin(); s != readers.end(); ++s )
|
||||
{
|
||||
if ( (*s)->id == id )
|
||||
{
|
||||
i = (*s);
|
||||
break;
|
||||
}
|
||||
}
|
||||
bool Manager::RemoveStream(const string &name) {
|
||||
Filter *i = FindFilter(name);
|
||||
|
||||
if ( i == 0 ) {
|
||||
return false; // not found
|
||||
|
@ -601,38 +553,29 @@ bool Manager::RemoveStream(const EnumVal* id) {
|
|||
i->reader->Finish();
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Successfully queued removal of stream %s",
|
||||
desc.Description());
|
||||
name.c_str());
|
||||
#endif
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
bool Manager::RemoveStreamContinuation(const ReaderFrontend* reader) {
|
||||
ReaderInfo *i = 0;
|
||||
bool Manager::RemoveStreamContinuation(ReaderFrontend* reader) {
|
||||
Filter *i = FindFilter(reader);
|
||||
|
||||
|
||||
for ( vector<ReaderInfo *>::iterator s = readers.begin(); s != readers.end(); ++s )
|
||||
{
|
||||
if ( (*s)->reader && (*s)->reader == reader )
|
||||
{
|
||||
i = *s;
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
i->id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Successfully executed removal of stream %s",
|
||||
desc.Description());
|
||||
#endif
|
||||
delete(i);
|
||||
readers.erase(s);
|
||||
return true;
|
||||
}
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Stream not found in RemoveStreamContinuation");
|
||||
return false;
|
||||
}
|
||||
|
||||
reporter->Error("Stream not found in RemoveStreamContinuation");
|
||||
return false;
|
||||
|
||||
|
||||
#ifdef DEBUG
|
||||
DBG_LOG(DBG_INPUT, "Successfully executed removal of stream %s",
|
||||
i->name.c_str());
|
||||
#endif
|
||||
readers.erase(reader);
|
||||
delete(i);
|
||||
return true;
|
||||
}
|
||||
|
||||
bool Manager::UnrollRecordType(vector<Field*> *fields, const RecordType *rec, const string& nameprepend) {
|
||||
|
@ -680,132 +623,24 @@ bool Manager::UnrollRecordType(vector<Field*> *fields, const RecordType *rec, co
|
|||
return true;
|
||||
}
|
||||
|
||||
bool Manager::ForceUpdate(const EnumVal* id)
|
||||
bool Manager::ForceUpdate(const string &name)
|
||||
{
|
||||
ReaderInfo *i = FindReader(id);
|
||||
Filter *i = FindFilter(name);
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Reader not found");
|
||||
reporter->Error("Stream %s not found", name.c_str());
|
||||
return false;
|
||||
}
|
||||
|
||||
i->reader->Update();
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Forcing update of stream %s",
|
||||
desc.Description());
|
||||
name.c_str());
|
||||
#endif
|
||||
|
||||
return true; // update is async :(
|
||||
}
|
||||
|
||||
bool Manager::RemoveTableFilter(EnumVal* id, const string &name) {
|
||||
ReaderInfo *i = FindReader(id);
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Reader not found");
|
||||
return false;
|
||||
}
|
||||
|
||||
bool found = false;
|
||||
int filterId;
|
||||
|
||||
for ( map<int, Manager::Filter*>::iterator it = i->filters.begin(); it != i->filters.end(); ++it ) {
|
||||
if ( (*it).second->name == name ) {
|
||||
found = true;
|
||||
filterId = (*it).first;
|
||||
|
||||
if ( (*it).second->filter_type != TABLE_FILTER ) {
|
||||
reporter->Error("Trying to remove filter %s of wrong type", name.c_str());
|
||||
return false;
|
||||
}
|
||||
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if ( !found ) {
|
||||
reporter->Error("Trying to remove nonexisting filter %s", name.c_str());
|
||||
return false;
|
||||
}
|
||||
|
||||
i->reader->RemoveFilter(filterId);
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Queued removal of tablefilter %s for stream %s",
|
||||
name.c_str(), desc.Description());
|
||||
#endif
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
bool Manager::RemoveFilterContinuation(const ReaderFrontend* reader, const int filterId) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Reader not found");
|
||||
return false;
|
||||
}
|
||||
|
||||
map<int, Manager::Filter*>::iterator it = i->filters.find(filterId);
|
||||
if ( it == i->filters.end() ) {
|
||||
reporter->Error("Got RemoveFilterContinuation where filter nonexistant for %d", filterId);
|
||||
return false;
|
||||
}
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
i->id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Executed removal of (table|event)-filter %s for stream %s",
|
||||
(*it).second->name.c_str(), desc.Description());
|
||||
#endif
|
||||
|
||||
delete (*it).second;
|
||||
i->filters.erase(it);
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
bool Manager::RemoveEventFilter(EnumVal* id, const string &name) {
|
||||
ReaderInfo *i = FindReader(id);
|
||||
if ( i == 0 ) {
|
||||
reporter->Error("Reader not found");
|
||||
return false;
|
||||
}
|
||||
|
||||
bool found = false;
|
||||
int filterId;
|
||||
for ( map<int, Manager::Filter*>::iterator it = i->filters.begin(); it != i->filters.end(); ++it ) {
|
||||
if ( (*it).second->name == name ) {
|
||||
found = true;
|
||||
filterId = (*it).first;
|
||||
|
||||
if ( (*it).second->filter_type != EVENT_FILTER ) {
|
||||
reporter->Error("Trying to remove filter %s of wrong type", name.c_str());
|
||||
return false;
|
||||
}
|
||||
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if ( !found ) {
|
||||
reporter->Error("Trying to remove nonexisting filter %s", name.c_str());
|
||||
return false;
|
||||
}
|
||||
|
||||
i->reader->RemoveFilter(filterId);
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Queued removal of eventfilter %s for stream %s",
|
||||
name.c_str(), desc.Description());
|
||||
#endif
|
||||
return true;
|
||||
}
|
||||
|
||||
Val* Manager::ValueToIndexVal(int num_fields, const RecordType *type, const Value* const *vals) {
|
||||
Val* idxval;
|
||||
int position = 0;
|
||||
|
@ -833,24 +668,19 @@ Val* Manager::ValueToIndexVal(int num_fields, const RecordType *type, const Valu
|
|||
}
|
||||
|
||||
|
||||
void Manager::SendEntry(const ReaderFrontend* reader, const int id, Value* *vals) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
void Manager::SendEntry(ReaderFrontend* reader, Value* *vals) {
|
||||
Filter *i = FindFilter(reader);
|
||||
if ( i == 0 ) {
|
||||
reporter->InternalError("Unknown reader");
|
||||
return;
|
||||
}
|
||||
|
||||
if ( !i->HasFilter(id) ) {
|
||||
reporter->InternalError("Unknown filter");
|
||||
reporter->InternalError("Unknown reader in SendEntry");
|
||||
return;
|
||||
}
|
||||
|
||||
int readFields;
|
||||
if ( i->filters[id]->filter_type == TABLE_FILTER ) {
|
||||
readFields = SendEntryTable(reader, id, vals);
|
||||
} else if ( i->filters[id]->filter_type == EVENT_FILTER ) {
|
||||
if ( i->filter_type == TABLE_FILTER ) {
|
||||
readFields = SendEntryTable(i, vals);
|
||||
} else if ( i->filter_type == EVENT_FILTER ) {
|
||||
EnumVal *type = new EnumVal(BifEnum::Input::EVENT_NEW, BifType::Enum::Input::Event);
|
||||
readFields = SendEventFilterEvent(reader, type, id, vals);
|
||||
readFields = SendEventFilterEvent(i, type, vals);
|
||||
} else {
|
||||
assert(false);
|
||||
}
|
||||
|
@ -861,16 +691,13 @@ void Manager::SendEntry(const ReaderFrontend* reader, const int id, Value* *vals
|
|||
delete [] vals;
|
||||
}
|
||||
|
||||
int Manager::SendEntryTable(const ReaderFrontend* reader, const int id, const Value* const *vals) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
|
||||
int Manager::SendEntryTable(Filter* i, const Value* const *vals) {
|
||||
bool updated = false;
|
||||
|
||||
assert(i);
|
||||
assert(i->HasFilter(id));
|
||||
|
||||
assert(i->filters[id]->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i->filters[id];
|
||||
assert(i->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i;
|
||||
|
||||
HashKey* idxhash = HashValues(filter->num_idx_fields, vals);
|
||||
|
||||
|
@ -1006,29 +833,25 @@ int Manager::SendEntryTable(const ReaderFrontend* reader, const int id, const Va
|
|||
}
|
||||
|
||||
|
||||
void Manager::EndCurrentSend(const ReaderFrontend* reader, int id) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
void Manager::EndCurrentSend(ReaderFrontend* reader) {
|
||||
Filter *i = FindFilter(reader);
|
||||
if ( i == 0 ) {
|
||||
reporter->InternalError("Unknown reader");
|
||||
reporter->InternalError("Unknown reader in EndCurrentSend");
|
||||
return;
|
||||
}
|
||||
|
||||
assert(i->HasFilter(id));
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
i->id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Got EndCurrentSend for filter %d and stream %s",
|
||||
id, desc.Description());
|
||||
DBG_LOG(DBG_INPUT, "Got EndCurrentSend stream %s",
|
||||
i->name.c_str());
|
||||
#endif
|
||||
|
||||
if ( i->filters[id]->filter_type == EVENT_FILTER ) {
|
||||
if ( i->filter_type == EVENT_FILTER ) {
|
||||
// nothing to do..
|
||||
return;
|
||||
}
|
||||
|
||||
assert(i->filters[id]->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i->filters[id];
|
||||
assert(i->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i;
|
||||
|
||||
// lastdict contains all deleted entries and should be empty apart from that
|
||||
IterCookie *c = filter->lastDict->InitForIteration();
|
||||
|
@ -1093,8 +916,8 @@ void Manager::EndCurrentSend(const ReaderFrontend* reader, int id) {
|
|||
filter->currDict = new PDict(InputHash);
|
||||
|
||||
#ifdef DEBUG
|
||||
DBG_LOG(DBG_INPUT, "EndCurrentSend complete for filter %d and stream %s, queueing update_finished event",
|
||||
id, desc.Description());
|
||||
DBG_LOG(DBG_INPUT, "EndCurrentSend complete for stream %s, queueing update_finished event",
|
||||
i->name.c_str());
|
||||
#endif
|
||||
|
||||
// Send event that the current update is indeed finished.
|
||||
|
@ -1104,42 +927,40 @@ void Manager::EndCurrentSend(const ReaderFrontend* reader, int id) {
|
|||
}
|
||||
|
||||
|
||||
Ref(i->id);
|
||||
SendEvent(handler, 1, i->id);
|
||||
SendEvent(handler, 1, new BroString(i->name));
|
||||
}
|
||||
|
||||
void Manager::Put(const ReaderFrontend* reader, int id, Value* *vals) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
void Manager::Put(ReaderFrontend* reader, Value* *vals) {
|
||||
Filter *i = FindFilter(reader);
|
||||
if ( i == 0 ) {
|
||||
reporter->InternalError("Unknown reader");
|
||||
reporter->InternalError("Unknown reader in Put");
|
||||
return;
|
||||
}
|
||||
|
||||
if ( !i->HasFilter(id) ) {
|
||||
reporter->InternalError("Unknown filter");
|
||||
return;
|
||||
}
|
||||
|
||||
if ( i->filters[id]->filter_type == TABLE_FILTER ) {
|
||||
PutTable(reader, id, vals);
|
||||
} else if ( i->filters[id]->filter_type == EVENT_FILTER ) {
|
||||
int readFields;
|
||||
if ( i->filter_type == TABLE_FILTER ) {
|
||||
readFields = PutTable(i, vals);
|
||||
} else if ( i->filter_type == EVENT_FILTER ) {
|
||||
EnumVal *type = new EnumVal(BifEnum::Input::EVENT_NEW, BifType::Enum::Input::Event);
|
||||
SendEventFilterEvent(reader, type, id, vals);
|
||||
readFields = SendEventFilterEvent(i, type, vals);
|
||||
} else {
|
||||
assert(false);
|
||||
}
|
||||
|
||||
for ( int i = 0; i < readFields; i++ ) {
|
||||
delete vals[i];
|
||||
}
|
||||
delete [] vals;
|
||||
|
||||
}
|
||||
|
||||
int Manager::SendEventFilterEvent(const ReaderFrontend* reader, EnumVal* type, int id, const Value* const *vals) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
|
||||
int Manager::SendEventFilterEvent(Filter* i, EnumVal* type, const Value* const *vals) {
|
||||
bool updated = false;
|
||||
|
||||
assert(i);
|
||||
assert(i->HasFilter(id));
|
||||
|
||||
assert(i->filters[id]->filter_type == EVENT_FILTER);
|
||||
EventFilter* filter = (EventFilter*) i->filters[id];
|
||||
assert(i->filter_type == EVENT_FILTER);
|
||||
EventFilter* filter = (EventFilter*) i;
|
||||
|
||||
Val *val;
|
||||
list<Val*> out_vals;
|
||||
|
@ -1170,14 +991,11 @@ int Manager::SendEventFilterEvent(const ReaderFrontend* reader, EnumVal* type, i
|
|||
|
||||
}
|
||||
|
||||
int Manager::PutTable(const ReaderFrontend* reader, int id, const Value* const *vals) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
|
||||
int Manager::PutTable(Filter* i, const Value* const *vals) {
|
||||
assert(i);
|
||||
assert(i->HasFilter(id));
|
||||
|
||||
assert(i->filters[id]->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i->filters[id];
|
||||
assert(i->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i;
|
||||
|
||||
Val* idxval = ValueToIndexVal(filter->num_idx_fields, filter->itype, vals);
|
||||
Val* valval;
|
||||
|
@ -1274,43 +1092,37 @@ int Manager::PutTable(const ReaderFrontend* reader, int id, const Value* const *
|
|||
}
|
||||
|
||||
// Todo:: perhaps throw some kind of clear-event?
|
||||
void Manager::Clear(const ReaderFrontend* reader, int id) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
void Manager::Clear(ReaderFrontend* reader) {
|
||||
Filter *i = FindFilter(reader);
|
||||
if ( i == 0 ) {
|
||||
reporter->InternalError("Unknown reader");
|
||||
reporter->InternalError("Unknown reader in Clear");
|
||||
return;
|
||||
}
|
||||
|
||||
#ifdef DEBUG
|
||||
ODesc desc;
|
||||
i->id->Describe(&desc);
|
||||
DBG_LOG(DBG_INPUT, "Got Clear for filter %d and stream %s",
|
||||
id, desc.Description());
|
||||
DBG_LOG(DBG_INPUT, "Got Clear for stream %s",
|
||||
i->name.c_str());
|
||||
#endif
|
||||
|
||||
assert(i->HasFilter(id));
|
||||
|
||||
assert(i->filters[id]->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i->filters[id];
|
||||
assert(i->filter_type == TABLE_FILTER);
|
||||
TableFilter* filter = (TableFilter*) i;
|
||||
|
||||
filter->tab->RemoveAll();
|
||||
}
|
||||
|
||||
// put interface: delete old entry from table.
|
||||
bool Manager::Delete(const ReaderFrontend* reader, int id, Value* *vals) {
|
||||
ReaderInfo *i = FindReader(reader);
|
||||
bool Manager::Delete(ReaderFrontend* reader, Value* *vals) {
|
||||
Filter *i = FindFilter(reader);
|
||||
if ( i == 0 ) {
|
||||
reporter->InternalError("Unknown reader");
|
||||
reporter->InternalError("Unknown reader in Delete");
|
||||
return false;
|
||||
}
|
||||
|
||||
assert(i->HasFilter(id));
|
||||
|
||||
bool success = false;
|
||||
int readVals = 0;
|
||||
|
||||
if ( i->filters[id]->filter_type == TABLE_FILTER ) {
|
||||
TableFilter* filter = (TableFilter*) i->filters[id];
|
||||
if ( i->filter_type == TABLE_FILTER ) {
|
||||
TableFilter* filter = (TableFilter*) i;
|
||||
Val* idxval = ValueToIndexVal(filter->num_idx_fields, filter->itype, vals);
|
||||
assert(idxval != 0);
|
||||
readVals = filter->num_idx_fields + filter->num_val_fields;
|
||||
|
@ -1352,9 +1164,9 @@ bool Manager::Delete(const ReaderFrontend* reader, int id, Value* *vals) {
|
|||
reporter->Error("Internal error while deleting values from input table");
|
||||
}
|
||||
}
|
||||
} else if ( i->filters[id]->filter_type == EVENT_FILTER ) {
|
||||
} else if ( i->filter_type == EVENT_FILTER ) {
|
||||
EnumVal *type = new EnumVal(BifEnum::Input::EVENT_REMOVED, BifType::Enum::Input::Event);
|
||||
readVals = SendEventFilterEvent(reader, type, id, vals);
|
||||
readVals = SendEventFilterEvent(i, type, vals);
|
||||
success = true;
|
||||
} else {
|
||||
assert(false);
|
||||
|
@ -1867,29 +1679,24 @@ Val* Manager::ValueToVal(const Value* val, BroType* request_type) {
|
|||
return NULL;
|
||||
}
|
||||
|
||||
Manager::ReaderInfo* Manager::FindReader(const ReaderFrontend* reader)
|
||||
Manager::Filter* Manager::FindFilter(const string &name)
|
||||
{
|
||||
for ( vector<ReaderInfo *>::iterator s = readers.begin(); s != readers.end(); ++s )
|
||||
for ( map<ReaderFrontend*, Filter*>::iterator s = readers.begin(); s != readers.end(); ++s )
|
||||
{
|
||||
if ( (*s)->reader && (*s)->reader == reader )
|
||||
if ( (*s).second->name == name )
|
||||
{
|
||||
return *s;
|
||||
return (*s).second;
|
||||
}
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
Manager::ReaderInfo* Manager::FindReader(const EnumVal* id)
|
||||
{
|
||||
for ( vector<ReaderInfo *>::iterator s = readers.begin(); s != readers.end(); ++s )
|
||||
{
|
||||
if ( (*s)->id && (*s)->id->AsEnum() == id->AsEnum() )
|
||||
{
|
||||
return *s;
|
||||
}
|
||||
}
|
||||
|
||||
return 0;
|
||||
Manager::Filter* Manager::FindFilter(ReaderFrontend* reader)
|
||||
{
|
||||
map<ReaderFrontend*, Filter*>::iterator s = readers.find(reader);
|
||||
if ( s != readers.end() ) {
|
||||
return s->second;
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
|
Loading…
Add table
Add a link
Reference in a new issue