Merge into create

This commit is contained in:
Tat Dat Duong
2025-03-19 20:42:48 +01:00
parent 972ab1a935
commit e779c8e0b1
+59 -103
View File
@@ -839,6 +839,8 @@ class ThreadsClient:
metadata: Json = None,
thread_id: Optional[str] = None,
if_exists: Optional[OnConflictBehavior] = None,
from_supersteps: Optional[Sequence[dict[str, Sequence[dict[str, Any]]]]] = None,
graph_id: Optional[str] = None,
) -> Thread:
"""Create a new thread.
@@ -848,13 +850,16 @@ class ThreadsClient:
If None, ID will be a randomly generated UUID.
if_exists: How to handle duplicate creation. Defaults to 'raise' under the hood.
Must be either 'raise' (raise error if duplicate), or 'do_nothing' (return existing thread).
from_supersteps: Apply a list of supersteps when creating a thread, each containing a sequence of updates.
Each update has `values` or `command` and `as_node`. Used for copying a thread between deployments.
graph_id: Optional graph ID to associate with the thread.
Returns:
Thread: The created thread.
Example Usage:
thread = await client.threads.create(
thread = client.threads.create(
metadata={"number":1},
thread_id="my-thread-id",
if_exists="raise"
@@ -863,10 +868,32 @@ class ThreadsClient:
payload: Dict[str, Any] = {}
if thread_id:
payload["thread_id"] = thread_id
if metadata:
payload["metadata"] = metadata
if metadata or graph_id:
payload["metadata"] = {
**(metadata or {}),
**({"graph_id": graph_id} if graph_id else {}),
}
if if_exists:
payload["if_exists"] = if_exists
if from_supersteps:
payload["supersteps"] = {
"supersteps": [
{
"updates": [
{
"values": u["values"],
"command": u.get("command"),
"as_node": u["as_node"],
}
for u in s["updates"]
]
}
for s in from_supersteps
],
}
if payload.get("supersteps") is not None:
return self.http.post("/threads/state/bulk", json=payload)
return await self.http.post("/threads", json=payload)
async def update(self, thread_id: str, *, metadata: dict[str, Any]) -> Thread:
@@ -1142,55 +1169,6 @@ class ThreadsClient:
payload["as_node"] = as_node
return await self.http.post(f"/threads/{thread_id}/state", json=payload)
async def bulk_update_state(
self,
supersteps: Sequence[dict[str, Sequence[dict[str, Any]]]],
*,
graph_id: Optional[str] = None,
thread_id: Optional[str] = None,
metadata: Optional[dict[str, Any]] = None,
if_exists: Optional[OnConflictBehavior] = None,
) -> Thread:
"""Create a new thread from a batch of states.
Args:
supersteps: A sequence of supersteps, each containing a sequence of updates.
Each update has `values` or `command` and `as_node`.
graph_id: Optional graph ID to associate with the thread.
thread_id: Optional thread ID to use. If not provided, a new one will be generated.
metadata: Optional metadata to associate with the thread.
if_exists: Optional behavior when `thread_id` already exists.
Returns:
The created thread.
"""
payload: Dict[str, Any] = {
"supersteps": [
{
"updates": [
{
"values": u["values"],
"command": u.get("command"),
"as_node": u["as_node"],
}
for u in s["updates"]
]
}
for s in supersteps
],
}
if thread_id:
payload["thread_id"] = thread_id
if metadata or graph_id:
payload["metadata"] = {
**(metadata or {}),
**({"graph_id": graph_id} if graph_id else {}),
}
if if_exists:
payload["if_exists"] = if_exists
return await self.http.post("/threads/state/bulk", json=payload)
async def get_history(
self,
thread_id: str,
@@ -3085,6 +3063,8 @@ class SyncThreadsClient:
metadata: Json = None,
thread_id: Optional[str] = None,
if_exists: Optional[OnConflictBehavior] = None,
from_supersteps: Optional[Sequence[dict[str, Sequence[dict[str, Any]]]]] = None,
graph_id: Optional[str] = None,
) -> Thread:
"""Create a new thread.
@@ -3094,6 +3074,9 @@ class SyncThreadsClient:
If None, ID will be a randomly generated UUID.
if_exists: How to handle duplicate creation. Defaults to 'raise' under the hood.
Must be either 'raise' (raise error if duplicate), or 'do_nothing' (return existing thread).
from_supersteps: Apply a list of supersteps when creating a thread, each containing a sequence of updates.
Each update has `values` or `command` and `as_node`. Used for copying a thread between deployments.
graph_id: Optional graph ID to associate with the thread.
Returns:
Thread: The created thread.
@@ -3109,10 +3092,32 @@ class SyncThreadsClient:
payload: Dict[str, Any] = {}
if thread_id:
payload["thread_id"] = thread_id
if metadata:
payload["metadata"] = metadata
if metadata or graph_id:
payload["metadata"] = {
**(metadata or {}),
**({"graph_id": graph_id} if graph_id else {}),
}
if if_exists:
payload["if_exists"] = if_exists
if from_supersteps:
payload["supersteps"] = {
"supersteps": [
{
"updates": [
{
"values": u["values"],
"command": u.get("command"),
"as_node": u["as_node"],
}
for u in s["updates"]
]
}
for s in from_supersteps
],
}
if payload.get("supersteps") is not None:
return self.http.post("/threads/state/bulk", json=payload)
return self.http.post("/threads", json=payload)
def update(self, thread_id: str, *, metadata: dict[str, Any]) -> Thread:
@@ -3386,55 +3391,6 @@ class SyncThreadsClient:
payload["as_node"] = as_node
return self.http.post(f"/threads/{thread_id}/state", json=payload)
def bulk_update_state(
self,
supersteps: Sequence[dict[str, Sequence[dict[str, Any]]]],
*,
graph_id: Optional[str] = None,
thread_id: Optional[str] = None,
metadata: Optional[dict[str, Any]] = None,
if_exists: Optional[OnConflictBehavior] = None,
) -> Thread:
"""Create a new thread from a batch of states.
Args:
supersteps: A sequence of supersteps, each containing a sequence of updates.
Each update has `values` or `command` and `as_node`.
graph_id: Optional graph ID to associate with the thread.
thread_id: Optional thread ID to use. If not provided, a new one will be generated.
metadata: Optional metadata to associate with the thread.
if_exists: Optional behavior when `thread_id` already exists.
Returns:
The created thread.
"""
payload: Dict[str, Any] = {
"supersteps": [
{
"updates": [
{
"values": u["values"],
"command": u.get("command"),
"as_node": u["as_node"],
}
for u in s["updates"]
]
}
for s in supersteps
],
}
if thread_id:
payload["thread_id"] = thread_id
if metadata or graph_id:
payload["metadata"] = {
**(metadata or {}),
**({"graph_id": graph_id} if graph_id else {}),
}
if if_exists:
payload["if_exists"] = if_exists
return self.http.post("/threads/state/bulk", json=payload)
def get_history(
self,
thread_id: str,