@@ -10,25 +10,63 @@ import (
1010 "io"
1111 "os"
1212 "sync"
13+ "sync/atomic"
14+
15+ `github.com/agentcodinglab/aicodingagentteam/internal/types`
1316)
1417
1518// Session represents an ACP agent session.
1619type Session struct {
1720 ID string `json:"id"`
1821 Status string `json:"status"` // active / paused / stopped
22+ mu sync.Mutex
23+ tasks map [string ]* Task // active task IDs for this session
24+ }
25+
26+ // Task is a single newTask execution within a session (ADR-0020).
27+ type Task struct {
28+ ID string `json:"taskId"`
29+ SessionID string `json:"sessionId"`
30+ Status string `json:"status"` // pending / running / done / error
31+ Events []TaskEvent `json:"events"`
32+ }
33+
34+ // TaskEvent is a single streamed event for a newTask.
35+ type TaskEvent struct {
36+ Type string `json:"type"` // start / message / tool_call / done / error
37+ Content string `json:"content"`
38+ }
39+
40+ // DirectorLike is the minimal contract the ACP server needs from the coordinator.
41+ // It lets us dispatch user requests without an import cycle.
42+ type DirectorLike interface {
43+ Handle (ctx context.Context , req types.UserRequest ) (* types.Delivery , error )
1944}
2045
2146// Server handles ACP session lifecycle for standard agent clients.
2247type Server struct {
23- mu sync.Mutex
24- sessions map [string ]* Session
48+ mu sync.Mutex
49+ sessions map [string ]* Session
50+ director DirectorLike // optional: when set, session/newTask dispatches to it
51+ taskCounter atomic.Int64
52+ notifier func (method string , params interface {}) // optional: write notifications to client (e.g. notifications/session/update)
2553}
2654
27- // New creates an ACP Server.
55+ // New creates an ACP Server without a director (lifecycle-only) .
2856func New () * Server {
2957 return & Server {sessions : make (map [string ]* Session )}
3058}
3159
60+ // NewWithDirector creates an ACP Server backed by a DirectorLike (ADR-0020).
61+ // notifier is invoked for each streamed event; it must be safe to call from any goroutine.
62+ func NewWithDirector (d DirectorLike , notifier func (method string , params interface {})) * Server {
63+ return & Server {
64+ sessions : make (map [string ]* Session ),
65+ director : d ,
66+ notifier : notifier ,
67+ }
68+ }
69+
3270// jsonRPCRequest is a JSON-RPC 2.0 request envelope.
3371type jsonRPCRequest struct {
3472 JSONRPC string `json:"jsonrpc"`
@@ -121,6 +159,8 @@ func (s *Server) handleMethod(ctx context.Context, req jsonRPCRequest) jsonRPCRe
121159 return s .handleSessionStop (ctx , req )
122160 case "session/list" :
123161 return s .handleSessionList (ctx , req )
162+ case "session/newTask" :
163+ return s .handleSessionNewTask (ctx , req )
124164 default :
125165 return jsonRPCResponse {
126166 JSONRPC : "2.0" ,
@@ -191,6 +231,102 @@ func (s *Server) handleSessionList(ctx context.Context, req jsonRPCRequest) json
191231 }
192232}
193233
234+
235+ // sessionNewTaskParams holds arguments for session/newTask (ADR-0020).
236+ type sessionNewTaskParams struct {
237+ SessionID string `json:"sessionId"`
238+ AgentID string `json:"agentId"`
239+ Prompt string `json:"prompt"`
240+ }
241+
242+ // sessionUpdateParams is the payload of notifications/session/update events.
243+ type sessionUpdateParams struct {
244+ SessionID string `json:"sessionId"`
245+ TaskID string `json:"taskId"`
246+ Event TaskEvent `json:"event"`
247+ }
248+
249+ // handleSessionNewTask creates a new task in a session and dispatches the
250+ // request to the wired Director in a goroutine. Events are streamed via the
251+ // notifier (if set) as JSON-RPC notifications/session/update.
252+ func (s * Server ) handleSessionNewTask (ctx context.Context , req jsonRPCRequest ) jsonRPCResponse {
253+ var params sessionNewTaskParams
254+ if err := json .Unmarshal (req .Params , & params ); err != nil {
255+ return jsonRPCResponse {JSONRPC : "2.0" , ID : req .ID , Error : & jsonRPCErr {Code : - 32602 , Message : "invalid params" }}
256+ }
257+ if s .director == nil {
258+ return jsonRPCResponse {JSONRPC : "2.0" , ID : req .ID , Error : & jsonRPCErr {Code : - 32000 , Message : "session/newTask requires Director" }}
259+ }
260+ s .mu .Lock ()
261+ sess , ok := s .sessions [params .SessionID ]
262+ if ! ok {
263+ s .mu .Unlock ()
264+ return jsonRPCResponse {JSONRPC : "2.0" , ID : req .ID , Error : & jsonRPCErr {Code : - 32001 , Message : "session not found: " + params .SessionID }}
265+ }
266+ taskID := fmt .Sprintf ("task-%d" , s .taskCounter .Add (1 ))
267+ task := & Task {ID : taskID , SessionID : params .SessionID , Status : "pending" }
268+ if sess .tasks == nil {
269+ sess .tasks = make (map [string ]* Task )
270+ }
271+ sess .tasks [taskID ] = task
272+ s .mu .Unlock ()
273+
274+ go s .dispatchTaskAsync (ctx , sess , task , params )
275+
276+ return jsonRPCResponse {
277+ JSONRPC : "2.0" ,
278+ ID : req .ID ,
279+ Result : map [string ]interface {}{
280+ "taskId" : taskID ,
281+ "status" : "pending" ,
282+ },
283+ }
284+ }
285+
286+ // dispatchTaskAsync runs the task and emits streamed events.
287+ func (s * Server ) dispatchTaskAsync (ctx context.Context , sess * Session , task * Task , params sessionNewTaskParams ) {
288+ s .emitUpdate (sess , task , "start" , params .Prompt )
289+
290+ delivery , err := s .director .Handle (ctx , types.UserRequest {
291+ Message : params .Prompt ,
292+ Backend : params .AgentID ,
293+ })
294+ if err != nil {
295+ task .Status = "error"
296+ s .emitUpdate (sess , task , "error" , err .Error ())
297+ return
298+ }
299+
300+ s .emitUpdate (sess , task , "message" , fmt .Sprintf ("planID=%s score=%d passed=%v" , delivery .PlanID , delivery .Score , delivery .Passed ))
301+ for _ , a := range delivery .Artifacts {
302+ s .emitUpdate (sess , task , "tool_call" , a )
303+ }
304+ if delivery .Passed {
305+ task .Status = "done"
306+ s .emitUpdate (sess , task , "done" , "ok" )
307+ } else {
308+ task .Status = "error"
309+ s .emitUpdate (sess , task , "error" , "delivery did not pass" )
310+ }
311+ }
312+
313+ // emitUpdate records the event on the task and (if wired) pushes a JSON-RPC
314+ // notifications/session/update to the client.
315+ func (s * Server ) emitUpdate (sess * Session , task * Task , eventType , content string ) {
316+ ev := TaskEvent {Type : eventType , Content : content }
317+ sess .mu .Lock ()
318+ task .Events = append (task .Events , ev )
319+ sess .mu .Unlock ()
320+ if s .notifier != nil {
321+ s .notifier ("notifications/session/update" , sessionUpdateParams {
322+ SessionID : sess .ID ,
323+ TaskID : task .ID ,
324+ Event : ev ,
325+ })
326+ }
327+ }
328+
329+
194330// Serve starts the ACP server over stdio JSON-RPC.
195331func (s * Server ) Serve (ctx context.Context ) error {
196332 return s .serveReader (ctx , os .Stdin , os .Stdout )
0 commit comments