Changed apps to use milliseconds instead of seconds for timestamps to be more accurate
This commit is contained in:
@@ -70,7 +70,7 @@ SHUFFLE_SWARM_BRIDGE_DEFAULT_MTU=1500 # 1500 by default
|
|||||||
# Used for auto-cleanup of containers. REALLY important at scale. Set to false to see all container info.
|
# Used for auto-cleanup of containers. REALLY important at scale. Set to false to see all container info.
|
||||||
SHUFFLE_MEMCACHED=
|
SHUFFLE_MEMCACHED=
|
||||||
SHUFFLE_CONTAINER_AUTO_CLEANUP=true
|
SHUFFLE_CONTAINER_AUTO_CLEANUP=true
|
||||||
SHUFFLE_ORBORUS_EXECUTION_CONCURRENCY=7 # The amount of concurrent executions Orborus can handle. This is a soft limit, but it's recommended to keep it low.
|
SHUFFLE_ORBORUS_EXECUTION_CONCURRENCY=5 # The amount of concurrent executions Orborus can handle. This is a soft limit, but it's recommended to keep it low.
|
||||||
SHUFFLE_HEALTHCHECK_DISABLED=false
|
SHUFFLE_HEALTHCHECK_DISABLED=false
|
||||||
SHUFFLE_ELASTIC=true
|
SHUFFLE_ELASTIC=true
|
||||||
SHUFFLE_LOGS_DISABLED=false
|
SHUFFLE_LOGS_DISABLED=false
|
||||||
|
|||||||
@@ -302,9 +302,11 @@ class AppBase:
|
|||||||
self.authorization = os.getenv("AUTHORIZATION", "")
|
self.authorization = os.getenv("AUTHORIZATION", "")
|
||||||
self.current_execution_id = os.getenv("EXECUTIONID", "")
|
self.current_execution_id = os.getenv("EXECUTIONID", "")
|
||||||
self.full_execution = os.getenv("FULL_EXECUTION", "")
|
self.full_execution = os.getenv("FULL_EXECUTION", "")
|
||||||
self.start_time = int(time.time())
|
|
||||||
self.result_wrapper_count = 0
|
self.result_wrapper_count = 0
|
||||||
|
|
||||||
|
# Make start time with milliseconds
|
||||||
|
self.start_time = int(time.time_ns())
|
||||||
|
|
||||||
self.action_result = {
|
self.action_result = {
|
||||||
"action": self.action,
|
"action": self.action,
|
||||||
"authorization": self.authorization,
|
"authorization": self.authorization,
|
||||||
@@ -312,7 +314,7 @@ class AppBase:
|
|||||||
"result": f"",
|
"result": f"",
|
||||||
"started_at": self.start_time,
|
"started_at": self.start_time,
|
||||||
"status": "",
|
"status": "",
|
||||||
"completed_at": int(time.time()),
|
"completed_at": int(time.time_ns()),
|
||||||
}
|
}
|
||||||
|
|
||||||
if isinstance(self.action, str):
|
if isinstance(self.action, str):
|
||||||
@@ -468,7 +470,7 @@ class AppBase:
|
|||||||
|
|
||||||
# Try it with some magic
|
# Try it with some magic
|
||||||
|
|
||||||
action_result["completed_at"] = int(time.time())
|
action_result["completed_at"] = int(time.time_ns())
|
||||||
self.logger.info(f"""[DEBUG] Inside Send result with status {action_result["status"]}""")
|
self.logger.info(f"""[DEBUG] Inside Send result with status {action_result["status"]}""")
|
||||||
#if isinstance(action_result,
|
#if isinstance(action_result,
|
||||||
|
|
||||||
@@ -1011,7 +1013,7 @@ class AppBase:
|
|||||||
"result": f"All {len(param_multiplier)} values were non-unique",
|
"result": f"All {len(param_multiplier)} values were non-unique",
|
||||||
"started_at": self.start_time,
|
"started_at": self.start_time,
|
||||||
"status": "SKIPPED",
|
"status": "SKIPPED",
|
||||||
"completed_at": int(time.time()),
|
"completed_at": int(time.time_ns()),
|
||||||
}
|
}
|
||||||
|
|
||||||
self.send_result(self.action_result, {"Content-Type": "application/json", "Authorization": "Bearer %s" % self.authorization}, "/api/v1/streams")
|
self.send_result(self.action_result, {"Content-Type": "application/json", "Authorization": "Bearer %s" % self.authorization}, "/api/v1/streams")
|
||||||
@@ -1457,7 +1459,7 @@ class AppBase:
|
|||||||
"authorization": self.authorization,
|
"authorization": self.authorization,
|
||||||
"execution_id": self.current_execution_id,
|
"execution_id": self.current_execution_id,
|
||||||
"result": "",
|
"result": "",
|
||||||
"started_at": int(time.time()),
|
"started_at": int(time.time_ns()),
|
||||||
"status": "EXECUTING"
|
"status": "EXECUTING"
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2469,7 +2471,7 @@ class AppBase:
|
|||||||
self.action_result["result"] = f"Failed to parse LiquidPy: {error_msg}"
|
self.action_result["result"] = f"Failed to parse LiquidPy: {error_msg}"
|
||||||
print("[WARNING] Failed to set LiquidPy result")
|
print("[WARNING] Failed to set LiquidPy result")
|
||||||
|
|
||||||
self.action_result["completed_at"] = int(time.time())
|
self.action_result["completed_at"] = int(time.time_ns())
|
||||||
self.send_result(self.action_result, headers, stream_path)
|
self.send_result(self.action_result, headers, stream_path)
|
||||||
|
|
||||||
self.logger.info(f"[ERROR] Sent FAILURE response to backend due to : {e}")
|
self.logger.info(f"[ERROR] Sent FAILURE response to backend due to : {e}")
|
||||||
@@ -3071,7 +3073,7 @@ class AppBase:
|
|||||||
self.logger.info("Failed one or more branch conditions.")
|
self.logger.info("Failed one or more branch conditions.")
|
||||||
self.action_result["result"] = tmpresult
|
self.action_result["result"] = tmpresult
|
||||||
self.action_result["status"] = "SKIPPED"
|
self.action_result["status"] = "SKIPPED"
|
||||||
self.action_result["completed_at"] = int(time.time())
|
self.action_result["completed_at"] = int(time.time_ns())
|
||||||
|
|
||||||
self.send_result(self.action_result, headers, stream_path)
|
self.send_result(self.action_result, headers, stream_path)
|
||||||
return
|
return
|
||||||
@@ -3557,7 +3559,7 @@ class AppBase:
|
|||||||
self.logger.info("[WARNING] SHOULD STOP EXECUTION BECAUSE FIELDS AREN'T UNIQUE")
|
self.logger.info("[WARNING] SHOULD STOP EXECUTION BECAUSE FIELDS AREN'T UNIQUE")
|
||||||
self.action_result["status"] = "SKIPPED"
|
self.action_result["status"] = "SKIPPED"
|
||||||
self.action_result["result"] = f"A non-unique value was found"
|
self.action_result["result"] = f"A non-unique value was found"
|
||||||
self.action_result["completed_at"] = int(time.time())
|
self.action_result["completed_at"] = int(time.time_ns())
|
||||||
self.send_result(self.action_result, headers, stream_path)
|
self.send_result(self.action_result, headers, stream_path)
|
||||||
return
|
return
|
||||||
|
|
||||||
@@ -3885,7 +3887,7 @@ class AppBase:
|
|||||||
})
|
})
|
||||||
|
|
||||||
# Send the result :)
|
# Send the result :)
|
||||||
self.action_result["completed_at"] = int(time.time())
|
self.action_result["completed_at"] = int(time.time_ns())
|
||||||
self.send_result(self.action_result, headers, stream_path)
|
self.send_result(self.action_result, headers, stream_path)
|
||||||
|
|
||||||
#try:
|
#try:
|
||||||
|
|||||||
@@ -3706,6 +3706,10 @@ func runInitEs(ctx context.Context) {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// FIXME: Add a randomized timer to avoid all schedules running at the same time
|
||||||
|
// Many are at 5 minutes / 1 hour. The point is to spread these out
|
||||||
|
// a bit instead of all of them starting at the exact same time
|
||||||
|
|
||||||
//log.Printf("Schedule: %#v", schedule)
|
//log.Printf("Schedule: %#v", schedule)
|
||||||
//log.Printf("Schedule time: every %d seconds", schedule.Seconds)
|
//log.Printf("Schedule time: every %d seconds", schedule.Seconds)
|
||||||
jobret, err := newscheduler.Every(schedule.Seconds).Seconds().NotImmediately().Run(job(schedule))
|
jobret, err := newscheduler.Every(schedule.Seconds).Seconds().NotImmediately().Run(job(schedule))
|
||||||
@@ -3984,7 +3988,7 @@ func runInitEs(ctx context.Context) {
|
|||||||
|
|
||||||
r, err := git.Clone(storer, fs, cloneOptions)
|
r, err := git.Clone(storer, fs, cloneOptions)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[WARNING] Failed loading repo into memory (init): %s", err)
|
log.Printf("[ERROR] Failed loading repo into memory (init): %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
dir, err := fs.ReadDir("")
|
dir, err := fs.ReadDir("")
|
||||||
@@ -4021,7 +4025,7 @@ func runInitEs(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
_, err = git.Clone(storer, fs, cloneOptions)
|
_, err = git.Clone(storer, fs, cloneOptions)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[WARNING] Failed loading repo %s into memory: %s", apis, err)
|
log.Printf("[ERROR] Failed loading repo %s into memory: %s", apis, err)
|
||||||
} else if err == nil && len(workflowapps) < 10 {
|
} else if err == nil && len(workflowapps) < 10 {
|
||||||
log.Printf("[INFO] Finished git clone. Looking for updates to the repo.")
|
log.Printf("[INFO] Finished git clone. Looking for updates to the repo.")
|
||||||
dir, err := fs.ReadDir("")
|
dir, err := fs.ReadDir("")
|
||||||
|
|||||||
@@ -43,6 +43,7 @@ services:
|
|||||||
- /var/run/docker.sock:/var/run/docker.sock
|
- /var/run/docker.sock:/var/run/docker.sock
|
||||||
environment:
|
environment:
|
||||||
- SHUFFLE_APP_SDK_TIMEOUT=300
|
- SHUFFLE_APP_SDK_TIMEOUT=300
|
||||||
|
- SHUFFLE_ORBORUS_EXECUTION_CONCURRENCY=5 # The amount of concurrent executions Orborus can handle.
|
||||||
#- DOCKER_HOST=tcp://docker-socket-proxy:2375
|
#- DOCKER_HOST=tcp://docker-socket-proxy:2375
|
||||||
- ENVIRONMENT_NAME=${ENVIRONMENT_NAME}
|
- ENVIRONMENT_NAME=${ENVIRONMENT_NAME}
|
||||||
- BASE_URL=http://${OUTER_HOSTNAME}:5001
|
- BASE_URL=http://${OUTER_HOSTNAME}:5001
|
||||||
|
|||||||
@@ -63,7 +63,7 @@ var sleepTime = 2
|
|||||||
|
|
||||||
// Making it work on low-end machines even during busy times :)
|
// Making it work on low-end machines even during busy times :)
|
||||||
// May cause some things to run slowly
|
// May cause some things to run slowly
|
||||||
var maxConcurrency = 3
|
var maxConcurrency = 5
|
||||||
|
|
||||||
// Timeout if something rashes
|
// Timeout if something rashes
|
||||||
var workerTimeoutEnv = os.Getenv("SHUFFLE_ORBORUS_EXECUTION_TIMEOUT")
|
var workerTimeoutEnv = os.Getenv("SHUFFLE_ORBORUS_EXECUTION_TIMEOUT")
|
||||||
|
|||||||
Reference in New Issue
Block a user