v3.1: clean NDJSON streaming — no SSE prefix, no [DONE], content-type text/plain

Backend: yields raw JSON+\n, stream_response with text/plain
Frontend: XHR.onprogress splits on \n, parses each line as JSON
This commit is contained in:
ymq 2026-08-08 21:32:04 +08:00
parent 9822d18eeb
commit 5ff3c35a91
2 changed files with 22 additions and 39 deletions

View File

@ -1,4 +1,4 @@
# cockpit_chat.dspy - SDLC Agent Loop v3.1 (SSE streaming)
# cockpit_chat.dspy - SDLC Agent Loop v3.1 (streaming NDJSON)
import aiohttp
action = (params_kw or {}).get('action', 'list_messages')
@ -95,19 +95,15 @@ def _parse(raw):
except: pass
return {"action":"reply","message":raw}
# SSE helper: format as "data: <json>\n\n"
def _sse(data):
return f"data: {data}\n\n"
def _w_text(t): return {"widgettype":"Text","options":{"text":t,"css":"agent-text"}}
def _w_card(title, body, kind="success"):
colors = {"success":"#10b981","error":"#ef4444","reply":"#6366f1","progress":"#f59e0b"}
c = colors.get(kind,"#6b7280")
return {"widgettype":"VBox","options":{"css":"agent-card","style":{"borderLeft":f"3px solid {c}","padding":"8px 12px","margin":"4px 0"}},"subwidgets":[
{"widgettype":"Text","options":{"text":title,"css":"agent-card-title","style":{"fontWeight":"bold","fontSize":"13px","color":c}}},
return {"widgettype":"VBox","options":{"style":{"borderLeft":f"3px solid {c}","padding":"8px 12px","margin":"4px 0"}},"subwidgets":[
{"widgettype":"Text","options":{"text":title,"style":{"fontWeight":"bold","fontSize":"13px","color":c}}},
body
]}
def _w_progress(text): return {"widgettype":"Text","options":{"text":text,"css":"agent-progress","style":{"color":"#f59e0b","fontSize":"12px","padding":"2px 8px"}}}
def _w_progress(text): return {"widgettype":"Text","options":{"text":text,"style":{"color":"#f59e0b","fontSize":"12px","padding":"2px 8px"}}}
_tool_labels = {'switch_project':'切换项目','create_task':'创建任务','check_progress':'检查进度','list_tasks':'查询任务','add_repo':'关联仓库','list_repos':'查看仓库','add_bug':'报告Bug'}
@ -122,7 +118,6 @@ async def _exec_tool(sor, tool, params, ctx, uid, org_id):
return 'OK: 已切换到「'+proj.name+'」'
allp = await sor.sqlExe("SELECT name FROM sd_projects ORDER BY created_at DESC LIMIT 15",{})
return 'FAIL: 未找到。可用: '+', '.join([getattr(r,'name','') for r in (allp or [])])
elif tool == 'create_task':
if not ctx['pid']: return 'FAIL: 请先切换到项目'
title = p.get('title','')[:100]
@ -133,7 +128,6 @@ async def _exec_tool(sor, tool, params, ctx, uid, org_id):
rd = json.loads(r)
if rd.get('success'): return f"OK: 任务「{title}」已创建(角色:{role})"
return f"FAIL: {rd.get('message','')}"
elif tool == 'check_progress':
if not ctx['pid']: return 'FAIL: 请先切换到项目'
tid = p.get('task_id','')
@ -146,9 +140,8 @@ async def _exec_tool(sor, tool, params, ctx, uid, org_id):
ags = await sor.sqlExe("SELECT role_name,status FROM pipeline_project_agents WHERE project_id=${p}$",{"p":ctx['pid']})
al = ', '.join([f"{r.role_name}({r.status})" for r in (ags or [])]) or '未配置'
ds = await sor.sqlExe("SELECT deliverable_type,title,review_status FROM pipeline_deliverables WHERE project_id=${p}$ ORDER BY created_at DESC LIMIT 3",{"p":ctx['pid']})
dl = '\n'.join([f"[{d.deliverable_type}] {d.title[:40]}({d.review_status})" for d in (ds or [])]) or '无'
dl = '\n'.join([str(d.deliverable_type)+' '+str(d.title)[:40] for d in (ds or [])]) or '无'
return f"任务: {summary}\nAgent: {al}\n交付:\n{dl}"
elif tool == 'list_tasks':
if not ctx['pid']: return 'FAIL: 请先切换到项目'
sf = p.get('state','')
@ -158,7 +151,6 @@ async def _exec_tool(sor, tool, params, ctx, uid, org_id):
ts = await sor.sqlExe(sql, {"p":ctx['pid'],"s":sf} if sf else {"p":ctx['pid']})
em = {'submitted':'⏳','running':'🔄','review':'👀','approved':'✅','completed':'🏁','failed':'❌'}
return '\n'.join([f"{em.get(t.state,'?')}[{t.state}][{t.role}]{t.title}" for t in ts]) if ts else '无任务'
elif tool == 'add_repo':
if not ctx['pid']: return 'FAIL: 请先切换到项目'
url = p.get('url','')
@ -169,19 +161,16 @@ async def _exec_tool(sor, tool, params, ctx, uid, org_id):
await sor.C('sd_project_repos',{'id':getID(),'project_id':ctx['pid'],'repo_name':name,'repo_url':url,'default_branch':'main','local_path':'','org_id':org_id})
ctx['repos'].append({'n':name,'u':url})
return f"OK: 已关联 {name}"
elif tool == 'list_repos':
if not ctx['pid']: return 'FAIL: 请先切换到项目'
rs = await sor.sqlExe("SELECT repo_name,repo_url FROM sd_project_repos WHERE project_id=${p}$",{"p":ctx['pid']})
return '\n'.join([f"{r.repo_name}: {r.repo_url}" for r in rs]) if rs else '无仓库'
elif tool == 'add_bug':
if not ctx['pid']: return 'FAIL: 请先切换到项目'
title = p.get('title','')[:100]
if not title: return 'FAIL: 缺少标题'
await sor.C('sd_bugs',{'id':getID(),'iteration_id':'','title':title,'description':p.get('description',''),'severity':'major','priority':'P1','status':'open','reporter_type':'human','reporter_id':uid})
return f"OK: Bug已记录 {title}"
return f'未实现: {tool}'
except Exception as e:
return f'ERROR: {str(e)[:300]}'
@ -198,27 +187,23 @@ if action == 'send_message':
uid = await get_user()
org_id = await get_userorgid() or '0'
async def agent_sse():
async def agent_stream():
async with DBPools().sqlorContext(dbname) as sor:
ctx = await _load_ctx(sor, uid)
model = await _sel_model(sor, user_model_id)
if not model:
d = json.dumps({"type":"widget","widget":_w_card("❌ 错误",_w_text("无可用LLM"),"error")},ensure_ascii=False)
debug(f"SSE yield error: {d[:80]}")
yield _sse(d)
yield _sse("[DONE]")
d = json.dumps({"type":"widget","widget":_w_card("❌ 错误",_w_text("无可用LLM"),"error")},ensure_ascii=False)+'\n'
debug(f"YIELD error: {d[:80]}")
yield d
return
repos_str = ', '.join([r['n'] for r in ctx['repos']]) or '无'
env_text = f"项目: {ctx['pname'] or '未选择'}\n仓库: {repos_str}"
system = AGENT_PROMPT.replace('__TOOLS__',TOOLS_TEXT).replace('__ENV__',env_text)
msgs = [{"role":"system","content":system}]
history = await sor.sqlExe("SELECT role,content FROM pipeline_conversations ORDER BY created_at ASC LIMIT 20",{})
for h in (history or []):
for h in (await sor.sqlExe("SELECT role,content FROM pipeline_conversations ORDER BY created_at ASC LIMIT 20",{}) or []):
msgs.append({"role":'user' if getattr(h,'role','')=='user' else 'assistant',"content":getattr(h,'content','')})
msgs.append({"role":"user","content":message_text})
await sor.C('pipeline_conversations',{'id':getID(),'role':'user','content':message_text,'msg_type':'text','org_id':org_id,'created_by':uid})
agent_reply = ''
@ -230,28 +215,26 @@ if action == 'send_message':
if act.get('action') == 'tool_call':
tool = act.get('tool','')
label = _tool_labels.get(tool,tool)
d = json.dumps({"type":"widget","widget":_w_progress(f"🔄 {label}")},ensure_ascii=False)
debug(f"SSE yield progress: {d[:80]}")
yield _sse(d)
d = json.dumps({"type":"widget","widget":_w_progress(f"🔄 {label}")},ensure_ascii=False)+'\n'
debug(f"YIELD progress: {d[:60]}")
yield d
result = await _exec_tool(sor, tool, act.get('params',{}), ctx, uid, org_id)
ok = result.startswith("OK:")
d = json.dumps({"type":"widget","widget":_w_card(f"{'✅' if ok else '❌'} {label}",_w_text(result),"success" if ok else "error")},ensure_ascii=False)
debug(f"SSE yield result: {d[:80]}")
yield _sse(d)
d = json.dumps({"type":"widget","widget":_w_card(f"{'✅' if ok else '❌'} {label}",_w_text(result),"success" if ok else "error")},ensure_ascii=False)+'\n'
debug(f"YIELD result: {d[:60]}")
yield d
msgs.append({"role":"assistant","content":raw})
msgs.append({"role":"user","content":f"工具 {tool} 结果:\n{result}"})
continue
agent_reply = raw; break
if not agent_reply: agent_reply = "处理超时,请简化需求。"
d = json.dumps({"type":"widget","widget":_w_card("🤖 驾驶舱 Agent",_w_text(agent_reply),"reply")},ensure_ascii=False)
debug(f"SSE yield final: {d[:80]}")
yield _sse(d)
yield _sse("[DONE]")
if not agent_reply: agent_reply = "处理超时"
d = json.dumps({"type":"widget","widget":_w_card("🤖 驾驶舱 Agent",_w_text(agent_reply),"reply")},ensure_ascii=False)+'\n'
debug(f"YIELD final: {d[:60]}")
yield d
await sor.C('pipeline_conversations',{'id':getID(),'role':'agent','content':agent_reply,'msg_type':'text','org_id':org_id,'created_by':'system'})
return await stream_response(request, agent_sse, 'text/event-stream')
return await stream_response(request, agent_stream, 'text/plain')
elif action == 'list_messages':
async with DBPools().sqlorContext(dbname) as sor:

View File

@ -191,7 +191,7 @@
"event": "inputed",
"actiontype": "script",
"target": "self",
"script": "var p=params.prompt;var mid=params.model_id||'';if(typeof mid==='object'&&mid!==null)mid=mid.model_id||'';var ctx=bricks.getWidgetById('context_data',bricks.app);var cd=ctx?ctx.getValue():{};var iid=cd.iteration_id||'';var pid=cd.project_id||'';var fnames=params.file_names||[];var chat=bricks.getWidgetById('chat_scroll',bricks.app);var at=null;var _makeWidget=function(desc){if(!desc||!desc.widgettype)return null;var Klass=bricks.Factory.get(desc.widgettype);if(!Klass)return null;var opts=JSON.parse(JSON.stringify(desc.options||{}));if(desc.subwidgets&&desc.subwidgets.length>0){opts.subwidgets=desc.subwidgets.map(function(s){return _makeWidget(s)}).filter(function(s){return s!==null})}return new Klass(opts)};if(chat){var ub=new bricks.HBox({width:'100%'});var um=new bricks.VBox({width:'85%',alignSelf:'flex-end',bgcolor:'#dbeafe',borderRadius:'12px',padding:'12px 16px',marginBottom:'10px',gap:'4px'});um.add_widget(new bricks.Text({text:'你',cfontsize:0.75,color:'#2563eb',fontWeight:'bold'}));um.add_widget(new bricks.Text({text:p,cfontsize:0.95,color:'#1e293b',whiteSpace:'pre-wrap'}));for(var i=0;i<fnames.length;i++){var fn=fnames[i];var fl=new bricks.Text({text:'📎 '+fn,cfontsize:0.75,color:'#3b82f6',css:'file-tag',style:{background:'#eff6ff',padding:'2px 8px',borderRadius:'4px',display:'inline-block',marginTop:'4px'}});um.add_widget(fl)}ub.add_widget(new bricks.VBox({css:'filler'}));ub.add_widget(um);ub.add_widget(new bricks.Svg({rate:2,url:bricks_resource('imgs/chat-user.svg')}));chat.add_widget(ub);var ab=new bricks.HBox({width:'100%'});var am=new bricks.VBox({width:'85%',alignSelf:'flex-start',bgcolor:'#e8f0fe',borderRadius:'12px',padding:'12px 16px',marginBottom:'10px',gap:'4px'});am.add_widget(new bricks.Text({text:'Agent',cfontsize:0.75,color:'#3b82f6',fontWeight:'bold'}));at=new bricks.Text({text:'正在分析需求...',cfontsize:0.85,color:'#64748b'});am.add_widget(at);ab.add_widget(new bricks.Svg({rate:2,url:bricks_resource('imgs/llm.svg')}));ab.add_widget(am);ab.add_widget(new bricks.VBox({css:'filler'}));chat.add_widget(ab);}var fd=new FormData();fd.append('action','send_message');fd.append('message_text',p);fd.append('model_id',mid);fd.append('iteration_id',iid);fd.append('project_id',pid);fd.append('file_paths',JSON.stringify(fnames));var xhr=new XMLHttpRequest();xhr.open('POST','/pipeline-sdlc/api/cockpit_chat.dspy');var lastLen=0;var first=true;xhr.onprogress=function(){var t=xhr.responseText;var d=t.substring(lastLen);lastLen=t.length;var blocks=d.split('data: ');for(var j=1;j<blocks.length;j++){var block=blocks[j];var nn=block.indexOf('\n\n');var line=(nn>=0?block.substring(0,nn):block).trim();if(!line||line==='[DONE]')continue;try{var obj=JSON.parse(line);if(obj.type==='widget'&&obj.widget){var w=_makeWidget(obj.widget);if(w&&chat){if(first){am.remove_widget(at);first=false}chat.add_widget(w);chat.scrollToBottom()}}}catch(e){}}};xhr.onload=function(){if(at)at.set_text('已完成');var s=bricks.getWidgetById('stats_row',bricks.app);if(s)s.render({});};xhr.onerror=function(){if(at)at.set_text('网络错误')};xhr.send(fd)"
"script": "var p=params.prompt;var mid=params.model_id||'';if(typeof mid==='object'&&mid!==null)mid=mid.model_id||'';var ctx=bricks.getWidgetById('context_data',bricks.app);var cd=ctx?ctx.getValue():{};var iid=cd.iteration_id||'';var pid=cd.project_id||'';var fnames=params.file_names||[];var chat=bricks.getWidgetById('chat_scroll',bricks.app);var at=null;var _makeWidget=function(desc){if(!desc||!desc.widgettype)return null;var Klass=bricks.Factory.get(desc.widgettype);if(!Klass)return null;var opts=JSON.parse(JSON.stringify(desc.options||{}));if(desc.subwidgets&&desc.subwidgets.length>0){opts.subwidgets=desc.subwidgets.map(function(s){return _makeWidget(s)}).filter(function(s){return s!==null})}return new Klass(opts)};if(chat){var ub=new bricks.HBox({width:'100%'});var um=new bricks.VBox({width:'85%',alignSelf:'flex-end',bgcolor:'#dbeafe',borderRadius:'12px',padding:'12px 16px',marginBottom:'10px',gap:'4px'});um.add_widget(new bricks.Text({text:'你',cfontsize:0.75,color:'#2563eb',fontWeight:'bold'}));um.add_widget(new bricks.Text({text:p,cfontsize:0.95,color:'#1e293b',whiteSpace:'pre-wrap'}));for(var i=0;i<fnames.length;i++){var fn=fnames[i];var fl=new bricks.Text({text:'📎 '+fn,cfontsize:0.75,color:'#3b82f6',css:'file-tag',style:{background:'#eff6ff',padding:'2px 8px',borderRadius:'4px',display:'inline-block',marginTop:'4px'}});um.add_widget(fl)}ub.add_widget(new bricks.VBox({css:'filler'}));ub.add_widget(um);ub.add_widget(new bricks.Svg({rate:2,url:bricks_resource('imgs/chat-user.svg')}));chat.add_widget(ub);var ab=new bricks.HBox({width:'100%'});var am=new bricks.VBox({width:'85%',alignSelf:'flex-start',bgcolor:'#e8f0fe',borderRadius:'12px',padding:'12px 16px',marginBottom:'10px',gap:'4px'});am.add_widget(new bricks.Text({text:'Agent',cfontsize:0.75,color:'#3b82f6',fontWeight:'bold'}));at=new bricks.Text({text:'正在分析需求...',cfontsize:0.85,color:'#64748b'});am.add_widget(at);ab.add_widget(new bricks.Svg({rate:2,url:bricks_resource('imgs/llm.svg')}));ab.add_widget(am);ab.add_widget(new bricks.VBox({css:'filler'}));chat.add_widget(ab);}var fd=new FormData();fd.append('action','send_message');fd.append('message_text',p);fd.append('model_id',mid);fd.append('iteration_id',iid);fd.append('project_id',pid);fd.append('file_paths',JSON.stringify(fnames));var xhr=new XMLHttpRequest();xhr.open('POST','/pipeline-sdlc/api/cockpit_chat.dspy');var lastLen=0;var first=true;xhr.onprogress=function(){var t=xhr.responseText;var d=t.substring(lastLen);lastLen=t.length;var lines=d.split('\n');for(var j=0;j<lines.length;j++){var line=lines[j].trim();if(!line)continue;try{var obj=JSON.parse(line);if(obj.type=='widget'&&obj.widget){var w=_makeWidget(obj.widget);if(w&&chat){if(first){am.remove_widget(at);first=false}chat.add_widget(w);chat.scrollToBottom()}}}catch(e){}}};xhr.onload=function(){if(at)at.set_text('已完成');var s=bricks.getWidgetById('stats_row',bricks.app);if(s)s.render({})};xhr.onerror=function(){if(at)at.set_text('网络错误')};xhr.send(fd)"
}
]
}