()
| 739 | logger.debug(f"Mem-reader mode is: {sync_mode}") |
| 740 | |
| 741 | def process_textual_memory(): |
| 742 | if ( |
| 743 | (messages is not None) |
| 744 | and self.config.enable_textual_memory |
| 745 | and self.mem_cubes[mem_cube_id].text_mem |
| 746 | ): |
| 747 | if self.mem_cubes[mem_cube_id].config.text_mem.backend != "tree_text": |
| 748 | add_memory = [] |
| 749 | metadata = TextualMemoryMetadata( |
| 750 | user_id=target_user_id, session_id=target_session_id, source="conversation" |
| 751 | ) |
| 752 | for message in messages: |
| 753 | add_memory.append( |
| 754 | TextualMemoryItem(memory=message["content"], metadata=metadata) |
| 755 | ) |
| 756 | self.mem_cubes[mem_cube_id].text_mem.add(add_memory) |
| 757 | else: |
| 758 | messages_list = [messages] |
| 759 | memories = self.mem_reader.get_memory( |
| 760 | messages_list, |
| 761 | type="chat", |
| 762 | info={"user_id": target_user_id, "session_id": target_session_id}, |
| 763 | mode="fast" if sync_mode == "async" else "fine", |
| 764 | ) |
| 765 | memories_flatten = [m for m_list in memories for m in m_list] |
| 766 | mem_ids: list[str] = self.mem_cubes[mem_cube_id].text_mem.add(memories_flatten) |
| 767 | logger.info( |
| 768 | f"Added memory user {target_user_id} to memcube {mem_cube_id}: {mem_ids}" |
| 769 | ) |
| 770 | # submit messages for scheduler |
| 771 | if self.enable_mem_scheduler and self.mem_scheduler is not None: |
| 772 | if sync_mode == "async": |
| 773 | message_item = ScheduleMessageItem( |
| 774 | user_id=target_user_id, |
| 775 | mem_cube_id=mem_cube_id, |
| 776 | label=MEM_READ_TASK_LABEL, |
| 777 | content=json.dumps(mem_ids), |
| 778 | timestamp=datetime.utcnow(), |
| 779 | task_id=task_id, |
| 780 | ) |
| 781 | self.mem_scheduler.submit_messages(messages=[message_item]) |
| 782 | else: |
| 783 | message_item = ScheduleMessageItem( |
| 784 | user_id=target_user_id, |
| 785 | mem_cube_id=mem_cube_id, |
| 786 | label=ADD_TASK_LABEL, |
| 787 | content=json.dumps(mem_ids), |
| 788 | timestamp=datetime.utcnow(), |
| 789 | task_id=task_id, |
| 790 | ) |
| 791 | logger.info( |
| 792 | f"[DIAGNOSTIC] core.add: Submitting message to scheduler: {message_item.model_dump_json(indent=2)}" |
| 793 | ) |
| 794 | self.mem_scheduler.submit_messages(messages=[message_item]) |
| 795 | |
| 796 | def process_preference_memory(): |
| 797 | if ( |
nothing calls this directly
no test coverage detected