Merge pull request #1459 from satti-hari-krishna-reddy/new-sigma

sigma detection ui
This commit is contained in:
Frikky
2024-07-31 19:34:02 +02:00
committed by GitHub
9 changed files with 1482 additions and 294 deletions
+1
View File
@@ -40,6 +40,7 @@ BACKEND_HOSTNAME=shuffle-backend
BACKEND_PORT=5001
FRONTEND_PORT=3001
FRONTEND_PORT_HTTPS=3443
AUTH_FOR_ORBORUS =
# CHANGE THIS IF YOU WANT GOOD LOCAL EXECUTIONS:
OUTER_HOSTNAME=shuffle-backend
+110 -7
View File
@@ -1978,7 +1978,6 @@ func handleWebhookCallback(resp http.ResponseWriter, request *http.Request) {
}
func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
if request.Method != "POST" {
request.Method = "POST"
}
@@ -1999,7 +1998,7 @@ func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
location := strings.Split(request.URL.String(), "/")
var pipelineId string
if location[1] == "api" {
if len(location) <= 4 {
log.Printf("[INFO] Couldn't handle location. Too short in pipeline: %d", len(location))
@@ -2013,7 +2012,7 @@ func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
userAgent := request.Header.Get("User-Agent")
if strings.Contains(strings.ToLower(userAgent), "microsoftpreview") || strings.Contains(strings.ToLower(userAgent), "googlebot") {
log.Printf("[AUDIT] Blocking googlebot and microsoftbot for pielines. UA: '%s'", userAgent)
log.Printf("[AUDIT] Blocking googlebot and microsoftbot for pipelines. UA: '%s'", userAgent)
resp.WriteHeader(400)
resp.Write([]byte(`{"success": false, "reason": "Google/Microsoft preview bots not allowed. Please change the useragent."}`))
return
@@ -2058,11 +2057,27 @@ func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
return
}
parsedBody := shuffle.GetExecutionbody(body)
// Parse concatenated JSON logs
jsonList, err := parseConcatenatedJSONLogs(string(body))
if err != nil {
log.Printf("[DEBUG] JSON parsing error: %s", err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
parsedBody, err := json.Marshal(jsonList)
if err != nil {
log.Printf("[ERROR] Failed to marshal jsonList: %s", err)
resp.WriteHeader(500)
resp.Write([]byte(`{"success": false}`))
return
}
newBody := shuffle.ExecutionStruct{
Start: pipeline.StartNode,
ExecutionSource: "pipeline",
ExecutionArgument: parsedBody,
ExecutionArgument: string(parsedBody),
}
workflow, err := shuffle.GetWorkflow(ctx, pipeline.WorkflowId)
@@ -2093,8 +2108,7 @@ func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
}
if len(pipeline.StartNode) == 0 {
log.Printf("[WARNING] No start node for pipeline %s - running with workflow default.", pipeline.TriggerId)
log.Printf("[WARNING] No start node for pipeline %s - running with workflow default.")
}
newRequest := &http.Request{
@@ -2108,6 +2122,9 @@ func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
if err == nil {
resp.WriteHeader(200)
resp.Write([]byte(fmt.Sprintf(`{"success": true, "execution_id": "%s"}`, workflowExecution.ExecutionId)))
// Track Sigma rules
trackSigmaRules(ctx, pipeline.OrgId, jsonList)
return
}
@@ -2115,6 +2132,82 @@ func handlePipelineCallback(resp http.ResponseWriter, request *http.Request) {
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, executionResp)))
}
func parseConcatenatedJSONLogs(logs string) ([]map[string]interface{}, error) {
var jsonList []map[string]interface{}
decoder := json.NewDecoder(strings.NewReader(logs))
for decoder.More() {
var jsonObject map[string]interface{}
if err := decoder.Decode(&jsonObject); err != nil {
log.Printf("[WARNING] JSON decoding error: %s. Skipping this object.", err)
continue
}
jsonList = append(jsonList, jsonObject)
}
if err := decoder.Decode(&struct{}{}); err != io.EOF {
return nil, fmt.Errorf("error after decoding all JSON objects: %v", err)
}
return jsonList, nil
}
func trackSigmaRules(ctx context.Context, orgId string, jsonList []map[string]interface{}) {
ruleCount := make(map[string]int)
for _, logEntry := range jsonList {
if rule, ok := logEntry["rule"].(map[string]interface{}); ok {
if ruleName, ok := rule["title"].(string); ok {
ruleCount[ruleName]++
}
}
}
for ruleName, count := range ruleCount {
shuffle.IncrementCache(ctx, orgId, ruleName, count)
log.Printf("[INFO] Rule %s incremented by %d", ruleName, count)
}
}
func handleTenzirHealthUpdate(resp http.ResponseWriter, request *http.Request) {
if request.Method != "POST" {
request.Method = "POST"
}
type HealthUpdate struct {
Status string `json:"status"`
}
var healthUpdate HealthUpdate
err := json.NewDecoder(request.Body).Decode(&healthUpdate)
if err != nil {
resp.WriteHeader(http.StatusBadRequest)
fmt.Fprintf(resp, "Failed to decode JSON: %v", err)
return
}
ctx := context.Background()
status := healthUpdate.Status
result, err := shuffle.GetDisabledRules(ctx)
if (err != nil && err.Error() == "rules doesn't exist") || err == nil {
result.IsTenzirActive = status
result.LastActive = time.Now().Unix()
err = shuffle.StoreDisabledRules(ctx, *result)
if err != nil {
resp.WriteHeader(500)
resp.Write([]byte(`{"success": false}`))
return
}
resp.WriteHeader(200)
resp.Write([]byte(fmt.Sprintf(`{"success": true}`)))
return
}
resp.WriteHeader(500)
resp.Write([]byte(`{"success": false}`))
return
}
func executeCloudAction(action shuffle.CloudSyncJob, apikey string) error {
data, err := json.Marshal(action)
if err != nil {
@@ -5114,6 +5207,7 @@ func initHandlers() {
// PS: For cloud, this has to use cloud storage.
// https://developer.box.com/reference/get-files-id-content/
r.HandleFunc("/api/v1/files/download_remote", shuffle.HandleDownloadRemoteFiles).Methods("POST", "OPTIONS")
r.HandleFunc("/api/v1/files/download_remote_enhanced", shuffle.HandleEnhancedDownloadRemoteFiles).Methods("POST", "OPTIONS")
r.HandleFunc("/api/v1/files/namespaces/{namespace}", shuffle.HandleGetFileNamespace).Methods("GET", "OPTIONS")
r.HandleFunc("/api/v1/files/{fileId}/content", shuffle.HandleGetFileContent).Methods("GET", "OPTIONS")
r.HandleFunc("/api/v1/files/create", shuffle.HandleCreateFile).Methods("POST", "OPTIONS")
@@ -5122,6 +5216,15 @@ func initHandlers() {
r.HandleFunc("/api/v1/files/{fileId}", shuffle.HandleGetFileMeta).Methods("GET", "OPTIONS")
r.HandleFunc("/api/v1/files/{fileId}", shuffle.HandleDeleteFile).Methods("DELETE", "OPTIONS")
r.HandleFunc("/api/v1/files", shuffle.HandleGetFiles).Methods("GET", "OPTIONS")
r.HandleFunc("/api/v1/files/detection/sigma_rules", shuffle.HandleGetSigmaRules).Methods("GET", "OPTIONS")
r.HandleFunc("/api/v1/files/detection/{fileId}/{action}", shuffle.HandleToggleRule).Methods("PUT", "OPTIONS")
r.HandleFunc("/api/v1/files/detection/{action}", shuffle.HandleFolderToggle).Methods("PUT", "OPTIONS")
r.HandleFunc("/api/v1/detection/siem/connect", shuffle.HandleConnectSiem).Methods("GET", "OPTIONS")
r.HandleFunc("/api/v1/detection/siem/node_health", handleTenzirHealthUpdate).Methods("POST","OPTIONS")
r.HandleFunc("/api/v1/detection/{triggerId}/selected_rules", shuffle.HandleGetSelectedRules).Methods("GET","OPTIONS")
r.HandleFunc("/api/v1/detection/{triggerId}/selected_rules/save", shuffle.HandleSaveSelectedRules).Methods("POST","OPTIONS")
// Introduced in 0.9.21 to handle notifications for e.g. failed Workflow
r.HandleFunc("/api/v1/notifications", shuffle.HandleCreateNotification).Methods("POST", "OPTIONS")
+6
View File
@@ -15,6 +15,7 @@ import HealthPage from "./components/HealthPage.jsx";
import theme from "./theme";
import Apps from "./views/Apps";
import AppCreator from "./views/AppCreator";
import DetectionDashBoard from "./views/DetectionDashboard.jsx";
import Welcome from "./views/Welcome.jsx";
import Dashboard from "./views/Dashboard.jsx";
@@ -414,6 +415,11 @@ const App = (message, props) => {
/>
}
/>
<Route
exact
path="/detections/sigma"
element={<DetectionDashBoard globalUrl={globalUrl} />}
/>
<Route
exact
path="/workflows"
+294 -172
View File
@@ -535,6 +535,8 @@ const AngularWorkflow = (defaultprops) => {
const [listCache, setListCache] = React.useState([]);
const [selectedOption, setSelectedOption] = React.useState("");
const [tenzirConfigModalOpen, setTenzirConfigModalOpen] = React.useState(false);
const [rules, setRules] = React.useState([]);
const [sigmaFilesNames, setSigmaFileNames] = React.useState("")
const [distributedFromParent, setDistributedFromParent] = React.useState("")
const [suborgWorkflows, setSuborgWorkflows] = React.useState([])
@@ -1438,7 +1440,7 @@ const releaseToConnectLabel = "Release to Connect"
});
};
const handleKafkaSubmit = (trigger) => {
const handleCommandSubmit = (trigger) => {
if (trigger.trigger_type !== "PIPELINE") {
toast("Unable to save the configuration");
return;
@@ -1446,18 +1448,48 @@ const releaseToConnectLabel = "Release to Connect"
trigger.parameters = []
const topic = document.getElementById('topic')?.value;
const bootstrapServers = document.getElementById('bootstrap_servers')?.value;
const groupId = document.getElementById('group_id')?.value;
//const autoOffsetReset = document.getElementById('auto_offset_reset')?.value;
const command = document.getElementById('sigma')?.value
if(command) {
trigger.parameters.push({
name: "command",
value: command
})
} else {
toast("Please enter the comamnd");
return;
}
// if (autoOffsetReset) {
// trigger.parameters.push({
// name: "auto_offset_reset",
// value: autoOffsetReset
// });
// }
setTenzirConfigModalOpen(false);
};
const handleSubmit = (trigger) => {
if (trigger.trigger_type !== "PIPELINE") {
toast("Unable to save the configuration");
return;
}
if (selectedOption === "Kafka Queue") {
trigger.parameters = []
const topic = document.getElementById('topic')?.value
const bootstrapServers = document.getElementById('bootstrap_servers')?.value
const groupId = document.getElementById('group_id')?.value
const autoOffsetReset = document.getElementById('auto_offset_reset')?.value;
if(topic) {
trigger.parameters.push({
name: "topic",
value: topic
});
})
} else {
toast("please enter the topic name");
toast("Please enter the topic name");
return;
}
@@ -1478,16 +1510,33 @@ const releaseToConnectLabel = "Release to Connect"
});
}
// if (autoOffsetReset) {
// trigger.parameters.push({
// name: "auto_offset_reset",
// value: autoOffsetReset
// });
// }
if (autoOffsetReset) {
trigger.parameters.push({
name: "auto_offset_reset",
value: autoOffsetReset
});
}
setTenzirConfigModalOpen(false);
} else if (selectedOption === "Syslog listener") {
trigger.parameters = []
const endpoint = document.getElementById('endpoint')?.value
if(endpoint) {
trigger.parameters.push({
name: "endpoint",
value: endpoint
})
} else {
toast("Please enter your endpoint");
return;
}
}
};
const handleColoring = (actionId, status, label) => {
if (cy === undefined) {
return
@@ -8245,11 +8294,11 @@ const releaseToConnectLabel = "Release to Connect"
toast("Pipeline deleted!")
return
}
if (trigger.parameters){
trigger.parameters.push({
name: data.name,
value: data.command,
});
});}
if (data.type === "stop") trigger.status = "stopped";
else trigger.status = "running";
@@ -8257,7 +8306,6 @@ const releaseToConnectLabel = "Release to Connect"
setSelectedTrigger(trigger);
setWorkflow(workflow);
console.log("Should set the status to running and save");
saveWorkflow(workflow);
}
})
@@ -8332,6 +8380,32 @@ const releaseToConnectLabel = "Release to Connect"
});
};
const getSigmaInfo = () => {
const url = globalUrl + "/api/v1/files/detection/sigma_rules";
fetch(url, {
method: "GET",
credentials: "include",
headers: {
"Content-Type": "application/json",
},
})
.then((response) =>
response.json().then((responseJson) => {
if (responseJson["success"] === false) {
toast("Failed to get sigma rules");
} else {
setRules(responseJson.sigma_info);
}
})
)
.catch((error) => {
console.log("Error in getting sigma files: ", error);
toast("An error occurred while fetching sigma rules");
});
};
const parsedHeight = isMobile ? bodyHeight - appBarSize * 4 : bodyHeight - appBarSize - 50
const appViewStyle = {
marginLeft: 5,
@@ -15277,6 +15351,14 @@ const releaseToConnectLabel = "Release to Connect"
}
</div>
const defaultEnvironment = environments.find(
(env) => env.default && env.Name.toLowerCase() !== "cloud"
);
if (selectedTrigger.trigger_type === "PIPELINE" && selectedTrigger.environment === "onprem" && defaultEnvironment !== undefined) {
selectedTrigger.environment = defaultEnvironment.Name
setSelectedTrigger(selectedTrigger) }
const PipelineSidebar = Object.getOwnPropertyNames(selectedTrigger).length === 0 || workflow.triggers[selectedTriggerIndex] === undefined && selectedTrigger.trigger_type !== "SCHEDULE" ? null :
<div style={appApiViewStyle}>
<h3 style={{ marginBottom: "5px" }}>
@@ -15375,14 +15457,32 @@ const releaseToConnectLabel = "Release to Connect"
<div
key="syslogListener"
onClick={() => {
// setSelectedOption("Syslog listener")
// setTenzirConfigModalOpen(true);
}}
if(selectedTrigger.status === "running"){
//toast("please stop the trigger to edit the configuration");
return;
} else {
setSelectedOption("Syslog listener");
const url = `${globalUrl}/api/v1/pipelines/pipeline_${selectedTrigger.id}`
const command = `from tcp://192.168.1.100:5162 read syslog | import`
const pipelineConfig = {
command: command,
name: selectedTrigger.label,
type: "create",
environment: selectedTrigger.environment,
workflow_id: workflow.id,
trigger_id: selectedTrigger.id,
start_node: "",
url:url,
};
submitPipeline(selectedTrigger, selectedTriggerIndex, pipelineConfig);
}}}
style={{
border: "1px solid rgba(255,255,255,0.3)",
borderRadius: theme.palette.borderRadius,
padding: 10,
cursor: "not-allowed",
cursor: "pointer",
marginTop: 5,
display: "flex",
alignItems: "center",
@@ -15392,27 +15492,46 @@ const releaseToConnectLabel = "Release to Connect"
control={
<Radio
checked={selectedOption === "Syslog listener"}
onChange={() => setSelectedOption("Syslog listener")}
onChange={() => {
if (selectedTrigger.status !== "running"){
setSelectedOption("Syslog listener")}}
}
value={"Syslog listener"}
name="option"
disabled={true}
/>
}
label="Start Syslog listener"
label= {selectedOption === "Syslog listener" && selectedTrigger.status === "running" ? "listening at 192.168.1.100:5162" : "Start Syslog listener"}
/>
</div>
<div
key="sigmaRulesearch"
onClick={() => {
// setSelectedOption("Sigma Rulesearch")
// setTenzirConfigModalOpen(true);
}}
if(selectedTrigger.status === "running"){
// toast("please stop the trigger to edit the configuration");
return;
} else {
setSelectedOption("SigmaRule");
const url = `${globalUrl}/api/v1/pipelines/pipeline_${selectedTrigger.id}`
const command = `export | sigma /var/lib/tenzir/sigma_rules | to ${url}`
const pipelineConfig = {
command: command,
name: selectedTrigger.label,
type: "create",
environment: selectedTrigger.environment,
workflow_id: workflow.id,
trigger_id: selectedTrigger.id,
start_node: "",
url:url,
};
submitPipeline(selectedTrigger, selectedTriggerIndex, pipelineConfig);
}}}
style={{
border: "1px solid rgba(255,255,255,0.3)",
borderRadius: theme.palette.borderRadius,
padding: 10,
cursor: "not-allowed",
cursor: "pointer",
marginTop: 5,
display: "flex",
alignItems: "center",
@@ -15421,11 +15540,12 @@ const releaseToConnectLabel = "Release to Connect"
<FormControlLabel
control={
<Radio
checked={selectedOption === "Sigma Rulesearch"}
onChange={() => setSelectedOption("Sigma Rulesearch")}
checked={selectedOption === "SigmaRule"}
onChange={() => {
if (selectedTrigger.status !== "running"){
setSelectedOption("SigmaRule")}}}
value={"Sigma Rulesearch"}
name="option"
disabled={true}
/>
}
label="Run Sigma Rulesearch"
@@ -15456,7 +15576,9 @@ const releaseToConnectLabel = "Release to Connect"
control={
<Radio
checked={selectedOption === "Kafka Queue"}
onChange={() => setSelectedOption("Kafka Queue")}
onChange={() => {
if (selectedTrigger.status !== "running"){
setSelectedOption("Kafka Queue")}}}
value={"Kafka Queue"}
name="option"
/>
@@ -15472,10 +15594,12 @@ const releaseToConnectLabel = "Release to Connect"
disabled={selectedTrigger.status === "running"}
onClick={() => {
if (selectedOption === "Kafka Queue"){
const url = `${globalUrl}/api/v1/pipelines/pipeline_${selectedTrigger.id}`
const topic = (selectedTrigger?.parameters?.find(param => param.name === "topic")?.value) || ''
const bootstrapServers = (selectedTrigger?.parameters?.find(param => param.name === "bootstrap_servers")?.value) || ''
const groupId = (selectedTrigger?.parameters?.find(param => param.name === "group_id")?.value) || ''
// const autoOffsetReset = (selectedTrigger?.parameters?.find(param => param.name === "auto_offset_reset")?.value) || ''
const autoOffsetReset = (selectedTrigger?.parameters?.find(param => param.name === "auto_offset_reset")?.value) || ''
let command = "from kafka"
if(topic) {
@@ -15496,15 +15620,15 @@ const releaseToConnectLabel = "Release to Connect"
} else {
command = `${command},group.id=${selectedTrigger.id}`
}
// if(autoOffsetReset) {
// command = `${command},auto.offset.reset=${autoOffsetReset}`
// } else {
// command = `${command},auto.offset.reset=earliest`
if(autoOffsetReset) {
command = `${command},auto.offset.reset=${autoOffsetReset}`
} else {
command = `${command},auto.offset.reset=earliest`
// }
}
command = `${command},auto.offset.reset=earliest`
command = `${command},client.id=${selectedTrigger.id},enable.auto.commit=true,auto.commit.interval.ms=1`
command = `${command} read json | to ${globalUrl}/api/v1/pipelines/pipeline_${selectedTrigger.id}`
command = `${command} read json | to ${url}`
const pipelineConfig = {
command: command,
@@ -15514,9 +15638,11 @@ const releaseToConnectLabel = "Release to Connect"
workflow_id: workflow.id,
trigger_id: selectedTrigger.id,
start_node: "",
url: url,
};
submitPipeline(selectedTrigger, selectedTriggerIndex, pipelineConfig);
}}
}
}}
color="primary"
>
Start
@@ -21229,146 +21355,142 @@ const releaseToConnectLabel = "Release to Connect"
</Dialog>
) : null;
const tenzirConfigModal = tenzirConfigModalOpen ? (
<Dialog
PaperComponent={PaperComponent}
hideBackdrop={true}
disableEnforceFocus={true}
disableBackdropClick={true}
style={{ pointerEvents: "none" }}
open={tenzirConfigModalOpen}
PaperProps={{
style: {
pointerEvents: "auto",
color: "white",
minWidth: 600,
minHeight: 450,
maxHeight: 450,
padding: 15,
overflow: "hidden",
zIndex: 10012,
border: theme.palette.defaultBorder,
},
}}
>
<div
style={{
flex: 2,
padding: 0,
minHeight: isMobile ? "90%" : 700,
maxHeight: isMobile ? "90%" : 700,
overflowY: "auto",
overflowX: isMobile ? "auto" : "hidden",
}}
>
<DialogTitle id="tenzir-config-modal" style={{ cursor: "move" }}>
<div style={{ color: "white" }}>Configuration options for {selectedOption}</div>
</DialogTitle>
<DialogContent>
{selectedOption === "Kafka Queue" && (
<>
<b>Topic</b>
<TextField
id="topic"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"topic name"}
defaultValue={(selectedTrigger?.parameters?.find(param => param.name === "topic")?.value) || ''}
/>
<b>bootstrap.servers</b>
<TextField
id="bootstrap_servers"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"broker1.example.com:9092,192.168.1.100:9092"}
defaultValue={(selectedTrigger?.parameters?.find(param => param.name === "bootstrap_servers")?.value) || ''}
/>
<b>group.id</b>
<TextField
id="group_id"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"tenzir"}
defaultValue={(selectedTrigger?.parameters?.find(param => param.name === "group_id")?.value) || ''}
/>
{/* <b>auto.offest.reset</b>
<TextField
id="auto_offset_reset"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"earliest"}
defaultValue={(selectedTrigger?.parameters?.find(param => param.name === "auto_offset_reset")?.value) || ''}
/> */}
</>
)}
</DialogContent>
<DialogActions>
<Button
style={{ borderRadius: "0px" }}
onClick={() => {
setTenzirConfigModalOpen(false);
const TenzirConfigModal = () => {
if (!tenzirConfigModalOpen) return null;
return (
<Dialog
PaperComponent={PaperComponent}
hideBackdrop={true}
disableEnforceFocus={true}
disableBackdropClick={true}
style={{ pointerEvents: "none" }}
open={tenzirConfigModalOpen}
PaperProps={{
style: {
pointerEvents: "auto",
color: "white",
minWidth: 600,
minHeight: 550,
maxHeight: 550,
padding: 15,
overflow: "hidden",
zIndex: 10012,
border: theme.palette.defaultBorder,
},
}}
>
<DialogTitle id="tenzir-config-modal" style={{ cursor: "move" }}>
<div style={{ color: "white" }}>Configuration options for Kafka</div>
</DialogTitle>
<DialogContent>
{selectedOption === "Kafka Queue" ? (
<div>
<b>Topic</b>
<TextField
id="topic"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
color="primary"
>
Cancel
</Button>
<Button
style={{ borderRadius: "0px" }}
onClick={() => {
handleKafkaSubmit(selectedTrigger);
InputProps={{
style: {},
}}
fullWidth
color="primary"
>
Submit
</Button>
</DialogActions>
</div>
placeholder={"topic name"}
defaultValue={
selectedTrigger?.parameters?.find(
(param) => param.name === "topic",
)?.value || ""
}
/>
<b>bootstrap.servers</b>
<TextField
id="bootstrap_servers"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"broker1.example.com:9092,192.168.1.100:9092"}
defaultValue={
selectedTrigger?.parameters?.find(
(param) => param.name === "bootstrap_servers",
)?.value || ""
}
/>
<b>group.id</b>
<TextField
id="group_id"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"tenzir"}
defaultValue={
selectedTrigger?.parameters?.find(
(param) => param.name === "group_id",
)?.value || ""
}
/>
<b>auto.offest.reset</b>
<TextField
id="auto_offset_reset"
style={{
backgroundColor: theme.palette.inputColor,
borderRadius: theme.palette.borderRadius,
}}
InputProps={{
style: {},
}}
fullWidth
color="primary"
placeholder={"earliest"}
defaultValue={
selectedTrigger?.parameters?.find(
(param) => param.name === "auto_offset_reset",
)?.value || ""
}
/>
</div>
) : null}{" "}
<IconButton
style={{
zIndex: 5000,
position: "absolute",
top: 14,
right: 18,
color: "grey",
}}
</DialogContent>
<DialogActions>
<Button
style={{ borderRadius: "0px" }}
onClick={() => {
setTenzirConfigModalOpen(false);
}}
color="primary"
>
<CloseIcon />
</IconButton>
</Dialog>
) : null;
Cancel
</Button>
<Button
style={{ borderRadius: "0px" }}
onClick={() => {
handleSubmit(selectedTrigger);
}}
color="primary"
>
Submit
</Button>
</DialogActions>
</Dialog>
);
};
const SuggestionBoxUi = () => {
@@ -21947,7 +22069,7 @@ const releaseToConnectLabel = "Release to Connect"
{codePopoutModal}
{workflowRevisions}
{authenticationModal}
{tenzirConfigModal}
{<TenzirConfigModal/>}
{/*editWorkflowModal*/}
{authgroupModal}
{executionArgumentModal}
+198
View File
@@ -0,0 +1,198 @@
import React, { useState } from "react";
import {
Container,
Box,
TextField,
Switch,
Typography,
Button,
} from "@mui/material";
import { toast } from "react-toastify";
import RuleCard from "./RuleCard";
import CircularProgress from "@material-ui/core/CircularProgress";
const handleDirectoryChange = (folderDisabled, setFolderDisabled, globalUrl, isTenzirActive) => {
if (!isTenzirActive) {
toast("connect to siem first for global enable/disable to work");
return;
}
const action = folderDisabled ? "enable_folder" : "disable_folder";
const url = `${globalUrl}/api/v1/files/detection/${action}`;
fetch(url, {
method: "PUT",
credentials: "include",
headers: {
"Content-Type": "application/json",
},
})
.then((response) =>
response.json().then((responseJson) => {
if (responseJson["success"] === true) {
if (action === "enable_folder") setFolderDisabled(false);
else setFolderDisabled(true);
} else {
//toast(`failed to disable rule`);
}
})
)
.catch((error) => {
console.log(`Error in ${action} the rule: `, error);
toast(`An error occurred while ${action} the rule`);
});
};
const Detection = ({
globalUrl,
ruleInfo,
folderDisabled,
setFolderDisabled,
isTenzirActive,
}) => {
const [searchQuery, setSearchQuery] = useState("");
const [loading, setLoading] = useState(false);
const handleConnectClick = () => {
if (!isTenzirActive) {
setLoading(true);
const url = `${globalUrl}/api/v1/detection/siem/connect`;
fetch(url, {
method: "GET",
credentials: "include",
headers: {
"Content-Type": "application/json",
},
})
.then((response) =>
response.json().then((responseJson) => {
if (responseJson["success"] === true) {
setTimeout(() => {
setLoading(false);
window.location.reload();
}, 15000);
} else {
setLoading(false);
toast("Failed to connect to SIEM");
}
})
)
.catch((error) => {
setLoading(false);
console.log(`Error in connecting to SIEM: `, error);
toast("An error occurred while connecting to SIEM");
});
} else {
console.log("Already connected to SIEM");
}
};
const filteredRules = ruleInfo?.filter((rule) =>
rule.title.toLowerCase().includes(searchQuery.toLowerCase()) ||
rule.description.toLowerCase().includes(searchQuery.toLowerCase())
);
return (
<Container sx={{ mt: 4 }}>
<Box sx={{ border: "1px solid #ccc", borderRadius: 2, p: 3 }}>
<Box
sx={{
display: "flex",
justifyContent: "space-between",
alignItems: "center",
mb: 2,
}}
>
<Typography variant="h6" component="div">
Sigma Detection Rules
</Typography>
<Button
variant="contained"
onClick={handleConnectClick}
disabled={loading} // Disable the button while loading
style={{ backgroundColor: isTenzirActive ? "green" : "red"}}
>
{loading ? <CircularProgress size={24} /> : isTenzirActive ? "Connected to siem" : "Connect to siem"}
</Button>
</Box>
<Box
sx={{
display: "flex",
justifyContent: "space-between",
alignItems: "center",
mb: 2,
}}
>
<Box
sx={{
display: "flex",
}}
>
<TextField
label="Search rules"
variant="outlined"
size="small"
sx={{ mr: 2 }}
value={searchQuery}
onChange={(e) => setSearchQuery(e.target.value)}
/>
{/* <Button
color="primary"
variant="contained"
onClick={() => uploadRef.current.click()}
>
<PublishIcon /> Upload sigma file
</Button>
<input
hidden
type="file"
multiple
ref={uploadRef}
onChange={(event) => {
uploadFiles(event.target.files);
}}
/> */}
</Box>
<Box sx={{ display: "flex", alignItems: "center" }}>
<Typography variant="body2" sx={{ mr: 1 }}>
Global disable/enable
</Typography>
<Switch
checked={!folderDisabled}
onChange={() =>
handleDirectoryChange(folderDisabled, setFolderDisabled, globalUrl, isTenzirActive)
}
disabled={!isTenzirActive}
/>
</Box>
</Box>
<Box
sx={{
height: "500px",
width: "100%",
overflowY: "auto",
border: "1px solid #ddd",
p: 1,
}}
>
{filteredRules?.length > 0 &&
filteredRules.map((card) => (
<RuleCard
key={card.file_id}
ruleName={card.title}
description={card.description}
file_id={card.file_id}
globalUrl={globalUrl}
folderDisabled={folderDisabled}
isTenzirActive={isTenzirActive}
{...card}
/>
))}
</Box>
</Box>
</Container>
);
};
export default Detection;
+162
View File
@@ -0,0 +1,162 @@
import React, { useState, useEffect } from "react";
import { Container, CircularProgress, Typography } from "@mui/material";
import { toast } from "react-toastify";
import Detection from "./Detection";
const DetectionDashBoard = (props) => {
const { globalUrl } = props;
const [ruleInfo, setRuleInfo] = useState(null);
const [, setSelectedRule] = useState(null);
const [, setFileData] = useState("");
const [isTenzirActive, setIsTenzirActive] = useState(false);
const [folderDisabled, setFolderDisabled] = useState(false);
const [isLoading, setIsLoading] = useState(false);
const [importAttempts, setImportAttempts] = useState(0);
const maxImportAttempts = 2;
useEffect(() => {
const fetchTimeout = setTimeout(() => {
fetchSigmaInfo();
}, 1000); // Delay by 1 second
return () => clearTimeout(fetchTimeout);
}, [globalUrl]);
useEffect(() => {
if (ruleInfo && ruleInfo.length === 0 && importAttempts < maxImportAttempts) {
importSigmaFromUrl();
}
}, [ruleInfo]);
const openEditBar = (rule) => {
setSelectedRule(rule);
fetchFileContent(rule.file_id);
};
const handleSave = (updatedContent) => {
toast("This will be saved");
};
const fetchFileContent = (file_id) => {
setFileData("");
fetch(`${globalUrl}/api/v1/files/${file_id}/content`, {
method: "GET",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
},
credentials: "include",
})
.then((response) => {
if (response.status !== 200) {
console.log("Status not 200 for file :O!");
return "";
}
return response.text();
})
.then((respdata) => {
if (respdata.length === 0) {
toast("Failed getting file. Is it deleted?");
return;
}
setFileData(respdata);
})
.catch((error) => {
toast(error.toString());
});
};
const fetchSigmaInfo = () => {
const url = `${globalUrl}/api/v1/files/detection/sigma_rules`;
setIsLoading(true);
fetch(url, {
method: "GET",
credentials: "include",
headers: {
"Content-Type": "application/json",
},
})
.then((response) => response.json())
.then((responseJson) => {
if (responseJson["success"] === false) {
toast("Failed to get sigma rules");
} else {
setRuleInfo(responseJson.sigma_info || []);
setFolderDisabled(responseJson.folder_disabled);
setIsTenzirActive(responseJson.is_tenzir_active);
}
setIsLoading(false);
})
.catch((error) => {
setIsLoading(false);
console.log("Error in getting sigma files: ", error);
toast("An error occurred while fetching sigma rules");
setRuleInfo([]);
});
};
const importSigmaFromUrl = () => {
setIsLoading(true);
setImportAttempts((prevAttempts) => prevAttempts + 1);
const url = "https://github.com/satti-hari-krishna-reddy/shuffle_sigma";
const folder = "sigma";
const parsedData = {
url: url,
path: folder,
field_3: "main",
};
toast(`Getting files from url ${url}. This may take a while if the repository is large. Please wait...`);
fetch(`${globalUrl}/api/v1/files/download_remote_enhanced`, {
method: "POST",
mode: "cors",
headers: {
Accept: "application/json",
},
body: JSON.stringify(parsedData),
credentials: "include",
})
.then((response) => response.json())
.then((responseJson) => {
if (responseJson.success) {
toast("Successfully loaded files from " + url);
fetchSigmaInfo(); // Fetch again after successful import
} else {
toast(responseJson.reason ? `Failed loading: ${responseJson.reason}` : "Failed loading");
}
setIsLoading(false);
})
.catch((error) => {
toast(error.toString());
setIsLoading(false);
});
};
if (isLoading && (!ruleInfo || ruleInfo.length === 0)) {
return (
<Container style={{ display: "flex", justifyContent: "center", alignItems: "center", height: "100vh" }}>
<div>
<CircularProgress />
<Typography variant="h6" style={{ marginTop: 20 }}>Downloading rules, please wait...</Typography>
</div>
</Container>
);
}
return (
<Container style={{ display: "flex" }}>
<Detection
globalUrl={globalUrl}
ruleInfo={ruleInfo}
folderDisabled={folderDisabled}
setFolderDisabled={setFolderDisabled}
isTenzirActive={isTenzirActive}
/>
</Container>
);
};
export default DetectionDashBoard;
+43
View File
@@ -0,0 +1,43 @@
import React from 'react';
import { Box, Typography, Button, TextField } from '@mui/material';
const EditComponent = ({ ruleName, description, content, setContent, lastEdited, editedBy, onSave }) => {
const handleSave = () => {
onSave(content);
};
return (
<Box sx={{ p: 2, border: '1px solid #ccc', borderRadius: 2, height: '100%', width: '100%', marginTop:'30px'}}>
<Box sx={{ display: 'flex', justifyContent: 'space-between', alignItems: 'center' }}>
<Typography variant="h6">{ruleName}</Typography>
</Box>
<Typography variant="body2" style={{ marginTop: '2%' }}>
{description}
</Typography>
<Typography variant="body2" sx={{ mt: 1 }}>
Last edited: {lastEdited}
</Typography>
<Typography variant="body2">
Edited By: {editedBy}
</Typography>
<Box sx={{ mt: 2 }}>
<TextField
multiline
rows={12}
value={content}
onChange={(e) => setContent(e.target.value)}
variant="outlined"
fullWidth
/>
</Box>
<Box sx={{ display: 'flex', justifyContent: 'flex-end', mt: 2 }}>
<Button variant="contained" color="primary" onClick={handleSave}>
Save
</Button>
</Box>
</Box>
);
};
export default EditComponent;
+168
View File
@@ -0,0 +1,168 @@
import React from "react";
import {
Card,
CardContent,
IconButton,
Typography,
Switch,
} from "@mui/material";
import EditIcon from "@mui/icons-material/Edit";
import { toast } from "react-toastify";
import ShuffleCodeEditor from "../components/ShuffleCodeEditor1.jsx";
const RuleCard = ({ ruleName, description, file_id, globalUrl, folderDisabled, isTenzirActive, ...otherProps }) => {
const [openCodeEditor, setOpenCodeEditor] = React.useState(false);
const [fileData, setFileData] = React.useState("");
const [isEnabled, setIsEnabled] = React.useState(otherProps.is_enabled);
const isCloud = ["localhost:3002", "shuffler.io"].includes(window.location.host);
const handleSwitchChange = (event) => {
if (folderDisabled) {
toast("enable the directory to enable individual rules");
return;
}
if (!isTenzirActive) {
toast("connect to the siem to enable/disable the rule");
return;
}
const newIsEnabled = event.target.checked;
toggleRule(file_id, !newIsEnabled, globalUrl, () => {
setIsEnabled(newIsEnabled);
});
};
const UpdateText = (text) => {
fetch(`${globalUrl}/api/v1/files/${file_id}/edit`, {
method: "PUT",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
},
body: text,
credentials: "include",
})
.then((response) => {
if (response.status !== 200) {
console.log("Can't update file");
}
return response.json();
})
.then((responseJson) => {
if (responseJson.success === true) {
toast("Successfully updated file");
}
})
.catch((error) => {
toast("Error updating file: " + error.toString());
});
};
return (
<Card variant="outlined" sx={{ mb: 2 }}>
<CardContent>
<div
style={{
display: 'flex',
justifyContent: 'space-between',
alignItems: 'center',
marginBottom: 16,
}}
>
<Typography variant="h6">{ruleName}</Typography>
<div style={{ display: 'flex', alignItems: 'center' }}>
<IconButton onClick={() => openEditBar(file_id, setOpenCodeEditor, setFileData, globalUrl)}>
<EditIcon />
</IconButton>
<Switch
checked={isEnabled && !folderDisabled}
onChange={handleSwitchChange}
disabled={!isTenzirActive}
/>
</div>
</div>
<Typography variant="body2" style={{ marginTop: '2%' }}>
{description}
</Typography>
<ShuffleCodeEditor
isCloud={isCloud}
expansionModalOpen={openCodeEditor}
setExpansionModalOpen={setOpenCodeEditor}
setcodedata={setFileData}
codedata={fileData}
isFileEditor={true}
key={fileData} // https://reactjs.org/docs/reconciliation.html#recursing-on-children
runUpdateText={UpdateText}
/>
</CardContent>
</Card>
);
}
const toggleRule = (fileId, isCurrentlyEnabled, globalUrl, callback) => {
const action = isCurrentlyEnabled ? "disable" : "enable";
const url = `${globalUrl}/api/v1/files/detection/${fileId}/${action}_rule`;
fetch(url, {
method: "PUT",
credentials: "include",
headers: {
"Content-Type": "application/json",
},
})
.then((response) =>
response.json().then((responseJson) => {
if (responseJson["success"] === false) {
toast(`Failed to ${action} the rule`);
} else {
toast(`Rule ${action}d successfully`);
callback();
}
})
)
.catch((error) => {
console.log(`Error in ${action}ing the rule: `, error);
toast(`An error occurred while ${action}ing the rule`);
});
};
const openEditBar = (file_id, setOpenCodeEditor, setFileData, globalUrl) => {
getFileContent(file_id, setFileData, globalUrl)
setOpenCodeEditor(true);
};
const getFileContent = (file_id, setFileData, globalUrl) => {
setFileData("");
fetch(globalUrl + "/api/v1/files/" + file_id + "/content", {
method: "GET",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
},
credentials: "include",
})
.then((response) => {
if (response.status !== 200) {
console.log("Status not 200 for file :O!");
return "";
}
return response.text();
})
.then((respdata) => {
if (respdata.length === 0) {
toast("Failed getting file. Is it deleted?");
return;
}
return respdata
})
.then((responseData) => {
setFileData(responseData);
})
.catch((error) => {
toast(error.toString());
});
};
export default RuleCard;
+500 -115
View File
@@ -10,7 +10,7 @@ package main
import (
"github.com/shuffle/shuffle-shared"
"archive/zip"
"bytes"
"context"
"encoding/json"
@@ -29,6 +29,7 @@ import (
"strings"
"sync"
"time"
"path/filepath"
//"os/signal"
//"syscall"
@@ -101,6 +102,7 @@ var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL")
var memcached = os.Getenv("SHUFFLE_MEMCACHED")
var tenzirUrl = os.Getenv("SHUFFLE_TENZIR_URL")
var apiKey = os.Getenv("AUTH_FOR_ORBORUS")
var executionIds = []string{}
var namespacemade = false // For K8s
@@ -1805,6 +1807,12 @@ func main() {
log.Printf("[WARNING] Defaulting to environment name %s. Set environment variable ENVIRONMENT_NAME to change. This should be the same as in the frontend action.", environment)
}
if tenzirUrl == "" {
tenzirUrl = "http://localhost:5160"
log.Printf("[WARNING] SHUFFLE_TENZIR_URL not set, falling back to default URL: %s",tenzirUrl)
}
// FIXME - during init, BUILD and/or LOAD worker and app_sdk
// Build/load app_sdk so it can be loaded as 127.0.0.1:5000/walkoff_app_sdk
log.Printf("[INFO] Setting up Docker environment. Downloading worker and App SDK!")
@@ -1906,6 +1914,7 @@ func main() {
log.Printf("[INFO] Waiting for executions at %s with Environment %#v", fullUrl, environment)
hasStarted := false
for {
_ = sendTenzirHealthStatus()
if req.Method == "POST" {
// Should find data to send (memory etc.)
@@ -2021,6 +2030,71 @@ func main() {
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "CATEGORY_UPDATE" {
err := deployTenzirNode()
if err != nil{
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
}
err = handleFileCategoryChange()
if err != nil {
log.Printf("[ERROR] Failed to download the file category: %s", err)
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "DISABLE_SIGMA_FILE" {
fileName := incRequest.ExecutionArgument
err := deployTenzirNode()
if err != nil{
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
}
err = disableRule(fileName)
if err != nil {
log.Printf("[ERROR] Failed to disable the sigma file %s, reason: %s", fileName, err)
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "ENABLE_SIGMA_FILE" {
fileName := incRequest.ExecutionArgument
err := deployTenzirNode()
if err != nil{
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
}
err = enableRule(fileName)
if err != nil {
log.Printf("[ERROR] Failed to disable the sigma file %s, reason: %s", fileName, err)
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "DISABLE_SIGMA_FOLDER" {
err := deployTenzirNode()
if err != nil{
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
}
err = removeAllFiles()
if err != nil {
log.Printf("[ERROR] Failed to disable the sigma rules: %s", err)
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "START_TENZIR" {
err := deployTenzirNode()
if err != nil{
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else {
newrequests = append(newrequests, incRequest)
}
@@ -2211,7 +2285,6 @@ func main() {
}
}
time.Sleep(time.Duration(sleepTime) * time.Second)
}
}
@@ -2371,11 +2444,6 @@ func main() {
// docker run tenzir/tenzir:latest 'from http://192.168.86.44:5002/api/v1/orgs/7e9b9007-5df2-4b47-bca5-c4d267ef2943/cache/CIDR%20ranges?type=text&authorization=cec9d01f-09b2-4419-8a0a-76c6046e3fef read lines | to http://192.168.86.44:5002/api/v1/hooks/webhook_665ace5f-f27b-496a-a365-6e07eb61078c write lines'
func handlePipeline(incRequest shuffle.ExecutionRequest) error {
if tenzirUrl == "" {
tenzirUrl = "http://localhost:5160"
log.Printf("[WARNING] SHUFFLE_TENZIR_URL not set, falling back to default URL: %s", tenzirUrl)
}
err := deployTenzirNode()
if err != nil {
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
@@ -2419,7 +2487,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error {
log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
return err
}
_, err = updatePipelineState(pipelineId, "stop")
_, err = updatePipelineState(command, pipelineId, "stop")
if err != nil {
log.Printf("[ERROR] Failed to stop Pipeline: %s reason:%s ", pipelineId, err)
return err
@@ -2439,7 +2507,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error {
log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
return err
}
_, err = updatePipelineState(pipelineId, "start")
_, err = updatePipelineState(command, pipelineId, "start")
if err != nil {
log.Printf("[ERROR] Failed to start Pipeline: %s reason:%s ", pipelineId, err)
return err
@@ -2472,9 +2540,19 @@ func deployTenzirNode() error {
return nil
}
containerInfo, err := dockercli.ContainerInspect(ctx, containerName)
if err != nil {
if dockerclient.IsErrNotFound(err) {
containerInfo, err := dockercli.ContainerInspect(ctx, containerName)
if err != nil {
if dockerclient.IsErrNotFound(err) {
// Create network if it doesn't exist
networkName := "tenzir-network"
networkSubnet := "192.168.1.0/24"
networkGateway := "192.168.1.1"
err = createNetworkIfNotExists(ctx, networkName, networkSubnet, networkGateway)
if err != nil {
log.Printf("[ERROR] Failed to create network: %s", err)
return err
}
// Check if image exists
_, _, err := dockercli.ImageInspectWithRaw(ctx, imageName)
@@ -2535,30 +2613,6 @@ func deployTenzirNode() error {
return nil
}
func checkTenzirNode() error {
retries := 20
retryInterval := 3 * time.Second
url := fmt.Sprintf("%s/api/v0/ping", tenzirUrl)
forwardMethod := "POST"
client := http.Client{}
req, err := http.NewRequest(forwardMethod, url, nil)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err
}
for i := 0; i < retries; i++ {
resp, err := client.Do(req)
if err == nil && resp.StatusCode == http.StatusOK {
return nil
}
time.Sleep(retryInterval)
}
return fmt.Errorf("tenzir node is not available")
}
func createAndStartTenzirNode(ctx context.Context, containerName, imageName string, containerStartOptions container.StartOptions) error {
healthconfig := &container.HealthConfig{
Test: []string{"tenzir --connection-timeout=30s --connection-retry-delay=1s 'api /ping'"},
@@ -2574,23 +2628,34 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri
Entrypoint: []string{containerName},
}
hostConfig := &container.HostConfig{
PortBindings: nat.PortMap{
"5160/tcp": []nat.PortBinding{{HostPort: "5160"}},
},
Mounts: []mount.Mount{
{
Type: mount.TypeVolume,
Source: containerName,
Target: "/var/lib/tenzir/",
},
},
VolumeDriver: "local",
}
_, err := dockercli.ContainerCreate(ctx, config, hostConfig, nil, nil, containerName)
if err != nil {
return err
}
hostConfig := &container.HostConfig{
PortBindings: nat.PortMap{
"5160/tcp": []nat.PortBinding{{HostPort: "5160"}},
},
Mounts: []mount.Mount{
{
Type: mount.TypeVolume,
Source: containerName,
Target: "/var/lib/tenzir/",
},
},
VolumeDriver: "local",
}
networkingConfig := &network.NetworkingConfig{
EndpointsConfig: map[string]*network.EndpointSettings{
"tenzir-network": {
IPAMConfig: &network.EndpointIPAMConfig{
IPv4Address: "192.168.1.100",
},
},
},
}
_, err := dockercli.ContainerCreate(ctx, config, hostConfig, networkingConfig, nil, containerName)
if err != nil {
return err
}
err = dockercli.ContainerStart(ctx, containerName, containerStartOptions)
if err != nil {
@@ -2599,16 +2664,76 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri
}
log.Printf("[INFO] Tenzir Node container started successfully")
log.Printf("[INFO] Waiting for Tenzir to become available ...")
err = checkTenzirNode()
if err != nil {
return err
}
log.Printf("[INFO] Successfully deployed Tenzir Node !")
log.Printf("[INFO] Waiting for Tenzir to become available ...")
err = checkTenzirNode()
if err != nil {
return err
}
log.Printf("[INFO] Successfully deployed Tenzir Node!")
return nil
}
func createNetworkIfNotExists(ctx context.Context, networkName, subnet, gateway string) error {
networks, err := dockercli.NetworkList(ctx, types.NetworkListOptions{})
if err != nil {
return err
}
for _, network := range networks {
if network.Name == networkName {
// Network exists
return nil
}
}
ipamConfig := &network.IPAM{
Config: []network.IPAMConfig{
{
Subnet: subnet,
Gateway: gateway,
},
},
}
networkCreate := types.NetworkCreate{
CheckDuplicate: true,
Driver: "bridge",
IPAM: ipamConfig,
}
_, err = dockercli.NetworkCreate(ctx, networkName, networkCreate)
if err != nil {
return err
}
return nil
}
func checkTenzirNode() error {
retries := 5
retryInterval := 3 * time.Second
url := fmt.Sprintf("%s/api/v0/ping",tenzirUrl)
forwardMethod := "POST"
client := http.Client{}
req, err := http.NewRequest(forwardMethod, url, nil)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err
}
for i := 0; i < retries; i++ {
resp, err := client.Do(req)
if err == nil && resp.StatusCode == http.StatusOK {
return nil
}
time.Sleep(retryInterval)
}
return fmt.Errorf("tenzir node is not available")
}
func createPipeline(command, identifier string) (string, error) {
toBeDeleted := false
@@ -2627,32 +2752,36 @@ func createPipeline(command, identifier string) (string, error) {
log.Printf("[INFO] an existing pipeline found with ID: %s. it will be deleted", pipelineId)
toBeDeleted = true
}
if strings.Contains(command, "shuffler.io") {
// if strings.Contains(command, "shuffler.io") {
} else {
var scheme string
if strings.Contains(command, "http://") {
scheme = "http://"
} else if strings.Contains(command, "https://") {
scheme = "https://"
}
// } else {
// var scheme string
// if strings.Contains(command, "http://") {
// scheme = "http://"
// } else if strings.Contains(command, "https://") {
// scheme = "https://"
// }
startIndex := strings.Index(command, scheme)
if startIndex != -1 {
endIndex := startIndex + len(scheme)
endIndex += strings.Index(command[endIndex:], "/")
// startIndex := strings.Index(command, scheme)
// if startIndex != -1 {
// endIndex := startIndex + len(scheme)
// endIndex += strings.Index(command[endIndex:], "/")
// command = command[:startIndex] + baseUrl + command[endIndex:]
// }
// }
//command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | sigma /var/lib/tenzir/rule.yaml"
//command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | import"
command = command[:startIndex] + baseUrl + command[endIndex:]
}
}
requestBody := map[string]interface{}{
"definition": command,
"name": identifier,
"hidden": false,
"autostart": map[string]bool{
"created": true,
"completed": true,
"failed": true,
"completed": false,
"failed": false,
},
"autodelete": map[string]bool{
"completed": false,
@@ -2718,13 +2847,14 @@ func createPipeline(command, identifier string) (string, error) {
return id, nil
}
func updatePipelineState(pipelineId, action string) (string, error) {
func updatePipelineState(command, pipelineId, action string) (string, error) {
url := fmt.Sprintf("%s/api/v0/pipeline/update", tenzirUrl)
forwardMethod := "POST"
requestBody := map[string]interface{}{
"id": pipelineId,
"definition": command,
"action": action,
"autostart": map[string]bool{
"created": true,
@@ -2870,53 +3000,308 @@ func searchPipeline(identifier string) (string, error) {
return "", errors.New("no existing pipeline found with name")
}
// func savePipelineData(pipelineId, identifier, status string) error {
func handleFileCategoryChange() error {
apiEndpoint := baseUrl + "/api/v1/files/namespaces/sigma"
req, err := http.NewRequest("GET", apiEndpoint, nil)
if err != nil {
return err
}
// url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", baseUrl)
// identifierWithoutPrefix := strings.TrimPrefix(identifier, "shuffle-")
req.Header.Add("Authorization", "Bearer "+apiKey)
// forwardMethod := "PUT"
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
// payload := map[string]interface{}{
// "pipeline_id": pipelineId,
// "trigger_id": identifierWithoutPrefix,
// "status": status,
// }
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("received non-200 response: %s", resp.Status)
}
// payloadBytes, err := json.Marshal(payload)
// if err != nil {
// log.Printf("[ERROR] Failed to marshal payload: %s", err)
// return err
// }
out, err := os.Create("files.zip")
if err != nil {
return err
}
// forwardData := bytes.NewBuffer(payloadBytes)
defer out.Close()
defer os.Remove("files.zip")
// req, err := http.NewRequest(
// forwardMethod,
// url,
// forwardData,
// )
// if err != nil {
// log.Printf("[ERROR] Failed to create HTTP request: %s", err)
// return err
// }
// req.Header.Set("Content-Type", "application/json")
_, err = io.Copy(out, resp.Body)
if err != nil {
return err
}
// client := &http.Client{Timeout: 10 * time.Second}
// resp, err := client.Do(req)
// if err != nil {
// log.Printf("[ERROR] Failed to send HTTP request: %s", err)
// return err
// }
// defer resp.Body.Close()
log.Println("ZIP file downloaded successfully.")
// if resp.StatusCode != 200 {
// log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode)
// return fmt.Errorf("unexpected HTTP status code: %d", resp.StatusCode)
// }
err = extractZIP("files.zip", "sigma_rules")
if err != nil {
return err
}
// return nil
// }
destPath := "/var/lib/tenzir/sigma_rules"
err = copyToTenzir("sigma_rules", destPath)
if err != nil {
return err
}
log.Println("Files copied to container successfully.")
checkDisabledDirCmd := exec.Command("docker", "exec", "tenzir-node", "sh", "-c", "test -d /var/lib/tenzir/disabled_rules")
if err := checkDisabledDirCmd.Run(); err != nil {
if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
// Directory does not exist, nothing to do
log.Println("[DEBUG] /var/lib/tenzir/disabled_rules does not exist.")
return nil
}
return fmt.Errorf("error checking disabled rules directory: %v", err)
}
// List files in /var/lib/tenzir/disabled_rules
listFilesCmd := exec.Command("docker", "exec", "tenzir-node", "sh", "-c", "ls /var/lib/tenzir/disabled_rules")
output, err := listFilesCmd.CombinedOutput()
if err != nil {
return fmt.Errorf("error listing files in disabled rules directory: %v, output: %s", err, output)
}
files := strings.Split(strings.TrimSpace(string(output)), "\n")
for _, file := range files {
disabledFilePath := fmt.Sprintf("/var/lib/tenzir/sigma_rules/%s", file)
checkFileCmd := exec.Command("docker", "exec", "tenzir-node", "sh", "-c", fmt.Sprintf("test -f %s", disabledFilePath))
if err := checkFileCmd.Run(); err != nil {
if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
log.Printf("[ERROR] File does not exist: %s, moving on.\n", disabledFilePath)
continue
}
return fmt.Errorf("error checking file: %v", err)
}
deleteFileCmd := exec.Command("docker", "exec", "-u", "root", "tenzir-node", "sh", "-c", fmt.Sprintf("rm -f %s", disabledFilePath))
if err := deleteFileCmd.Run(); err != nil {
return fmt.Errorf("error deleting file: %v", err)
}
log.Printf("[INFO] Deleted file: %s\n", disabledFilePath)
}
return nil
}
func extractZIP(zipFile, destDir string) error {
r, err := zip.OpenReader(zipFile)
if err != nil {
return err
}
defer r.Close()
if err := os.MkdirAll(destDir, 0755); err != nil {
return err
}
for _, f := range r.File {
err := extractFile(f, destDir)
if err != nil {
return err
}
}
return nil
}
func extractFile(f *zip.File, destDir string) error {
rc, err := f.Open()
if err != nil {
return err
}
defer rc.Close()
path := filepath.Join(destDir, f.Name)
out, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, f.Mode())
if err != nil {
return err
}
defer out.Close()
_, err = io.Copy(out, rc)
return err
}
func copyToTenzir(srcPath, destPath string) error {
containerName := "tenzir-node"
checkCmd := exec.Command("docker", "exec", containerName, "test", "-d", destPath)
if err := checkCmd.Run(); err == nil {
rmCmd := exec.Command("docker", "exec", "-u", "root", containerName, "rm", "-rf", destPath)
if err := rmCmd.Run(); err != nil {
return fmt.Errorf("error removing existing directory in container: %v", err)
}
}
cpCmd := exec.Command("docker", "cp", srcPath, fmt.Sprintf("%s:%s", containerName, destPath))
var out bytes.Buffer
cpCmd.Stdout = &out
cpCmd.Stderr = &out
err := cpCmd.Run()
if err != nil {
return fmt.Errorf("error copying files: %v, output: %s", err, out.String())
}
return nil
}
func removeAllFiles() error {
containerName := "tenzir-node"
sigmaPath := "/var/lib/tenzir/sigma_rules/*"
checkCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("ls %s", sigmaPath))
checkOutput, checkErr := checkCmd.CombinedOutput()
if checkErr != nil {
if strings.Contains(string(checkOutput), "No such file or directory") {
return nil // nothing to delete
}
return fmt.Errorf("error checking files: %v, output: %s", checkErr, checkOutput)
}
cmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("rm -rf %s", sigmaPath))
output, err := cmd.CombinedOutput()
if err != nil {
return fmt.Errorf("error removing files: %v, output: %s", err, output)
}
return nil
}
func removeFile(fileName string) error {
containerName := "tenzir-node"
srcPath := fmt.Sprintf("/var/lib/tenzir/sigma_rules/%s", fileName)
checkSrcCmd := exec.Command("docker", "exec", containerName, "sh", "-c", fmt.Sprintf("test -f %s", srcPath))
if err := checkSrcCmd.Run(); err != nil {
// If the file does not exist, simply return nil
if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
log.Printf("[ERROR] No such file: %s, nothing to delete\n", srcPath)
return nil
}
return fmt.Errorf("error checking source file: %v", err)
}
return removePath(containerName, srcPath)
}
func removePath(containerName, path string) error {
rmCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("rm -rf %s", path))
output, err := rmCmd.CombinedOutput()
if err != nil {
return fmt.Errorf("error removing path: %v, output: %s", err, output)
}
return nil
}
func sendTenzirHealthStatus() error {
var status string
url := fmt.Sprintf("%s/api/v1/detection/siem/node_health", baseUrl)
err := checkTenzirNode()
if err != nil {
return err
} else {
status = "active"
}
forwardMethod := "POST"
payload := map[string]interface{}{
"status": status,
}
payloadBytes, err := json.Marshal(payload)
if err != nil {
log.Printf("[ERROR] Failed to marshal payload: %s", err)
return err
}
forwardData := bytes.NewBuffer(payloadBytes)
req, err := http.NewRequest(
forwardMethod,
url,
forwardData,
)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err
}
req.Header.Set("Content-Type", "application/json")
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
if err != nil {
log.Printf("[ERROR] Failed to send HTTP request: %s", err)
return err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode)
return fmt.Errorf("unexpected HTTP status code: %d", resp.StatusCode)
}
return nil
}
func disableRule(fileName string) error {
containerName := "tenzir-node"
srcPath := fmt.Sprintf("/var/lib/tenzir/sigma_rules/%s", fileName)
destDir := "/var/lib/tenzir/disabled_rules"
destPath := fmt.Sprintf("%s/%s", destDir, fileName)
checkSrcCmd := exec.Command("docker", "exec", containerName, "sh", "-c", fmt.Sprintf("test -f %s", srcPath))
if err := checkSrcCmd.Run(); err != nil {
if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
fmt.Printf("File does not exist: %s\n", srcPath)
return nil // Nothing to disable
}
return fmt.Errorf("error checking source file: %v", err)
}
checkDestDirCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mkdir -p %s", destDir))
if err := checkDestDirCmd.Run(); err != nil {
return fmt.Errorf("error ensuring destination directory exists: %v", err)
}
moveCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mv %s %s", srcPath, destPath))
if err := moveCmd.Run(); err != nil {
return fmt.Errorf("error moving file: %v", err)
}
fmt.Printf("File %s moved to %s successfully.\n", fileName, destDir)
return nil
}
func enableRule(fileName string) error {
containerName := "tenzir-node"
srcPath := fmt.Sprintf("/var/lib/tenzir/disabled_rules/%s", fileName)
destDir := "/var/lib/tenzir/sigma_rules"
destPath := fmt.Sprintf("%s/%s", destDir, fileName)
checkSrcCmd := exec.Command("docker", "exec", containerName, "sh", "-c", fmt.Sprintf("test -f %s", srcPath))
if err := checkSrcCmd.Run(); err != nil {
if exitErr, ok := err.(*exec.ExitError); ok && exitErr.ExitCode() == 1 {
fmt.Printf("File does not exist: %s\n", srcPath)
return nil // Nothing to enable
}
return fmt.Errorf("error checking source file: %v", err)
}
checkDestDirCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mkdir -p %s", destDir))
if err := checkDestDirCmd.Run(); err != nil {
return fmt.Errorf("error ensuring destination directory exists: %v", err)
}
moveCmd := exec.Command("docker", "exec", "-u", "root", containerName, "sh", "-c", fmt.Sprintf("mv %s %s", srcPath, destPath))
if err := moveCmd.Run(); err != nil {
return fmt.Errorf("error moving file: %v", err)
}
fmt.Printf("File %s moved to %s successfully.\n", fileName, destDir)
return nil
}
// Is this ok to do with Docker? idk :)
func getRunningWorkers(ctx context.Context, workerTimeout int) int {