Private
Public Access
Merge branch 'tier2/phase2_4_5_call_site_completion_20260621' into tier2/code_path_audit_20260607
This commit is contained in:
+12
-6
@@ -2568,6 +2568,7 @@ def _send_grok(md_content: str, user_message: str, base_dir: str,
|
||||
Runs synchronously in the caller thread; synchronizes Grok history using _grok_history_lock.
|
||||
"""
|
||||
from src.openai_compatible import OpenAICompatibleRequest, _classify_openai_compatible_error
|
||||
from src.openai_schemas import ChatMessage
|
||||
try:
|
||||
client = _ensure_grok_client()
|
||||
tools: list[Metadata] | None = _get_deepseek_tools() or None
|
||||
@@ -2584,8 +2585,9 @@ def _send_grok(md_content: str, user_message: str, base_dir: str,
|
||||
_grok_history.append({"role": "user", "content": user_content})
|
||||
def _build_grok_request(_round_idx: int) -> OpenAICompatibleRequest:
|
||||
with _grok_history_lock:
|
||||
messages: list[Metadata] = [{"role": "system", "content": f"{_get_combined_system_prompt()}\n\n<context>\n{md_content}\n</context>"}]
|
||||
messages.extend(_grok_history)
|
||||
history_msgs: list[ChatMessage] = [ChatMessage(role=m["role"], content=m["content"]) for m in _grok_history]
|
||||
messages: list[ChatMessage] = [ChatMessage(role="system", content=f"{_get_combined_system_prompt()}\n\n<context>\n{md_content}\n</context>")]
|
||||
messages.extend(history_msgs)
|
||||
extra_body: Metadata = {}
|
||||
if caps.web_search:
|
||||
extra_body["search_parameters"] = {"mode": "auto"}
|
||||
@@ -2653,6 +2655,7 @@ def _send_minimax(md_content: str, user_message: str, base_dir: str,
|
||||
Runs synchronously in the caller thread; synchronizes MiniMax history using _minimax_history_lock.
|
||||
"""
|
||||
from src.openai_compatible import OpenAICompatibleRequest
|
||||
from src.openai_schemas import ChatMessage
|
||||
try:
|
||||
_ensure_minimax_client()
|
||||
tools: list[Metadata] | None = _get_deepseek_tools() or None
|
||||
@@ -2663,8 +2666,9 @@ def _send_minimax(md_content: str, user_message: str, base_dir: str,
|
||||
_minimax_history.append({"role": "user", "content": user_message})
|
||||
def _build_minimax_request(_round_idx: int) -> OpenAICompatibleRequest:
|
||||
with _minimax_history_lock:
|
||||
messages: list[Metadata] = [{"role": "system", "content": f"{_get_combined_system_prompt()}\n\n<context>\n{md_content}\n</context>"}]
|
||||
messages.extend(_minimax_history)
|
||||
history_msgs: list[ChatMessage] = [ChatMessage(role=m["role"], content=m["content"]) for m in _minimax_history]
|
||||
messages: list[ChatMessage] = [ChatMessage(role="system", content=f"{_get_combined_system_prompt()}\n\n<context>\n{md_content}\n</context>")]
|
||||
messages.extend(history_msgs)
|
||||
return OpenAICompatibleRequest(
|
||||
messages=messages, model=_model, temperature=_temperature, top_p=_top_p,
|
||||
max_tokens=min(_max_tokens, 8192), stream=stream, stream_callback=stream_callback,
|
||||
@@ -2893,6 +2897,7 @@ def _send_llama(md_content: str, user_message: str, base_dir: str,
|
||||
Runs synchronously in the caller thread; synchronizes history using _llama_history_lock.
|
||||
"""
|
||||
from src.openai_compatible import OpenAICompatibleRequest, _classify_openai_compatible_error
|
||||
from src.openai_schemas import ChatMessage
|
||||
try:
|
||||
if "localhost" in _llama_base_url or "127.0.0.1" in _llama_base_url:
|
||||
return _send_llama_native(md_content, user_message, base_dir, file_items, discussion_history, stream, pre_tool_callback, qa_callback, stream_callback, patch_callback)
|
||||
@@ -2910,8 +2915,9 @@ def _send_llama(md_content: str, user_message: str, base_dir: str,
|
||||
_llama_history.append({"role": "user", "content": user_content})
|
||||
def _build_llama_request(_round_idx: int) -> OpenAICompatibleRequest:
|
||||
with _llama_history_lock:
|
||||
messages: list[Metadata] = [{"role": "system", "content": f"{_get_combined_system_prompt()}\n\n<context>\n{md_content}\n</context>"}]
|
||||
messages.extend(_llama_history)
|
||||
history_msgs: list[ChatMessage] = [ChatMessage(role=m["role"], content=m["content"]) for m in _llama_history]
|
||||
messages: list[ChatMessage] = [ChatMessage(role="system", content=f"{_get_combined_system_prompt()}\n\n<context>\n{md_content}\n</context>")]
|
||||
messages.extend(history_msgs)
|
||||
return OpenAICompatibleRequest(
|
||||
messages=messages, model=_model, temperature=_temperature, top_p=_top_p,
|
||||
max_tokens=_max_tokens, stream=stream, stream_callback=stream_callback,
|
||||
|
||||
@@ -1841,12 +1841,13 @@ class AppController:
|
||||
|
||||
def _process_pending_gui_tasks(self) -> None:
|
||||
"""Processes pending GUI tasks from the queue on the main render thread."""
|
||||
from src.api_hooks import WebSocketMessage
|
||||
now = time.time()
|
||||
if hasattr(self, 'event_queue') and hasattr(self.event_queue, 'websocket_server') and self.event_queue.websocket_server:
|
||||
if now - self._last_telemetry_time >= 1.0:
|
||||
self._last_telemetry_time = now
|
||||
metrics = self.perf_monitor.get_metrics()
|
||||
self.event_queue.websocket_server.broadcast("telemetry", metrics)
|
||||
self.event_queue.websocket_server.broadcast(WebSocketMessage(channel="telemetry", payload=metrics))
|
||||
|
||||
if not self._pending_gui_tasks: return
|
||||
|
||||
|
||||
+3
-1
@@ -34,6 +34,8 @@ import queue
|
||||
from pathlib import Path
|
||||
from typing import Callable, Any, Dict, List, Tuple, Optional
|
||||
|
||||
from src.api_hooks import WebSocketMessage
|
||||
|
||||
|
||||
class EventEmitter:
|
||||
"""
|
||||
@@ -112,7 +114,7 @@ class AsyncEventQueue:
|
||||
elif hasattr(payload, '__dict__'):
|
||||
serializable_payload = vars(payload)
|
||||
|
||||
self.websocket_server.broadcast("events", {"event": event_name, "payload": serializable_payload})
|
||||
self.websocket_server.broadcast(WebSocketMessage(channel="events", payload={"event": event_name, "payload": serializable_payload}))
|
||||
|
||||
def get(self) -> Tuple[str, Any]:
|
||||
"""
|
||||
|
||||
@@ -193,6 +193,36 @@ class LogRegistry:
|
||||
self.data[session_id]['whitelisted'] = whitelisted
|
||||
self.save_registry() # Save after update
|
||||
|
||||
def set_session_start_time(self, session_id: str, start_time: datetime | str) -> None:
|
||||
"""
|
||||
Updates the start_time of an existing session.
|
||||
|
||||
Used by tests and maintenance tools to backdate a session for pruning
|
||||
verification. Creates a new Session with the updated start_time while
|
||||
preserving all other fields (Session is frozen).
|
||||
|
||||
Args:
|
||||
session_id (str): Unique identifier for the session.
|
||||
start_time (datetime|str): The new start timestamp.
|
||||
[C: tests/test_logging_e2e.py:test_logging_e2e]
|
||||
"""
|
||||
if session_id not in self.data:
|
||||
print(f"Error: Session ID '{session_id}' not found for start_time update.")
|
||||
return
|
||||
if isinstance(start_time, datetime):
|
||||
start_time_str: str = start_time.isoformat()
|
||||
else:
|
||||
start_time_str = start_time
|
||||
existing = self.data[session_id]
|
||||
self.data[session_id] = Session(
|
||||
session_id=existing.session_id,
|
||||
path=existing.path,
|
||||
start_time=start_time_str,
|
||||
whitelisted=existing.whitelisted,
|
||||
metadata=existing.metadata,
|
||||
)
|
||||
self.save_registry()
|
||||
|
||||
def is_session_whitelisted(self, session_id: str) -> bool:
|
||||
"""
|
||||
Checks if a specific session is marked as whitelisted.
|
||||
|
||||
Reference in New Issue
Block a user