返回 DeepSeek-Reasonix
delegation_parent_claim_test.go
根目录 / internal / agent / delegation_parent_claim_test.go
1 package agent
2
3 import (
4 "context"
5 "encoding/json"
6 "path/filepath"
7 "testing"
8 "time"
9
10 "reasonix/internal/event"
11 "reasonix/internal/provider"
12 "reasonix/internal/tool"
13 )
14
15 // wholeWorkspaceDelegate stands in for run_skill/task at depth 0: its Execute
16 // takes a whole-workspace writer slot and drives a depth-1 child that writes.
17 type wholeWorkspaceDelegate struct {
18 sched *SubagentScheduler
19 root string
20 acquireErr error
21 strayErr error
22 child toolOutcome
23 }
24
25 func (*wholeWorkspaceDelegate) Name() string { return "delegate" }
26 func (*wholeWorkspaceDelegate) Description() string { return "delegate" }
27 func (*wholeWorkspaceDelegate) Schema() json.RawMessage { return json.RawMessage(`{"type":"object"}`) }
28 func (*wholeWorkspaceDelegate) ReadOnly() bool { return false }
29
30 func (d *wholeWorkspaceDelegate) Execute(ctx context.Context, _ json.RawMessage) (string, error) {
31 whole, err := WholeWorkspaceWriteClaim(d.root)
32 if err != nil {
33 return "", err
34 }
35 release, id, err := d.sched.AcquireWithID(ctx, AcquireRequest{Writer: true, WritePaths: whole})
36 d.acquireErr = err
37 if err != nil {
38 return "", err
39 }
40 defer release()
41 _, d.strayErr = d.sched.Acquire(context.Background(), AcquireRequest{Writer: true, WritePaths: whole, Nested: true})
42
43 reg := tool.NewRegistry()
44 reg.Add(&recordingWriter{name: "write_file"})
45 child := New(nil, reg, NewSession(""), Options{
46 WriteScheduler: d.sched,
47 WriteWorkspaceRoot: d.root,
48 SubagentDepth: 1,
49 }, event.Discard)
50 d.child = child.executeOne(WithSubagentClaimID(ctx, id), &child.turn, provider.ToolCall{
51 ID: "child-write",
52 Name: "write_file",
53 Arguments: `{"path":"` + filepath.ToSlash(filepath.Join(d.root, "out.md")) + `","content":"x"}`,
54 })
55 return "delegated", nil
56 }
57
58 func TestDepthZeroDelegationIsNotBlockedByItsOwnHookParentClaim(t *testing.T) {
59 root := t.TempDir()
60 sched := NewSubagentScheduler(4, 2)
61 delegate := &wholeWorkspaceDelegate{sched: sched, root: root}
62 reg := tool.NewRegistry()
63 reg.Add(delegate)
64 a := New(nil, reg, NewSession(""), Options{
65 Hooks: &parentClaimProbeHooks{scheduler: sched},
66 WriteScheduler: sched,
67 WriteWorkspaceRoot: root,
68 }, event.Discard)
69
70 done := make(chan toolOutcome, 1)
71 go func() {
72 done <- a.executeOne(context.Background(), &a.turn, provider.ToolCall{ID: "d-1", Name: "delegate", Arguments: `{}`})
73 }()
74 var out toolOutcome
75 select {
76 case out = <-done:
77 case <-time.After(5 * time.Second):
78 t.Fatalf("depth-0 delegation still waiting on the parent claim its own tool call holds; queued claims=%d", len(sched.ActiveWriterClaims()))
79 }
80 if delegate.acquireErr != nil || out.errMsg != "" {
81 t.Fatalf("delegation failed: acquire=%v outcome=%+v", delegate.acquireErr, out)
82 }
83 if delegate.child.blocked || delegate.child.errMsg != "" {
84 t.Fatalf("child write under the delegated slot was blocked: %+v", delegate.child)
85 }
86 if delegate.strayErr == nil {
87 t.Fatal("an unrelated whole-workspace writer started while the delegation held the workspace")
88 }
89 if n := len(sched.ActiveWriterClaims()); n != 0 {
90 t.Fatalf("claims after delegation = %d, want 0", n)
91 }
92 }
93
94 func TestDelegationInsideParentClaimIsNotQueuedBehindWritersThatClaimBlocks(t *testing.T) {
95 root := t.TempDir()
96 sched := NewSubagentScheduler(4, 2)
97 whole, err := WholeWorkspaceWriteClaim(root)
98 if err != nil {
99 t.Fatal(err)
100 }
101 releaseParent, parentID, err := sched.ReserveParentWriteWithID(whole)
102 if err != nil {
103 t.Fatal(err)
104 }
105 queued := make(chan error, 1)
106 go func() {
107 release, err := sched.Acquire(context.Background(), AcquireRequest{Writer: true, WritePaths: whole})
108 if err == nil {
109 release()
110 }
111 queued <- err
112 }()
113 deadline := time.Now().Add(5 * time.Second)
114 for {
115 sched.mu.Lock()
116 n := len(sched.waiters)
117 sched.mu.Unlock()
118 if n == 1 {
119 break
120 }
121 if time.Now().After(deadline) {
122 t.Fatal("background writer never queued")
123 }
124 time.Sleep(time.Millisecond)
125 }
126
127 ctx, cancel := context.WithTimeout(WithParentWriteClaimID(context.Background(), parentID), 5*time.Second)
128 defer cancel()
129 releaseDelegate, err := sched.Acquire(ctx, AcquireRequest{Writer: true, WritePaths: whole})
130 if err != nil {
131 t.Fatalf("delegation queued behind a writer its own parent claim blocks: %v", err)
132 }
133 releaseDelegate()
134 releaseParent()
135 if err := <-queued; err != nil {
136 t.Fatalf("queued writer after parent release: %v", err)
137 }
138 }
139
139 lines GO