1   /*
2    * Copyright 2019 LINE Corporation
3    *
4    * LINE Corporation licenses this file to you under the Apache License,
5    * version 2.0 (the "License"); you may not use this file except in compliance
6    * with the License. You may obtain a copy of the License at:
7    *
8    *   https://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12   * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13   * License for the specific language governing permissions and limitations
14   * under the License.
15   */
16  package com.linecorp.centraldogma.server.internal.mirror;
17  
18  import static java.util.Objects.requireNonNull;
19  
20  import java.io.File;
21  import java.util.concurrent.CompletableFuture;
22  import java.util.concurrent.CompletionStage;
23  
24  import org.jspecify.annotations.Nullable;
25  import org.slf4j.Logger;
26  import org.slf4j.LoggerFactory;
27  
28  import com.google.common.base.MoreObjects;
29  
30  import com.linecorp.centraldogma.server.CentralDogmaConfig;
31  import com.linecorp.centraldogma.server.ZoneConfig;
32  import com.linecorp.centraldogma.server.mirror.MirroringServicePluginConfig;
33  import com.linecorp.centraldogma.server.plugin.Plugin;
34  import com.linecorp.centraldogma.server.plugin.PluginContext;
35  import com.linecorp.centraldogma.server.plugin.PluginTarget;
36  
37  public final class DefaultMirroringServicePlugin implements Plugin {
38  
39      private static final Logger logger = LoggerFactory.getLogger(DefaultMirroringServicePlugin.class);
40  
41      @Nullable
42      public static MirroringServicePluginConfig mirrorConfig(CentralDogmaConfig config) {
43          return (MirroringServicePluginConfig) config.pluginConfigMap().get(MirroringServicePluginConfig.class);
44      }
45  
46      @Nullable
47      private volatile MirrorSchedulingService mirroringService;
48  
49      @Nullable
50      private PluginTarget pluginTarget;
51  
52      @Override
53      public PluginTarget target(CentralDogmaConfig config) {
54          requireNonNull(config, "config");
55          if (pluginTarget != null) {
56              return pluginTarget;
57          }
58  
59          final MirroringServicePluginConfig mirrorConfig = mirrorConfig(config);
60          if (mirrorConfig != null && mirrorConfig.zonePinned()) {
61              pluginTarget = PluginTarget.ZONE_LEADER_ONLY;
62          } else {
63              pluginTarget = PluginTarget.LEADER_ONLY;
64          }
65          return pluginTarget;
66      }
67  
68      @Override
69      public synchronized CompletionStage<Void> start(PluginContext context) {
70          requireNonNull(context, "context");
71  
72          MirrorSchedulingService mirroringService = this.mirroringService;
73          if (mirroringService == null) {
74              final CentralDogmaConfig cfg = context.config();
75              final MirroringServicePluginConfig mirroringServicePluginConfig = mirrorConfig(cfg);
76              final int numThreads;
77              final int maxNumFilesPerMirror;
78              final long maxNumBytesPerMirror;
79              final ZoneConfig zoneConfig;
80  
81              if (mirroringServicePluginConfig != null) {
82                  numThreads = mirroringServicePluginConfig.numMirroringThreads();
83                  maxNumFilesPerMirror = mirroringServicePluginConfig.maxNumFilesPerMirror();
84                  maxNumBytesPerMirror = mirroringServicePluginConfig.maxNumBytesPerMirror();
85                  if (mirroringServicePluginConfig.zonePinned()) {
86                      zoneConfig = cfg.zone();
87                      assert zoneConfig != null : "zonePinned is enabled but no zone configuration found";
88                  } else {
89                      zoneConfig = null;
90                  }
91                  if (mirroringServicePluginConfig.trustedHostKeys().isEmpty()) {
92                      logger.warn("No 'trustedHostKeys' configured in the mirroring service plugin config. " +
93                                  "SSH mirror connections will accept any host key without verification.");
94                  }
95              } else {
96                  numThreads = MirroringServicePluginConfig.INSTANCE.numMirroringThreads();
97                  maxNumFilesPerMirror = MirroringServicePluginConfig.INSTANCE.maxNumFilesPerMirror();
98                  maxNumBytesPerMirror = MirroringServicePluginConfig.INSTANCE.maxNumBytesPerMirror();
99                  zoneConfig = null;
100                 logger.warn("No 'trustedHostKeys' configured in the mirroring service plugin config. " +
101                             "SSH mirror connections will accept any host key without verification.");
102             }
103             mirroringService = new MirrorSchedulingService(new File(cfg.dataDir(), "_mirrors"),
104                                                            context.projectManager(),
105                                                            context.meterRegistry(),
106                                                            numThreads,
107                                                            maxNumFilesPerMirror,
108                                                            maxNumBytesPerMirror, zoneConfig,
109                                                            context.mirrorAccessController());
110             this.mirroringService = mirroringService;
111         }
112         mirroringService.start(context.commandExecutor());
113         return CompletableFuture.completedFuture(null);
114     }
115 
116     @Override
117     public synchronized CompletionStage<Void> stop(PluginContext context) {
118         final MirrorSchedulingService mirroringService = this.mirroringService;
119         if (mirroringService != null && mirroringService.isStarted()) {
120             mirroringService.stop();
121         }
122         return CompletableFuture.completedFuture(null);
123     }
124 
125     @Override
126     public Class<?> configType() {
127         return MirroringServicePluginConfig.class;
128     }
129 
130     @Nullable
131     public MirrorSchedulingService mirroringService() {
132         return mirroringService;
133     }
134 
135     @Override
136     public String toString() {
137         return MoreObjects.toStringHelper(this)
138                           .omitNullValues()
139                           .add("configType", configType().getName())
140                           .add("target", pluginTarget)
141                           .toString();
142     }
143 }