Skip to content

Commit e2d9041

Browse files
committed
black
1 parent f39886c commit e2d9041

63 files changed

Lines changed: 5604 additions & 4017 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

backends/advanced/scripts/create_plugin.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
Usage:
88
uv run python scripts/create_plugin.py my_awesome_plugin
99
"""
10+
1011
import argparse
1112
import os
1213
import shutil

backends/advanced/scripts/delete_all_conversations_api.py

Lines changed: 37 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,12 @@
44
This uses the proper API authentication and endpoints.
55
"""
66

7+
import argparse
78
import asyncio
89
import os
910
import sys
10-
import argparse
1111
from pathlib import Path
12+
1213
import aiohttp
1314
from dotenv import load_dotenv
1415

@@ -20,113 +21,111 @@ async def get_auth_token():
2021
"""Get admin authentication token."""
2122
admin_email = os.getenv("ADMIN_EMAIL")
2223
admin_password = os.getenv("ADMIN_PASSWORD")
23-
24+
2425
if not admin_email or not admin_password:
2526
print("Error: ADMIN_EMAIL and ADMIN_PASSWORD must be set in .env file")
2627
sys.exit(1)
27-
28+
2829
base_url = "http://localhost:8000"
29-
30+
3031
async with aiohttp.ClientSession() as session:
3132
# Login to get token
32-
login_data = {
33-
"username": admin_email,
34-
"password": admin_password
35-
}
36-
33+
login_data = {"username": admin_email, "password": admin_password}
34+
3735
async with session.post(
38-
f"{base_url}/auth/jwt/login",
39-
data=login_data
36+
f"{base_url}/auth/jwt/login", data=login_data
4037
) as response:
4138
if response.status != 200:
4239
print(f"Failed to login: {response.status}")
4340
text = await response.text()
4441
print(f"Response: {text}")
4542
sys.exit(1)
46-
43+
4744
result = await response.json()
4845
return result["access_token"]
4946

5047

5148
async def delete_all_conversations(skip_prompt=False):
5249
"""Delete all conversations using the API."""
53-
50+
5451
base_url = "http://localhost:8000"
55-
52+
5653
# Get auth token
5754
print("Getting admin authentication token...")
5855
token = await get_auth_token()
59-
60-
headers = {
61-
"Authorization": f"Bearer {token}"
62-
}
63-
56+
57+
headers = {"Authorization": f"Bearer {token}"}
58+
6459
async with aiohttp.ClientSession() as session:
6560
# First, get all conversations
6661
print("Fetching all conversations...")
6762
async with session.get(
68-
f"{base_url}/api/conversations",
69-
headers=headers
63+
f"{base_url}/api/conversations", headers=headers
7064
) as response:
7165
if response.status != 200:
7266
print(f"Failed to fetch conversations: {response.status}")
7367
text = await response.text()
7468
print(f"Response: {text}")
7569
return
76-
70+
7771
data = await response.json()
78-
72+
7973
# Extract conversations from nested structure
8074
conversations_dict = data.get("conversations", {})
8175
conversations = []
8276
for client_id, client_conversations in conversations_dict.items():
8377
conversations.extend(client_conversations)
84-
78+
8579
print(f"Found {len(conversations)} conversations")
86-
80+
8781
if len(conversations) == 0:
8882
print("No conversations to delete")
8983
return
90-
84+
9185
# Confirm deletion unless --yes flag is used
9286
if not skip_prompt:
93-
response = input(f"Are you sure you want to delete ALL {len(conversations)} conversations? (yes/no): ")
87+
response = input(
88+
f"Are you sure you want to delete ALL {len(conversations)} conversations? (yes/no): "
89+
)
9490
if response.lower() != "yes":
9591
print("Deletion cancelled")
9692
return
97-
93+
9894
# Delete each conversation
9995
deleted_count = 0
10096
failed_count = 0
101-
97+
10298
for conv in conversations:
10399
audio_uuid = conv.get("audio_uuid")
104100
if not audio_uuid:
105101
print(f"Skipping conversation without audio_uuid: {conv.get('_id')}")
106102
continue
107-
103+
108104
# Delete the conversation
109105
async with session.delete(
110-
f"{base_url}/api/conversations/{audio_uuid}",
111-
headers=headers
106+
f"{base_url}/api/conversations/{audio_uuid}", headers=headers
112107
) as response:
113108
if response.status == 200:
114109
deleted_count += 1
115-
print(f"Deleted conversation {audio_uuid} ({deleted_count}/{len(conversations)})")
110+
print(
111+
f"Deleted conversation {audio_uuid} ({deleted_count}/{len(conversations)})"
112+
)
116113
else:
117114
failed_count += 1
118115
text = await response.text()
119116
print(f"Failed to delete {audio_uuid}: {response.status} - {text}")
120-
117+
121118
print(f"\nDeletion complete:")
122119
print(f" Successfully deleted: {deleted_count}")
123120
print(f" Failed: {failed_count}")
124121

