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
3 changes: 2 additions & 1 deletion app.py
Original file line number Diff line number Diff line change
Expand Up @@ -330,7 +330,8 @@ def redirect_settings_to_system(subpath=None):
# Smoke-test import (used by the updater to validate a new tree before
# restarting): import all modules but skip channel/scheduler startup.
_smoke_test = _os.environ.get('EVONIC_SMOKE_TEST') == '1'
if (not _reloader_active or _is_reloader_child) and not _smoke_test:
_testing = _os.environ.get('EVONIC_TESTING') == '1'
if (not _reloader_active or _is_reloader_child) and not _smoke_test and not _testing:
# Run SYSTEM.md migration eagerly (not lazily on first GET /api/agents).
# Agents that predate the on-disk SYSTEM.md feature need their file written
# before they start processing messages — otherwise read_file("/_self/SYSTEM.md")
Expand Down
17 changes: 15 additions & 2 deletions backend/agent_runtime/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -167,9 +167,18 @@ def _send_free_notification(agent_id: str):

from models.db import db
notify_msg = "Hey! I'm done and ready to help again. Is there anything I can do?"
message_id = None
try:
db.add_chat_message(session_id, 'assistant', notify_msg,
agent_id=agent_id, metadata={"free_notification": True})
message_id = db.add_chat_message(
session_id, 'assistant', notify_msg,
agent_id=agent_id, metadata={"free_notification": True},
)
message_id = message_id if type(message_id) in (int, str) else None
from models.chatlog import chatlog_manager
chatlog_manager.get(agent_id, session_id).append({
'type': 'final', 'session_id': session_id, 'content': notify_msg,
'metadata': {'free_notification': True}, 'message_id': message_id,
})
except Exception as e:
log.error("[AgentFreeNotify] Failed to save notification message: %s", e)

Expand All @@ -181,6 +190,10 @@ def _send_free_notification(agent_id: str):
'session_id': session_id,
'external_user_id': external_user_id,
'channel_id': channel_id,
'message': notify_msg,
'message_id': message_id,
'metadata': {"free_notification": True},
'role': 'assistant',
})
except Exception as e:
log.error("[AgentFreeNotify] Failed to emit message_received event: %s", e)
Expand Down
87 changes: 39 additions & 48 deletions backend/agent_runtime/llm_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -419,6 +419,9 @@ def run_tool_loop(agent: Dict[str, Any],
from backend.event_stream import event_stream
from models.chatlog import chatlog_manager

def _message_id(value):
return value if type(value) in (int, str) else None

agent_id = agent['id']
db_agent_id = session_db_agent_id or agent_id # which per-agent DB owns this session
external_user_id = agent_context.get('user_id')
Expand All @@ -440,7 +443,6 @@ def run_tool_loop(agent: Dict[str, Any],
_parent_agent_id = None
_loop_ts = int(time.time() * 1000)
chatlog.append({'type': 'turn_begin', 'session_id': session_id, 'ts': _loop_ts})
event_stream.emit('turn_begin', {'session_id': session_id, 'ts': _loop_ts})

tool_trace = []
timeline = []
Expand Down Expand Up @@ -510,10 +512,13 @@ def _emit_task_lifecycle_event(event_name, task_ids):
def _finalize_gate_response(response: str, source: str):
duration = round(time.time() - _loop_start_time, 1)
metadata = {'plugin_gate': source, 'thinking_duration': duration}
db.add_chat_message(session_id, 'assistant', response,
agent_id=db_agent_id, metadata=metadata)
message_id = _message_id(db.add_chat_message(
session_id, 'assistant', response,
agent_id=db_agent_id, metadata=metadata,
))
chatlog.append({'type': 'final', 'session_id': session_id,
'content': response, 'metadata': metadata})
'content': response, 'metadata': metadata,
'message_id': message_id})
chatlog.append({'type': 'turn_end', 'session_id': session_id,
'thinking_duration': duration})
event_stream.emit('final_answer', {
Expand Down Expand Up @@ -949,11 +954,14 @@ def _get_agent_config_ig(agt_id: str) -> dict:
_logger.info("Stop signal received during ATG execution for session %s", session_id)
stop_msg = "Agent stopped by user request."
_atg_stop_dur = round(time.time() - _loop_start_time, 1)
db.add_chat_message(session_id, 'assistant', stop_msg, agent_id=db_agent_id,
metadata={"timeline": timeline, "stopped": True,
"thinking_duration": _atg_stop_dur})
message_id = _message_id(db.add_chat_message(
session_id, 'assistant', stop_msg, agent_id=db_agent_id,
metadata={"timeline": timeline, "stopped": True,
"thinking_duration": _atg_stop_dur},
))
chatlog.append({'type': 'final', 'session_id': session_id, 'content': stop_msg,
'metadata': {'stopped': True, 'thinking_duration': _atg_stop_dur}})
'metadata': {'stopped': True, 'thinking_duration': _atg_stop_dur},
'message_id': message_id})
chatlog.append({'type': 'turn_end', 'session_id': session_id,
'thinking_duration': _atg_stop_dur})
event_stream.emit('final_answer', {
Expand Down Expand Up @@ -1183,10 +1191,14 @@ def _get_agent_config_ig(agt_id: str) -> dict:
_logger.info("Stop signal received for session %s — aborting loop", session_id)
stop_msg = "Agent stopped by user request."
_stop_dur = round(time.time() - _loop_start_time, 1)
db.add_chat_message(session_id, 'assistant', stop_msg, agent_id=db_agent_id,
metadata={"timeline": timeline, "stopped": True, "thinking_duration": _stop_dur})
message_id = _message_id(db.add_chat_message(
session_id, 'assistant', stop_msg, agent_id=db_agent_id,
metadata={"timeline": timeline, "stopped": True,
"thinking_duration": _stop_dur},
))
chatlog.append({'type': 'final', 'session_id': session_id, 'content': stop_msg,
'metadata': {'stopped': True, 'thinking_duration': _stop_dur}})
'metadata': {'stopped': True, 'thinking_duration': _stop_dur},
'message_id': message_id})
_stop_inj = ("[SYSTEM] Your previous reasoning and response were forcefully "
"interrupted by the user via /stop before completion. "
"Await the user's next instruction.")
Expand Down Expand Up @@ -1983,9 +1995,12 @@ def _get_agent_config_ig(agt_id: str) -> dict:
'content': _display_content, 'is_final': True,
'send_as_message': True,
})
db.add_chat_message(session_id, 'assistant', _display_content, agent_id=db_agent_id, metadata=meta)
message_id = _message_id(db.add_chat_message(
session_id, 'assistant', _display_content,
agent_id=db_agent_id, metadata=meta,
))
chatlog.append({'type': 'final', 'session_id': session_id, 'content': _display_content,
'metadata': _cl_meta})
'metadata': _cl_meta, 'message_id': message_id})
chatlog.append({'type': 'turn_end', 'session_id': session_id, 'thinking_duration': _final_dur})
# Archive sub-agent session at turn-end — single-turn only. Explorer &
# kb-organizer are single-shot, so they archive on completion (no need to
Expand Down Expand Up @@ -2344,13 +2359,8 @@ def _normalize(s):
})

