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
130 changes: 119 additions & 11 deletions src/powermem/intelligence/ebbinghaus_algorithm.py
Original file line number Diff line number Diff line change
Expand Up @@ -255,16 +255,8 @@ def should_forget(self, memory: Dict[str, Any]) -> bool:
True if memory should be forgotten
"""
try:
# Check decay factor
created_at = memory.get("created_at")
if created_at:
decay_factor = self.calculate_decay(
created_at,
decay_rate=self._resolve_decay_rate(memory),
)
if decay_factor < self.working_threshold:
return True

if self.calculate_current_retention(memory) < self.working_threshold:
return True
return False

except Exception as e:
Expand Down Expand Up @@ -303,6 +295,66 @@ def should_archive(self, memory: Dict[str, Any]) -> bool:
logger.error(f"Failed to check archiving: {e}")
return False

def reinforce(self, memory: Dict[str, Any]) -> Dict[str, Any]:
"""Boost current_retention on review and advance the review schedule.

Called when a memory is accessed at or after its ``next_review`` time.
Uses diminishing-returns formula so retention approaches but never
exceeds 1.0.

Returns:
Dict with updated intelligence fields to merge back.
"""
_, intelligence = self._resolve_metadata_sections(memory)
current_retention = self.calculate_current_retention(memory)
reinforcement_factor = self._resolve_reinforcement_factor(memory)
review_count = int(intelligence.get("review_count") or 0)

new_retention = min(
1.0, current_retention + reinforcement_factor * (1.0 - current_retention)
)
new_review_count = review_count + 1

review_schedule = intelligence.get("review_schedule") or []
next_review = (
review_schedule[new_review_count]
if new_review_count < len(review_schedule)
else None
)

now = get_current_datetime()
return {
"current_retention": new_retention,
"review_count": new_review_count,
"last_reviewed": now.isoformat(),
"next_review": next_review,
}

def calculate_current_retention(self, memory: Dict[str, Any]) -> float:
"""Return the real-time effective retention for display/ranking.

