loopable / platform /modules /tasks.py
fsanyoto's picture
Deploy AIOS web (React glide grid + FastAPI slice)
c14ceee verified
Raw
History Blame Contribute Delete
8.74 kB
"""Store-backed task / worklist lifecycle — the writable PROCESS layer over the read-only ERP
(the Palantir split: subject layer = Odoo, read-only; process layer = this, in the HF store).
Store key 'tasks': {"seq": int, "items": {id: task}}. A task:
id 't<seq>' — stable
title verb-first next action ("Call Smith Floral about reorder")
detail one-line evidence/why (frozen at flag time for rule tasks)
entity {'kind': customer|sku|agent|account, 'id': …, 'label': …} | None
owner_user username in the users registry
due 'YYYY-MM-DD' | None
state open | done | snoozed | dismissed
source 'user' | 'rule:<rule_id>'
cond_key rule tasks: fingerprint of the flagging condition — a dismissed task is NOT
re-created until the condition RESETS (cond_key changes). The anti-fatigue rule.
dismiss_reason enumerated (required on dismiss); reviewed to tune/kill rules
outcome 'self-resolved' when a rule flag cleared on its own (auto-closed)
created_by / created_at / status_changed / snooze_until / history (append-only)
All mutations go through store.update (strict read-modify-write — never get()+put()).
"""
import datetime as dt
import core.store as store
KEY = 'tasks'
DISMISS_REASONS = ['Not relevant', 'Already handled', 'Data is wrong', 'Duplicate', 'Other']
STATES = ('open', 'done', 'snoozed', 'dismissed')
def _now():
return dt.datetime.now().strftime('%Y-%m-%d %H:%M')
def _today():
return dt.date.today().isoformat()
def _entity_key(entity):
if not entity:
return None
return f"{entity.get('kind')}:{entity.get('id')}"
def all_tasks():
d = store.get(KEY)
return list((d.get('items') or {}).values())
def effective_state(t, today=None):
"""Snoozed past its wake date reads as open (Linear semantics — wake early on due)."""
s = t.get('state', 'open')
if s == 'snoozed' and (t.get('snooze_until') or '') <= (today or _today()):
return 'open'
return s
def for_owner(owner=None, states=('open',)):
"""Open queue for a user (None = everyone), due-date order, snooze-awakened."""
today = _today()
rows = [t for t in all_tasks()
if (owner is None or t.get('owner_user') == owner)
and effective_state(t, today) in states]
return sorted(rows, key=lambda t: (t.get('due') or '9999-12-31', t.get('id', '')))
def counts_by_user():
"""{user: {'open': n, 'overdue': n, 'stale7': n}} — feeds the owner rollup + digests."""
today = _today()
week_ago = (dt.date.today() - dt.timedelta(days=7)).isoformat()
out = {}
for t in all_tasks():
if effective_state(t, today) != 'open':
continue
u = t.get('owner_user') or '—'
rec = out.setdefault(u, {'open': 0, 'overdue': 0, 'stale7': 0})
rec['open'] += 1
if (t.get('due') or '9999') < today:
rec['overdue'] += 1
if (t.get('status_changed') or t.get('created_at') or '')[:10] <= week_ago:
rec['stale7'] += 1
return out
def create(title, owner_user, created_by, detail='', entity=None, due=None,
source='user', cond_key=None):
"""Create one task (human producer). Returns the new task id."""
created = {}
def _add(d):
d.setdefault('seq', 0)
d.setdefault('items', {})
d['seq'] += 1
tid = f"t{d['seq']}"
d['items'][tid] = {
'id': tid, 'title': str(title).strip(), 'detail': str(detail or '').strip(),
'entity': entity or None, 'owner_user': owner_user, 'due': due,
'state': 'open', 'source': source, 'cond_key': cond_key,
'dismiss_reason': None, 'outcome': None,
'created_by': created_by, 'created_at': _now(), 'status_changed': _now(),
'snooze_until': None,
'history': [{'at': _now(), 'by': created_by, 'event': 'created'}],
}
created['id'] = tid
return d
store.update(KEY, _add)
return created.get('id')
def set_state(tid, state, user, dismiss_reason=None, snooze_until=None):
"""Move a task through its lifecycle. Dismiss REQUIRES an enumerated reason."""
if state not in STATES:
raise ValueError(f'bad state {state}')
if state == 'dismissed' and not dismiss_reason:
raise ValueError('dismiss requires a reason')
def _set(d):
t = (d.get('items') or {}).get(tid)
if not t:
return d
t['state'] = state
t['status_changed'] = _now()
if state == 'dismissed':
t['dismiss_reason'] = dismiss_reason
if state == 'snoozed':
t['snooze_until'] = snooze_until
t.setdefault('history', []).append(
{'at': _now(), 'by': user, 'event': state,
**({'reason': dismiss_reason} if dismiss_reason else {}),
**({'until': snooze_until} if snooze_until else {})})
return d
store.update(KEY, _set)
def sync_rule_tasks(rule_id, flags, default_owner=None):
"""Rule producer (the nightly job): reconcile the store against the CURRENT flag set.
flags: [{entity, title, detail, owner_user?, due?, cond_key}] — cond_key fingerprints the
condition (e.g. bucket/threshold hit), so:
- open/snoozed task for the same (rule, entity): refresh detail (evidence stays current);
- done/dismissed task with the SAME cond_key: skip (never re-nag until the condition resets);
- cond_key changed: create a NEW task (the condition reset and re-fired);
- open rule task whose entity no longer flags: auto-close as self-resolved (the
stale-list/paid-invoice trust-killer).
One store commit for the whole reconcile.
"""
src = f'rule:{rule_id}'
def _sync(d):
d.setdefault('seq', 0)
d.setdefault('items', {})
items = d['items']
by_entity = {}
for t in items.values():
if t.get('source') == src and t.get('entity'):
by_entity.setdefault(_entity_key(t['entity']), []).append(t)
seen = set()
for f in flags:
ek = _entity_key(f.get('entity'))
if not ek:
continue
seen.add(ek)
existing = sorted(by_entity.get(ek, []), key=lambda t: t['id'])
live = [t for t in existing if t.get('state') in ('open', 'snoozed')]
if live:
live[-1]['detail'] = f.get('detail', live[-1].get('detail'))
live[-1]['cond_key'] = f.get('cond_key', live[-1].get('cond_key'))
continue
closed_same = [t for t in existing
if t.get('cond_key') and t.get('cond_key') == f.get('cond_key')]
if closed_same:
continue # dismissed/done and unchanged — stay quiet
d['seq'] += 1
tid = f"t{d['seq']}"
items[tid] = {
'id': tid, 'title': f.get('title', rule_id), 'detail': f.get('detail', ''),
'entity': f.get('entity'), 'owner_user': f.get('owner_user') or default_owner,
'due': f.get('due'), 'state': 'open', 'source': src,
'cond_key': f.get('cond_key'), 'dismiss_reason': None, 'outcome': None,
'created_by': 'system', 'created_at': _now(), 'status_changed': _now(),
'snooze_until': None,
'history': [{'at': _now(), 'by': 'system', 'event': 'created'}],
}
for ek, ts in by_entity.items():
if ek in seen:
continue
for t in ts:
if t.get('state') in ('open', 'snoozed'):
t['state'] = 'done'
t['outcome'] = 'self-resolved'
t['status_changed'] = _now()
t.setdefault('history', []).append(
{'at': _now(), 'by': 'system', 'event': 'self-resolved'})
return d
store.update(KEY, _sync)
def rule_health():
"""fired / actioned(done by a person) / dismissed / self-resolved per rule — the noise budget.
<20% actioned = the rule is noise: tighten or kill (SRE actionability test)."""
out = {}
for t in all_tasks():
src = t.get('source') or ''
if not src.startswith('rule:'):
continue
r = out.setdefault(src[5:], {'fired': 0, 'done': 0, 'dismissed': 0, 'self_resolved': 0})
r['fired'] += 1
if t.get('outcome') == 'self-resolved':
r['self_resolved'] += 1
elif t.get('state') == 'done':
r['done'] += 1
elif t.get('state') == 'dismissed':
r['dismissed'] += 1
return out