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
7 changes: 6 additions & 1 deletion backend/src/agents/main_agent/config/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,12 @@ class Defaults:
AWS_REGION = "us-west-2"

# --- Memory Retrieval ---
MEMORY_RELEVANCE_SCORE = 0.7
# Retrieved long-term memory records scoring below this are dropped. 0.7
# dropped everything realistic: in dev, a natural question about a stored
# fact scored 0.57-0.67 while unrelated records scored <= 0.40, so no turn
# ever received memory context (docs/specs/memory-baseline-decision.md).
# Override per environment with AGENTCORE_MEMORY_RELEVANCE_SCORE.
MEMORY_RELEVANCE_SCORE = 0.5
MEMORY_TOP_K = 10

# --- Compaction ---
Expand Down
172 changes: 115 additions & 57 deletions backend/src/apis/app_api/sessions/services/session_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@

import logging
import os
from typing import Optional
from typing import Any, List, Optional
from datetime import datetime, timezone
from decimal import Decimal

Expand Down Expand Up @@ -336,69 +336,127 @@ def delete_agentcore_memory(self, session_id: str, user_id: str) -> None:

client = boto3.client('bedrock-agentcore', region_name=config.region)

# List all events for this session with pagination (max 100 per request)
all_event_ids = []
next_token = None

try:
while True:
list_params = {
'memoryId': config.memory_id,
'actorId': user_id,
'sessionId': session_id,
'maxResults': 100 # API max is 100
}
if next_token:
list_params['nextToken'] = next_token

events_response = client.list_events(**list_params)
events = events_response.get('events', [])

# Extract event IDs from this page
for event in events:
if event.get('eventId'):
all_event_ids.append(event['eventId'])

# Check for more pages
next_token = events_response.get('nextToken')
if not next_token:
break

except client.exceptions.ResourceNotFoundException:
# Session doesn't exist in AgentCore Memory - nothing to delete
logger.debug("Session not found in AgentCore Memory")
return
except Exception as e:
logger.warning("Failed to list events for session")
return

if not all_event_ids:
logger.debug("No events found for session in AgentCore Memory")
return

# Delete events sequentially - this runs in background so no need
# for parallel execution overhead
deleted_count = 0
for event_id in all_event_ids:
try:
client.delete_event(
memoryId=config.memory_id,
actorId=user_id,
sessionId=session_id,
eventId=event_id
)
deleted_count += 1
except Exception as e:
logger.warning("Failed to delete event from AgentCore Memory")

logger.info("Deleted events from AgentCore Memory")
self._delete_session_events(client, config.memory_id, session_id, user_id)
# Events expire after 90 days; the summaries extracted from them do
# not, so purge runs whether or not any events were left.
self._purge_session_summaries(client, config.memory_id, session_id, user_id)

except ImportError:
logger.debug("AgentCore Memory SDK not available, skipping content deletion")
except Exception as e:
# Log but don't raise - content deletion failures shouldn't block session deletion
logger.error("Failed to delete AgentCore Memory content for session")

def _delete_session_events(self, client: Any, memory_id: str, session_id: str, user_id: str) -> None:
"""Delete every short-term event of one session (list, then delete one by one)."""
# List all events for this session with pagination (max 100 per request)
all_event_ids = []
next_token = None

try:
while True:
list_params = {
'memoryId': memory_id,
'actorId': user_id,
'sessionId': session_id,
'maxResults': 100 # API max is 100
}
if next_token:
list_params['nextToken'] = next_token

events_response = client.list_events(**list_params)
events = events_response.get('events', [])

# Extract event IDs from this page
for event in events:
if event.get('eventId'):
all_event_ids.append(event['eventId'])

# Check for more pages
next_token = events_response.get('nextToken')
if not next_token:
break

except client.exceptions.ResourceNotFoundException:
# Session doesn't exist in AgentCore Memory - nothing to delete
logger.debug("Session not found in AgentCore Memory")
return
except Exception as e:
logger.warning("Failed to list events for session")
return

if not all_event_ids:
logger.debug("No events found for session in AgentCore Memory")
return

# Delete events sequentially - this runs in background so no need
# for parallel execution overhead
deleted_count = 0
for event_id in all_event_ids:
try:
client.delete_event(
memoryId=memory_id,
actorId=user_id,
sessionId=session_id,
eventId=event_id
)
deleted_count += 1
except Exception as e:
logger.warning("Failed to delete event from AgentCore Memory")

logger.info("Deleted events from AgentCore Memory")

def _purge_session_summaries(self, client: Any, memory_id: str, session_id: str, user_id: str) -> None:
"""Delete the long-term SUMMARIZATION records extracted from this session.

Summaries live under a per-session namespace
(``/strategies/{summaryId}/actors/{userId}/sessions/{sessionId}/``), so
they can be found and removed exactly. Semantic facts and preferences
live under the actor namespace, are consolidated across sessions and
carry no source-session metadata, so they are not attributable to one
session and are left alone (docs/specs/memory-baseline-decision.md).
"""
from apis.app_api.memory.services.memory_service import _get_strategy_namespaces

