Skip to content
Merged
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 VERSION
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion learning_observer/VERSION
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
103 changes: 103 additions & 0 deletions learning_observer/learning_observer/stream_analytics/time_on_task.py
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion modules/writing_observer/VERSION
Original file line number Diff line number Diff line change
@@ -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
106 changes: 12 additions & 94 deletions modules/writing_observer/writing_observer/writing_analysis.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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,
Expand Down Expand Up @@ -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)
Expand All @@ -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


Expand Down
Loading