From 8e8db67df0821ced5fd549578e79af85117681ce Mon Sep 17 00:00:00 2001 From: Bradley Erickson Date: Tue, 13 Jan 2026 13:33:25 -0500 Subject: [PATCH] abstracted time on task reducers to common area to allow for more modules to use --- VERSION | 2 +- learning_observer/VERSION | 2 +- .../stream_analytics/__init__.py | 2 + .../stream_analytics/time_on_task.py | 103 +++++++++++++++++ modules/writing_observer/VERSION | 2 +- .../writing_observer/writing_analysis.py | 106 ++---------------- 6 files changed, 120 insertions(+), 97 deletions(-) create mode 100644 learning_observer/learning_observer/stream_analytics/time_on_task.py diff --git a/VERSION b/VERSION index 55d096bf..32e666db 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.1.0+2026.01.06T16.20.35.371Z.b80005f9.master +0.1.0+2026.01.13T18.33.25.519Z.0984e08f.berickson.20260113.abstract.time.on.task.reducers diff --git a/learning_observer/VERSION b/learning_observer/VERSION index 55d096bf..32e666db 100644 --- a/learning_observer/VERSION +++ b/learning_observer/VERSION @@ -1 +1 @@ -0.1.0+2026.01.06T16.20.35.371Z.b80005f9.master +0.1.0+2026.01.13T18.33.25.519Z.0984e08f.berickson.20260113.abstract.time.on.task.reducers diff --git a/learning_observer/learning_observer/stream_analytics/__init__.py b/learning_observer/learning_observer/stream_analytics/__init__.py index e52cf138..36a809ea 100644 --- a/learning_observer/learning_observer/stream_analytics/__init__.py +++ b/learning_observer/learning_observer/stream_analytics/__init__.py @@ -20,6 +20,8 @@ from learning_observer.log_event import debug_log +import learning_observer.stream_analytics.time_on_task + REDUCER_MODULES = None LAST_UPDATED = None diff --git a/learning_observer/learning_observer/stream_analytics/time_on_task.py b/learning_observer/learning_observer/stream_analytics/time_on_task.py new file mode 100644 index 00000000..a31391f8 --- /dev/null +++ b/learning_observer/learning_observer/stream_analytics/time_on_task.py @@ -0,0 +1,103 @@ +''' +Helpers for time-on-task reducers. +''' +import pmss + +pmss.register_field( + name='time_on_task_threshold', + type=pmss.pmsstypes.TYPES.integer, + description='Maximum time to pass before marking a session as over. '\ + 'Should be 60-300 seconds in production, but 5 seconds is nice for '\ + 'debugging in a local deployment.', + default=60 +) +pmss.register_field( + name='binned_time_on_task_bin_size', + type=pmss.pmsstypes.TYPES.integer, + description='How large (in seconds) to make timestamp bins when '\ + 'recording binned time on task.', + default=600 +) + + +def default_time_on_task_state(): + return { + 'saved_ts': None, + 'total_time_on_task': 0 + } + + +def apply_time_on_task(internal_state, current_timestamp, time_delta_threshold): + if internal_state is None: + internal_state = default_time_on_task_state() + last_ts = internal_state['saved_ts'] + internal_state['saved_ts'] = current_timestamp + + if last_ts is None: + last_ts = internal_state['saved_ts'] + if last_ts is not None: + delta_t = min( + time_delta_threshold, + internal_state['saved_ts'] - last_ts + ) + internal_state['total_time_on_task'] += delta_t + return internal_state + + +def default_binned_time_on_task_state(): + return { + 'saved_ts': None, + 'binned_time_on_task': {}, + 'current_bin': None + } + + +def get_time_bin(timestamp, bin_size): + b = (timestamp // bin_size) * bin_size + return int(b) + + +def update_binned_time_on_task(internal_state, current_bin, last_timestamp, delta_time, bin_size): + '''Handle updating the internal state for binned time on task.''' + next_bin = current_bin + bin_size + next_bin_str = str(next_bin) + + current_bin_str = str(current_bin) + if current_bin_str not in internal_state['binned_time_on_task']: + internal_state['binned_time_on_task'][current_bin_str] = 0 + + if last_timestamp + delta_time >= next_bin: + internal_state['binned_time_on_task'][current_bin_str] += next_bin - last_timestamp + if next_bin_str not in internal_state['binned_time_on_task']: + internal_state['binned_time_on_task'][next_bin_str] = 0 + internal_state['binned_time_on_task'][next_bin_str] += last_timestamp + delta_time - next_bin + else: + internal_state['binned_time_on_task'][current_bin_str] += delta_time + + +def apply_binned_time_on_task( + internal_state, + current_timestamp, + time_delta_threshold, + bin_size +): + if internal_state is None: + internal_state = default_binned_time_on_task_state() + last_timestamp = internal_state['saved_ts'] + current_bin = internal_state['current_bin'] + internal_state['saved_ts'] = current_timestamp + + if last_timestamp is None: + last_timestamp = internal_state['saved_ts'] + if current_bin is None: + current_bin = get_time_bin(last_timestamp, bin_size) + + if last_timestamp is not None: + delta_time = min( + time_delta_threshold, + internal_state['saved_ts'] - last_timestamp + ) + update_binned_time_on_task(internal_state, current_bin, last_timestamp, delta_time, bin_size) + + internal_state['current_bin'] = get_time_bin(internal_state['saved_ts'], bin_size) + return internal_state diff --git a/modules/writing_observer/VERSION b/modules/writing_observer/VERSION index de21e3a8..32e666db 100644 --- a/modules/writing_observer/VERSION +++ b/modules/writing_observer/VERSION @@ -1 +1 @@ -0.1.0+2025.10.13T19.31.35.177Z.f12616b4.master +0.1.0+2026.01.13T18.33.25.519Z.0984e08f.berickson.20260113.abstract.time.on.task.reducers diff --git a/modules/writing_observer/writing_observer/writing_analysis.py b/modules/writing_observer/writing_observer/writing_analysis.py index 7ebe5dd1..4d1a836e 100644 --- a/modules/writing_observer/writing_observer/writing_analysis.py +++ b/modules/writing_observer/writing_observer/writing_analysis.py @@ -16,6 +16,7 @@ import learning_observer.adapters import learning_observer.communication_protocol.integration from learning_observer.stream_analytics.helpers import student_event_reducer, kvs_pipeline, KeyField, EventField, Scope +import learning_observer.stream_analytics.time_on_task import learning_observer.settings import learning_observer.util @@ -34,21 +35,6 @@ # (e.g. all the numbers would go up/down 20%, but behavior was # substantatively identical). -pmss.register_field( - name='time_on_task_threshold', - type=pmss.pmsstypes.TYPES.integer, - description='Maximum time to pass before marking a session as over. '\ - 'Should be 60-300 seconds in production, but 5 seconds is nice for '\ - 'debugging in a local deployment.', - default=60 -) -pmss.register_field( - name='binned_time_on_task_bin_size', - type=pmss.pmsstypes.TYPES.integer, - description='How large (in seconds) to make timestamp bins when '\ - 'recording binned time on task.', - default=600 -) pmss.register_field( name='activity_threshold', type=pmss.pmsstypes.TYPES.integer, @@ -93,64 +79,12 @@ async def time_on_task(event, internal_state): goes away for 2 hours without typing, we only add e.g. 5 minutes if `time_threshold` is set to 300. ''' - if internal_state is None: - internal_state = { - 'saved_ts': None, - 'total_time_on_task': 0 - } - last_ts = internal_state['saved_ts'] - internal_state['saved_ts'] = event['server']['time'] - - # Initial conditions - if last_ts is None: - last_ts = internal_state['saved_ts'] - if last_ts is not None: - delta_t = min( - learning_observer.settings.module_setting('writing_obersver', 'time_on_task_threshold'), # Maximum time step - internal_state['saved_ts'] - last_ts # Time step - ) - internal_state['total_time_on_task'] += delta_t - return internal_state, internal_state - - -def _get_time_delta(last_event_timestamp, current_event_timestamp): - return min( - learning_observer.settings.module_setting('writing_obersver', 'time_on_task_threshold'), # Maximum time step - last_event_timestamp - current_event_timestamp # Time step + internal_state = learning_observer.stream_analytics.time_on_task.apply_time_on_task( + internal_state, + event['server']['time'], + learning_observer.settings.module_setting('writing_obersver', 'time_on_task_threshold') ) - - -def _get_time_bin(timestamp): - bin_size = learning_observer.settings.module_setting('writing_obersver', 'binned_time_on_task_bin_size') - b = (timestamp // bin_size) * bin_size - b = int(b) - return b - - -def _update_binned_time_on_task(internal_state, current_bin, last_timestamp, delta_time): - '''Handle updating the internal state for binned time on task. - ''' - next_bin = current_bin + learning_observer.settings.module_setting('writing_obersver', 'binned_time_on_task_bin_size') - next_bin_str = str(next_bin) - - # default current_bin to 0 if it doesn't exist - current_bin_str = str(current_bin) - if current_bin_str not in internal_state['binned_time_on_task']: - internal_state['binned_time_on_task'][current_bin_str] = 0 - - # time-on-task overflows to the next bin - # first add a portion of the time to the current bin - # default the next bin to 0 if it doesn't exist - # add remaining time to next bin - if last_timestamp + delta_time >= next_bin: - internal_state['binned_time_on_task'][current_bin_str] += next_bin - last_timestamp - if next_bin_str not in internal_state['binned_time_on_task']: - internal_state['binned_time_on_task'][next_bin_str] = 0 - internal_state['binned_time_on_task'][next_bin_str] += last_timestamp + delta_time - next_bin - # process normal within bin time on task update - else: - internal_state['binned_time_on_task'][current_bin_str] += delta_time - + return internal_state, internal_state @kvs_pipeline(scope=gdoc_scope) @@ -159,28 +93,12 @@ async def binned_time_on_task(event, internal_state): Similar to the `time_on_task` reducer defined above, except it bins the time spent. ''' - if internal_state is None: - internal_state = { - 'saved_ts': None, - 'binned_time_on_task': {}, - 'current_bin': None - } - last_timestamp = internal_state['saved_ts'] - current_bin = internal_state['current_bin'] - internal_state['saved_ts'] = event['server']['time'] - - # Initialization - if last_timestamp is None: - last_timestamp = internal_state['saved_ts'] - if current_bin is None: - current_bin = _get_time_bin(last_timestamp) - - if last_timestamp is not None: - delta_time = _get_time_delta(internal_state['saved_ts'], last_timestamp) - _update_binned_time_on_task(internal_state, current_bin, last_timestamp, delta_time) - - # update our current bin with the current event's timestamp - internal_state['current_bin'] = _get_time_bin(internal_state['saved_ts']) + internal_state = learning_observer.stream_analytics.time_on_task.apply_binned_time_on_task( + internal_state, + event['server']['time'], + learning_observer.settings.module_setting('writing_obersver', 'time_on_task_threshold'), + learning_observer.settings.module_setting('writing_obersver', 'binned_time_on_task_bin_size') + ) return internal_state, internal_state