# Escalation: ensure a human can see the approval.
# We always fan-out to BOTH web SSE AND messaging channels,
# because has_web_listener() is unreliable — it only checks
# listener registration, not actual SSE delivery. If SSE
# disconnects and reconnects, the approval event may already
# be gone from the ring buffer. Web SSE delivers the approval
# modal in the browser; messaging channels deliver a fallback
# notification via Telegram/WhatsApp.
# Always fan out to BOTH durable web SSE and messaging channels.
# Messaging remains the out-of-browser fallback.
# List of (session_id, external_user_id, channel_id) that received
# approval_required — used to fan-out approval_resolved to all of them.
_escalation_targets: list = []
Expand Down Expand Up @@ -2385,31 +2395,8 @@ def _normalize(s):
})
_escalation_targets.append((_human_session_id, _web_uid, _web_cid))

# Channel (Telegram/WhatsApp): always notify via the super
# agent's messaging channel as a fallback. We no longer gate
# this behind has_web_listener() because the registration
# may exist while the SSE connection is not delivering.

# Verify web SSE delivery with heartbeat-aware check.
# has_web_listener() only confirms a callback is registered;
# this confirms the SSE connection is actually sending heartbeats.
_web_sse_active = False
try:
from routes.realtime import has_active_web_sse
_web_sse_active = (
has_active_web_sse(session_id) or
(_human_session_id and has_active_web_sse(_human_session_id))
)
except Exception:
pass # routes.realtime may not be importable in all contexts

if not _web_sse_active:
_logger.info(
"approval %s: web SSE appears inactive for session %s "
"(heartbeat not received within window) — relying on "
"messaging channel fallback",
pending.approval_id, session_id,
)
# Channel (Telegram/WhatsApp) remains an unconditional
# fallback; browser connection liveness is not authoritative.
_super = db.get_super_agent()
if _super and _super['id'] != agent_id:
_su_uid, _su_cid = _resolve_agent_target(_super['id'])
Expand Down Expand Up @@ -2812,10 +2799,14 @@ def _normalize(s):
_logger.info("Stop signal received for session %s — aborting after tools", session_id)
stop_msg = "Agent stopped by user request."
_stopb_dur = round(time.time() - _loop_start_time, 1)
db.add_chat_message(session_id, 'assistant', stop_msg, agent_id=db_agent_id,
metadata={"timeline": timeline, "stopped": True, "thinking_duration": _stopb_dur})
message_id = _message_id(db.add_chat_message(
session_id, 'assistant', stop_msg, agent_id=db_agent_id,
metadata={"timeline": timeline, "stopped": True,
"thinking_duration": _stopb_dur},
))
chatlog.append({'type': 'final', 'session_id': session_id, 'content': stop_msg,
'metadata': {'stopped': True, 'thinking_duration': _stopb_dur}})
'metadata': {'stopped': True, 'thinking_duration': _stopb_dur},
'message_id': message_id})
_stopb_inj = ("[SYSTEM] Your previous reasoning and response were forcefully "
"interrupted by the user via /stop before completion. "
"Await the user's next instruction.")
Expand Down
11 changes: 10 additions & 1 deletion backend/agent_runtime/notifier.py
Original file line number Diff line number Diff line change
Expand Up @@ -190,10 +190,17 @@ def notify_agent(agent_id: str, tag: str, message: str,
)
else:
meta = dict(metadata) if metadata else {}
db.add_chat_message(
message_id = db.add_chat_message(
target_session_id, role='user', content=full_message,
agent_id=_db_agent_id, metadata=meta if meta else None,
)
message_id = message_id if type(message_id) in (int, str) else None
from models.chatlog import chatlog_manager
chatlog_manager.get(_db_agent_id, target_session_id).append({
'type': 'user', 'session_id': target_session_id,
'content': full_message, 'metadata': meta,
'message_id': message_id,
})
from backend.event_stream import event_stream
event_stream.emit('message_received', {
'agent_id': agent_id,
Expand All @@ -202,6 +209,8 @@ def notify_agent(agent_id: str, tag: str, message: str,
'channel_id': channel_id,
'message': full_message,
'metadata': meta,
'message_id': message_id,
'role': 'user',
})
if deliver_external and channel_id:
from backend.channels.registry import channel_manager
Expand Down
Loading
Loading