Skip to content

Commit 0a00a59

Browse files
committed
feat: implement SubstrateHarness for managing sandboxed actor execution via gRPC
1 parent d3dcd40 commit 0a00a59

2 files changed

Lines changed: 459 additions & 0 deletions

File tree

‎internal/harness/substrate.go‎

Lines changed: 185 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,185 @@
1+
// Copyright 2026 Google LLC
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
//go:build ate
16+
17+
package harness
18+
19+
import (
20+
"context"
21+
"errors"
22+
"fmt"
23+
"io"
24+
"log"
25+
"sync"
26+
"time"
27+
28+
"google.golang.org/grpc"
29+
"google.golang.org/grpc/credentials/insecure"
30+
"google.golang.org/grpc/metadata"
31+
32+
"github.com/google/ax/internal/experimental/k8s/ate"
33+
"github.com/google/ax/proto"
34+
"github.com/google/uuid"
35+
)
36+
37+
// SubstrateHarness manages execution in a SubstrATE sandboxed actor over gRPC HarnessService.
38+
type SubstrateHarness struct {
39+
ateClient *ate.Client
40+
port int
41+
dialOpts []grpc.DialOption
42+
}
43+
44+
// NewSubstrateHarness creates a new SubstrateHarness.
45+
func NewSubstrateHarness(endpoint string, namespace string, template string, port int, opts ...grpc.DialOption) (*SubstrateHarness, error) {
46+
if port == 0 {
47+
port = 50053 // Default HarnessService port
48+
}
49+
if namespace == "" {
50+
namespace = "ax"
51+
}
52+
if template == "" {
53+
template = "ax-harness-template"
54+
}
55+
client, err := ate.NewClient(namespace, template, endpoint)
56+
if err != nil {
57+
return nil, fmt.Errorf("failed to create ATE client: %w", err)
58+
}
59+
if len(opts) == 0 {
60+
opts = append(opts, grpc.WithTransportCredentials(insecure.NewCredentials()))
61+
}
62+
return &SubstrateHarness{
63+
ateClient: client,
64+
port: port,
65+
dialOpts: opts,
66+
}, nil
67+
}
68+
69+
// Start implements Harness interface. It creates/resumes the target actor.
70+
func (h *SubstrateHarness) Start(ctx context.Context, conversationID string) (Execution, error) {
71+
if conversationID == "" {
72+
return nil, errors.New("SubstrateHarness needs valid conversationID")
73+
}
74+
75+
resp, err := h.ateClient.CreateActor(ctx, conversationID)
76+
if err != nil {
77+
return nil, fmt.Errorf("failed to create substrate actor %s: %w", conversationID, err)
78+
}
79+
actor := resp.Actor
80+
if actor == nil {
81+
return nil, fmt.Errorf("received nil actor in response for %s", conversationID)
82+
}
83+
if actor.AteomPodIp == "" {
84+
return nil, fmt.Errorf("actor %s has no active worker IP address", conversationID)
85+
}
86+
87+
// 2. Establish connection to the actor's worker IP
88+
workerAddr := fmt.Sprintf("%s:%d", actor.AteomPodIp, h.port)
89+
conn, err := grpc.NewClient(workerAddr, h.dialOpts...)
90+
if err != nil {
91+
return nil, fmt.Errorf("failed to dial remote harness service at %s: %w", workerAddr, err)
92+
}
93+
94+
return &substrateExecution{
95+
harness: h,
96+
conversationID: conversationID,
97+
execID: uuid.NewString(),
98+
conn: conn,
99+
client: proto.NewHarnessServiceClient(conn),
100+
}, nil
101+
}
102+
103+
type substrateExecution struct {
104+
harness *SubstrateHarness
105+
conversationID string
106+
execID string
107+
conn *grpc.ClientConn
108+
client proto.HarnessServiceClient
109+
110+
mu sync.Mutex
111+
pending []*proto.Message
112+
}
113+
114+
func (e *substrateExecution) ID() string {
115+
return e.execID
116+
}
117+
118+
func (e *substrateExecution) Queue(ctx context.Context, msg ...*proto.Message) error {
119+
e.mu.Lock()
120+
defer e.mu.Unlock()
121+
e.pending = append(e.pending, msg...)
122+
return nil
123+
}
124+
125+
func (e *substrateExecution) Run(ctx context.Context, handler Handler) error {
126+
e.mu.Lock()
127+
inputs := e.pending
128+
e.pending = nil
129+
e.mu.Unlock()
130+
131+
// Append conversation ID to metadata for server tracking
132+
streamCtx := metadata.AppendToOutgoingContext(ctx, "ax-conversation-id", e.conversationID)
133+
stream, err := e.client.Connect(streamCtx)
134+
if err != nil {
135+
return fmt.Errorf("failed to open harness service stream: %w", err)
136+
}
137+
138+
// Send inputs
139+
err = stream.Send(&proto.HarnessMessage{
140+
Messages: inputs,
141+
})
142+
if err != nil {
143+
return fmt.Errorf("failed to send harness inputs: %w", err)
144+
}
145+
146+
// Close send direction to trigger server processing
147+
if err := stream.CloseSend(); err != nil {
148+
return fmt.Errorf("failed to close stream send direction: %w", err)
149+
}
150+
151+
// Receive response stream
152+
for {
153+
resp, err := stream.Recv()
154+
if err == io.EOF {
155+
break
156+
}
157+
if err != nil {
158+
return fmt.Errorf("error receiving from harness stream: %w", err)
159+
}
160+
for _, m := range resp.Messages {
161+
if err := handler.OnMessage(ctx, e.execID, m); err != nil {
162+
return err
163+
}
164+
}
165+
}
166+
167+
return handler.OnComplete(ctx, e.execID)
168+
}
169+
170+
func (e *substrateExecution) Close(ctx context.Context) error {
171+
// Close connection
172+
if e.conn != nil {
173+
e.conn.Close()
174+
}
175+
176+
// Suspend actor to return resource to standard standby pool
177+
log.Printf("Suspending SubstrATE actor for conversation %s (execution %s)", e.conversationID, e.execID)
178+
suspendCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
179+
defer cancel()
180+
if _, err := e.harness.ateClient.SuspendActor(suspendCtx, e.conversationID); err != nil {
181+
log.Printf("Failed to suspend actor %s: %v", e.conversationID, err)
182+
}
183+
184+
return nil
185+
}

0 commit comments

Comments
 (0)