(self, keys: List[BufferKey])
| 236 | return keys |
| 237 | |
| 238 | def _finalize_keys(self, keys: List[BufferKey]) -> List[SaveItem]: |
| 239 | out: List[SaveItem] = [] |
| 240 | for key in keys: |
| 241 | entry = self._buffers.get(key) |
| 242 | if not entry: |
| 243 | continue |
| 244 | payload = entry.snapshot_payload() |
| 245 | if payload is not None: |
| 246 | out.append( |
| 247 | SaveItem( |
| 248 | item_id=entry.item_id, |
| 249 | event=key[3], |
| 250 | conversation_id=key[0], |
| 251 | thread_id=key[1], |
| 252 | task_id=key[2], |
| 253 | payload=payload, |
| 254 | role=entry.role or Role.AGENT, |
| 255 | agent_name=entry.agent_name, |
| 256 | metadata=None, # Buffered entries don't have metadata |
| 257 | ) |
| 258 | ) |
| 259 | if key in self._buffers: |
| 260 | del self._buffers[key] |
| 261 | return out |
| 262 | |
| 263 | def flush_task( |
| 264 | self, |
no test coverage detected