Skip to content
Closed
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
249 changes: 188 additions & 61 deletions ovos_core/intent_services/fallback_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
import operator
import threading
import time
from _thread import LockType
from collections import namedtuple
from typing import Callable, Dict, List, Optional, Tuple, Union

Expand All @@ -41,18 +42,17 @@ def __init__(self, bus: Optional[Union[MessageBusClient, FakeBus]] = None,
config = config if config is not None else Configuration().get("skills", {}).get("fallbacks", {})
super().__init__(bus, config)
self.registered_fallbacks: Dict[str, int] = {} # skill_id: priority
self._registered_fallbacks_lock = threading.RLock()
self._fallback_session_locks: Dict[str, Tuple[LockType, int]] = {}
self._fallback_session_locks_lock = threading.Lock()
# skill_id -> (start_handler, response_handler) wired for the
# done-signal translation, so they can be removed on deregister
self._lifecycle_handlers: Dict[str, Tuple[Callable, Callable]] = {}
self._fallback_response_event = threading.Event()
self.bus.on("ovos.skills.fallback.register", self.handle_register_fallback)
self.bus.on("ovos.skills.fallback.deregister", self.handle_deregister_fallback)

def _wire_lifecycle(self, skill_id: str) -> None:
"""Translate lifecycle done-signal for a fallback skill."""
if skill_id in self._lifecycle_handlers:
return

def _on_start(message: Message) -> None:
HandlerLifecycle(self.bus, message, skill_id=skill_id,
handler_name=f"{skill_id}.fallback").start()
Expand All @@ -64,17 +64,21 @@ def _on_response(message: Message) -> None:
HandlerLifecycle(self.bus, message, skill_id=skill_id,
handler_name=f"{skill_id}.fallback").complete()

self.bus.on(f"ovos.skills.fallback.{skill_id}.start", _on_start)
self.bus.on(f"ovos.skills.fallback.{skill_id}.response", _on_response)
self._lifecycle_handlers[skill_id] = (_on_start, _on_response)
with self._registered_fallbacks_lock:
if skill_id in self._lifecycle_handlers:
return
self.bus.on(f"ovos.skills.fallback.{skill_id}.start", _on_start)
self.bus.on(f"ovos.skills.fallback.{skill_id}.response", _on_response)
self._lifecycle_handlers[skill_id] = (_on_start, _on_response)

def _unwire_lifecycle(self, skill_id: str) -> None:
handlers = self._lifecycle_handlers.pop(skill_id, None)
if not handlers:
return
start_handler, response_handler = handlers
self.bus.remove(f"ovos.skills.fallback.{skill_id}.start", start_handler)
self.bus.remove(f"ovos.skills.fallback.{skill_id}.response", response_handler)
with self._registered_fallbacks_lock:
handlers = self._lifecycle_handlers.pop(skill_id, None)
if not handlers:
return
start_handler, response_handler = handlers
self.bus.remove(f"ovos.skills.fallback.{skill_id}.start", start_handler)
self.bus.remove(f"ovos.skills.fallback.{skill_id}.response", response_handler)

def handle_register_fallback(self, message: Message) -> None:
skill_id = message.data.get("skill_id")
Expand All @@ -84,12 +88,13 @@ def handle_register_fallback(self, message: Message) -> None:

# check if .conf is overriding the priority for this skill
priority_overrides = self.config.get("fallback_priorities", {})
if skill_id in priority_overrides:
new_priority = priority_overrides.get(skill_id)
LOG.info(f"forcing {skill_id} fallback priority from {priority} to {new_priority}")
self.registered_fallbacks[skill_id] = new_priority
else:
self.registered_fallbacks[skill_id] = priority
with self._registered_fallbacks_lock:
if skill_id in priority_overrides:
new_priority = priority_overrides.get(skill_id)
LOG.info(f"forcing {skill_id} fallback priority from {priority} to {new_priority}")
self.registered_fallbacks[skill_id] = new_priority
else:
self.registered_fallbacks[skill_id] = priority

# report this skill's fallback dispatch lifecycle as the framework
# done-signal so an orchestrator can resolve it (no skill_id -> skip)
Expand All @@ -98,10 +103,36 @@ def handle_register_fallback(self, message: Message) -> None:

def handle_deregister_fallback(self, message: Message) -> None:
skill_id = message.data.get("skill_id")
if skill_id in self.registered_fallbacks:
self.registered_fallbacks.pop(skill_id)
with self._registered_fallbacks_lock:
if skill_id in self.registered_fallbacks:
self.registered_fallbacks.pop(skill_id)
self._unwire_lifecycle(skill_id)

def _fallback_registry_snapshot(self) -> Dict[str, int]:
"""Return a stable fallback registry view for one match operation."""
with self._registered_fallbacks_lock:
return dict(self.registered_fallbacks)

def _acquire_fallback_session_lock(self, session_id: str) -> LockType:
"""Serialize overlapping fallback polls for the same bus session."""
with self._fallback_session_locks_lock:
lock, users = self._fallback_session_locks.get(
session_id, (threading.Lock(), 0))
self._fallback_session_locks[session_id] = (lock, users + 1)
lock.acquire()
return lock

def _release_fallback_session_lock(self, session_id: str,
lock: LockType) -> None:
lock.release()
with self._fallback_session_locks_lock:
current_lock, users = self._fallback_session_locks[session_id]
if users == 1:
self._fallback_session_locks.pop(session_id)
else:
self._fallback_session_locks[session_id] = (
current_lock, users - 1)

def _fallback_allowed(self, skill_id: str) -> bool:
"""Checks if a skill_id is allowed to fallback

Expand Down Expand Up @@ -131,47 +162,142 @@ def _collect_fallback_skills(self, message: Message,
"""
if fb_range is None:
fb_range = FallbackRange(0, 100)
skill_ids = [] # skill_ids that already answered to ping
fallback_skills = [] # skill_ids that want to handle fallback

sess = SessionManager.get(message)
if sess is None:
return fallback_skills
# filter skills outside the fallback_range
in_range = [s for s, p in self.registered_fallbacks.items()
if fb_range.start < p <= fb_range.stop
and s not in (sess.blacklisted_skills or [])]
skill_ids += [s for s in self.registered_fallbacks if s not in in_range]

def handle_ack(msg):
skill_id = msg.data["skill_id"]
if msg.data.get("can_handle", True):
if skill_id in in_range:
fallback_skills.append(skill_id)
LOG.info(f"{skill_id} will try to handle fallback")
else:
LOG.debug(f"{skill_id} is out of range, skipping")
else:
LOG.debug(f"{skill_id} does NOT WANT to try to handle fallback")
skill_ids.append(skill_id)
self._fallback_response_event.set()

if in_range: # no need to search if no skills available
self.bus.on("ovos.skills.fallback.pong", handle_ack)

return []

registered_fallbacks = self._fallback_registry_snapshot()
pool = [
skill_id for skill_id, priority in sorted(
registered_fallbacks.items(), key=operator.itemgetter(1))
if fb_range.start < priority <= fb_range.stop
and skill_id not in (sess.blacklisted_skills or [])
and self._fallback_allowed(skill_id)
]
if not pool:
return []

session_id = sess.session_id
session_lock = self._acquire_fallback_session_lock(session_id)
responses: Dict[str, Optional[bool]] = {
skill_id: None for skill_id in pool
}
response_event = threading.Event()
response_lock = threading.Lock()
handlers: Dict[str, Callable] = {}

def _record(expected_skill_id: str, can_handle) -> None:
"""First answer for a skill_id wins, regardless of which pong
topic (addressed or broadcast) it arrived on -- a skill running
fixed ovos-workshop (#465) answers BOTH ping families during the
migration window and must only count once."""
valid = isinstance(can_handle, bool)
with response_lock:
if responses[expected_skill_id] is not None:
return
responses[expected_skill_id] = can_handle if valid else False
response_event.set()

def make_handler(expected_skill_id: str) -> Callable:
def handle_ack(msg: Message) -> None:
response_session = SessionManager.get(msg)
if response_session is None or \
response_session.session_id != session_id:
return
skill_id = msg.data.get("skill_id")
can_handle = msg.data.get("can_handle")
if skill_id != expected_skill_id:
_record(expected_skill_id, False)
return
_record(expected_skill_id, can_handle)

return handle_ack

def handle_broadcast_pong(msg: Message) -> None:
# DEPRECATION WINDOW (ovos-core kill-switch #837 conventions):
# the general `ovos.skills.fallback.pong` collector is kept
# alongside the skill-addressed one for one deprecation window,
# so that a released ovos-workshop (pre-#465, only answering the
# broadcast ping) still gets picked up. Removable once the
# ovos-workshop floor pin guarantees dual-binding (#465).
response_session = SessionManager.get(msg)
if response_session is None or \
response_session.session_id != session_id:
return
skill_id = msg.data.get("skill_id")
can_handle = msg.data.get("can_handle")
if skill_id not in responses:
return
_record(skill_id, can_handle)

try:
LOG.info("checking for FallbackSkill candidates")
message.data["range"] = (fb_range.start, fb_range.stop)
# wait for all skills to acknowledge they want to answer fallback queries
self.bus.emit(message.forward("ovos.skills.fallback.ping",
message.data))
start = time.time()
while not all(s in skill_ids for s in self.registered_fallbacks) \
and time.time() - start <= 0.5:
self._fallback_response_event.clear()
self._fallback_response_event.wait(0.02)

self.bus.remove("ovos.skills.fallback.pong", handle_ack)
return fallback_skills
for skill_id in pool:
pong_type = f"{skill_id}.fallback.pong"
handler = make_handler(skill_id)
handlers[pong_type] = handler
self.bus.on(pong_type, handler)

# DEPRECATION WINDOW (ovos-core kill-switch #837 conventions):
# bind the broadcast pong collector once per poll round, with
# the same session filter as the addressed collectors above.
broadcast_pong_type = "ovos.skills.fallback.pong"
handlers[broadcast_pong_type] = handle_broadcast_pong
self.bus.on(broadcast_pong_type, handle_broadcast_pong)

query_data = {
"utterances": list(message.data.get("utterances", [])),
"lang": message.data.get("lang")
}
for skill_id in pool:
# FALLBACK-1 section 6.1 defines this as a dotted-addressed
# reply derived from the inbound utterance envelope.
self.bus.emit(message.reply(
f"{skill_id}.fallback.ping", query_data))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
# DEPRECATION WINDOW: also broadcast the legacy general ping once
# per poll round, so a released ovos-workshop (pre-#465) that
# only binds `ovos.skills.fallback.ping` still answers. Removable
# when the ovos-workshop floor pin guarantees dual-binding
# (#465) -- see ovos-core kill-switch #837 conventions.
self.bus.emit(message.forward(
"ovos.skills.fallback.ping", query_data))

try:
timeout = max(0.0, float(self.config.get(
"fallback_query_timeout", 0.5)))
except (TypeError, ValueError):
LOG.warning("Invalid fallback_query_timeout; using 0.5 seconds")
timeout = 0.5
deadline = time.monotonic() + timeout
while True:
response_event.clear()
with response_lock:
ordered_responses = [responses[skill_id]
for skill_id in pool]
for index, response in enumerate(ordered_responses):
if response is None:
break
if response:
selected = pool[index]
LOG.info(f"{selected} will try to handle fallback")
return [selected]
else:
return []

remaining = deadline - time.monotonic()
if remaining <= 0:
with response_lock:
final_responses = [responses[skill_id]
for skill_id in pool]
for index, response in enumerate(final_responses):
if response:
return [pool[index]]
return []
response_event.wait(remaining)
finally:
for pong_type, handler in handlers.items():
self.bus.remove(pong_type, handler)
self._release_fallback_session_lock(session_id, session_lock)

def _fallback_range(self, utterances: List[str], lang: str,
message: Message, fb_range: FallbackRange) -> Optional[IntentHandlerMatch]:
Expand All @@ -198,11 +324,12 @@ def _fallback_range(self, utterances: List[str], lang: str,
return None
# new style bus api
available_skills = self._collect_fallback_skills(message, fb_range)
fallbacks = [(k, v) for k, v in self.registered_fallbacks.items()
registered_fallbacks = self._fallback_registry_snapshot()
fallbacks = [(k, v) for k, v in registered_fallbacks.items()
if k in available_skills]
sorted_handlers = sorted(fallbacks, key=operator.itemgetter(1))

for skill_id, prio in sorted_handlers:
for skill_id, _priority in sorted_handlers:
if skill_id in (sess.blacklisted_skills or []):
LOG.debug(f"ignoring match, skill_id '{skill_id}' blacklisted by Session '{sess.session_id}'")
continue
Expand Down
8 changes: 5 additions & 3 deletions ovos_core/intent_services/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -252,9 +252,11 @@ def disambiguate_lang(message):
for k in lang_keys:
if k in message.context:
v = standardize_lang(message.context[k])
# closest_lang already applies the "distance below 10" threshold
# and returns None when no candidate is close enough
best_lang = closest_lang(v, valid_langs, max_distance=10)
# closest_lang applies the language-distance threshold and
# returns None when no candidate is close enough. The bound is
# inclusive, so a member language still matches its
# macrolanguage (distance 10, eg. "arz" against "ar")
best_lang = closest_lang(v, valid_langs)
if best_lang is None:
LOG.warning(f"ignoring {k}, {v} is not in enabled languages: {valid_langs}")
continue
Expand Down
27 changes: 19 additions & 8 deletions test/end2end/test_fallback.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,22 +69,31 @@ def _run_fallback_match(self, namespace: str) -> None:
minicroft=minicroft,
skill_ids=[self.skill_id],
eof_msgs=[UTTERANCE_HANDLED],
flip_points=[utt_topic],
flip_points=[
utt_topic,
],
entry_points=[utt_topic],
final_session=final_session,
keep_original_src=[
"ovos.skills.fallback.ping",
# "ovos.skills.fallback.pong", # TODO
],
ignore_messages=["recognizer_loop:audio_output_start",
"recognizer_loop:audio_output_end"],
activation_points=[f"ovos.skills.fallback.{self.skill_id}.request"],
source_message=message,
expected_messages=[
message,
# DEPRECATION WINDOW (kill-switch #837 conventions): core
# still polls both ping families. The released ovos-workshop
# installed by this test run (pre-#465, dev floor pin) only
# binds the legacy broadcast ping, so the skill-addressed
# ping is emitted but goes unanswered here -- it is only
# honored once ovos-workshop >=#465 is the floor pin.
Message(f"{self.skill_id}.fallback.ping",
{"utterances": ["hello world"],
"lang": session.lang}),
Message("ovos.skills.fallback.ping",
{"utterances": ["hello world"], "lang": session.lang, "range": [90, 101]}),
Message("ovos.skills.fallback.pong", {"skill_id": self.skill_id, "can_handle": True}),
{"utterances": ["hello world"],
"lang": session.lang}),
Message("ovos.skills.fallback.pong",
{"skill_id": self.skill_id, "can_handle": True}),
# PIPELINE-1 §9.2: matched notification precedes the dispatch. The
# fallback match_type is the .request topic; it bears no ':' so
# skill_id/intent_name resolve to that topic.
Expand All @@ -95,7 +104,9 @@ def _run_fallback_match(self, namespace: str) -> None:
Message(HANDLER_START,
data={"intent_name": f"ovos.skills.fallback.{self.skill_id}.request"}),
Message(f"ovos.skills.fallback.{self.skill_id}.request",
{"utterances": ["hello world"], "lang": session.lang, "range": [90, 101], "skill_id": self.skill_id}),
{"utterances": ["hello world"],
"lang": session.lang,
"skill_id": self.skill_id}),
Message(f"ovos.skills.fallback.{self.skill_id}.start", {}),
# core reports the fallback dispatch lifecycle as the framework
# done-signal by translating the skill's own .start/.response
Expand Down
Loading
Loading