Metadata-Version: 2.5
Name: weft-publish
Version: 0.0.7
Summary: Python Agent Runtime — long-lived session orchestration, event replay, ask-user interrupt/resume
Author: GodweiLL
License-Expression: MIT
License-File: LICENSE
Keywords: agent,asyncio,langgraph,llm,runtime,websocket
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Topic :: Software Development :: Libraries :: Application Frameworks
Requires-Python: >=3.11
Requires-Dist: pydantic>=2.6
Provides-Extra: dev
Requires-Dist: pyright>=1.1; extra == 'dev'
Requires-Dist: ruff>=0.6; extra == 'dev'
Provides-Extra: fastapi
Requires-Dist: fastapi>=0.110; extra == 'fastapi'
Requires-Dist: websockets>=12; extra == 'fastapi'
Provides-Extra: langchain
Requires-Dist: langchain>=1.0; extra == 'langchain'
Requires-Dist: langgraph>=1.1; extra == 'langchain'
Provides-Extra: test
Requires-Dist: pytest-asyncio>=0.23; extra == 'test'
Requires-Dist: pytest-cov>=5; extra == 'test'
Requires-Dist: pytest>=8; extra == 'test'
Description-Content-Type: text/markdown

# weft

> Python Agent Runtime — long-lived session orchestration, event replay, ask-user interrupt/resume.

`weft` 把"一次 HTTP 调用"模型升级为**长生命周期、可订阅、可重连**的 Agent 会话。Transport 层和 Agent 框架解耦——LangGraph / OpenAI Agents SDK / 裸 LLM SDK 都能接入。

## 核心能力

| 能力 | 实现 |
|---|---|
| 长生命周期 turn task | 客户端断连只 detach listener,agent 后台续跑 |
| 断网/换 tab/F5 不丢事件 | Ring buffer + `last_seq` 重连补帧协议 |
| 多 listener 同时订阅 | 每条 listener 独立有界 queue + 后台 pump,慢/死 listener 不反压 emit |
| ask-user 中断 → resume | sync 工具线程内调 ContextVar future,async 主 loop 收答复后续跑 |
| Cancel-safe | turn cancel 时把当前 user input 抢救进 agent 自己的 checkpoint |
| Janitor GC | idle runner 自动回收;turn 跑中即使无 listener 也保留 |

## 跟生态的关系

| 对比 | 关系 |
|---|---|
| **LangGraph / OpenAI Agents SDK** | weft 在它们之上一层,管 graph 之外的 session/lifecycle/transport;通过 `AgentProtocol` 接入,核心包不 import 它们 |
| **Temporal / Sidekiq** | 单进程 asyncio 范围内的 agent 会话;不抢跨进程 workflow engine 位置 |
| **Langfuse / OpenTelemetry** | weft 暴露 middleware 数据 + emit hook,observability backend 自己接 |

## 30 秒 demo

```python
import asyncio
from weft import (
    AskUserHandler, CancelToken, EventEmitter,
    MainBlockStart, MainBlockDelta, MainBlockEnd,
    ThreadRunner, WSListener,
)


class EchoAgent:
    async def run_turn(
        self, user_input: str, *,
        emit: EventEmitter, askuser: AskUserHandler, cancel_token: CancelToken,
    ) -> None:
        await emit(MainBlockStart(block_id="b1", block_type="text"))
        for ch in user_input:
            if cancel_token.cancelled:
                break
            await emit(MainBlockDelta(block_id="b1", delta=ch))
        await emit(MainBlockEnd(block_id="b1"))


async def main():
    runner = ThreadRunner("t-1", EchoAgent())

    captured = []
    async def send(payload): captured.append(payload)

    listener = WSListener(send)
    await runner.attach_listener(listener)
    await runner.start_turn("hi", listener)
    assert runner._turn_task is not None
    await runner._turn_task

asyncio.run(main())
```

完整 WS server demo:`examples/hello_echo/`。

## FastAPI 适配器

`weft.adapters.fastapi.run_ws_session` 把上面那段"收 client → 路由到 runner"接收循环抽成一行调用; 业务侧扩展走鸭子类型 `WSSessionHooks`,全部方法可选,不实现等于 no-op。

```python
from fastapi import FastAPI, WebSocket
from weft import RunnerJanitor, RunnerRegistry
from weft.adapters.fastapi import run_ws_session

app = FastAPI()
registry = RunnerRegistry(my_agent_factory)
janitor = RunnerJanitor(registry); janitor.start()

class Hooks:
    async def on_attach(self, listener):
        listener.enqueue(my_usage_snapshot())   # 补一帧客户端 hydrate 用的快照
    async def on_resume_submitted(self, answers):
        await persist_clarify(answers)          # clarify 答案落业务库
    async def on_config(self, msg):
        await update_role_models(msg)           # ConfigUpdate 业务字段由 hook 解释

@app.websocket("/ws/{thread_id}")
async def ws(ws: WebSocket, thread_id: str):
    runner = await registry.get_or_create(thread_id)
    await run_ws_session(ws, thread_id, runner, hooks=Hooks())
```

`run_ws_session` 负责: ws.accept / 推 ReadyEvent / 收 hello 触发 attach 补帧 / 路由 user_message / cancel / resume / compact 到 runner / detach 关 listener / pump 死亡时主动 close ws。安装: `pip install "weft[fastapi]"`。

## 架构

```
┌─────────────────────────────────────────────────────────────┐
│  Transport adapter (FastAPI WS / SSE / 自定义)             │
│  ↓ attach_listener  ↑ user_message/resume/cancel            │
├─────────────────────────────────────────────────────────────┤
│  ThreadRunner                                              │
│  ├─ state machine: idle / streaming / awaiting_resume       │
│  ├─ ring buffer + compute_replay (last_seq 补帧)            │
│  ├─ WSListener[] (per-transport queue + pump, 背压隔离)     │
│  └─ ask-user future / cancel salvage                        │
├─────────────────────────────────────────────────────────────┤
│  AgentProtocol  (你的实现, 或 weft.adapters.langgraph)      │
│  └─ run_turn(emit, askuser, cancel_token)                   │
└─────────────────────────────────────────────────────────────┘
```

`RunnerRegistry` 按 key 索引 runner,`RunnerJanitor` 周期回收 idle。

## 安装

```bash
pip install weft                    # 核心 (零 langgraph/langchain 依赖)
pip install "weft[langchain]"       # + LangGraph adapter (规划中)
pip install "weft[fastapi]"         # + FastAPI WS adapter (规划中)
```

## Status

**v0.0.1 — alpha**。核心 transport/lifecycle 层稳定,middleware 套装 + LangGraph adapter 在 v0.1 完成。

详细设计见 [`docs/architecture.md`](docs/architecture.md),事件协议见 [`docs/protocol.md`](docs/protocol.md)。

## License

MIT
