539 changed files with 1232 additions and 18413 deletions
@ -1,159 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import akka.actor.ActorContext; |
|
||||
import akka.actor.ActorRef; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.service.ComponentActor; |
|
||||
import org.thingsboard.server.actors.service.ContextBasedCreator; |
|
||||
import org.thingsboard.server.actors.stats.StatsPersistTick; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.msg.TbActorMsg; |
|
||||
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; |
|
||||
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; |
|
||||
import org.thingsboard.server.common.msg.timeout.TimeoutMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.ToPluginRpcResponseDeviceMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.rest.PluginRestMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.rpc.PluginRpcMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg; |
|
||||
import org.thingsboard.server.extensions.api.rules.RuleException; |
|
||||
|
|
||||
public class PluginActor extends ComponentActor<PluginId, PluginActorMessageProcessor> { |
|
||||
|
|
||||
private PluginActor(ActorSystemContext systemContext, TenantId tenantId, PluginId pluginId) { |
|
||||
super(systemContext, tenantId, pluginId); |
|
||||
setProcessor(new PluginActorMessageProcessor(tenantId, pluginId, systemContext, |
|
||||
logger, context().parent(), context().self())); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected boolean process(TbActorMsg msg) { |
|
||||
//TODO Move everything here, to work with TbActorMsg
|
|
||||
return false; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onReceive(Object msg) throws Exception { |
|
||||
if (msg instanceof PluginWebsocketMsg) { |
|
||||
onWebsocketMsg((PluginWebsocketMsg<?>) msg); |
|
||||
} else if (msg instanceof PluginRestMsg) { |
|
||||
onRestMsg((PluginRestMsg) msg); |
|
||||
} else if (msg instanceof PluginCallbackMessage) { |
|
||||
onPluginCallback((PluginCallbackMessage) msg); |
|
||||
} else if (msg instanceof RuleToPluginMsgWrapper) { |
|
||||
onRuleToPluginMsg((RuleToPluginMsgWrapper) msg); |
|
||||
} else if (msg instanceof PluginRpcMsg) { |
|
||||
onRpcMsg((PluginRpcMsg) msg); |
|
||||
} else if (msg instanceof ClusterEventMsg) { |
|
||||
onClusterEventMsg((ClusterEventMsg) msg); |
|
||||
} else if (msg instanceof ComponentLifecycleMsg) { |
|
||||
onComponentLifecycleMsg((ComponentLifecycleMsg) msg); |
|
||||
} else if (msg instanceof ToPluginRpcResponseDeviceMsg) { |
|
||||
onRpcResponse((ToPluginRpcResponseDeviceMsg) msg); |
|
||||
} else if (msg instanceof PluginTerminationMsg) { |
|
||||
logger.info("[{}][{}] Going to terminate plugin actor.", tenantId, id); |
|
||||
context().parent().tell(msg, ActorRef.noSender()); |
|
||||
context().stop(self()); |
|
||||
} else if (msg instanceof TimeoutMsg) { |
|
||||
onTimeoutMsg(context(), (TimeoutMsg) msg); |
|
||||
} else if (msg instanceof StatsPersistTick) { |
|
||||
onStatsPersistTick(id); |
|
||||
} else { |
|
||||
logger.debug("[{}][{}] Unknown msg type.", tenantId, id, msg.getClass().getName()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void onPluginCallback(PluginCallbackMessage msg) { |
|
||||
try { |
|
||||
processor.onPluginCallbackMsg(msg); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onPluginCallbackMsg", e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void onTimeoutMsg(ActorContext context, TimeoutMsg msg) { |
|
||||
processor.onTimeoutMsg(context, msg); |
|
||||
} |
|
||||
|
|
||||
private void onRpcResponse(ToPluginRpcResponseDeviceMsg msg) { |
|
||||
processor.onDeviceRpcMsg(msg.getResponse()); |
|
||||
} |
|
||||
|
|
||||
private void onRuleToPluginMsg(RuleToPluginMsgWrapper msg) throws RuleException { |
|
||||
logger.debug("[{}] Going to process rule msg: {}", id, msg.getMsg()); |
|
||||
try { |
|
||||
processor.onRuleToPluginMsg(msg); |
|
||||
increaseMessagesProcessedCount(); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onRuleMsg", e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void onWebsocketMsg(PluginWebsocketMsg<?> msg) { |
|
||||
logger.debug("[{}] Going to process web socket msg: {}", id, msg); |
|
||||
try { |
|
||||
processor.onWebsocketMsg(msg); |
|
||||
increaseMessagesProcessedCount(); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onWebsocketMsg", e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void onRestMsg(PluginRestMsg msg) { |
|
||||
logger.debug("[{}] Going to process rest msg: {}", id, msg); |
|
||||
try { |
|
||||
processor.onRestMsg(msg); |
|
||||
increaseMessagesProcessedCount(); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onRestMsg", e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void onRpcMsg(PluginRpcMsg msg) { |
|
||||
try { |
|
||||
logger.debug("[{}] Going to process rpc msg: {}", id, msg); |
|
||||
processor.onRpcMsg(msg); |
|
||||
} catch (Exception e) { |
|
||||
logAndPersist("onRpcMsg", e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public static class ActorCreator extends ContextBasedCreator<PluginActor> { |
|
||||
private static final long serialVersionUID = 1L; |
|
||||
|
|
||||
private final TenantId tenantId; |
|
||||
private final PluginId pluginId; |
|
||||
|
|
||||
public ActorCreator(ActorSystemContext context, TenantId tenantId, PluginId pluginId) { |
|
||||
super(context); |
|
||||
this.tenantId = tenantId; |
|
||||
this.pluginId = pluginId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public PluginActor create() throws Exception { |
|
||||
return new PluginActor(context, tenantId, pluginId); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected long getErrorPersistFrequency() { |
|
||||
return 0; |
|
||||
// return systemContext.getPluginErrorPersistFrequency();
|
|
||||
} |
|
||||
} |
|
||||
@ -1,252 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import akka.actor.ActorContext; |
|
||||
import akka.actor.ActorRef; |
|
||||
import akka.event.LoggingAdapter; |
|
||||
import com.fasterxml.jackson.core.JsonProcessingException; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.shared.ComponentMsgProcessor; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleState; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentType; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg; |
|
||||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|
||||
import org.thingsboard.server.common.msg.core.BasicStatusCodeResponse; |
|
||||
import org.thingsboard.server.common.msg.session.FromDeviceRequestMsg; |
|
||||
import org.thingsboard.server.common.msg.session.SessionMsgType; |
|
||||
import org.thingsboard.server.common.msg.session.SessionMsgType; |
|
||||
import org.thingsboard.server.extensions.api.plugins.Plugin; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginInitializationException; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.FromDeviceRpcResponse; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.ResponsePluginToRuleMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.RuleToPluginMsg; |
|
||||
import org.thingsboard.server.common.msg.timeout.TimeoutMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.rest.PluginRestMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.rpc.PluginRpcMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg; |
|
||||
import org.thingsboard.server.extensions.api.rules.RuleException; |
|
||||
|
|
||||
/** |
|
||||
* @author Andrew Shvayka |
|
||||
*/ |
|
||||
public class PluginActorMessageProcessor extends ComponentMsgProcessor<PluginId> { |
|
||||
|
|
||||
private final SharedPluginProcessingContext pluginCtx; |
|
||||
private final PluginProcessingContext trustedCtx; |
|
||||
private PluginMetaData pluginMd; |
|
||||
private Plugin pluginImpl; |
|
||||
private ComponentLifecycleState state; |
|
||||
|
|
||||
|
|
||||
protected PluginActorMessageProcessor(TenantId tenantId, PluginId pluginId, ActorSystemContext systemContext |
|
||||
, LoggingAdapter logger, ActorRef parent, ActorRef self) { |
|
||||
super(systemContext, logger, tenantId, pluginId); |
|
||||
this.pluginCtx = new SharedPluginProcessingContext(systemContext, tenantId, pluginId, parent, self); |
|
||||
this.trustedCtx = new PluginProcessingContext(pluginCtx, null); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void start(ActorContext context) throws Exception { |
|
||||
logger.info("[{}] Going to start plugin actor.", entityId); |
|
||||
pluginMd = systemContext.getPluginService().findPluginById(entityId); |
|
||||
if (pluginMd == null) { |
|
||||
throw new PluginInitializationException("Plugin not found!"); |
|
||||
} |
|
||||
if (pluginMd.getConfiguration() == null) { |
|
||||
throw new PluginInitializationException("Plugin metadata is empty!"); |
|
||||
} |
|
||||
state = pluginMd.getState(); |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
logger.info("[{}] Plugin is active. Going to initialize plugin.", entityId); |
|
||||
initComponent(); |
|
||||
} else { |
|
||||
logger.info("[{}] Plugin is suspended. Skipping plugin initialization.", entityId); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void stop(ActorContext context) throws Exception { |
|
||||
onStop(); |
|
||||
} |
|
||||
|
|
||||
private void initComponent() { |
|
||||
try { |
|
||||
pluginImpl = initComponent(pluginMd.getClazz(), ComponentType.PLUGIN, mapper.writeValueAsString(pluginMd.getConfiguration())); |
|
||||
} catch (InstantiationException e) { |
|
||||
throw new PluginInitializationException("No default constructor for plugin implementation!", e); |
|
||||
} catch (IllegalAccessException e) { |
|
||||
throw new PluginInitializationException("Illegal Access Exception during plugin initialization!", e); |
|
||||
} catch (ClassNotFoundException e) { |
|
||||
throw new PluginInitializationException("Plugin Class not found!", e); |
|
||||
} catch (JsonProcessingException e) { |
|
||||
throw new PluginInitializationException("Plugin Configuration is invalid!", e); |
|
||||
} catch (Exception e) { |
|
||||
throw new PluginInitializationException(e.getMessage(), e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void onRuleToPluginMsg(RuleToPluginMsgWrapper msg) throws RuleException { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
try { |
|
||||
pluginImpl.process(trustedCtx, msg.getRuleTenantId(), msg.getRuleId(), msg.getMsg()); |
|
||||
} catch (Exception ex) { |
|
||||
logger.debug("[{}] Failed to process RuleToPlugin msg: [{}] [{}]", tenantId, msg.getMsg(), ex); |
|
||||
RuleToPluginMsg ruleMsg = msg.getMsg(); |
|
||||
SessionMsgType responceMsgType = SessionMsgType.RULE_ENGINE_ERROR; |
|
||||
Integer requestId = 0; |
|
||||
if (ruleMsg.getPayload() instanceof FromDeviceRequestMsg) { |
|
||||
requestId = ((FromDeviceRequestMsg) ruleMsg.getPayload()).getRequestId(); |
|
||||
} |
|
||||
trustedCtx.reply( |
|
||||
new ResponsePluginToRuleMsg(ruleMsg.getUid(), tenantId, msg.getRuleId(), |
|
||||
BasicStatusCodeResponse.onError(responceMsgType, requestId, ex))); |
|
||||
} |
|
||||
} else { |
|
||||
//TODO: reply with plugin suspended message
|
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void onWebsocketMsg(PluginWebsocketMsg<?> msg) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
pluginImpl.process(new PluginProcessingContext(pluginCtx, msg.getSecurityCtx()), msg); |
|
||||
} else { |
|
||||
//TODO: reply with plugin suspended message
|
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void onRestMsg(PluginRestMsg msg) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
pluginImpl.process(new PluginProcessingContext(pluginCtx, msg.getSecurityCtx()), msg); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void onRpcMsg(PluginRpcMsg msg) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
pluginImpl.process(trustedCtx, msg.getRpcMsg()); |
|
||||
} else { |
|
||||
//TODO: reply with plugin suspended message
|
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void onPluginCallbackMsg(PluginCallbackMessage msg) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
if (msg.isSuccess()) { |
|
||||
msg.getCallback().onSuccess(trustedCtx, msg.getV()); |
|
||||
} else { |
|
||||
msg.getCallback().onFailure(trustedCtx, msg.getE()); |
|
||||
} |
|
||||
} else { |
|
||||
//TODO: reply with plugin suspended message
|
|
||||
} |
|
||||
} |
|
||||
|
|
||||
|
|
||||
public void onTimeoutMsg(ActorContext context, TimeoutMsg<?> msg) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
pluginImpl.process(trustedCtx, msg); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
|
|
||||
public void onDeviceRpcMsg(FromDeviceRpcResponse response) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
pluginImpl.process(trustedCtx, response); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onClusterEventMsg(ClusterEventMsg msg) { |
|
||||
if (state == ComponentLifecycleState.ACTIVE) { |
|
||||
ServerAddress address = msg.getServerAddress(); |
|
||||
if (msg.isAdded()) { |
|
||||
logger.debug("[{}] Going to process server add msg: {}", entityId, address); |
|
||||
pluginImpl.onServerAdded(trustedCtx, address); |
|
||||
} else { |
|
||||
logger.debug("[{}] Going to process server remove msg: {}", entityId, address); |
|
||||
pluginImpl.onServerRemoved(trustedCtx, address); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onCreated(ActorContext context) { |
|
||||
logger.info("[{}] Going to process onCreated plugin.", entityId); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onUpdate(ActorContext context) throws Exception { |
|
||||
PluginMetaData oldPluginMd = pluginMd; |
|
||||
pluginMd = systemContext.getPluginService().findPluginById(entityId); |
|
||||
boolean requiresRestart = false; |
|
||||
logger.info("[{}] Plugin configuration was updated from {} to {}.", entityId, oldPluginMd, pluginMd); |
|
||||
if (!oldPluginMd.getClazz().equals(pluginMd.getClazz())) { |
|
||||
logger.info("[{}] Plugin requires restart due to clazz change from {} to {}.", |
|
||||
entityId, oldPluginMd.getClazz(), pluginMd.getClazz()); |
|
||||
requiresRestart = true; |
|
||||
} else if (!oldPluginMd.getConfiguration().equals(pluginMd.getConfiguration())) { |
|
||||
logger.info("[{}] Plugin requires restart due to configuration change from {} to {}.", |
|
||||
entityId, oldPluginMd.getConfiguration(), pluginMd.getConfiguration()); |
|
||||
requiresRestart = true; |
|
||||
} |
|
||||
if (requiresRestart) { |
|
||||
this.state = ComponentLifecycleState.SUSPENDED; |
|
||||
if (pluginImpl != null) { |
|
||||
pluginImpl.stop(trustedCtx); |
|
||||
} |
|
||||
start(context); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onStop(ActorContext context) { |
|
||||
onStop(); |
|
||||
// scheduleMsgWithDelay(context, new PluginTerminationMsg(entityId), systemContext.getPluginActorTerminationDelay());
|
|
||||
} |
|
||||
|
|
||||
private void onStop() { |
|
||||
logger.info("[{}] Going to process onStop plugin.", entityId); |
|
||||
this.state = ComponentLifecycleState.SUSPENDED; |
|
||||
if (pluginImpl != null) { |
|
||||
pluginImpl.stop(trustedCtx); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onActivate(ActorContext context) throws Exception { |
|
||||
logger.info("[{}] Going to process onActivate plugin.", entityId); |
|
||||
this.state = ComponentLifecycleState.ACTIVE; |
|
||||
if (pluginImpl != null) { |
|
||||
pluginImpl.resume(trustedCtx); |
|
||||
logger.info("[{}] Plugin resumed.", entityId); |
|
||||
} else { |
|
||||
start(context); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onSuspend(ActorContext context) { |
|
||||
logger.info("[{}] Going to process onSuspend plugin.", entityId); |
|
||||
this.state = ComponentLifecycleState.SUSPENDED; |
|
||||
if (pluginImpl != null) { |
|
||||
pluginImpl.suspend(trustedCtx); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,53 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import lombok.Data; |
|
||||
import lombok.Getter; |
|
||||
import lombok.ToString; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginCallback; |
|
||||
|
|
||||
import java.util.Optional; |
|
||||
|
|
||||
/** |
|
||||
* @author Andrew Shvayka |
|
||||
*/ |
|
||||
@ToString |
|
||||
public final class PluginCallbackMessage<V> { |
|
||||
@Getter |
|
||||
private final PluginCallback<V> callback; |
|
||||
@Getter |
|
||||
private final boolean success; |
|
||||
@Getter |
|
||||
private final V v; |
|
||||
@Getter |
|
||||
private final Exception e; |
|
||||
|
|
||||
public static <V> PluginCallbackMessage<V> onSuccess(PluginCallback<V> callback, V data) { |
|
||||
return new PluginCallbackMessage<V>(true, callback, data, null); |
|
||||
} |
|
||||
|
|
||||
public static <V> PluginCallbackMessage<V> onError(PluginCallback<V> callback, Exception e) { |
|
||||
return new PluginCallbackMessage<V>(false, callback, null, e); |
|
||||
} |
|
||||
|
|
||||
private PluginCallbackMessage(boolean success, PluginCallback<V> callback, V v, Exception e) { |
|
||||
this.success = success; |
|
||||
this.callback = callback; |
|
||||
this.v = v; |
|
||||
this.e = e; |
|
||||
} |
|
||||
} |
|
||||
@ -1,576 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import akka.actor.ActorRef; |
|
||||
import com.google.common.base.Function; |
|
||||
import com.google.common.util.concurrent.FutureCallback; |
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.thingsboard.server.common.data.Customer; |
|
||||
import org.thingsboard.server.common.data.Device; |
|
||||
import org.thingsboard.server.common.data.EntityType; |
|
||||
import org.thingsboard.server.common.data.Tenant; |
|
||||
import org.thingsboard.server.common.data.asset.Asset; |
|
||||
import org.thingsboard.server.common.data.audit.ActionType; |
|
||||
import org.thingsboard.server.common.data.id.*; |
|
||||
import org.thingsboard.server.common.data.kv.AttributeKey; |
|
||||
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
||||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
||||
import org.thingsboard.server.common.data.kv.TsKvQuery; |
|
||||
import org.thingsboard.server.common.data.page.TextPageLink; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.common.data.relation.EntityRelation; |
|
||||
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|
||||
import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; |
|
||||
import org.thingsboard.server.common.data.rule.RuleChain; |
|
||||
import org.thingsboard.server.common.data.rule.RuleMetaData; |
|
||||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|
||||
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; |
|
||||
import org.thingsboard.server.common.msg.timeout.TimeoutMsg; |
|
||||
import org.thingsboard.server.extensions.api.device.DeviceAttributesEventNotificationMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginApiCallSecurityContext; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginCallback; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginContext; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.*; |
|
||||
import org.thingsboard.server.extensions.api.plugins.rpc.PluginRpcMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.rpc.RpcMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.PluginWebsocketSessionRef; |
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg; |
|
||||
|
|
||||
import javax.annotation.Nullable; |
|
||||
import java.io.IOException; |
|
||||
import java.util.*; |
|
||||
import java.util.concurrent.Executor; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.stream.Collectors; |
|
||||
|
|
||||
@Slf4j |
|
||||
public final class PluginProcessingContext implements PluginContext { |
|
||||
|
|
||||
private static final Executor executor = Executors.newSingleThreadExecutor(); |
|
||||
public static final String CUSTOMER_USER_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION = "Customer user is not allowed to perform this operation!"; |
|
||||
public static final String SYSTEM_ADMINISTRATOR_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION = "System administrator is not allowed to perform this operation!"; |
|
||||
public static final String DEVICE_WITH_REQUESTED_ID_NOT_FOUND = "Device with requested id wasn't found!"; |
|
||||
|
|
||||
private final SharedPluginProcessingContext pluginCtx; |
|
||||
private final Optional<PluginApiCallSecurityContext> securityCtx; |
|
||||
|
|
||||
public PluginProcessingContext(SharedPluginProcessingContext pluginCtx, PluginApiCallSecurityContext securityCtx) { |
|
||||
super(); |
|
||||
this.pluginCtx = pluginCtx; |
|
||||
this.securityCtx = Optional.ofNullable(securityCtx); |
|
||||
} |
|
||||
|
|
||||
public void persistError(String method, Exception e) { |
|
||||
pluginCtx.persistError(method, e); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void sendPluginRpcMsg(RpcMsg msg) { |
|
||||
//ToDO is this a cluster messsage?
|
|
||||
// this.pluginCtx.rpcService.tell(new PluginRpcMsg(pluginCtx.tenantId, pluginCtx.pluginId, msg));
|
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void send(PluginWebsocketMsg<?> wsMsg) throws IOException { |
|
||||
pluginCtx.msgEndpoint.send(wsMsg); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void close(PluginWebsocketSessionRef sessionRef) throws IOException { |
|
||||
pluginCtx.msgEndpoint.close(sessionRef); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void saveAttributes(final TenantId tenantId, final EntityId entityId, final String scope, final List<AttributeKvEntry> attributes, final PluginCallback<Void> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<Void>> futures = pluginCtx.attributesService.save(entityId, scope, attributes); |
|
||||
Futures.addCallback(futures, getListCallback(callback, v -> { |
|
||||
if (entityId.getEntityType() == EntityType.DEVICE) { |
|
||||
onDeviceAttributesChanged(tenantId, new DeviceId(entityId.getId()), scope, attributes); |
|
||||
} |
|
||||
return null; |
|
||||
}), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void removeAttributes(final TenantId tenantId, final EntityId entityId, final String scope, final List<String> keys, final PluginCallback<Void> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<Void>> futures = pluginCtx.attributesService.removeAll(entityId, scope, keys); |
|
||||
Futures.addCallback(futures, getCallback(callback, v -> null), executor); |
|
||||
if (entityId.getEntityType() == EntityType.DEVICE) { |
|
||||
onDeviceAttributesDeleted(tenantId, new DeviceId(entityId.getId()), keys.stream().map(key -> new AttributeKey(scope, key)).collect(Collectors.toSet())); |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadAttribute(EntityId entityId, String attributeType, String attributeKey, final PluginCallback<Optional<AttributeKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<Optional<AttributeKvEntry>> future = pluginCtx.attributesService.find(entityId, attributeType, attributeKey); |
|
||||
Futures.addCallback(future, getCallback(callback, v -> v), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadAttributes(EntityId entityId, String attributeType, Collection<String> attributeKeys, final PluginCallback<List<AttributeKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<AttributeKvEntry>> future = pluginCtx.attributesService.find(entityId, attributeType, attributeKeys); |
|
||||
Futures.addCallback(future, getCallback(callback, v -> v), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadAttributes(EntityId entityId, String attributeType, PluginCallback<List<AttributeKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<AttributeKvEntry>> future = pluginCtx.attributesService.findAll(entityId, attributeType); |
|
||||
Futures.addCallback(future, getCallback(callback, v -> v), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadAttributes(final EntityId entityId, final Collection<String> attributeTypes, final PluginCallback<List<AttributeKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
List<ListenableFuture<List<AttributeKvEntry>>> futures = new ArrayList<>(); |
|
||||
attributeTypes.forEach(attributeType -> futures.add(pluginCtx.attributesService.findAll(entityId, attributeType))); |
|
||||
convertFuturesAndAddCallback(callback, futures); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadAttributes(final EntityId entityId, final Collection<String> attributeTypes, final Collection<String> attributeKeys, final PluginCallback<List<AttributeKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
List<ListenableFuture<List<AttributeKvEntry>>> futures = new ArrayList<>(); |
|
||||
attributeTypes.forEach(attributeType -> futures.add(pluginCtx.attributesService.find(entityId, attributeType, attributeKeys))); |
|
||||
convertFuturesAndAddCallback(callback, futures); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void saveTsData(final EntityId entityId, final TsKvEntry entry, final PluginCallback<Void> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<Void>> rsListFuture = pluginCtx.tsService.save(entityId, entry); |
|
||||
Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void saveTsData(final EntityId entityId, final List<TsKvEntry> entries, final PluginCallback<Void> callback) { |
|
||||
saveTsData(entityId, entries, 0L, callback); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void saveTsData(final EntityId entityId, final List<TsKvEntry> entries, long ttl, final PluginCallback<Void> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<Void>> rsListFuture = pluginCtx.tsService.save(entityId, entries, ttl); |
|
||||
Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
|
|
||||
@Override |
|
||||
public void loadTimeseries(final EntityId entityId, final List<TsKvQuery> queries, final PluginCallback<List<TsKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<TsKvEntry>> future = pluginCtx.tsService.findAll(entityId, queries); |
|
||||
Futures.addCallback(future, getCallback(callback, v -> v), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadLatestTimeseries(final EntityId entityId, final PluginCallback<List<TsKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<TsKvEntry>> future = pluginCtx.tsService.findAllLatest(entityId); |
|
||||
Futures.addCallback(future, getCallback(callback, v -> v), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void logAttributesUpdated(PluginApiCallSecurityContext ctx, EntityId entityId, String attributeType, |
|
||||
List<AttributeKvEntry> attributes, Exception e) { |
|
||||
pluginCtx.auditLogService.logEntityAction( |
|
||||
ctx.getTenantId(), |
|
||||
ctx.getCustomerId(), |
|
||||
ctx.getUserId(), |
|
||||
ctx.getUserName(), |
|
||||
(UUIDBased & EntityId)entityId, |
|
||||
null, |
|
||||
ActionType.ATTRIBUTES_UPDATED, |
|
||||
e, |
|
||||
attributeType, |
|
||||
attributes); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void logAttributesDeleted(PluginApiCallSecurityContext ctx, EntityId entityId, String attributeType, List<String> keys, Exception e) { |
|
||||
pluginCtx.auditLogService.logEntityAction( |
|
||||
ctx.getTenantId(), |
|
||||
ctx.getCustomerId(), |
|
||||
ctx.getUserId(), |
|
||||
ctx.getUserName(), |
|
||||
(UUIDBased & EntityId)entityId, |
|
||||
null, |
|
||||
ActionType.ATTRIBUTES_DELETED, |
|
||||
e, |
|
||||
attributeType, |
|
||||
keys); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void logAttributesRead(PluginApiCallSecurityContext ctx, EntityId entityId, String attributeType, List<String> keys, Exception e) { |
|
||||
pluginCtx.auditLogService.logEntityAction( |
|
||||
ctx.getTenantId(), |
|
||||
ctx.getCustomerId(), |
|
||||
ctx.getUserId(), |
|
||||
ctx.getUserName(), |
|
||||
(UUIDBased & EntityId)entityId, |
|
||||
null, |
|
||||
ActionType.ATTRIBUTES_READ, |
|
||||
e, |
|
||||
attributeType, |
|
||||
keys); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void loadLatestTimeseries(final EntityId entityId, final Collection<String> keys, final PluginCallback<List<TsKvEntry>> callback) { |
|
||||
validate(entityId, new ValidationCallback(callback, ctx -> { |
|
||||
ListenableFuture<List<TsKvEntry>> rsListFuture = pluginCtx.tsService.findLatest(entityId, keys); |
|
||||
Futures.addCallback(rsListFuture, getCallback(callback, v -> v), executor); |
|
||||
})); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void reply(PluginToRuleMsg<?> msg) { |
|
||||
pluginCtx.parentActor.tell(msg, ActorRef.noSender()); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public PluginId getPluginId() { |
|
||||
return pluginCtx.pluginId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Optional<PluginApiCallSecurityContext> getSecurityCtx() { |
|
||||
return securityCtx; |
|
||||
} |
|
||||
|
|
||||
private void onDeviceAttributesDeleted(TenantId tenantId, DeviceId deviceId, Set<AttributeKey> keys) { |
|
||||
pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onDelete(tenantId, deviceId, keys)); |
|
||||
} |
|
||||
|
|
||||
private void onDeviceAttributesChanged(TenantId tenantId, DeviceId deviceId, String scope, List<AttributeKvEntry> values) { |
|
||||
pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onUpdate(tenantId, deviceId, scope, values)); |
|
||||
} |
|
||||
|
|
||||
private <T, R> FutureCallback<List<T>> getListCallback(final PluginCallback<R> callback, Function<List<T>, R> transformer) { |
|
||||
return new FutureCallback<List<T>>() { |
|
||||
@Override |
|
||||
public void onSuccess(@Nullable List<T> result) { |
|
||||
pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, transformer.apply(result)), ActorRef.noSender()); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(Throwable t) { |
|
||||
if (t instanceof Exception) { |
|
||||
pluginCtx.self().tell(PluginCallbackMessage.onError(callback, (Exception) t), ActorRef.noSender()); |
|
||||
} else { |
|
||||
log.error("Critical error: {}", t.getMessage(), t); |
|
||||
} |
|
||||
} |
|
||||
}; |
|
||||
} |
|
||||
|
|
||||
private <T, R> FutureCallback<R> getCallback(final PluginCallback<T> callback, Function<R, T> transformer) { |
|
||||
return new FutureCallback<R>() { |
|
||||
@Override |
|
||||
public void onSuccess(@Nullable R result) { |
|
||||
try { |
|
||||
pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, transformer.apply(result)), ActorRef.noSender()); |
|
||||
} catch (Exception e) { |
|
||||
pluginCtx.self().tell(PluginCallbackMessage.onError(callback, e), ActorRef.noSender()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(Throwable t) { |
|
||||
if (t instanceof Exception) { |
|
||||
pluginCtx.self().tell(PluginCallbackMessage.onError(callback, (Exception) t), ActorRef.noSender()); |
|
||||
} else { |
|
||||
log.error("Critical error: {}", t.getMessage(), t); |
|
||||
} |
|
||||
} |
|
||||
}; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void checkAccess(DeviceId deviceId, PluginCallback<Void> callback) { |
|
||||
validate(deviceId, new ValidationCallback(callback, ctx -> callback.onSuccess(ctx, null))); |
|
||||
} |
|
||||
|
|
||||
private void validate(EntityId entityId, ValidationCallback callback) { |
|
||||
if (securityCtx.isPresent()) { |
|
||||
final PluginApiCallSecurityContext ctx = securityCtx.get(); |
|
||||
switch (entityId.getEntityType()) { |
|
||||
case DEVICE: |
|
||||
validateDevice(ctx, entityId, callback); |
|
||||
return; |
|
||||
case ASSET: |
|
||||
validateAsset(ctx, entityId, callback); |
|
||||
return; |
|
||||
case RULE: |
|
||||
validateRule(ctx, entityId, callback); |
|
||||
return; |
|
||||
case RULE_CHAIN: |
|
||||
validateRuleChain(ctx, entityId, callback); |
|
||||
return; |
|
||||
case PLUGIN: |
|
||||
validatePlugin(ctx, entityId, callback); |
|
||||
return; |
|
||||
case CUSTOMER: |
|
||||
validateCustomer(ctx, entityId, callback); |
|
||||
return; |
|
||||
case TENANT: |
|
||||
validateTenant(ctx, entityId, callback); |
|
||||
return; |
|
||||
default: |
|
||||
//TODO: add support of other entities
|
|
||||
throw new IllegalStateException("Not Implemented!"); |
|
||||
} |
|
||||
} else { |
|
||||
callback.onSuccess(this, ValidationResult.ok(null)); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void validateDevice(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isSystemAdmin()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(SYSTEM_ADMINISTRATOR_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else { |
|
||||
ListenableFuture<Device> deviceFuture = pluginCtx.deviceService.findDeviceByIdAsync(new DeviceId(entityId.getId())); |
|
||||
Futures.addCallback(deviceFuture, getCallback(callback, device -> { |
|
||||
if (device == null) { |
|
||||
return ValidationResult.entityNotFound(DEVICE_WITH_REQUESTED_ID_NOT_FOUND); |
|
||||
} else { |
|
||||
if (!device.getTenantId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Device doesn't belong to the current Tenant!"); |
|
||||
} else if (ctx.isCustomerUser() && !device.getCustomerId().equals(ctx.getCustomerId())) { |
|
||||
return ValidationResult.accessDenied("Device doesn't belong to the current Customer!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void validateAsset(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isSystemAdmin()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(SYSTEM_ADMINISTRATOR_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else { |
|
||||
ListenableFuture<Asset> assetFuture = pluginCtx.assetService.findAssetByIdAsync(new AssetId(entityId.getId())); |
|
||||
Futures.addCallback(assetFuture, getCallback(callback, asset -> { |
|
||||
if (asset == null) { |
|
||||
return ValidationResult.entityNotFound("Asset with requested id wasn't found!"); |
|
||||
} else { |
|
||||
if (!asset.getTenantId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Asset doesn't belong to the current Tenant!"); |
|
||||
} else if (ctx.isCustomerUser() && !asset.getCustomerId().equals(ctx.getCustomerId())) { |
|
||||
return ValidationResult.accessDenied("Asset doesn't belong to the current Customer!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void validateRule(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isCustomerUser()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(CUSTOMER_USER_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else { |
|
||||
ListenableFuture<RuleMetaData> ruleFuture = pluginCtx.ruleService.findRuleByIdAsync(new RuleId(entityId.getId())); |
|
||||
Futures.addCallback(ruleFuture, getCallback(callback, rule -> { |
|
||||
if (rule == null) { |
|
||||
return ValidationResult.entityNotFound("Rule with requested id wasn't found!"); |
|
||||
} else { |
|
||||
if (ctx.isTenantAdmin() && !rule.getTenantId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Rule doesn't belong to the current Tenant!"); |
|
||||
} else if (ctx.isSystemAdmin() && !rule.getTenantId().isNullUid()) { |
|
||||
return ValidationResult.accessDenied("Rule is not in system scope!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void validateRuleChain(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isCustomerUser()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(CUSTOMER_USER_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else { |
|
||||
ListenableFuture<RuleChain> ruleChainFuture = pluginCtx.ruleChainService.findRuleChainByIdAsync(new RuleChainId(entityId.getId())); |
|
||||
Futures.addCallback(ruleChainFuture, getCallback(callback, ruleChain -> { |
|
||||
if (ruleChain == null) { |
|
||||
return ValidationResult.entityNotFound("Rule chain with requested id wasn't found!"); |
|
||||
} else { |
|
||||
if (ctx.isTenantAdmin() && !ruleChain.getTenantId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Rule chain doesn't belong to the current Tenant!"); |
|
||||
} else if (ctx.isSystemAdmin() && !ruleChain.getTenantId().isNullUid()) { |
|
||||
return ValidationResult.accessDenied("Rule chain is not in system scope!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
|
|
||||
private void validatePlugin(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isCustomerUser()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(CUSTOMER_USER_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else { |
|
||||
ListenableFuture<PluginMetaData> pluginFuture = pluginCtx.pluginService.findPluginByIdAsync(new PluginId(entityId.getId())); |
|
||||
Futures.addCallback(pluginFuture, getCallback(callback, plugin -> { |
|
||||
if (plugin == null) { |
|
||||
return ValidationResult.entityNotFound("Plugin with requested id wasn't found!"); |
|
||||
} else { |
|
||||
if (ctx.isTenantAdmin() && !plugin.getTenantId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Plugin doesn't belong to the current Tenant!"); |
|
||||
} else if (ctx.isSystemAdmin() && !plugin.getTenantId().isNullUid()) { |
|
||||
return ValidationResult.accessDenied("Plugin is not in system scope!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void validateCustomer(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isSystemAdmin()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(SYSTEM_ADMINISTRATOR_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else { |
|
||||
ListenableFuture<Customer> customerFuture = pluginCtx.customerService.findCustomerByIdAsync(new CustomerId(entityId.getId())); |
|
||||
Futures.addCallback(customerFuture, getCallback(callback, customer -> { |
|
||||
if (customer == null) { |
|
||||
return ValidationResult.entityNotFound("Customer with requested id wasn't found!"); |
|
||||
} else { |
|
||||
if (!customer.getTenantId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Customer doesn't belong to the current Tenant!"); |
|
||||
} else if (ctx.isCustomerUser() && !customer.getId().equals(ctx.getCustomerId())) { |
|
||||
return ValidationResult.accessDenied("Customer doesn't relate to the currently authorized customer user!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void validateTenant(final PluginApiCallSecurityContext ctx, EntityId entityId, ValidationCallback callback) { |
|
||||
if (ctx.isCustomerUser()) { |
|
||||
callback.onSuccess(this, ValidationResult.accessDenied(CUSTOMER_USER_IS_NOT_ALLOWED_TO_PERFORM_THIS_OPERATION)); |
|
||||
} else if (ctx.isSystemAdmin()) { |
|
||||
callback.onSuccess(this, ValidationResult.ok(null)); |
|
||||
} else { |
|
||||
ListenableFuture<Tenant> tenantFuture = pluginCtx.tenantService.findTenantByIdAsync(new TenantId(entityId.getId())); |
|
||||
Futures.addCallback(tenantFuture, getCallback(callback, tenant -> { |
|
||||
if (tenant == null) { |
|
||||
return ValidationResult.entityNotFound("Tenant with requested id wasn't found!"); |
|
||||
} else if (!tenant.getId().equals(ctx.getTenantId())) { |
|
||||
return ValidationResult.accessDenied("Tenant doesn't relate to the currently authorized user!"); |
|
||||
} else { |
|
||||
return ValidationResult.ok(null); |
|
||||
} |
|
||||
})); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<List<EntityRelation>> findByFromAndType(EntityId from, String relationType) { |
|
||||
return this.pluginCtx.relationService.findByFromAndTypeAsync(from, relationType, RelationTypeGroup.COMMON); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<List<EntityRelation>> findByToAndType(EntityId from, String relationType) { |
|
||||
return this.pluginCtx.relationService.findByToAndTypeAsync(from, relationType, RelationTypeGroup.COMMON); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Optional<ServerAddress> resolve(EntityId entityId) { |
|
||||
return pluginCtx.routingService.resolveById(entityId); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void getDevice(DeviceId deviceId, PluginCallback<Device> callback) { |
|
||||
ListenableFuture<Device> deviceFuture = pluginCtx.deviceService.findDeviceByIdAsync(deviceId); |
|
||||
Futures.addCallback(deviceFuture, getCallback(callback, v -> v)); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void getCustomerDevices(TenantId tenantId, CustomerId customerId, int limit, PluginCallback<List<Device>> callback) { |
|
||||
//TODO: add caching here with async api.
|
|
||||
List<Device> devices = pluginCtx.deviceService.findDevicesByTenantIdAndCustomerId(tenantId, customerId, new TextPageLink(limit)).getData(); |
|
||||
pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, devices), ActorRef.noSender()); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void sendRpcRequest(ToDeviceRpcRequest msg) { |
|
||||
pluginCtx.sendRpcRequest(msg); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void logRpcRequest(PluginApiCallSecurityContext ctx, DeviceId deviceId, ToDeviceRpcRequestBody body, boolean oneWay, Optional<RpcError> rpcError, Exception e) { |
|
||||
String rpcErrorStr = ""; |
|
||||
if (rpcError.isPresent()) { |
|
||||
rpcErrorStr = "RPC Error: " + rpcError.get().name(); |
|
||||
} |
|
||||
String method = body.getMethod(); |
|
||||
String params = body.getParams(); |
|
||||
pluginCtx.auditLogService.logEntityAction( |
|
||||
ctx.getTenantId(), |
|
||||
ctx.getCustomerId(), |
|
||||
ctx.getUserId(), |
|
||||
ctx.getUserName(), |
|
||||
deviceId, |
|
||||
null, |
|
||||
ActionType.RPC_CALL, |
|
||||
e, |
|
||||
rpcErrorStr, |
|
||||
new Boolean(oneWay), |
|
||||
method, |
|
||||
params); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void scheduleTimeoutMsg(TimeoutMsg msg) { |
|
||||
pluginCtx.scheduleTimeoutMsg(msg); |
|
||||
} |
|
||||
|
|
||||
|
|
||||
private void convertFuturesAndAddCallback(PluginCallback<List<AttributeKvEntry>> callback, List<ListenableFuture<List<AttributeKvEntry>>> futures) { |
|
||||
ListenableFuture<List<AttributeKvEntry>> future = Futures.transform(Futures.successfulAsList(futures), |
|
||||
(Function<? super List<List<AttributeKvEntry>>, ? extends List<AttributeKvEntry>>) input -> { |
|
||||
List<AttributeKvEntry> result = new ArrayList<>(); |
|
||||
input.forEach(r -> result.addAll(r)); |
|
||||
return result; |
|
||||
}, executor); |
|
||||
Futures.addCallback(future, getCallback(callback, v -> v), executor); |
|
||||
} |
|
||||
} |
|
||||
@ -1,30 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import org.thingsboard.server.actors.shared.ActorTerminationMsg; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.SessionId; |
|
||||
|
|
||||
/** |
|
||||
* @author Andrew Shvayka |
|
||||
*/ |
|
||||
public class PluginTerminationMsg extends ActorTerminationMsg<PluginId> { |
|
||||
|
|
||||
public PluginTerminationMsg(PluginId id) { |
|
||||
super(id); |
|
||||
} |
|
||||
} |
|
||||
@ -1,66 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.msg.aware.RuleAwareMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.RuleToPluginMsg; |
|
||||
import org.thingsboard.server.extensions.api.plugins.msg.ToPluginActorMsg; |
|
||||
|
|
||||
public class RuleToPluginMsgWrapper implements ToPluginActorMsg, RuleAwareMsg { |
|
||||
|
|
||||
private final TenantId pluginTenantId; |
|
||||
private final PluginId pluginId; |
|
||||
private final TenantId ruleTenantId; |
|
||||
private final RuleId ruleId; |
|
||||
private final RuleToPluginMsg<?> msg; |
|
||||
|
|
||||
public RuleToPluginMsgWrapper(TenantId pluginTenantId, PluginId pluginId, TenantId ruleTenantId, RuleId ruleId, RuleToPluginMsg<?> msg) { |
|
||||
super(); |
|
||||
this.pluginTenantId = pluginTenantId; |
|
||||
this.pluginId = pluginId; |
|
||||
this.ruleTenantId = ruleTenantId; |
|
||||
this.ruleId = ruleId; |
|
||||
this.msg = msg; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public TenantId getPluginTenantId() { |
|
||||
return pluginTenantId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public PluginId getPluginId() { |
|
||||
return pluginId; |
|
||||
} |
|
||||
|
|
||||
public TenantId getRuleTenantId() { |
|
||||
return ruleTenantId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public RuleId getRuleId() { |
|
||||
return ruleId; |
|
||||
} |
|
||||
|
|
||||
|
|
||||
public RuleToPluginMsg getMsg() { |
|
||||
return msg; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,142 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import akka.actor.ActorRef; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.common.data.id.DeviceId; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.msg.TbActorMsg; |
|
||||
import org.thingsboard.server.common.msg.cluster.ServerAddress; |
|
||||
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; |
|
||||
import org.thingsboard.server.common.msg.timeout.TimeoutMsg; |
|
||||
import org.thingsboard.server.controller.plugin.PluginWebSocketMsgEndpoint; |
|
||||
import org.thingsboard.server.dao.asset.AssetService; |
|
||||
import org.thingsboard.server.dao.attributes.AttributesService; |
|
||||
import org.thingsboard.server.dao.audit.AuditLogService; |
|
||||
import org.thingsboard.server.dao.customer.CustomerService; |
|
||||
import org.thingsboard.server.dao.device.DeviceService; |
|
||||
import org.thingsboard.server.dao.plugin.PluginService; |
|
||||
import org.thingsboard.server.dao.relation.RelationService; |
|
||||
import org.thingsboard.server.dao.rule.RuleChainService; |
|
||||
import org.thingsboard.server.dao.rule.RuleService; |
|
||||
import org.thingsboard.server.dao.tenant.TenantService; |
|
||||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
||||
import org.thingsboard.server.extensions.api.device.DeviceAttributesEventNotificationMsg; |
|
||||
import org.thingsboard.server.service.cluster.routing.ClusterRoutingService; |
|
||||
import org.thingsboard.server.service.cluster.rpc.ClusterRpcService; |
|
||||
import scala.concurrent.duration.Duration; |
|
||||
|
|
||||
import java.util.Optional; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.function.BiConsumer; |
|
||||
|
|
||||
@Slf4j |
|
||||
public final class SharedPluginProcessingContext { |
|
||||
final ActorRef parentActor; |
|
||||
final ActorRef currentActor; |
|
||||
final ActorSystemContext systemContext; |
|
||||
final PluginWebSocketMsgEndpoint msgEndpoint; |
|
||||
final AssetService assetService; |
|
||||
final DeviceService deviceService; |
|
||||
final RuleService ruleService; |
|
||||
final RuleChainService ruleChainService; |
|
||||
final PluginService pluginService; |
|
||||
final CustomerService customerService; |
|
||||
final TenantService tenantService; |
|
||||
final TimeseriesService tsService; |
|
||||
final AttributesService attributesService; |
|
||||
final ClusterRpcService rpcService; |
|
||||
final ClusterRoutingService routingService; |
|
||||
final RelationService relationService; |
|
||||
final AuditLogService auditLogService; |
|
||||
final PluginId pluginId; |
|
||||
final TenantId tenantId; |
|
||||
|
|
||||
public SharedPluginProcessingContext(ActorSystemContext sysContext, TenantId tenantId, PluginId pluginId, |
|
||||
ActorRef parentActor, ActorRef self) { |
|
||||
super(); |
|
||||
this.tenantId = tenantId; |
|
||||
this.pluginId = pluginId; |
|
||||
this.parentActor = parentActor; |
|
||||
this.currentActor = self; |
|
||||
this.systemContext = sysContext; |
|
||||
this.msgEndpoint = sysContext.getWsMsgEndpoint(); |
|
||||
this.tsService = sysContext.getTsService(); |
|
||||
this.attributesService = sysContext.getAttributesService(); |
|
||||
this.assetService = sysContext.getAssetService(); |
|
||||
this.deviceService = sysContext.getDeviceService(); |
|
||||
this.rpcService = sysContext.getRpcService(); |
|
||||
this.routingService = sysContext.getRoutingService(); |
|
||||
this.ruleService = sysContext.getRuleService(); |
|
||||
this.ruleChainService = sysContext.getRuleChainService(); |
|
||||
this.pluginService = sysContext.getPluginService(); |
|
||||
this.customerService = sysContext.getCustomerService(); |
|
||||
this.tenantService = sysContext.getTenantService(); |
|
||||
this.relationService = sysContext.getRelationService(); |
|
||||
this.auditLogService = sysContext.getAuditLogService(); |
|
||||
} |
|
||||
|
|
||||
public PluginId getPluginId() { |
|
||||
return pluginId; |
|
||||
} |
|
||||
|
|
||||
public TenantId getPluginTenantId() { |
|
||||
return tenantId; |
|
||||
} |
|
||||
|
|
||||
public void toDeviceActor(DeviceAttributesEventNotificationMsg msg) { |
|
||||
forward(msg.getDeviceId(), msg); |
|
||||
} |
|
||||
|
|
||||
public void sendRpcRequest(ToDeviceRpcRequest msg) { |
|
||||
log.trace("[{}] Forwarding msg {} to device actor!", pluginId, msg); |
|
||||
// ToDeviceRpcRequestPluginMsg rpcMsg = new ToDeviceRpcRequestPluginMsg(pluginId, tenantId, msg);
|
|
||||
// forward(msg.getDeviceId(), rpcMsg, rpcService::tell);
|
|
||||
} |
|
||||
|
|
||||
private <T extends TbActorMsg> void forward(DeviceId deviceId, T msg) { |
|
||||
Optional<ServerAddress> instance = routingService.resolveById(deviceId); |
|
||||
if (instance.isPresent()) { |
|
||||
log.trace("[{}] Forwarding msg {} to remote device actor!", pluginId, msg); |
|
||||
rpcService.tell(systemContext.getEncodingService().convertToProtoDataMessage(instance.get(), msg)); |
|
||||
} else { |
|
||||
log.trace("[{}] Forwarding msg {} to local device actor!", pluginId, msg); |
|
||||
parentActor.tell(msg, ActorRef.noSender()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void scheduleTimeoutMsg(TimeoutMsg msg) { |
|
||||
log.debug("Scheduling msg {} with delay {} ms", msg, msg.getTimeout()); |
|
||||
systemContext.getScheduler().scheduleOnce( |
|
||||
Duration.create(msg.getTimeout(), TimeUnit.MILLISECONDS), |
|
||||
currentActor, |
|
||||
msg, |
|
||||
systemContext.getActorSystem().dispatcher(), |
|
||||
ActorRef.noSender()); |
|
||||
|
|
||||
} |
|
||||
|
|
||||
public void persistError(String method, Exception e) { |
|
||||
systemContext.persistError(tenantId, pluginId, method, e); |
|
||||
} |
|
||||
|
|
||||
public ActorRef self() { |
|
||||
return currentActor; |
|
||||
} |
|
||||
} |
|
||||
@ -1,27 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import akka.actor.ActorRef; |
|
||||
|
|
||||
/** |
|
||||
* @author Andrew Shvayka |
|
||||
*/ |
|
||||
public interface TimeoutScheduler { |
|
||||
|
|
||||
void scheduleMsgWithDelay(Object msg, long delayInMs); |
|
||||
|
|
||||
} |
|
||||
@ -1,71 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.plugin; |
|
||||
|
|
||||
import java.util.function.Consumer; |
|
||||
import org.thingsboard.server.extensions.api.exception.AccessDeniedException; |
|
||||
import org.thingsboard.server.extensions.api.exception.EntityNotFoundException; |
|
||||
import org.thingsboard.server.extensions.api.exception.InternalErrorException; |
|
||||
import org.thingsboard.server.extensions.api.exception.UnauthorizedException; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginCallback; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginContext; |
|
||||
|
|
||||
/** |
|
||||
* Created by ashvayka on 21.02.17. |
|
||||
*/ |
|
||||
public class ValidationCallback implements PluginCallback<ValidationResult> { |
|
||||
|
|
||||
private final PluginCallback<?> callback; |
|
||||
private final Consumer<PluginContext> action; |
|
||||
|
|
||||
public ValidationCallback(PluginCallback<?> callback, Consumer<PluginContext> action) { |
|
||||
this.callback = callback; |
|
||||
this.action = action; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onSuccess(PluginContext ctx, ValidationResult result) { |
|
||||
ValidationResultCode resultCode = result.getResultCode(); |
|
||||
if (resultCode == ValidationResultCode.OK) { |
|
||||
action.accept(ctx); |
|
||||
} else { |
|
||||
Exception e; |
|
||||
switch (resultCode) { |
|
||||
case ENTITY_NOT_FOUND: |
|
||||
e = new EntityNotFoundException(result.getMessage()); |
|
||||
break; |
|
||||
case UNAUTHORIZED: |
|
||||
e = new UnauthorizedException(result.getMessage()); |
|
||||
break; |
|
||||
case ACCESS_DENIED: |
|
||||
e = new AccessDeniedException(result.getMessage()); |
|
||||
break; |
|
||||
case INTERNAL_ERROR: |
|
||||
e = new InternalErrorException(result.getMessage()); |
|
||||
break; |
|
||||
default: |
|
||||
e = new UnauthorizedException("Permission denied."); |
|
||||
break; |
|
||||
} |
|
||||
onFailure(ctx, e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(PluginContext ctx, Exception e) { |
|
||||
callback.onFailure(ctx, e); |
|
||||
} |
|
||||
} |
|
||||
@ -1,24 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.service; |
|
||||
|
|
||||
import org.thingsboard.server.extensions.api.plugins.rest.PluginRestMsg; |
|
||||
|
|
||||
public interface RestMsgProcessor { |
|
||||
|
|
||||
void process(PluginRestMsg msg); |
|
||||
|
|
||||
} |
|
||||
@ -1,24 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.service; |
|
||||
|
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg; |
|
||||
|
|
||||
public interface WebSocketMsgProcessor { |
|
||||
|
|
||||
void process(PluginWebsocketMsg<?> msg); |
|
||||
|
|
||||
} |
|
||||
@ -1,42 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.shared.plugin; |
|
||||
|
|
||||
import akka.japi.Creator; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.plugin.PluginActor; |
|
||||
import org.thingsboard.server.actors.shared.EntityActorsManager; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.dao.plugin.PluginService; |
|
||||
|
|
||||
@Slf4j |
|
||||
public abstract class PluginManager extends EntityActorsManager<PluginId, PluginActor, PluginMetaData> { |
|
||||
|
|
||||
protected final PluginService pluginService; |
|
||||
|
|
||||
public PluginManager(ActorSystemContext systemContext) { |
|
||||
super(systemContext); |
|
||||
this.pluginService = systemContext.getPluginService(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Creator<PluginActor> creator(PluginId entityId){ |
|
||||
return new PluginActor.ActorCreator(systemContext, getTenantId(), entityId); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,45 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.shared.plugin; |
|
||||
|
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.service.DefaultActorService; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.dao.plugin.BasePluginService; |
|
||||
|
|
||||
public class SystemPluginManager extends PluginManager { |
|
||||
|
|
||||
public SystemPluginManager(ActorSystemContext systemContext) { |
|
||||
super(systemContext); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected FetchFunction<PluginMetaData> getFetchEntitiesFunction() { |
|
||||
return pluginService::findSystemPlugins; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected TenantId getTenantId() { |
|
||||
return BasePluginService.SYSTEM_TENANT; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected String getDispatcherName() { |
|
||||
return DefaultActorService.SYSTEM_PLUGIN_DISPATCHER_NAME; |
|
||||
} |
|
||||
} |
|
||||
@ -1,57 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.actors.shared.plugin; |
|
||||
|
|
||||
import akka.actor.ActorContext; |
|
||||
import org.thingsboard.server.actors.ActorSystemContext; |
|
||||
import org.thingsboard.server.actors.service.DefaultActorService; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable; |
|
||||
import org.thingsboard.server.common.data.page.PageDataIterable.FetchFunction; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
|
|
||||
public class TenantPluginManager extends PluginManager { |
|
||||
|
|
||||
private final TenantId tenantId; |
|
||||
|
|
||||
public TenantPluginManager(ActorSystemContext systemContext, TenantId tenantId) { |
|
||||
super(systemContext); |
|
||||
this.tenantId = tenantId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void init(ActorContext context) { |
|
||||
if (systemContext.isTenantComponentsInitEnabled()) { |
|
||||
super.init(context); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected FetchFunction<PluginMetaData> getFetchEntitiesFunction() { |
|
||||
return link -> pluginService.findTenantPlugins(tenantId, link); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected TenantId getTenantId() { |
|
||||
return tenantId; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected String getDispatcherName() { |
|
||||
return DefaultActorService.TENANT_PLUGIN_DISPATCHER_NAME; |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,240 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.controller; |
|
||||
|
|
||||
import org.springframework.http.HttpStatus; |
|
||||
import org.springframework.security.access.prepost.PreAuthorize; |
|
||||
import org.springframework.web.bind.annotation.*; |
|
||||
import org.thingsboard.server.common.data.EntityType; |
|
||||
import org.thingsboard.server.common.data.audit.ActionType; |
|
||||
import org.thingsboard.server.common.data.id.PluginId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.TextPageData; |
|
||||
import org.thingsboard.server.common.data.page.TextPageLink; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.common.data.security.Authority; |
|
||||
import org.thingsboard.server.dao.model.ModelConstants; |
|
||||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|
||||
|
|
||||
import java.util.List; |
|
||||
|
|
||||
@RestController |
|
||||
@RequestMapping("/api") |
|
||||
public class PluginController extends BaseController { |
|
||||
|
|
||||
public static final String PLUGIN_ID = "pluginId"; |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/{pluginId}", method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public PluginMetaData getPluginById(@PathVariable(PLUGIN_ID) String strPluginId) throws ThingsboardException { |
|
||||
checkParameter(PLUGIN_ID, strPluginId); |
|
||||
try { |
|
||||
PluginId pluginId = new PluginId(toUUID(strPluginId)); |
|
||||
return checkPlugin(pluginService.findPluginById(pluginId)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/token/{pluginToken}", method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public PluginMetaData getPluginByToken(@PathVariable("pluginToken") String pluginToken) throws ThingsboardException { |
|
||||
checkParameter("pluginToken", pluginToken); |
|
||||
try { |
|
||||
return checkPlugin(pluginService.findPluginByApiToken(pluginToken)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin", method = RequestMethod.POST) |
|
||||
@ResponseBody |
|
||||
public PluginMetaData savePlugin(@RequestBody PluginMetaData source) throws ThingsboardException { |
|
||||
try { |
|
||||
boolean created = source.getId() == null; |
|
||||
source.setTenantId(getCurrentUser().getTenantId()); |
|
||||
PluginMetaData plugin = checkNotNull(pluginService.savePlugin(source)); |
|
||||
actorService.onEntityStateChange(plugin.getTenantId(), plugin.getId(), |
|
||||
created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); |
|
||||
|
|
||||
logEntityAction(plugin.getId(), plugin, |
|
||||
null, |
|
||||
created ? ActionType.ADDED : ActionType.UPDATED, null); |
|
||||
|
|
||||
return plugin; |
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.PLUGIN), source, |
|
||||
null, source.getId() == null ? ActionType.ADDED : ActionType.UPDATED, e); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/{pluginId}/activate", method = RequestMethod.POST) |
|
||||
@ResponseStatus(value = HttpStatus.OK) |
|
||||
public void activatePluginById(@PathVariable(PLUGIN_ID) String strPluginId) throws ThingsboardException { |
|
||||
checkParameter(PLUGIN_ID, strPluginId); |
|
||||
try { |
|
||||
PluginId pluginId = new PluginId(toUUID(strPluginId)); |
|
||||
PluginMetaData plugin = checkPlugin(pluginService.findPluginById(pluginId)); |
|
||||
pluginService.activatePluginById(pluginId); |
|
||||
actorService.onEntityStateChange(plugin.getTenantId(), plugin.getId(), ComponentLifecycleEvent.ACTIVATED); |
|
||||
|
|
||||
logEntityAction(plugin.getId(), plugin, |
|
||||
null, |
|
||||
ActionType.ACTIVATED, null, strPluginId); |
|
||||
|
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.PLUGIN), |
|
||||
null, |
|
||||
null, |
|
||||
ActionType.ACTIVATED, e, strPluginId); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/{pluginId}/suspend", method = RequestMethod.POST) |
|
||||
@ResponseStatus(value = HttpStatus.OK) |
|
||||
public void suspendPluginById(@PathVariable(PLUGIN_ID) String strPluginId) throws ThingsboardException { |
|
||||
checkParameter(PLUGIN_ID, strPluginId); |
|
||||
try { |
|
||||
PluginId pluginId = new PluginId(toUUID(strPluginId)); |
|
||||
PluginMetaData plugin = checkPlugin(pluginService.findPluginById(pluginId)); |
|
||||
pluginService.suspendPluginById(pluginId); |
|
||||
actorService.onEntityStateChange(plugin.getTenantId(), plugin.getId(), ComponentLifecycleEvent.SUSPENDED); |
|
||||
|
|
||||
logEntityAction(plugin.getId(), plugin, |
|
||||
null, |
|
||||
ActionType.SUSPENDED, null, strPluginId); |
|
||||
|
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.PLUGIN), |
|
||||
null, |
|
||||
null, |
|
||||
ActionType.SUSPENDED, e, strPluginId); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAuthority('SYS_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/system", params = {"limit"}, method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public TextPageData<PluginMetaData> getSystemPlugins( |
|
||||
@RequestParam int limit, |
|
||||
@RequestParam(required = false) String textSearch, |
|
||||
@RequestParam(required = false) String idOffset, |
|
||||
@RequestParam(required = false) String textOffset) throws ThingsboardException { |
|
||||
try { |
|
||||
TextPageLink pageLink = createPageLink(limit, textSearch, idOffset, textOffset); |
|
||||
return checkNotNull(pluginService.findSystemPlugins(pageLink)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAuthority('SYS_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/tenant/{tenantId}", params = {"limit"}, method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public TextPageData<PluginMetaData> getTenantPlugins( |
|
||||
@PathVariable("tenantId") String strTenantId, |
|
||||
@RequestParam int limit, |
|
||||
@RequestParam(required = false) String textSearch, |
|
||||
@RequestParam(required = false) String idOffset, |
|
||||
@RequestParam(required = false) String textOffset) throws ThingsboardException { |
|
||||
checkParameter("tenantId", strTenantId); |
|
||||
try { |
|
||||
TenantId tenantId = new TenantId(toUUID(strTenantId)); |
|
||||
TextPageLink pageLink = createPageLink(limit, textSearch, idOffset, textOffset); |
|
||||
return checkNotNull(pluginService.findTenantPlugins(tenantId, pageLink)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugins", method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public List<PluginMetaData> getPlugins() throws ThingsboardException { |
|
||||
try { |
|
||||
if (getCurrentUser().getAuthority() == Authority.SYS_ADMIN) { |
|
||||
return checkNotNull(pluginService.findSystemPlugins()); |
|
||||
} else { |
|
||||
TenantId tenantId = getCurrentUser().getTenantId(); |
|
||||
List<PluginMetaData> plugins = checkNotNull(pluginService.findAllTenantPluginsByTenantId(tenantId)); |
|
||||
plugins.stream() |
|
||||
.filter(plugin -> plugin.getTenantId().getId().equals(ModelConstants.NULL_UUID)) |
|
||||
.forEach(plugin -> plugin.setConfiguration(null)); |
|
||||
return plugins; |
|
||||
} |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin", params = {"limit"}, method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public TextPageData<PluginMetaData> getTenantPlugins( |
|
||||
@RequestParam int limit, |
|
||||
@RequestParam(required = false) String textSearch, |
|
||||
@RequestParam(required = false) String idOffset, |
|
||||
@RequestParam(required = false) String textOffset) throws ThingsboardException { |
|
||||
try { |
|
||||
TenantId tenantId = getCurrentUser().getTenantId(); |
|
||||
TextPageLink pageLink = createPageLink(limit, textSearch, idOffset, textOffset); |
|
||||
return checkNotNull(pluginService.findTenantPlugins(tenantId, pageLink)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/plugin/{pluginId}", method = RequestMethod.DELETE) |
|
||||
@ResponseStatus(value = HttpStatus.OK) |
|
||||
public void deletePlugin(@PathVariable(PLUGIN_ID) String strPluginId) throws ThingsboardException { |
|
||||
checkParameter(PLUGIN_ID, strPluginId); |
|
||||
try { |
|
||||
PluginId pluginId = new PluginId(toUUID(strPluginId)); |
|
||||
PluginMetaData plugin = checkPlugin(pluginService.findPluginById(pluginId)); |
|
||||
pluginService.deletePluginById(pluginId); |
|
||||
actorService.onEntityStateChange(plugin.getTenantId(), plugin.getId(), ComponentLifecycleEvent.DELETED); |
|
||||
|
|
||||
logEntityAction(pluginId, plugin, |
|
||||
null, |
|
||||
ActionType.DELETED, null, strPluginId); |
|
||||
|
|
||||
} catch (Exception e) { |
|
||||
logEntityAction(emptyId(EntityType.PLUGIN), |
|
||||
null, |
|
||||
null, |
|
||||
ActionType.DELETED, e, strPluginId); |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
|
|
||||
} |
|
||||
@ -1,239 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.controller; |
|
||||
|
|
||||
import org.springframework.http.HttpStatus; |
|
||||
import org.springframework.security.access.prepost.PreAuthorize; |
|
||||
import org.springframework.web.bind.annotation.*; |
|
||||
import org.thingsboard.server.common.data.EntityType; |
|
||||
import org.thingsboard.server.common.data.audit.ActionType; |
|
||||
import org.thingsboard.server.common.data.id.RuleId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.TextPageData; |
|
||||
import org.thingsboard.server.common.data.page.TextPageLink; |
|
||||
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
|
||||
import org.thingsboard.server.common.data.plugin.PluginMetaData; |
|
||||
import org.thingsboard.server.common.data.rule.RuleMetaData; |
|
||||
import org.thingsboard.server.common.data.security.Authority; |
|
||||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|
||||
|
|
||||
import java.util.List; |
|
||||
|
|
||||
@RestController |
|
||||
@RequestMapping("/api") |
|
||||
public class RuleController extends BaseController { |
|
||||
|
|
||||
public static final String RULE_ID = "ruleId"; |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/{ruleId}", method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public RuleMetaData getRuleById(@PathVariable(RULE_ID) String strRuleId) throws ThingsboardException { |
|
||||
checkParameter(RULE_ID, strRuleId); |
|
||||
try { |
|
||||
RuleId ruleId = new RuleId(toUUID(strRuleId)); |
|
||||
return checkRule(ruleService.findRuleById(ruleId)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/token/{pluginToken}", method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public List<RuleMetaData> getRulesByPluginToken(@PathVariable("pluginToken") String pluginToken) throws ThingsboardException { |
|
||||
checkParameter("pluginToken", pluginToken); |
|
||||
try { |
|
||||
PluginMetaData plugin = checkPlugin(pluginService.findPluginByApiToken(pluginToken)); |
|
||||
return ruleService.findPluginRules(plugin.getApiToken()); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule", method = RequestMethod.POST) |
|
||||
@ResponseBody |
|
||||
public RuleMetaData saveRule(@RequestBody RuleMetaData source) throws ThingsboardException { |
|
||||
try { |
|
||||
boolean created = source.getId() == null; |
|
||||
source.setTenantId(getCurrentUser().getTenantId()); |
|
||||
RuleMetaData rule = checkNotNull(ruleService.saveRule(source)); |
|
||||
actorService.onEntityStateChange(rule.getTenantId(), rule.getId(), |
|
||||
created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); |
|
||||
|
|
||||
logEntityAction(rule.getId(), rule, |
|
||||
null, |
|
||||
created ? ActionType.ADDED : ActionType.UPDATED, null); |
|
||||
|
|
||||
return rule; |
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.RULE), source, |
|
||||
null, source.getId() == null ? ActionType.ADDED : ActionType.UPDATED, e); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/{ruleId}/activate", method = RequestMethod.POST) |
|
||||
@ResponseStatus(value = HttpStatus.OK) |
|
||||
public void activateRuleById(@PathVariable(RULE_ID) String strRuleId) throws ThingsboardException { |
|
||||
checkParameter(RULE_ID, strRuleId); |
|
||||
try { |
|
||||
RuleId ruleId = new RuleId(toUUID(strRuleId)); |
|
||||
RuleMetaData rule = checkRule(ruleService.findRuleById(ruleId)); |
|
||||
ruleService.activateRuleById(ruleId); |
|
||||
actorService.onEntityStateChange(rule.getTenantId(), rule.getId(), ComponentLifecycleEvent.ACTIVATED); |
|
||||
|
|
||||
logEntityAction(rule.getId(), rule, |
|
||||
null, |
|
||||
ActionType.ACTIVATED, null, strRuleId); |
|
||||
|
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.RULE), |
|
||||
null, |
|
||||
null, |
|
||||
ActionType.ACTIVATED, e, strRuleId); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/{ruleId}/suspend", method = RequestMethod.POST) |
|
||||
@ResponseStatus(value = HttpStatus.OK) |
|
||||
public void suspendRuleById(@PathVariable(RULE_ID) String strRuleId) throws ThingsboardException { |
|
||||
checkParameter(RULE_ID, strRuleId); |
|
||||
try { |
|
||||
RuleId ruleId = new RuleId(toUUID(strRuleId)); |
|
||||
RuleMetaData rule = checkRule(ruleService.findRuleById(ruleId)); |
|
||||
ruleService.suspendRuleById(ruleId); |
|
||||
actorService.onEntityStateChange(rule.getTenantId(), rule.getId(), ComponentLifecycleEvent.SUSPENDED); |
|
||||
|
|
||||
logEntityAction(rule.getId(), rule, |
|
||||
null, |
|
||||
ActionType.SUSPENDED, null, strRuleId); |
|
||||
|
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.RULE), |
|
||||
null, |
|
||||
null, |
|
||||
ActionType.SUSPENDED, e, strRuleId); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAuthority('SYS_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/system", params = {"limit"}, method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public TextPageData<RuleMetaData> getSystemRules( |
|
||||
@RequestParam int limit, |
|
||||
@RequestParam(required = false) String textSearch, |
|
||||
@RequestParam(required = false) String idOffset, |
|
||||
@RequestParam(required = false) String textOffset) throws ThingsboardException { |
|
||||
try { |
|
||||
TextPageLink pageLink = createPageLink(limit, textSearch, idOffset, textOffset); |
|
||||
return checkNotNull(ruleService.findSystemRules(pageLink)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAuthority('SYS_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/tenant/{tenantId}", params = {"limit"}, method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public TextPageData<RuleMetaData> getTenantRules( |
|
||||
@PathVariable("tenantId") String strTenantId, |
|
||||
@RequestParam int limit, |
|
||||
@RequestParam(required = false) String textSearch, |
|
||||
@RequestParam(required = false) String idOffset, |
|
||||
@RequestParam(required = false) String textOffset) throws ThingsboardException { |
|
||||
checkParameter("tenantId", strTenantId); |
|
||||
try { |
|
||||
TenantId tenantId = new TenantId(toUUID(strTenantId)); |
|
||||
TextPageLink pageLink = createPageLink(limit, textSearch, idOffset, textOffset); |
|
||||
return checkNotNull(ruleService.findTenantRules(tenantId, pageLink)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rules", method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public List<RuleMetaData> getRules() throws ThingsboardException { |
|
||||
try { |
|
||||
if (getCurrentUser().getAuthority() == Authority.SYS_ADMIN) { |
|
||||
return checkNotNull(ruleService.findSystemRules()); |
|
||||
} else { |
|
||||
TenantId tenantId = getCurrentUser().getTenantId(); |
|
||||
return checkNotNull(ruleService.findAllTenantRulesByTenantId(tenantId)); |
|
||||
} |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule", params = {"limit"}, method = RequestMethod.GET) |
|
||||
@ResponseBody |
|
||||
public TextPageData<RuleMetaData> getTenantRules( |
|
||||
@RequestParam int limit, |
|
||||
@RequestParam(required = false) String textSearch, |
|
||||
@RequestParam(required = false) String idOffset, |
|
||||
@RequestParam(required = false) String textOffset) throws ThingsboardException { |
|
||||
try { |
|
||||
TenantId tenantId = getCurrentUser().getTenantId(); |
|
||||
TextPageLink pageLink = createPageLink(limit, textSearch, idOffset, textOffset); |
|
||||
return checkNotNull(ruleService.findTenantRules(tenantId, pageLink)); |
|
||||
} catch (Exception e) { |
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") |
|
||||
@RequestMapping(value = "/rule/{ruleId}", method = RequestMethod.DELETE) |
|
||||
@ResponseStatus(value = HttpStatus.OK) |
|
||||
public void deleteRule(@PathVariable(RULE_ID) String strRuleId) throws ThingsboardException { |
|
||||
checkParameter(RULE_ID, strRuleId); |
|
||||
try { |
|
||||
RuleId ruleId = new RuleId(toUUID(strRuleId)); |
|
||||
RuleMetaData rule = checkRule(ruleService.findRuleById(ruleId)); |
|
||||
ruleService.deleteRuleById(ruleId); |
|
||||
actorService.onEntityStateChange(rule.getTenantId(), rule.getId(), ComponentLifecycleEvent.DELETED); |
|
||||
|
|
||||
logEntityAction(ruleId, rule, |
|
||||
null, |
|
||||
ActionType.DELETED, null, strRuleId); |
|
||||
|
|
||||
} catch (Exception e) { |
|
||||
|
|
||||
logEntityAction(emptyId(EntityType.RULE), |
|
||||
null, |
|
||||
null, |
|
||||
ActionType.DELETED, e, strRuleId); |
|
||||
|
|
||||
throw handleException(e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,85 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.controller.plugin; |
|
||||
|
|
||||
|
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.web.bind.annotation.RequestMapping; |
|
||||
import org.springframework.web.bind.annotation.RestController; |
|
||||
import org.thingsboard.server.controller.BaseController; |
|
||||
import org.thingsboard.server.extensions.api.plugins.PluginConstants; |
|
||||
|
|
||||
@RestController |
|
||||
@RequestMapping(PluginConstants.PLUGIN_URL_PREFIX) |
|
||||
@Slf4j |
|
||||
public class PluginApiController extends BaseController { |
|
||||
|
|
||||
// @SuppressWarnings("rawtypes")
|
|
||||
// @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
|
|
||||
// @RequestMapping(value = "/{pluginToken}/**")
|
|
||||
// @ResponseStatus(value = HttpStatus.OK)
|
|
||||
// public DeferredResult<ResponseEntity> processRequest(
|
|
||||
// @PathVariable("pluginToken") String pluginToken,
|
|
||||
// RequestEntity<byte[]> requestEntity,
|
|
||||
// HttpServletRequest request)
|
|
||||
// throws ThingsboardException {
|
|
||||
// log.debug("[{}] Going to process requst uri: {}", pluginToken, requestEntity.getUrl());
|
|
||||
// DeferredResult<ResponseEntity> result = new DeferredResult<ResponseEntity>();
|
|
||||
// PluginMetaData pluginMd = pluginService.findPluginByApiToken(pluginToken);
|
|
||||
// if (pluginMd == null) {
|
|
||||
// result.setErrorResult(new PluginNotFoundException("Plugin with token: " + pluginToken + " not found!"));
|
|
||||
// } else {
|
|
||||
// TenantId tenantId = getCurrentUser().getTenantId();
|
|
||||
// CustomerId customerId = getCurrentUser().getCustomerId();
|
|
||||
// if (validatePluginAccess(pluginMd, tenantId, customerId)) {
|
|
||||
// if(tenantId != null && ModelConstants.NULL_UUID.equals(tenantId.getId())){
|
|
||||
// tenantId = null;
|
|
||||
// }
|
|
||||
// UserId userId = getCurrentUser().getId();
|
|
||||
// String userName = getCurrentUser().getName();
|
|
||||
// PluginApiCallSecurityContext securityCtx = new PluginApiCallSecurityContext(pluginMd.getTenantId(), pluginMd.getId(),
|
|
||||
// tenantId, customerId, userId, userName);
|
|
||||
// actorService.process(new BasicPluginRestMsg(securityCtx, new RestRequest(requestEntity, request), result));
|
|
||||
// } else {
|
|
||||
// result.setResult(new ResponseEntity<>(HttpStatus.FORBIDDEN));
|
|
||||
// }
|
|
||||
//
|
|
||||
// }
|
|
||||
// return result;
|
|
||||
// }
|
|
||||
//
|
|
||||
// public static boolean validatePluginAccess(PluginMetaData pluginMd, TenantId tenantId, CustomerId customerId) {
|
|
||||
// boolean systemAdministrator = tenantId == null || ModelConstants.NULL_UUID.equals(tenantId.getId());
|
|
||||
// boolean tenantAdministrator = !systemAdministrator && (customerId == null || ModelConstants.NULL_UUID.equals(customerId.getId()));
|
|
||||
// boolean systemPlugin = ModelConstants.NULL_UUID.equals(pluginMd.getTenantId().getId());
|
|
||||
//
|
|
||||
// boolean validUser = false;
|
|
||||
// if (systemPlugin) {
|
|
||||
// if (pluginMd.isPublicAccess() || systemAdministrator) {
|
|
||||
// // All users can access public system plugins. Only system
|
|
||||
// // users can access private system plugins
|
|
||||
// validUser = true;
|
|
||||
// }
|
|
||||
// } else {
|
|
||||
// if ((pluginMd.isPublicAccess() || tenantAdministrator) && tenantId != null && tenantId.equals(pluginMd.getTenantId())) {
|
|
||||
// // All tenant users can access public tenant plugins. Only tenant
|
|
||||
// // administrator can access private tenant plugins
|
|
||||
// validUser = true;
|
|
||||
// }
|
|
||||
// }
|
|
||||
// return validUser;
|
|
||||
// }
|
|
||||
} |
|
||||
@ -1,30 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.controller.plugin; |
|
||||
|
|
||||
import org.springframework.http.HttpStatus; |
|
||||
import org.springframework.web.bind.annotation.ResponseStatus; |
|
||||
|
|
||||
@ResponseStatus(HttpStatus.NOT_FOUND) |
|
||||
public class PluginNotFoundException extends RuntimeException { |
|
||||
|
|
||||
private static final long serialVersionUID = 1L; |
|
||||
|
|
||||
public PluginNotFoundException(String message){ |
|
||||
super(message); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,28 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2018 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.controller.plugin; |
|
||||
|
|
||||
import java.io.IOException; |
|
||||
|
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.PluginWebsocketSessionRef; |
|
||||
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg; |
|
||||
|
|
||||
public interface PluginWebSocketMsgEndpoint { |
|
||||
|
|
||||
void send(PluginWebsocketMsg<?> wsMsg) throws IOException; |
|
||||
|
|
||||
void close(PluginWebsocketSessionRef sessionRef) throws IOException; |
|
||||
} |
|
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue