forked from BasedHardware/omi
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtools.py
More file actions
361 lines (301 loc) · 13.5 KB
/
Copy pathtools.py
File metadata and controls
361 lines (301 loc) · 13.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
"""
Platform tools router — exposes backend tools as REST endpoints for any client.
Unlike /v1/agent/execute-tool (which wraps LangChain tools for VM agents),
these endpoints are direct REST with proper HTTP semantics, designed for
desktop, web, and mobile agent clients.
Endpoints:
- GET /v1/tools/conversations — list conversations
- POST /v1/tools/conversations/search — semantic search conversations
- GET /v1/tools/memories — list memories/facts
- POST /v1/tools/memories/search — semantic search memories
- GET /v1/tools/action-items — list action items
- POST /v1/tools/action-items — create action item
- PATCH /v1/tools/action-items/{id} — update action item
- POST /v1/tools/calendar-events — create calendar event
"""
import logging
from datetime import datetime, timezone
from typing import Any, Optional
from urllib.parse import urlsplit
from fastapi import APIRouter, Depends, Query
from pydantic import BaseModel, Field, field_validator
import database.vector_db as vector_db
from utils.other.endpoints import get_current_user_uid, with_rate_limit
from utils.conversations.transcript_chunks import hydrate_chunk_texts
from utils.retrieval.safety import safe_isoformat
from utils.retrieval.tool_services.conversations import get_conversations_text, search_conversations_text
from utils.retrieval.tool_services.memories import get_memories_text, search_memories_text
from utils.retrieval.tool_result_boundaries import preserve_chat_memory_tool_result_boundary
from utils.retrieval.tool_services.action_items import (
get_action_items_text,
create_action_item_text,
update_action_item_text,
)
from utils.retrieval.tools.calendar_tools import create_calendar_event_tool
logger = logging.getLogger(__name__)
router = APIRouter()
# --------------- response envelope ---------------
class ToolSource(BaseModel):
kind: str = Field(max_length=32)
source_id: str = Field(max_length=512)
title: str = Field(default='', max_length=160)
preview: str = Field(default='', max_length=600)
created_at: Optional[str] = Field(default=None, max_length=80)
moment_timestamp_ms: Optional[int] = None
app_name: Optional[str] = Field(default=None, max_length=80)
url: Optional[str] = Field(default=None, max_length=2048)
@field_validator('url')
@classmethod
def require_http_url(cls, value: Optional[str]) -> Optional[str]:
if value is None:
return None
parsed = urlsplit(value)
if parsed.scheme.lower() not in {'http', 'https'} or not parsed.netloc:
raise ValueError('url must be an absolute HTTP(S) URL')
return value
class ToolResponse(BaseModel):
tool_name: str
result_text: str
is_error: bool = False
sources: list[ToolSource] = Field(default_factory=list)
def _ok(tool_name: str, text: str, sources: Optional[list[dict]] = None) -> dict:
is_error = text.startswith("Error")
return {
"tool_name": tool_name,
"result_text": text,
"is_error": is_error,
"sources": [] if is_error else (sources or []),
}
# --------------- request models ---------------
class SearchConversationsRequest(BaseModel):
query: str = Field(description="Natural-language topic, canonical conversation UUID, or h.omi.me share URL")
start_date: Optional[str] = Field(default=None, description="ISO date with timezone")
end_date: Optional[str] = Field(default=None, description="ISO date with timezone")
limit: int = Field(default=5, ge=1, le=20)
include_transcript: bool = Field(default=True)
class SearchMemoriesRequest(BaseModel):
query: str = Field(description="Semantic search query")
limit: int = Field(default=5, ge=1, le=20)
class CreateActionItemRequest(BaseModel):
description: str = Field(description="Action item description")
due_at: Optional[str] = Field(default=None, description="ISO date with timezone")
conversation_id: Optional[str] = Field(default=None, description="Source conversation ID")
class UpdateActionItemRequest(BaseModel):
completed: Optional[bool] = Field(default=None)
description: Optional[str] = Field(default=None)
due_at: Optional[str] = Field(default=None, description="ISO date with timezone")
class CreateCalendarEventRequest(BaseModel):
title: str = Field(description="Event title")
start_time: datetime = Field(description="ISO date/time with timezone")
end_time: datetime = Field(description="ISO date/time with timezone")
description: Optional[str] = Field(default=None, description="Event description")
location: Optional[str] = Field(default=None, description="Event location")
attendees: Optional[str] = Field(default=None, description="Comma-separated attendee names or email addresses")
@field_validator('start_time', 'end_time')
@classmethod
def require_timezone(cls, value: datetime) -> datetime:
if value.tzinfo is None or value.tzinfo.utcoffset(value) is None:
raise ValueError('datetime must include timezone')
return value
# --------------- conversation endpoints ---------------
@router.get("/v1/tools/conversations", response_model=ToolResponse)
def get_conversations(
start_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
end_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
limit: int = Query(default=20, ge=1, le=5000),
offset: int = Query(default=0, ge=0),
include_transcript: bool = Query(default=True),
uid: str = Depends(get_current_user_uid),
):
sources: list[dict] = []
result = get_conversations_text(
uid=uid,
start_date=start_date,
end_date=end_date,
limit=limit,
offset=offset,
include_transcript=include_transcript,
source_sink=sources,
)
return _ok("get_conversations", result, sources)
@router.post("/v1/tools/conversations/search", response_model=ToolResponse)
def search_conversations(
body: SearchConversationsRequest,
uid: str = Depends(with_rate_limit(get_current_user_uid, "tools:search")),
):
sources: list[dict] = []
result = search_conversations_text(
uid=uid,
query=body.query,
start_date=body.start_date,
end_date=body.end_date,
limit=body.limit,
include_transcript=body.include_transcript,
source_sink=sources,
)
return _ok("search_conversations", result, sources)
class SearchChunksRequest(BaseModel):
query: str = Field(description="Semantic search query")
limit: int = Field(default=20, ge=1, le=30)
def _transcript_chunk_source(row: dict[str, Any]) -> dict[str, Any]:
"""Typed source for one hydrated chunk row, shaped exactly like the sibling
conversation sources (`_append_conversation_source`): kind 'conversation' with
the PARENT conversation id, so chunk citations share the summary results' ref
namespace and no client citation validation changes. The preview is the
verbatim excerpt flattened to one line — it is quoted into client prompts, so
it must not be able to forge line-oriented prompt structure."""
created_at: Optional[str] = None
ts = row.get('created_at')
if isinstance(ts, (int, float)) and ts > 0:
created_at = datetime.fromtimestamp(ts, tz=timezone.utc).isoformat()
else:
created_at = safe_isoformat(row.get('conversation_started_at'))
return {
'kind': 'conversation',
'source_id': str(row['conversation_id']),
'title': str(row.get('conversation_title') or 'Conversation')[:160],
'preview': ' '.join(str(row.get('text') or '').split())[:600],
'created_at': created_at,
}
@router.post("/v1/tools/conversations/search-chunks", response_model=ToolResponse)
def search_conversation_chunks(
body: SearchChunksRequest,
uid: str = Depends(with_rate_limit(get_current_user_uid, "tools:search")),
):
"""Semantic search over RAW transcript chunks (verbatim evidence with dates).
Complements /conversations/search, which matches against conversation summaries:
summaries drop specifics (exact dates, names, numbers), so detail questions need
this verbatim layer. Returns chunks newest-relevant with their conversation date,
plus typed sources (one per parent conversation, best chunk first) so clients can
cite verbatim evidence the same way they cite summary results.
"""
rows = vector_db.search_transcript_chunks(uid, body.query, limit=body.limit)
rows = hydrate_chunk_texts(uid, rows)
if not rows:
return _ok("search_conversation_chunks", f"No transcript excerpts found matching '{body.query}'.")
parts = []
sources: list[dict] = []
seen_conversation_ids: set[str] = set()
for i, r in enumerate(rows, 1):
parts.append(f"Excerpt {i} (relevance: {r['score']:.2f}):\n{r['text']}")
conversation_id = r.get('conversation_id')
if not conversation_id or conversation_id in seen_conversation_ids:
continue
seen_conversation_ids.add(conversation_id)
sources.append(_transcript_chunk_source(r))
return _ok("search_conversation_chunks", "\n\n".join(parts), sources)
# --------------- memory endpoints ---------------
@router.get("/v1/tools/memories", response_model=ToolResponse)
def get_memories(
limit: int = Query(default=50, ge=1, le=5000),
offset: int = Query(default=0, ge=0),
start_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
end_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
uid: str = Depends(get_current_user_uid),
):
sources: list[dict] = []
result = get_memories_text(
uid=uid,
limit=limit,
offset=offset,
start_date=start_date,
end_date=end_date,
source_sink=sources,
)
bounded_result = preserve_chat_memory_tool_result_boundary('get_memories_tool', result)
if bounded_result != result:
sources = []
result = bounded_result
return _ok("get_memories", result, sources)
@router.post("/v1/tools/memories/search", response_model=ToolResponse)
def search_memories(
body: SearchMemoriesRequest,
uid: str = Depends(with_rate_limit(get_current_user_uid, "tools:search")),
):
sources: list[dict] = []
result = search_memories_text(
uid=uid,
query=body.query,
limit=body.limit,
source_sink=sources,
)
bounded_result = preserve_chat_memory_tool_result_boundary('search_memories_tool', result)
if bounded_result != result:
sources = []
result = bounded_result
return _ok("search_memories", result, sources)
# --------------- action item endpoints ---------------
@router.get("/v1/tools/action-items", response_model=ToolResponse)
def get_action_items(
limit: int = Query(default=50, ge=1, le=500),
offset: int = Query(default=0, ge=0),
completed: Optional[bool] = Query(default=None),
conversation_id: Optional[str] = Query(default=None),
start_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
end_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
due_start_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
due_end_date: Optional[str] = Query(default=None, description="ISO date with timezone"),
uid: str = Depends(get_current_user_uid),
):
sources: list[dict] = []
result = get_action_items_text(
uid=uid,
limit=limit,
offset=offset,
completed=completed,
conversation_id=conversation_id,
start_date=start_date,
end_date=end_date,
due_start_date=due_start_date,
due_end_date=due_end_date,
source_sink=sources,
)
return _ok("get_action_items", result, sources)
@router.post("/v1/tools/action-items", response_model=ToolResponse)
def create_action_item(
body: CreateActionItemRequest,
uid: str = Depends(with_rate_limit(get_current_user_uid, "tools:mutate")),
):
result = create_action_item_text(
uid=uid,
description=body.description,
due_at=body.due_at,
conversation_id=body.conversation_id,
)
return _ok("create_action_item", result)
@router.patch("/v1/tools/action-items/{action_item_id}", response_model=ToolResponse)
def update_action_item(
action_item_id: str,
body: UpdateActionItemRequest,
uid: str = Depends(with_rate_limit(get_current_user_uid, "tools:mutate")),
):
result = update_action_item_text(
uid=uid,
action_item_id=action_item_id,
completed=body.completed,
description=body.description,
due_at=body.due_at,
)
return _ok("update_action_item", result)
# --------------- calendar endpoints ---------------
@router.post("/v1/tools/calendar-events", response_model=ToolResponse)
async def create_calendar_event(
body: CreateCalendarEventRequest,
uid: str = Depends(with_rate_limit(get_current_user_uid, "tools:mutate")),
):
result = await create_calendar_event_tool.ainvoke(
{
"title": body.title,
"start_time": body.start_time.isoformat(),
"end_time": body.end_time.isoformat(),
"description": body.description,
"location": body.location,
"attendees": body.attendees,
},
config={"configurable": {"user_id": uid}},
)
return {
"tool_name": "create_calendar_event",
"result_text": result,
"is_error": not result.startswith("✅ Successfully created calendar event:"),
}