""" Support to interface with Alexa Devices. SPDX-License-Identifier: Apache-2.0 For more details about this platform, please refer to the documentation at https://community.home-assistant.io/t/echo-devices-alexa-as-media-player-testers-needed/58639 """ import asyncio from datetime import datetime, timedelta from json import JSONDecodeError, loads import logging import os import random import time from typing import Optional from urllib.parse import urlparse import aiohttp from alexapy import ( AlexaAPI, AlexaLogin, AlexapyConnectionError, AlexapyLoginError, HTTP2EchoClient, __version__ as alexapy_version, hide_serial, obfuscate, ) from alexapy.errors import AlexapyTooManyRequestsError from alexapy.helpers import delete_cookie as alexapy_delete_cookie import async_timeout from homeassistant.components.persistent_notification import ( async_create as async_create_persistent_notification, async_dismiss as async_dismiss_persistent_notification, ) from homeassistant.config_entries import SOURCE_REAUTH from homeassistant.const import ( CONF_EMAIL, CONF_PASSWORD, CONF_URL, EVENT_HOMEASSISTANT_STARTED, EVENT_HOMEASSISTANT_STOP, ) from homeassistant.core import HomeAssistant, callback from homeassistant.data_entry_flow import UnknownFlow from homeassistant.exceptions import ConfigEntryNotReady from homeassistant.helpers import config_validation as cv, device_registry as dr from homeassistant.helpers.dispatcher import async_dispatcher_send from homeassistant.helpers.issue_registry import IssueSeverity, async_create_issue from homeassistant.helpers.update_coordinator import UpdateFailed from homeassistant.loader import async_get_integration from homeassistant.util import dt, slugify import voluptuous as vol from .alexa_entity import AlexaEntityData, get_entity_data, parse_alexa_entities from .config_flow import in_progress_instances from .const import ( ALEXA_COMPONENTS, CONF_ACCOUNTS, CONF_DEBUG, CONF_EXCLUDE_DEVICES, CONF_EXTENDED_ENTITY_DISCOVERY, CONF_INCLUDE_DEVICES, CONF_OAUTH, CONF_OTPSECRET, CONF_PUBLIC_URL, CONF_QUEUE_DELAY, CONF_SCAN_INTERVAL, DATA_ALEXAMEDIA, DATA_LISTENER, DEFAULT_EXTENDED_ENTITY_DISCOVERY, DEFAULT_PUBLIC_URL, DEFAULT_QUEUE_DELAY, DEFAULT_SCAN_INTERVAL, DEPENDENT_ALEXA_COMPONENTS, DOMAIN, HTTP2_ERROR_THRESHOLD, ISSUE_URL, LAST_CALLED_429_BACKOFF_INITIAL_S, LAST_CALLED_429_BACKOFF_MAX_S, LAST_CALLED_COALESCE_WINDOW_MS, LAST_CALLED_CONN_BACKOFF_S, LAST_CALLED_DEBOUNCE_S, LAST_CALLED_ITEMS, LAST_CALLED_LOGIN_BACKOFF_S, LAST_CALLED_LOOKBACK_MS, LAST_CALLED_RETRY_DELAY_S, LAST_CALLED_RETRY_LIMIT, LAST_CALLED_STALE_FUDGE_MS, LAST_CALLED_SUCCESS_PACE_S, LAST_PING_MAX_AGE_SECONDS, LAST_PUSH_INACTIVITY_SECONDS, MIN_TIME_BETWEEN_FORCED_SCANS, MIN_TIME_BETWEEN_SCANS, NOTIFICATION_COOLDOWN, NOTIFY_REFRESH_BACKOFF, NOTIFY_REFRESH_MAX_RETRIES, SCAN_INTERVAL, STARTUP_MESSAGE, ) from .coordinator import AlexaMediaCoordinator from .exceptions import TimeoutException from .helpers import ( _catch_login_errors, _existing_serials, alarm_just_dismissed, calculate_uuid, hide_email, report_relogin_required, safe_get, ) from .metrics import AlexaMetrics, get_metrics from .notify import async_unload_entry as notify_async_unload_entry from .runtime_data import AlexaRuntimeData from .services import AlexaMediaServices _LOGGER = logging.getLogger(__name__) ACCOUNT_CONFIG_SCHEMA = vol.Schema( { vol.Required(CONF_EMAIL): cv.string, vol.Required(CONF_PASSWORD): cv.string, vol.Required(CONF_URL): cv.string, vol.Optional(CONF_INCLUDE_DEVICES, default=[]): vol.All( cv.ensure_list, [cv.string] ), vol.Optional(CONF_EXCLUDE_DEVICES, default=[]): vol.All( cv.ensure_list, [cv.string] ), vol.Optional(CONF_SCAN_INTERVAL, default=SCAN_INTERVAL): cv.time_period, vol.Optional(CONF_QUEUE_DELAY, default=DEFAULT_QUEUE_DELAY): cv.positive_float, vol.Optional(CONF_EXTENDED_ENTITY_DISCOVERY, default=False): cv.boolean, vol.Optional(CONF_DEBUG, default=False): cv.boolean, } ) CONFIG_SCHEMA = vol.Schema( { DOMAIN: vol.Schema( { vol.Optional(CONF_ACCOUNTS): vol.All( cv.ensure_list, [ACCOUNT_CONFIG_SCHEMA] ) } ) }, extra=vol.ALLOW_EXTRA, ) def _valid_voice_summary(summary: object) -> bool: """Return True if summary looks like a real spoken utterance.""" if not isinstance(summary, str): return False summary = summary.strip() return bool(summary) and any(ch.isalnum() for ch in summary) def _queue_last_called_activity( account: dict, *, device_serial: str, customer_id: str | None, activity_ts: int | None, command: str, ) -> None: """Queue or refresh a last_called activity candidate keyed by device/customer.""" if not device_serial: return try: ts = int(activity_ts) if activity_ts is not None else 0 except (TypeError, ValueError): ts = 0 queue: list[dict] = account.setdefault("last_called_activity_queue", []) for item in queue: if ( item.get("serial") == device_serial and item.get("customer_id") == customer_id ): # Keep earliest timestamp like alexa-remote does for the queued burst. prev_ts = int(item.get("activity_ts") or 0) if ts and (prev_ts == 0 or ts < prev_ts): item["activity_ts"] = ts item["command"] = command return queue.append( { "serial": device_serial, "customer_id": customer_id, "activity_ts": ts, "command": command, } ) def _snapshot_last_called_activity_queue(account: dict) -> list[dict]: """Return a shallow snapshot of queued activity entries.""" queue = account.get("last_called_activity_queue") or [] return [dict(item) for item in queue if isinstance(item, dict)] def _remove_last_called_activity_queue_entries( account: dict, resolved_keys: set[tuple[str, str | None]] ) -> None: """Remove resolved queue entries by (serial, customer_id).""" queue = account.get("last_called_activity_queue") or [] account["last_called_activity_queue"] = [ item for item in queue if ( item.get("serial"), item.get("customer_id"), ) not in resolved_keys ] def _valid_utterance_type(record: dict) -> bool: """Filter utterance types similar to Node-RED/alexa-remote logic.""" utterance_type = record.get("utteranceType") if utterance_type in { "DEVICE_ARBITRATION", "ASR_TIMEOUT", "WAKE_WORD_ONLY", }: return False return True def _push_healthy(account: dict) -> bool: """Return True if HTTP2 push is likely usable (enough to skip polling last_called).""" http2 = account.get("http2") if not http2: return False # Hard negative: the underlying transport is closed. client = getattr(http2, "client", None) if client is not None and getattr(client, "is_closed", False): return False # If alexapy has already driven error count to "give up", treat as down. if int(account.get("http2error") or 0) >= HTTP2_ERROR_THRESHOLD: return False last_push = float(account.get("last_push_activity") or 0.0) if last_push and (time.time() - last_push) > LAST_PUSH_INACTIVITY_SECONDS: return False # If we have a recent ping, that's a strong positive. last_ping_dt = getattr(http2, "_last_ping", None) # private, best-effort if last_ping_dt: try: age = time.time() - last_ping_dt.timestamp() # ping is ~299s; allow generous slack for scheduler jitter. if age <= LAST_PING_MAX_AGE_SECONDS: return True # If ping is *very* stale, treat as suspicious but not definitive. # Don't force False here unless you also have other negative signals. except Exception as exc: # pylint: disable=broad-except _LOGGER.debug("Could not evaluate http2 ping age: %s", exc) # Unknown state: object exists and client isn't closed -> assume usable. return True def _network_allowed(login_obj) -> bool: if login_obj.close_requested: return False if login_obj.session.closed: return False if not login_obj.status.get("login_successful"): return False return True def _entity_backed_device_identifiers(account_dict: dict) -> set[str]: """Collect device identifier strings for devices that are backed by entities. Alexa Media Player historically prunes stale HA Device Registry entries by comparing device identifiers against the current *media_player* serials. That works for Echoes, but fails for entity-only devices (e.g., Amazon Indoor Air Quality Monitor) which have no media_player entity. Those devices would get pruned unless we also consider the identifiers referenced by entities we created. """ identifiers: set[str] = set() def _collect_from_device_info(device_info) -> None: if not device_info: return try: # dr.DeviceInfo di_idents = getattr(device_info, "identifiers", None) if di_idents: for ident in di_idents: if isinstance(ident, tuple) and len(ident) == 2: identifiers.add(ident[1]) return except Exception as exc: # pylint: disable=broad-except _LOGGER.debug("Could not extract identifiers from device_info: %s", exc) # dict-style device_info if isinstance(device_info, dict): di_idents = device_info.get("identifiers") if di_idents: for ident in di_idents: if isinstance(ident, tuple) and len(ident) == 2: identifiers.add(ident[1]) # Recursively walk nested entity structures (dict/list/tuple/set) and collect any device_info found. def _walk(obj) -> None: if obj is None: return # Entity-ish object _collect_from_device_info(getattr(obj, "device_info", None)) if isinstance(obj, dict): for v in obj.values(): _walk(v) elif isinstance(obj, (list, tuple, set)): for v in obj: _walk(v) _walk(account_dict.get("entities", {})) return identifiers def _entity_backed_serials(account: dict) -> set[str]: """Return serials that exist only because we created entity-backed devices. These serials may not be discoverable via media player inventory, but should still be considered 'current' so we don't prune/ignore them. """ serials: set[str] = set() entities = account.get("entities") if not isinstance(entities, dict): return serials # Sensors are stored keyed by serial; this is where AIAQM lives. sensors = entities.get("sensor") if isinstance(sensors, dict): serials.update(s for s in sensors.keys() if isinstance(s, str) and s) return serials def _select_last_called_payload_from_records( records: list[dict], queue_snapshot: list[dict], account: dict, existing_serials_local: set[str], ) -> tuple[dict | None, set[tuple[str, str | None]]]: """Select the best last_called payload from raw customer history records.""" if not records or not queue_snapshot: return None, set() watermark = int(account.get("last_called_customer_history_ts") or 0) last_pushed_activity = account.get("last_called_last_pushed_activity") or {} queue_by_key = { (item.get("serial"), item.get("customer_id")): item for item in queue_snapshot if item.get("serial") } def _record_ts(record: dict) -> int: try: return int(record.get("creationTimestamp") or 0) except (TypeError, ValueError): return 0 sorted_records = sorted( (r for r in records if isinstance(r, dict)), key=_record_ts, reverse=True, ) for record in sorted_records: serial = record.get("deviceSerialNumber") if not serial or serial not in existing_serials_local: continue if not _valid_utterance_type(record): continue summary = ((record.get("description") or {}).get("summary") or "").strip() if not _valid_voice_summary(summary): continue try: ts = int(record.get("creationTimestamp") or 0) except (TypeError, ValueError): continue if ts <= watermark: continue if ts <= int(last_pushed_activity.get(serial) or 0): continue for key, queued in queue_by_key.items(): queued_serial, _queued_customer = key queued_ts = int(queued.get("activity_ts") or 0) if serial != queued_serial: continue # Future enhancement: customer/user matching. # The push payload includes destinationUserId, but the current # get_customer_history_records() helper does not expose a matching # user identifier from the raw RVH records, so this check is currently # ineffective. Leave in place in case the API helper is extended later. # # if queued_customer and record.get("customerId") not in (None, queued_customer): # continue if queued_ts and ts < (queued_ts - LAST_CALLED_STALE_FUDGE_MS): continue payload = { "serialNumber": serial, "timestamp": ts, "summary": summary, "response": (record.get("alexaResponse") or "").strip() or None, } return payload, {key} return None, set() async def async_setup(hass, config): """Set up the Alexa domain.""" # Initialize metrics hass.data.setdefault(DOMAIN, {}) hass.data[DOMAIN]["metrics"] = AlexaMetrics(hass) metrics = hass.data[DOMAIN]["metrics"] metrics.start_boot_tracking() integration = await async_get_integration(hass, DOMAIN) integration_name = integration.name or "" _LOGGER.info( STARTUP_MESSAGE.format( name=integration_name, ISSUE_URL=ISSUE_URL, DOMAIN=DOMAIN, version=integration.version, alexapy_version=alexapy_version, ) ) metrics.record_boot_stage("domain_setup") if DOMAIN in config: async_create_issue( hass, DOMAIN, "deprecated_yaml_configuration", is_fixable=False, issue_domain=DOMAIN, severity=IssueSeverity.ERROR, translation_key="deprecated_yaml_configuration", learn_more_url="https://github.com/alandtse/alexa_media_player/wiki/Configuration#configurationyaml", ) _LOGGER.error( "YAML configuration of Alexa Media Player is no longer supported. " "Please remove `alexa_media` from your configuration, " "restart Home Assistant and use the UI to configure it instead. " "Settings > Devices & services > Integrations > ADD INTEGRATION" ) return False return True def _store_and_dispatch_last_called( hass: HomeAssistant, email: str, last_called: dict, force: bool = False, ) -> None: """Store last_called data and dispatch change event if needed. Shared helper used by both the closure-based update_last_called and the module-level _async_update_last_called_global. """ accounts = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}) if email not in accounts: _LOGGER.debug("%s: Account removed during update, skipping", hide_email(email)) return stored_data = accounts[email] payload = dict(last_called) if isinstance(last_called, dict) else {} ts_raw = payload.get("timestamp") try: ts = int(ts_raw or 0) except (TypeError, ValueError): ts = 0 if 0 < ts < 10_000_000_000: ts *= 1000 if ts > 0: payload["timestamp"] = ts prev = stored_data.get("last_called") changed = prev != payload if ts > 0: stored_data["last_called_customer_history_ts"] = max( int(stored_data.get("last_called_customer_history_ts") or 0), ts, ) stored_data["last_called"] = payload if ( force or (prev is None and last_called is not None) or (prev is not None and changed) ): _LOGGER.debug( "%s: last_called changed", hide_email(email), ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"last_called_change": payload}, ) async def _async_update_last_called_global( hass: HomeAssistant, login_obj, email: str, last_called: dict | None = None, force: bool = False, ) -> None: """Update the last called device globally (standalone function). This version is defined at module level for use in background tasks. It delegates storage/dispatch to _store_and_dispatch_last_called. """ if not _network_allowed(login_obj): return accounts = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}) account = accounts.get(email) # If we're doing a "refresh" (no payload) and not forcing, prefer the probe worker # whenever it's available, regardless of http2 connection. if ( (not isinstance(last_called, dict) or not last_called.get("summary")) and not force and account ): trigger = account.get("last_called_probe_trigger") if callable(trigger): trigger("GLOBAL_REFRESH", None) return if not isinstance(last_called, dict) or not last_called.get("summary"): try: # Do not timebox this call; alexapy may back off/sleep on 429. # Serialize RVH calls per account to avoid parallel rate-limited requests. api_lock = None if account: api_lock = account.get("last_called_api_lock") if api_lock is None: last_called = await AlexaAPI.get_last_device_serial(login_obj) else: async with api_lock: last_called = await AlexaAPI.get_last_device_serial(login_obj) except asyncio.CancelledError: # Task cancelled during unload/shutdown; propagate cancellation. raise except AlexapyTooManyRequestsError: _LOGGER.debug( "%s: Rate limited during last_called update; skipping", hide_email(email), ) return except AlexapyConnectionError as exc: _LOGGER.debug( "%s: Connection error during last_called update: %s", hide_email(email), exc, ) return except AlexapyLoginError: _LOGGER.debug( "%s: Login error during last_called update", hide_email(email) ) report_relogin_required(hass, login_obj, email) return except TypeError: _LOGGER.debug( "%s: Error updating last_called: %s", hide_email(email), repr(last_called), ) return if not isinstance(last_called, dict): _LOGGER.debug( "%s: Error updating last_called: unexpected response %s", hide_email(email), repr(last_called), ) return if not _valid_voice_summary(last_called.get("summary")): _LOGGER.debug( "%s: Ignoring last_called with invalid summary", hide_email(email), ) return _LOGGER.debug( "%s: Updated last_called: %s", hide_email(email), hide_serial(last_called) ) _store_and_dispatch_last_called(hass, email, last_called, force) # @retry_async(limit=5, delay=5, catch_exceptions=True) async def async_setup_entry(hass, config_entry): """Set up Alexa Media Player as config entry. This function uses the new runtime_data pattern for type-safe data storage. Legacy hass.data[DATA_ALEXAMEDIA] is maintained for backward compatibility during the migration period. """ _boot_start = time.monotonic() async def close_alexa_media(event=None) -> None: """Clean up Alexa connections.""" _LOGGER.debug("Received shutdown request: %s", event) if accounts := safe_get(hass.data, [DATA_ALEXAMEDIA, "accounts"], {}): for email, _ in accounts.items(): await close_connections(hass, email) async def complete_startup(event=None) -> None: # pylint: disable=unused-argument """Run final tasks after startup.""" _LOGGER.debug("Completing remaining startup tasks.") await asyncio.sleep(10) if hass.data[DATA_ALEXAMEDIA].get("notify_service"): notify = hass.data[DATA_ALEXAMEDIA].get("notify_service") _LOGGER.debug("Refreshing notify targets") await notify.async_register_services() async def relogin(event=None) -> None: """Relogin to Alexa.""" if hide_email(email) == event.data.get("email"): _LOGGER.debug("%s: Received relogin request: %s", hide_email(email), event) login_obj: AlexaLogin = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get( "login_obj" ) uuid = (await calculate_uuid(hass, email, url))["uuid"] if login_obj is None: login_obj = AlexaLogin( url=url, email=email, password=password, outputpath=hass.config.path, debug=account.get(CONF_DEBUG), otp_secret=account.get(CONF_OTPSECRET, ""), oauth=account.get(CONF_OAUTH, {}), uuid=uuid, oauth_login=True, ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["login_obj"] = login_obj else: login_obj.oauth_login = True await login_obj.reset() # await login_obj.login() if await test_login_status(hass, config_entry, login_obj): await setup_alexa(hass, config_entry, login_obj) async def login_success(event=None) -> None: """Relogin to Alexa.""" if hide_email(email) == event.data.get("email"): _LOGGER.debug("Received Login success: %s", event) login_obj: AlexaLogin = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get( "login_obj" ) await setup_alexa(hass, config_entry, login_obj) hass.data.setdefault(DATA_ALEXAMEDIA, {}) hass.data[DATA_ALEXAMEDIA].setdefault("accounts", {}) hass.data[DATA_ALEXAMEDIA].setdefault("config_flows", {}) hass.data[DATA_ALEXAMEDIA].setdefault("notify_service", None) account = config_entry.data email = account.get(CONF_EMAIL) password = account.get(CONF_PASSWORD) url = account.get(CONF_URL) hass.data[DATA_ALEXAMEDIA]["accounts"].setdefault( email, { "coordinator": None, "config_entry": config_entry, "setup_alexa": setup_alexa, "devices": { "media_player": {}, "switch": {}, "guard": [], "light": [], "binary_sensor": [], "temperature": [], "smart_switch": [], }, "entities": { "media_player": {}, "switch": {}, "sensor": {}, "light": [], "binary_sensor": [], "alarm_control_panel": {}, "smart_switch": [], }, "excluded": {}, "new_devices": True, "http2_lastattempt": 0, "http2error": 0, "http2_commands": {}, "http2_activity": {"serials": {}, "refreshed": {}}, "http2": None, "auth_info": None, "second_account_index": 0, "should_get_network": True, "first_run": True, "notifications": {}, # already used for the raw notifications dict "notifications_pending": set(), # doppler serials that need a refresh "notifications_refresh_task": None, # running task or None "notifications_retry_count": 0, # simple backoff counter "options": { CONF_INCLUDE_DEVICES: config_entry.data.get(CONF_INCLUDE_DEVICES, ""), CONF_EXCLUDE_DEVICES: config_entry.data.get(CONF_EXCLUDE_DEVICES, ""), CONF_QUEUE_DELAY: config_entry.data.get( CONF_QUEUE_DELAY, DEFAULT_QUEUE_DELAY ), CONF_SCAN_INTERVAL: config_entry.data.get( CONF_SCAN_INTERVAL, DEFAULT_SCAN_INTERVAL ), CONF_PUBLIC_URL: config_entry.data.get( CONF_PUBLIC_URL, DEFAULT_PUBLIC_URL ), CONF_EXTENDED_ENTITY_DISCOVERY: config_entry.data.get( CONF_EXTENDED_ENTITY_DISCOVERY, DEFAULT_EXTENDED_ENTITY_DISCOVERY ), CONF_DEBUG: config_entry.data.get(CONF_DEBUG, False), }, DATA_LISTENER: [config_entry.add_update_listener(update_listener)], }, ) uuid_dict = await calculate_uuid(hass, email, url) uuid = uuid_dict["uuid"] hass.data[DATA_ALEXAMEDIA]["accounts"][email]["second_account_index"] = uuid_dict[ "index" ] login: AlexaLogin = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get( "login_obj", AlexaLogin( url=url, email=email, password=password, outputpath=hass.config.path, debug=account.get(CONF_DEBUG), otp_secret=account.get(CONF_OTPSECRET, ""), oauth=account.get(CONF_OAUTH, {}), uuid=uuid, oauth_login=True, ), ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["login_obj"] = login hass.data[DATA_ALEXAMEDIA]["accounts"][email]["last_push_activity"] = 0 # Create runtime_data for optimized architecture # This provides type-safe access to account data if not hasattr(config_entry, "runtime_data") or config_entry.runtime_data is None: config_entry.runtime_data = AlexaRuntimeData( login_obj=login, config_entry=config_entry, second_account_index=uuid_dict["index"], ) if not hass.data[DATA_ALEXAMEDIA]["accounts"][email]["second_account_index"]: hass.bus.async_listen_once(EVENT_HOMEASSISTANT_STOP, close_alexa_media) hass.bus.async_listen_once(EVENT_HOMEASSISTANT_STARTED, complete_startup) hass.bus.async_listen("alexa_media_relogin_required", relogin) hass.bus.async_listen("alexa_media_relogin_success", login_success) try: _t = time.monotonic() cookies = await login.load_cookie() cookie_login_ok = False if cookies: try: if login._session is None or getattr(login._session, "closed", False): login._create_session(True) async with login._session.get( "https://alexa.amazon.com/api/bootstrap", cookies=cookies, ssl=login._ssl, allow_redirects=False, ) as response: if response.status == 200: data = loads(await response.text()) auth = (data or {}).get("authentication") or {} customer_email = (auth.get("customerEmail") or "").lower() if ( auth.get("authenticated") and customer_email == email.lower() ): _LOGGER.debug( "[BOOT] Cookie auth confirmed via /api/bootstrap" ) login.status["login_successful"] = True login.customer_id = auth.get("customerId") login.stats["login_timestamp"] = datetime.now() login.stats["api_calls"] = 0 await login.check_domain() await login.finalize_login() cookie_login_ok = True except (JSONDecodeError, ValueError, aiohttp.ClientError) as ex: _LOGGER.debug("[BOOT] Bootstrap cookie auth check failed: %s", ex) if not cookie_login_ok: await login.login(cookies=cookies) _LOGGER.debug("[BOOT] login completed in %.2fs", time.monotonic() - _t) _t = time.monotonic() if await test_login_status(hass, config_entry, login): _LOGGER.debug("[BOOT] test_login_status in %.2fs", time.monotonic() - _t) _t = time.monotonic() await setup_alexa(hass, config_entry, login) _LOGGER.debug( "[BOOT] setup_entry total: %.2fs", time.monotonic() - _boot_start ) return True return False except AlexapyConnectionError as err: raise ConfigEntryNotReady(str(err) or "Connection Error during login") from err async def setup_alexa(hass, config_entry, login_obj: AlexaLogin): # pylint: disable=too-many-statements,too-many-locals """Set up a alexa api based on host parameter.""" debug = config_entry.data.get(CONF_DEBUG, False) # Record metrics metrics = get_metrics(hass) email = login_obj.email if metrics: metrics.record_boot_stage(f"setup_alexa_start_{hide_email(email)}") # Initialize throttling state and lock dnd_update_lock = asyncio.Lock() last_dnd_update_times: dict[str, datetime] = {} pending_dnd_updates: dict[str, bool] = {} scheduled_dnd_tasks: dict[str, asyncio.Task] = {} async def async_update_data() -> Optional[AlexaEntityData]: # noqa pylint: disable=too-many-branches """Fetch data from API endpoint. This is the place to pre-process the data to lookup tables so entities can quickly look up their data. This will ping Alexa API to identify all devices, bluetooth, and the last called device. If any guards, sensors, switches or lights are configured, their current state will be acquired. This data is returned directly so that it is available on the coordinator. This will add new devices and services when discovered. By default this runs every SCAN_INTERVAL seconds unless another method calls it. if push is connected, it will increase the delay 10-fold between updates. While throttled at MIN_TIME_BETWEEN_SCANS, care should be taken to reduce the number of runs to avoid flooding. Slow changing states should be checked here instead of in spawned components like media_player since this object is one per account. Each AlexaAPI call generally results in two webpage requests. """ email = config.get(CONF_EMAIL) accounts = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}) account = accounts.get(email) if not account: return None login_obj = account.get("login_obj") if not login_obj or not _network_allowed(login_obj): return None account = hass.data[DATA_ALEXAMEDIA]["accounts"][email] existing_serials = set(_existing_serials(hass, login_obj)) existing_serials |= _entity_backed_serials(account) existing_entities = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"][ "media_player" ].values() auth_info = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get("auth_info") new_devices = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["new_devices"] extended_entity_discovery = hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "options" ].get(CONF_EXTENDED_ENTITY_DISCOVERY) should_get_network = hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "should_get_network" ] first_run = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["first_run"] devices = {} bluetooth = {} preferences = {} dnd = {} entity_state = {} # Try to get cached data for faster boot entry_id = config_entry.entry_id if config_entry else "" cache_key_prefix = f"{email}_{entry_id}" cached_devices = None _used_cached_devices = False if metrics: cached_devices = metrics.api_cache.get(f"{cache_key_prefix}_devices") if cached_devices and not new_devices: _LOGGER.debug("%s: Using cached devices data", hide_email(email)) # NOTE: DataCache returns direct references. We intentionally enrich device dicts # in-place each refresh cycle (bluetooth_state/locale/dnd/etc.). devices = cached_devices _used_cached_devices = True # Still need fresh bluetooth, preferences, and DND tasks = [ AlexaAPI.get_bluetooth(login_obj), AlexaAPI.get_device_preferences(login_obj), AlexaAPI.get_dnd_state(login_obj), ] else: tasks = [ AlexaAPI.get_devices(login_obj), AlexaAPI.get_bluetooth(login_obj), AlexaAPI.get_device_preferences(login_obj), AlexaAPI.get_dnd_state(login_obj), ] if new_devices: tasks.append(AlexaAPI.get_authentication(login_obj)) entities_to_monitor = set() # Temperature sensors (stored as entities["sensor"][serial]["Temperature"] = sensor) for per_serial in hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"][ "sensor" ].values(): if not isinstance(per_serial, dict): continue temp = per_serial.get("Temperature") if temp and temp.enabled: entities_to_monitor.add(temp.alexa_entity_id) # Air Quality sensors: # entities["sensor"][serial]["Air_Quality"][unique_id] = sensor airq = per_serial.get("Air_Quality") if isinstance(airq, dict): for aq_sensor in airq.values(): if aq_sensor and aq_sensor.enabled: entities_to_monitor.add(aq_sensor.alexa_entity_id) elif airq and getattr(airq, "enabled", False): # Backwards compat if some installs still have a single sensor stored entities_to_monitor.add(airq.alexa_entity_id) for light in hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"].get( "light", [] ): if light.enabled: entities_to_monitor.add(light.alexa_entity_id) for binary_sensor in hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "entities" ].get("binary_sensor", []): if binary_sensor.enabled: entities_to_monitor.add(binary_sensor.alexa_entity_id) for guard in ( hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"] .get("alarm_control_panel", {}) .values() ): if guard.enabled: entities_to_monitor.add(guard.unique_id) for smart_switch in hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "entities" ].get("smart_switch", []): if smart_switch.enabled: entities_to_monitor.add(smart_switch.alexa_entity_id) if entities_to_monitor: tasks.append(get_entity_data(login_obj, list(entities_to_monitor))) if should_get_network: tasks.append(AlexaAPI.get_network_details(login_obj)) optional_task_results = [] try: if tasks: # Note: asyncio.TimeoutError and aiohttp.ClientError are already # handled by the data update coordinator. # Increase timeout from 30s to 45s to permit # get_network_details() retries which could up to 30s. async with async_timeout.timeout(45): start_fetch = time.monotonic() if _used_cached_devices: ( bluetooth, preferences, dnd, *optional_task_results, ) = await asyncio.gather(*tasks) else: ( devices, bluetooth, preferences, dnd, *optional_task_results, ) = await asyncio.gather(*tasks) fetch_time = time.monotonic() - start_fetch _LOGGER.debug( "[BOOT] API fetch (%d tasks, cached=%s) in %.2fs", len(tasks) if tasks else 0, _used_cached_devices, fetch_time, ) # Record API call metrics if metrics: metrics.record_api_call("initial_fetch", fetch_time) # Cache the devices for faster next boot (only freshly fetched) if not _used_cached_devices: metrics.api_cache.cache_set( f"{cache_key_prefix}_devices", devices ) _t_post = time.monotonic() if should_get_network and optional_task_results: # First run is a special case. Get the state of all entities(including disabled) # This ensures all entities have state during startup without needing to request coordinator refresh _LOGGER.info( "%s: Network Discovery: Checking", hide_email(email) ) api_devices = optional_task_results.pop() if not api_devices: _LOGGER.warning( "%s: Network Discovery: AlexaAPI returned an unexpected response. Retrying on next polling cycle", hide_email(email), ) else: _LOGGER.debug( "%s: Network Discovery: Success, processing response", hide_email(email), ) # Only process this once after success hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "should_get_network" ] = False # Discard the entities_to_monitor results since we now have full network details if entities_to_monitor and optional_task_results: optional_task_results.pop() entities_to_monitor.clear() alexa_entities = parse_alexa_entities( api_devices, debug=hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "options" ].get(CONF_DEBUG, False), ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["devices"].update( alexa_entities ) _entities_to_monitor = set() for type_of_entity, entities in alexa_entities.items(): if ( type_of_entity in {"guard", "temperature", "air_quality", "aiaqm"} or extended_entity_discovery ): for entity in entities: _entities_to_monitor.add(entity.get("id")) _LOGGER.debug("Monitoring: %s", entity.get("name")) _LOGGER.debug( "%s: Network Discovery: %s entities will be monitored", hide_email(email), len(list(_entities_to_monitor)), ) # Use shorter timeout for entity data to avoid blocking _t_ed = time.monotonic() try: entity_state = await asyncio.wait_for( get_entity_data(login_obj, list(_entities_to_monitor)), timeout=10.0, ) except asyncio.TimeoutError: _LOGGER.warning( "%s: get_entity_data timed out after 10s, " "entity states will be fetched on next cycle", hide_email(email), ) _LOGGER.debug( "[BOOT] get_entity_data (network) in %.2fs", time.monotonic() - _t_ed, ) if entities_to_monitor and optional_task_results: entity_state = optional_task_results.pop() _LOGGER.debug( "%s: Processing %s entities to monitor", hide_email(email), len(list(entities_to_monitor)), ) if new_devices and optional_task_results: auth_info = optional_task_results.pop() _LOGGER.debug( "%s: Found %s devices, %s bluetooth", hide_email(email), len(devices) if devices is not None else "", ( len(bluetooth.get("bluetoothStates", [])) if bluetooth is not None else "" ), ) # Process notifications in background to avoid blocking boot # (process_notifications has a 4s sleep + API call) _LOGGER.debug( "[BOOT] post-fetch processing in %.2fs", time.monotonic() - _t_post ) existing_notif_task = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get( "notifications_init_task" ) if existing_notif_task and not existing_notif_task.done(): _LOGGER.debug( "%s: Notifications background task already running, skipping", hide_email(email), ) else: async def _bg_process_notifications(): try: await process_notifications(login_obj) except ( AlexapyConnectionError, AlexapyLoginError, asyncio.TimeoutError, JSONDecodeError, ): _LOGGER.debug( "%s: Background notifications failed, retrying once", hide_email(email), ) try: await asyncio.sleep(5) await process_notifications(login_obj) except ( AlexapyConnectionError, AlexapyLoginError, asyncio.TimeoutError, JSONDecodeError, ): _LOGGER.debug( "%s: Background notifications retry failed", hide_email(email), ) hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "notifications_init_task" ] = hass.async_create_background_task( _bg_process_notifications(), f"{DOMAIN}_notifications_init", ) except (AlexapyLoginError, JSONDecodeError): _LOGGER.debug( "%s: Alexa API disconnected; attempting to relogin : status %s", hide_email(email), login_obj.status, ) if login_obj.status: hass.bus.async_fire( "alexa_media_relogin_required", event_data={"email": hide_email(email), "url": login_obj.url}, ) return None except asyncio.CancelledError: # Task cancelled during unload/shutdown; propagate cancellation. raise _t_proc = time.monotonic() new_alexa_clients = [] # list of newly discovered device names exclude_filter = [] include_filter = [] for device in devices: serial = device["serialNumber"] dev_name = device["accountName"] if include and dev_name not in include: include_filter.append(dev_name) if "appDeviceList" in device: for app in device["appDeviceList"]: ( hass.data[DATA_ALEXAMEDIA]["accounts"][email]["excluded"][ app["serialNumber"] ] ) = device hass.data[DATA_ALEXAMEDIA]["accounts"][email]["excluded"][ serial ] = device continue if exclude and dev_name in exclude: exclude_filter.append(dev_name) if "appDeviceList" in device: for app in device["appDeviceList"]: ( hass.data[DATA_ALEXAMEDIA]["accounts"][email]["excluded"][ app["serialNumber"] ] ) = device hass.data[DATA_ALEXAMEDIA]["accounts"][email]["excluded"][ serial ] = device continue if ( dev_name not in include_filter and device.get("capabilities") and not any( x in device["capabilities"] for x in ["MUSIC_SKILL", "TIMERS_AND_ALARMS", "REMINDERS"] ) ): # skip devices without music or notification skill _LOGGER.debug("Excluding %s for lacking capability", dev_name) continue if bluetooth is not None and "bluetoothStates" in bluetooth: for b_state in bluetooth["bluetoothStates"]: if serial == b_state["deviceSerialNumber"]: device["bluetooth_state"] = b_state break if preferences is not None and "devicePreferences" in preferences: for dev in preferences["devicePreferences"]: if dev["deviceSerialNumber"] == serial: device["locale"] = dev["locale"] device["timeZoneId"] = dev["timeZoneId"] _LOGGER.debug( "%s: Locale %s timezone %s", dev_name, device["locale"], device["timeZoneId"], ) break if dnd is not None and "doNotDisturbDeviceStatusList" in dnd: for dev in dnd["doNotDisturbDeviceStatusList"]: if dev["deviceSerialNumber"] == serial: device["dnd"] = dev["enabled"] _LOGGER.debug("%s: DND %s", dev_name, device["dnd"]) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["devices"][ "switch" ].setdefault(serial, {"dnd": True}) break hass.data[DATA_ALEXAMEDIA]["accounts"][email]["auth_info"] = device[ "auth_info" ] = auth_info hass.data[DATA_ALEXAMEDIA]["accounts"][email]["devices"]["media_player"][ serial ] = device if serial not in existing_serials: new_alexa_clients.append(dev_name) elif ( serial in existing_serials and hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"][ "media_player" ].get(serial) and hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"][ "media_player" ] .get(serial) .enabled ): await ( hass.data[DATA_ALEXAMEDIA]["accounts"][email]["entities"][ "media_player" ] .get(serial) .refresh(device, skip_api=True) ) _LOGGER.debug( "%s: Existing: %s New: %s;" " Filtered out by not being in include: %s " "or in exclude: %s", hide_email(email), list(existing_entities), new_alexa_clients, include_filter, exclude_filter, ) _LOGGER.debug("[BOOT] device processing in %.2fs", time.monotonic() - _t_proc) if new_alexa_clients: cleaned_config = config.copy() cleaned_config.pop(CONF_PASSWORD, None) # CONF_PASSWORD contains sensitive info which is no longer needed # Load multiple platforms in parallel using async_forward_entry_setups _LOGGER.debug("Loading platforms: %s", ", ".join(ALEXA_COMPONENTS)) try: _t = time.monotonic() await hass.config_entries.async_forward_entry_setups( config_entry, ALEXA_COMPONENTS ) _LOGGER.debug("[BOOT] platform loading in %.2fs", time.monotonic() - _t) if metrics: metrics.record_boot_stage(f"platforms_loaded_{hide_email(email)}") except (asyncio.TimeoutError, TimeoutException) as ex: _LOGGER.error(f"Error while loading platforms: {ex}") raise ConfigEntryNotReady( f"Timeout while loading platforms: {ex}" ) from ex hass.data[DATA_ALEXAMEDIA]["accounts"][email]["new_devices"] = False # prune stale devices device_registry = dr.async_get(hass) entity_backed_ids = _entity_backed_device_identifiers( hass.data[DATA_ALEXAMEDIA]["accounts"][email] ) for device_entry in dr.async_entries_for_config_entry( device_registry, config_entry.entry_id ): for _, identifier in device_entry.identifiers: if ( identifier in hass.data[DATA_ALEXAMEDIA]["accounts"][email]["devices"][ "media_player" ] or identifier in map( lambda x: slugify(f"{x}_{email}"), hass.data[DATA_ALEXAMEDIA]["accounts"][email]["devices"][ "media_player" ].keys(), ) or identifier in entity_backed_ids ): break else: device_registry.async_remove_device(device_entry.id) _LOGGER.debug( "%s: Removing stale device %s", hide_email(email), device_entry.name, ) await login_obj.save_cookiefile() if login_obj.access_token: hass.config_entries.async_update_entry( config_entry, data={ **config_entry.data, CONF_OAUTH: { "access_token": login_obj.access_token, "refresh_token": login_obj.refresh_token, "expires_in": login_obj.expires_in, "mac_dms": login_obj.mac_dms, "code_verifier": login_obj.code_verifier, "authorization_code": login_obj.authorization_code, }, }, ) if first_run or not _push_healthy(account): if _network_allowed(login_obj): trigger = account.get("last_called_probe_trigger") if callable(trigger): trigger("POLL_REFRESH", None) else: # fallback if probe not initialized for some reason hass.async_create_background_task( _async_update_last_called_global(hass, login_obj, email), f"{DOMAIN}_last_called_poll_{hide_email(email)}", ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["first_run"] = False return entity_state @_catch_login_errors async def process_notifications(login_obj, raw_notifications=None) -> bool: """Process raw notifications json. Returns True if notifications were updated, False if we skipped (e.g. due to cooldown or alexapy returned None). """ email: str = login_obj.email account_dict = hass.data[DATA_ALEXAMEDIA]["accounts"][email] if raw_notifications is None: now = time.time() last = account_dict.get("last_notif_poll", 0.0) delta = now - last if delta < NOTIFICATION_COOLDOWN: _LOGGER.debug( "%s: Skipping get_notifications; last poll %.1fs ago " "(cooldown %ss).", hide_email(email), delta, NOTIFICATION_COOLDOWN, ) return False account_dict["last_notif_poll"] = now # Small delay to let Alexa settle if we're polling explicitly await asyncio.sleep(4) raw_notifications = await AlexaAPI.get_notifications(login_obj) previous = account_dict.get("notifications", {}) notifications = {"process_timestamp": dt.utcnow()} if raw_notifications is not None: for notification in raw_notifications: n_dev_id = notification.get("deviceSerialNumber") if n_dev_id is None: # skip notifications untied to a device for now # https://github.com/alandtse/alexa_media_player/issues/633#issuecomment-610705651 continue n_type = notification.get("type") if n_type is None: continue if n_type == "MusicAlarm": n_type = "Alarm" n_id = notification["notificationIndex"] if n_type == "Alarm": n_date = notification.get("originalDate") n_time = notification.get("originalTime") notification["date_time"] = ( f"{n_date} {n_time}" if n_date and n_time else None ) previous_alarm = safe_get(previous, [n_dev_id, "Alarm", n_id], {}) if previous_alarm and alarm_just_dismissed( notification, previous_alarm.get("status"), previous_alarm.get("version"), ): hass.bus.async_fire( "alexa_media_alarm_dismissal_event", event_data={ "device": {"id": n_dev_id}, "event": notification, }, ) if n_dev_id not in notifications: notifications[n_dev_id] = {} if n_type not in notifications[n_dev_id]: notifications[n_dev_id][n_type] = {} notifications[n_dev_id][n_type][n_id] = notification account_dict["notifications"] = notifications _LOGGER.debug( "%s: Updated %s notifications for %s devices at %s", hide_email(email), len(raw_notifications) if raw_notifications is not None else 0, len(notifications), dt.as_local(account_dict["notifications"]["process_timestamp"]), ) # Notify sensors that the notifications snapshot has been refreshed async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"notifications_refreshed": True}, ) return True # --------------------------------------------------------------------- # Near-real-time last_called probing (customer history), single worker # Initialize ONCE per account (setup_alexa), not inside http2_handler. # --------------------------------------------------------------------- def _init_last_called_probe_worker(account: dict) -> None: """Initialize per-account last_called probe worker trigger function.""" account.setdefault("last_called_api_lock", asyncio.Lock()) account.setdefault( "last_called_customer_history_ts", 0 ) # ms epoch last applied account.setdefault("last_called_probe_backoff_s", 0.0) account.setdefault("last_called_probe_event", asyncio.Event()) account.setdefault("last_called_probe_next_allowed", 0.0) # monotonic seconds account.setdefault("last_called_probe_task", None) account.setdefault("last_called_probe_trigger_cmd", "") account.setdefault("last_called_probe_trigger_serial", None) account.setdefault("last_called_probe_trigger_ts", 0) # newest push ts (ms) account.setdefault("last_called_probe_trigger", None) account.setdefault("last_called_activity_queue", []) account.setdefault("last_called_last_pushed_activity", {}) account.setdefault("last_volumes", {}) account.setdefault("last_equalizer", {}) if callable(account.get("last_called_probe_trigger")): return async def _last_called_probe_worker() -> None: """Single worker per account: debounce bursts, then quick-retry until history >= trigger_ts.""" skip_debounce = False try: while True: accounts = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}) account_live = accounts.get(email) if not account_live: return def _debug(msg, *args): if debug: _LOGGER.debug("%s: " + msg, hide_email(email), *args) login_live = account_live.get("login_obj") if not login_live or not _network_allowed(login_live): await asyncio.sleep(5) continue # Stop this worker if the login object has been torn down (entry # reload/unload) so it cannot orphan onto a dead login and keep # hammering the API across reloads. if getattr(login_live, "close_requested", False) or ( getattr(login_live, "session", None) is not None and login_live.session.closed ): return # Skip probing while the login is unhealthy: the request would be # built with a missing csrf / None header and raise a serialization # TypeError, and hammering the history API while logged-out can # hinder re-auth. Poll quietly until the login recovers. _status = getattr(login_live, "status", None) if not (_status and _status.get("login_successful")): await asyncio.sleep(15) continue # Bounded wait so that if the login is torn down while we are # idle here, the next iteration's teardown guard exits promptly # instead of parking on the event forever. try: await asyncio.wait_for( account_live["last_called_probe_event"].wait(), timeout=15, ) except asyncio.TimeoutError: continue account_live["last_called_probe_event"].clear() if not skip_debounce: await asyncio.sleep( LAST_CALLED_DEBOUNCE_S + random.uniform(0.0, 0.05) # nosec B311 # noqa: S311 ) if account_live["last_called_probe_event"].is_set(): continue skip_debounce = False preempted = False while True: now = time.monotonic() next_allowed = float( account_live.get("last_called_probe_next_allowed", 0.0) or 0.0 ) delay = max(0.0, next_allowed - now) if delay <= 0: break try: await asyncio.wait_for( account_live["last_called_probe_event"].wait(), timeout=delay, ) account_live["last_called_probe_event"].clear() preempted = True break except asyncio.TimeoutError: break if preempted: continue trigger_cmd = str( account_live.get("last_called_probe_trigger_cmd") or "push" ) for attempt in range(LAST_CALLED_RETRY_LIMIT + 1): if account_live["last_called_probe_event"].is_set(): break try: async with account_live["last_called_api_lock"]: queue_snapshot = _snapshot_last_called_activity_queue( account_live ) if not queue_snapshot: if trigger_cmd in ( "GLOBAL_REFRESH", "SERVICE_REFRESH", "POLL_REFRESH", ): last_called = ( await AlexaAPI.get_last_device_serial( login_live, items=LAST_CALLED_ITEMS ) ) if isinstance( last_called, dict ) and _valid_voice_summary( last_called.get("summary") ): await update_last_called( login_live, last_called ) account_live["last_called_probe_trigger_ts"] = 0 break earliest_ts = min( ( int(item.get("activity_ts") or 0) for item in queue_snapshot if item.get("activity_ts") ), default=0, ) rvh_window_ms = max( LAST_CALLED_LOOKBACK_MS, 15 * 60 * 1000 ) start_time = ( max(0, earliest_ts - rvh_window_ms) if earliest_ts else int((time.time() * 1000) - rvh_window_ms) ) end_time = int(time.time() * 1000) + rvh_window_ms max_record_size = max( LAST_CALLED_ITEMS, len(queue_snapshot) + 2 ) try: records = await AlexaAPI.get_customer_history_records( login_live, start_time=start_time, end_time=end_time, max_record_size=max_record_size, ) except TypeError as exc: # A serialization TypeError (e.g. a None header key # when the login is half-torn-down) is not transient # within a session: back off and do NOT re-arm. The # old behavior crash-looped ~every 11s, hammering the # API. Wait for the next genuine push trigger instead. backoff = max(LAST_CALLED_CONN_BACKOFF_S, 60.0) account_live["last_called_probe_next_allowed"] = ( time.monotonic() + backoff ) _LOGGER.warning( "%s: last_called probe API TypeError (%s): %s; " "backing off, not re-arming", hide_email(email), trigger_cmd, exc, ) break if records is None: _LOGGER.warning( "%s: last_called probe API returned None (%s)", hide_email(email), trigger_cmd, ) records = [] # 🔎 DEBUG: inspect raw history result _LOGGER.debug( "%s: last_called probe retrieved %s history records", hide_email(email), len(records or []), ) if isinstance(records, list): _debug( "last_called probe raw records (%s): %s", trigger_cmd, [ { "summary": ( item.get("description") or {} ).get("summary"), "response": item.get("alexaResponse"), "serial": item.get("deviceSerialNumber"), "ts": item.get("creationTimestamp"), "utteranceType": item.get("utteranceType"), } for item in records[:10] if isinstance(item, dict) ], ) else: _debug( "last_called probe returned unexpected type (%s): %r", trigger_cmd, records, ) except asyncio.CancelledError: raise except AlexapyTooManyRequestsError: uk_floor = random.uniform( # noqa: S311 30.0, 63.0 ) # nosec B311 prev = float( account_live.get("last_called_probe_backoff_s", 0.0) or 0.0 ) backoff = ( LAST_CALLED_429_BACKOFF_INITIAL_S if prev <= 0.0 else min(prev * 2.0, LAST_CALLED_429_BACKOFF_MAX_S) ) backoff = max(backoff, uk_floor) jitter = random.uniform( # noqa: S311 0.0, min(5.0, backoff * 0.1) ) # nosec B311 account_live["last_called_probe_backoff_s"] = backoff account_live["last_called_probe_next_allowed"] = ( time.monotonic() + backoff + jitter ) _LOGGER.debug( "%s: last_called probe rate-limited (%s); backing off %.1fs (%.1fs jitter) then self-retry", hide_email(email), trigger_cmd, backoff, jitter, ) skip_debounce = True account_live["last_called_probe_event"].set() break except AlexapyLoginError: account_live["last_called_probe_next_allowed"] = ( time.monotonic() + LAST_CALLED_LOGIN_BACKOFF_S ) _LOGGER.debug( "%s: last_called probe login error (%s); skipping", hide_email(email), trigger_cmd, ) report_relogin_required(hass, login_live, email) break except AlexapyConnectionError as exc: account_live["last_called_probe_next_allowed"] = ( time.monotonic() + LAST_CALLED_CONN_BACKOFF_S ) _LOGGER.debug( "%s: last_called probe connection error (%s): %s", hide_email(email), trigger_cmd, exc, ) break existing_serials_local = set( _existing_serials(hass, login_live) ) existing_serials_local |= _entity_backed_serials(account_live) payload, resolved_keys = ( _select_last_called_payload_from_records( records, queue_snapshot, account_live, existing_serials_local, ) ) _debug( "last_called queue match (%s): queued=%s matched=%s resolved=%s", trigger_cmd, [ { "serial": item.get("serial"), "customer_id": item.get("customer_id"), "activity_ts": item.get("activity_ts"), "command": item.get("command"), } for item in queue_snapshot ], ( { "serialNumber": payload.get("serialNumber"), "timestamp": payload.get("timestamp"), "summary": payload.get("summary"), "response": payload.get("response"), } if payload else None ), sorted(resolved_keys), ) if not payload: if attempt < LAST_CALLED_RETRY_LIMIT: await asyncio.sleep(LAST_CALLED_RETRY_DELAY_S) continue if trigger_cmd in ( "GLOBAL_REFRESH", "SERVICE_REFRESH", "POLL_REFRESH", ): _LOGGER.debug( "%s: queued activity unresolved after retries; falling back to direct refresh", hide_email(email), ) unresolved_keys = { (item.get("serial"), item.get("customer_id")) for item in queue_snapshot if item.get("serial") } _remove_last_called_activity_queue_entries( account_live, unresolved_keys ) try: async with account_live["last_called_api_lock"]: last_called = ( await AlexaAPI.get_last_device_serial( login_live, items=LAST_CALLED_ITEMS, ) ) if isinstance( last_called, dict ) and _valid_voice_summary( last_called.get("summary") ): await update_last_called( login_live, last_called ) except asyncio.CancelledError: raise except AlexapyLoginError as exc: _LOGGER.debug( "%s: fallback last_called refresh failed (%s): %s", hide_email(email), trigger_cmd, exc, ) report_relogin_required(hass, login_live, email) except ( AlexapyTooManyRequestsError, AlexapyConnectionError, ) as exc: _LOGGER.debug( "%s: fallback last_called refresh failed (%s): %s", hide_email(email), trigger_cmd, exc, ) account_live["last_called_probe_trigger_ts"] = 0 account_live["last_called_probe_event"].clear() break account_live["last_called_probe_backoff_s"] = 0.0 account_live["last_called_probe_next_allowed"] = ( time.monotonic() + LAST_CALLED_SUCCESS_PACE_S + random.uniform(0.0, 0.25) # nosec B311 # noqa: S311 ) trigger_serial = account_live.get( "last_called_probe_trigger_serial" ) _LOGGER.debug( "%s: Updating last_called via %s (triggered by %s): %s", hide_email(email), trigger_cmd, ( hide_serial(trigger_serial) if trigger_serial else "unknown" ), hide_serial(payload["serialNumber"]), ) await update_last_called(login_live, payload) account_live["last_called_last_pushed_activity"][ payload["serialNumber"] ] = payload["timestamp"] _remove_last_called_activity_queue_entries( account_live, resolved_keys ) account_live["last_called_probe_trigger_ts"] = 0 account_live["last_called_probe_event"].clear() break except asyncio.CancelledError: raise except Exception: # pylint: disable=broad-except _LOGGER.exception( "%s: last_called probe worker crashed", hide_email(email), ) def _trigger_last_called_probe( trigger_command: str, trigger_ts_ms: int | None ) -> None: """Record newest trigger + wake worker. Does NOT cancel running worker.""" accounts = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}) account_live = accounts.get(email) if not account_live: return if trigger_ts_ms is not None: try: ts = int(trigger_ts_ms) except (TypeError, ValueError): ts = 0 prev = int(account_live.get("last_called_probe_trigger_ts") or 0) if ts > prev: account_live["last_called_probe_trigger_ts"] = ts account_live["last_called_probe_trigger_cmd"] = trigger_command else: # Manual refresh triggers clear any push watermark if trigger_command in ( "GLOBAL_REFRESH", "SERVICE_REFRESH", "POLL_REFRESH", ): account_live["last_called_probe_trigger_ts"] = 0 account_live["last_called_probe_trigger_serial"] = None account_live["last_called_probe_trigger_cmd"] = trigger_command task = account_live.get("last_called_probe_task") if task is None or task.done(): account_live["last_called_probe_task"] = ( hass.async_create_background_task( _last_called_probe_worker(), name=f"{DOMAIN}_last_called_probe_{hide_email(email)}", ) ) account_live["last_called_probe_event"].set() # Store the trigger on the live account as well (so reload swaps don't strand it) account["last_called_probe_trigger"] = _trigger_last_called_probe def _is_dnd_voice_toggle(last_called: dict) -> bool: summary = " ".join(((last_called.get("summary") or "").strip().lower()).split()) response = " ".join( ((last_called.get("response") or "").strip().lower()).split() ) # Normalize apostrophes (replace smart quotes with ASCII) summary = summary.replace("\u2019", "'").replace("\u2018", "'") response = response.replace("\u2019", "'").replace("\u2018", "'") return ( "do not disturb" in summary or "won't disturb you" in response or "do not disturb is now off" in response ) @_catch_login_errors async def update_last_called(login_obj, last_called=None, force=False): """Update the last called device for the login_obj. Stores the last_called in hass.data and fires an event to notify listeners. Delegates storage/dispatch to the module-level _store_and_dispatch_last_called helper. """ if not isinstance(last_called, dict) or not last_called.get("summary"): try: # Serialize calls per account to avoid parallel rate-limited requests. account = ( hass.data.get(DATA_ALEXAMEDIA, {}) .get("accounts", {}) .get(email, {}) ) api_lock = account.get("last_called_api_lock") if account else None if api_lock is None: last_called = await AlexaAPI.get_last_device_serial(login_obj) else: async with api_lock: last_called = await AlexaAPI.get_last_device_serial(login_obj) except asyncio.CancelledError: raise except AlexapyTooManyRequestsError: _LOGGER.debug( "%s: Rate limited during last_called update; skipping", hide_email(email), ) return except AlexapyLoginError: _LOGGER.debug( "%s: Login error during last_called update", hide_email(email), ) report_relogin_required(hass, login_obj, email) return except AlexapyConnectionError as exc: _LOGGER.debug( "%s: Connection error during last_called update: %s", hide_email(email), exc, ) return except TypeError: _LOGGER.debug( "%s: Error updating last_called: %s", hide_email(email), repr(last_called), ) return if not isinstance(last_called, dict): _LOGGER.debug( "%s: Error updating last_called: unexpected response %s", hide_email(email), repr(last_called), ) return # Central voice-only gate if not _valid_voice_summary(last_called.get("summary")): _LOGGER.debug( "%s: Ignoring last_called with invalid/non-voice summary: %s", hide_email(email), repr(last_called.get("summary")), ) return _LOGGER.debug( "%s: Updated last_called: %s", hide_email(email), hide_serial(last_called) ) _store_and_dispatch_last_called(hass, email, last_called, force) if _is_dnd_voice_toggle(last_called): _LOGGER.debug( "%s: last_called indicates DND voice toggle", hide_email(email) ) await update_dnd_state(login_obj) @_catch_login_errors async def update_bluetooth_state(login_obj, device_serial): """Update the bluetooth state on ws bluetooth event.""" bluetooth = await AlexaAPI.get_bluetooth(login_obj) device = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["devices"][ "media_player" ][device_serial] if bluetooth is not None and "bluetoothStates" in bluetooth: for b_state in bluetooth["bluetoothStates"]: if device_serial == b_state["deviceSerialNumber"]: _LOGGER.debug( "%s: setting value for: %s to %s", hide_email(email), hide_serial(device_serial), hide_serial(b_state), ) device["bluetooth_state"] = b_state return device["bluetooth_state"] _LOGGER.debug( "%s: get_bluetooth for: %s failed with %s", hide_email(email), hide_serial(device_serial), hide_serial(bluetooth), ) return None async def schedule_update_dnd_state(email: str) -> None: """Run one deferred DND refresh after the cooldown expires.""" try: while True: async with dnd_update_lock: if not pending_dnd_updates.get(email, False): scheduled_dnd_tasks.pop(email, None) return last_run = last_dnd_update_times.get(email) now = datetime.utcnow() remaining = 0.0 if last_run is not None: elapsed = now - last_run if elapsed < MIN_TIME_BETWEEN_SCANS: remaining = ( MIN_TIME_BETWEEN_SCANS - elapsed ).total_seconds() if remaining > 0: _LOGGER.debug( "%s: Deferred DND update sleeping %.3fs until cooldown expires", hide_email(email), remaining, ) await asyncio.sleep(remaining) async with dnd_update_lock: if not pending_dnd_updates.get(email, False): scheduled_dnd_tasks.pop(email, None) return last_run = last_dnd_update_times.get(email) now = datetime.utcnow() if ( last_run is not None and (now - last_run) < MIN_TIME_BETWEEN_SCANS ): # Another update snuck in or timing was slightly early; loop and re-evaluate. continue pending_dnd_updates[email] = False scheduled_dnd_tasks.pop(email, None) login_obj = ( hass.data.get(DATA_ALEXAMEDIA, {}) .get("accounts", {}) .get(email, {}) .get("login_obj") ) if not login_obj: _LOGGER.debug( "%s: Skipping scheduled forced DND update: login_obj missing", hide_email(email), ) return _LOGGER.debug( "%s: Executing scheduled forced DND update", hide_email(email), ) await update_dnd_state(login_obj) return except asyncio.CancelledError: _LOGGER.debug("%s: Deferred DND update task cancelled", hide_email(email)) raise finally: async with dnd_update_lock: task = scheduled_dnd_tasks.get(email) if task is asyncio.current_task(): scheduled_dnd_tasks.pop(email, None) @_catch_login_errors async def update_dnd_state(login_obj) -> None: """Update the DND state on websocket DND combo event.""" email = login_obj.email now = datetime.utcnow() async with dnd_update_lock: last_run = last_dnd_update_times.get(email) if last_run is not None and (now - last_run) < MIN_TIME_BETWEEN_SCANS: pending_dnd_updates[email] = True if ( email not in scheduled_dnd_tasks or scheduled_dnd_tasks[email].done() ): _LOGGER.debug( "%s: Throttling active; scheduling deferred DND update.", hide_email(email), ) scheduled_dnd_tasks[email] = asyncio.create_task( schedule_update_dnd_state(email) ) else: _LOGGER.debug( "%s: Throttling active; deferred DND update already scheduled.", hide_email(email), ) return last_dnd_update_times[email] = now _LOGGER.debug("%s: Updating DND state", hide_email(email)) try: dnd = await AlexaAPI.get_dnd_state(login_obj) except asyncio.TimeoutError: _LOGGER.error( "%s: Timeout occurred while fetching DND state", hide_email(email), ) return except Exception as err: # pylint: disable=broad-except _LOGGER.error( "%s: Unexpected error while fetching DND state: %s", hide_email(email), err, ) return if dnd is not None and "doNotDisturbDeviceStatusList" in dnd: async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"dnd_update": dnd["doNotDisturbDeviceStatusList"]}, ) return _LOGGER.debug("%s: get_dnd_state failed: dnd:%s", hide_email(email), dnd) def _schedule_notifications_refresh( hass, email: str, device_serial: str | None = None, reason: str = "", ) -> None: """Mark notifications as needing refresh and ensure worker task is running. device_serial is just for debug; we track a set of pending devices but we always refresh the full notifications payload once. """ account = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}).get(email) if not account: return if device_serial: account["notifications_pending"].add(device_serial) else: # Special marker for "global" changes if you want one account["notifications_pending"].add("*") if reason: _LOGGER.debug( "%s: Scheduling notifications refresh (reason=%s, pending=%s)", hide_email(email), reason, account["notifications_pending"], ) task = account.get("notifications_refresh_task") if task is not None and not task.done(): # Already have a running worker; it'll see the new pending set return # Start new worker account["notifications_refresh_task"] = hass.async_create_task( _run_notifications_refresh(hass, email) ) async def _run_notifications_refresh(hass, email: str) -> None: """Worker task: refresh notifications for an account if pending. - Uses alexapy.AlexaAPI.get_notifications(login) - Retries a few times if we only get None (cooldown/throttle) - Clears notifications_pending when successful or when we give up """ account = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}).get(email) if not account: return login = account.get("login_obj") if not login: return try: retries = 0 while ( account.get("notifications_pending", set()) and retries <= NOTIFY_REFRESH_MAX_RETRIES ): try: data = await AlexaAPI.get_notifications(login) except Exception as ex: _LOGGER.warning( "%s: get_notifications raised %s; treating as None. This may indicate an unexpected error.", hide_email(email), ex, ) data = None if data is not None: # Success: update through the normal processing path await process_notifications(login, raw_notifications=data) account["notifications_retry_count"] = 0 account["notifications_pending"].clear() _LOGGER.debug( "%s: Refreshed notifications snapshot (pending cleared)", hide_email(email), ) return # If we get here, alexapy side returned None (cooldown / throttle) retries += 1 account["notifications_retry_count"] = retries if not account["notifications_pending"]: # Nothing to do anymore, bail early break _LOGGER.debug( "%s: Notifications refresh returned None (retry %s/%s); " "pending=%s; sleeping %.1fs", hide_email(email), retries, NOTIFY_REFRESH_MAX_RETRIES, account["notifications_pending"], NOTIFY_REFRESH_BACKOFF, ) await asyncio.sleep(NOTIFY_REFRESH_BACKOFF) # If we fall through, give up for now but leave pending set alone if account["notifications_pending"]: _LOGGER.debug( "%s: Giving up notifications refresh after %s attempts; " "still pending=%s", hide_email(email), retries, account["notifications_pending"], ) finally: # Always clear the task pointer so future pushes can schedule again account["notifications_refresh_task"] = None async def http2_connect() -> HTTP2EchoClient: """Open HTTP2 Push connection. This will only attempt one login before failing. """ http2: Optional[HTTP2EchoClient] = None email = login_obj.email try: if login_obj.session.closed: _LOGGER.debug( "%s: HTTP2 creation aborted. Session is closed.", hide_email(email), ) return http2 = HTTP2EchoClient( login_obj, msg_callback=http2_handler, open_callback=http2_open_handler, close_callback=http2_close_handler, error_callback=http2_error_handler, loop=hass.loop, ) _LOGGER.debug("%s: Starting HTTP2: %s", hide_email(email), http2) await http2.async_run() except AlexapyLoginError as exception_: _LOGGER.debug( "%s: Login Error detected from http2: %s", hide_email(email), exception_, ) hass.bus.async_fire( "alexa_media_relogin_required", event_data={"email": hide_email(email), "url": login_obj.url}, ) return except BaseException as exception_: # pylint: disable=broad-except _LOGGER.debug( "%s: HTTP2 creation failed: %s", hide_email(email), exception_ ) return _LOGGER.debug("%s: HTTP2 created: %s", hide_email(email), http2) return http2 @callback async def http2_handler(message_obj): # pylint: disable=too-many-branches,too-many-statements """Handle http2 push messages. This allows push notifications from Alexa to update last_called and media state. """ coordinator = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get("coordinator") account = hass.data[DATA_ALEXAMEDIA]["accounts"][email] def _now_ms() -> int: return int(time.time() * 1000) def simulate_activity( device_serial: str, customer_id: str | None, trigger_command: str, trigger_ts_ms: int | None, ) -> None: _queue_last_called_activity( account, device_serial=device_serial, customer_id=customer_id, activity_ts=trigger_ts_ms, command=trigger_command, ) account["last_called_probe_trigger_serial"] = device_serial trigger = account.get("last_called_probe_trigger") if callable(trigger): trigger(trigger_command, trigger_ts_ms) def _handle_volume_change_activity( serial: str, json_payload: dict, trigger_ts_ms: int | None, ) -> None: """Mirror alexa-remote.ts PUSH_VOLUME_CHANGE -> simulateActivity() conditions.""" last_volumes: dict = account["last_volumes"] last_equalizer: dict = account["last_equalizer"] vol = json_payload.get("volumeSetting") muted = json_payload.get("isMuted") prev = last_volumes.get(serial) eq_prev = last_equalizer.get(serial) should_simulate = ( eq_prev is not None and abs(_now_ms() - int(eq_prev.get("updated", 0))) < LAST_CALLED_COALESCE_WINDOW_MS ) if should_simulate: _LOGGER.debug( "[_handle_volume_change_activity] Simulating activity", ) simulate_activity( serial, json_payload.get("destinationUserId"), "PUSH_VOLUME_CHANGE", trigger_ts_ms, ) else: _LOGGER.debug( "[_handle_volume_change_activity] Not simulating activity", ) last_volumes[serial] = { "volumeSetting": vol, "isMuted": muted, "updated": _now_ms(), } def _handle_equalizer_change_activity( serial: str, json_payload: dict, trigger_ts_ms: int | None, ) -> None: """Mirror alexa-remote.ts PUSH_EQUALIZER_STATE_CHANGE -> simulateActivity() conditions.""" last_volumes: dict = account["last_volumes"] last_equalizer: dict = account["last_equalizer"] bass = json_payload.get("bass") treble = json_payload.get("treble") midrange = json_payload.get("midrange") prev = last_equalizer.get(serial) vol_prev = last_volumes.get(serial) should_simulate = ( prev is not None and prev.get("bass") == bass and prev.get("treble") == treble and prev.get("midrange") == midrange ) or ( vol_prev is not None and abs(_now_ms() - int(vol_prev.get("updated", 0))) < LAST_CALLED_COALESCE_WINDOW_MS ) if should_simulate: simulate_activity( serial, json_payload.get("destinationUserId"), "PUSH_EQUALIZER_STATE_CHANGE", trigger_ts_ms, ) last_equalizer[serial] = { "bass": bass, "treble": treble, "midrange": midrange, "updated": _now_ms(), } # --------------------------------------------------------------------- # Main http2push parsing / dispatch # --------------------------------------------------------------------- updates = ( message_obj.get("directive", {}) .get("payload", {}) .get("renderingUpdates", []) ) existing_serials = set(_existing_serials(hass, login_obj)) existing_serials |= _entity_backed_serials(account) for item in updates: try: resource = loads(item.get("resourceMetadata", "")) except JSONDecodeError: continue command = ( resource["command"] if isinstance(resource, dict) and "command" in resource else None ) try: json_payload = ( loads(resource["payload"]) if isinstance(resource, dict) and "payload" in resource else None ) except (JSONDecodeError, TypeError): _LOGGER.debug( "%s: Skipping malformed push payload for command %s", hide_email(email), command, ) continue seen_commands = hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "http2_commands" ] if command and json_payload: _LOGGER.debug( "%s: Received http2push command: %s : %s", hide_email(email), command, hide_serial(json_payload), ) account["last_push_activity"] = time.time() serial = None command_time = time.time() if command not in seen_commands: _LOGGER.debug( "Adding %s to seen_commands: %s", command, seen_commands ) seen_commands[command] = command_time if ( "dopplerId" in json_payload and "deviceSerialNumber" in json_payload["dopplerId"] ): serial = json_payload["dopplerId"]["deviceSerialNumber"] elif ( "key" in json_payload and "entryId" in json_payload["key"] and json_payload["key"]["entryId"].find("#") != -1 ): serial = (json_payload["key"]["entryId"]).split("#")[2] json_payload["key"]["serialNumber"] = serial else: serial = None if command in ( "PUSH_AUDIO_PLAYER_STATE", "PUSH_MEDIA_CHANGE", "PUSH_MEDIA_PROGRESS_CHANGE", "NotifyMediaSessionsUpdated", "NotifyNowPlayingUpdated", ): if serial and serial in existing_serials: _LOGGER.debug( "Updating media_player: %s", hide_serial(json_payload) ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"player_state": json_payload}, ) elif command == "NotifyNowPlayingUpdated": _LOGGER.debug("Send NowPlaying: %s", hide_serial(json_payload)) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"now_playing": json_payload}, ) elif command == "PUSH_VOLUME_CHANGE": try: ts = resource.get("timeStamp") ts_ms = int(ts) if ts else None if serial and isinstance(json_payload, dict): _handle_volume_change_activity(serial, json_payload, ts_ms) except Exception: _LOGGER.exception( "%s: http2_handler failed processing %s", hide_email(email), command, ) if serial and serial in existing_serials: _LOGGER.debug( "Updating media_player volume: %s", hide_serial(json_payload), ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"player_state": json_payload}, ) elif command == "PUSH_DOPPLER_CONNECTION_CHANGE": # Player availability update if serial and serial in existing_serials: _LOGGER.debug( "Updating media_player availability %s", hide_serial(json_payload), ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"player_state": json_payload}, ) elif command == "PUSH_EQUALIZER_STATE_CHANGE": # Player equalizer update try: ts = resource.get("timeStamp") ts_ms = int(ts) if ts else None if serial and isinstance(json_payload, dict): _handle_equalizer_change_activity( serial, json_payload, ts_ms ) except Exception: _LOGGER.exception( "%s: http2_handler failed processing %s", hide_email(email), command, ) if serial and serial in existing_serials: _LOGGER.debug( "Updating media_player equalizer state %s", hide_serial(json_payload), ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"player_state": json_payload}, ) elif command == "PUSH_BLUETOOTH_STATE_CHANGE": # Player bluetooth update bt_event = ( json_payload.get("bluetoothEvent") if isinstance(json_payload, dict) else None ) _LOGGER.debug("bt_event: %s", bt_event) bt_success = ( json_payload.get("bluetoothEventSuccess") if isinstance(json_payload, dict) else None ) _LOGGER.debug("bt_success: %s", bt_success) if ( serial and serial in existing_serials and bt_success and bt_event and bt_event in { "DEVICE_CONNECTED", "DEVICE_DISCONNECTED", "STREAMING_STATE_CHANGED", } ): _LOGGER.debug( "Updating media_player bluetooth %s", hide_serial(json_payload), ) bluetooth_state = await update_bluetooth_state( login_obj, serial ) _LOGGER.debug( "bluetooth_state %s", hide_serial(bluetooth_state) ) if bluetooth_state: async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"bluetooth_change": bluetooth_state}, ) elif command == "PUSH_MEDIA_QUEUE_CHANGE": # Player availability update if serial and serial in existing_serials: _LOGGER.debug( "Updating media_player queue %s", hide_serial(json_payload), ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"queue_state": json_payload}, ) elif command == "PUSH_NOTIFICATION_CHANGE": # Notification/alarm state changed on this device. # Queue a refresh with backoff to ride out alexa-side cooldowns. _schedule_notifications_refresh( hass, email, device_serial=serial, reason="PUSH_NOTIFICATION_CHANGE", ) if serial and serial in existing_serials: _LOGGER.debug( "Updating mediaplayer notifications: %s", hide_serial(json_payload), ) async_dispatcher_send( hass, f"{DOMAIN}_{hide_email(email)}"[0:32], {"notification_update": json_payload}, ) elif command in [ "PUSH_DELETE_DOPPLER_ACTIVITIES", # Delete Alexa history, ]: pass elif command in [ "PUSH_TODO_CHANGE", # Update To-Do List "PUSH_LIST_CHANGE", # Clear a shopping list https://github.com/alandtse/alexa_media_player/issues/1190 "PUSH_LIST_ITEM_CHANGE", # Update shopping list ]: # To-do _LOGGER.debug("%s currently not supported", command) pass elif command in [ "PUSH_CONTENT_FOCUS_CHANGE", # Likely prime related refocus "PUSH_DEVICE_SETUP_STATE_CHANGE", # Likely device changes mid setup ]: _LOGGER.debug("%s currently not supported", command) pass elif command in [ "PUSH_MEDIA_PREFERENCE_CHANGE", # Disliking or liking songs, https://github.com/alandtse/alexa_media_player/issues/1599 ]: _LOGGER.debug("%s currently not supported", command) pass elif command in [ "MATTER_SETUP_NOTIFICATION", # New command observed 2026-02-20 ]: _LOGGER.debug("%s: New command; currently not supported", command) pass else: _LOGGER.debug( "Unhandled command: %s with data %s. Please report at %s", command, hide_serial(json_payload), ISSUE_URL, ) # Preserve existing http2 activity tracking + new-device discovery if serial in existing_serials: history = hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "http2_activity" ]["serials"].get(serial) if history is None or ( history and command_time - history[len(history) - 1][1] > 2 ): history = [(command, command_time)] else: history.append([command, command_time]) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2_activity"][ "serials" ][serial] = history events = [] for old_command, old_command_time in history: if ( old_command in {"PUSH_VOLUME_CHANGE", "PUSH_EQUALIZER_STATE_CHANGE"} and command_time - old_command_time < 0.25 ): events.append( (old_command, round(command_time - old_command_time, 2)) ) elif old_command in {"PUSH_AUDIO_PLAYER_STATE"}: events = [] if len(events) >= 4: _LOGGER.debug( "%s: Detected potential DND http2push change with %s events %s", hide_serial(serial), len(events), events, ) await update_dnd_state(login_obj) if ( serial and serial not in existing_serials and serial not in hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "excluded" ].keys() ): _LOGGER.debug("Discovered new media_player %s", hide_serial(serial)) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["new_devices"] = True if coordinator: await coordinator.async_request_refresh() @callback async def http2_open_handler(): """Handle http2 open.""" email: str = login_obj.email _LOGGER.debug("%s: HTTP2push successfully connected", hide_email(email)) hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "http2error" ] = 0 # set errors to 0 hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2_lastattempt"] = time.time() @callback async def http2_close_handler(): """Handle http2 close. This should attempt to reconnect up to 5 times """ email: str = login_obj.email hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2"] = None if login_obj.close_requested: _LOGGER.debug( "%s: Close requested; will not reconnect http2", hide_email(email) ) return if not login_obj.status.get("login_successful"): _LOGGER.debug( "%s: Login error; will not reconnect http2", hide_email(email) ) return errors: int = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] delay: int = 5 * 2**errors last_attempt = hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "http2_lastattempt" ] now = time.time() if (now - last_attempt) < delay: return http2_client = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2"] http2_enabled = bool(http2_client) while errors < 5 and not (http2_enabled): _LOGGER.debug( "%s: HTTP2 push closed; reconnect #%i in %is", hide_email(email), errors, delay, ) hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "http2_lastattempt" ] = time.time() http2_client = await http2_connect() hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2"] = http2_client http2_enabled = bool(http2_client) if http2_enabled: break errors = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] = ( hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] + 1 ) delay = 5 * 2**errors errors = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] await asyncio.sleep(delay) if not http2_enabled: _LOGGER.debug( "%s: HTTP2Push connection closed; retries exceeded; polling", hide_email(email), ) coordinator = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get("coordinator") if coordinator: coordinator.update_interval = timedelta( seconds=scan_interval * 10 if http2_enabled else scan_interval ) _LOGGER.debug( "HTTP2push: %s, Polling interval: %s", http2_enabled, coordinator.update_interval, ) await coordinator.async_request_refresh() @callback async def http2_error_handler(message): """Handle http2push error. This currently logs the error. In the future, this should invalidate the http2push and determine if a reconnect should be done. """ email: str = login_obj.email errors = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] _LOGGER.debug( "%s: Received http2push error #%i %s: type %s", hide_email(email), errors, message, type(message), ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2"] = None if not login_obj.close_requested and ( login_obj.session.closed or isinstance(message, AlexapyLoginError) ): hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] = 5 _LOGGER.debug("%s: Login error detected.", hide_email(email)) hass.bus.async_fire( "alexa_media_relogin_required", event_data={"email": hide_email(email), "url": login_obj.url}, ) return hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2error"] = errors + 1 _LOGGER.debug("Setting up Alexa devices for %s", hide_email(login_obj.email)) config = config_entry.data email = config.get(CONF_EMAIL) include = ( cv.ensure_list_csv(config[CONF_INCLUDE_DEVICES]) if config[CONF_INCLUDE_DEVICES] else "" ) _LOGGER.debug("include: %s", include) exclude = ( cv.ensure_list_csv(config[CONF_EXCLUDE_DEVICES]) if config[CONF_EXCLUDE_DEVICES] else "" ) _LOGGER.debug("exclude: %s", exclude) scan_interval: float = ( config.get(CONF_SCAN_INTERVAL).total_seconds() if isinstance(config.get(CONF_SCAN_INTERVAL), timedelta) else config.get(CONF_SCAN_INTERVAL) ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["login_obj"] = login_obj # Initialize the per-account probe worker exactly once here (not per push message). _init_last_called_probe_worker(hass.data[DATA_ALEXAMEDIA]["accounts"][email]) _t = time.monotonic() http2_enabled = hass.data[DATA_ALEXAMEDIA]["accounts"][email]["http2"] = ( await http2_connect() ) _LOGGER.debug("[BOOT] http2_connect in %.2fs", time.monotonic() - _t) coordinator = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get("coordinator") # Get runtime_data for optimized coordinator runtime_data = ( config_entry.runtime_data if hasattr(config_entry, "runtime_data") else None ) if coordinator is None: _LOGGER.debug("%s: Creating optimized coordinator", hide_email(email)) # Use optimized coordinator with debouncing coordinator = AlexaMediaCoordinator( hass=hass, runtime_data=runtime_data, update_method=async_update_data, scan_interval=scan_interval, ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["coordinator"] = coordinator # Set correct interval now that http2 status is known coordinator.set_http2_status(bool(http2_enabled)) # Also store in runtime_data for type-safe access if runtime_data: runtime_data.coordinator = coordinator else: _LOGGER.debug("%s: setup_alexa: Reusing coordinator", hide_email(email)) # Use the optimized set_http2_status method if available if isinstance(coordinator, AlexaMediaCoordinator): coordinator.set_http2_status(bool(http2_enabled)) else: coordinator.update_interval = timedelta( seconds=scan_interval * 10 if http2_enabled else scan_interval ) # Fetch initial data _LOGGER.debug("%s: setup_alexa: Starting coordinator refresh", hide_email(email)) _t = time.monotonic() await coordinator.async_config_entry_first_refresh() _LOGGER.debug("[BOOT] first_refresh in %.2fs", time.monotonic() - _t) # Register services (fast - just registers callbacks) hass.data[DATA_ALEXAMEDIA]["services"] = alexa_services = AlexaMediaServices( hass, functions={"update_last_called": update_last_called} ) await alexa_services.register() # Update last_called in background to avoid blocking _LOGGER.debug("%s: setup_alexa: Scheduling last_called update", hide_email(email)) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["last_called_init_task"] = ( hass.async_create_background_task( _async_update_last_called_background(hass, login_obj, email), f"{DOMAIN}_last_called_init", ) ) return True async def _async_update_last_called_background( hass: HomeAssistant, login_obj, email: str ) -> None: """Update last_called in background to avoid blocking startup.""" try: await _async_update_last_called_global(hass, login_obj, email) _LOGGER.debug("%s: Background last_called update completed", hide_email(email)) # Record metrics metrics = get_metrics(hass) if metrics: metrics.record_boot_stage(f"last_called_{hide_email(email)}") except asyncio.CancelledError: # Task cancelled during unload/shutdown; propagate cancellation. raise except Exception as err: _LOGGER.debug( "%s: Background last_called update failed: %s", hide_email(email), err ) async def async_unload_entry(hass, entry) -> bool: """Unload a config entry""" email = entry.data["email"] _LOGGER.debug("Unloading entry: %s", hide_email(email)) for task_key in ( "notifications_refresh_task", "notifications_init_task", "last_called_init_task", "service_update_last_called_task", ): accounts = hass.data.get(DATA_ALEXAMEDIA, {}).get("accounts", {}) account = accounts.get(email) if not account: return True task = account.get(task_key) if task and not task.done(): task.cancel() try: await task except asyncio.CancelledError: # Expected during unload/shutdown pass last_called_task = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get( "last_called_probe_task" ) if last_called_task and not last_called_task.done(): last_called_task.cancel() try: await last_called_task except asyncio.CancelledError: # Expected during unload/shutdown pass except Exception as err: # pragma: no cover _LOGGER.debug( "%s: Exception while cancelling last_called_probe_task: %s", hide_email(email), err, ) hass.data[DATA_ALEXAMEDIA]["accounts"][email]["last_called_probe_task"] = None debouncer = hass.data[DATA_ALEXAMEDIA]["accounts"][email].get( "confirm_refresh_debouncer" ) if debouncer: debouncer.async_cancel() hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "confirm_refresh_debouncer" ] = None for component in ALEXA_COMPONENTS + DEPENDENT_ALEXA_COMPONENTS: try: if component == "notify": await notify_async_unload_entry(hass, entry) else: _LOGGER.debug("Forwarding unload entry to %s", component) await hass.config_entries.async_forward_entry_unload(entry, component) except Exception: _LOGGER.error("Error unloading: %s", component) await close_connections(hass, email) for listener in hass.data[DATA_ALEXAMEDIA]["accounts"][email][DATA_LISTENER]: listener() hass.data[DATA_ALEXAMEDIA]["accounts"].pop(email) # Clean up config flows in progress flows_to_remove = [] if hass.data[DATA_ALEXAMEDIA].get("config_flows"): for key, flow in hass.data[DATA_ALEXAMEDIA]["config_flows"].items(): if key.startswith(email) and flow: _LOGGER.debug("Aborting flow %s %s", key, flow) flows_to_remove.append(key) try: hass.config_entries.flow.async_abort(flow.get("flow_id")) except UnknownFlow: pass for flow in flows_to_remove: hass.data[DATA_ALEXAMEDIA]["config_flows"].pop(flow) # Clean up hass.data if not hass.data[DATA_ALEXAMEDIA].get("accounts"): _LOGGER.debug("Removing accounts data and services") hass.data[DATA_ALEXAMEDIA].pop("accounts") alexa_services = hass.data[DATA_ALEXAMEDIA].get("services") if alexa_services: await alexa_services.unregister() hass.data[DATA_ALEXAMEDIA].pop("services") if hass.data[DATA_ALEXAMEDIA].get("config_flows") == {}: _LOGGER.debug("Removing config_flows data") async_dismiss_persistent_notification( hass, f"alexa_media_{slugify(email)}{slugify((entry.data['url'])[7:])}" ) hass.data[DATA_ALEXAMEDIA].pop("config_flows") if not hass.data[DATA_ALEXAMEDIA]: _LOGGER.debug("Removing alexa_media data structure") if hass.data.get(DATA_ALEXAMEDIA): hass.data.pop(DATA_ALEXAMEDIA) else: _LOGGER.debug( "Unable to remove alexa_media data structure: %s", hass.data.get(DATA_ALEXAMEDIA), ) _LOGGER.debug("Unloaded entry for %s", hide_email(email)) return True async def async_remove_entry(hass, entry) -> bool: """Handle removal of an entry.""" email = entry.data["email"] obfuscated_email = hide_email(email) _LOGGER.debug("Removing config entry: %s", hide_email(email)) login_obj = AlexaLogin( url="", email=email, password="", # nosec outputpath=hass.config.path, ) # Delete cookiefile cookiefile = hass.config.path(f".storage/{DOMAIN}.{email}.pickle") obfuscated_cookiefile = hass.config.path( f".storage/{DOMAIN}.{obfuscated_email}.pickle" ) if callable(getattr(AlexaLogin, "delete_cookiefile", None)): try: await login_obj.delete_cookiefile() _LOGGER.debug("Cookiefile %s deleted.", obfuscated_cookiefile) except Exception as ex: _LOGGER.error( "delete_cookiefile() exception: %s;" " Manually delete cookiefile before re-adding the integration: %s", ex, obfuscated_cookiefile, ) else: if os.path.exists(cookiefile): try: await alexapy_delete_cookie(cookiefile) _LOGGER.debug( "Successfully deleted cookiefile: %s", obfuscated_cookiefile ) except (OSError, EOFError, TypeError, AttributeError) as ex: _LOGGER.error( "alexapy_delete_cookie() exception: %s;" " Manually delete cookiefile before re-adding the integration: %s", ex, obfuscated_cookiefile, ) else: _LOGGER.error("Cookiefile not found: %s", obfuscated_cookiefile) _LOGGER.debug("Config entry %s removed.", obfuscated_email) return True async def close_connections(hass, email: str) -> None: """Clear open aiohttp connections for email.""" if ( email not in hass.data[DATA_ALEXAMEDIA]["accounts"] or "login_obj" not in hass.data[DATA_ALEXAMEDIA]["accounts"][email] ): return account_dict = hass.data[DATA_ALEXAMEDIA]["accounts"][email] login_obj = account_dict["login_obj"] await login_obj.save_cookiefile() await login_obj.close() _LOGGER.debug( "%s: Connection closed: %s", hide_email(email), login_obj.session.closed ) async def update_listener(hass, config_entry): """Update when config_entry options update.""" account = config_entry.data email = account.get(CONF_EMAIL) reload_needed: bool = False for key, old_value in hass.data[DATA_ALEXAMEDIA]["accounts"][email][ "options" ].items(): new_value = config_entry.data.get(key) if new_value is not None and new_value != old_value: hass.data[DATA_ALEXAMEDIA]["accounts"][email]["options"][key] = new_value _LOGGER.debug( "Option %s changed from %s to %s", key, old_value, hass.data[DATA_ALEXAMEDIA]["accounts"][email]["options"][key], ) reload_needed = True if reload_needed: await hass.config_entries.async_reload(config_entry.entry_id) _LOGGER.debug( "%s options reloaded", hass.data[DATA_ALEXAMEDIA]["accounts"][email], ) async def test_login_status(hass, config_entry, login) -> bool: """Test the login status and spawn requests for info.""" _LOGGER.debug("Testing login status: %s", login.status) if login.status and login.status.get("login_successful"): return True account = config_entry.data _LOGGER.debug("Logging in: %s %s", obfuscate(account), in_progress_instances(hass)) _LOGGER.debug("Login stats: %s", login.stats) message: str = ( f"Reauthenticate {login.email} on the [Integrations](/config/integrations) page. " ) if login.stats.get("login_timestamp") != datetime(1, 1, 1): elaspsed_time: str = str(datetime.now() - login.stats.get("login_timestamp")) api_calls: int = login.stats.get("api_calls") message += f"Relogin required after {elaspsed_time} and {api_calls} api calls." host = urlparse(login.url).hostname or login.url async_create_persistent_notification( hass, title="Alexa Media Reauthentication Required", message=message, notification_id=f"alexa_media_{slugify(login.email)}_{slugify(host)}", ) flow = hass.data[DATA_ALEXAMEDIA]["config_flows"].get( f"{account[CONF_EMAIL]} - {account[CONF_URL]}" ) if flow: if flow.get("flow_id") in in_progress_instances(hass): _LOGGER.debug("Existing config flow detected") return False _LOGGER.debug("Stopping orphaned config flow %s", flow.get("flow_id")) try: hass.config_entries.flow.async_abort(flow.get("flow_id")) except UnknownFlow: pass hass.data[DATA_ALEXAMEDIA]["config_flows"][ f"{account[CONF_EMAIL]} - {account[CONF_URL]}" ] = None _LOGGER.debug("Creating new config flow to login") config_entry.async_start_reauth( hass, context={"source": SOURCE_REAUTH}, data={ CONF_EMAIL: account[CONF_EMAIL], CONF_PASSWORD: account[CONF_PASSWORD], CONF_URL: account[CONF_URL], CONF_DEBUG: account[CONF_DEBUG], CONF_INCLUDE_DEVICES: account[CONF_INCLUDE_DEVICES], CONF_EXCLUDE_DEVICES: account[CONF_EXCLUDE_DEVICES], CONF_SCAN_INTERVAL: ( account[CONF_SCAN_INTERVAL].total_seconds() if isinstance(account[CONF_SCAN_INTERVAL], timedelta) else account[CONF_SCAN_INTERVAL] ), CONF_OTPSECRET: account.get(CONF_OTPSECRET, ""), }, ) try: flow_obj = config_entry.async_get_active_flows(hass, {SOURCE_REAUTH}).__next__() hass.data[DATA_ALEXAMEDIA]["config_flows"][ f"{account[CONF_EMAIL]} - {account[CONF_URL]}" ] = flow_obj except StopIteration: _LOGGER.debug("A new config flow could not be created.") return False