This repository has been archived by the owner on Mar 21, 2023. It is now read-only.
add remove_from_default boolean option to route_to_stream function #220
Merged
Changes from all commits
Commits
Show all changes
2 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Jump to
Jump to file
Failed to load files.
Diff view
Diff view
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
105 changes: 105 additions & 0 deletions
105
.../main/java/org/graylog/plugins/pipelineprocessor/functions/messages/RemoveFromStream.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,105 @@ | ||
/** | ||
* This file is part of Graylog Pipeline Processor. | ||
* | ||
* Graylog Pipeline Processor is free software: you can redistribute it and/or modify | ||
* it under the terms of the GNU General Public License as published by | ||
* the Free Software Foundation, either version 3 of the License, or | ||
* (at your option) any later version. | ||
* | ||
* Graylog Pipeline Processor is distributed in the hope that it will be useful, | ||
* but WITHOUT ANY WARRANTY; without even the implied warranty of | ||
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | ||
* GNU General Public License for more details. | ||
* | ||
* You should have received a copy of the GNU General Public License | ||
* along with Graylog Pipeline Processor. If not, see <http://www.gnu.org/licenses/>. | ||
*/ | ||
package org.graylog.plugins.pipelineprocessor.functions.messages; | ||
|
||
import com.google.inject.Inject; | ||
import org.graylog.plugins.pipelineprocessor.EvaluationContext; | ||
import org.graylog.plugins.pipelineprocessor.ast.functions.AbstractFunction; | ||
import org.graylog.plugins.pipelineprocessor.ast.functions.FunctionArgs; | ||
import org.graylog.plugins.pipelineprocessor.ast.functions.FunctionDescriptor; | ||
import org.graylog.plugins.pipelineprocessor.ast.functions.ParameterDescriptor; | ||
import org.graylog2.plugin.Message; | ||
import org.graylog2.plugin.streams.DefaultStream; | ||
import org.graylog2.plugin.streams.Stream; | ||
|
||
import javax.inject.Provider; | ||
import java.util.Collection; | ||
import java.util.Collections; | ||
import java.util.Optional; | ||
|
||
import static com.google.common.collect.ImmutableList.of; | ||
import static org.graylog.plugins.pipelineprocessor.ast.functions.ParameterDescriptor.string; | ||
import static org.graylog.plugins.pipelineprocessor.ast.functions.ParameterDescriptor.type; | ||
|
||
public class RemoveFromStream extends AbstractFunction<Void> { | ||
|
||
public static final String NAME = "remove_from_stream"; | ||
private static final String ID_ARG = "id"; | ||
private static final String NAME_ARG = "name"; | ||
private final StreamCacheService streamCacheService; | ||
private final Provider<Stream> defaultStreamProvider; | ||
private final ParameterDescriptor<Message, Message> messageParam; | ||
private final ParameterDescriptor<String, String> nameParam; | ||
private final ParameterDescriptor<String, String> idParam; | ||
|
||
@Inject | ||
public RemoveFromStream(StreamCacheService streamCacheService, @DefaultStream Provider<Stream> defaultStreamProvider) { | ||
this.streamCacheService = streamCacheService; | ||
this.defaultStreamProvider = defaultStreamProvider; | ||
|
||
messageParam = type("message", Message.class).optional().description("The message to use, defaults to '$message'").build(); | ||
nameParam = string(NAME_ARG).optional().description("The name of the stream to remove the message from, must match exactly").build(); | ||
idParam = string(ID_ARG).optional().description("The ID of the stream").build(); | ||
} | ||
|
||
@Override | ||
public Void evaluate(FunctionArgs args, EvaluationContext context) { | ||
Optional<String> id = idParam.optional(args, context); | ||
|
||
Collection<Stream> streams; | ||
if (!id.isPresent()) { | ||
final Optional<Collection<Stream>> foundStreams = nameParam.optional(args, context).map(streamCacheService::getByName); | ||
|
||
if (!foundStreams.isPresent()) { | ||
// TODO signal error somehow | ||
return null; | ||
} else { | ||
streams = foundStreams.get(); | ||
} | ||
} else { | ||
final Stream stream = streamCacheService.getById(id.get()); | ||
if (stream == null) { | ||
return null; | ||
} | ||
streams = Collections.singleton(stream); | ||
} | ||
final Message message = messageParam.optional(args, context).orElse(context.currentMessage()); | ||
streams.forEach(stream -> { | ||
if (!stream.isPaused()) { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Why is it important whether a stream is paused or not? There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. It's not particularly, but paused streams aren't taken into account when adding either. |
||
message.removeStream(stream); | ||
} | ||
}); | ||
// always leave a message at least on the default stream if we removed the last stream it was on | ||
if (message.getStreams().isEmpty()) { | ||
message.addStream(defaultStreamProvider.get()); | ||
} | ||
return null; | ||
} | ||
|
||
@Override | ||
public FunctionDescriptor<Void> descriptor() { | ||
return FunctionDescriptor.<Void>builder() | ||
.name(NAME) | ||
.returnType(Void.class) | ||
.params(of( | ||
nameParam, | ||
idParam, | ||
messageParam)) | ||
.description("Removes a message from a stream. Removing the last stream will put the message back onto the default stream. To complete drop a message use the drop_message function.") | ||
.build(); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
5 changes: 5 additions & 0 deletions
5
...n/src/test/resources/org/graylog/plugins/pipelineprocessor/functions/removeFromStream.txt
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,5 @@ | ||
rule "stream routing" | ||
when true | ||
then | ||
remove_from_stream(name: "some name"); | ||
end |
7 changes: 7 additions & 0 deletions
7
...sources/org/graylog/plugins/pipelineprocessor/functions/removeFromStreamRetainDefault.txt
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,7 @@ | ||
rule "stream routing" | ||
when true | ||
then | ||
remove_from_stream(name: "some name"); | ||
// if a message is taken off all stream it was on, the default stream will be added back to avoid dropping the message | ||
remove_from_stream(id: "000000000000000000000001"); | ||
end |
5 changes: 5 additions & 0 deletions
5
.../resources/org/graylog/plugins/pipelineprocessor/functions/routeToStreamRemoveDefault.txt
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,5 @@ | ||
rule "stream routing" | ||
when true | ||
then | ||
route_to_stream(name: "some name", remove_from_default: true); | ||
end |
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Uh oh, that's not a good sign.
That means someone forgot to add these functions…
Refs #190