125122

126123
if __name__ == "__main__":
127124
# Parse command line arguments
128-
parser = argparse.ArgumentParser(description='Delete all conversations')
129-
parser.add_argument('--yes', '-y', action='store_true', help='Skip confirmation prompt')
125+
parser = argparse.ArgumentParser(description="Delete all conversations")
126+
parser.add_argument(
127+
"--yes", "-y", action="store_true", help="Skip confirmation prompt"
128+
)
130129
args = parser.parse_args()
131-
132-
asyncio.run(delete_all_conversations(skip_prompt=args.yes))
130+
131+
asyncio.run(delete_all_conversations(skip_prompt=args.yes))

backends/advanced/src/advanced_omi_backend/controllers/client_controller.py

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -6,10 +6,7 @@
66

77
from fastapi.responses import JSONResponse
88

9-
from advanced_omi_backend.client_manager import (
10-
ClientManager,
11-
get_user_clients_active,
12-
)
9+
from advanced_omi_backend.client_manager import ClientManager, get_user_clients_active
1310
from advanced_omi_backend.users import User
1411

1512
logger = logging.getLogger(__name__)
@@ -37,7 +34,9 @@ async def get_active_clients(user: User, client_manager: ClientManager):
3734

3835
# Filter to only the user's clients
3936
user_clients = [
40-
client for client in all_clients if client["client_id"] in user_active_clients
37+
client
38+
for client in all_clients
39+
if client["client_id"] in user_active_clients
4140
]
4241

4342
return {

backends/advanced/src/advanced_omi_backend/models/__init__.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,4 +7,4 @@
77

88
# Models can be imported directly from their files
99
# e.g. from .job import TranscriptionJob
10-
# e.g. from .conversation import Conversation, create_conversation
10+
# e.g. from .conversation import Conversation, create_conversation

backends/advanced/src/advanced_omi_backend/routers/modules/user_routes.py

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -24,13 +24,17 @@ async def get_users(current_user: User = Depends(current_superuser)):
2424

2525

2626
@router.post("")
27-
async def create_user(user_data: UserCreate, current_user: User = Depends(current_superuser)):
27+
async def create_user(
28+
user_data: UserCreate, current_user: User = Depends(current_superuser)
29+
):
2830
"""Create a new user. Admin only."""
2931
return await user_controller.create_user(user_data)
3032

3133

3234
@router.put("/{user_id}")
33-
async def update_user(user_id: str, user_data: UserUpdate, current_user: User = Depends(current_superuser)):
35+
async def update_user(
36+
user_id: str, user_data: UserUpdate, current_user: User = Depends(current_superuser)
37+
):
3438
"""Update a user. Admin only."""
3539
return await user_controller.update_user(user_id, user_data)
3640

@@ -43,4 +47,6 @@ async def delete_user(
4347
delete_memories: bool = False,
4448
):
4549
"""Delete a user and optionally their associated data. Admin only."""
46-
return await user_controller.delete_user(user_id, delete_conversations, delete_memories)
50+
return await user_controller.delete_user(
51+
user_id, delete_conversations, delete_memories
52+
)

backends/advanced/src/advanced_omi_backend/services/audio_service.py

Lines changed: 17 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -42,9 +42,9 @@ def __init__(self, redis_url: Optional[str] = None):
4242
self.memory_events_stream = "memory:events"
4343

4444
# Consumer group names (action verbs - what they DO)
45-
self.audio_writer = "audio-file-writer" # Writes audio chunks to WAV files
46-
self.memory_enqueuer = "memory-job-enqueuer" # Enqueues memory extraction jobs
47-
self.event_listener = "event-listener" # Listens for completion events
45+
self.audio_writer = "audio-file-writer" # Writes audio chunks to WAV files
46+
self.memory_enqueuer = "memory-job-enqueuer" # Enqueues memory extraction jobs
47+
self.event_listener = "event-listener" # Listens for completion events
4848

4949
async def connect(self):
5050
"""Connect to Redis with connection pooling."""
@@ -55,7 +55,7 @@ async def connect(self):
5555
max_connections=20, # Allow multiple concurrent operations
5656
socket_keepalive=True,
5757
socket_connect_timeout=5,
58-
retry_on_timeout=True
58+
retry_on_timeout=True,
5959
)
6060
logger.info(f"Audio stream service connected to Redis at {self.redis_url}")
6161

