Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions webapp/app/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,16 @@ def prepare_export_jobs_folder(export_jobs_folder):
db = SQLA(app)
appbuilder = AppBuilder(app, db.session, security_manager_class=CustomSecurityManager)

from flask_limiter import Limiter
from flask_limiter.util import get_remote_address

limiter = Limiter(
get_remote_address,
app=app,
default_limits=["500/hour"],
storage_uri="memory://",
)

# Import views and APIs
# pylint: disable=C0413
from . import views
Expand Down
84 changes: 50 additions & 34 deletions webapp/app/apis.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
import os
import json
import logging
import threading
import time
from flask_appbuilder.models.sqla.interface import SQLAInterface
from flask_appbuilder import ModelRestApi, has_access
Expand Down Expand Up @@ -119,8 +120,12 @@ def validate_agent_key(self, value, **_kwargs):
)
if check_password_hash(agentkey, value):
return True # If the key is existing
except NoResultFound as error:
raise ValidationError("Invalid AGENT_KEY") from error
except NoResultFound:
# Dummy comparison to prevent timing oracle on key index existence
check_password_hash(
"pbkdf2:sha256:260000$plum_timing_dummy$" + "a" * 64, value
)
raise ValidationError("Invalid AGENT_KEY")
raise ValidationError("Invalid AGENT_KEY")


Expand Down Expand Up @@ -259,41 +264,45 @@ def _get_available_job_priorities():
return priorities


_priority_lock = threading.Lock()


def _select_weighted_priority(available_priorities):
"""
Smooth weighted round-robin over currently non-empty priority queues.
"""
available_priorities = set(available_priorities or [])
if not available_priorities:
return []

state = db.app.config.setdefault(
"priority_weighted_round_robin",
{priority: 0 for priority in PRIORITY_WEIGHTS},
)
for priority in PRIORITY_WEIGHTS:
state.setdefault(priority, 0)
if priority not in available_priorities:
state[priority] = 0

total_weight = 0
for priority in sorted(available_priorities, reverse=True):
weight = PRIORITY_WEIGHTS[priority]
state[priority] += weight
total_weight += weight

selected_priority = max(
available_priorities,
key=lambda priority: (state[priority], priority),
)
state[selected_priority] -= total_weight
db.app.config["priority_weighted_round_robin"] = state
with _priority_lock:
available_priorities = set(available_priorities or [])
if not available_priorities:
return []

state = db.app.config.setdefault(
"priority_weighted_round_robin",
{priority: 0 for priority in PRIORITY_WEIGHTS},
)
for priority in PRIORITY_WEIGHTS:
state.setdefault(priority, 0)
if priority not in available_priorities:
state[priority] = 0

total_weight = 0
for priority in sorted(available_priorities, reverse=True):
weight = PRIORITY_WEIGHTS[priority]
state[priority] += weight
total_weight += weight

selected_priority = max(
available_priorities,
key=lambda priority: (state[priority], priority),
)
state[selected_priority] -= total_weight
db.app.config["priority_weighted_round_robin"] = state

fallback_priorities = sorted(
available_priorities - {selected_priority},
reverse=True,
)
return [selected_priority] + fallback_priorities
fallback_priorities = sorted(
available_priorities - {selected_priority},
reverse=True,
)
return [selected_priority] + fallback_priorities


class PublicTargetsApi(ModelRestApi):
Expand Down Expand Up @@ -605,11 +614,18 @@ def sndjobs(self):
return self.response(404, message="job not found")

if job_bot.finished and not job_bot.active:
if job_bot.bot_id != submitting_bot.id:
logger.warning(
"Bot %s tried to submit already-completed job %s owned by bot_id %s",
botinfo.get("UID"),
job_bot.uid,
job_bot.bot_id,
)
return self.response(403, message="forbidden")
logger.info(
"Bot %s resubmitted already completed job %s assigned to bot_id %s; returning idempotent success",
"Bot %s resubmitted already completed job %s; returning idempotent success",
botinfo.get("UID"),
job_bot.uid,
job_bot.bot_id,
)
submitting_bot.running = False
submitting_bot.last_seen = utcnow_naive()
Expand Down