From 5ff3c35a915aa6c86c54869c98e67fd0790e0be5 Mon Sep 17 00:00:00 2001 From: ymq Date: Sat, 8 Aug 2026 21:32:04 +0800 Subject: [PATCH] =?UTF-8?q?v3.1:=20clean=20NDJSON=20streaming=20=E2=80=94?= =?UTF-8?q?=20no=20SSE=20prefix,=20no=20[DONE],=20content-type=20text/plai?= =?UTF-8?q?n?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Backend: yields raw JSON+\n, stream_response with text/plain Frontend: XHR.onprogress splits on \n, parses each line as JSON --- wwwroot/api/cockpit_chat.dspy | 59 +++++++++++++---------------------- wwwroot/index_cockpit.ui | 2 +- 2 files changed, 22 insertions(+), 39 deletions(-) diff --git a/wwwroot/api/cockpit_chat.dspy b/wwwroot/api/cockpit_chat.dspy index 3febc1d..e1324ba 100644 --- a/wwwroot/api/cockpit_chat.dspy +++ b/wwwroot/api/cockpit_chat.dspy @@ -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: \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: diff --git a/wwwroot/index_cockpit.ui b/wwwroot/index_cockpit.ui index fdf7398..a162eb1 100644 --- a/wwwroot/index_cockpit.ui +++ b/wwwroot/index_cockpit.ui @@ -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=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