@@ -84,7 +84,7 @@ async def publish_audio_chunk(
8484
user_email: str,
8585
audio_chunk: AudioChunk,
8686
audio_uuid: Optional[str] = None,
87-
timestamp: Optional[int] = None
87+
timestamp: Optional[int] = None,
8888
) -> str:
8989
"""
9090
Publish audio chunk to Redis Stream.
@@ -135,12 +135,11 @@ async def publish_audio_chunk(
135135
# Ensure consumer group exists for this stream
136136
try:
137137
await self.redis.xgroup_create(
138-
stream_name,
139-
self.audio_writer,
140-
id="0",
141-
mkstream=True
138+
stream_name, self.audio_writer, id="0", mkstream=True
139+
)
140+
audio_logger.debug(
141+
f"Created consumer group {self.audio_writer} for {stream_name}"
142142
)
143-
audio_logger.debug(f"Created consumer group {self.audio_writer} for {stream_name}")
144143
except aioredis.ResponseError as e:
145144
if "BUSYGROUP" not in str(e):
146145
raise
@@ -152,7 +151,7 @@ async def publish_transcript_event(
152151
audio_uuid: str,
153152
conversation_id: str,
154153
status: str,
155-
error: Optional[str] = None
154+
error: Optional[str] = None,
156155
):
157156
"""
158157
Publish transcript completion event.
@@ -189,7 +188,7 @@ async def publish_transcript_event(
189188
self.transcript_events_stream,
190189
self.memory_enqueuer,
191190
id="0",
192-
mkstream=True
191+
mkstream=True,
193192
)
194193
except aioredis.ResponseError as e:
195194
if "BUSYGROUP" not in str(e):
@@ -200,7 +199,7 @@ async def publish_memory_event(
200199
conversation_id: str,
201200
status: str,
202201
memory_count: int = 0,
203-
error: Optional[str] = None
202+
error: Optional[str] = None,
204203
):
205204
"""
206205
Publish memory processing event.
@@ -234,21 +233,14 @@ async def publish_memory_event(
234233
# Ensure consumer group exists
235234
try:
236235
await self.redis.xgroup_create(
237-
self.memory_events_stream,
238-
self.event_listener,
239-
id="0",
240-
mkstream=True
236+
self.memory_events_stream, self.event_listener, id="0", mkstream=True
241237
)
242238
except aioredis.ResponseError as e:
243239
if "BUSYGROUP" not in str(e):
244240
raise
245241

246242
async def consume_audio_stream(
247-
self,
248-
consumer_name: str,
249-
callback,
250-
block_ms: int = 5000,
251-
count: int = 10
243+
self, consumer_name: str, callback, block_ms: int = 5000, count: int = 10
252244
):
253245
"""
254246
Consume audio chunks from all client streams.
@@ -288,7 +280,7 @@ async def consume_audio_stream(
288280
consumer_name,
289281
streams_dict,
290282
count=count,
291-
block=block_ms
283+
block=block_ms,
292284
)
293285

294286
for stream_name, stream_messages in messages:
@@ -299,15 +291,13 @@ async def consume_audio_stream(
299291

300292
# Acknowledge message
301293
await self.redis.xack(
302-
stream_name,
303-
self.audio_writer,
304-
message_id
294+
stream_name, self.audio_writer, message_id
305295
)
306296

307297
except Exception as e:
308298
logger.error(
309299
f"Error processing audio message {message_id.decode()}: {e}",
310-
exc_info=True
300+
exc_info=True,
311301
)
312302

313303
except Exception as e:

0 commit comments

Comments
 (0)