-
Notifications
You must be signed in to change notification settings - Fork 1.6k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
使用eventbus来传递协议变更事件
- Loading branch information
Showing
2 changed files
with
20 additions
and
83 deletions.
There are no files selected for viewing
81 changes: 11 additions & 70 deletions
81
...ent/src/main/java/org/jetlinks/community/protocol/LazyInitManagementProtocolSupports.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,96 +1,37 @@ | ||
package org.jetlinks.community.protocol; | ||
|
||
import lombok.AccessLevel; | ||
import lombok.Getter; | ||
import lombok.Setter; | ||
import lombok.extern.slf4j.Slf4j; | ||
import org.jetlinks.core.ProtocolSupport; | ||
import org.jetlinks.core.cluster.ClusterManager; | ||
import org.jetlinks.supports.protocol.StaticProtocolSupports; | ||
import org.jetlinks.supports.protocol.management.ProtocolSupportDefinition; | ||
import org.jetlinks.core.event.EventBus; | ||
import org.jetlinks.supports.protocol.management.DefaultProtocolSupportManager; | ||
import org.jetlinks.supports.protocol.management.ProtocolSupportLoader; | ||
import org.jetlinks.supports.protocol.management.ProtocolSupportManager; | ||
import org.springframework.boot.CommandLineRunner; | ||
import org.springframework.core.Ordered; | ||
import org.springframework.core.annotation.Order; | ||
import reactor.core.publisher.Mono; | ||
|
||
import java.time.Duration; | ||
import java.util.Map; | ||
import java.util.concurrent.ConcurrentHashMap; | ||
import java.util.function.Consumer; | ||
|
||
@Slf4j | ||
@Getter | ||
@Setter | ||
@Order(Ordered.HIGHEST_PRECEDENCE) | ||
public class LazyInitManagementProtocolSupports extends StaticProtocolSupports implements CommandLineRunner { | ||
|
||
private ProtocolSupportManager manager; | ||
|
||
private ProtocolSupportLoader loader; | ||
|
||
private ClusterManager clusterManager; | ||
|
||
@Setter(AccessLevel.PRIVATE) | ||
private Map<String, String> configProtocolIdMapping = new ConcurrentHashMap<>(); | ||
|
||
private Duration loadTimeOut = Duration.ofSeconds(30); | ||
|
||
public void init() { | ||
|
||
clusterManager.<ProtocolSupportDefinition>getTopic("_protocol_changed") | ||
.subscribe() | ||
.subscribe(protocol -> this.init(protocol).subscribe()); | ||
|
||
try { | ||
manager | ||
.loadAll() | ||
.filter(de -> de.getState() == 1) | ||
.flatMap(this::init) | ||
.blockLast(loadTimeOut); | ||
} catch (Throwable e) { | ||
log.error("load protocol error", e); | ||
} | ||
public class LazyInitManagementProtocolSupports extends DefaultProtocolSupportManager implements CommandLineRunner { | ||
|
||
public LazyInitManagementProtocolSupports(EventBus eventBus, | ||
ClusterManager clusterManager, | ||
ProtocolSupportLoader loader) { | ||
super(eventBus, clusterManager.getCache("__protocol_supports"), loader); | ||
} | ||
|
||
public Mono<Void> init(ProtocolSupportDefinition definition) { | ||
try { | ||
if (definition.getState() != 1) { | ||
String protocol = configProtocolIdMapping.get(definition.getId()); | ||
if (protocol != null) { | ||
log.debug("uninstall protocol:{}", definition); | ||
unRegister(protocol); | ||
return Mono.empty(); | ||
} | ||
} | ||
String operation = definition.getState() != 1 ? "uninstall" : "install"; | ||
Consumer<ProtocolSupport> consumer = definition.getState() != 1 ? this::unRegister : this::register; | ||
|
||
log.debug("{} protocol:{}", operation, definition); | ||
|
||
return loader | ||
.load(definition) | ||
.doOnNext(e -> { | ||
e.init(definition.getConfiguration()); | ||
log.debug("{} protocol[{}] success: {}", operation, definition.getId(), e); | ||
configProtocolIdMapping.put(definition.getId(), e.getId()); | ||
consumer.accept(e); | ||
}) | ||
.onErrorResume((e) -> { | ||
log.error("{} protocol[{}] error: {}", operation, definition.getId(), e.getLocalizedMessage()); | ||
return Mono.empty(); | ||
}) | ||
.then(); | ||
} catch (Throwable err) { | ||
log.error("init protocol error", err); | ||
} | ||
return Mono.empty(); | ||
public void init() { | ||
super.init(); | ||
} | ||
|
||
@Override | ||
public void run(String... args) { | ||
init(); | ||
} | ||
|
||
|
||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters