์ด SDK๋ ProcessGPT ์์ด์ ํธ ์๋ฒ๋ฅผ ๋ง๋ค ๋ ํ์ํ ๊ณตํต ๊ธฐ๋ฅ์ ์ ๊ณตํฉ๋๋ค.
- DB์์ ์์ (todo) ํด๋ง โ ์ฒ๋ฆฌํ ์ผ๊ฐ ๊ฐ์ ธ์ค๊ธฐ
- ์ปจํ ์คํธ ์ค๋น (์ฌ์ฉ์ ์ ๋ณด, ํผ ์ ์, MCP ์ค์ ๋ฑ ์๋์ผ๋ก ์กฐํ)
- ๋ค์ํ ์์ด์ ํธ ์ค์ผ์คํธ๋ ์ด์ (A2A) ๊ณผ ํธํ
- ์ด๋ฒคํธ(Event) ์ ์ก ๊ท๊ฒฉ ํต์ผํ โ ๊ฒฐ๊ณผ๋ฅผ DB์ ์์ ํ๊ฒ ์ ์ฅ
- ์ฑํ (SSE) ์ ์ก ๊ณ์ธต โ ํํธ๋นํธ ยท ์ฌ์ ์(attach) ยท ์ค์ง(stop) ๋ฅผ ํ๋ ์์ํฌ๊ฐ ์ ๊ณต
- ํ
๋ํธ ์ธ์ฆ โ ์์ฒญ์ด ๋ณด๋ธ
tenant_id๋ฅผ ์์ฒญ์์ JWT ๋ก ๊ฒ์ฆ - ์์ด์ ํธ ์ฐ์ถ๋ฌผ โ ๋ง๋ ํ์ผ์ ๋น๊ณต๊ฐ ๋ฒํท์ ๋ณด๊ดํ๊ณ ๋ง๋ฃ๋๋ ์๋ช ์ฃผ์๋ก ๋ด์ค
๐ ์ฝ๊ฒ ๋งํ๋ฉด: ์ฌ๋ฌ ์ข ๋ฅ์ AI ์์ด์ ํธ๋ฅผ ๊ฐ์ ๊ท์น์ผ๋ก ์คํ/์ ์ฅ/ํธ์ถํ ์ ์๊ฒ ํด์ฃผ๋ ํตํฉ SDK ์ ๋๋ค.
0.8.0 ์์ ๋ฌ๋ผ์ง ์
์์ด์ ํธ๊ฐ ๋ง๋ ํ์ผ์ ์ฌ์ฉ์์๊ฒ ๋ด์ฃผ๋ ์ผ์ด SDK ๋ก ์ฌ๋ผ์์ต๋๋ค. ๊ทธ๋์ ์์ด์ ํธ๋ง๋ค ๋ฐ๋ก ๋ง๋ค์ด ์ฐ๋ ์์งยท๋ณด๊ดยท๋ณธ๋ฌธ ๋งํฌ ์นํ์processgpt_agent_sdk.artifactsํ๋๋ก ๋์ฒดํฉ๋๋ค. ์ฐ์ถ๋ฌผ์ ๋น๊ณต๊ฐ ๋ฒํท์ ๋ค์ด๊ฐ๊ณ ์ฃผ์๋ ํ ์๊ฐ์ง๋ฆฌ ์๋ช ์ฃผ์์ด๋ฉฐ, ๋ง๋ฃ๋๋ฉดfile_id๋ก ๋ค์ ๋ฐ๊ธ๋ฐ์ต๋๋ค. ์์ธํ ๋ด์ฉ์ 7 ์ ๋ณด์ธ์. 0.8.1 ๋ถํฐ๋ ์ฐ์ถ๋ฌผ์ด SSEdone์๋ ์ค๋ ค ๋๊ฐ๋๋ค(7.5).
0.5.0 ์์ ๋ฌ๋ผ์ง ์
์ฑํ SSE ์ ์ก ๊ณ์ธต๊ณผ ํ ๋ํธ ์ธ์ฆ์ด SDK ๋ก ์ฌ๋ผ์์ต๋๋ค. ๊ทธ๋์ ๊ฐ ์์ด์ ํธ ์ ์ฅ์๊ฐ ๋ฐ๋ก ๋ง๋ค์ด ์ฐ๋ ํํธ๋นํธยท์ฌ์ ์ยท์ค์งยท์ธ์ฆ์mount_chat_routes()ํ ๋ฒ์ผ๋ก ๋์ฒดํ ์ ์์ต๋๋ค. ์์ธํ ๋ด์ฉ์ 4.5 ยท 4.6 ์ ๋ณด์ธ์. ๊ธฐ์กดmount_chat_sse()๋จ๋ ํธ์ถ์ ๋์์ด ๊ทธ๋๋ก๋ผ ๊ณง๋ฐ๋ก ์ฌ๋ ค๋ ๊นจ์ง์ง ์์ต๋๋ค.
flowchart TD
subgraph DB["Postgres / Supabase"]
T["todolist"]:::db
E["events"]:::db
CH["chats"]:::db
end
subgraph SDK["SDK โ ํ๋ก์ธ์ค(ํด๋ง)"]
P["Polling<br/>(fetch_pending_task)"] --> C["Context ์ค๋น<br/>(fetch_context_bundle ๋ฑ)"]
C --> X["Executor<br/>(MinimalExecutor)"]
end
subgraph CHAT["SDK โ ์ฑํ
(SSE)"]
G["ํ
๋ํธ ๊ฐ๋<br/>(JWT ๊ฒ์ฆ)"] --> S["POST /chat/stream"]
R["๋ฐ ๋ ์ง์คํธ๋ฆฌ"] --> A["POST /chat/stream/attach"]
K["POST /chat/stop"]
ST["POST /chat/steer"] --> Q["์์ ์ง์ ๋๊ธฐ์ด"]
end
S --> X
K -->|์ทจ์| X
Q -->|์์ ์ง์ ์์ ๋ฐ์| X
X -->|TaskStatusUpdateEvent| E
X -->|TaskArtifactUpdateEvent| T
X -->|ํ ํฐยทdone| R
R --> CH
classDef db fill:#f2f2f2,stroke:#333,stroke-width:1px;
- todolist: ๊ฐ ์์ (Task)์ ์งํ ์ํ, ๊ฒฐ๊ณผ๋ฌผ ์ ์ฅ
- events: ์คํ ์ค๊ฐ์ ๋ฐ์ํ ์ด๋ฒคํธ ๋ก๊ทธ ์ ์ฅ
- chats: ์ฑํ ํด์ ์ต์ข ์๋ต ์ ์ฅ
- SDK๋ ์ธ ํ ์ด๋ธ์ ์๋์ผ๋ก ์ฐ๊ฒฐํด ์ค๋๋ค.
- ์ฑํ ๊ฒฝ๋ก๋ ์์ฒญ์ด Executor ์ ๋ฟ๊ธฐ ์ ์ ํ ๋ํธ๋ฅผ ๊ฒ์ฆํ๊ณ , ํด์ด ๋ด๋ณด๋ด๋ ์ด๋ฒคํธ๋ฅผ ๋ฐ ๋ ์ง์คํธ๋ฆฌ์ ๋จ๊ฒจ ์ฌ์ ์ยท์ค์ง๊ฐ ๊ฐ๋ฅํ๊ฒ ํฉ๋๋ค.
| A2A ํ์ | ์ค๋ช | ๋งค์นญ ํ ์ด๋ธ |
|---|---|---|
| TaskStatusUpdateEvent | ์์ ์ํ ์ ๋ฐ์ดํธ | events ํ
์ด๋ธ |
| TaskArtifactUpdateEvent | ์์ ๊ฒฐ๊ณผ๋ฌผ ์ ๋ฐ์ดํธ | todolist ํ
์ด๋ธ |
a2a-sdk v1.0๋ถํฐ A2A ์คํ(ProtoJSON) ์ ํฉ์ฑ์ ์ํด ๋ชจ๋ enum ๊ฐ์ด ๋๋ฌธ์ ์ค๋ค์ดํฌ ์ผ์ด์ค๋ก ํ์คํ๋์์ต๋๋ค.
-
TaskState
TaskState.submittedโTaskState.TASK_STATE_SUBMITTEDTaskState.workingโTaskState.TASK_STATE_WORKINGTaskState.completedโTaskState.TASK_STATE_COMPLETEDTaskState.failedโTaskState.TASK_STATE_FAILEDTaskState.canceledโTaskState.TASK_STATE_CANCELEDTaskState.input_requiredโTaskState.TASK_STATE_INPUT_REQUIREDTaskState.auth_requiredโTaskState.TASK_STATE_AUTH_REQUIREDTaskState.rejectedโTaskState.TASK_STATE_REJECTED- (์ถ๊ฐ)
TaskState.TASK_STATE_UNSPECIFIED
-
Role
Role.userโRole.ROLE_USERRole.agentโRole.ROLE_AGENT- (์ถ๊ฐ)
Role.ROLE_UNSPECIFIED
DB ์ events.event_type ์ปฌ๋ผ์ enum ์
๋๋ค. Executor ๊ฐ emitํ TaskStatusUpdateEvent ๊ฐ events ํ
์ด๋ธ์ ์ ์ฅ๋ ๋ ์ด๋ค enum ๊ฐ์ผ๋ก ๋ค์ด๊ฐ๋์ง๋ ๋ค์๊ณผ ๊ฐ์ด ๊ฒฐ์ ๋ฉ๋๋ค.
| event_type (DB enum) | ๋ฐํ ์ฃผ์ฒด | A2A ์ด๋ฒคํธ ํํ | ๋งคํ ๋ฐฉ์ |
|---|---|---|---|
task_started |
Executor | TaskStatusUpdateEvent(state=SUBMITTED) |
์๋ (state ๊ธฐ๋ฐ) |
task_completed |
Executor | TaskStatusUpdateEvent(state=COMPLETED) |
์๋ (state ๊ธฐ๋ฐ) |
error |
Executor | TaskStatusUpdateEvent(state=FAILED) |
์๋ (state ๊ธฐ๋ฐ) |
human_asked |
Executor | TaskStatusUpdateEvent(state=INPUT_REQUIRED) |
์๋ (state ๊ธฐ๋ฐ) |
task_working |
Executor | TaskStatusUpdateEvent(state=WORKING) + metadata["event_type"]="task_working" |
๋ช ์ |
tool_usage_started / tool_usage_finished |
Executor | TaskStatusUpdateEvent(state=WORKING) + metadata["event_type"]="tool_usage_*" |
๋ช ์ (sub-event, ์๋ ์ฐธ์กฐ) |
crew_completed |
SDK | (Executor ๊ฐ emit X) | ์๋ โ TaskArtifactUpdateEvent(last_chunk=True) ์ฒ๋ฆฌ ์์ ์ SDK ๊ฐ ๋ฐํ (์์ ๋ง: framework ์ task_done()) |
์๋ ๋งคํ ๊ท์น: SDK ๋ ๋ค์ lifecycle state ๋ฅผ ์๋์ผ๋ก enum ๊ฐ์ผ๋ก ๋งคํํฉ๋๋ค.
TASK_STATE_SUBMITTEDโtask_startedTASK_STATE_COMPLETEDโtask_completedTASK_STATE_FAILEDโerrorTASK_STATE_INPUT_REQUIREDโhuman_asked
TASK_STATE_WORKING์ ์๋์ ์ผ๋ก ์๋ ๋งคํ ๋์์ด ์๋๋๋ค. WORKING ์ ๋๋ฌด ๊ด๋ฒ์ํ๊ณ ๋๋ฉ์ธ sub-event(tool_usage_*๋ฑ) ์ ๋ฒ ์ด์ค๋ก๋ ์ฌ์ฌ์ฉ๋๋ฏ๋ก, sub-event ์๋ฏธ์ ์ถฉ๋ํ์ง ์๋๋ก NULL ๋ก ๋๊ฑฐ๋metadata["event_type"]์ผ๋ก ๋ช ์ํ์ธ์. metadata ๊ฐ ์์ผ๋ฉดevent_type์ปฌ๋ผ์ NULL ๋ก ์ ์ฅ๋ฉ๋๋ค(ํ์ฉ๋จ).๋ช ์ vs ์๋ ์ฐ์ ์์:
metadata["event_type"]๊ฐ ์์ผ๋ฉด ์๋ ๋งคํ๋ณด๋ค ์ฐ์ ํฉ๋๋ค (explicit > implicit). ์:state=WORKING + metadata["event_type"]="tool_usage_started"โtool_usage_started๋ก ์ ์ฅ.์ฌ๋์๊ฒ ๋ฌป๊ณ ๋๋ ์คํ (
INPUT_REQUIRED): ์คํ ์คINPUT_REQUIRED๋ฅผ ๋ธ ๋ค (WORKING/COMPLETED๋ก ์ด์ด์ง์ง ์๊ณ ) ๋๋๋ฉด SDK ๋ ๊ทธ ์คํ์ ์๋ฃ๋ก ๋ณด์ง ์์ต๋๋ค. ๋ค๋ฐ๋ฅด๋TaskArtifactUpdateEvent(last_chunk=True)๋ ์ง๋ฌธ ๋ณธ๋ฌธ์ผ๋ก ๋ณด๊ณ todolist ๊ฒฐ๊ณผ(output/draft)์ ์ ์ฅํ์ง ์์ผ๋ฉฐ,crew_completed๋ ๋ด์ง ์์ต๋๋ค. ๋์ ์์ ์draft_status='HUMAN_ASKED'๋ก ๋๊ณ ์ ์ (consumer, lease)๋ฅผ ํ๋๋ค โstatus๋IN_PROGRESS๊ทธ๋๋ก์ ๋๋ค. ์ฌ์ฉ์๊ฐ ๋ตํ๋ฉด ํ๋ฉด์ดFB_REQUESTED๋ก ๋ฐ๊พธ๊ณ ์์ปค๊ฐ ๋ค์ ์ง์ต๋๋ค. ์ํฐํฉํธ ์์ด ์ํ๋ง ๋ด๊ณ ๋๋๋ ๊ฐ์ต๋๋ค. ์ด ๊ท์น์ด ์์ผ๋ฉด COMPLETE ๋ชจ๋์์ ์ง๋ฌธ์ด ์ฐ์ถ๋ฌผ๋กSUBMITTED๋์ด ํ๋ก์ธ์ค๊ฐ ๋ค์ ๋จ๊ณ๋ก ๋์ด๊ฐ๋๋ค.
task_completedvsTaskArtifactUpdateEvent: ๋์ ๋ณ๊ฐ์ ๋๋ค.task_completed๋ events ํ ์ด๋ธ์ lifecycle ํ์์ด๊ณ , ์ค์ ๊ฒฐ๊ณผ๋ฌผ ์ ์ฅ์TaskArtifactUpdateEvent(last_chunk=True)๊ฐ todolist ํ ์ด๋ธ์ ์ํํฉ๋๋ค.
์์น: Executor ๋ A2A ํ์ค ์ด๋ฒคํธ์ ํ์ค ํ๋๋ง emit. SDK ๋ ๋งค์ง ๋ฉํ๋ฐ์ดํฐ ์์ด A2A ์ด๋ฒคํธ ํ์ ์์ฒด๋ฅผ ๋ผ์ฐํ ํค๋ก ์ฌ์ฉํฉ๋๋ค. ํํฐ๋ง์ Executor ์ฑ ์์ด๊ณ SDK ๋ ๋ฐ์ ๋๋ก ๋ผ์ฐํ ํฉ๋๋ค.
๋ผ์ฐํ ๋งคํธ๋ฆญ์ค:
| A2A ์ด๋ฒคํธ ํ์ | ChatEventQueue | ProcessEventQueue |
|---|---|---|
Task (๋ผ์ดํ์ฌ์ดํด ๋ง์ปค) |
silently ignore | silently ignore |
Message |
SSE {"type":"token","content":...} (๋ฐ๋ณต ํ์ฉ = ํ ํฐ ์คํธ๋ฆฌ๋ฐ) |
silently ignore |
TaskStatusUpdateEvent |
silently ignore | events ํ ์ด๋ธ ์ ์ฅ (state/text ๊ทธ๋๋ก) |
TaskArtifactUpdateEvent(last_chunk=True) |
SSE done + chats ์ ์ฅ |
todolist ์ ์ฅ (is_final=True) |
Executor ์ ํ์ค ํ๋ฆ (LLM ์คํธ๋ฆฌ๋ฐ ์์):
Task(state=SUBMITTED) โ ๋ผ์ดํ์ฌ์ดํด ์์- ํ ํฐ๋ง๋ค:
Message(text=token)โ ์ฑํ ์ฉTaskStatusUpdateEvent(state=WORKING, text=token)โ ํ๋ก์ธ์ค์ฉ (ํ์ ์ JSON payload)
TaskStatusUpdateEvent(state=COMPLETED, text=full)โ ์ข ๋ฃ ์๋ฆผ (์ ํ)TaskArtifactUpdateEvent(last_chunk=True, text=full)โ ์ต์ข ๊ฒฐ๊ณผ
๋ถํ ์ฐ๋ ค๊ฐ ์๋ค๋ฉด (์: events ํ ์ด๋ธ์ ํ ํฐ row ๊ฐ ๋๋ฌด ๋ง์ด ์์ผ ๊ฒฝ์ฐ) Executor ๊ฐ ์ง์ ํํฐ๋ง/์ง๊ณํ์ธ์. SDK ๋ ์ ์ฑ ์ ๊ฐ์ ํ์ง ์์ต๋๋ค.
tool_usage_started, tool_usage_finished ๊ฐ์ ๋๋ฉ์ธ sub-event ๋ A2A TaskState ์ ์ง์ ๋งคํ๋์ง ์์ต๋๋ค โ TaskState ๋ ์์
์ ์ฒด์ lifecycle (SUBMITTED โ WORKING โ COMPLETED/FAILED) ์ถ์ํ์ด๊ณ , ๋๊ตฌ ํธ์ถ์ ๊ทธ ์์์ ์ผ์ด๋๋ ์ธ๋ถ ์ฌ๊ฑด์
๋๋ค.
์ด๋ฐ sub-event ๋ state=TASK_STATE_WORKING ๊ทธ๋๋ก ๋๊ณ , metadata["event_type"] ๋ก enum ๊ฐ์ ๋ช
์ํ์ธ์. SDK ๊ฐ ๊ทธ ๊ฐ์ events.event_type ์ปฌ๋ผ์ ๊ทธ๋๋ก ๊ธฐ๋กํฉ๋๋ค. data ์ปฌ๋ผ์๋ text ๋ก ์ค์ JSON payload (๋๊ตฌ ์ด๋ฆ, ์ธ์, ๊ฒฐ๊ณผ ๋ฑ) ๊ฐ ์ ์ฅ๋ฉ๋๋ค.
import json
from a2a.helpers import new_text_status_update_event
from a2a.types import TaskState
# ๋๊ตฌ ํธ์ถ ์์
evt_start = new_text_status_update_event(
task_id=task_id, context_id=context_id,
state=TaskState.TASK_STATE_WORKING,
text=json.dumps(
{"tool": "web_search", "args": {"query": "process-gpt"}},
ensure_ascii=False,
),
)
evt_start.metadata.update({"event_type": "tool_usage_started"})
await event_queue.enqueue_event(evt_start)
# ... ๋๊ตฌ ์ค์ ํธ์ถ ...
# ๋๊ตฌ ํธ์ถ ์ข
๋ฃ
evt_end = new_text_status_update_event(
task_id=task_id, context_id=context_id,
state=TaskState.TASK_STATE_WORKING,
text=json.dumps(
{"tool": "web_search", "result_summary": "...", "elapsed_ms": 312},
ensure_ascii=False,
),
)
evt_end.metadata.update({"event_type": "tool_usage_finished"})
await event_queue.enqueue_event(evt_end)TaskState ๋ lifecycle, metadata ๋ ๋๋ฉ์ธ ๋ถ๋ฅ: A2A ํ์ค envelope ์์ ๋จธ๋ฌด๋ฅด๋ฉด์ ๋๋ฉ์ธ ์ด๋ฒคํธ๋ enum ์ผ๋ก ์ ํํ ๊ธฐ๋กํ ์ ์๋ ๋ฐฉ์์ ๋๋ค.
metadata์์ฒด๋ A2ATaskStatusUpdateEvent์ ํ์ค free-form ํ๋๋ผ "A2A ํ์ค๋ง ์ฌ์ฉ" ์์น๊ณผ ์ถฉ๋ํ์ง ์์ต๋๋ค.๋น๋ ์ฃผ์: tool_usage ๋ trace ์ฑ๊ฒฉ์ด๋ผ LLM ํ ๋ฒ์ ๋๊ตฌ 5๋ฒ ํธ์ถํ๋ฉด events row 10๊ฐ๊ฐ ์์ ๋๋ค. ๋๋ฌด ๋น๋ฒํ๋ฉด Executor ์ธก์์ sampling/aggregation ์ ์ ์ฉํ๊ฑฐ๋, ๊ฒฐ๊ณผ๋ง ํ ๋ฒ์ ๋ฌถ์ด emit ํ์ธ์.
tool_usage_*๋ ๋๊ตฌ ํธ์ถ ๋จ์์์๋ง emitํ๊ณ , ํ ํฐ ์คํธ๋ฆฌ๋ฐ ๋ฃจํ ์์๋ ์ ๋ ๋ฃ์ง ๋ง์ธ์.
์ด SDK๋ โํ๋์ ์์ ํ ์๋น์คโ๊ฐ ์๋๋ผ, ๋ด ์๋น์ค์ ๋ถ์ฌ์ ์ฌ์ฉํ๋ ํ๋ ์์ํฌ/๋ผ์ด๋ธ๋ฌ๋ฆฌ์ ๋๋ค.
์๋ ์์๋ ํ ํ๋ก์ธ์ค์์ ๋ค์์ ๋์์ ์ ๊ณตํฉ๋๋ค.
- ํ๋ก์ธ์ค(ํด๋ง):
await server.run()๋ก DB์์ todo๋ฅผ ๊ฐ์ ธ์ ์ฒ๋ฆฌ - ์ฑํ
(SSE):
/chat/stream์๋ํฌ์ธํธ๋ก ์์ฒญ์ ๋ฐ์Message-only๋ก ์๋ต +chats์ ์ ์ฅ - ์ฑํ
๋ถ๊ฐ ๋ผ์ฐํธ:
/chat/stream/attach(์ฌ์ ์) ยท/chat/stop(์ค์ง) ยท/chat/steer(์์ ์ง์)
import asyncio
import uvicorn
from starlette.applications import Starlette
from processgpt_agent_sdk import ProcessGPTAgentServer
from my_service.my_executor import MyExecutor
async def main():
server = ProcessGPTAgentServer(
agent_executor=MyExecutor(),
agent_type="langchain-react",
# ์์ฒญ ๋ณธ๋ฌธ์ tenant_id ๋ฅผ ์์ฒญ์์ JWT ๋ก ๊ฒ์ฆํ๋ค(๊ธฐ๋ณธ๊ฐ). 4.6 ์ฐธ๊ณ .
tenant_auth=True,
)
app = Starlette()
# /chat/stream ยท /chat/stream/attach ยท /chat/stop ยท /chat/steer ๋ฅผ ํ ๋ฒ์ ๋ถ์ด๊ณ ,
# tenant_auth ๊ฐ ์ผ์ ธ ์์ผ๋ฉด ์คํธ๋ฆผ ๊ฒฝ๋ก์ ๊ฒ์ฆ ๋ฏธ๋ค์จ์ด๋ ๊ฑด๋ค.
server.mount_chat_routes(app)
uvicorn_server = uvicorn.Server(
uvicorn.Config(app, host="127.0.0.1", port=8010, log_level="info")
)
await asyncio.gather(
server.run(),
uvicorn_server.serve(),
)
if __name__ == "__main__":
asyncio.run(main())์คํธ๋ฆผ ํ๋๋ง ํ์ํ๋ฉด
mount_chat_sse(app, path="/chat/stream")๋ฅผ ๊ทธ๋๋ก ์ธ ์ ์์ต๋๋ค(ํํธ๋นํธ์ ๋ฐ ๋ ์ง์คํธ๋ฆฌ๋ ์ด ๊ฒฝ๋ก์๋ ์ ์ฉ๋ฉ๋๋ค). ๋ค๋ง ์ฌ์ ์ยท์ค์ง ๋ผ์ฐํธ์ ํ ๋ํธ ๊ฒ์ฆ ๋ฏธ๋ค์จ์ด๋mount_chat_routes()๋ง ๋ถ์ฌ ์ค๋๋ค.
์ 1์์น: Executor ๋ A2A ํ์ค ์ด๋ฒคํธ๋ง emit. SDK ๋งค์ง ๋ฉํ๋ฐ์ดํฐ(
metadata.update({"type":"token",...})๊ฐ์) ์ผ์ ์ฌ์ฉ ๊ธ์ง. A2A ์ด๋ฒคํธ ํ์ ์์ฒด๊ฐ ๋ผ์ฐํ ํค.
์ด๋ฒคํธ ํ๋ฆ:
Task(state=SUBMITTED)โ ๋ผ์ดํ์ฌ์ดํด ์์- ํ ํฐ๋ง๋ค:
Message(text=token)โ ์ฑํ (SSE) ํ ํฐ ์ฒญํฌ. ChatEventQueue ๊ฐ ์ฒ๋ฆฌ.TaskStatusUpdateEvent(state=WORKING, text=token)โ ํ๋ก์ธ์ค ์งํ. ProcessEventQueue ๊ฐ events ํ ์ด๋ธ์ ์ ์ฅ.
TaskStatusUpdateEvent(state=COMPLETED, text=full)โ ์ข ๋ฃ ์๋ฆผ (์ ํ)TaskArtifactUpdateEvent(last_chunk=True, text=full)โ ์ต์ข ๊ฒฐ๊ณผ. ๋ ๋ค ์ฒ๋ฆฌ (chats / todolist).
๋๊ตฌ ํธ์ถ ๊ฐ์ ๋๋ฉ์ธ sub-event ๋ ์ ์ฝ๋ ํ๋ฆ๊ณผ ๋ณ๊ฐ๋ก, "๋๋ฉ์ธ sub-event ๊ธฐ๋ก" ์น์ ์ ํจํด (
metadata["event_type"]) ์ ์ฐธ๊ณ ํด์ ๋๊ตฌ ํธ์ถ ๋จ์์์ emit ํ์ธ์.
import os
from a2a.helpers import (
new_task,
new_text_artifact_update_event,
new_text_message,
new_text_status_update_event,
)
from a2a.types import Role, TaskState
import litellm
class MyExecutor(...):
async def execute(self, context, event_queue):
model = os.environ.get("LLM_MODEL")
proxy_url = (os.environ.get("LLM_PROXY_URL") or "").rstrip("/")
api_key = os.environ.get("LLM_PROXY_API_KEY")
api_base = proxy_url if proxy_url.endswith("/v1") else f"{proxy_url}/v1"
task_id = str(context.task_id)
context_id = str(context.context_id)
# 1) ๋ผ์ดํ์ฌ์ดํด ์์
await event_queue.enqueue_event(
new_task(task_id=task_id, context_id=context_id, state=TaskState.TASK_STATE_SUBMITTED)
)
# 2) LLM ์คํธ๋ฆฌ๋ฐ
stream = await litellm.acompletion(
model=model,
messages=[
{"role": "system", "content": "You are a helpful assistant. Reply in Korean."},
{"role": "user", "content": context.get_user_input()},
],
temperature=0, stream=True,
api_base=api_base, api_key=api_key,
)
full = ""
async for chunk in stream:
try:
token = chunk.choices[0].delta.content
except Exception:
token = None
if not token:
continue
full += token
# ์ฑํ
์ฉ โ Message 1๊ฐ = SSE token 1๊ฐ
await event_queue.enqueue_event(
new_text_message(text=token, role=Role.ROLE_AGENT)
)
# ํ๋ก์ธ์ค์ฉ โ events ํ
์ด๋ธ์ ์งํ row 1๊ฐ์ฉ.
# ๋ถํ๊ฐ ์ฐ๋ ค๋๋ฉด ์ฌ๊ธฐ์ ์ง์ ํํฐ๋ง/์ง๊ณ (์: JSON payload ๋จ์๋ก๋ง emit).
await event_queue.enqueue_event(
new_text_status_update_event(
task_id=task_id, context_id=context_id,
state=TaskState.TASK_STATE_WORKING,
text=token,
)
)
# 3) ์ข
๋ฃ ์๋ฆผ (์ ํ)
await event_queue.enqueue_event(
new_text_status_update_event(
task_id=task_id, context_id=context_id,
state=TaskState.TASK_STATE_COMPLETED,
text=full,
)
)
# 4) ์ต์ข
๊ฒฐ๊ณผ โ chats / todolist ์์ชฝ์ด ๋์ผ ์ด๋ฒคํธ๋ก ์ ์ฅ
await event_queue.enqueue_event(
new_text_artifact_update_event(
task_id=task_id, context_id=context_id,
name="assistant_response",
text=full,
last_chunk=True,
)
)์ ์ฅ ํ์: chats ํ ์ด๋ธ์ ์ ์ฅ๋๋ payload ๊ตฌ์กฐ๋ ํ๋ ์์ํฌ๊ฐ ์ฑ ์์ง๋๋ค. Executor๋ raw A2A ์ด๋ฒคํธ๋ง emitํ๋ฉด ๋ฉ๋๋ค. ์ ์ฅ ์คํค๋ง๋ฅผ ๋ฐ๊พธ๊ณ ์ถ๋ค๋ฉด
mount_chat_sse(persist=...)๋ก ์ปค์คํ persist ํจ์๋ฅผ ์ฃผ์ ํ์ธ์.
ํ์ํธํ: ๊ธฐ์กด์ ์ฑํ ๊ฒฝ๋ก์์
Message๋ก ์ต์ข ์๋ต์ emitํ๋ Executor๋ ๊ทธ๋๋ก ๋์ํฉ๋๋ค.ChatEventQueue๋TaskArtifactUpdateEvent(๊ถ์ฅ) ๋๋Message๋ ๋ค ์ต์ข ์๋ต์ผ๋ก ๋ฐ์๋ค์ ๋๋ค.
์์ฒญ ๋ฐ๋ ์์:
curl -N -X POST http://127.0.0.1:8010/chat/stream \
-H 'content-type: application/json' \
-H 'authorization: Bearer <supabase-jwt>' \
-d '{"message":"hello","conversation_id":"conv-1","tenant_id":"t1","user_uid":"u1"}'tenant_auth=True(๊ธฐ๋ณธ๊ฐ) ๋ฉด Authorization: Bearer โฆ ๊ฐ ํ์ํฉ๋๋ค. ํ ํฐ์ด ์์ผ๋ฉด
401, ๋ณธ๋ฌธ์ tenant_id ๊ฐ ๊ทธ ์ฌ์ฉ์์ ์์์ด ์๋๋ฉด 403 ์
๋๋ค. ๋ณธ๋ฌธ์ tenant_id ๋ฅผ
๋ฃ์ง ์์ผ๋ฉด ๊ฒ์ฆ๋ ๊ฐ์ด ์๋์ผ๋ก ์ฑ์์ง๋๋ค. ํ์ํธํ์ผ๋ก ๋ณธ๋ฌธ user_jwt ๋ ๋ฐ์ต๋๋ค.
์ฑํ (SSE)์ ํฌํจํด ์ฌ์ฉํ๋ ค๋ฉด extras๊ฐ ํ์ํฉ๋๋ค.
pip install "process-gpt-agent-sdk[sse]"ํ ๋ํธ ์ธ์ฆ๊น์ง ์ฐ๋ ค๋ฉด JWT ๊ฒ์ฆ์ฉ extras๋ฅผ ํจ๊ป ์ค์นํฉ๋๋ค.
pip install "process-gpt-agent-sdk[sse,auth]"| extras | ๋ค์ด์ค๋ ๊ฒ | ์ธ์ ํ์ํ๊ฐ |
|---|---|---|
sse |
starlette, uvicorn | ์ฑํ ๋ผ์ฐํธ๋ฅผ ๋ง์ดํธํ ๋ |
auth |
pyjwt[crypto] | tenant_auth=True ๋ก ๋ ๋ |
auth ๋ ํจ์ ์์์ import ํ๋ฏ๋ก, ์ธ์ฆ์ ๋ ์๋ฒ๋ ์ค์นํ์ง ์์๋ ๊ธฐ๋์ ์ํฅ์ด
์์ต๋๋ค(๊ฒ์ฆ์ ์ค์ ๋ก ์๋ํ๋ ์๊ฐ ์ค์น ์๋ด์ ํจ๊ป 500 ์ด ๋ฉ๋๋ค).
์ฐธ๊ณ ๋ก, ๋ ํฌ์๋ ๋น ๋ฅด๊ฒ ํ์ธํ ์ ์๋ ์ํ(sample_server/minimal_server.py, sample_server/minimal_executor.py)๋ ํฌํจ๋์ด ์์ต๋๋ค.
mount_chat_routes() ๋ฅผ ์ฐ๋ฉด ์๋ ์ธ ๊ฐ์ง๊ฐ ์๋์ผ๋ก ๋ฐ๋ผ์ต๋๋ค. Executor ๋ ๋ฐ๋์ง
์์ต๋๋ค โ ์ง๊ธ๊น์ง์ฒ๋ผ A2A ์ด๋ฒคํธ๋ง emit ํ๋ฉด ๋ฉ๋๋ค.
| ๊ธฐ๋ฅ | ๋ฌด์์ ํด๊ฒฐํ๋ |
|---|---|
| ํํธ๋นํธ | ํ ํด์ LLM ์ด ์ค๋ ์๊ฐํ๋ ๋์ ์ ๋ถ์ฉ ์๋ฌด ์ด๋ฒคํธ๋ ๋ด๋ณด๋ด์ง ์๋๋ค. ์ค๊ฐ ํ๋ก์(Cloudflare ๋ฑ)๊ฐ ์ ํด ์ปค๋ฅ์
์ 100์ด ์ํ์์ ๋์ผ๋ฉด ๋ฐฑ์๋๋ ๊ณ์ ๋๋๋ฐ ํ๋ฉด๋ง "์๊ฐ ์คโฆ" ์์ ๋ฉ์ถ๋ค. 15์ด๋ง๋ค SSE ์ฃผ์(: keep-alive)์ ๋ผ์ ์ปค๋ฅ์
์ ์ด๋ ค ๋๋ค. ์ฃผ์์ด๋ผ ํ๋ก ํธ ํ์(data: ๋ง ์ฒ๋ฆฌ)๋ ๋ฌด์ํ๋ค. |
| ์ฌ์ ์ | ์๋ก๊ณ ์นจํ๊ฑฐ๋ ๋ฐฉ์ ๋ค์ ์ด๋ฉด ์งํ ์ค์ธ ํด์ ๊ฒฐ๊ณผ๋ฅผ ์์ ๋ชป ๋ฐ์๋ค. /chat/stream/attach ๊ฐ ์ง๊ธ๊น์ง ์์ธ ๋ณธ๋ฌธ์ snapshot 1๊ฑด์ผ๋ก ์ฃผ๊ณ ์ดํ ํ ํฐ์ ์ค์๊ฐ์ผ๋ก ์๋๋ค. |
| ์ค์ง | ํ๋ก ํธ์ ์ค์ง ๋ฒํผ์ ์๊ธฐ ์ชฝ fetch ๋ง ๋์ ๋ฟ์ด๋ผ ์๋ฒ ์คํ์ ๊ณ์ ๋๋ฉฐ ํ ํฐยท๋๊ตฌ ํธ์ถ์ ๊ทธ๋๋ก ์๋นํ๋ค. /chat/stop ์ด ์คํ task ๋ฅผ ์ค์ ๋ก ์ทจ์ํ๋ค. |
์ฌ๊ธฐ์ ๋ํด, ๊ฐ์ ๋ฐฉ์ ์ ํด์ด ๋ค์ด์ค๋ฉด ์ด์ ํด์ ๋จผ์ ๋์ต๋๋ค. ํ๋ก ํธ๊ฐ ๋ก๋ฉ ์ค
์ ๋ฉ์์ง๋ฅผ ๋ณด๋ผ ๋ ์๋ฒ์๋ ์ทจ์ ์ ํธ๋ฅผ ์ฃผ์ง ์์, ์ด๊ฒ ์์ผ๋ฉด ๊ฐ์
conversation_id ์ ๋ ์คํ์ด ๊ฒน์นฉ๋๋ค.
์ฌ์ ์ ์์ฒญ
curl -N -X POST http://127.0.0.1:8010/chat/stream/attach \
-H 'content-type: application/json' \
-H 'authorization: Bearer <supabase-jwt>' \
-d '{"conversation_id":"conv-1","tenant_id":"t1"}'ํ์ฑ ํด์ด ์์ผ๋ฉด text/event-stream ์ผ๋ก ์๋ตํฉ๋๋ค.
event: message
data: {"type": "snapshot", "content": "์ง๊ธ๊น์ง ์ด ๋ณธ๋ฌธ"}
event: message
data: {"type": "token", "content": "์ด์ด์"}
ํ์ฑ ํด์ด ์์ผ๋ฉด SSE ๊ฐ ์๋๋ผ 200 {"active": false} ์
๋๋ค. 404 ๊ฐ ์๋ ์ด์ ๋,
์ฒซ attach ์๋๋ ํญ์ "์์ง ์คํธ๋ฆผ ์์" ์ด๋ผ 404 ๋ก ๋๋ฉด ๋ฉ์์ง๋ฅผ ๋ณด๋ผ ๋๋ง๋ค ๋ธ๋ผ์ฐ์
๋คํธ์ํฌ ํญ์ ์คํจ ์์ฒญ์ด ์์ด๊ธฐ ๋๋ฌธ์
๋๋ค. ํ๋ก ํธ๋ content-type ์ด
text/event-stream ์ด ์๋๋ฉด ์กฐ์ฉํ ์ข
๋ฃํ๋ฉด ๋ฉ๋๋ค.
์ค์ง ์์ฒญ
curl -X POST http://127.0.0.1:8010/chat/stop \
-H 'content-type: application/json' \
-H 'authorization: Bearer <supabase-jwt>' \
-d '{"conversation_id":"conv-1","tenant_id":"t1"}'| ์๋ต | ์๋ฏธ |
|---|---|
{"stopped": true} |
์งํ ์ค์ด๋ ํด์ ์ทจ์ํ๋ค |
{"stopped": false, "reason": "no_active_turn"} |
์ทจ์ํ ์คํ์ด ์๋ค(์ด๋ฏธ ๋๋ฌ๊ฑฐ๋ HITL ๋๊ธฐ ์ค) |
403 {"stopped": false, "reason": "forbidden"} |
์์ฒญ์์ ํ ๋ํธ์ ๋ฐฉ ์์ ํ ๋ํธ๊ฐ ๋ค๋ฅด๋ค |
์ฌ์ ์ยท์ค์ง ๋ชจ๋ ๋ฐฉ ์์ ํ
๋ํธ(chat_rooms.tenant_id)์ ์์ฒญ์์ ํ
๋ํธ๊ฐ ๊ฐ์ ๋๋ง
๋์ํฉ๋๋ค. ๋ฐฉ์ ์กฐํํ์ง ๋ชปํ๋ฉด ๊ฑฐ๋ถํฉ๋๋ค(fail-closed).
ํํธ๋นํธ ์ฃผ๊ธฐ๋ SSE_HEARTBEAT_SECONDS ํ๊ฒฝ๋ณ์๋ก ๋ฐ๊ฟ๋๋ค(๊ธฐ๋ณธ 15์ด).
์์ ์ด ๋๋๊ธฐ ์ ์ ๋ฐฉํฅ์ด ์ด๊ธ๋ ๊ฒ์ ๋ฐ๊ฒฌํ์ ๋, ์คํ์ ์ทจ์ํ๊ณ ์ฒ์๋ถํฐ ๋ค์ ์ํค๋ ๋์ ์ง๊ธ๊น์ง์ ๋งฅ๋ฝ์ ์ ์งํ ์ฑ ์ง์๋ง ๋ฐ๊ฟ๋๋ค. ์ฌ์ฉ์๋ ๊ธด ์์ ์ ์์จ์ ์ผ๋ก ๋๋ ค ๋๊ณ ํ์ํ ๋๋ง ๊ฐ์ ํฉ๋๋ค.
์ด๊ฑด ํน์ ์์ด์ ํธ์ ๊ธฐ๋ฅ์ด ์๋๋ผ ํ์ค ๋์์ ๋๋ค. SDK ๊ฐ ์์ฒญ ํํ์ ์ด๋ฒคํธ๊น์ง๋ฅผ ์์ ํ๊ณ , ์ค์ ์คํ ์ ํ์ ์์ด์ ํธ๋ณ ์ด๋ํฐ๊ฐ ๋งก์ต๋๋ค โ ์ง์๋ฅผ ์ธ์ ์ง์ด๋ฃ์ด์ผ ์์ ํ์ง๋ ๋ฐํ์๋ง๋ค ๋ค๋ฅด๊ธฐ ๋๋ฌธ์ ๋๋ค.
curl -X POST http://127.0.0.1:8010/chat/steer \
-H 'content-type: application/json' \
-H 'authorization: Bearer <supabase-jwt>' \
-d '{"conversation_id":"conv-1","tenant_id":"t1","message":"ํ ๋์ ๊ธ๋ก ์จ ์ค"}'/chat/stream ์ผ๋ก {"action":"steer", ...} ๋ฅผ ๋ณด๋ด๋ ๊ฐ์ต๋๋ค(์๋ํฌ์ธํธ๋ฅผ ํ๋๋ง ์๋
ํด๋ผ์ด์ธํธ๋ฅผ ์ํด). action ์ด ์๋ ๊ธฐ์กด ์์ฒญ์ ์ข
์ ๋๋ก ์ ํด์ ๋๋ฆฝ๋๋ค. ๋ชจ๋ฅด๋
action ์ 400 ์
๋๋ค โ ์คํ(steeer)๋ฅผ ํ๋ฒํ ๋ฉ์์ง๋ก ํ๋ฆฌ๋ฉด ๋ฐฉํฅ์ ๋ฐ๊พธ๋ ค๋ ์์ฒญ์ด
์งํ ์ค์ธ ํด์ ๋์ฒดํด ์์
์ ๋ ๋ฆฝ๋๋ค.
| ์๋ต | ์๋ฏธ |
|---|---|
{"accepted": true, "directive_id": "โฆ"} |
๋ฐ์๋ค. ๋ฐ์์ ์๋๋ค(์๋ ์ฐธ๊ณ ) |
{"accepted": true, "duplicate": true, "directive_id": "โฆ"} |
๊ฐ์ ๋ฌธ์ฅ์ ์ฐ์์ผ๋ก ๋ฐ์๋ค. ์ฒ์ ์ ์์ id ๋ฅผ ๊ทธ๋๋ก ์ค๋ค |
409 {"accepted": false, "reason": "no_active_turn"} |
๋๊ณ ์๋ ํด์ด ์๋ค(์ด๋ฏธ ๋๋ฌ๊ฑฐ๋ ๋ง๋ฌด๋ฆฌ์ ๋ค์ด๊ฐ๋ค) |
409 {"accepted": false, "reason": "awaiting_human_input"} |
์ฌ๋์ ๋ต์ ๊ธฐ๋ค๋ฆฌ๋ฉฐ ๋ฉ์ถฐ ์๋ค. ๊ทธ ์ง๋ฌธ์ ๋ตํ๋ ๊ฒ์ด ๋ฐฉํฅ ์ ํ์ด๋ค |
501 {"accepted": false, "reason": "unsupported"} |
์ด Executor ๊ฐ steer() ๋ฅผ ๊ตฌํํ์ง ์์๋ค |
400 {"accepted": false, "reason": "empty_message"} |
๋ณด๋ผ ์ง์๊ฐ ์๋ค |
403 {"accepted": false, "reason": "forbidden"} |
์์ฒญ์์ ํ ๋ํธ์ ๋ฐฉ ์์ ํ ๋ํธ๊ฐ ๋ค๋ฅด๋ค |
์ ์์ ๋ฐ์์ ๋ค๋ฅธ ์ด๋ฒคํธ์ ๋๋ค. ์ ์ ์์ ์ ์์ด์ ํธ๋ ์์ง ์๋ ์ง์๋๋ก ๋๊ตฌ๋ฅผ ๋๋ฆฌ๊ณ ์์ต๋๋ค. ๋์ ํ ์ด๋ฒคํธ๋ก ํฉ์น๋ฉด ํ๋ฉด์ ์ ์๋ง์ผ๋ก "๋ฐ์ ์๋ฃ" ๋ฅผ ํ์ํ๊ณ , ์ฌ์ฉ์๋ ๋ฐ์๋์ง ์์ ๊ฒฐ๊ณผ๋ฅผ ๋ฐ์๋ ๊ฒ์ผ๋ก ์ฝ์ต๋๋ค.
event: message
data: {"type": "steer_accepted", "directive_id": "โฆ", "content": "ํ ๋์ ๊ธ๋ก ์จ ์ค"}
event: message
data: {"type": "steer_applied", "directive_id": "โฆ", "content": "ํ ๋์ ๊ธ๋ก ์จ ์ค"}
๋ ์ด๋ฒคํธ๋ ํด์ ์ถ๋ ฅ ํ๋ก ๋๊ฐ๋๋ค โ ์๋ ํด๋ผ์ด์ธํธ์ ์ฌ์ ์ํ ํด๋ผ์ด์ธํธ๊ฐ ๊ฐ์
๊ฒ์ ๋ด
๋๋ค. ์ ์๋ง ๋๊ณ ์์ง ๋ฐ์๋์ง ์์ ์ง์๋ ์ฌ์ ์ ์ค๋
์ท์ pending_steers ๋ก๋
์ค๋ฆฝ๋๋ค.
| ์ํฉ | ์ฒ๋ฆฌ |
|---|---|
| ๋๊ตฌ ์คํ ์ค | ์ ์๋ง ํ๊ณ ๋๊ธฐ์ด์ ๋ฃ๋๋ค. ๋๊ตฌ๋ฅผ ์ค๊ฐ์ ๋์ง ์๋๋ค โ ์ฐ๋ค ๋ง ํ์ผ์ด๋ ๊ฒฐ๊ณผ ์๋ ๋๊ตฌ ํธ์ถ์ ๋จ๊ธฐ๋ ํธ์ด ๋ ๋์๋ค. ์ด๋ํฐ๊ฐ ๋ค์ ์์ ์ง์ ์์ ์ง์ด ๊ฐ๋ค |
| ์๋ฃ ์ง์ | ์ด๋ํฐ๊ฐ ๋ง๋ฌด๋ฆฌ ์ ์ ๋๊ธฐ์ด์ ๋ซ๊ณ ๋จ์ ์ง์๋ฅผ ๋ง์ง๋ง์ผ๋ก ์ง์ด ๊ฐ๋ค. ๋ซํ ๋ค์ ์ง์๋ no_active_turn ์ผ๋ก ๊ฑฐ์ ํ๋ค โ ๋ฐ์ ๋๊ณ ์๋ฌด ๋ฐ๋ ๋ฐ์ํ์ง ์๋ ๊ฒ๋ณด๋ค ๊ฑฐ์ ์ด ์ ์งํ๋ค |
| ์ค๋ณต ์ฐ์ ์์ | ๋ ๋ฒ์งธ๋ถํฐ๋ ์์ง ์๊ณ duplicate: true ์ ์ฒ์ ์ ์์ id ๋ฅผ ์ค๋ค. ์ ์ ์ด๋ฒคํธ๋ ๋ค์ ๋ด๋ณด๋ด์ง ์๋๋ค |
| ์ฌ๋ ํ์ธ ๋๊ธฐ ์ค | ๋๊ณ ์๋ ์คํ์ด ์์ด ๋ฃ์ ๊ณณ์ด ์๋ค. awaiting_human_input ์ผ๋ก ๊ฑฐ์ ํ๋ค |
| ์ฌ์ ์ | ์ ์ยท๋ฐ์์ด ๋ค๋ฅธ ์ด๋ฒคํธ์ ๊ฐ์ ํ๋ก ๋๊ฐ๋ฏ๋ก ๊ทธ๋๋ก ๋ฐ๋๋ค + ์ค๋
์ท์ pending_steers |
์ด๋ํฐ ์ชฝ(์์ด์ ํธ ์ ์ฅ์) ์ด ๊ตฌํํ ๊ฒ์ ๋ ๊ฐ์ง์ ๋๋ค โ ์ง๊ธ ๋ฐ์ ์ ์๋์ง ํ์ ํ๋ ๊ฒ๊ณผ, ์์ ์ง์ ์์ ์ค์ ๋ก ์น๋ ๊ฒ.
from processgpt_agent_sdk import SteerDirective, SteerResult, get_steering_inbox, mark_applied
class MyExecutor(AgentExecutor):
async def steer(self, directive: SteerDirective) -> SteerResult | None:
"""์ง๊ธ ๋ฐฉํฅ์ ๋ฐ๊ฟ ์ ์๋๊ฐ. None ์ด๋ฉด SDK ์ ํ์ค ํ์ ์ ๊ทธ๋๋ก ์ด๋ค."""
if await self._parked_on_question(directive.conversation_id):
return SteerResult.reject("awaiting_human_input")
return None๊ทธ๋ฆฌ๊ณ ์คํ ์ชฝ์์๋, ๋ค์ ํ๋จ์ด ์์๋๊ธฐ ์ง์ (๋๊ตฌ๊ฐ ๋๋๊ณ ๋ชจ๋ธ์ ๋ถ๋ฅด๊ธฐ ์ )๋ง๋ค:
taken = await get_steering_inbox().take(cid) # ์ง์ด ๊ฐ๊ธฐ โ ๋ฐ์
for directive in taken:
await mark_applied(directive) # ์ค์ ๋ก ๋ค์ ํ๋จ์ ๋ฃ๋ ์๊ฐ
... # directive.message ๋ฅผ ์
๋ ฅ์ ์น๋๋ค
# ๋ ์ด์ ์น์ ์ง์ ์ด ์๋ค๋ฉด(ํด์ด ๋๋๋ ค ํ๋ค๋ฉด) ๋๊ธฐ์ด์ ๋ซ๋๋ค.
# ๋ซ์ ๋ค์๋ ์คํ์ด ์ด์ด์ง๊ฒ ๋๋ค๋ฉด reopen(cid) ์ผ๋ก ๋ค์ ์ฐ๋ค.
await get_steering_inbox().close(cid)๊ทธ "์์ ์ง์ " ์ด ์ด๋์ธ์ง๋ ๋ฐํ์์ด ์ ํฉ๋๋ค. deepagents ๋ LangChain ๋ฏธ๋ค์จ์ด์
๋ชจ๋ธ ํธ์ถ ์ง์ ํ
(abefore_model)์์ ์ง์ด ๊ฐ๊ณ , ํด์ ๋๋ด๋ ค๋ ์์
(aafter_model)์ ํ ๋ฒ ๋ ํ์ธํด ๋จ์ ์ง์๊ฐ ์์ผ๋ฉด ๋ชจ๋ธ๋ก ๋๋๋ฆฝ๋๋ค.
steer() ๊ฐ ์์ผ๋ฉด ๊ทธ ์์ด์ ํธ๋ ๋ฏธ์ง์(501)์
๋๋ค. ํ์ค ๋์์ ์ ์ํ๋ ๊ฒ๊ณผ ๋ชจ๋
์์ด์ ํธ๊ฐ ๊ทธ๊ฒ์ ํ ์ ์๋ค๊ณ ์ฃผ์ฅํ๋ ๊ฒ์ ๋ค๋ฆ
๋๋ค.
์ฑํ
์์ฒญ์ tenant_id ๋ Executor ๋ฅผ ์ง๋ ์คํฌ ๋๋ ํฐ๋ฆฌ ๋ก๋ยท์๋๋ฐ์ค ๋ง์ดํธยท์์
๊ณต๊ฐ
๊ฒฝ๋ก๊น์ง ๊ทธ๋๋ก ํ๋ฌ๊ฐ๋๋ค. ๊ฒ์ฆํ์ง ์์ผ๋ฉด ๊ฐ๋ง ๋ฐ๊ฟ ํธ์ถํด ๋จ์ ํ
๋ํธ ์์์ ๋ฟ์ ์
์์ต๋๋ค. ๊ทธ๋์ tenant_auth ์ ๊ธฐ๋ณธ๊ฐ์ ์ผ์ง์
๋๋ค.
server = ProcessGPTAgentServer(
agent_executor=MyExecutor(),
agent_type="crewai-action",
tenant_auth=True, # ๊ธฐ๋ณธ๊ฐ
)
server.mount_chat_routes(app)๊ฒ์ฆ ๊ฒฝ๋ก๋ ํ ํฐ ํค๋์ alg ์ ๋ฐ๋ผ ๊ฐ๋ฆฝ๋๋ค.
- ๋น๋์นญ(ES256/RS256 โฆ) โ Supabase JWKS(
/auth/v1/.well-known/jwks.json) ๋ก์ปฌ ๊ฒ์ฆ. ์ต์ Supabase ํ๋ก์ ํธ(JWT signing keys)๊ฐ ์ฌ๊ธฐ ํด๋นํฉ๋๋ค. - HS256 +
SUPABASE_JWT_SECRETโ ๋ ๊ฑฐ์ ๋์นญํค ํ๋ก์ ํธยท์ปค์คํ SSO ํ ํฐ ๋ก์ปฌ ๊ฒ์ฆ. - ๋ ๋ค ๋ถ๊ฐํ๋ฉด GoTrue
/auth/v1/user์ ์์.
์์ ํ
๋ํธ๋ JWT ํด๋ ์(tenant_id / app_metadata.tenant_id / tenant_ids)์ ๋จผ์
๋ณด๊ณ , ์์ผ๋ฉด users ํ
์ด๋ธ์์ ์กฐํํฉ๋๋ค(๋ฉํฐ ํ
๋ํธ ์์ ๋์). ๊ฒ์ฆ ๊ฒฐ๊ณผ๋ ํ ํฐ
๋จ์๋ก 60์ด ์บ์ํ๊ณ , ํ ํฐ ๋ง๋ฃ๊ฐ ๋ ์ด๋ฅด๋ฉด ๊ทธ ์์ ๊น์ง๋ง ์บ์ํฉ๋๋ค.
| ํ๊ฒฝ๋ณ์ | ์ฐ์ |
|---|---|
SUPABASE_URL |
JWKS ์ฃผ์๋ฅผ ๋ง๋ ๋ค(๋น๋์นญ ๊ฒ์ฆ) |
SUPABASE_JWT_SECRET |
HS256 ๋์นญํค ๊ฒ์ฆ(์์ ๋๋ง) |
ํ์ ๊ท์น
| ์ํฉ | ๊ฒฐ๊ณผ |
|---|---|
| ํ ํฐ ์์ | 401 |
์์ฒญํ tenant_id ๊ฐ ์์์ด ์๋ |
403 |
์์ฒญ์ tenant_id ์๊ณ ์์์ด ํ๋ |
๊ทธ ๊ฐ์ผ๋ก ์ฑ์ ํต๊ณผ |
์์ฒญ์ tenant_id ์๊ณ ์์์ด ์ฌ๋ฟ |
400 (tenant_id is required) |
| ์์ ํ ๋ํธ๊ฐ ํ๋๋ ์์ | 403 |
์ง์ ๋ง๋ ๋ผ์ฐํธ์๋ ๊ฐ์ ๊ท์น์ ๊ฑธ ์ ์์ต๋๋ค.
from processgpt_agent_sdk import tenant_guard, request_tenant_id
@tenant_guard
async def my_handler(request):
# ์์ฒญ์ด ๋ณด๋ธ ๊ฐ์ด ์๋๋ผ ๊ฒ์ฆ๋ ๊ฐ๋ง ์ด๋ค
tenant_id = request_tenant_id(request)
...ํ
๋ํธ ์ค์ฝํ๊ฐ ์๋ ์๋ํฌ์ธํธ๋ auth_guard ๋ก ์ธ์ฆ๋ง ํ์ธํฉ๋๋ค.
๋๋ ๊ฒฝ์ฐ: ์๋จ์ ๋ณ๋ ์ธ์ฆ ๊ฒ์ดํธ์จ์ด๊ฐ ์๊ฑฐ๋ ๋ก์ปฌ ๊ฐ๋ฐ์ผ ๋๋ง
tenant_auth=False๋ก ๋ก๋๋ค. ์ด๋๋ ๊ธฐ๋ ๋ก๊ทธ์ ๊ฒฝ๊ณ ๊ฐ ๋จ๊ณ , ์์ฒญ ๋ณธ๋ฌธ์tenant_id๊ฐ ๊ทธ๋๋ก Executor ์ ์ ๋ฌ๋ฉ๋๋ค.
๋ฐ๋์ json.dumps()๋ก ์ง๋ ฌํํด์ผ ํฉ๋๋ค.
-
โ ์ด๋ ๊ฒ ํ๋ฉด ์๋จ:
text = str({"key": "value"}) # Python dict string โ JSON ์๋
DB์
"'{key: value}'"๊ผด๋ก ๋ฌธ์์ด ์ ์ฅ๋จ โ ํ์ฑ ์คํจ -
โ ์ด๋ ๊ฒ ํด์ผ ํจ:
text = json.dumps({"key": "value"}, ensure_ascii=False)
DB์
{"key": "value"}JSON ์ ์ฅ๋จ โ ํ์ฑ ์ฑ๊ณต
๐ SDK๋ ๋ด๋ถ์์ json.loads๋ก ์ฌํ์ฑํ๊ธฐ ๋๋ฌธ์, ํ์ค JSON ๋ฌธ์์ด์ด ์๋๋ฉด ๋ฌด์กฐ๊ฑด ๋ฌธ์์ด๋ก๋ง ๋จ์ต๋๋ค.
ํต์ฌ์ Executor ์์์ ๋ชจ๋๋ฅผ ๋ถ๊ธฐํ์ง ์๋ ๊ฒ์ ๋๋ค. ๋์ผํ A2A ์ด๋ฒคํธ ์ํ์ค๋ฅผ emitํ๋ฉด, ํ๋ ์์ํฌ์ EventQueue ๊ตฌํ์ฒด๊ฐ ํ๋ก์ธ์ค/์ฑํ ์ ๋ง๊ฒ ๋ผ์ฐํ ํฉ๋๋ค.
- ๊ณตํต ์ด๋ฒคํธ ํ๋ฆ:
TaskโTaskStatusUpdateEvent[..]โTaskArtifactUpdateEvent(last_chunk=True) - ํ๋ก์ธ์ค(ํด๋ง) ๊ฒฝ๋ก:
ProcessEventQueue๊ฐ status๋ events ํ ์ด๋ธ, artifact๋ todolist ํ ์ด๋ธ์ ์ ์ฅ - ์ฑํ
(SSE) ๊ฒฝ๋ก:
ChatEventQueue๊ฐ status๋ SSE message ์ฒญํฌ๋ก, ์ต์ข artifact๋ SSE done + chats ํ ์ด๋ธ ์ ์ฅ์ผ๋ก ๋ณํ. ์ต์ข ์ด๋ฒคํธ์pdfFiles๋ฅผ ์ค์ผ๋ฉด ์ฐ์ถ๋ฌผ๋ ๊ฐ์ด ๋ฐ๋ผ๊ฐ๋๋ค(7.5) - ์ ์ก ๊ณ์ธต: ํํธ๋นํธยท์ฌ์ ์ยท์ค์งยทํ ๋ํธ ๊ฒ์ฆ์ ํ๋ ์์ํฌ๊ฐ ์ฒ๋ฆฌํฉ๋๋ค. Executor ๋ ์ด๋ค์ ์ ํ์๊ฐ ์์ต๋๋ค (4.5 ยท 4.6)
from processgpt_agent_sdk import ProcessGPTAgentServer
server = ProcessGPTAgentServer(agent_executor=MyExecutor(), agent_type="crewai-action")
await server.run()SSE๋ฅผ ์ฐ๋ ค๋ฉด extras ์ค์น๊ฐ ํ์ํฉ๋๋ค.
pip install "process-gpt-agent-sdk[sse,auth]"from starlette.applications import Starlette
from processgpt_agent_sdk import ProcessGPTAgentServer
server = ProcessGPTAgentServer(agent_executor=MyExecutor(), agent_type="crewai-action")
app = Starlette()
server.mount_chat_routes(app)
# POST /chat/stream ยท /chat/stream/attach ยท /chat/stop ยท /chat/steer๊ฒฝ๋ก๋ฅผ ๋ฐ๊พธ๊ฑฐ๋ ์ผ๋ถ๋ง ๋ถ์ผ ์๋ ์์ต๋๋ค. ์์ ํ๋ก ํธ๊ฐ ์ฐ๋ ๋ณ์นญ ๊ฒฝ๋ก๊ฐ ์์ผ๋ฉด
stream_paths ์ ํจ๊ป ๋๊น๋๋ค.
server.mount_chat_routes(
app,
stream_paths=("/chat/stream", "/{agent_id}/chat/stream"),
attach_path="/chat/stream/attach", # None ์ด๋ฉด ๋ถ์ด์ง ์๋๋ค
stop_path="/chat/stop", # None ์ด๋ฉด ๋ถ์ด์ง ์๋๋ค
steer_path="/chat/steer", # None ์ด๋ฉด ๋ถ์ด์ง ์๋๋ค
)FastAPI ์ฑ์๋ ๊ทธ๋๋ก ๋๊ธธ ์ ์์ต๋๋ค. ๋ผ์ฐํธ๋ ํญ์ Starlette ๋ผ์ฐํธ๋ก ๋ฑ๋ก๋๋๋ฐ,
FastAPI ์ add_api_route() ์ ๋๊ธฐ๋ฉด ํธ๋ค๋ฌ ์๊ทธ๋์ฒ(ํ์
์ฃผ์ ์๋ request)๋ฅผ ํ์
์ฟผ๋ฆฌ ํ๋ผ๋ฏธํฐ๋ก ํด์ํด ๋ชจ๋ ์์ฒญ์ด 422 ๋ก ๋จ์ด์ง๊ธฐ ๋๋ฌธ์
๋๋ค.
์์ฒญ ๋ฐ๋ ์์:
{
"message": "์๋
",
"tenant_id": "t1",
"user_uid": "u1",
"user_email": "user@example.com",
"user_name": "ํ๊ธธ๋",
"user_jwt": "",
"conversation_id": "conv-1",
"file": null,
"files": [],
"file_count": 0,
"stream": true,
"metadata": {}
}import asyncio
import uvicorn
from starlette.applications import Starlette
from processgpt_agent_sdk import ProcessGPTAgentServer
async def main():
server = ProcessGPTAgentServer(agent_executor=MyExecutor(), agent_type="crewai-action")
app = Starlette()
server.mount_chat_routes(app)
uvicorn_server = uvicorn.Server(
uvicorn.Config(app, host="127.0.0.1", port=8010, log_level="info")
)
await asyncio.gather(
server.run(),
uvicorn_server.serve(),
)
if __name__ == "__main__":
asyncio.run(main())์ด์ ํ๊ฒฝ์์๋ ํด๋ง ํ๋ก์ธ์ค์ HTTP API ํ๋ก์ธ์ค๋ฅผ ๋ถ๋ฆฌ ์ด์ํ๋ ๊ฒฝ์ฐ๋ ๋ง์ต๋๋ค.
์ฌ์ ์ยท์ค์ง ๋ ์ง์คํธ๋ฆฌ์ ๊ธฐ๋ณธ ๊ตฌํ์ ํ๋ก์ธ์ค ๋ก์ปฌ dict ์ ๋๋ค(๋จ์ผ uvicorn ์์ปค ์ ์ ). ํ๋๋ฅผ ์ฌ๋ฌ ๊ฐ ๋์ฐ๋ฉด "A ํ๋๊ฐ ๋๋ฆฌ๋ ํด์ B ํ๋๊ฐ ๋ฐ์ attach ์์ฒญ" ์ด ๋ง์ง ์์ผ๋ฏ๋ก, ๊ณต์ ๋ฐฑ์๋ ๊ตฌํ์ผ๋ก ๊ฐ์๋ผ์๋๋ค. SDK ๋ ๊ณ์ฝ๋ง ์์ ํ๊ณ ๊ณต์ ์ํ ๋ฐฑ์๋๋ ์ ํ๋ฆฌ์ผ์ด์ ์ด ๊ณ ๋ฆ ๋๋ค.
from processgpt_agent_sdk import (
ChatRunRegistry,
set_run_registry,
set_inflight_registry,
)
class RedisRunRegistry(ChatRunRegistry):
async def start_run(self, conversation_id): ...
async def record(self, conversation_id, payload): ...
async def mark_done(self, conversation_id): ...
async def subscribe(self, conversation_id): ...
async def unsubscribe(self, conversation_id, q): ...
set_run_registry(RedisRunRegistry())asyncio.Task.cancel() ์ ๊ทธ task ๋ฅผ ๋ง๋ ํ๋ก์ธ์ค ์์์๋ง ๊ฐ๋ฅํ๋ฏ๋ก, ์ค์ง์ ๊ฒฝ์ฐ
๊ณต์ ๋ฐฑ์๋๋ "์ด ๋ฐฉ์ ์ด๋ ํ๋๊ฐ ๋ค๊ณ ์๋์ง" ์์ ๊ถ ๊ธฐ๋ก๊ณผ ์ทจ์ ์ ํธ ์ ๋ฌ์๋ง ์ฐ๊ณ
์ทจ์ ์์ฒด๋ ํญ์ ์์ ํ๋์์ ์ผ์ด๋์ผ ํฉ๋๋ค.
supersede ์ cancel ์ ๋ค๋ฆ
๋๋ค. InflightRegistry ๋ ๋ ๋ฉ์๋๋ฅผ ๋ฐ๋ก ๋ก๋๋ค.
| ๋ฉ์๋ | ์ธ์ ๋ถ๋ฆฌ๋ | ๊ธฐ๋ณธ ๋์ |
|---|---|---|
supersede(cid) |
๊ฐ์ ๋ฐฉ์ ์ ํด์ด ์์๋ ๋ | cancel() ์ ๋ถ๋ฅธ๋ค |
cancel(cid) |
/chat/stop ์ผ๋ก ๋ช
์์ ์ค์งํ ๋ |
์คํ task ๋ฅผ ์ทจ์ํ๋ค |
์ ์์ฒญ์ด ์ด์ ์๋ต์ ๋์ฒดํ๋์ง ์ด์ด๊ฐ๋์ง๋ ์๋ฒ๋ง๋ค ๋ค๋ฆ
๋๋ค. ๋ํ์ ์ผ๋ก
HITL(์ฌ๋ ํ์ธ) ์๋ต์ ์ด์ ํด์ด ๋จ๊ธด interrupt ๋ฅผ ์ฌ๊ฐํ๋ ๊ฒ์ด๋ผ, ๊ทธ ํด์ ์ทจ์ํ๋ฉด
์ฒดํฌํฌ์ธํธ๊ฐ ์ฌ๋ผ์ ธ ์ฌ๊ฐ๊ฐ ๋ถ๊ฐ๋ฅํด์ง๋๋ค. ์ด๋ ์ชฝ์ธ์ง๋ ์์ฒญ ๋ณธ๋ฌธ์ ํด์ํด์ผ ์ ์
์๊ณ ๊ทธ๊ฑด Executor ์ ๋ชซ์ด๋ฏ๋ก, ๊ทธ๋ฐ ์๋ฒ๋ supersede() ๋ฅผ no-op ์ผ๋ก ์ฌ์ ์ํ๊ณ
Executor ์์์ ์ง์ ํ๋จํฉ๋๋ค.
from processgpt_agent_sdk import InflightRegistry, set_inflight_registry
class ExecutorDecides(InflightRegistry):
async def supersede(self, conversation_id):
return False # ๋์ฒด/์ด์ด๊ฐ๊ธฐ ํ๋จ์ Executor ๊ฐ ํ๋ค
set_inflight_registry(ExecutorDecides())์์ด์ ํธ๊ฐ ๋ง๋ ํ์ผ์ ์ฌ์ฉ์์๊ฒ ๋ด์ฃผ๋ ์ผ์ SDK ๊ฐ ๊ณต์ฉ์ผ๋ก ๊ฐ์ต๋๋ค. ์์ด์ ํธ๋ง๋ค ๋ค์ ๋ง๋ค ๊ฒ์ด ์์ต๋๋ค.
์ด ๋ชจ๋์ด ์๊ธฐ๊ธฐ ์ ์๋ ์์ด์ ํธ๋ง๋ค ์๊ธฐ ๋ฐฉ์์ด ์์์ต๋๋ค. ํ์ชฝ์ ์์ฒด store ์ collector ๋ฅผ, ๋ค๋ฅธ ์ชฝ์ ํ ์๊ฐ์ง๋ฆฌ ํ ํฐ ์ฃผ์๋ฅผ ๋ค๊ณ ์์๊ณ , ๋ ๋ค ์ฃผ์๊ฐ ์ฃฝ์ผ๋ฉด ํ์ผ์ ์์ ๋ฐ์ ์ ์์์ต๋๋ค. ๊ณต๊ฐ ๋ฒํท์ ์๊ตฌ ์ฃผ์๋ฅผ ์ฐ๋ ์ชฝ์ ๋ฐ๋ ๋ฌธ์ ๊ฐ ์์์ต๋๋ค โ ๊ณ์ฝ์ ๊ฒํ ๊ฒฐ๊ณผ๋ ์ฌ๋ด ๋ณด๊ณ ์๊ฐ ์ฃผ์๋ง ์๋ฉด ๋๊ตฌ๋, ์ธ์ ๊น์ง๋ ์ด๋ฆฌ๋ ์๋ฆฌ์ ๋์์ต๋๋ค.
์ด๋ฒ ํด์ outputs/ ์ ๋์ธ ๊ฒ๋ง ์ฐ์ถ๋ฌผ์
๋๋ค.
๋ต๋ณ ๋ณธ๋ฌธ์ ํ์ด ๊ฒฝ๋ก๋ฅผ ์ฐพ์๋ด๋ ๋ฐฉ์์ ์ฐ์ง ์์ต๋๋ค. ๊ทธ ๊ฒฝ๋ก๋ฅผ ์ ๋ ์ฃผ์ฒด๊ฐ ์์ด์ ํธ๋ผ์, ๋ฌธ์ฅ์ด ๋ฐ๋ ๋๋ง๋ค ๊ท์น์ ๊ณ ์ณ์ผ ํ๊ณ ๋ชป ์ก์ผ๋ฉด ์ฌ์ฉ์๋ ํ์ผ์ ๋ฐ์ง ๋ชปํฉ๋๋ค. ๋๋ ํฐ๋ฆฌ๋ ์์ด์ ํธ๊ฐ ์ด๋ป๊ฒ ๋งํ๋ ๊ฐ์ต๋๋ค.
๊ฑฐ๋ ๋๋ ์ด๋ฒ ํด์ ์๋ก ์๊ธฐ๊ฑฐ๋ ๋ฐ๋ ๊ฒ๋ง ๋ด
๋๋ค(snapshot). ๊ทธ๋ฌ์ง ์์ผ๋ฉด
์ด์ด์ง๋ ํด๋ง๋ค ๊ฐ์ ํ์ผ์ด ๋ค์ ์ฌ๋ผ๊ฐ ํ๋ฉด์ ๊ฐ์ ๊ฒ์ด ์ฌ๋ฌ ๋ฒ ๋น๋๋ค.
| ๋ชจ๋ | ํ๋ ์ผ |
|---|---|
ArtifactCollector |
outputs/ ๋ฅผ ๊ฑฐ๋ฌ ์ฐ์ถ๋ฌผ ๋ ์ฝ๋๋ก ๋ง๋ ๋ค |
ArtifactStore (Protocol) |
๋ณด๊ดํ๊ณ ์ฃผ์๋ฅผ ๋ฐ๊ธํ๋ ๊ณ์ฝ |
MementoArtifactStore |
๊ธฐ๋ณธ ๊ตฌํ โ Memento ๋ฅผ ๊ฑฐ์ณ ๋น๊ณต๊ฐ ๋ฒํท์ ๋ณด๊ด |
build_artifact ยท download_file |
์๋ฒ ์์ ์ ๋ณธ โ ํ๋ฉด์ด ์ฝ๋ ๋ชจ์ |
strip_local_paths |
๋ณธ๋ฌธ์ ๋จ์ ๋ด๋ถ ๊ฒฝ๋ก๋ฅผ ์ฐ์ถ๋ฌผ ์ฃผ์๋ก ์นํ |
๋ชจ๋ from processgpt_agent_sdk.artifacts import ... ๋ก ๊ฐ์ ธ์ต๋๋ค.
์ฐ์ถ๋ฌผ์ ๋น๊ณต๊ฐ ๋ฒํท์ ๋ค์ด๊ฐ๊ณ , ์ฃผ์๋ ํ ์๊ฐ์ง๋ฆฌ ์๋ช
์ฃผ์์
๋๋ค. ์ฃผ์๊ฐ ์ฃฝ์ด๋
ํ์ผ์ ๋จ์ ์์ผ๋ฏ๋ก ๋ ์ฝ๋์ ์ค๋ฆฐ file_id ๋ก ๋ค์ ๋ฐ๊ธ๋ฐ์ต๋๋ค.
stored = await store.url_for(tenant_id="acme", file_id="artifacts/<uuid>.docx")
stored.url # ์ ์๋ช
์ฃผ์
stored.expires_at # ์ ๋ง๋ฃ ์๊ฐ๋ ์ฝ๋์๋ file_id ์ url_expires_at ์ด ํจ๊ป ์ค๋ฆฝ๋๋ค. ์ด ๋์ด ์์ผ๋ฉด ํ๋ฉด์ ์ฃผ์๊ฐ
์ฃฝ์ ๋ค ํ ์ ์๋ ์ผ์ด ์์ต๋๋ค โ ๋ง๋ฃ ์๊ฐ์ด ๋น์ด ์์ผ๋ฉด "๋ง๋ฃ๊ฐ ์๋ค"๋ ๋ป์
๋๋ค
(์ ๊ณต๊ฐ ๋ฒํท์ ์๊ตฌ ์ฃผ์).
from processgpt_agent_sdk.artifacts import (
ArtifactCollector, MementoArtifactStore, download_file, strip_local_paths,
)
collector = ArtifactCollector(MementoArtifactStore(MEMENTO_BASE_URL))
# ํด์ ์์ํ ๋ ์ฐ๋๋ค โ ์ด๋ฒ ํด์ด ๋ง๋ ๊ฒ๋ง ๊ฐ๋ ค๋ด๊ธฐ ์ํด.
before = collector.snapshot(outputs_dir)
# ... ์์ด์ ํธ๊ฐ outputs/ ์ ์ต์ข
๋ณธ์ ๋๋๋ค ...
files = await collector.collect(
outputs_dir,
tenant_id=tenant_id,
conversation_id=conversation_id,
before=before,
turn_id=turn_id,
)
# ๋ณธ๋ฌธ์ ๋จ์ ๋ด๋ถ ๊ฒฝ๋ก๋ฅผ ๋ฐ์ ์ ์๋ ์ฃผ์๋ก ๋ฐ๊พผ๋ค.
answer = strip_local_paths(answer, files)๊ฑฐ๋ ๋ ์ฝ๋๋ ์ต์ข
์ด๋ฒคํธ์ metadata ์ pdfFiles ๋ก ์ฃ์ต๋๋ค.
ParseDict({"role": "assistant", "pdfFiles": [download_file(f) for f in files]},
artifact_evt.metadata)๊ทธ๋ฌ๋ฉด SDK ๊ฐ ๋ ๊ฐ์ง๋ฅผ ํจ๊ป ์ฒ๋ฆฌํฉ๋๋ค.
chats.messages.pdfFiles๋ก ์ ์ฅ โ ๋ฐฉ์ ๋ค์ ์ด์ด๋ ํ์ผ์ด ๋จ์ต๋๋ค.- SSE
done์files๋ก ์ค์ด ๋ณด๋ โ ํ๋ฉด์ด ์คํธ๋ฆฌ๋ฐ ๋์ ์ฐ๋ ํ์๋ ๋ค์ด๊ฐ๋๋ค.
โ ๏ธ done์ ์ฃ์ง ์์ผ๋ฉด ํ๋ฉด์ด ์๊ธฐ๊ฐ ๋ณธ ๊ฒ์ผ๋ก ์ด ํ์๋ ์ฐ์ถ๋ฌผ์ด ์์ต๋๋ค. ๋ฐฉ์ ๋ค์ ์ด์์ ๋ ๊ทธ ํ์ด ์ด๊ธฐ๋ฉด ๋ง๋ ํ์ผ์ด ์ฌ๋ผ์ง ๊ฒ์ฒ๋ผ ๋ณด์ ๋๋ค. ์ด์์์ ์ค์ ๋ก ๊ทธ๋ ๊ฒ ๋๊ฐ์ต๋๋ค.
| ๊ธฐ์ค | ๊ฐ | ์ด์ |
|---|---|---|
| ํ์ | DEFAULT_EXTENSIONS |
์คํ ํ์ผยท์์คยท๋ก๊ทธ๋ ๊ฒฐ๊ณผ๋ฌผ์ด ์๋๋ค |
| ํฌ๊ธฐ | 50MB (DEFAULT_MAX_BYTES) |
์ด๋ณด๋ค ํฌ๋ฉด ์ ๋ก๋๊ฐ ํด์ ๋ถ์ก๊ณ ํ๋ฉด์์๋ ๋ฐ๋ค๊ฐ ๋๊ธด๋ค |
| ๋น ํ์ผ | ์ ์ธ |
collect(..., index=False) ๋ ์์ธ์ ๊ฑด๋๋๋๋ค. ๋ฏธ๋ฆฌ๋ณด๊ธฐ์ฒ๋ผ ์ฌ๋์ด ํ ๋ฒ ๋ณด๊ณ ๋ง๋
๋ถ์ฐ๋ฌผ์ ์๋๋ค โ ์ด์์์ ๋ฏธ๋ฆฌ๋ณด๊ธฐ PDF ๋ฅผ ์์ธํ๋ค๊ฐ ํด์ด 5๋ถ์ฉ ๋ฉ์ถ ์ ์ด ์์ต๋๋ค.
decorate ๋ก ๋ ์ฝ๋๋ง๋ค ์ ์ฅ์๋ณ ํ์(๋ฏธ๋ฆฌ๋ณด๊ธฐ ๋ ๋ยท๊ฒ์ ํ์ )๋ฅผ ๋ง๋ถ์ผ ์ ์์ต๋๋ค.
Memento ์ชฝ์ ๊ณต๊ฐ ์ ์ฑ ์ด ์๋ ๋ฒํท์ด ํ๋ ์์ด์ผ ํฉ๋๋ค. ์์ผ๋ฉด ์ฐ์ถ๋ฌผ ์ ๋ก๋๊ฐ ์ ๋ถ ์คํจํฉ๋๋ค.
| ํ๊ฒฝ๋ณ์ | ๊ธฐ๋ณธ๊ฐ | ๋ป |
|---|---|---|
ARTIFACT_BUCKET |
artifacts |
์ฐ์ถ๋ฌผ ์ ์ฉ ๋น๊ณต๊ฐ ๋ฒํท |
ARTIFACT_URL_TTL_SECONDS |
3600 |
์๋ช ์ฃผ์ ์๋ช (์ด) |
์์ธํ ๋ด์ฉ์ Memento ์ ์ฅ์์ docs/artifact-bucket.md ์ ์์ต๋๋ค.
์์ปค๊ฐ todolist ํ ๊ฑด์ ์ง์ผ๋ฉด ๊ทธ ํ์ draft_status='STARTED' ๊ฐ ๋ฉ๋๋ค. ์์ ์๋
๊ฑฐ๊ธฐ์ ๋์ด์์ต๋๋ค โ ๋ง๋ฃ๊ฐ ์์์ต๋๋ค. ์์ปค๊ฐ kill -9 ๋ก ์ฃฝ์ผ๋ฉด(OOM, ๋
ธ๋ ์ถ์ถ,
KEDA ์ถ์) ๊ทธ ํ์ STARTED ๋ก ์๊ตฌํ ๋จ๊ณ , ์ง๋ ์กฐ๊ฑด(draft_status IS NULL ๋๋
FB_REQUESTED)์ ๋ค์ ๊ฑธ๋ฆฌ์ง ์์ ์๋ฌด๋ ์ง์ง ์์ต๋๋ค. ์ฌ์ฉ์์๊ฒ๋ ์์ํ
๋๋์ง ์๋ ์์
์ผ๋ก ๋ณด์
๋๋ค.
์์ด์ ํธ ์ฝ๋๊ฐ ๋ฐ๋ก ํ ์ผ์ ์์ต๋๋ค. ProcessGPTAgentServer ๊ฐ ์์์ ํฉ๋๋ค.
- ์ง์ ๋ lease ๋ฅผ ๊ฑด๋ค โ
fetch_pending_task์p_lease_seconds๋ฅผ ๋๊น๋๋ค. - ์ํ ์ค ์ฐ์ฅํ๋ค โ
LeaseKeeper๊ฐ ๋ณ๋ OS ์ค๋ ๋์์renew_task_lease๋ฅผ ์ฃผ๊ธฐ์ ์ผ๋ก ๋ถ๋ฆ ๋๋ค. ์ต์คํํฐ๊ฐ ๋๊ธฐ ํธ์ถ๋ก ์ด๋ฒคํธ ๋ฃจํ๋ฅผ ๋ถ์ก๊ณ ์์ด๋ ์ฐ์ฅ์ด ๋ฉ์ถ์ง ์์์ผ ํ๊ธฐ ๋๋ฌธ์ ๋๋ค(๋ฉ์ถ๋ฉด ์ด์์ ์ผํ๋ ์ค์ธ ์์ ์ด ํ์๋์ด ๋ ๋ฒ ์ํ๋ฉ๋๋ค). - ํ์๋นํ๋ฉด ๋ฒ๋ฆฐ๋ค โ ์ฐ์ฅ์ด
not_owner๋ก ๊ฑฐ์ ๋๋ฉด(๋ค๋ฅธ ์์ปค๊ฐ ์ด๋ฏธ ๊ฐ์ ธ๊ฐ๋ค) ์งํ ์ค์ธexecute()๋ฅผ ์ทจ์ํ๊ณ , FAILED ๋ก ๋งํนํ์ง ์์ต๋๋ค. ๊ทธ ์์ ์ ์ด์ ๋จ์ ๊ฒ์ด๊ณ , FAILED ๋ก ๋ฎ์ผ๋ฉด ๊ทธ์ชฝ์ด ๋๋ธ ๊ฒฐ๊ณผ๋ฅผ ์คํจ๋ก ๋ฐ๊ฟ๋๋ค. - ๋ ๋ ๋ ํด์ ํ๋ค โ ์ ์ ์ข
๋ฃ ์
release_task_lease๋ก ์ ์ ๋ฅผ ์ฆ์ ๋น์๋๋ค.
COMPLETEDยทHUMAN_ASKEDยทCANCELLED ๋ก ๋์ด๊ฐ ์์
์ ์ฐ์ฅ ์คํจ๋ "๋ฒ๋ ค๋ผ" ๊ฐ
์๋๋๋ค(not_started). heartbeat ๋ง ๋ฉ์ถ๊ณ ์์
์ ๊ทธ๋๋ก ๋ก๋๋ค โ ์ฌ๋ ๋ต๋ณ์
๊ธฐ๋ค๋ฆฌ๋ ์์
์ด lease ๋ง๋ฃ๋ก ํ์๋์ง ์๋ ๊ฒ๋ ๊ฐ์ ์ด์ ์
๋๋ค(RPC ๊ฐ STARTED
๋ง ํ์ํฉ๋๋ค).
| ํ๊ฒฝ๋ณ์ | ๊ธฐ๋ณธ๊ฐ | ๋ป |
|---|---|---|
TASK_LEASE_SECONDS |
120 | ์ ์ ๊ฐ ์ ์ง๋๋ ์๊ฐ |
TASK_LEASE_HEARTBEAT_SECONDS |
lease/4 (=30) | ์ฐ์ฅ ์ฃผ๊ธฐ |
TASK_MAX_CLAIMS |
3 | ํ ์์ ์ด ์ ์ ๋ ์ ์๋ ํ์(์ต์ด + ํ์ 2ํ). ๋์ผ๋ฉด RPC ๊ฐ FAILED ๋ก ์ข ๊ฒฐ |
์ฃผ๊ธฐ๋ฅผ lease ์ 1/4 ๋ก ๋ ๊ฒ์, ์ผ์์ DB ์ค๋ฅ๋ก ์ฐ์ ์ธ ๋ฒ ๋์ณ๋ lease ๊ฐ ๋จ์ ์๊ฒ ํ๊ธฐ ์ํด์์ ๋๋ค. 1/2 ๋ก ๋๋ฉด ํ ๋ฒ ๋์น๋ ๊ฒ๋ง์ผ๋ก ๋ง๋ฃ์ ๋ฟ์ ๋ฉ์ฉกํ ์์ ์ด ํ์๋ฉ๋๋ค.
todolist ์ lease_until timestamptz ์ claim_count integer ๊ฐ ์์ด์ผ ํ๊ณ ,
fetch_pending_task ๊ฐ p_lease_seconds/p_max_claims ๋ฅผ ๋ฐ์์ผ ํฉ๋๋ค. ์คํค๋ง์
RPC, ๊ทธ๋ฆฌ๊ณ kind ํด๋ฌ์คํฐ ์ค์ธก ๊ฒฐ๊ณผ๋ process-gpt-infra-docker ์ ์ฅ์์ ์์ต๋๋ค
(volumes/db/{init.sql,migration.sql}, tests/lease/README.md).
๊ตฌ๋ฒ์ SDK ์ ์์ฌ ๋์๋ ๋ฉ๋๋ค. p_lease_seconds ๋ฅผ ๋๊ธฐ์ง ์๋ ํธ์ถ์ lease ์์ด
์ง๊ณ (= ์์ ๋์), lease ๊ฐ ๋น์ด ์๋ ์ ์ ๋ ํ์ ๋์์ด ์๋๋๋ค.
- ./release.sh ๋ฒ์
- ์ค๋ฅ ๋ฐ์์ : python -m ensurepip --upgrade
- ์คํ ๋ฆฌ์ง ์
๋ก๋ ์ ํธ์
processgpt_agent_sdk.integrations.storage๋ก ๋ถ๋ฆฌ๋์์ต๋๋ค. - ๊ธฐ์กด
processgpt_agent_sdk.utils.upload_file_to_bucket,upload_files_to_bucket๋ ํ์ํธํ์ฉ์ผ๋ก ์ ์ง๋์ง๋ง deprecated ์ ๋๋ค. - ์ ๊ท ์ฝ๋๋ ์๋ ๊ฒฝ๋ก๋ฅผ ์ฌ์ฉํ์ธ์:
from processgpt_agent_sdk.integrations.storage import upload_file_to_bucket, upload_files_to_bucket
| ๋ชจ๋ | ๋ค์ด ์๋ ๊ฒ |
|---|---|
processgpt_agent_sdk.tenant_auth |
authorize_tenant, tenant_guard, auth_guard, request_tenant_id, ChatTenantGuardMiddleware, TenantAuthError |
processgpt_agent_sdk.chat_registry |
ChatRunRegistry, InflightRegistry, set_run_registry, set_inflight_registry |
processgpt_agent_sdk.chat_sse |
with_heartbeat, apply_heartbeat, format_sse_message, make_attach_handler, make_stop_handler |
์
๋ค ์ต์์(from processgpt_agent_sdk import ...)๋ก๋ ๋
ธ์ถ๋ฉ๋๋ค. ๋๋ถ๋ถ์ ๊ฒฝ์ฐ
mount_chat_routes() ํ๋๋ฉด ์ถฉ๋ถํ๊ณ , ์ด ๋ชจ๋๋ค์ ๋ผ์ฐํธ๋ฅผ ์ง์ ์กฐ๋ฆฝํ๊ฑฐ๋ ๋ ์ง์คํธ๋ฆฌ๋ฅผ
๊ต์ฒดํ ๋๋ง ์๋๋ค.