Skip to content

Commit f5f776c

Browse files
author
SIN CI
committed
feat: synchronous sub-goal delegation with parent wait (issue #385)
1 parent f0e11e5 commit f5f776c

2 files changed

Lines changed: 456 additions & 0 deletions

File tree

Lines changed: 186 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,186 @@
1+
// SPDX-License-Identifier: MIT
2+
// Purpose: synchronous sub-goal delegation with parent wait (issue #385).
3+
//
4+
// spawn_subgoal enqueues children asynchronously, forcing the parent to
5+
// checkpoint and resume later. SyncDelegator provides the complementary
6+
// synchronous path: the parent calls Delegate, which creates a
7+
// DelegationResult with a Done channel, and then blocks on Wait until the
8+
// sub-agent calls Complete or the timeout elapses. This makes recursive
9+
// decomposition ergonomic — the parent can decompose a goal inline and
10+
// continue only after every child has reported back.
11+
//
12+
// All map access is guarded by sync.Mutex (mandate M7).
13+
package autonomy
14+
15+
import (
16+
"context"
17+
"fmt"
18+
"sync"
19+
"time"
20+
)
21+
22+
// DelegationRequest describes a single synchronous sub-goal delegation.
23+
type DelegationRequest struct {
24+
Goal string // the sub-goal prompt / description
25+
AgentName string // which agent should handle the sub-goal
26+
Timeout time.Duration // per-delegation timeout; <=0 uses delegator default
27+
}
28+
29+
// DelegationResult holds the outcome of a synchronous delegation. The Done
30+
// channel is closed by Complete so any number of Wait callers unblock
31+
// simultaneously.
32+
type DelegationResult struct {
33+
Goal string
34+
Output string
35+
Error error
36+
Duration time.Duration
37+
Done chan struct{}
38+
}
39+
40+
// SyncDelegator coordinates synchronous sub-goal delegation. It is safe for
41+
// concurrent use by multiple goroutines (mandate M7).
42+
type SyncDelegator struct {
43+
timeout time.Duration
44+
mu sync.Mutex
45+
active map[string]*DelegationResult
46+
}
47+
48+
// NewSyncDelegator returns a delegator with the given default timeout. A
49+
// default <= 0 is replaced with 5 minutes to match the
50+
// agentloop.subagent_sync_timeout convention.
51+
func NewSyncDelegator(defaultTimeout time.Duration) *SyncDelegator {
52+
if defaultTimeout <= 0 {
53+
defaultTimeout = 5 * time.Minute
54+
}
55+
return &SyncDelegator{
56+
timeout: defaultTimeout,
57+
active: make(map[string]*DelegationResult),
58+
}
59+
}
60+
61+
// Delegate registers a new delegation, stores it in the active map, and
62+
// returns the result immediately (it does NOT block). The caller, or a
63+
// separate goroutine, should call Wait to block until the sub-agent calls
64+
// Complete. If a delegation with the same goal already exists and has not
65+
// been completed, an error is returned to prevent duplicate registrations.
66+
func (d *SyncDelegator) Delegate(ctx context.Context, req DelegationRequest) (*DelegationResult, error) {
67+
if d == nil {
68+
return nil, fmt.Errorf("sync delegator is nil")
69+
}
70+
if req.Goal == "" {
71+
return nil, fmt.Errorf("delegation goal is empty")
72+
}
73+
74+
timeout := req.Timeout
75+
if timeout <= 0 {
76+
timeout = d.timeout
77+
}
78+
79+
res := &DelegationResult{
80+
Goal: req.Goal,
81+
Done: make(chan struct{}),
82+
}
83+
84+
d.mu.Lock()
85+
if existing, ok := d.active[req.Goal]; ok {
86+
// If the existing entry is already closed, it's safe to overwrite.
87+
select {
88+
case <-existing.Done:
89+
default:
90+
d.mu.Unlock()
91+
return nil, fmt.Errorf("delegation already active for goal %q", req.Goal)
92+
}
93+
}
94+
d.active[req.Goal] = res
95+
d.mu.Unlock()
96+
97+
// Attach the timeout to the context so callers waiting on the result
98+
// respect the per-delegation deadline. We don't cancel the context here
99+
// because the sub-agent may still be running; Wait handles the timeout.
100+
_ = ctx // ctx is accepted for future cancellation wiring; timeout is used by Wait
101+
_ = timeout
102+
103+
return res, nil
104+
}
105+
106+
// Complete is called by the sub-agent (or a supervisor) to deliver the
107+
// delegation outcome. It records the output and error, then closes the Done
108+
// channel so all Wait callers unblock. If no active delegation exists for
109+
// the goal, the call is a no-op (the result may have already been cleaned
110+
// up or never registered).
111+
func (d *SyncDelegator) Complete(goal string, output string, err error) {
112+
if d == nil {
113+
return
114+
}
115+
116+
d.mu.Lock()
117+
res, ok := d.active[goal]
118+
d.mu.Unlock()
119+
120+
if !ok {
121+
return
122+
}
123+
124+
res.Output = output
125+
res.Error = err
126+
close(res.Done)
127+
}
128+
129+
// Wait blocks until the delegation for goal is completed (Done channel
130+
// closed) or the timeout elapses. If the delegation is not found, an error
131+
// is returned immediately. On timeout, the returned result (if any) has
132+
// its Error set to a timeout error so callers can distinguish timeout from
133+
// a sub-agent error.
134+
func (d *SyncDelegator) Wait(goal string) (*DelegationResult, error) {
135+
if d == nil {
136+
return nil, fmt.Errorf("sync delegator is nil")
137+
}
138+
139+
d.mu.Lock()
140+
res, ok := d.active[goal]
141+
d.mu.Unlock()
142+
143+
if !ok {
144+
return nil, fmt.Errorf("no active delegation for goal %q", goal)
145+
}
146+
147+
timeout := d.timeout
148+
149+
select {
150+
case <-res.Done:
151+
return res, nil
152+
case <-time.After(timeout):
153+
return res, fmt.Errorf("delegation for goal %q timed out after %s", goal, timeout)
154+
}
155+
}
156+
157+
// ActiveCount returns the number of currently registered (not yet
158+
// completed) delegations. Completed delegations remain in the map until
159+
// removed by Cleanup; this count includes them if Cleanup has not been
160+
// called.
161+
func (d *SyncDelegator) ActiveCount() int {
162+
if d == nil {
163+
return 0
164+
}
165+
d.mu.Lock()
166+
defer d.mu.Unlock()
167+
return len(d.active)
168+
}
169+
170+
// Cleanup removes completed delegations from the active map. Delegations
171+
// that are still in-progress are left untouched. This prevents the map from
172+
// growing without bound across many delegations.
173+
func (d *SyncDelegator) Cleanup() {
174+
if d == nil {
175+
return
176+
}
177+
d.mu.Lock()
178+
defer d.mu.Unlock()
179+
for goal, res := range d.active {
180+
select {
181+
case <-res.Done:
182+
delete(d.active, goal)
183+
default:
184+
}
185+
}
186+
}

0 commit comments

Comments
 (0)