@@ -346,6 +346,9 @@ async def _build_memory_diff(
346346 adds = adds ,
347347 updates = updates ,
348348 deletes = deletes ,
349+ skipped_operations = _serialize_skipped_operations (
350+ getattr (result , "skipped_operations" , [])
351+ ),
349352 )
350353
351354 @tracer (ignore_result = True )
@@ -363,6 +366,7 @@ async def extract_long_term_memories(
363366 allow_self_memory : bool = True ,
364367 allowed_peer_ids : Optional [set [str ]] = None ,
365368 event_search_tags : Optional [List [str ]] = None ,
369+ peer_memory_enabled : bool = True ,
366370 ):
367371 if not agent_evolution_enabled :
368372 effective_types = (
@@ -400,6 +404,7 @@ async def extract_long_term_memories(
400404 archive_uri = archive_uri ,
401405 allowed_memory_types = allowed_memory_types ,
402406 allow_self_memory = allow_self_memory ,
407+ peer_memory_enabled = peer_memory_enabled ,
403408 allowed_peer_ids = allowed_peer_ids ,
404409 event_search_tags = event_search_tags ,
405410 )
@@ -601,6 +606,7 @@ async def _extract_user_memories(
601606 archive_uri : Optional [str ] = None ,
602607 allowed_memory_types : Optional [set [str ]] = None ,
603608 allow_self_memory : bool = True ,
609+ peer_memory_enabled : bool = True ,
604610 allowed_peer_ids : Optional [set [str ]] = None ,
605611 event_search_tags : Optional [List [str ]] = None ,
606612 ) -> "_V3ExtractionResult" :
@@ -642,6 +648,7 @@ async def _extract_user_memories(
642648 allowed_memory_types = allowed_memory_types ,
643649 allow_self = allow_self_memory ,
644650 allowed_peer_ids = allowed_peer_ids ,
651+ peer_memory_enabled = peer_memory_enabled ,
645652 )
646653 isolation_handler .prepare_messages ()
647654 context_provider ._isolation_handler = isolation_handler
@@ -682,6 +689,7 @@ async def _extract_user_memories(
682689 "allowed_memory_types" : allowed_memory_types ,
683690 "allow_self" : allow_self_memory ,
684691 "allowed_peer_ids" : allowed_peer_ids ,
692+ "peer_memory_enabled" : peer_memory_enabled ,
685693 },
686694 metadata = {
687695 "source_extraction_id" : extraction_id ,
@@ -719,6 +727,9 @@ async def _extract_user_memories(
719727 cases = canonical_cases ,
720728 memory_diff = memory_diff ,
721729 case_uri_by_name = _case_uri_by_name (canonical_cases , patch_operations , result ),
730+ skipped_operations = _serialize_skipped_operations (
731+ getattr (result , "skipped_operations" , [])
732+ ),
722733 )
723734
724735 def _session_skill_extraction_enabled (self ) -> bool :
@@ -1223,6 +1234,7 @@ class _V3ExtractionResult:
12231234 cases : list [Case ] = field (default_factory = list )
12241235 memory_diff : dict [str , Any ] | None = None
12251236 case_uri_by_name : dict [str , str ] = field (default_factory = dict )
1237+ skipped_operations : list [dict [str , Any ]] = field (default_factory = list )
12261238
12271239
12281240@dataclass (slots = True )
@@ -2041,6 +2053,22 @@ def _same_memory_file(before: Optional[MemoryFile], after: Optional[MemoryFile])
20412053 )
20422054
20432055
2056+ def _serialize_skipped_operations (items : Any ) -> list [dict [str , Any ]]:
2057+ serialized : list [dict [str , Any ]] = []
2058+ for item in list (items or []):
2059+ if isinstance (item , dict ):
2060+ payload = dict (item )
2061+ else :
2062+ model_dump = getattr (item , "model_dump" , None )
2063+ if not callable (model_dump ):
2064+ continue
2065+ payload = model_dump (mode = "json" , exclude_none = True )
2066+ if isinstance (payload , dict ):
2067+ payload .pop ("source" , None )
2068+ serialized .append (payload )
2069+ return serialized
2070+
2071+
20442072def _v3_extraction_response (
20452073 * ,
20462074 contexts : list [Context ],
@@ -2051,10 +2079,9 @@ def _v3_extraction_response(
20512079
20522080 Historically ``extract_long_term_memories`` returned ``list[Context]`` and
20532081 a number of direct callers still index/compare the return value as a list.
2054- Commit orchestration now also understands the execution-memory style
2055- ``{"contexts": ..., "session_skills": ...}`` shape so it can count
2056- session skills. Preserve the old list shape unless there are actual
2057- session skills to report.
2082+ Commit orchestration also understands a structured response for session
2083+ skills. Preserve the old list shape unless there are actual session skills
2084+ to report.
20582085 """
20592086 skill_dicts : list [dict [str , Any ]] = []
20602087 seen : set [str ] = set ()
@@ -2075,7 +2102,9 @@ def _make_memory_diff(
20752102 adds : list [dict [str , Any ]],
20762103 updates : list [dict [str , Any ]],
20772104 deletes : list [dict [str , Any ]],
2105+ skipped_operations : Optional [list [dict [str , Any ]]] = None ,
20782106) -> dict [str , Any ]:
2107+ skipped = list (skipped_operations or [])
20792108 return {
20802109 "archive_uri" : archive_uri ,
20812110 "trace_id" : tracer .get_trace_id () or None ,
@@ -2085,10 +2114,12 @@ def _make_memory_diff(
20852114 "updates" : list (updates ),
20862115 "deletes" : list (deletes ),
20872116 },
2117+ "skipped_operations" : skipped ,
20882118 "summary" : {
20892119 "total_adds" : len (adds ),
20902120 "total_updates" : len (updates ),
20912121 "total_deletes" : len (deletes ),
2122+ "total_skipped" : len (skipped ),
20922123 },
20932124 }
20942125
@@ -2101,12 +2132,16 @@ def _merge_memory_diffs(
21012132 adds : list [dict [str , Any ]] = []
21022133 updates : list [dict [str , Any ]] = []
21032134 deletes : list [dict [str , Any ]] = []
2135+ skipped_operations : list [dict [str , Any ]] = []
21042136 trace_id = tracer .get_trace_id () or None
21052137 for diff in diffs :
21062138 if not isinstance (diff , dict ):
21072139 continue
21082140 if trace_id is None and diff .get ("trace_id" ):
21092141 trace_id = str (diff .get ("trace_id" ))
2142+ skipped_operations .extend (
2143+ item for item in diff .get ("skipped_operations" , []) if isinstance (item , dict )
2144+ )
21102145 operations = diff .get ("operations" )
21112146 if not isinstance (operations , dict ):
21122147 continue
@@ -2118,6 +2153,7 @@ def _merge_memory_diffs(
21182153 adds = adds ,
21192154 updates = updates ,
21202155 deletes = deletes ,
2156+ skipped_operations = skipped_operations ,
21212157 )
21222158 merged ["trace_id" ] = trace_id
21232159 return merged
@@ -2130,7 +2166,8 @@ def _memory_diff_has_changes(diff: Any) -> bool:
21302166 if not isinstance (summary , dict ):
21312167 return False
21322168 return any (
2133- int (summary .get (key ) or 0 ) > 0 for key in ("total_adds" , "total_updates" , "total_deletes" )
2169+ int (summary .get (key ) or 0 ) > 0
2170+ for key in ("total_adds" , "total_updates" , "total_deletes" , "total_skipped" )
21342171 )
21352172
21362173
0 commit comments