_, _, summary_strategy_id = _get_strategy_namespaces()
if not summary_strategy_id:
logger.debug("No summary strategy discovered, skipping summary purge")
return

session_ns = f"/strategies/{summary_strategy_id}/actors/{user_id}/sessions/{session_id}"
record_ids: List[str] = []
next_token = None
try:
while True:
params = {"memoryId": memory_id, "namespace": session_ns, "maxResults": 100}
if next_token:
params["nextToken"] = next_token
page = client.list_memory_records(**params)
for record in page.get("memoryRecordSummaries", []):
# The namespace filter is a prefix; keep only this exact session.
if any(ns.rstrip("/") == session_ns for ns in record.get("namespaces") or []):
record_ids.append(record["memoryRecordId"])
next_token = page.get("nextToken")
if not next_token:
break
except Exception:
logger.warning("Failed to list summary records for session")
return

deleted = 0
for i in range(0, len(record_ids), 100):
chunk = record_ids[i:i + 100]
try:
resp = client.batch_delete_memory_records(
memoryId=memory_id,
records=[{"memoryRecordId": rid} for rid in chunk],
)
deleted += len(resp.get("successfulRecords", chunk))
except Exception:
logger.warning("Failed to delete summary records for session")
if record_ids:
logger.info("Purged %d of %d summary records for deleted session", deleted, len(record_ids))

def delete_session_files(self, session_id: str) -> None:
"""
Delete all files associated with a session (sync, for background tasks).
Expand Down
28 changes: 28 additions & 0 deletions backend/src/apis/app_api/shares/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,27 @@
logger = logging.getLogger(__name__)



def _skip_long_term_extraction(mgr: Any) -> None:
"""Make every event ``mgr`` writes skip long-term memory extraction.

