Skip to content
59 changes: 48 additions & 11 deletions aw_query/functions.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import logging
import warnings
from datetime import timedelta
from functools import wraps
from inspect import signature
Expand All @@ -12,6 +14,7 @@
import iso8601
from aw_core.models import Event
from aw_datastore import Datastore
from aw_transform.chunk_events_by_key import CHUNK_DEPRECATION
from aw_transform import (
Rule,
categorize,
Expand All @@ -36,6 +39,11 @@

from .exceptions import QueryFunctionException

logger = logging.getLogger(__name__)

# Logged once per process: queries run repeatedly (e.g. on every dashboard refresh)
_chunk_deprecation_logged = False


def _verify_bucket_exists(datastore, bucketname):
if bucketname in datastore.buckets():
Expand Down Expand Up @@ -256,7 +264,14 @@ def q2_merge_subwatcher_fields(
@q2_function(chunk_events_by_key)
@q2_typecheck
def q2_chunk_events_by_key(events: list, key: str) -> List[Event]:
return chunk_events_by_key(events, key)
global _chunk_deprecation_logged
if not _chunk_deprecation_logged:
_chunk_deprecation_logged = True
logger.warning(CHUNK_DEPRECATION)
with warnings.catch_warnings():
# Logged once above; don't also emit the transform's DeprecationWarning
warnings.simplefilter("ignore", DeprecationWarning)
return chunk_events_by_key(events, key)


"""
Expand Down Expand Up @@ -344,21 +359,43 @@ def q2_nop():
"""


def _parse_rules(function: str, classes: list) -> list:
"""[[name, rule_dict], ...] into (name, Rule) pairs, with a
QueryFunctionException (like aw-server-rust's query error) for malformed
entries instead of a crash."""
rules = []
for entry in classes:
if not isinstance(entry, list) or len(entry) != 2:
raise QueryFunctionException(
f"{function} expects a list of [name, rule] pairs, got {entry!r}"
)
name, rule_dict = entry
if not isinstance(rule_dict, dict):
raise QueryFunctionException(
f"{function} rule must be a dict, got {type(rule_dict).__name__}: {rule_dict!r}"
)
try:
rules.append((name, Rule(rule_dict)))
except ValueError as exc:
raise QueryFunctionException(str(exc)) from None
return rules


@q2_function(categorize)
@q2_typecheck
def q2_categorize(events: list, classes: list):
try:
classes = [(_cls, Rule(rule_dict)) for _cls, rule_dict in classes]
except ValueError as exc:
raise QueryFunctionException(str(exc)) from None
return categorize(_copy_events(events), classes)
return categorize(_copy_events(events), _parse_rules("categorize", classes))


@q2_function(tag)
@q2_typecheck
def q2_tag(events: list, classes: list):
try:
classes = [(_cls, Rule(rule_dict)) for _cls, rule_dict in classes]
except ValueError as exc:
raise QueryFunctionException(str(exc)) from None
return tag(_copy_events(events), classes)
# Tag names are strings, like in aw-server-rust. Category-style list names
# belong to categorize (ActivityWatch/activitywatch#1466).
rules = _parse_rules("tag", classes)
for name, _ in rules:
if not isinstance(name, str):
raise QueryFunctionException(
f"tag name must be a string, got {type(name).__name__}: {name!r}"
)
return tag(_copy_events(events), rules)
15 changes: 15 additions & 0 deletions aw_transform/chunk_events_by_key.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import logging
import warnings
from datetime import timedelta
from typing import List

Expand All @@ -7,13 +8,27 @@
logger = logging.getLogger(__name__)


CHUNK_DEPRECATION = (
"chunk_events_by_key is deprecated and will be removed. There is no drop-in "
"replacement: merge_events_by_keys merges all events with the same value "
"(also across gaps) and doesn't produce subevents."
)


def chunk_events_by_key(
events: List[Event], key: str, pulsetime: float = 5.0
) -> List[Event]:
"""
"Chunks" adjacent events together which have the same value for a key, and stores the
original events in the :code:`subevents` key of the new event.

.. deprecated::
Will be removed (ActivityWatch/activitywatch#1466). There is no
drop-in replacement: aw-server-rust never supported ``subevents``, and
:func:`merge_events_by_keys` merges all events with the same value,
also across gaps, instead of adjacent runs.
"""
warnings.warn(CHUNK_DEPRECATION, DeprecationWarning, stacklevel=2)
chunked_events: List[Event] = []
for event in events:
if key not in event.data:
Expand Down
13 changes: 11 additions & 2 deletions aw_transform/classify.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,21 +79,30 @@ def _categorize_one(e: Event, classes: List[Tuple[Category, Rule]]) -> Event:
return e


def _matching_tags(e: Event, classes: List[Tuple[Tag, Rule]]) -> List[Tag]:
# Sorted and deduplicated, like aw-server-rust (ActivityWatch/activitywatch#1466)
return sorted({_cls for _cls, rule in classes if rule.match(e)})


def tag(events: List[Event], classes: List[Tuple[Tag, Rule]]) -> List[Event]:
"""
Adds the names of all matching rules to ``$tags`` (sorted, without
duplicates). Unlike categories, an event can have several tags.
"""
cache: Dict[str, List[Tag]] = {}
for e in events:
try:
key = json.dumps(e.data, sort_keys=True)
except TypeError:
key = str(id(e.data))
if key not in cache:
cache[key] = [_cls for _cls, rule in classes if rule.match(e)]
cache[key] = _matching_tags(e, classes)
e.data["$tags"] = list(cache[key])
return events


def _tag_one(e: Event, classes: List[Tuple[Tag, Rule]]) -> Event:
e.data["$tags"] = [_cls for _cls, rule in classes if rule.match(e)]
e.data["$tags"] = _matching_tags(e, classes)
return e


Expand Down
63 changes: 37 additions & 26 deletions aw_transform/merge_events_by_keys.py
Original file line number Diff line number Diff line change
@@ -1,40 +1,51 @@
import copy
import json
import logging
from typing import List, Dict, Tuple
from typing import Any, Dict, List

from aw_core.models import Event

logger = logging.getLogger(__name__)


def merge_events_by_keys(events, keys) -> List[Event]:
def _non_json_key(value: Any) -> Dict[str, str]:
# Event data from the datastore is always JSON, like in aw-server-rust.
# Other values (only possible through the Python API) are tagged with their
# type, so e.g. a datetime can't merge with a string that looks the same.
return {
"$non-json": f"{type(value).__module__}.{type(value).__qualname__}:{value!r}"
}


def merge_events_by_keys(events: List[Event], keys: List[str]) -> List[Event]:
"""
Sums the duration of all events which share a value for a key and returns a new event for each value.
Merges all events that share the same values for all of ``keys``, whether
they are adjacent or not, summing their durations.

.. note: The result will be a list of events without timestamp since they are merged.
Each merged event keeps the timestamp and the whole ``data`` of the first
event in its group (not only the merge keys), so fields that are the same
across the group, like ``$category`` for ``["app", "title"]``, stay
available. Events missing any of the keys are dropped, and an empty key list
returns no events. This matches aw-server-rust (ActivityWatch/activitywatch#1466).
"""
# Call recursively until all keys are consumed
if len(keys) < 1:
return events
merged_events: Dict[Tuple, Event] = {}
if not keys:
return []
merged_events: Dict[str, Event] = {}
for event in events:
composite_key: Tuple = ()
for key in keys:
if key in event.data:
val = event["data"][key]
# Needed for when the value is a list, such as for categories
if isinstance(val, list):
val = tuple(val)
composite_key = composite_key + (val,)
if composite_key not in merged_events:
try:
values = [event.data[key] for key in keys]
except KeyError:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

❌ P1 — In merge_events_by_keys, the function now drops events that are missing any of the keys. This is a breaking change from the previous behavior, which kept events missing keys (they were grouped under an empty composite key). The PR description says this matches aw-server-rust, and the first-party consumers rely on it. However, this is a contract change for any existing aw-core users who used merge_events_by_keys with events that may lack some keys. The test test_merge_events_by_keys_1 was updated to reflect this. This is intentional, but it is a breaking change. I'll report it as a contract finding.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Intentional, same decision (ActivityWatch/activitywatch#1466): aw-webui's Android and multidevice queries rely on events without the keys being dropped, as aw-server-rust does. Marked BREAKING CHANGE and described in the PR.

continue
# Group by the JSON values, like aw-server-rust (so 1 and 1.0 differ,
# and list values such as categories work).
composite_key = json.dumps(values, sort_keys=True, default=_non_json_key)
merged = merged_events.get(composite_key)
if merged is None:
merged_events[composite_key] = Event(
timestamp=event.timestamp, duration=event.duration, data={}
timestamp=event.timestamp,
duration=event.duration,
data=copy.deepcopy(event.data),
)
for key in keys:
if key in event.data:
merged_events[composite_key].data[key] = event.data[key]
else:
merged_events[composite_key].duration += event.duration
result = []
for key in merged_events:
result.append(Event(**merged_events[key]))
return result
merged.duration += event.duration
return list(merged_events.values())
Loading
Loading