``current_retention`` is a snapshot captured at ``last_reviewed`` (or
creation time for initial metadata), so it must decay before runtime
consumers use it. This avoids treating initialized retention as a
permanent floor while still reflecting recent review reinforcement.
"""
stored = self._resolve_current_retention(memory)
if stored is not None:
anchor = self._resolve_retention_anchor(memory)
decay = self.calculate_decay(
anchor,
decay_rate=self._resolve_decay_rate(memory),
)
return max(0.0, min(1.0, stored * decay))

initial = self._resolve_initial_retention(memory)
created_at = self._resolve_created_at(memory)
decay = self.calculate_decay(
created_at,
decay_rate=self._resolve_decay_rate(memory),
)
return max(0.0, min(1.0, initial * decay))

def get_review_schedule(
self, memory: Dict[str, Any], *, prefer_stored: bool = True
) -> list:
Expand Down Expand Up @@ -455,7 +507,63 @@ def _resolve_reinforcement_factor(self, memory: Dict[str, Any]) -> float:
logger.warning("Invalid reinforcement_factor: %s", raw)
return 0.0
return max(factor, 0.0)


def _resolve_initial_retention(self, memory: Dict[str, Any]) -> float:
"""Resolve initial_retention from memory metadata with config fallback."""
meta, intelligence = self._resolve_metadata_sections(memory)
raw = self._first_present(
memory.get("initial_retention"),
meta.get("initial_retention"),
intelligence.get("initial_retention"),
)
if raw is not None:
try:
val = float(raw)
except (TypeError, ValueError):
logger.warning("Invalid initial_retention: %s", raw)
else:
if 0.0 < val <= 1.0:
return val
return self.initial_retention

def _resolve_current_retention(
self, memory: Dict[str, Any]
) -> Optional[float]:
"""Resolve stored current_retention as a bounded snapshot value."""
meta, intelligence = self._resolve_metadata_sections(memory)
raw = self._first_present(
memory.get("current_retention"),
meta.get("current_retention"),
intelligence.get("current_retention"),
)
if raw is None:
return None
try:
retention = float(raw)
except (TypeError, ValueError):
logger.warning("Invalid current_retention: %s", raw)
return None
return max(0.0, min(1.0, retention))

def _resolve_created_at(self, memory: Dict[str, Any]) -> Any:
"""Resolve creation timestamp from supported memory layouts."""
meta, intelligence = self._resolve_metadata_sections(memory)
return self._first_present(
memory.get("created_at"),
meta.get("created_at"),
intelligence.get("created_at"),
)

def _resolve_retention_anchor(self, memory: Dict[str, Any]) -> Any:
"""Resolve the timestamp for decaying current_retention snapshots."""
meta, intelligence = self._resolve_metadata_sections(memory)
return self._first_present(
intelligence.get("last_reviewed"),
memory.get("last_reviewed"),
meta.get("last_reviewed"),
self._resolve_created_at(memory),
)

def _parse_datetime(self, value: Any) -> datetime:
"""Parse datetime from object or ISO string."""
if isinstance(value, datetime):
Expand Down
12 changes: 4 additions & 8 deletions src/powermem/intelligence/intelligent_memory_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -173,14 +173,10 @@ def process_search_results(
Processed and ranked results
"""
try:
# Apply Ebbinghaus decay to results
processed_results = []
for result in results:
# Apply decay based on age and memory type
decay_rate = self.ebbinghaus_algorithm._resolve_decay_rate(result)
decay_factor = self.ebbinghaus_algorithm.calculate_decay(
result.get("created_at", get_current_datetime()),
decay_rate=decay_rate,
effective_retention = (
self.ebbinghaus_algorithm.calculate_current_retention(result)
)

processed_result = result.copy()
Expand All @@ -193,12 +189,12 @@ def process_search_results(
)
if original_score is not None:
processed_result["original_score"] = original_score
processed_result["decay_factor"] = decay_factor
processed_result["effective_retention"] = effective_retention
processed_result["forgotten_score_multiplier"] = (
forgotten_score_multiplier
)
processed_result["final_score"] = (
base_score * decay_factor * forgotten_score_multiplier
base_score * effective_retention * forgotten_score_multiplier
)
processed_result["score"] = processed_result["final_score"]

Expand Down
53 changes: 43 additions & 10 deletions src/powermem/intelligence/plugin.py
Original file line number Diff line number Diff line change
Expand Up @@ -122,9 +122,8 @@ def on_get(self, memory: Dict[str, Any]) -> Tuple[Optional[Dict[str, Any]], bool
if not self.enabled or not self._algo:
return None, False
try:
# Normalize: intelligence fields may be stored inside the metadata
# JSON column and not exposed at the top level of the memory dict.
meta = memory.get("metadata") or {}
intelligence = meta.get("intelligence") or memory.get("intelligence") or {}
memory_type = memory.get("memory_type") or meta.get("memory_type")
access_count_old = memory.get("access_count")
if access_count_old is None:
Expand All @@ -136,21 +135,39 @@ def on_get(self, memory: Dict[str, Any]) -> Tuple[Optional[Dict[str, Any]], bool
importance_score = 0.5

new_access_count = access_count_old + 1
now = get_current_datetime()
updates: Dict[str, Any] = {
"access_count": new_access_count,
"updated_at": get_current_datetime(),
"updated_at": now,
}
# Track which fields need updating inside the metadata JSON column
meta_updates: Dict[str, Any] = {"access_count": new_access_count}
intel_updates: Dict[str, Any] = {}

# Provide normalized values to algorithm checks
normalized = {
**memory,
"memory_type": memory_type,
"access_count": access_count_old,
"importance_score": importance_score,
}

# Review reinforcement: if the access happens at or after
# next_review, boost current_retention and advance the schedule.
next_review_raw = intelligence.get("next_review")
if next_review_raw:
next_review_dt = self._algo._parse_datetime(next_review_raw)
if now >= next_review_dt:
reinforcement_result = self._algo.reinforce(normalized)
intel_updates.update(reinforcement_result)

# Apply reinforcement to normalized so downstream should_forget()
# uses the boosted retention and its new timestamp anchor.
if intel_updates:
norm_intel = dict(normalized.get("metadata", {}).get("intelligence") or {})
norm_intel.update(intel_updates)
norm_meta = dict(normalized.get("metadata") or {})
norm_meta["intelligence"] = norm_intel
normalized = {**normalized, "metadata": norm_meta}

# Check promotion first — an accessed memory that qualifies for
# promotion should not be forgotten in the same on_get call.
new_memory_type = memory_type
Expand All @@ -164,8 +181,7 @@ def on_get(self, memory: Dict[str, Any]) -> Tuple[Optional[Dict[str, Any]], bool
updates["memory_type"] = "long_term"
meta_updates["memory_type"] = "long_term"

# Clear forget marker on promotion — a promoted memory should
# no longer carry the 0.1x search penalty from a prior soft-forget.
# Clear forget marker on promotion
if new_memory_type != memory_type:
meta_updates["should_forget"] = False
meta_updates["marked_for_forgetting_at"] = None
Expand All @@ -176,7 +192,6 @@ def on_get(self, memory: Dict[str, Any]) -> Tuple[Optional[Dict[str, Any]], bool
if new_memory_type == memory_type and self._algo.should_forget(normalized):
return None, True

# Check if memory should be archived
if self._algo.should_archive(normalized):
meta_updates["archived"] = True

Expand All @@ -194,8 +209,26 @@ def on_get(self, memory: Dict[str, Any]) -> Tuple[Optional[Dict[str, Any]], bool
original_content, importance_score, new_memory_type or "working"
)
if "intelligence" in intelligence_metadata:
updates["metadata"]["intelligence"] = intelligence_metadata["intelligence"]
updates["last_reprocessed_at"] = get_current_datetime()
reprocessed_intel = intelligence_metadata["intelligence"]
# Preserve current_retention and review progress from
# reinforcement; reprocessing should only refresh
# decay_rate and review_schedule, not reset retention.
for keep_key in (
"current_retention",
"review_count",
"last_reviewed",
"next_review",
):
if keep_key in intel_updates:
reprocessed_intel[keep_key] = intel_updates[keep_key]
elif keep_key in intelligence:
reprocessed_intel[keep_key] = intelligence[keep_key]
updates["metadata"]["intelligence"] = reprocessed_intel
updates["last_reprocessed_at"] = now
elif intel_updates:
existing_intel = dict(updates["metadata"].get("intelligence") or intelligence)
existing_intel.update(intel_updates)
updates["metadata"]["intelligence"] = existing_intel

return updates, False
except Exception as e:
Expand Down
10 changes: 5 additions & 5 deletions tests/unit/intelligence/test_ebbinghaus_decay_rate.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,8 +147,8 @@ def test_reinforced_memory_decays_slower_in_search_results():
by_id = {item["id"]: item for item in processed}

assert (
by_id["reinforced"]["decay_factor"]
> by_id["unreinforced"]["decay_factor"]
by_id["reinforced"]["effective_retention"]
> by_id["unreinforced"]["effective_retention"]
)
assert processed[0]["id"] == "reinforced"

Expand Down Expand Up @@ -216,7 +216,7 @@ def test_search_results_use_type_specific_decay_rate():
processed = manager.process_search_results(results, "keyword")
by_id = {item["id"]: item for item in processed}

assert by_id["working"]["decay_factor"] < by_id["long"]["decay_factor"]
assert by_id["working"]["effective_retention"] < by_id["long"]["effective_retention"]
assert processed[0]["id"] == "long"


Expand Down Expand Up @@ -338,7 +338,7 @@ def test_search_results_do_not_demote_unmarked_memories():

assert processed[0]["forgotten_score_multiplier"] == pytest.approx(1.0)
assert processed[0]["final_score"] == pytest.approx(
0.85 * processed[0]["decay_factor"]
0.85 * processed[0]["effective_retention"]
)


Expand All @@ -355,7 +355,7 @@ def test_search_results_use_storage_score_for_ranking():

assert processed[0]["original_score"] == pytest.approx(0.92)
assert processed[0]["final_score"] == pytest.approx(
0.92 * processed[0]["decay_factor"]
0.92 * processed[0]["effective_retention"]
)


Expand Down
Loading
Loading