Skip to content

Commit 66c5de7

Browse files
authored
refactor: move substrate runtime into v2 (#2596)
## Summary - move the Substrate client, gateway, config, and list helpers under `core/v2` - remove the legacy `sandboxbackend` implementation and dormant session actor lifecycle - update v2 controller and system/session services to use the v2 package ## Testing - `go test ./core/v2/substrate ./core/v2/agentinstance ./core/internal/service/session ./core/internal/service/system ./core/cmd/controller-v2` - `go test -skip 'TestE2E.*' ./core/...`\n- `git diff --check` Signed-off-by: Eitan Yarmush <eitan.yarmush@solo.io>
1 parent 221674d commit 66c5de7

37 files changed

Lines changed: 40 additions & 2849 deletions

go/core/cmd/controller-v2/main.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -37,12 +37,12 @@ import (
3737
taskservice "github.com/kagent-dev/kagent/go/core/internal/service/task"
3838
"github.com/kagent-dev/kagent/go/core/pkg/auth"
3939
"github.com/kagent-dev/kagent/go/core/pkg/migrations"
40-
legacysubstrate "github.com/kagent-dev/kagent/go/core/pkg/sandboxbackend/substrate"
4140
"github.com/kagent-dev/kagent/go/core/v2/a2agateway"
4241
"github.com/kagent-dev/kagent/go/core/v2/agentinstance"
4342
"github.com/kagent-dev/kagent/go/core/v2/checkpoint"
4443
v2controller "github.com/kagent-dev/kagent/go/core/v2/controller"
4544
v2mcp "github.com/kagent-dev/kagent/go/core/v2/mcp"
45+
"github.com/kagent-dev/kagent/go/core/v2/substrate"
4646
"golang.org/x/sync/errgroup"
4747
k8sruntime "k8s.io/apimachinery/pkg/runtime"
4848
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
@@ -105,7 +105,7 @@ func main() {
105105
log.Fatalf("add reconciler to controller manager: %v", err)
106106
}
107107

108-
actors, err := legacysubstrate.Dial(ctx, legacysubstrate.Config{
108+
actors, err := substrate.Dial(ctx, substrate.Config{
109109
AteAPIEndpoint: env("SUBSTRATE_ATE_API_ENDPOINT", "dns:///api.ate-system.svc:443"),
110110
CAFile: os.Getenv("SUBSTRATE_ATE_API_CA_FILE"),
111111
ClientCertFile: os.Getenv("SUBSTRATE_ATE_API_CLIENT_CERT_FILE"),
@@ -122,7 +122,7 @@ func main() {
122122
instances := agentinstance.NewService(store, authorizer, instanceWorkflow)
123123
checkpoints := checkpoint.NewService(store, authorizer, actors, instanceWorkflow)
124124
gatewayDialer, err := a2agateway.NewRuntimeDialer(
125-
env("SUBSTRATE_ATENET_ROUTER_URL", legacysubstrate.DefaultAtenetRouterURL),
125+
env("SUBSTRATE_ATENET_ROUTER_URL", substrate.DefaultAtenetRouterURL),
126126
authenticator,
127127
)
128128
if err != nil {

go/core/internal/service/session/service.go

Lines changed: 6 additions & 83 deletions
Original file line numberDiff line numberDiff line change
@@ -14,9 +14,6 @@ import (
1414
"github.com/kagent-dev/kagent/go/core/internal/service/serviceerrors"
1515
"github.com/kagent-dev/kagent/go/core/internal/utils"
1616
"github.com/kagent-dev/kagent/go/core/pkg/auth"
17-
apierrors "k8s.io/apimachinery/pkg/api/errors"
18-
"sigs.k8s.io/controller-runtime/pkg/client"
19-
ctrllog "sigs.k8s.io/controller-runtime/pkg/log"
2017
)
2118

2219
type Store interface {
@@ -34,15 +31,9 @@ type Store interface {
3431
DeleteSessionShare(context.Context, string, string, string) error
3532
}
3633

37-
type SandboxActorCleaner interface {
38-
DeleteSandboxAgentSessionActor(context.Context, *v1alpha3.SandboxAgent, string) (bool, error)
39-
}
40-
4134
type Service struct {
42-
store Store
43-
kube client.Client
44-
actorCleaner SandboxActorCleaner
45-
token func() (string, error)
35+
store Store
36+
token func() (string, error)
4637
}
4738

4839
type Option func(*Service)
@@ -80,13 +71,6 @@ func NewService(store Store, options ...Option) *Service {
8071
return service
8172
}
8273

83-
func WithSandboxLifecycle(kube client.Client, cleaner SandboxActorCleaner) Option {
84-
return func(service *Service) {
85-
service.kube = kube
86-
service.actorCleaner = cleaner
87-
}
88-
}
89-
9074
func WithShareTokenGenerator(generator func() (string, error)) Option {
9175
return func(service *Service) {
9276
if generator != nil {
@@ -155,18 +139,12 @@ func (s *Service) Create(ctx context.Context, request CreateRequest) (*database.
155139
return nil, serviceerrors.NewInvalidArgument(fmt.Sprintf("Agent ref is invalid, please check the agent ref %s", request.AgentRef), err)
156140
}
157141
if agent.WorkloadType == v1alpha3.WorkloadModeSandbox {
158-
_, isSubstrateSandbox, err := s.lookupSubstrateSandboxAgent(ctx, request.AgentRef)
142+
existing, err := s.store.ListSessionsForAgentAllUsers(ctx, agentID)
159143
if err != nil {
160-
return nil, serviceerrors.NewInternal("Failed to inspect sandbox agent", err)
144+
return nil, serviceerrors.NewInternal("Failed to list sessions for agent", err)
161145
}
162-
if !isSubstrateSandbox {
163-
existing, err := s.store.ListSessionsForAgentAllUsers(ctx, agentID)
164-
if err != nil {
165-
return nil, serviceerrors.NewInternal("Failed to list sessions for agent", err)
166-
}
167-
if len(existing) > 0 {
168-
return nil, serviceerrors.NewAlreadyExists("Sandbox agents support only one chat session", fmt.Errorf("a session already exists for this agent"))
169-
}
146+
if len(existing) > 0 {
147+
return nil, serviceerrors.NewAlreadyExists("Sandbox agents support only one chat session", fmt.Errorf("a session already exists for this agent"))
170148
}
171149
}
172150

@@ -261,23 +239,9 @@ func (s *Service) Delete(ctx context.Context, sessionID string) error {
261239
if s.store == nil {
262240
return serviceerrors.NewInternal("Failed to delete session", fmt.Errorf("database client is not configured"))
263241
}
264-
265-
var cleanup *v1alpha3.SandboxAgent
266-
if s.actorCleaner != nil {
267-
if session, getErr := s.store.GetSession(ctx, sessionID, userID); getErr == nil && session != nil && session.AgentID != nil {
268-
if sandboxAgent, lookupErr := s.substrateSandboxAgentForSession(ctx, session); lookupErr == nil {
269-
cleanup = sandboxAgent
270-
}
271-
}
272-
}
273242
if err := s.store.DeleteSession(ctx, sessionID, userID); err != nil {
274243
return serviceerrors.NewInternal("Failed to delete session", err)
275244
}
276-
if cleanup != nil {
277-
if _, err := s.actorCleaner.DeleteSandboxAgentSessionActor(ctx, cleanup, sessionID); err != nil {
278-
ctrllog.FromContext(ctx).Error(err, "failed to delete substrate session actor", "sessionID", sessionID)
279-
}
280-
}
281245
return nil
282246
}
283247

@@ -393,47 +357,6 @@ func (s *Service) DeleteShare(ctx context.Context, sessionID, token string) erro
393357
return nil
394358
}
395359

396-
func (s *Service) lookupSubstrateSandboxAgent(ctx context.Context, agentRef string) (*v1alpha3.SandboxAgent, bool, error) {
397-
if s.kube == nil {
398-
return nil, false, nil
399-
}
400-
ref := strings.TrimSpace(agentRef)
401-
if ref == "" {
402-
return nil, false, nil
403-
}
404-
kubernetesRef := utils.ConvertToKubernetesIdentifier(ref)
405-
namespacedName, err := utils.ParseRefString(kubernetesRef, "")
406-
if err != nil {
407-
return nil, false, nil
408-
}
409-
sandboxAgent := &v1alpha3.SandboxAgent{}
410-
if err := s.kube.Get(ctx, namespacedName, sandboxAgent); err != nil {
411-
if apierrors.IsNotFound(err) {
412-
return nil, false, nil
413-
}
414-
return nil, false, err
415-
}
416-
return sandboxAgent, true, nil
417-
}
418-
419-
func (s *Service) substrateSandboxAgentForSession(ctx context.Context, session *database.Session) (*v1alpha3.SandboxAgent, error) {
420-
if session == nil || session.AgentID == nil {
421-
return nil, nil
422-
}
423-
agent, err := s.store.GetAgent(ctx, *session.AgentID)
424-
if err != nil {
425-
return nil, err
426-
}
427-
if agent.WorkloadType != v1alpha3.WorkloadModeSandbox {
428-
return nil, nil
429-
}
430-
sandboxAgent, isSubstrate, err := s.lookupSubstrateSandboxAgent(ctx, *session.AgentID)
431-
if err != nil || !isSubstrate {
432-
return nil, err
433-
}
434-
return sandboxAgent, nil
435-
}
436-
437360
func authenticatedPrincipal(ctx context.Context) (auth.Principal, error) {
438361
session, ok := auth.AuthSessionFrom(ctx)
439362
if !ok || session == nil {

go/core/internal/service/system/service.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ import (
1212
"github.com/kagent-dev/kagent/go/core/internal/service/serviceerrors"
1313
"github.com/kagent-dev/kagent/go/core/internal/version"
1414
"github.com/kagent-dev/kagent/go/core/pkg/auth"
15-
"github.com/kagent-dev/kagent/go/core/pkg/sandboxbackend/substrate"
15+
"github.com/kagent-dev/kagent/go/core/v2/substrate"
1616
corev1 "k8s.io/api/core/v1"
1717
apierrors "k8s.io/apimachinery/pkg/api/errors"
1818
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -330,8 +330,8 @@ func (s *Service) listSubstrateCRs(ctx context.Context, namespace string) ([]Sub
330330
WorkerSelector: labelSelectorString(ctx, actorTemplate.Spec.WorkerSelector),
331331
ManagedByKagent: actorTemplate.Labels["app.kubernetes.io/managed-by"] == "kagent",
332332
}
333-
if agentName := substrate.SandboxAgentNameFromLabels(actorTemplate.Labels); agentName != "" {
334-
entry.HarnessName = agentName
333+
if harnessName := actorTemplate.Labels[substrate.RevisionHarnessLabel]; harnessName != "" {
334+
entry.HarnessName = harnessName
335335
}
336336
actorTemplates = append(actorTemplates, entry)
337337
}

go/core/internal/service/system/service_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ import (
1111
"github.com/kagent-dev/kagent/go/core/internal/service/serviceerrors"
1212
"github.com/kagent-dev/kagent/go/core/internal/service/system"
1313
pkgAuth "github.com/kagent-dev/kagent/go/core/pkg/auth"
14-
"github.com/kagent-dev/kagent/go/core/pkg/sandboxbackend/substrate"
14+
"github.com/kagent-dev/kagent/go/core/v2/substrate"
1515
"github.com/stretchr/testify/assert"
1616
"github.com/stretchr/testify/require"
1717
corev1 "k8s.io/api/core/v1"
@@ -126,7 +126,7 @@ func TestGetSubstrateStatus(t *testing.T) {
126126
&atev1alpha1.ActorTemplate{
127127
ObjectMeta: metav1.ObjectMeta{Namespace: "team", Name: "template", Labels: map[string]string{
128128
"app.kubernetes.io/managed-by": "kagent",
129-
substrate.SandboxAgentLabelKey: "agent",
129+
substrate.RevisionHarnessLabel: "agent",
130130
}},
131131
Spec: atev1alpha1.ActorTemplateSpec{SandboxClass: atev1alpha1.SandboxClassGvisor},
132132
Status: atev1alpha1.ActorTemplateStatus{Phase: atev1alpha1.PhaseReady},

go/core/pkg/sandboxbackend/async.go

Lines changed: 0 additions & 56 deletions
This file was deleted.

go/core/pkg/sandboxbackend/backend.go

Lines changed: 0 additions & 36 deletions
This file was deleted.

go/core/pkg/sandboxbackend/filter_translator_owned.go

Lines changed: 0 additions & 75 deletions
This file was deleted.

0 commit comments

Comments
 (0)