Skip to content

Multiple user limits #118

New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Closed
wants to merge 7 commits into from
Closed
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
23 changes: 3 additions & 20 deletions cads_broker/database.py
Original file line number Diff line number Diff line change
Expand Up @@ -427,26 +427,9 @@ def count_users(status: str, entry_point: str, session: sa.orm.Session) -> int:
)


def update_dismissed_requests(session: sa.orm.Session) -> Iterable[str]:
stmt_dismissed = (
sa.update(SystemRequest)
.where(SystemRequest.status == "dismissed")
.returning(SystemRequest.request_uid)
.values(status="failed", response_error={"reason": "dismissed request"})
)
dismissed_uids = session.scalars(stmt_dismissed).fetchall()
session.execute( # type: ignore
sa.insert(Events),
map(
lambda x: {
"request_uid": x,
"message": DISMISSED_MESSAGE,
"event_type": "user_visible_error",
},
dismissed_uids,
),
)
return dismissed_uids
def get_dismissed_requests(session: sa.orm.Session) -> Iterable[SystemRequest]:
stmt_dismissed = sa.select(SystemRequest).where(SystemRequest.status == "dismissed")
return session.scalars(stmt_dismissed).fetchall()


def get_events_from_request(
Expand Down
30 changes: 17 additions & 13 deletions cads_broker/dispatcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -170,9 +170,9 @@ def values(self) -> Iterable[Any]:
with self._lock:
return self.queue_dict.values()

def pop(self, key: str) -> Any:
def pop(self, key: str, default=None) -> Any:
with self._lock:
return self.queue_dict.pop(key, None)
return self.queue_dict.pop(key, default)

def len(self) -> int:
with self._lock:
Expand All @@ -192,9 +192,12 @@ def __init__(self, number_of_workers) -> None:
if os.path.exists(self.rules_path):
self.rules = self.rules_path
else:
parser = QoS.RulesParser(io.StringIO(os.getenv("DEFAULT_RULES", "")))
logger.info("rules file not found", rules_path=self.rules_path)
parser = QoS.RulesParser(
io.StringIO(os.getenv("DEFAULT_RULES", "")), logger=logger
)
self.rules = QoS.RuleSet()
parser.parse_rules(self.rules, self.environment)
parser.parse_rules(self.rules, self.environment, raise_exception=False)


@attrs.define
Expand Down Expand Up @@ -235,6 +238,7 @@ def from_address(
qos_config.rules,
qos_config.environment,
rules_hash=rules_hash,
logger=logger,
)
with session_maker_write() as session:
qos.environment.set_session(session)
Expand Down Expand Up @@ -336,16 +340,16 @@ def sync_database(self, session: sa.orm.Session) -> None:
- If the task is not in the dask scheduler, it is re-queued.
This behaviour can be changed with an environment variable.
"""
# the retrieve API sets the status to "dismissed", here the broker deletes the request
# this is to better control the status of the QoS
dismissed_uids = db.update_dismissed_requests(session)
for uid in dismissed_uids:
if future := self.futures.pop(uid, None):
# the retrieve API sets the status to "dismissed",
# here the broker fixes the QoS and queue status accordingly
dismissed_requests = db.get_dismissed_requests(session)
for request in dismissed_requests:
if future := self.futures.pop(request.request_uid, None):
future.cancel()
if dismissed_uids:
self.queue.reset()
self.qos.reload_rules(session)
db.reset_qos_rules(session, self.qos)
self.qos.notify_end_of_request(
request, session, scheduler=self.internal_scheduler
)
self.queue.pop(request.request_uid, None)
session.commit()

statement = sa.select(db.SystemRequest).where(
Expand Down
1 change: 0 additions & 1 deletion cads_broker/entry_points.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,6 @@ def add_dummy_requests(
entry_point="cads_adaptors:DummyAdaptor",
)
session.add(request)
if i % 100 == 0:
session.commit()
session.commit()

Expand Down
3 changes: 2 additions & 1 deletion cads_broker/expressions/Parser.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,14 +41,15 @@ class Parser:
This class must be sub-classed.
"""

def __init__(self, path, comments=True):
def __init__(self, path, logger, comments=True):
if isinstance(path, str):
self.reader = Reader(open(path))
else:
self.reader = Reader(path)
self.comments = comments
self.eof = False
self.line = 0
self.logger = logger

def read(self):
if self.eof:
Expand Down
47 changes: 27 additions & 20 deletions cads_broker/expressions/RulesParser.py
Original file line number Diff line number Diff line change
Expand Up @@ -327,7 +327,7 @@ def parse(self):

return result

def parse_rules(self, rules, environment):
def parse_rules(self, rules, environment, raise_exception=True):
"""Parse the text provided in the constructor.

Args:
Expand All @@ -340,22 +340,29 @@ def parse_rules(self, rules, environment):
ParserError: [description]
"""
while self.peek():
ident = self.parse_ident()

if ident == "limit":
self.parse_global_limit(rules, environment)
continue

if ident == "priority":
self.parse_priority(rules, environment)
continue

if ident == "permission":
self.parse_permission(rules, environment)
continue

if ident == "user":
self.parse_user_limit(rules, environment)
continue

raise ParserError(f"Unknown rule: '{ident}'", self.line + 1)
try:
ident = self.parse_ident()

if ident == "limit":
self.parse_global_limit(rules, environment)
continue

if ident == "priority":
self.parse_priority(rules, environment)
continue

if ident == "permission":
self.parse_permission(rules, environment)
continue

if ident == "user":
self.parse_user_limit(rules, environment)
continue

raise ParserError(f"Unknown rule: '{ident}'", self.line + 1)
except ParserError as e:
if raise_exception:
raise e
else:
self.logger.info(e)
return
28 changes: 14 additions & 14 deletions cads_broker/qos/QoS.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,13 @@ def wrapped(self, *args, **kwargs):


class QoS:
def __init__(self, rules, environment, rules_hash):
def __init__(self, rules, environment, rules_hash, logger):
self.lock = threading.RLock()

self.rules_hash = rules_hash

self.environment = environment
self.logger = logger
# The list of active requests

# Cache associating Request and their Properties
Expand All @@ -53,13 +54,13 @@ def __init__(self, rules, environment, rules_hash):
def read_rules(self):
"""Read the rule files and populate the rule_set."""
# Create a parser to parse the rules file
parser = RulesParser(self.path)
parser = RulesParser(self.path, logger=self.logger)

# The rules will be stored in self.rules
self.rules = RuleSet()

# Parse the rules
parser.parse_rules(self.rules, self.environment)
parser.parse_rules(self.rules, self.environment, raise_exception=False)

# Print the rules
self.rules.dump()
Expand Down Expand Up @@ -156,9 +157,9 @@ def _properties(self, request, session):
properties.limits.append(rule)

# Add per-user limits
limit = self.user_limit(request)
if limit is not None:
properties.limits.append(limit)
limits = self.user_limit(request)
if limits != []:
properties.limits.extend(limits)

# Add priorities and compute starting priority
priority = 0
Expand Down Expand Up @@ -251,10 +252,10 @@ def user_limit(self, request):
"""Return the per-user limit for the user associated with the request."""
user = request.user_uid

limit = self.per_user_limits.get(user)
if limit is not None:
print(user, limit)
return limit
limits = self.per_user_limits.get(user, [])
if limits != []:
print(user, limits)
return limits

for limit in self.rules.user_limits:
if limit.match(request):
Expand All @@ -263,10 +264,9 @@ def user_limit(self, request):
user otherwise all users will share that limit
"""
limit = limit.clone()
self.per_user_limits[user] = limit
return limit
return None
# raise Exception(f"Not rules matching user '{user}'")
limits.append(limit)
self.per_user_limits[user] = limits
return limits

@locked
def pick(self, queue, session):
Expand Down
5 changes: 4 additions & 1 deletion tests/test_01_expressions.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import io
import logging

from cads_broker import Environment
from cads_broker.expressions import FunctionFactory
Expand All @@ -14,6 +15,8 @@
lambda context, *args: context.request.adaptor,
)

logger = logging.getLogger("test")


class TestRequest:
user_uid = "david"
Expand All @@ -30,7 +33,7 @@ class TestRequest:


def compile(text):
parser = RulesParser(io.StringIO(text))
parser = RulesParser(io.StringIO(text), logger=logger)
return parser.parse()


Expand Down
38 changes: 24 additions & 14 deletions tests/test_02_database.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,11 +14,13 @@


class MockRule:
def __init__(self, name, conclusion, info, condition):
def __init__(self, name, conclusion, info, condition, queued=[], running=0):
self.name = name
self.conclusion = conclusion
self.info = info
self.condition = condition
self.queued = queued
self.value = running

def evaluate(self, request):
return 10
Expand Down Expand Up @@ -517,8 +519,12 @@ def test_add_qos_rule(session_obj: sa.orm.sessionmaker) -> None:


def test_add_request_qos_status(session_obj: sa.orm.sessionmaker) -> None:
rule1 = MockRule("name1", "conclusion1", "info1", "condition1")
rule2 = MockRule("name2", "conclusion2", "info2", "condition2")
rule1 = MockRule(
"name1", "conclusion1", "info1", "condition1", queued=list(range(5))
)
rule2 = MockRule(
"name2", "conclusion2", "info2", "condition2", queued=list(range(1))
)
adaptor_properties = mock_config()
request = mock_system_request(adaptor_properties_hash=adaptor_properties.hash)
request_uid = request.request_uid
Expand All @@ -535,16 +541,18 @@ def test_add_request_qos_status(session_obj: sa.orm.sessionmaker) -> None:
with session_obj() as session:
request = db.get_request(request_uid, session=session)
assert db.get_qos_status_from_request(request) == {
"name1": [
{"info": "info1", "queued": 5 + 1, "running": 0, "conclusion": "10"}
],
"name1": [{"info": "info1", "queued": 5, "running": 0, "conclusion": "10"}],
"name2": [{"info": "info2", "queued": 1, "running": 0, "conclusion": "10"}],
}


def test_delete_request_qos_status(session_obj: sa.orm.sessionmaker) -> None:
rule1 = MockRule("name1", "conclusion1", "info1", "condition1")
rule2 = MockRule("name2", "conclusion2", "info2", "condition2")
rule1 = MockRule(
"name1", "conclusion1", "info1", "condition1", queued=list(range(5)), running=2
)
rule2 = MockRule(
"name2", "conclusion2", "info2", "condition2", queued=list(range(3)), running=2
)
adaptor_properties = mock_config()
request = mock_system_request(adaptor_properties_hash=adaptor_properties.hash)
request_uid = request.request_uid
Expand Down Expand Up @@ -574,14 +582,16 @@ def test_delete_request_qos_status(session_obj: sa.orm.sessionmaker) -> None:
rule1 = db.get_qos_rule(str(rule1.__hash__()), session=session)
rule2 = db.get_qos_rule(str(rule2.__hash__()), session=session)
assert rule1.queued == rule1_queued
assert rule1.running == rule1_running + 1
assert rule2.queued == rule2_queued
assert rule2.running == rule2_running + 1


def test_decrement_qos_rule_running(session_obj: sa.orm.sessionmaker) -> None:
rule1 = MockRule("name1", "conclusion1", "info1", "condition1")
rule2 = MockRule("name2", "conclusion2", "info2", "condition2")
rule1 = MockRule(
"name1", "conclusion1", "info1", "condition1", queued=list(range(5)), running=2
)
rule2 = MockRule(
"name2", "conclusion2", "info2", "condition2", queued=list(range(3)), running=4
)
rule1_queued = 5
rule1_running = 2
rule2_queued = 3
Expand All @@ -602,11 +612,11 @@ def test_decrement_qos_rule_running(session_obj: sa.orm.sessionmaker) -> None:
with session_obj() as session:
assert (
db.get_qos_rule(str(rule1.__hash__()), session=session).running
== rule1_running - 1
== rule1_running
)
assert (
db.get_qos_rule(str(rule2.__hash__()), session=session).running
== rule2_running - 1
== rule2_running
)


Expand Down
4 changes: 3 additions & 1 deletion tests/test_03_qos.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import io
import logging

from cads_broker import Environment
from cads_broker.expressions import FunctionFactory
Expand All @@ -13,6 +14,7 @@
"adaptor",
lambda context, *args: context.request.adaptor,
)
logger = logging.getLogger("test")


class TestRequest:
Expand All @@ -30,7 +32,7 @@ class TestRequest:


def compile(text):
parser = RulesParser(io.StringIO(text))
parser = RulesParser(io.StringIO(text), logger=logger)
rules = RuleSet()
parser.parse_rules(rules, environment)
return rules
Expand Down
Loading
Loading