1
2
3
4
5
6
7
8
9
10
11
12
13
14
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 }