@@ -605,63 +605,95 @@ private BiConsumer<Action, Pod> reconcilePodDistributedLogs() {
605605 }
606606
607607 private BiConsumer <Action , Pod > reconcilePodBackups () {
608- String backupNameKey =
609- StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .BACKUP_NAME_KEY ;
608+ String clusterNameKey =
609+ StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .CLUSTER_NAME_KEY ;
610610 return (action , pod ) -> synchronizedCopyOfValues (backups )
611611 .stream ()
612- .filter (cluster -> Objects .equals (
613- cluster .getMetadata ().getNamespace (),
612+ .filter (backup -> Objects .equals (
613+ backup .getMetadata ().getNamespace (),
614614 pod .getMetadata ().getNamespace ()))
615615 .filter (backup -> pod .getMetadata ().getLabels () != null )
616- .filter (backup -> Objects .equals (
617- pod .getMetadata ().getLabels ().get (backupNameKey ),
618- backup .getMetadata ().getName ()))
619- .forEach (backup -> reconcileBackup ().accept (action , backup ));
616+ .forEach (backup -> synchronizedCopyOfValues (clusters )
617+ .stream ()
618+ .filter (cluster -> Objects .equals (
619+ backup .getMetadata ().getNamespace (),
620+ cluster .getMetadata ().getNamespace ()))
621+ .filter (cluster -> Objects .equals (
622+ backup .getSpec ().getSgCluster (),
623+ cluster .getMetadata ().getName ()))
624+ .filter (cluster -> Objects .equals (
625+ pod .getMetadata ().getLabels ().get (clusterNameKey ),
626+ cluster .getMetadata ().getName ()))
627+ .forEach (cluster -> reconcileBackup ().accept (action , backup )));
620628 }
621629
622630 private BiConsumer <Action , Pod > reconcilePodDbOps () {
623- String dbOpsNameKey =
624- StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .DBOPS_NAME_KEY ;
631+ String clusterNameKey =
632+ StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .CLUSTER_NAME_KEY ;
625633 return (action , pod ) -> synchronizedCopyOfValues (dbOps )
626634 .stream ()
627635 .filter (cluster -> Objects .equals (
628636 cluster .getMetadata ().getNamespace (),
629637 pod .getMetadata ().getNamespace ()))
630638 .filter (dbOps -> pod .getMetadata ().getLabels () != null )
631- .filter (dbOps -> Objects .equals (
632- pod .getMetadata ().getLabels ().get (dbOpsNameKey ),
633- dbOps .getMetadata ().getName ()))
634- .forEach (dbOps -> reconcileDbOps ().accept (action , dbOps ));
639+ .forEach (dbOps -> synchronizedCopyOfValues (clusters )
640+ .stream ()
641+ .filter (cluster -> Objects .equals (
642+ dbOps .getMetadata ().getNamespace (),
643+ cluster .getMetadata ().getNamespace ()))
644+ .filter (cluster -> Objects .equals (
645+ dbOps .getSpec ().getSgCluster (),
646+ cluster .getMetadata ().getName ()))
647+ .filter (cluster -> Objects .equals (
648+ pod .getMetadata ().getLabels ().get (clusterNameKey ),
649+ cluster .getMetadata ().getName ()))
650+ .forEach (cluster -> reconcileDbOps ().accept (action , dbOps )));
635651 }
636652
637653 private BiConsumer <Action , Pod > reconcilePodShardedBackups () {
638- String backupNameKey =
639- StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .SHARDED_BACKUP_NAME_KEY ;
654+ String clusterNameKey =
655+ StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .SHARDED_CLUSTER_NAME_KEY ;
640656 return (action , pod ) -> synchronizedCopyOfValues (shardedBackups )
641657 .stream ()
642658 .filter (cluster -> Objects .equals (
643659 cluster .getMetadata ().getNamespace (),
644660 pod .getMetadata ().getNamespace ()))
645661 .filter (backup -> pod .getMetadata ().getLabels () != null )
646- .filter (backup -> Objects .equals (
647- pod .getMetadata ().getLabels ().get (backupNameKey ),
648- backup .getMetadata ().getName ()))
649- .forEach (backup -> reconcileShardedBackup ().accept (action , backup ));
662+ .forEach (backup -> synchronizedCopyOfValues (shardedClusters )
663+ .stream ()
664+ .filter (cluster -> Objects .equals (
665+ backup .getMetadata ().getNamespace (),
666+ cluster .getMetadata ().getNamespace ()))
667+ .filter (cluster -> Objects .equals (
668+ backup .getSpec ().getSgShardedCluster (),
669+ cluster .getMetadata ().getName ()))
670+ .filter (cluster -> Objects .equals (
671+ pod .getMetadata ().getLabels ().get (clusterNameKey ),
672+ cluster .getMetadata ().getName ()))
673+ .forEach (cluster -> reconcileShardedBackup ().accept (action , backup )));
650674 }
651675
652676 private BiConsumer <Action , Pod > reconcilePodShardedDbOps () {
653- String dbOpsNameKey =
654- StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .SHARDED_DBOPS_NAME_KEY ;
677+ String clusterNameKey =
678+ StackGresContext .STACKGRES_KEY_PREFIX + StackGresContext .SHARDED_CLUSTER_NAME_KEY ;
655679 return (action , pod ) -> synchronizedCopyOfValues (shardedDbOps )
656680 .stream ()
657681 .filter (cluster -> Objects .equals (
658682 cluster .getMetadata ().getNamespace (),
659683 pod .getMetadata ().getNamespace ()))
660684 .filter (dbOps -> pod .getMetadata ().getLabels () != null )
661- .filter (dbOps -> Objects .equals (
662- pod .getMetadata ().getLabels ().get (dbOpsNameKey ),
663- dbOps .getMetadata ().getName ()))
664- .forEach (dbOps -> reconcileShardedDbOps ().accept (action , dbOps ));
685+ .forEach (dbOps -> synchronizedCopyOfValues (shardedClusters )
686+ .stream ()
687+ .filter (cluster -> Objects .equals (
688+ dbOps .getMetadata ().getNamespace (),
689+ cluster .getMetadata ().getNamespace ()))
690+ .filter (cluster -> Objects .equals (
691+ dbOps .getSpec ().getSgShardedCluster (),
692+ cluster .getMetadata ().getName ()))
693+ .filter (cluster -> Objects .equals (
694+ pod .getMetadata ().getLabels ().get (clusterNameKey ),
695+ cluster .getMetadata ().getName ()))
696+ .forEach (cluster -> reconcileShardedDbOps ().accept (action , dbOps )));
665697 }
666698
667699 private BiConsumer <Action , Pod > reconcilePodStreams () {
0 commit comments