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
2 changes: 1 addition & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ test:
@# Note that extensive integration tests are also run in the bundle repo,
@# for both aw-server and aw-server-rust, but without code coverage.
python -c 'import aw_server'
python -m pytest tests/test_server.py tests/test_profile.py tests/test_profile_config.py
python -m pytest tests/test_server.py tests/test_profile.py tests/test_profile_config.py tests/test_performance.py

typecheck:
python -m mypy aw_server tests --ignore-missing-imports
Expand Down
234 changes: 156 additions & 78 deletions aw_server/api.py

Large diffs are not rendered by default.

192 changes: 192 additions & 0 deletions aw_server/query_cache.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,192 @@
"""Bounded query results with dependencies on the actual datastore reads."""

import copy
import sys
from collections import OrderedDict
from contextlib import contextmanager
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from threading import RLock
from time import monotonic
from typing import Any, List, Optional, Tuple

Read = Tuple[str, Optional[datetime], Optional[datetime]]


def _size(value, seen=None):
"""Account for retained Python objects, including keys and nested payloads."""
if seen is None:
seen = set()
if id(value) in seen:
return 0
seen.add(id(value))
size = sys.getsizeof(value)
if isinstance(value, dict):
size += sum(_size(k, seen) + _size(v, seen) for k, v in value.items())
elif isinstance(value, (list, tuple, set)):
size += sum(_size(v, seen) for v in value)
return size


@dataclass
class Change:
span: Optional[Tuple[datetime, datetime]] = None


@dataclass
class Entry:
result: Any
reads: List[Read]
metadata: bool
expires: float
size: int


class QueryCache:
def __init__(self, max_entries=128, max_bytes=8 * 1024 * 1024, ttl=300):
self.max_entries = max_entries
self.max_bytes = max_bytes
self.ttl = ttl
self.entries: OrderedDict = OrderedDict()
self.bytes = 0
self.revision = 0
self.writers = 0
self.lock = RLock()

def _remove(self, key):
self.bytes -= self.entries.pop(key).size

def lookup(self, key):
with self.lock:
for expired in [
k for k, e in self.entries.items() if e.expires <= monotonic()
]:
self._remove(expired)
if self.writers:
return False, None, None
entry = self.entries.get(key)
if entry is not None:
self.entries.move_to_end(key)
return True, copy.deepcopy(entry.result), self.revision
return False, None, self.revision

def store(self, key, result, tracker, revision):
if revision is None or not tracker.cacheable or self.max_entries <= 0:
return
# Do the potentially expensive copy outside the cache lock.
size = _size((key, result, tracker.reads))
if size > self.max_bytes:
return
snapshot = copy.deepcopy(result)
with self.lock:
if self.writers or self.revision != revision:
return
if key in self.entries:
self._remove(key)
while self.entries and (
len(self.entries) >= self.max_entries
or self.bytes + size > self.max_bytes
):
self._remove(next(iter(self.entries)))
self.entries[key] = Entry(
snapshot,
list(tracker.reads),
tracker.metadata,
monotonic() + self.ttl,
size,
)
self.bytes += size

@contextmanager
def mutation(self, bucket_id, metadata=False):
change = Change()
with self.lock:
self.writers += 1
self.revision += 1
try:
yield change
except BaseException:
# A failed bulk operation can still have made partial changes.
change.span = None
raise
finally:
with self.lock:
for key, entry in list(self.entries.items()):
affected = metadata and entry.metadata
for bid, start, end in entry.reads:
if bid != bucket_id:
continue
if metadata or change.span is None:
affected = True
else:
lo, hi = sorted(change.span)
affected |= (end is None or lo <= end) and (
start is None or start <= hi
)
if affected:
self._remove(key)
self.writers -= 1
self.revision += 1


@dataclass
class ReadTracker:
datastore: Any
reads: List[Read] = field(default_factory=list)
metadata: bool = False
cacheable: bool = True

def buckets(self):
self.metadata = True
return self.datastore.buckets()

def __getitem__(self, bucket_id):
return ReadBucket(self.datastore[bucket_id], self, bucket_id)

def __getattr__(self, name):
# Custom query functions using other datastore operations must not cache
# a result whose dependencies we cannot describe.
self.cacheable = False
return getattr(self.datastore, name)


class ReadBucket:
def __init__(self, bucket, tracker, bucket_id):
self.bucket = bucket
self.tracker = tracker
self.bucket_id = bucket_id

def _record(self, starttime, endtime):
# Bucket.get rounds start down and end up to millisecond boundaries.
padding = timedelta(milliseconds=1)
try:
start = starttime - padding if starttime else None
except OverflowError:
start = None
try:
end = endtime + padding if endtime else None
except OverflowError:
end = None
self.tracker.reads.append(
(
self.bucket_id,
start,
end,
)
)

