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 aw_core/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,6 @@ def assert_version(required_version: Tuple[int, ...] = (3, 5)): # pragma: no co
(
"Python version {} not supported, you need to upgrade your Python"
+ " version to at least {}."
).format(required_version)
).format(actual_version, required_version)
)
logger.debug(f"Python version: {_version_info_tuple()}")
12 changes: 10 additions & 2 deletions aw_datastore/datastore.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ def __getitem__(self, bucket_id: str) -> "Bucket":
# If this bucket doesn't have a initialized object, create it
if bucket_id not in self.bucket_instances:
# If the bucket exists in the database, create an object representation of it
if bucket_id in self.buckets():
if self.has_bucket(bucket_id):
bucket = Bucket(self, bucket_id)
self.bucket_instances[bucket_id] = bucket
else:
Expand Down Expand Up @@ -72,7 +72,12 @@ def delete_bucket(self, bucket_id: str):
del self.bucket_instances[bucket_id]
return self.storage_strategy.delete_bucket(bucket_id)

def buckets(self):
def has_bucket(self, bucket_id: str) -> bool:
return self.storage_strategy.has_bucket(bucket_id)

def buckets(self, include_last_updated: bool = False):
if include_last_updated:
return self.storage_strategy.buckets_with_last_updated()
return self.storage_strategy.buckets()


Expand All @@ -85,6 +90,9 @@ def __init__(self, datastore: Datastore, bucket_id: str) -> None:
def metadata(self) -> dict:
return self.ds.storage_strategy.get_metadata(self.bucket_id)

def iter_events(self):
return self.ds.storage_strategy.iter_events(self.bucket_id)

