stack_controller.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416
  1. /*
  2. Copyright 2026 LocoStack.
  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. http://www.apache.org/licenses/LICENSE-2.0
  7. Unless required by applicable law or agreed to in writing, software
  8. distributed under the License is distributed on an "AS IS" BASIS,
  9. WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  10. See the License for the specific language governing permissions and
  11. limitations under the License.
  12. */
  13. package controller
  14. import (
  15. "context"
  16. "fmt"
  17. "maps"
  18. "strings"
  19. corev1 "k8s.io/api/core/v1"
  20. apierrors "k8s.io/apimachinery/pkg/api/errors"
  21. apimeta "k8s.io/apimachinery/pkg/api/meta"
  22. "k8s.io/apimachinery/pkg/api/resource"
  23. metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
  24. "k8s.io/apimachinery/pkg/runtime"
  25. "k8s.io/apimachinery/pkg/types"
  26. "k8s.io/utils/ptr"
  27. ctrl "sigs.k8s.io/controller-runtime"
  28. "sigs.k8s.io/controller-runtime/pkg/client"
  29. "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
  30. "sigs.k8s.io/controller-runtime/pkg/handler"
  31. logf "sigs.k8s.io/controller-runtime/pkg/log"
  32. "sigs.k8s.io/controller-runtime/pkg/reconcile"
  33. "github.com/LocoStack/loco-operator/api/v1alpha1"
  34. "github.com/LocoStack/loco-operator/internal/reconciler"
  35. "github.com/LocoStack/loco-operator/pkg/templates"
  36. )
  37. // StackReconciler reconciles a Stack object
  38. type StackReconciler struct {
  39. client.Client
  40. Scheme *runtime.Scheme
  41. }
  42. // +kubebuilder:rbac:groups=locostack.com,resources=stacks,verbs=get;list;watch;create;update;patch;delete
  43. // +kubebuilder:rbac:groups=locostack.com,resources=stacks/status,verbs=get;update;patch
  44. // +kubebuilder:rbac:groups=locostack.com,resources=stacks/finalizers,verbs=update
  45. // +kubebuilder:rbac:groups=locostack.com,resources=components,verbs=get;list;watch;create;update;patch;delete
  46. // +kubebuilder:rbac:groups="",resources=persistentvolumes,verbs=get;list;watch;create;update;patch
  47. // +kubebuilder:rbac:groups="",resources=persistentvolumeclaims,verbs=get;list;watch;create;update;patch
  48. func (r *StackReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
  49. log := logf.FromContext(ctx)
  50. s := &v1alpha1.Stack{}
  51. if err := r.Get(ctx, req.NamespacedName, s); err != nil {
  52. if apierrors.IsNotFound(err) {
  53. return ctrl.Result{}, nil
  54. }
  55. return ctrl.Result{}, err
  56. }
  57. patch := client.MergeFrom(s.DeepCopy())
  58. if err := r.reconcileSharedVolumes(ctx, s); err != nil {
  59. return ctrl.Result{}, err
  60. }
  61. if err := r.reconcileComponents(ctx, s); err != nil {
  62. return ctrl.Result{}, err
  63. }
  64. if err := r.reconcileGateways(ctx, s); err != nil {
  65. return ctrl.Result{}, err
  66. }
  67. s.Status.ObservedGeneration = s.Generation
  68. if err := r.Status().Patch(ctx, s, patch); err != nil {
  69. return ctrl.Result{}, client.IgnoreNotFound(err)
  70. }
  71. log.Info("Reconciled Stack", "namespace", s.Namespace, "name", s.Name)
  72. return ctrl.Result{}, nil
  73. }
  74. // reconcileSharedVolumes ensures all PVC-backed shared volumes declared in the Stack spec exist.
  75. func (r *StackReconciler) reconcileSharedVolumes(ctx context.Context, s *v1alpha1.Stack) error {
  76. failedVolumes := make([]string, 0)
  77. for _, volume := range s.Spec.SharedVolumes {
  78. claimName := fmt.Sprintf("%s-%s", s.Name, volume)
  79. pvc := &corev1.PersistentVolumeClaim{
  80. ObjectMeta: metav1.ObjectMeta{
  81. Name: claimName,
  82. Namespace: s.Namespace,
  83. },
  84. }
  85. pv := &corev1.PersistentVolume{}
  86. if err := r.Get(ctx, client.ObjectKey{Name: volume}, pv); err != nil {
  87. if apierrors.IsNotFound(err) {
  88. apimeta.SetStatusCondition(&s.Status.Conditions, metav1.Condition{
  89. Type: "SharedVolumesResolved",
  90. Status: metav1.ConditionFalse,
  91. Reason: "ReconcileFailed",
  92. Message: fmt.Sprintf("Could not find shared volume %s: %v", volume, err),
  93. ObservedGeneration: s.Generation,
  94. })
  95. return nil
  96. }
  97. return err
  98. }
  99. _, err := controllerutil.CreateOrUpdate(ctx, r.Client, pvc, func() error {
  100. if pvc.Labels == nil {
  101. pvc.Labels = map[string]string{}
  102. }
  103. maps.Copy(pvc.Labels, reconciler.ResourceLabels(s.Name, "SharedVolume", volume, ""))
  104. pvc.Spec.VolumeName = volume
  105. pvc.Spec.StorageClassName = ptr.To("")
  106. pvc.Spec.AccessModes = []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}
  107. if pvc.Spec.Resources.Requests == nil {
  108. pvc.Spec.Resources.Requests = corev1.ResourceList{}
  109. }
  110. if _, ok := pvc.Spec.Resources.Requests[corev1.ResourceStorage]; !ok {
  111. pvc.Spec.Resources.Requests[corev1.ResourceStorage] = resource.MustParse("1Gi")
  112. }
  113. return controllerutil.SetControllerReference(s, pvc, r.Scheme)
  114. })
  115. if err != nil {
  116. failedVolumes = append(failedVolumes, volume)
  117. }
  118. }
  119. if len(failedVolumes) > 0 {
  120. apimeta.SetStatusCondition(&s.Status.Conditions, metav1.Condition{
  121. Type: "SharedVolumesResolved",
  122. Status: metav1.ConditionFalse,
  123. Reason: "ReconcileFailed",
  124. Message: fmt.Sprintf("Could not reconcile shared volume(s): %v", failedVolumes),
  125. ObservedGeneration: s.Generation,
  126. })
  127. } else {
  128. apimeta.SetStatusCondition(&s.Status.Conditions, metav1.Condition{
  129. Type: "SharedVolumesResolved",
  130. Status: metav1.ConditionTrue,
  131. Reason: "ReconcileSucceeded",
  132. Message: "All shared volumes reconciled successfully",
  133. ObservedGeneration: s.Generation,
  134. })
  135. }
  136. return nil
  137. }
  138. func (r *StackReconciler) reconcileComponents(ctx context.Context, s *v1alpha1.Stack) error {
  139. components := getComponents(s)
  140. failedComponents := make([]string, 0)
  141. errList := make([]error, 0)
  142. for _, comp := range components {
  143. resolvedTemplateName, reconciled, err := r.reconcileComponent(ctx, s, comp)
  144. setResolvedTemplateName(&s.Status, comp.name, resolvedTemplateName)
  145. if !reconciled {
  146. failedComponents = append(failedComponents, comp.name)
  147. }
  148. if err != nil {
  149. errList = append(errList, err)
  150. }
  151. }
  152. if len(failedComponents) > 0 {
  153. apimeta.SetStatusCondition(&s.Status.Conditions, metav1.Condition{
  154. Type: "ComponentsReconciled",
  155. Status: metav1.ConditionFalse,
  156. Reason: "ReconcileFailed",
  157. Message: fmt.Sprintf("%d component(s) failed to reconcile: %v", len(failedComponents), failedComponents),
  158. ObservedGeneration: s.Generation,
  159. })
  160. } else {
  161. apimeta.SetStatusCondition(&s.Status.Conditions, metav1.Condition{
  162. Type: "ComponentsReconciled",
  163. Status: metav1.ConditionTrue,
  164. Reason: "ReconcileSucceeded",
  165. Message: "All enabled components were reconciled successfully",
  166. ObservedGeneration: s.Generation,
  167. })
  168. }
  169. if len(errList) > 0 {
  170. return fmt.Errorf("failed with errors: %v", errList)
  171. }
  172. return nil
  173. }
  174. type componentConfig struct {
  175. name string
  176. defaultTmpl string
  177. template *v1alpha1.Template
  178. enabled bool
  179. }
  180. func newComponentConfig(name, defaultTmpl string, enabled bool, optional *v1alpha1.OptionalComponent) componentConfig {
  181. var tmpl *v1alpha1.Template
  182. if optional != nil {
  183. enabled = optional.Enabled
  184. tmpl = optional.Template
  185. }
  186. return componentConfig{
  187. name: name,
  188. defaultTmpl: defaultTmpl,
  189. template: tmpl,
  190. enabled: enabled,
  191. }
  192. }
  193. func getComponents(s *v1alpha1.Stack) []componentConfig {
  194. return []componentConfig{
  195. newComponentConfig("Gateway", "litellm", true, s.Spec.Gateway),
  196. newComponentConfig("VectorStore", "qdrant", true, s.Spec.VectorStore),
  197. newComponentConfig("GraphStore", "neo4j", false, s.Spec.GraphStore),
  198. newComponentConfig("Database", "postgresql", true, s.Spec.Database),
  199. newComponentConfig("Observability", "phoenix", false, s.Spec.Observability),
  200. }
  201. }
  202. func setResolvedTemplateName(status *v1alpha1.StackStatus, name string, templateName string) {
  203. switch name {
  204. case "Gateway":
  205. status.ComponentStatus.Gateway = templateName
  206. case "VectorStore":
  207. status.ComponentStatus.VectorStore = templateName
  208. case "GraphStore":
  209. status.ComponentStatus.GraphStore = templateName
  210. case "Database":
  211. status.ComponentStatus.Database = templateName
  212. case "Observability":
  213. status.ComponentStatus.Observability = templateName
  214. }
  215. }
  216. func (r *StackReconciler) reconcileComponent(ctx context.Context, s *v1alpha1.Stack, compCfg componentConfig) (string, bool, error) {
  217. compName := fmt.Sprintf("%s-%s", strings.ToLower(compCfg.name), s.Name)
  218. comp := &v1alpha1.Component{
  219. ObjectMeta: metav1.ObjectMeta{
  220. Name: compName,
  221. Namespace: s.Namespace,
  222. },
  223. }
  224. err := r.Get(ctx, client.ObjectKey{Namespace: s.Namespace, Name: compName}, comp)
  225. if !compCfg.enabled {
  226. tmpl := "disabled"
  227. if err == nil {
  228. if err := r.Delete(ctx, comp); err != nil {
  229. return tmpl, false, err
  230. }
  231. return tmpl, true, nil
  232. }
  233. if apierrors.IsNotFound(err) {
  234. return tmpl, true, nil
  235. }
  236. return tmpl, false, err
  237. }
  238. tmpl, err := templates.Manager.ResolveTemplate(compCfg.template, compCfg.defaultTmpl)
  239. if err != nil {
  240. return "unknown", false, err
  241. }
  242. _, err = controllerutil.CreateOrUpdate(ctx, r.Client, comp, func() error {
  243. if comp.Labels == nil {
  244. comp.Labels = make(map[string]string)
  245. }
  246. maps.Copy(comp.Labels, reconciler.ResourceLabels(s.Name, compCfg.name, s.Name, ""))
  247. comp.Labels["stack.locostack.com/default"] = "true"
  248. comp.Spec.Category = compCfg.name
  249. comp.Spec.Template = &tmpl
  250. comp.Spec.StackRef = &corev1.LocalObjectReference{Name: s.Name}
  251. return controllerutil.SetControllerReference(s, comp, r.Scheme)
  252. })
  253. available := apimeta.IsStatusConditionTrue(comp.Status.Conditions, "Available")
  254. return tmpl.Name, available, err
  255. }
  256. func (r *StackReconciler) reconcileGateways(ctx context.Context, s *v1alpha1.Stack) error {
  257. log := logf.FromContext(ctx)
  258. gwList := &v1alpha1.ComponentList{}
  259. err := r.List(ctx, gwList, client.InNamespace(s.Namespace), client.MatchingLabels{"locostack.com/stack": s.Name, "locostack.com/component": "Gateway"})
  260. if err != nil {
  261. return fmt.Errorf("failed to list gateway components: %v", err)
  262. }
  263. if len(gwList.Items) == 0 {
  264. log.Info("No gateways found for stack", "namespace", s.Namespace, "stack", s.Name)
  265. return nil
  266. }
  267. dependencies := make([]corev1.ObjectReference, 0)
  268. emList := &v1alpha1.ExternalModelList{}
  269. if err := r.List(ctx, emList, client.InNamespace(s.Namespace)); err != nil {
  270. return fmt.Errorf("Failed to list external models: %v", err)
  271. }
  272. for _, em := range emList.Items {
  273. if em.Spec.StackRef.Name == s.Name {
  274. dependencies = append(dependencies, corev1.ObjectReference{
  275. Kind: "ExternalModel",
  276. Name: em.Name,
  277. Namespace: em.Namespace,
  278. })
  279. }
  280. }
  281. mmList := &v1alpha1.ManagedModelList{}
  282. if err := r.List(ctx, mmList, client.InNamespace(s.Namespace)); err != nil {
  283. return fmt.Errorf("Failed to list managed models: %v", err)
  284. }
  285. for _, em := range mmList.Items {
  286. if em.Spec.StackRef.Name == s.Name {
  287. dependencies = append(dependencies, corev1.ObjectReference{
  288. Kind: "ManagedModel",
  289. Name: em.Name,
  290. Namespace: em.Namespace,
  291. })
  292. }
  293. }
  294. etList := &v1alpha1.ExternalToolList{}
  295. if err := r.List(ctx, etList, client.InNamespace(s.Namespace)); err != nil {
  296. return fmt.Errorf("Failed to list external tools: %v", err)
  297. }
  298. for _, em := range etList.Items {
  299. if em.Spec.StackRef.Name == s.Name {
  300. dependencies = append(dependencies, corev1.ObjectReference{
  301. Kind: "ExternalTool",
  302. Name: em.Name,
  303. Namespace: em.Namespace,
  304. })
  305. }
  306. }
  307. mtList := &v1alpha1.ManagedToolList{}
  308. if err := r.List(ctx, mtList, client.InNamespace(s.Namespace)); err != nil {
  309. return fmt.Errorf("Failed to list managed tools: %v", err)
  310. }
  311. for _, em := range mtList.Items {
  312. if em.Spec.StackRef.Name == s.Name {
  313. dependencies = append(dependencies, corev1.ObjectReference{
  314. Kind: "ManagedTool",
  315. Name: em.Name,
  316. Namespace: em.Namespace,
  317. })
  318. }
  319. }
  320. if len(dependencies) == 0 {
  321. return nil
  322. }
  323. for _, gw := range gwList.Items {
  324. gw.Spec.Dependencies = dependencies
  325. if err := r.Update(ctx, &gw); err != nil {
  326. return fmt.Errorf("Failed to update gateway %s with dependencies: %v", gw.Name, err)
  327. }
  328. }
  329. return nil
  330. }
  331. // SetupWithManager sets up the controller with the Manager.
  332. func (r *StackReconciler) SetupWithManager(mgr ctrl.Manager) error {
  333. mapToStack := func(ctx context.Context, obj client.Object) []reconcile.Request {
  334. var stackName string
  335. switch typed := obj.(type) {
  336. case *v1alpha1.ExternalModel:
  337. stackName = typed.Spec.StackRef.Name
  338. case *v1alpha1.ManagedModel:
  339. stackName = typed.Spec.StackRef.Name
  340. case *v1alpha1.ExternalTool:
  341. stackName = typed.Spec.StackRef.Name
  342. case *v1alpha1.ManagedTool:
  343. stackName = typed.Spec.StackRef.Name
  344. default:
  345. return nil
  346. }
  347. stackList := &v1alpha1.StackList{}
  348. if err := mgr.GetClient().List(ctx, stackList, client.InNamespace(obj.GetNamespace())); err != nil {
  349. return nil
  350. }
  351. var reqs []reconcile.Request
  352. for _, stack := range stackList.Items {
  353. if stackName == stack.Name {
  354. reqs = append(reqs, reconcile.Request{
  355. NamespacedName: types.NamespacedName{Name: stack.Name, Namespace: stack.Namespace},
  356. })
  357. }
  358. }
  359. return reqs
  360. }
  361. return ctrl.NewControllerManagedBy(mgr).
  362. For(&v1alpha1.Stack{}).
  363. Owns(&v1alpha1.Component{}).
  364. Owns(&corev1.PersistentVolumeClaim{}).
  365. Watches(&v1alpha1.ExternalModel{}, handler.EnqueueRequestsFromMapFunc(mapToStack)).
  366. Watches(&v1alpha1.ManagedModel{}, handler.EnqueueRequestsFromMapFunc(mapToStack)).
  367. Watches(&v1alpha1.ExternalTool{}, handler.EnqueueRequestsFromMapFunc(mapToStack)).
  368. Watches(&v1alpha1.ManagedTool{}, handler.EnqueueRequestsFromMapFunc(mapToStack)).
  369. Named("stack").
  370. Complete(r)
  371. }