Compare commits

..
Author SHA1 Message Date
SMKRV c13ef1921d feat(performance): Optimize system resources and token estimation
- Improve JSON history file processing
- Add memory and disk space validation
- Enhance parallel request handling
- Refine token counting heuristics
2024-12-06 02:51:06 +03:00
SMKRV 0fb9acfa8f docs:(README) 2024-12-05 02:00:24 +03:00
SMKRV 2ac4389e1a docs:(README) 2024-12-05 01:55:31 +03:00
SMKRV de194e425c docs:(README) 2024-12-05 01:54:09 +03:00
SMKRV 95f3d6506e docs:(README) 2024-12-05 01:49:26 +03:00
SMKRV 3e877243b9 docs:(README) 2024-12-05 01:46:28 +03:00
SMKRV fd795c9e90 docs:(readme) 2024-12-05 00:36:08 +03:00
SMKRV 5f5dc041b9 docs: structure 2024-12-04 23:32:05 +03:00
SMKRV 3d3885d43d Update README.md 2024-12-04 23:28:14 +03:00
smkrvandGitHub b059744716 Update README.md 2024-12-04 23:23:46 +03:00
SMKRV f2a41aaa2c docs: logo 2024-12-04 20:13:08 +03:00
SMKRV 1fcf751d0d docs: logo 2024-12-04 19:24:30 +03:00
SMKRV a45c9407fd docs: logo 2024-12-04 19:19:40 +03:00
SMKRV 21fdf108cf docs: logo 2024-12-04 19:17:09 +03:00
SMKRV c7ec2b3ea4 docs: logo 2024-12-04 19:13:57 +03:00
SMKRV 6c4da6cea5 docs: logo 2024-12-04 19:12:44 +03:00
SMKRV 3029c5dd26 docs: logo 2024-12-04 19:12:21 +03:00
SMKRV 3e97028094 docs: logo 2024-12-04 19:11:54 +03:00
SMKRV 1194c87134 docs: logo 2024-12-04 17:51:40 +03:00
SMKRV e03b5315c2 logo update 2024-12-04 17:51:02 +03:00
SMKRV fad887492f docs: logo 2024-12-04 17:48:04 +03:00
SMKRV f26df29937 docs: logo 2024-12-04 17:47:12 +03:00
SMKRV 69b54a08ba logo update 2024-12-04 17:42:32 +03:00
SMKRV 7ce77eb18e docs: logo 2024-12-04 17:36:18 +03:00
SMKRV 1480669dc8 icons update 2024-12-04 17:34:19 +03:00
16 changed files with 807 additions and 424 deletions
+58 -24
View File
@@ -4,7 +4,7 @@
![GitHub release](https://img.shields.io/github/release/smkrv/ha-text-ai.svg?style=flat-square) ![GitHub last commit](https://img.shields.io/github/last-commit/smkrv/ha-text-ai.svg?style=flat-square) [![License: CC BY-NC-SA 4.0](https://img.shields.io/badge/License-CC%20BY--NC--SA%204.0-lightgrey.svg?style=flat-square)](https://creativecommons.org/licenses/by-nc-sa/4.0/) [![hacs_badge](https://img.shields.io/badge/HACS-Custom-41BDF5.svg?style=flat-square)](https://github.com/hacs/integration) ![/README.md](https://img.shields.io/badge/language-English-green?style=flat-square) ![/README_RU.md](https://img.shields.io/badge/language-Russian-green?style=flat-square) ![](https://img.shields.io/badge/language-Deutch-green?style=flat-square) ![GitHub release](https://img.shields.io/github/release/smkrv/ha-text-ai.svg?style=flat-square) ![GitHub last commit](https://img.shields.io/github/last-commit/smkrv/ha-text-ai.svg?style=flat-square) [![License: CC BY-NC-SA 4.0](https://img.shields.io/badge/License-CC%20BY--NC--SA%204.0-lightgrey.svg?style=flat-square)](https://creativecommons.org/licenses/by-nc-sa/4.0/) [![hacs_badge](https://img.shields.io/badge/HACS-Custom-41BDF5.svg?style=flat-square)](https://github.com/hacs/integration) ![/README.md](https://img.shields.io/badge/language-English-green?style=flat-square) ![/README_RU.md](https://img.shields.io/badge/language-Russian-green?style=flat-square) ![](https://img.shields.io/badge/language-Deutch-green?style=flat-square)
<img src="https://github.com/smkrv/ha-text-ai/blob/524849f6a945ec62c2cf6a6b7ecd9a28b37bf0fa/misc/icons/logo.jpg" alt="HA Text AI" height="160"/> <img src="https://github.com/smkrv/ha-text-ai/blob/e03b5315c2cc7b489a6e1175da728075275ed614/misc/icons/logo.png" alt="HA Text AI" style="width: 80%; max-width: 640px; max-height: 150px; aspect-ratio: 2/1; object-fit: contain;"/>
### Advanced AI Integration for [Home Assistant](https://www.home-assistant.io/) with LLM multi-provider support ### Advanced AI Integration for [Home Assistant](https://www.home-assistant.io/) with LLM multi-provider support
</div> </div>
@@ -27,54 +27,68 @@ Transform your smart home experience with powerful AI assistance powered by mult
## 🌟 Features ## 🌟 Features
- 🧠 **Multi-Provider AI Integration**: - 🧠 **Multi-Provider AI Integration**: Support for OpenAI GPT and Anthropic Claude models
- 💬 **Advanced Language Processing**: Context-aware, multi-turn conversations
- 📝 **Enhanced Memory Management**: Secure file-based history storage
-**Performance Optimization**: Efficient token usage and smart rate limiting
- 🎯 **Advanced Customization**: Per-request model and parameter selection
- 🔒 **Enhanced Security**: Secure API key management and usage monitoring
- 🎨 **Improved User Experience**: Intuitive configuration and rich interfaces
- 🔄 **Automation Integration**: Event-driven responses and template compatibility
<details>
<summary>📦 Detailed Feature Breakdown</summary>
### 🧠 **Multi-Provider AI Integration**
- Support for OpenAI GPT models - Support for OpenAI GPT models
- Anthropic Claude integration - Anthropic Claude integration
- Custom API endpoints - Custom API endpoints
- Flexible model selection - Flexible model selection
- 💬 **Advanced Language Processing**: ### 💬 **Advanced Language Processing**
- Context-aware responses - Context-aware responses
- Multi-turn conversations - Multi-turn conversations
- Custom system instructions - Custom system instructions
- Natural conversation flow - Natural conversation flow
- 📝 **Enhanced Memory Management**: ### 📝 **Enhanced Memory Management**
- File-based conversation history storage - File-based conversation history storage
- Automatic history rotation - Automatic history rotation
- Configurable history size limits - Configurable history size limits
- Secure storage in Home Assistant - Secure storage in Home Assistant
-**Performance Optimization**: ### ⚡ **Performance Optimization**
- Efficient token usage - Efficient token usage
- Smart rate limiting - Smart rate limiting
- Response caching - Response caching
- Request interval control - Request interval control
- 🎯 **Advanced Customization**: ### 🎯 **Advanced Customization**
- Per-request model selection - Per-request model selection
- Adjustable parameters - Adjustable parameters
- Custom system prompts - Custom system prompts
- Temperature control - Temperature control
- 🔒 **Enhanced Security**: ### 🔒 **Enhanced Security**
- Secure API key storage - Secure API key storage
- Rate limiting protection - Rate limiting protection
- Error handling - Error handling
- Usage monitoring - Usage monitoring
- 🎨 **Improved User Experience**: ### 🎨 **Improved User Experience**
- Intuitive configuration UI - Intuitive configuration UI
- Detailed sensor attributes - Detailed sensor attributes
- Rich service interface - Rich service interface
- Model selection UI - Model selection UI
- 🔄 **Automation Integration**: ### 🔄 **Automation Integration**
- Event-driven responses - Event-driven responses
- Conditional logic support - Conditional logic support
- Template compatibility - Template compatibility
- Model-specific automation - Model-specific automation
</details>
## 📋 Prerequisites ## 📋 Prerequisites
- Home Assistant 2024.11 or later - Home Assistant 2024.11 or later
@@ -85,16 +99,22 @@ Transform your smart home experience with powerful AI assistance powered by mult
- Python 3.9 or newer - Python 3.9 or newer
- Stable internet connection - Stable internet connection
### Configuration Options ## Configuration Options
- API Provider (OpenAI/Anthropic)
- API Key (provider-specific)
- Model Selection (flexible, provider-specific models)
- Temperature (Creativity control, 0.0-2.0)
- Max Tokens (Response length limit)
- Request Interval (API call throttling)
- Custom API Endpoint (optional)
#### ⓘ Potentially Compatible Providers ### 🔧 **Core Configuration Settings**
- 🌐 **API Provider**: OpenAI/Anthropic
- 🔑 **API Key**: Provider-specific authentication
- 🤖 **Model Selection**: Flexible, provider-specific models
- 🌡️ **Temperature**: Creativity control (0.0-2.0)
- 📏 **Max Tokens**: Response length limit
- ⏱️ **Request Interval**: API call throttling
- 💾 **History Size**: Number of messages to retain
- 🌍 **Custom API Endpoint**: Optional advanced configuration
<details>
<summary>🌐 Potentially Compatible Providers</summary>
#### Flexible Provider Ecosystem
The integration is designed to be flexible and may work with other providers offering OpenAI-compatible APIs: The integration is designed to be flexible and may work with other providers offering OpenAI-compatible APIs:
- Groq - Groq
- Together AI - Together AI
@@ -104,18 +124,19 @@ The integration is designed to be flexible and may work with other providers off
- Local AI servers (like Ollama) - Local AI servers (like Ollama)
- Custom OpenAI-compatible endpoints - Custom OpenAI-compatible endpoints
#### Additional Notes #### 🚨 Compatibility Notes
- Not all providers guarantee full compatibility - Not all providers guarantee full compatibility
- Performance may vary between providers - Performance may vary between providers
- Check individual provider's documentation - Check individual provider's documentation
- Ensure your API key has sufficient credits/quota - Ensure your API key has sufficient credits/quota
#### Provider Compatibility Requirements #### 🔍 Provider Compatibility Requirements
To be compatible, a provider should support: To be compatible, a provider should support:
- OpenAI-like REST API structure - OpenAI-like REST API structure
- JSON request/response format - JSON request/response format
- Standard authentication method - Standard authentication method
- Similar model parameter handling - Similar model parameter handling
</details>
## ⚡ Installation ## ⚡ Installation
@@ -144,7 +165,8 @@ To be compatible, a provider should support:
3. Search for "HA Text AI" 3. Search for "HA Text AI"
4. Follow the configuration steps 4. Follow the configuration steps
### Via YAML <details>
<summary>📦 Via YAML (Advanced)</summary>
### Platform Configuration (Global Settings) ### Platform Configuration (Global Settings)
@@ -202,6 +224,8 @@ sensor:
| `temperature` | Float | ❌ | Platform setting | Override global temperature | | `temperature` | Float | ❌ | Platform setting | Override global temperature |
| `max_tokens` | Integer | ❌ | Platform setting | Override global max tokens | | `max_tokens` | Integer | ❌ | Platform setting | Override global max tokens |
</details>
## 🛠️ Available Services ## 🛠️ Available Services
### ask_question ### ask_question
@@ -258,7 +282,7 @@ sensor.ha_text_ai_YOUR_UNIQUE_SUFFIX
# Examples: # Examples:
sensor.ha_text_ai_gpt # GPT-based sensor sensor.ha_text_ai_gpt # GPT-based sensor
sensor.ha_text_ai_claude # Claude-based sensor sensor.ha_text_ai_claude # Claude-based sensor
sensor.ha_text_ai_gpt # Custom suffix sensor.ha_text_ai_abc # Custom suffix
``` ```
#### Response Retrieval #### Response Retrieval
@@ -291,6 +315,16 @@ automation:
### 🔍 HA Text AI Sensor Attributes ### 🔍 HA Text AI Sensor Attributes
- 🤖 **Model and Provider Information**: Tracking current AI model and service provider
- 🚦 **System Status**: Real-time API and processing readiness
- 📊 **Performance Metrics**: Request success rates and response times
- 💬 **Conversation Tracking**: Token usage and interaction history (for more precise token counting, install `tiktoken`; otherwise, a fallback estimation method is automatically used)
- 🕒 **Last Interaction Details**: Recent query and response tracking
- ❤️ **System Health**: Error monitoring and service uptime
<details>
<summary>📦 Detailed Sensor Attributes</summary>
#### Model and Provider Information #### Model and Provider Information
```yaml ```yaml
# Name of the AI model currently in use (e.g., latest version of GPT) # Name of the AI model currently in use (e.g., latest version of GPT)
@@ -352,7 +386,6 @@ automation:
# Last few conversation entries (limited to 3 for performance) # Last few conversation entries (limited to 3 for performance)
{{ state_attr('sensor.ha_text_ai_gpt', 'conversation_history') }} # [...] {{ state_attr('sensor.ha_text_ai_gpt', 'conversation_history') }} # [...]
``` ```
#### Last Interaction Details #### Last Interaction Details
@@ -391,6 +424,7 @@ Conversation history stored in `.storage/ha_text_ai_history/` directory:
- Use these attributes for monitoring and automation - Use these attributes for monitoring and automation
- Some values might be 0 or empty initially - Some values might be 0 or empty initially
</details>
## 📘 FAQ ## 📘 FAQ
@@ -427,7 +461,6 @@ A: Yes, archived history files are stored with timestamps and can be accessed ma
**Q: How much history is kept?** **Q: How much history is kept?**
A: By default, up to 100 conversations are stored, but this can be configured. Files are automatically rotated when they reach 1MB. A: By default, up to 100 conversations are stored, but this can be configured. Files are automatically rotated when they reach 1MB.
## 🤝 Contributing ## 🤝 Contributing
Contributions welcome! Please read our [Contributing Guide](CONTRIBUTING.md). Contributions welcome! Please read our [Contributing Guide](CONTRIBUTING.md).
@@ -474,6 +507,7 @@ If you want to say thanks financially, you can send a small token of appreciatio
--- ---
<div align="center"><img src="https://github.com/smkrv/ha-text-ai/blob/f2a41aaa2cd88adaf6929ba5f228205a315efdfa/custom_components/ha_text_ai/icons/dark_icon%402x.png" alt="HA Text AI" style="width: 128px; height: auto;"/></div>
<div align="center"> <div align="center">
Made with ❤️ and Claude 3.5 Sonnet for the Home Assistant Community Made with ❤️ and Claude 3.5 Sonnet for the Home Assistant Community
+63 -13
View File
@@ -11,8 +11,9 @@ from __future__ import annotations
import logging import logging
import os import os
import shutil import shutil
import hashlib
from datetime import datetime, timedelta from datetime import datetime, timedelta
from typing import Any, Dict from typing import Any, Dict, TypeVar
import voluptuous as vol import voluptuous as vol
from async_timeout import timeout from async_timeout import timeout
@@ -52,11 +53,13 @@ from .const import (
SERVICE_SET_SYSTEM_PROMPT, SERVICE_SET_SYSTEM_PROMPT,
DEFAULT_MAX_HISTORY, DEFAULT_MAX_HISTORY,
CONF_MAX_HISTORY_SIZE, CONF_MAX_HISTORY_SIZE,
ICONS_SUBDOMAIN,
) )
_LOGGER = logging.getLogger(__name__) _LOGGER = logging.getLogger(__name__)
CONFIG_SCHEMA = cv.config_entry_only_config_schema(DOMAIN) CONFIG_SCHEMA = cv.config_entry_only_config_schema(DOMAIN)
ConfigType = TypeVar("ConfigType", bound=Dict[str, Any])
SERVICE_SCHEMA_ASK_QUESTION = vol.Schema({ SERVICE_SCHEMA_ASK_QUESTION = vol.Schema({
vol.Required("instance"): cv.string, vol.Required("instance"): cv.string,
@@ -90,19 +93,18 @@ def get_coordinator_by_instance(hass: HomeAssistant, instance: str) -> HATextAIC
raise HomeAssistantError(f"Instance {instance} not found") raise HomeAssistantError(f"Instance {instance} not found")
async def async_setup(hass: HomeAssistant, config: Dict[str, Any]) -> bool: def get_file_hash(file_path: str) -> str:
"""Set up the HA Text AI component.""" """Calculate SHA256 hash of file."""
hass.data.setdefault(DOMAIN, {}) sha256_hash = hashlib.sha256()
with open(file_path, "rb") as f:
for byte_block in iter(lambda: f.read(4096), b""):
sha256_hash.update(byte_block)
return sha256_hash.hexdigest()
try: async def async_setup(hass: HomeAssistant, config: ConfigType) -> bool:
source = os.path.join(os.path.dirname(__file__), 'icons', 'icon.svg') """Set up the Home Assistant Text AI component."""
dest_dir = os.path.join(hass.config.path('www'), 'icons') # Initialize domain data storage
os.makedirs(dest_dir, exist_ok=True) hass.data.setdefault(DOMAIN, {})
dest = os.path.join(dest_dir, 'icon.png')
if not os.path.exists(dest):
shutil.copyfile(source, dest)
except Exception as ex:
_LOGGER.warning("Failed to copy custom icon: %s", str(ex))
async def async_ask_question(call: ServiceCall) -> None: async def async_ask_question(call: ServiceCall) -> None:
"""Handle ask_question service.""" """Handle ask_question service."""
@@ -150,6 +152,7 @@ async def async_setup(hass: HomeAssistant, config: Dict[str, Any]) -> bool:
_LOGGER.error("Error setting system prompt: %s", str(err)) _LOGGER.error("Error setting system prompt: %s", str(err))
raise HomeAssistantError(f"Failed to set system prompt: {str(err)}") raise HomeAssistantError(f"Failed to set system prompt: {str(err)}")
# Register services
hass.services.async_register( hass.services.async_register(
DOMAIN, DOMAIN,
SERVICE_ASK_QUESTION, SERVICE_ASK_QUESTION,
@@ -178,6 +181,53 @@ async def async_setup(hass: HomeAssistant, config: Dict[str, Any]) -> bool:
schema=SERVICE_SCHEMA_SET_SYSTEM_PROMPT schema=SERVICE_SCHEMA_SET_SYSTEM_PROMPT
) )
# Handle icons
try:
source_icon_path = os.path.join(
os.path.dirname(__file__),
ICONS_SUBDOMAIN,
'icon@2x.png'
)
destination_directory = os.path.join(
hass.config.path('www'),
DOMAIN,
ICONS_SUBDOMAIN
)
destination_icon_path = os.path.join(
destination_directory,
'icon.png'
)
if not os.path.exists(source_icon_path):
_LOGGER.error("Source icon not found: %s", source_icon_path)
return True
def create_directory():
os.makedirs(destination_directory, exist_ok=True)
await hass.async_add_executor_job(create_directory)
should_copy = True
if os.path.exists(destination_icon_path):
source_hash = await hass.async_add_executor_job(get_file_hash, source_icon_path)
dest_hash = await hass.async_add_executor_job(get_file_hash, destination_icon_path)
should_copy = source_hash != dest_hash
if should_copy:
def copy_file():
shutil.copyfile(source_icon_path, destination_icon_path)
await hass.async_add_executor_job(copy_file)
_LOGGER.debug("Icon updated: %s", destination_icon_path)
except PermissionError as e:
_LOGGER.error("Permission denied when managing icons: %s", str(e))
except Exception as e:
_LOGGER.error("Failed to manage icons: %s", str(e))
return True return True
async def async_check_api(session, endpoint: str, headers: dict, provider: str) -> bool: async def async_check_api(session, endpoint: str, headers: dict, provider: str) -> bool:
+3 -1
View File
@@ -41,7 +41,9 @@ CONF_IS_ANTHROPIC: Final = "is_anthropic"
CONF_CONTEXT_MESSAGES: Final = "context_messages" CONF_CONTEXT_MESSAGES: Final = "context_messages"
ABSOLUTE_MAX_HISTORY_SIZE = 500 ABSOLUTE_MAX_HISTORY_SIZE = 500
MAX_ENTRY_SIZE = 1 * 1024 * 1024 MAX_ATTRIBUTE_SIZE = 4 * 1024
MAX_HISTORY_FILE_SIZE = 1 * 1024 * 1024
ICONS_SUBDOMAIN = "icons"
# Default values # Default values
DEFAULT_MODEL: Final = "gpt-4o-mini" DEFAULT_MODEL: Final = "gpt-4o-mini"
+516 -247
View File
@@ -12,9 +12,13 @@ import logging
import traceback import traceback
import aiofiles import aiofiles
import os import os
import json
import asyncio import asyncio
import psutil
import re
import math
from datetime import datetime, timedelta from datetime import datetime, timedelta
from typing import Any, Dict, List, Optional from typing import Any, Dict, List, Optional, Union
from homeassistant.core import HomeAssistant from homeassistant.core import HomeAssistant
from homeassistant.helpers.update_coordinator import DataUpdateCoordinator from homeassistant.helpers.update_coordinator import DataUpdateCoordinator
@@ -35,11 +39,26 @@ from .const import (
DEFAULT_MAX_HISTORY, DEFAULT_MAX_HISTORY,
DEFAULT_CONTEXT_MESSAGES, DEFAULT_CONTEXT_MESSAGES,
ABSOLUTE_MAX_HISTORY_SIZE, ABSOLUTE_MAX_HISTORY_SIZE,
MAX_ENTRY_SIZE, MAX_ATTRIBUTE_SIZE,
MAX_HISTORY_FILE_SIZE,
) )
_LOGGER = logging.getLogger(__name__) _LOGGER = logging.getLogger(__name__)
def _check_memory_available(self):
"""Check if enough memory is available."""
memory = psutil.virtual_memory()
# Log the total and available memory
_LOGGER.debug("Total memory: %s, Available memory: %s", memory.total, memory.available)
if memory.available > 1024 * 1024 * 100: # 100MB
_LOGGER.debug("Sufficient memory available: %s bytes", memory.available)
return True
else:
_LOGGER.warning("Insufficient memory available: %s bytes", memory.available)
return False
class AsyncFileHandler: class AsyncFileHandler:
"""Async context manager for file operations.""" """Async context manager for file operations."""
@@ -62,7 +81,7 @@ class HATextAICoordinator(DataUpdateCoordinator):
hass: HomeAssistant, hass: HomeAssistant,
client: Any, client: Any,
model: str, model: str,
update_interval: int, # Moved up update_interval: int,
instance_name: str, instance_name: str,
max_tokens: int = DEFAULT_MAX_TOKENS, max_tokens: int = DEFAULT_MAX_TOKENS,
temperature: float = DEFAULT_TEMPERATURE, temperature: float = DEFAULT_TEMPERATURE,
@@ -78,26 +97,13 @@ class HATextAICoordinator(DataUpdateCoordinator):
from .config_flow import normalize_name from .config_flow import normalize_name
self.normalized_name = normalize_name(instance_name) self.normalized_name = normalize_name(instance_name)
self._metrics_file = os.path.join(
self.hass = hass hass.config.path(".storage"),
self.client = client "ha_text_ai_history",
self.model = model f"ha_text_ai_metrics_{normalize_name(instance_name)}.json"
self.temperature = temperature
self.max_tokens = max_tokens
self.max_history_size = min(
max(1, max_history_size),
ABSOLUTE_MAX_HISTORY_SIZE
) )
self.is_anthropic = is_anthropic
# Initialize essential attributes first
self._is_processing = False
self._is_rate_limited = False
self._is_maintenance = False
self.endpoint_status = "ready"
self._system_prompt = None
self._conversation_history = []
# Initialize performance metrics first
self._performance_metrics = { self._performance_metrics = {
"total_tokens": 0, "total_tokens": 0,
"prompt_tokens": 0, "prompt_tokens": 0,
@@ -110,6 +116,26 @@ class HATextAICoordinator(DataUpdateCoordinator):
"min_latency": float("inf"), "min_latency": float("inf"),
} }
# Continue with other initializations
self.hass = hass
self.client = client
self.model = model
self.temperature = temperature
self.max_tokens = max_tokens
self.max_history_size = min(
max(1, max_history_size),
ABSOLUTE_MAX_HISTORY_SIZE
)
self.is_anthropic = is_anthropic
# Initialize essential attributes
self._is_processing = False
self._is_rate_limited = False
self._is_maintenance = False
self.endpoint_status = "ready"
self._system_prompt = None
self._conversation_history = []
self._last_response = { self._last_response = {
"timestamp": dt_util.utcnow().isoformat(), "timestamp": dt_util.utcnow().isoformat(),
"question": "", "question": "",
@@ -122,7 +148,7 @@ class HATextAICoordinator(DataUpdateCoordinator):
update_interval_td = timedelta(seconds=update_interval) update_interval_td = timedelta(seconds=update_interval)
# Call super().__init__ BEFORE other property access # Call super().__init__
super().__init__( super().__init__(
hass, hass,
_LOGGER, _LOGGER,
@@ -130,11 +156,11 @@ class HATextAICoordinator(DataUpdateCoordinator):
update_interval=update_interval_td, update_interval=update_interval_td,
) )
# Now initialize _initial_state (after super().__init__) # Now initialize _initial_state
self._initial_state = { self._initial_state = {
"state": STATE_READY, "state": STATE_READY,
"metrics": self._performance_metrics.copy(), "metrics": self._performance_metrics.copy(),
"last_response": self.last_response.copy(), # Accessing here "last_response": self.last_response.copy(),
"is_processing": self._is_processing, "is_processing": self._is_processing,
"is_rate_limited": self._is_rate_limited, "is_rate_limited": self._is_rate_limited,
"is_maintenance": self._is_maintenance, "is_maintenance": self._is_maintenance,
@@ -156,17 +182,22 @@ class HATextAICoordinator(DataUpdateCoordinator):
hass.config.path(".storage"), hass.config.path(".storage"),
"ha_text_ai_history" "ha_text_ai_history"
) )
os.makedirs(self._history_dir, exist_ok=True) hass.async_create_task(self._create_history_dir())
hass.async_create_task(self._check_history_directory()) hass.async_create_task(self._check_history_directory())
hass.async_create_task(self._migrate_history_from_txt_to_json())
hass.async_create_task(self._initialize_metrics())
# History file path using instance name # History file path using instance name
self._history_file = os.path.join( self._history_file = os.path.join(
self._history_dir, self._history_dir,
f"{self.normalized_name}_history.txt" f"{self.normalized_name}_history.json"
) )
# Maximum history file size (10 MB)
self._max_history_file_size = 10 * 1024 * 1024 # Maximum history file size (1 MB) from const.py
self._max_history_file_size = MAX_HISTORY_FILE_SIZE
# Asynchronous file initialization # Asynchronous file initialization
hass.async_create_task(self.async_initialize_history_file()) hass.async_create_task(self.async_initialize_history_file())
@@ -221,46 +252,206 @@ class HATextAICoordinator(DataUpdateCoordinator):
""" """
self._last_response = value self._last_response = value
async def _file_exists(self, path: str) -> bool:
"""
Check if file exists asynchronously.
Args:
path (str): Full path to the file to check
Returns:
bool: True if file exists, False otherwise
"""
try:
result = await self.hass.async_add_executor_job(os.path.exists, path)
_LOGGER.debug(f"File existence check: {path} - {'Exists' if result else 'Not Found'}")
return result
except Exception as e:
_LOGGER.error(f"Error checking file existence for {path}: {e}")
return False
async def _create_history_dir(self):
"""
Asynchronously create history directory.
Creates the directory for storing history files
without blocking the event loop.
"""
try:
await self.hass.async_add_executor_job(
lambda: os.makedirs(self._history_dir, exist_ok=True)
)
_LOGGER.debug(f"Directory creation details: exist_ok=True")
except PermissionError:
_LOGGER.error(f"Permission denied when creating history directory: {self._history_dir}")
raise
except OSError as e:
_LOGGER.error(f"Error creating history directory {self._history_dir}: {e}")
raise
async def _initialize_metrics(self) -> None:
self._performance_metrics = await self._load_metrics() or {
"total_tokens": 0,
"prompt_tokens": 0,
"completion_tokens": 0,
"successful_requests": 0,
"failed_requests": 0,
"total_errors": 0,
"average_latency": 0,
"max_latency": 0,
"min_latency": float("inf"),
}
async def _load_metrics(self) -> Dict[str, Any]:
try:
if await self._file_exists(self._metrics_file):
def read_metrics():
with open(self._metrics_file, 'r') as f:
try:
return json.load(f)
except json.JSONDecodeError:
_LOGGER.warning("Metrics file corrupted, creating new")
return None
return await self.hass.async_add_executor_job(read_metrics)
except Exception as e:
_LOGGER.warning(f"Failed to load metrics: {e}")
return None
async def _save_metrics(self) -> None:
try:
def write_metrics():
with open(self._metrics_file, 'w') as f:
json.dump(self._performance_metrics, f)
await self.hass.async_add_executor_job(write_metrics)
except Exception as e:
_LOGGER.warning(f"Failed to save metrics: {e}")
async def _update_metrics(self, latency: float, response: dict) -> None:
"""Update performance metrics and save to storage."""
metrics = self._performance_metrics
tokens = response.get("tokens", {})
metrics["total_tokens"] += tokens.get("total", 0)
metrics["prompt_tokens"] += tokens.get("prompt", 0)
metrics["completion_tokens"] += tokens.get("completion", 0)
metrics["successful_requests"] += 1
metrics["average_latency"] = (
(metrics["average_latency"] * (metrics["successful_requests"] - 1) + latency)
/ metrics["successful_requests"]
)
metrics["max_latency"] = max(metrics["max_latency"], latency)
metrics["min_latency"] = min(metrics["min_latency"], latency)
# Save metrics after update
await self._save_metrics()
async def _get_current_metrics(self) -> Dict[str, Any]:
"""Get current performance metrics."""
metrics = self._performance_metrics.copy()
_LOGGER.debug(f"Current performance metrics: {metrics}")
return metrics
async def _handle_error(self, error: Exception) -> None:
"""Enhanced error handling with metric persistence."""
self._performance_metrics["total_errors"] += 1
self._performance_metrics["failed_requests"] += 1
# Save metrics after error
await self._save_metrics()
error_details = {
"timestamp": dt_util.utcnow().isoformat(),
"model": self.model,
"instance": self.instance_name,
"error_message": str(error),
"error_type": type(error).__name__,
"traceback": traceback.format_exc() if _LOGGER.isEnabledFor(logging.DEBUG) else None,
}
# Specific error type handling
error_mapping = {
HomeAssistantError: {"is_ha_error": True},
ConnectionError: {
"is_connection_error": True,
"is_rate_limited": True
},
TimeoutError: {"is_timeout": True},
PermissionError: {"is_permission_denied": True},
ValueError: {"is_validation_error": True}
}
for error_type, error_flags in error_mapping.items():
if isinstance(error, error_type):
error_details.update(error_flags)
break
# Update system state based on error type
if error_details.get("is_rate_limited"):
self._is_rate_limited = True
_LOGGER.warning(f"Rate limit detected for {self.instance_name}")
if error_details.get("is_connection_error"):
self.endpoint_status = "unavailable"
self.last_response = error_details
_LOGGER.error(f"AI Processing Error: {error_details}")
if _LOGGER.isEnabledFor(logging.DEBUG):
_LOGGER.debug(f"Full Error Traceback: {error_details['traceback']}")
async def async_initialize_history_file(self) -> None: async def async_initialize_history_file(self) -> None:
""" """
Asynchronously initialize history file. Asynchronously initialize history file.
Creates the file and writes an initialization timestamp Creates the file and writes an initialization timestamp
without blocking the event loop. without blocking the event loop using Home Assistant's executor.
""" """
try: try:
# Use asyncio.to_thread to perform file I/O in a separate thread await self.hass.async_add_executor_job(self._sync_initialize_history_file)
await asyncio.to_thread(self._sync_initialize_history_file)
except Exception as e: except Exception as e:
_LOGGER.error(f"Could not initialize history file: {e}") _LOGGER.error(f"Could not initialize history file: {e}")
_LOGGER.debug(traceback.format_exc()) _LOGGER.debug(traceback.format_exc())
def _sync_initialize_history_file(self) -> None: async def _sync_initialize_history_file(self) -> None:
"""
Synchronous method to create and initialize history file.
Runs in a separate thread to avoid blocking the event loop.
"""
try: try:
with open(self._history_file, 'a') as f: if not await self._file_exists(self._history_file):
f.write(f"History initialized at: {dt_util.utcnow().isoformat()}\n") def write_empty_history():
with open(self._history_file, 'w') as f:
json.dump([], f)
await self.hass.async_add_executor_job(write_empty_history)
except Exception as e: except Exception as e:
_LOGGER.error(f"Synchronous history file initialization failed: {e}") _LOGGER.error(f"History file initialization failed: {e}")
# Add size check to _update_history method # Size check to _update_history method
async def _update_history(self, question: str, response: dict) -> None: async def _update_history(self, question: str, response: dict) -> None:
"""Update conversation history with size limits.""" """Update conversation history with size validation."""
try: try:
# Limit entry size
question = question[:MAX_ENTRY_SIZE]
response_content = response.get("content", "")[:MAX_ENTRY_SIZE]
history_entry = { history_entry = {
"timestamp": dt_util.utcnow().isoformat(), "timestamp": dt_util.utcnow().isoformat(),
"question": question, "question": self._truncate_text(question, MAX_ATTRIBUTE_SIZE),
"response": response_content, "response": self._truncate_text(
response.get("content", ""),
MAX_ATTRIBUTE_SIZE
),
} }
# Check approximate entry size
entry_size = len(json.dumps(history_entry).encode('utf-8'))
current_size = await self._check_file_size(self._history_file)
if current_size + entry_size > MAX_HISTORY_FILE_SIZE:
_LOGGER.warning(
f"History size limit approaching. "
f"Current: {current_size}, "
f"Entry: {entry_size}, "
f"Max: {MAX_HISTORY_FILE_SIZE}"
)
await self._rotate_history()
self._conversation_history.append(history_entry) self._conversation_history.append(history_entry)
# Enforce size limit # Enforce size limit
@@ -269,48 +460,93 @@ class HATextAICoordinator(DataUpdateCoordinator):
await self._write_history_entry(history_entry) await self._write_history_entry(history_entry)
await self._rotate_history()
except Exception as e: except Exception as e:
_LOGGER.error(f"Error updating history: {e}") _LOGGER.error(f"Error updating history: {e}")
_LOGGER.debug(traceback.format_exc()) _LOGGER.debug(traceback.format_exc())
# Update _write_history_entry method
async def _write_history_entry(self, entry: dict) -> None: async def _write_history_entry(self, entry: dict) -> None:
"""Write a single history entry to file asynchronously.""" """Write history entry with file size checks."""
try: try:
if not os.path.exists(self._history_dir): if not await self._file_exists(self._history_dir):
os.makedirs(self._history_dir, exist_ok=True) await self._create_history_dir()
_LOGGER.debug(f"Writing history entry to {self._history_file}") # Check current file size
current_size = 0
async with AsyncFileHandler(self._history_file) as f: if await self._file_exists(self._history_file):
await f.write( current_size = await self.hass.async_add_executor_job(
f"{entry['timestamp']}: " os.path.getsize,
f"Question: {entry['question']} - " self._history_file
f"Response: {entry['response']}\n"
) )
await f.flush()
_LOGGER.debug(f"Successfully wrote history entry") # Calculate approximate entry size
entry_size = len(json.dumps(entry).encode('utf-8'))
if current_size + entry_size > MAX_HISTORY_FILE_SIZE:
_LOGGER.warning(
f"History file size limit reached. Current: {current_size}, "
f"Entry size: {entry_size}, Max: {MAX_HISTORY_FILE_SIZE}"
)
# Trigger rotation before writing
await self._rotate_history()
# Continue with writing
history = []
if await self._file_exists(self._history_file):
async with AsyncFileHandler(self._history_file, 'r') as f:
content = await f.read()
if content:
history = json.loads(content)
history.append(entry)
if len(history) > self.max_history_size:
history = history[-self.max_history_size:]
async with AsyncFileHandler(self._history_file, 'w') as f:
await f.write(json.dumps(history, indent=2))
except Exception as e: except Exception as e:
_LOGGER.error(f"Error writing history entry: {e}") _LOGGER.error(f"Error writing history entry: {e}")
_LOGGER.debug(traceback.format_exc()) _LOGGER.debug(traceback.format_exc())
def _sync_write_history_entry(self, entry: dict) -> None: async def _check_file_size(self, file_path: str) -> int:
""" """
Synchronous method to write history entry. Check file size asynchronously.
Runs in a separate thread to avoid blocking the event loop. Args:
file_path: Path to file to check
Returns:
int: File size in bytes
""" """
try: try:
with open(self._history_file, 'a') as f: if await self._file_exists(file_path):
f.write( size = await self.hass.async_add_executor_job(os.path.getsize, file_path)
f"{entry['timestamp']}: " _LOGGER.debug(f"File size check for {file_path}: {size} bytes")
f"Question: {entry['question']} - " return size
f"Response: {entry['response']}\n" return 0
) except Exception as e:
_LOGGER.error(f"Error checking file size for {file_path}: {e}")
return 0
def _sync_write_history_entry(self, entry: dict) -> None:
"""Synchronous method to write history entry."""
try:
history = []
if os.path.exists(self._history_file):
with open(self._history_file, 'r') as f:
content = f.read()
if content:
history = json.loads(content)
history.append(entry)
if len(history) > self.max_history_size:
history = history[-self.max_history_size:]
with open(self._history_file, 'w') as f:
json.dump(history, f, indent=2)
except Exception as e: except Exception as e:
_LOGGER.error(f"Synchronous history entry writing failed: {e}") _LOGGER.error(f"Synchronous history entry writing failed: {e}")
@@ -318,9 +554,8 @@ class HATextAICoordinator(DataUpdateCoordinator):
"""Rotate conversation history with file management.""" """Rotate conversation history with file management."""
try: try:
_LOGGER.debug(f"Starting history rotation for {self._history_file}") _LOGGER.debug(f"Starting history rotation for {self._history_file}")
await asyncio.to_thread(self._sync_rotate_history) await self.hass.async_add_executor_job(self._rotate_history_files)
_LOGGER.debug(f"Completed history rotation") _LOGGER.debug(f"Completed history rotation")
except Exception as e: except Exception as e:
_LOGGER.error(f"Error rotating history: {e}") _LOGGER.error(f"Error rotating history: {e}")
_LOGGER.debug(traceback.format_exc()) _LOGGER.debug(traceback.format_exc())
@@ -332,7 +567,7 @@ class HATextAICoordinator(DataUpdateCoordinator):
try: try:
# Test write permission in a separate thread # Test write permission in a separate thread
test_file_path = os.path.join(self._history_dir, ".write_test") test_file_path = os.path.join(self._history_dir, ".write_test")
await asyncio.to_thread(self._sync_test_directory_write, test_file_path) await self.hass.async_add_executor_job(self._sync_test_directory_write, test_file_path)
except PermissionError: except PermissionError:
_LOGGER.error(f"No write permissions for history directory: {self._history_dir}") _LOGGER.error(f"No write permissions for history directory: {self._history_dir}")
@@ -350,92 +585,128 @@ class HATextAICoordinator(DataUpdateCoordinator):
except Exception as e: except Exception as e:
_LOGGER.error(f"Directory write test failed: {e}") _LOGGER.error(f"Directory write test failed: {e}")
def _sync_rotate_history(self) -> None: async def _rotate_history_files(self) -> None:
""" """Rotate history files with size validation."""
Synchronous method to rotate history files.
Runs in a separate thread to avoid blocking the event loop.
"""
try: try:
# Check and manage file size if await self._file_exists(self._history_file):
if os.path.exists(self._history_file): current_size = await self._check_file_size(self._history_file)
file_size = os.path.getsize(self._history_file)
if file_size > self._max_history_file_size: if current_size > MAX_HISTORY_FILE_SIZE:
# Create timestamped archive _LOGGER.info(
f"Rotating history file. "
f"Current size: {current_size}, "
f"Max allowed: {MAX_HISTORY_FILE_SIZE}"
)
# Create archive filename with timestamp
archive_file = os.path.join( archive_file = os.path.join(
self._history_dir, self._history_dir,
f"{self.normalized_name}_history_{dt_util.utcnow().strftime('%Y%m%d_%H%M%S')}.txt" f"{self.normalized_name}_history_{dt_util.utcnow().strftime('%Y%m%d_%H%M%S')}.json"
) )
os.rename(self._history_file, archive_file)
# Optional: Log rotation event # Rotate files
_LOGGER.info(f"History file rotated: {archive_file}") await self.hass.async_add_executor_job(
os.rename,
# Ensure history doesn't exceed max size self._history_file,
if len(self._conversation_history) > self.max_history_size: archive_file
# Trim in-memory history
self._conversation_history = self._conversation_history[-self.max_history_size:]
# Write history entries
with open(self._history_file, 'a') as f:
for entry in self._conversation_history:
f.write(
f"{entry['timestamp']}: {entry['question']} - {entry['response']}\n"
) )
# Create new history file with recent entries
async with AsyncFileHandler(self._history_file, 'w') as f:
await f.write(json.dumps(
self._conversation_history[-self.max_history_size:],
indent=2
))
_LOGGER.info(f"History file rotated to: {archive_file}")
except Exception as e: except Exception as e:
_LOGGER.error(f"Synchronous history rotation failed: {e}") _LOGGER.error(f"History rotation failed: {e}")
_LOGGER.debug(traceback.format_exc()) _LOGGER.debug(traceback.format_exc())
async def _async_update_data(self) -> Dict[str, Any]: async def _async_update_data(self) -> Dict[str, Any]:
"""Update data via library.""" """Update coordinator data with improved error handling and performance."""
try: try:
async with asyncio.Semaphore(1):
current_state = self._get_current_state() current_state = self._get_current_state()
_LOGGER.debug(
f"Updating data for {self.instance_name}, current state: {current_state}"
)
limited_history = [] limited_history = await self._get_limited_history()
if self._conversation_history:
for entry in self._conversation_history[-3:]:
limited_entry = {
"timestamp": entry["timestamp"],
"question": entry["question"][:4096],
"response": entry["response"][:4096],
}
limited_history.append(limited_entry)
limited_last_response = self.last_response.copy() metrics = await self._get_current_metrics()
if "response" in limited_last_response:
limited_last_response["response"] = limited_last_response["response"][:4096] if metrics is None:
if "question" in limited_last_response: metrics = {}
limited_last_response["question"] = limited_last_response["question"][:4096]
data = { data = {
"state": current_state, "state": current_state,
"metrics": self._performance_metrics, "metrics": metrics,
"last_response": limited_last_response, "last_response": await self._get_sanitized_last_response(),
"is_processing": self._is_processing, "is_processing": self._is_processing,
"is_rate_limited": self._is_rate_limited, "is_rate_limited": self._is_rate_limited,
"is_maintenance": self._is_maintenance, "is_maintenance": self._is_maintenance,
"endpoint_status": self.endpoint_status, "endpoint_status": self.endpoint_status,
"uptime": (dt_util.utcnow() - self._start_time).total_seconds(), "uptime": self._calculate_uptime(),
"system_prompt": self._system_prompt[:4096] if self._system_prompt else None, "system_prompt": self._get_truncated_system_prompt(),
"history_size": len(self._conversation_history), "history_size": len(self._conversation_history),
"conversation_history": limited_history, "conversation_history": limited_history,
"normalized_name": self.normalized_name, "normalized_name": self.normalized_name,
} }
# Validate data await self._validate_update_data(data)
if not isinstance(data, dict):
raise ValueError("Invalid data format")
_LOGGER.debug(f"Updated data for {self.instance_name}: {data}")
return data return data
except asyncio.CancelledError:
_LOGGER.warning("Update was cancelled")
raise
except Exception as err: except Exception as err:
_LOGGER.error(f"Error updating data for {self.instance_name}: {err}") _LOGGER.error(f"Error updating data: {err}", exc_info=True)
return self._initial_state return self._get_safe_initial_state()
async def _get_limited_history(self) -> List[Dict[str, Any]]:
"""Get limited conversation history with size constraints."""
limited_history = []
if self._conversation_history:
for entry in self._conversation_history[-3:]:
limited_entry = {
"timestamp": entry["timestamp"],
"question": self._truncate_text(entry["question"]),
"response": self._truncate_text(entry["response"]),
}
limited_history.append(limited_entry)
return limited_history
def _truncate_text(self, text: str, max_length: int = MAX_ATTRIBUTE_SIZE) -> str:
"""Safely truncate text to maximum length."""
return text[:max_length] if text else ""
async def _get_sanitized_last_response(self) -> Dict[str, Any]:
"""Get sanitized version of last response."""
response = self.last_response.copy()
if "response" in response:
response["response"] = self._truncate_text(response["response"])
if "question" in response:
response["question"] = self._truncate_text(response["question"])
return response
def _calculate_uptime(self) -> float:
"""Calculate current uptime in seconds."""
return (dt_util.utcnow() - self._start_time).total_seconds()
def _get_truncated_system_prompt(self) -> Optional[str]:
"""Get truncated system prompt."""
return self._truncate_text(self._system_prompt, 4096) if self._system_prompt else None
async def _validate_update_data(self, data: Dict[str, Any]) -> None:
"""Validate update data structure and content."""
required_keys = ["state", "metrics", "last_response"]
for key in required_keys:
if key not in data:
raise ValueError(f"Missing required key: {key}")
if not isinstance(data["metrics"], dict):
raise ValueError("Invalid metrics format")
async def async_update_ha_state(self) -> None: async def async_update_ha_state(self) -> None:
"""Update Home Assistant state.""" """Update Home Assistant state."""
@@ -466,66 +737,61 @@ class HATextAICoordinator(DataUpdateCoordinator):
return STATE_ERROR return STATE_ERROR
return STATE_READY return STATE_READY
def _calculate_context_tokens(self, messages: List[Dict[str, str]], model: str = None) -> int: def _get_safe_initial_state(self) -> Dict[str, Any]:
""" """Return safe initial state when update fails."""
Estimate tokens for conversation context. return {
"state": STATE_ERROR,
"metrics": {},
"last_response": self.last_response,
"is_processing": False,
"is_rate_limited": False,
"is_maintenance": False,
"endpoint_status": "error",
"uptime": self._calculate_uptime(),
"system_prompt": None,
"history_size": 0,
"conversation_history": [],
"normalized_name": self.normalized_name,
}
Args: def _calculate_context_tokens(self, messages: List[Dict[str, str]]) -> int:
messages: List of message dictionaries total_tokens = 0
model: Optional model name for provider-specific estimation
# Compile regular expressions for performance
number_pattern = re.compile(r'[0-9]')
special_char_pattern = re.compile(r'[^\w\s]')
whitespace_pattern = re.compile(r'\s+')
Returns:
Estimated number of tokens
"""
try: try:
# Anthropic specific token counting for message in messages:
if self.is_anthropic and hasattr(self.client, 'count_tokens'): text = message.get('content', '')
return sum(self.client.count_tokens(msg['content']) for msg in messages) if text:
# Normalize whitespace
text = whitespace_pattern.sub(' ', text.strip())
def estimate_tokens(text: str) -> int: # Advanced token estimation heuristics
""" words = text.split()
Flexible token estimation algorithm. for word in words:
# Complexity-based token calculation
if len(word) > 8: # Long words
total_tokens += 2
elif number_pattern.search(word): # Words with numbers
total_tokens += 1.5
elif special_char_pattern.search(word): # Words with special characters
total_tokens += 1.5
elif word.isupper(): # Acronyms and technical terms
total_tokens += 1.5
else: # Regular words
total_tokens += 1
Heuristics: # Additional correction for technical texts
- Count words total_tokens = math.ceil(total_tokens * 1.2)
- Estimate special characters
- Fallback to character-based estimation
"""
# Word-based estimation
words = len(text.split())
# Special character handling return int(total_tokens)
special_chars = sum(1 for char in text if not char.isalnum())
# Character-based fallback
char_tokens = len(text) // 4
# Combine estimations with bias towards words
total_tokens = (words * 1.5) + (special_chars * 0.5) + char_tokens
return max(int(total_tokens), words)
# Calculate total tokens across all messages
total_tokens = sum(estimate_tokens(msg['content']) for msg in messages)
# Logging for debugging
_LOGGER.debug(
f"Token Estimation: "
f"Messages: {len(messages)}, "
f"Estimated Tokens: {total_tokens}"
)
return total_tokens
except Exception as e: except Exception as e:
# Safe fallback with detailed logging _LOGGER.error(f"Token calculation error: {e}. Processed {len(messages)} messages.")
_LOGGER.warning( # Safe fallback with minimal token estimation
f"Token estimation failed. "
f"Error: {e}. "
f"Using conservative estimation."
)
# Conservative token estimation
return len(messages) * 100 return len(messages) * 100
async def async_ask_question( async def async_ask_question(
@@ -571,6 +837,9 @@ class HATextAICoordinator(DataUpdateCoordinator):
""" """
Enhanced question processing with intelligent token management. Enhanced question processing with intelligent token management.
""" """
if self.client is None:
raise HomeAssistantError("AI client not initialized")
try: try:
self._is_processing = True self._is_processing = True
await self.async_update_ha_state() await self.async_update_ha_state()
@@ -650,7 +919,7 @@ class HATextAICoordinator(DataUpdateCoordinator):
# Update metrics # Update metrics
end_time = dt_util.utcnow() end_time = dt_util.utcnow()
latency = (end_time - start_time).total_seconds() latency = (end_time - start_time).total_seconds()
self._update_metrics(latency, response) await self._update_metrics(latency, response)
await self._update_history(question, response) await self._update_history(question, response)
@@ -667,6 +936,7 @@ class HATextAICoordinator(DataUpdateCoordinator):
async def async_process_message(self, question: str, **kwargs) -> dict: async def async_process_message(self, question: str, **kwargs) -> dict:
"""Process message using the AI client.""" """Process message using the AI client."""
try: try:
async with asyncio.timeout(60): # 60 second timeout
if self.is_anthropic: if self.is_anthropic:
response = await self._process_anthropic_message(question, **kwargs) response = await self._process_anthropic_message(question, **kwargs)
else: else:
@@ -678,11 +948,15 @@ class HATextAICoordinator(DataUpdateCoordinator):
"response": response["content"], "response": response["content"],
"model": kwargs.get("model", self.model), "model": kwargs.get("model", self.model),
"instance": self.instance_name, "instance": self.instance_name,
"normalized_name": self.normalized_name,
"error": None, "error": None,
} }
return response return response
except asyncio.TimeoutError:
raise HomeAssistantError("Request timed out")
except Exception as err: except Exception as err:
self._handle_error(err) self._handle_error(err)
raise raise
@@ -732,73 +1006,73 @@ class HATextAICoordinator(DataUpdateCoordinator):
_LOGGER.error(f"Error in OpenAI API call: {str(e)}") _LOGGER.error(f"Error in OpenAI API call: {str(e)}")
raise raise
def _update_metrics(self, latency: float, response: dict) -> None: async def _migrate_history_from_txt_to_json(self) -> None:
"""Update performance metrics.""" try:
metrics = self._performance_metrics old_history_file = os.path.join(
tokens = response.get("tokens", {}) self._history_dir,
f"{self.normalized_name}_history.txt"
metrics["total_tokens"] += tokens.get("total", 0)
metrics["prompt_tokens"] += tokens.get("prompt", 0)
metrics["completion_tokens"] += tokens.get("completion", 0)
metrics["successful_requests"] += 1
metrics["average_latency"] = (
(metrics["average_latency"] * (metrics["successful_requests"] - 1) + latency)
/ metrics["successful_requests"]
) )
metrics["max_latency"] = max(metrics["max_latency"], latency)
metrics["min_latency"] = min(metrics["min_latency"], latency)
def _handle_error(self, error: Exception) -> None: if not await self._file_exists(old_history_file):
""" return
Enhanced error handling with comprehensive diagnostics.
Captures detailed error information, tracks error metrics, _LOGGER.info(f"Found old history file for {self.instance_name}, starting migration to JSON format")
and provides context for troubleshooting AI processing issues.
"""
self._performance_metrics["total_errors"] += 1
self._performance_metrics["failed_requests"] += 1
error_details = { # Read old history
"timestamp": dt_util.utcnow().isoformat(), history_entries = []
"model": self.model, async with AsyncFileHandler(old_history_file, 'r') as f:
"instance": self.instance_name, content = await f.read()
"error_message": str(error),
"error_type": type(error).__name__,
"traceback": traceback.format_exc() if _LOGGER.isEnabledFor(logging.DEBUG) else None,
}
# Specific error type handling # Parse txt content
error_mapping = { for line in content.split('\n'):
HomeAssistantError: {"is_ha_error": True}, if not line or line.startswith("History initialized at:"):
ConnectionError: { continue
"is_connection_error": True,
"is_rate_limited": True
},
TimeoutError: {"is_timeout": True},
PermissionError: {"is_permission_denied": True},
ValueError: {"is_validation_error": True}
}
for error_type, error_flags in error_mapping.items(): try:
if isinstance(error, error_type): # Parse the old format: "timestamp: Question: question - Response: response"
error_details.update(error_flags) parts = line.split(': ', 1)
break if len(parts) != 2:
continue
# Update system state based on error type timestamp = parts[0]
if error_details.get("is_rate_limited"): content_parts = parts[1].split(' - ')
self._is_rate_limited = True if len(content_parts) != 2:
_LOGGER.warning(f"Rate limit detected for {self.instance_name}") continue
if error_details.get("is_connection_error"): question = content_parts[0].replace('Question: ', '')
self.endpoint_status = "unavailable" response = content_parts[1].replace('Response: ', '')
self.last_response = error_details history_entries.append({
_LOGGER.error(f"AI Processing Error: {error_details}") "timestamp": timestamp,
"question": question,
"response": response
})
# Optional: Add more sophisticated error tracking or notification logic except Exception as e:
if _LOGGER.isEnabledFor(logging.DEBUG): _LOGGER.warning(f"Error parsing history line: {line}. Error: {e}")
_LOGGER.debug(f"Full Error Traceback: {error_details['traceback']}") continue
if history_entries:
# Write to new JSON file
async with AsyncFileHandler(self._history_file, 'w') as f:
await f.write(json.dumps(history_entries, indent=2))
# Backup old file
backup_file = old_history_file + '.backup'
os.rename(old_history_file, backup_file)
_LOGGER.info(
f"Successfully migrated {len(history_entries)} entries "
f"from txt to JSON for {self.instance_name}. "
f"Old file backed up as {backup_file}"
)
# Update in-memory history
self._conversation_history = history_entries
except Exception as e:
_LOGGER.error(f"Error during history migration for {self.instance_name}: {e}")
_LOGGER.debug(traceback.format_exc())
async def async_clear_history(self) -> None: async def async_clear_history(self) -> None:
""" """
@@ -807,16 +1081,11 @@ class HATextAICoordinator(DataUpdateCoordinator):
Removes in-memory history and deletes history file. Removes in-memory history and deletes history file.
""" """
try: try:
# Clear in-memory history
self._conversation_history = [] self._conversation_history = []
if await self._file_exists(self._history_file):
# Safely remove history file await self.hass.async_add_executor_job(os.remove, self._history_file)
if os.path.exists(self._history_file):
await asyncio.to_thread(os.remove, self._history_file)
await self.async_update_ha_state() await self.async_update_ha_state()
_LOGGER.info(f"History for {self.instance_name} cleared successfully") _LOGGER.info(f"History for {self.instance_name} cleared successfully")
except Exception as e: except Exception as e:
_LOGGER.error(f"Error clearing history: {e}") _LOGGER.error(f"Error clearing history: {e}")
_LOGGER.debug(traceback.format_exc()) _LOGGER.debug(traceback.format_exc())
Binary file not shown.

After

Width:  |  Height:  |  Size: 121 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 351 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 87 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 257 KiB

Binary file not shown.

Before

Width:  |  Height:  |  Size: 16 KiB

After

Width:  |  Height:  |  Size: 117 KiB

Binary file not shown.

Before

Width:  |  Height:  |  Size: 45 KiB

After

Width:  |  Height:  |  Size: 325 KiB

Binary file not shown.

Before

Width:  |  Height:  |  Size: 12 KiB

After

Width:  |  Height:  |  Size: 86 KiB

Binary file not shown.

Before

Width:  |  Height:  |  Size: 34 KiB

After

Width:  |  Height:  |  Size: 259 KiB

+1 -1
View File
@@ -23,6 +23,6 @@
"single_config_entry": false, "single_config_entry": false,
"ssdp": [], "ssdp": [],
"usb": [], "usb": [],
"version": "2.0.5-beta", "version": "2.0.6-beta",
"zeroconf": [] "zeroconf": []
} }
+62 -41
View File
@@ -67,6 +67,7 @@ from .const import (
ENTITY_ICON_PROCESSING, ENTITY_ICON_PROCESSING,
DEFAULT_NAME_PREFIX, DEFAULT_NAME_PREFIX,
CONF_MAX_HISTORY_SIZE, CONF_MAX_HISTORY_SIZE,
MAX_ATTRIBUTE_SIZE,
) )
from .coordinator import HATextAICoordinator from .coordinator import HATextAICoordinator
@@ -178,12 +179,29 @@ class HATextAISensor(CoordinatorEntity, SensorEntity):
def _sanitize_attributes(self, attributes: Dict[str, Any]) -> Dict[str, Any]: def _sanitize_attributes(self, attributes: Dict[str, Any]) -> Dict[str, Any]:
"""Sanitize all attributes for JSON serialization.""" """Sanitize all attributes for JSON serialization."""
return { sanitized = {
key: self._sanitize_value(value) key: self._sanitize_value(value)
for key, value in attributes.items() for key, value in attributes.items()
if value is not None if value is not None
} }
# Log metrics for debugging
metrics_keys = [
METRIC_TOTAL_TOKENS,
METRIC_PROMPT_TOKENS,
METRIC_COMPLETION_TOKENS,
METRIC_SUCCESSFUL_REQUESTS,
METRIC_FAILED_REQUESTS,
METRIC_AVERAGE_LATENCY,
METRIC_MAX_LATENCY,
METRIC_MIN_LATENCY,
]
metrics_values = {k: sanitized.get(k) for k in metrics_keys if k in sanitized}
_LOGGER.debug(f"Metrics for {self.entity_id}: {metrics_values}")
return sanitized
@property @property
def native_value(self) -> StateType: def native_value(self) -> StateType:
"""Return the native value of the sensor.""" """Return the native value of the sensor."""
@@ -212,68 +230,65 @@ class HATextAISensor(CoordinatorEntity, SensorEntity):
try: try:
data = self.coordinator.data data = self.coordinator.data
metrics = data.get("metrics", {})
# Base attributes
attributes = { attributes = {
ATTR_MODEL: self._config_entry.data.get(CONF_MODEL, "Unknown"), ATTR_MODEL: self._config_entry.data.get(CONF_MODEL, "Unknown"),
ATTR_API_PROVIDER: self._config_entry.data.get( ATTR_API_PROVIDER: self._config_entry.data.get(CONF_API_PROVIDER, "Unknown"),
CONF_API_PROVIDER, "Unknown"
),
ATTR_API_STATUS: self._current_state, ATTR_API_STATUS: self._current_state,
ATTR_TOTAL_ERRORS: self._error_count, ATTR_TOTAL_ERRORS: metrics.get("total_errors", 0),
ATTR_LAST_ERROR: self._last_error,
"instance_name": self._instance_name, "instance_name": self._instance_name,
"normalized_name": self._normalized_name, "normalized_name": self._normalized_name,
ATTR_SYSTEM_PROMPT: data.get("system_prompt"), ATTR_SYSTEM_PROMPT: (data.get("system_prompt", "")[:MAX_ATTRIBUTE_SIZE]
if data.get("system_prompt") else None),
ATTR_IS_PROCESSING: data.get("is_processing", False), ATTR_IS_PROCESSING: data.get("is_processing", False),
ATTR_IS_RATE_LIMITED: data.get("is_rate_limited", False), ATTR_IS_RATE_LIMITED: data.get("is_rate_limited", False),
ATTR_IS_MAINTENANCE: data.get("is_maintenance", False), ATTR_IS_MAINTENANCE: data.get("is_maintenance", False),
ATTR_ENDPOINT_STATUS: data.get("endpoint_status", "unknown"), ATTR_ENDPOINT_STATUS: data.get("endpoint_status", "unknown"),
ATTR_UPTIME: data.get("uptime", 0), ATTR_UPTIME: round(data.get("uptime", 0), 2),
ATTR_HISTORY_SIZE: data.get("history_size", 0), ATTR_HISTORY_SIZE: data.get("history_size", 0),
ATTR_CONVERSATION_HISTORY: data.get("conversation_history", []),
} }
# Add metrics # History limit
metrics = data.get("metrics", {}) conversation_history = data.get("conversation_history", [])
if conversation_history:
limited_history = []
for entry in conversation_history:
limited_entry = {
"timestamp": entry["timestamp"],
"question": entry["question"][:MAX_ATTRIBUTE_SIZE],
"response": entry["response"][:MAX_ATTRIBUTE_SIZE]
}
limited_history.append(limited_entry)
attributes[ATTR_CONVERSATION_HISTORY] = limited_history
# Metrics
if isinstance(metrics, dict): if isinstance(metrics, dict):
self._metrics = metrics attributes.update({
attributes.update(
{
METRIC_TOTAL_TOKENS: metrics.get("total_tokens", 0), METRIC_TOTAL_TOKENS: metrics.get("total_tokens", 0),
METRIC_PROMPT_TOKENS: metrics.get("prompt_tokens", 0), METRIC_PROMPT_TOKENS: metrics.get("prompt_tokens", 0),
METRIC_COMPLETION_TOKENS: metrics.get("completion_tokens", 0), METRIC_COMPLETION_TOKENS: metrics.get("completion_tokens", 0),
METRIC_SUCCESSFUL_REQUESTS: metrics.get( METRIC_SUCCESSFUL_REQUESTS: metrics.get("successful_requests", 0),
"successful_requests", 0
),
METRIC_FAILED_REQUESTS: metrics.get("failed_requests", 0), METRIC_FAILED_REQUESTS: metrics.get("failed_requests", 0),
METRIC_AVERAGE_LATENCY: metrics.get("average_latency", 0), METRIC_AVERAGE_LATENCY: round(metrics.get("average_latency", 0), 2),
METRIC_MAX_LATENCY: metrics.get("max_latency", 0), METRIC_MAX_LATENCY: round(metrics.get("max_latency", 0), 2),
METRIC_MIN_LATENCY: metrics.get("min_latency", float("inf")), METRIC_MIN_LATENCY: (metrics.get("min_latency")
} if metrics.get("min_latency") != float("inf")
) else None),
})
# Add last response # Last response handling
last_response = data.get("last_response", {}) last_response = data.get("last_response", {})
if isinstance(last_response, dict): if isinstance(last_response, dict):
self._last_response = last_response attributes.update({
attributes.update( ATTR_RESPONSE: last_response.get("response", "")[:MAX_ATTRIBUTE_SIZE],
{ ATTR_QUESTION: last_response.get("question", "")[:MAX_ATTRIBUTE_SIZE],
ATTR_RESPONSE: last_response.get("response", ""),
ATTR_QUESTION: last_response.get("question", ""),
"last_model": last_response.get("model", ""), "last_model": last_response.get("model", ""),
"last_timestamp": last_response.get("timestamp", ""), "last_timestamp": last_response.get("timestamp", ""),
"last_error": last_response.get("error"), "last_error": (last_response.get("error", "")[:MAX_ATTRIBUTE_SIZE]
} if last_response.get("error") else None),
) })
# Add performance metrics if available
if ATTR_PERFORMANCE_METRICS in data:
attributes[ATTR_PERFORMANCE_METRICS] = data[
ATTR_PERFORMANCE_METRICS
]
# Add API version if available
if ATTR_API_VERSION in data:
attributes[ATTR_API_VERSION] = data[ATTR_API_VERSION]
return self._sanitize_attributes(attributes) return self._sanitize_attributes(attributes)
@@ -299,6 +314,12 @@ class HATextAISensor(CoordinatorEntity, SensorEntity):
self._is_processing = data.get("is_processing", False) self._is_processing = data.get("is_processing", False)
# Update metrics
metrics = data.get("metrics", {})
if isinstance(metrics, dict):
self._metrics.update(metrics)
_LOGGER.debug(f"Updated metrics for {self.entity_id}: {self._metrics}")
# Update conversation history and system prompt # Update conversation history and system prompt
self._conversation_history = data.get("conversation_history", []) self._conversation_history = data.get("conversation_history", [])
self._system_prompt = data.get("system_prompt") self._system_prompt = data.get("system_prompt")
Binary file not shown.

After

Width:  |  Height:  |  Size: 339 KiB

+29 -22
View File
@@ -1,25 +1,32 @@
``` ```
ha-text-ai/ ha-text-ai/
├── custom_components
├── custom_components/ │   └── ha_text_ai
├── ha_text_ai/ │   ├── __init__.py
├── __init__.py │   ├── api_client.py
├── config_flow.py │   ├── config_flow.py
├── coordinator.py │   ├── const.py
├── manifest.json │   ├── coordinator.py
├── sensor.py │   ├── icons
├── services.yaml │      ├── dark_icon.png
├── const.py │      ├── dark_icon@2x.png
└── api_client.py │      ├── dark_logo.png
│      ├── dark_logo@2x.png
├── translations/ │   │   ├── icon.png
│ ├── en.json │   │   ├── icon@2x.png
│ ├── de.json │   │   ├── logo.png
└── ru.json │      └── logo@2x.png
│   ├── manifest.json
── icons/ │   ── sensor.py
├── icon.png │   ├── services.yaml
├── icon@2x.png │   └── translations
├── logo.png │   ├── de.json
── logo@2x.png │   ── en.json
│   ├── es.json
│   ├── hi.json
│   ├── it.json
│   ├── ru.json
│   ├── sr.json
│   └── zh.json
``` ```