def get(
self,
limit: int = -1,
Expand Down
23 changes: 22 additions & 1 deletion aw_datastore/storages/abstract.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
from abc import ABCMeta, abstractmethod
from datetime import datetime
from typing import Dict, List, Optional
from typing import Dict, Iterator, List, Optional

from aw_core.models import Event

Expand All @@ -21,6 +21,27 @@ def __init__(self, testing: bool) -> None:
def buckets(self) -> Dict[str, dict]:
raise NotImplementedError

def has_bucket(self, bucket_id: str) -> bool:
return bucket_id in self.buckets()

def buckets_with_last_updated(self) -> Dict[str, dict]:
buckets = self.buckets()
for bucket_id, metadata in buckets.items():
events = self.get_events(bucket_id, 1)
if events:
metadata["last_updated"] = (
events[0].timestamp + events[0].duration
).isoformat()
return buckets

def iter_events(self, bucket_id: str) -> Iterator[Event]:
"""Iterate a bucket in the same order as an unbounded get_events call.

Disk backends override this to avoid materializing all events. Callers
must close the iterator if they stop consuming it early.
"""
yield from self.get_events(bucket_id, -1)

@abstractmethod
def create_bucket(
self,
Expand Down
12 changes: 11 additions & 1 deletion aw_datastore/storages/memory.py
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,14 @@ def buckets(self):
buckets[bucket_id] = self.get_metadata(bucket_id)
return buckets

def has_bucket(self, bucket_id: str) -> bool:
return bucket_id in self.db

def iter_events(self, bucket_id):
# Only sort references; copy each event when it is consumed.
for event in sorted(self.db[bucket_id], key=lambda e: e.timestamp)[::-1]:
yield copy.deepcopy(event)

def get_event(
self,
bucket_id: str,
Expand Down Expand Up @@ -132,7 +140,7 @@ def get_eventcount(

def get_metadata(self, bucket_id: str):
if bucket_id in self._metadata:
return self._metadata[bucket_id]
return copy.deepcopy(self._metadata[bucket_id])
else:
raise ValueError("Bucket did not exist, could not get metadata")

Expand Down Expand Up @@ -180,6 +188,8 @@ def replace(self, bucket_id, event_id, event):
event = copy.copy(event)
event.id = event_id
self.db[bucket_id][idx] = event
return True
return False

def replace_last(self, bucket_id, event):
# NOTE: This does not actually get the most recent event, only the last inserted
Expand Down
77 changes: 70 additions & 7 deletions aw_datastore/storages/peewee.py
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,63 @@ def update_bucket_keys(self) -> None:
def buckets(self) -> Dict[str, Dict[str, Any]]:
return {bucket.id: bucket.json() for bucket in BucketModel.select()}

def has_bucket(self, bucket_id: str) -> bool:
key = (
BucketModel.select(BucketModel.key)
.where(BucketModel.id == bucket_id)
.scalar()
)
if key is None:
self.bucket_keys.pop(bucket_id, None)
return False
# A cold lookup may discover a bucket created by another connection.
# Update the key used by subsequent event reads as well as existence.
self.bucket_keys[bucket_id] = key
return True

def buckets_with_last_updated(self) -> Dict[str, Dict[str, Any]]:
# The correlated seek uses (bucket_id, timestamp), including empty buckets.
latest = (
EventModel.select(EventModel.id)
.where(EventModel.bucket == BucketModel.key)
.order_by(EventModel.timestamp.desc())
.limit(1)
)
query = BucketModel.select(
BucketModel,
EventModel.timestamp.alias("last_timestamp"),
EventModel.duration.alias("last_duration"),
).join(EventModel, peewee.JOIN.LEFT_OUTER, on=(EventModel.id == latest))
buckets = {}
for row in query.objects().iterator():
metadata = row.json()
if row.last_timestamp is not None:
event = Event(
timestamp=row.last_timestamp, duration=float(row.last_duration)
)
metadata["last_updated"] = (
event.timestamp + event.duration
).isoformat()
buckets[row.id] = metadata
return buckets

def iter_events(self, bucket_id):
cursor = self.db.execute_sql(
"SELECT id, timestamp, duration, datastr FROM eventmodel "
"WHERE bucket_id = ? ORDER BY timestamp DESC",
(self.bucket_keys[bucket_id],),
)
try:
for event_id, timestamp, duration, datastr in cursor:
yield Event(
id=event_id,
timestamp=timestamp,
duration=float(duration),
data=json.loads(datastr),
)
finally:
cursor.close()

def create_bucket(
self,
bucket_id: str,
Expand Down Expand Up @@ -347,13 +404,19 @@ def delete(self, bucket_id, event_id):
)

def replace(self, bucket_id, event_id, event):
e = self._get_event(bucket_id, event_id)
e.timestamp = event.timestamp
e.duration = event.duration.total_seconds()
e.datastr = json.dumps(event.data)
e.save()
event.id = e.id
return event
updated = (
EventModel.update(
timestamp=event.timestamp,
duration=event.duration.total_seconds(),
datastr=json.dumps(event.data),
)
.where(EventModel.id == event_id)
.where(EventModel.bucket == self.bucket_keys[bucket_id])
.execute()
)
if updated:
event.id = event_id
return bool(updated)

def get_event(
self,
Expand Down
63 changes: 58 additions & 5 deletions aw_datastore/storages/sqlite.py
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,55 @@ def buckets(self):
}
return buckets

def has_bucket(self, bucket_id: str) -> bool:
return (
self.conn.execute(
"SELECT 1 FROM buckets WHERE id = ?", (bucket_id,)
).fetchone()
is not None
)

def buckets_with_last_updated(self):
# Match get_events' ordering by endtime for this backend.
rows = self.conn.execute(
"SELECT b.id, b.name, b.type, b.client, b.hostname, b.created, "
"b.datastr, e.endtime FROM buckets b LEFT JOIN events e ON e.id = "
"(SELECT id FROM events WHERE bucketrow = b.rowid "
"AND endtime >= 0 AND starttime <= ? ORDER BY endtime DESC LIMIT 1)",
(MAX_TIMESTAMP,),
)
try:
buckets = {}
for row in rows:
metadata = dict(
zip(
("id", "name", "type", "client", "hostname", "created"), row[:6]
)
)
metadata["data"] = json.loads(row[6] or "{}")
if row[7] is not None:
metadata["last_updated"] = datetime.fromtimestamp(
row[7] / 1000000, timezone.utc
).isoformat()
buckets[row[0]] = metadata
return buckets
finally:
rows.close()

def iter_events(self, bucket_id):
self.commit()
cursor = self.conn.execute(
"SELECT id, starttime, endtime, datastr FROM events "
"WHERE bucketrow = (SELECT rowid FROM buckets WHERE id = ?) "
"AND endtime >= 0 AND starttime <= ? ORDER BY endtime DESC",
(bucket_id, MAX_TIMESTAMP),
)
try:
for row in cursor:
yield _rows_to_events([row])[0]
finally:
cursor.close()

def create_bucket(
self,
bucket_id: str,
Expand Down Expand Up @@ -301,14 +350,18 @@ def replace(self, bucket_id, event_id, event) -> bool:
endtime = starttime + (event.duration.total_seconds() * 1000000)
datastr = json.dumps(event.data)
query = """UPDATE events
SET bucketrow = (SELECT rowid FROM buckets WHERE id = ?),
starttime = ?,
SET starttime = ?,
endtime = ?,
datastr = ?
WHERE id = ?"""
self.conn.execute(query, [bucket_id, starttime, endtime, datastr, event_id])
WHERE id = ? AND bucketrow =
(SELECT rowid FROM buckets WHERE id = ?)"""
cursor = self.conn.execute(
query, [starttime, endtime, datastr, event_id, bucket_id]
)
self.conditional_commit(1)
return True
if cursor.rowcount:
event.id = event_id
return bool(cursor.rowcount)

def get_event(
self,
Expand Down
Loading
Loading