``AgentCoreMemorySessionManager.create_message`` has no extraction-mode
argument, but all of its writes (conversational events via
``MemoryClient.create_event`` and oversized blob events) end in the data
plane client's ``create_event``. Wrapping that one call on this manager's
own client sets ``extractionMode="SKIP"`` without copying the SDK's
message conversion. The manager is built per fork, so nothing else shares
the wrapped client.
"""
gmdp = mgr.memory_client.gmdp_client
original = gmdp.create_event

def create_event(**kwargs: Any) -> Any:
kwargs.setdefault("extractionMode", "SKIP")
return original(**kwargs)

gmdp.create_event = create_event

class ShareService:
"""Handles share CRUD operations against the shared-conversations DynamoDB table."""

Expand Down Expand Up @@ -440,6 +461,12 @@ async def _copy_messages_to_memory(
Converts each MessageResponse dict to SessionMessage format and
persists via create_message to the "default" namespace.

The messages were written by someone else (the share's owner), but
they land under the forking user's actor. Every event is therefore
written with ``extractionMode="SKIP"``: it stays in short-term memory,
so the fork's history loads, but it never feeds long-term extraction,
so another person's content does not become the forker's "memories".

Returns:
Number of messages successfully written.
"""
Expand Down Expand Up @@ -475,6 +502,7 @@ async def _copy_messages_to_memory(
mgr = AgentCoreMemorySessionManager(
agentcore_memory_config=config, region_name=aws_region
)
_skip_long_term_extraction(mgr)

count = 0
for idx, msg_dict in enumerate(snapshot_messages):
Expand Down
15 changes: 8 additions & 7 deletions backend/src/apis/inference_api/chat/app_context_dispatch.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,14 @@
keyed by the App's bound resource URI.

Storage (decision #3): `agent.state` is the live Strands `AgentState` of
the cached conversation agent. Multi-turn continuity in cloud rides the
in-process LRU agent cache (AgentCore Memory is write-only for continuity —
see docs/specs/MAX_TOKENS_CONTINUE_SESSION_RESTORE_ANALYSIS.md), so the
same `agent.state` survives turn boundaries for free; a cold start /
eviction drops the *entire* conversation anyway, so a dropped pending
context there is consistent with existing behavior, not a new regression.
No `TurnBasedSessionManager` / Memory change is needed.
the cached conversation agent, so the same `agent.state` survives turn
boundaries while the in-process LRU agent cache holds that agent. A cold
start or eviction restores the conversation from AgentCore Memory (the
session manager's restore branch, which also restores agent state that was
synced after a turn), but this dispatch runs without a model turn, so a
payload stashed here is only synced by the next turn. Losing it to an
eviction before then is an accepted loss for a best-effort UI hint; no
`TurnBasedSessionManager` / Memory change is needed.

`AgentState` in strands 1.40 is a `.get()/.set()/.delete()` store whose
`.get()` returns a **deep copy** — nested in-place mutation does NOT
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ class TestRetrievalThresholdEnvVars:
def test_uses_default_thresholds(
self, mock_tbsm, mock_retrieval, mock_mem_config, mock_discover, monkeypatch
):
"""Default relevance_score=0.7 and top_k=10 when env vars not set."""
"""Default relevance_score=0.5 and top_k=10 when env vars not set."""
from agents.main_agent.session.session_factory import SessionFactory

monkeypatch.delenv("AGENTCORE_MEMORY_RELEVANCE_SCORE", raising=False)
Expand All @@ -208,7 +208,7 @@ def test_uses_default_thresholds(
# retrieved per message unless explicitly re-enabled.
assert mock_retrieval.call_count == 2
for c in mock_retrieval.call_args_list:
assert c == call(top_k=10, relevance_score=0.7)
assert c == call(top_k=10, relevance_score=0.5)

@patch("agents.main_agent.session.session_factory.AGENTCORE_MEMORY_AVAILABLE", True)
@patch("agents.main_agent.session.session_factory._discover_strategy_ids")
Expand Down
98 changes: 98 additions & 0 deletions backend/tests/apis/app_api/test_session_delete_memory_purge.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
"""Deleting a session purges its AgentCore short-term events AND the summary
records extracted from it (Shared Projects Phase 0.2).

Semantic facts and preferences are actor-scoped and carry no source session,
so they are deliberately left alone; these tests pin both halves.
"""

from types import SimpleNamespace
from unittest.mock import MagicMock, patch

from apis.app_api.sessions.services.session_service import SessionService

USER = "a1b2c3d4-0000-4000-8000-000000000001"
SESSION = "5b1f3c2e-8f7a-4d2b-9c1e-000000000001"
OTHER_SESSION = SESSION + "-sibling"
SUMMARY = "ConversationSummary-test"
SEMANTIC = "SemanticFactExtraction-test"


def _ns(strategy: str, session: str | None = None) -> str:
base = f"/strategies/{strategy}/actors/{USER}/"
return base + (f"sessions/{session}/" if session else "")


def _client(events=(), records=()):
client = MagicMock()
client.exceptions.ResourceNotFoundException = type("RNF", (Exception,), {})
client.list_events.return_value = {"events": [{"eventId": e} for e in events]}
client.list_memory_records.return_value = {"memoryRecordSummaries": list(records)}
client.batch_delete_memory_records.side_effect = lambda memoryId, records: {
"successfulRecords": records,
"failedRecords": [],
}
return client


def _run(client, summary_id=SUMMARY):
config = SimpleNamespace(is_cloud_mode=True, memory_id="mem-test", region="us-west-2")
with patch("agents.main_agent.session.memory_config.load_memory_config", return_value=config), \
patch("boto3.client", return_value=client), \
patch(
"apis.app_api.memory.services.memory_service._get_strategy_namespaces",
return_value=(SEMANTIC, "Pref-test", summary_id),
):
SessionService().delete_agentcore_memory(SESSION, USER)


def test_deletes_events_and_only_this_sessions_summaries():
client = _client(
events=["e1", "e2"],
records=[
{"memoryRecordId": "sum-this", "namespaces": [_ns(SUMMARY, SESSION)]},
# A prefix match on a different session id must not be deleted.
{"memoryRecordId": "sum-sibling", "namespaces": [_ns(SUMMARY, OTHER_SESSION)]},
],
)
_run(client)

assert client.delete_event.call_count == 2
listed = client.list_memory_records.call_args.kwargs
assert listed["namespace"] == f"/strategies/{SUMMARY}/actors/{USER}/sessions/{SESSION}"
client.batch_delete_memory_records.assert_called_once_with(
memoryId="mem-test", records=[{"memoryRecordId": "sum-this"}]
)


def test_purges_summaries_even_when_no_events_remain():
# Events expire after 90 days; the summaries extracted from them do not.
client = _client(events=[], records=[{"memoryRecordId": "sum-old", "namespaces": [_ns(SUMMARY, SESSION)]}])
_run(client)

client.delete_event.assert_not_called()
client.batch_delete_memory_records.assert_called_once()


def test_never_touches_actor_level_namespaces():
client = _client(events=["e1"], records=[])
_run(client)

for call in client.list_memory_records.call_args_list:
assert "/sessions/" in call.kwargs["namespace"]
client.batch_delete_memory_records.assert_not_called()


def test_skips_purge_without_a_summary_strategy():
client = _client(events=["e1"])
_run(client, summary_id=None)

client.list_memory_records.assert_not_called()
assert client.delete_event.call_count == 1


def test_purge_failure_does_not_raise():
client = _client(events=["e1"])
client.list_memory_records.side_effect = RuntimeError("boom")
_run(client) # background task: logs, never raises

assert client.delete_event.call_count == 1
Loading