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
42 changes: 1 addition & 41 deletions raven/proactive_engine/schedulers/cron/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -532,23 +532,13 @@ def add_job(
conversations; without dedup the user gets N near-identical
fires per scheduled tick.

Two cross-kind dedup layers also apply (in order):

1. **Message-equal dedup**: if any existing enabled job for the
A cross-kind dedup layer also applies: if any existing enabled job for the
same (channel, to) has a *byte-identical* ``payload.message``,
return it. Catches the case where the LLM creates the same
reminder N times across the simulation horizon (e.g. a
medication-reminder string appearing as both an ``at`` shot
today and a ``cron_expr`` recurring tomorrow — identical
text, different fire times).

2. **Time-window dedup**: if any existing enabled job for the
same (channel, to) is scheduled to fire within 15 minutes of
this new schedule's next fire (regardless of schedule kind),
return it. Catches the case where the LLM creates both a
recurring ``cron_expr`` AND a same-day ``at`` shot for the
same intent (e.g. "daily 8:00 take meds" + "today 8:00 take
meds").
"""
# One ``now`` snapshot for both validation and storage: the validate
# predicate and the stored next_run must agree on "now", or a boundary
Expand Down Expand Up @@ -626,36 +616,6 @@ def add_job(
self._arm_timer()
return j

# Cross-kind time-window dedup (covers caregiver-style
# "expr + at for the same intent" double-add). Window is
# generous (15min) because two genuinely-distinct reminders
# less than 15min apart are almost always an LLM mistake;
# the rare legitimate case (two distinct meds at 8:00 and
# 8:10) loses one fire — acceptable trade-off given the
# spam alternative.
new_next = _compute_next_run(schedule, now)
if new_next is not None:
for j in store.jobs:
if not j.enabled:
continue
if j.payload.channel != channel or j.payload.to != to:
continue
existing_next = j.state.next_run_at_ms
if existing_next is None:
continue
if abs(existing_next - new_next) <= 15 * 60 * 1000:
logger.info(
"Cron: skipped duplicate add — existing job '{}' "
"({}) fires within 15min of new request "
"(same channel/to, kinds={}/{})",
j.name,
j.id,
j.schedule.kind,
schedule.kind,
)
self._arm_timer()
return j

# Dedup: same recurring schedule + same channel + same recipient
# → update message in place rather than create a duplicate.
existing = self._find_duplicate_schedule(store.jobs, schedule, channel, to)
Expand Down
22 changes: 19 additions & 3 deletions tests/test_cron_service_dedup.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,9 @@
"""Tests for CronService.add_job dedup layers.

Three layers, applied in order at add time:
Three identity-based layers, applied in order at add time:
1. L7 topic_tag dedup — strictest, runs first when caller supplies tag.
2. Message-equal dedup — byte-identical payload.message + same channel/to.
3. Time-window dedup — within 15 min of an existing fire on same channel/to.
4. Schedule dedup — same recurring schedule + same channel/to.
3. Schedule dedup — same recurring schedule + same channel/to.

L7 is the load-bearing one for the caregiver-style failure: the LLM
re-asks for the same daily med reminder across the month with slightly
Expand Down Expand Up @@ -175,3 +174,20 @@ def test_topic_tag_dedup_runs_before_message_equal(svc):
topic_tag="meds_morning",
)
assert j1.id == j2.id


def test_distinct_messages_in_same_time_window_keep_separate(svc):
"""Nearby reminders with different content are distinct jobs."""
first = _add(
svc,
"prepare the daily brief",
CronSchedule(kind="at", at_ms=int(1e15)),
)
second = _add(
svc,
"review the presentation",
CronSchedule(kind="at", at_ms=int(1e15) + 10 * 60 * 1000),
)

assert first.id != second.id
assert len(svc.list_jobs()) == 2