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
41 changes: 0 additions & 41 deletions src/opencortex/context/commit_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -74,14 +74,6 @@ async def commit(
manager._committed_turns.setdefault(sk, set()).add(turn_id)

self._schedule_cited_rewards(cited_uris)
await self._record_valid_skill_citations(
sk=sk,
session_id=session_id,
turn_id=turn_id,
tenant_id=tenant_id,
user_id=user_id,
cited_uris=cited_uris,
)

buffer = manager._conversation_buffers.setdefault(
sk,
Expand Down Expand Up @@ -191,39 +183,6 @@ def _schedule_cited_rewards(self, cited_uris: Optional[List[str]]) -> None:
manager._pending_tasks.add(task)
task.add_done_callback(manager._pending_tasks.discard)

async def _record_valid_skill_citations(
self,
*,
sk: "SessionKey",
session_id: str,
turn_id: str,
tenant_id: str,
user_id: str,
cited_uris: Optional[List[str]],
) -> None:
manager = self._manager
if (
not cited_uris
or not hasattr(manager._orchestrator, "_skill_event_store")
or not manager._orchestrator._skill_event_store
):
return

skill_uris = [uri for uri in cited_uris if "/skills/" in uri]
server_selected = manager._selected_skill_uris.get((sk, turn_id), set())
for uri in skill_uris:
if uri not in server_selected:
logger.debug("[ContextManager] Dropped forged skill citation: %s", uri)
continue
await manager._append_skill_event(
session_id,
turn_id,
uri,
tenant_id,
user_id,
"cited",
)

def _build_write_items(
self,
*,
Expand Down
35 changes: 0 additions & 35 deletions src/opencortex/context/end_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@ class EndRunState:
start_time: float
total_turns: int
fail_fast: bool
session_owner_ids: List[str]
status: str = "closed"
traces: int = 0
knowledge_candidates: int = 0
Expand Down Expand Up @@ -69,9 +68,6 @@ async def end(
start_time=time.monotonic(),
total_turns=len(manager._committed_turns.get(sk, set())),
fail_fast=bool((config or {}).get("fail_fast_end", False)),
session_owner_ids=sorted(
manager._session_memory_owner_ids.get(sk, set())
),
)
try:
await self._wait_for_background_merge(
Expand Down Expand Up @@ -101,12 +97,6 @@ async def end(
tenant_id=tenant_id,
user_id=user_id,
)
self._schedule_autophagy(
state,
session_id=session_id,
tenant_id=tenant_id,
user_id=user_id,
)
await self._wait_for_merge_followups(
sk,
state,
Expand Down Expand Up @@ -283,31 +273,6 @@ async def _persist_source_and_end_session(
exc,
)

def _schedule_autophagy(
self,
state: EndRunState,
*,
session_id: str,
tenant_id: str,
user_id: str,
) -> None:
manager = self._manager
if (
not state.session_owner_ids
or getattr(manager._orchestrator, "_autophagy_kernel", None) is None
):
return
task = asyncio.create_task(
manager._run_autophagy_metabolism(
session_id=session_id,
tenant_id=tenant_id,
user_id=user_id,
owner_ids=state.session_owner_ids,
)
)
manager._pending_tasks.add(task)
task.add_done_callback(manager._pending_tasks.discard)

async def _wait_for_merge_followups(
self,
sk: "SessionKey",
Expand Down
Loading
Loading