def get(self, limit=-1, starttime=None, endtime=None):
self._record(starttime, endtime)
return self.bucket.get(limit, starttime, endtime)

def get_eventcount(self, starttime=None, endtime=None):
self._record(starttime, endtime)
return self.bucket.get_eventcount(starttime, endtime)

def metadata(self):
self.tracker.metadata = True
return self.bucket.metadata()

def __getattr__(self, name):
self.tracker.cacheable = False
return getattr(self.bucket, name)
43 changes: 26 additions & 17 deletions aw_server/rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@
Blueprint,
current_app,
jsonify,
make_response,
Response,
stream_with_context,
request,
)
from flask_restx import Api, Resource, fields
Expand Down Expand Up @@ -214,9 +215,9 @@ def get(self, bucket_id):
def post(self, bucket_id):
data = request.get_json()
logger.debug(
"Received post request for event in bucket '{}' and data: {}".format(
bucket_id, data
)
"Received post request for event in bucket '%s' and data: %s",
bucket_id,
data,
)

if isinstance(data, dict):
Expand Down Expand Up @@ -251,7 +252,9 @@ class EventResource(Resource):
@copy_doc(ServerAPI.get_event)
def get(self, bucket_id: str, event_id: int):
logger.debug(
f"Received get request for event with id '{event_id}' in bucket '{bucket_id}'"
"Received get request for event with id '%s' in bucket '%s'",
event_id,
bucket_id,
)
event = current_app.api.get_event(bucket_id, event_id)
if event:
Expand All @@ -262,9 +265,9 @@ def get(self, bucket_id: str, event_id: int):
@copy_doc(ServerAPI.delete_event)
def delete(self, bucket_id: str, event_id: int):
logger.debug(
"Received delete request for event with id '{}' in bucket '{}'".format(
event_id, bucket_id
)
"Received delete request for event with id '%s' in bucket '%s'",
event_id,
bucket_id,
)
success = current_app.api.delete_event(bucket_id, event_id)
return {"success": success}, 200
Expand Down Expand Up @@ -316,15 +319,19 @@ def post(self, bucket_id):
class QueryResource(Resource):
# TODO Docs
@api.expect(query, validate=True)
@api.param("name", "Name of the query (required if using cache)")
@api.param("name", "Name of the query")
@api.param("cache", "Cache query results (default: 1; set to 0 to bypass)")
def post(self):
name = ""
if "name" in request.args:
name = request.args["name"]
query = request.get_json()
cache_arg = request.args.get("cache", "1").lower()
if cache_arg not in ("0", "1", "false", "true"):
raise BadRequest("InvalidParameter", "cache must be 0, 1, false, or true")
try:
result = current_app.api.query2(
name, query["query"], query["timeperiods"], False
name, query["query"], query["timeperiods"], cache_arg in ("1", "true")
)
return jsonify(result)
except QueryException as qe:
Expand All @@ -340,9 +347,10 @@ class ExportAllResource(Resource):
@api.doc(model=buckets_export)
@copy_doc(ServerAPI.export_all)
def get(self):
buckets_export = current_app.api.export_all()
payload = {"buckets": buckets_export}
response = make_response(json.dumps(payload))
response = Response(
stream_with_context(current_app.api.stream_export()),
mimetype="application/json",
)
filename = "aw-buckets-export.json"
response.headers["Content-Disposition"] = "attachment; filename={}".format(
filename
Expand All @@ -356,10 +364,11 @@ class BucketExportResource(Resource):
@api.doc(model=buckets_export)
@copy_doc(ServerAPI.export_bucket)
def get(self, bucket_id):
bucket_export = current_app.api.export_bucket(bucket_id)
payload = {"buckets": {bucket_export["id"]: bucket_export}}
response = make_response(json.dumps(payload))
filename = "aw-bucket-export_{}.json".format(bucket_export["id"])
response = Response(
stream_with_context(current_app.api.stream_export(bucket_id)),
mimetype="application/json",
)
filename = "aw-bucket-export_{}.json".format(bucket_id)
response.headers["Content-Disposition"] = "attachment; filename={}".format(
filename
)
Expand Down
24 changes: 14 additions & 10 deletions poetry.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ aw-server = "aw_server:main"

[tool.poetry.dependencies]
python = "^3.8"
aw-core = "^0.5.18"
aw-core = { git = "https://github.com/0xbrayo/aw-core.git", rev = "a794efb9639b03c243c47c85819f701defbd2591" }
aw-client = "^0.5.8"
flask = "^2.2"
flask-restx = "^1.0.3"
Expand Down
Loading
Loading