This is an automated email from the ASF dual-hosted git repository.
riemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 6ea881fe7 Improve monitoring of adapters and pipelines (#1787) (#1788)
6ea881fe7 is described below
commit 6ea881fe74bda86d39bd9f924a4a01700e2158b5
Author: Dominik Riemer <[email protected]>
AuthorDate: Mon Jul 24 17:33:55 2023 +0200
Improve monitoring of adapters and pipelines (#1787) (#1788)
* Improve monitoring of adapters and pipelines (#1787)
* Improve layout of exception dialog
---
.../connect/context/IAdapterRuntimeContext.java | 4 +-
.../extensions/api/monitoring/FixedSizeList.java | 49 ++++++++------
...neElementLogger.java => IExtensionsLogger.java} | 17 ++++-
.../api/monitoring/SpMonitoringManager.java | 33 +++++----
.../extensions/api/pe/context/RuntimeContext.java | 4 +-
.../connect/AdapterWorkerManagement.java | 17 +++--
.../management/connect/PullAdapterScheduler.java | 4 +-
.../elements/SendToBrokerAdapterSink.java | 7 +-
.../context/AdapterContextGenerator.java | 9 ++-
.../context/SpAdapterRuntimeContext.java | 12 ++--
.../management/monitoring/ExtensionsLogger.java | 63 +++++++++++++++++
.../connect/AdapterWorkerManagementTest.java | 2 +-
.../iiot/adapters/iolink/IfmAlMqttAdapter.java | 6 +-
.../simulator/machine/MachineDataSimulator.java | 12 +---
.../iiot/protocol/stream/FileReplayAdapter.java | 5 +-
.../complex/TopologyValidationProcessor.java | 4 +-
.../streampipes/model/monitoring/SpLogEntry.java | 16 +++--
...{StreamPipesStatistics.java => SpLogLevel.java} | 6 +-
.../SpLogMessage.java} | 79 +++++++++++++++-------
.../monitoring/pipeline/ExtensionsLogProvider.java | 6 +-
.../apache/streampipes/ps/DataLakeResourceV4.java | 4 +-
.../extensions/connect/AdapterWorkerResource.java | 10 ++-
.../extensions/monitoring/MonitoringResource.java | 2 +-
.../rest/impl/connect/AdapterResource.java | 6 +-
.../rest/impl/connect/GuessResource.java | 8 +--
.../impl/connect/RuntimeResolvableResource.java | 8 +--
.../rest/impl/pe/DataProcessorResource.java | 4 +-
.../streampipes/rest/impl/pe/DataSinkResource.java | 4 +-
.../rest/impl/pe/DataStreamResource.java | 4 +-
.../standalone/function/FunctionContext.java | 10 +--
.../standalone/function/StreamPipesFunction.java | 4 +-
.../routing/StandaloneSpOutputCollector.java | 8 +--
.../runtime/StandaloneEventProcessorRuntime.java | 4 +-
.../runtime/StandaloneEventSinkRuntime.java | 4 +-
.../runtime/StandalonePipelineElementRuntime.java | 20 +++---
.../context/SpEventProcessorRuntimeContext.java | 19 ++----
.../wrapper/context/SpEventSinkRuntimeContext.java | 14 ++--
.../wrapper/context/SpRuntimeContext.java | 34 ++++------
.../generator/DataProcessorContextGenerator.java | 4 +-
.../generator/DataSinkContextGenerator.java | 4 +-
.../src/lib/model/gen/streampipes-model.ts | 55 +++++++--------
.../basic-nav-tabs/basic-nav-tabs.component.html | 7 +-
.../exception-details-dialog.component.html | 6 +-
.../exception-details-dialog.component.ts | 4 +-
.../sp-exception-message.component.html | 4 +-
.../sp-exception-message.component.ts | 4 +-
.../error-message/error-message.component.ts | 4 +-
.../event-schema/event-schema.component.ts | 4 +-
.../existing-adapters.component.ts | 7 +-
.../base-runtime-resolvable-input.ts | 7 +-
.../data-explorer-dashboard-widget.component.html | 1 +
.../data-explorer-dashboard-widget.component.ts | 4 +-
.../base/base-data-explorer-widget.directive.ts | 6 +-
.../widgets/base/data-explorer-widget-data.ts | 4 +-
.../components/startup/startup.component.scss | 4 +-
55 files changed, 379 insertions(+), 272 deletions(-)
diff --git
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/context/IAdapterRuntimeContext.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/context/IAdapterRuntimeContext.java
index 19f14b4d9..6b9035b40 100644
---
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/context/IAdapterRuntimeContext.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/connect/context/IAdapterRuntimeContext.java
@@ -18,10 +18,10 @@
package org.apache.streampipes.extensions.api.connect.context;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
public interface IAdapterRuntimeContext extends IAdapterGuessSchemaContext {
- SpMonitoringManager getLogger();
+ IExtensionsLogger getLogger();
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/StreamPipesRuntimeError.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/FixedSizeList.java
similarity index 52%
rename from
streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/StreamPipesRuntimeError.java
rename to
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/FixedSizeList.java
index eccf51de9..75f0a9f73 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/StreamPipesRuntimeError.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/FixedSizeList.java
@@ -16,36 +16,47 @@
*
*/
-package org.apache.streampipes.model.monitoring;
+package org.apache.streampipes.extensions.api.monitoring;
-public class StreamPipesRuntimeError {
+import java.util.ArrayList;
+import java.util.List;
- private long timestamp;
- private String title;
- private String message;
+public class FixedSizeList<T> {
+ private final int maxSize;
+ private final List<T> list;
- private String stackTrace;
+ public FixedSizeList(int maxSize) {
+ if (maxSize <= 0) {
+ throw new IllegalArgumentException("Max size must be greater than
zero.");
+ }
+ this.maxSize = maxSize;
+ this.list = new ArrayList<>();
+ }
- public StreamPipesRuntimeError(long timestamp, String title, String message,
String stackTrace) {
- this.timestamp = timestamp;
- this.title = title;
- this.message = message;
- this.stackTrace = stackTrace;
+ public void add(T element) {
+ list.add(0, element);
+ if (list.size() > maxSize) {
+ list.remove(list.size() - 1);
+ }
}
- public long getTimestamp() {
- return timestamp;
+ public T get(int index) {
+ return list.get(index);
}
- public String getTitle() {
- return title;
+ public int size() {
+ return list.size();
}
- public String getMessage() {
- return message;
+ public List<T> getAllItems() {
+ return this.list;
}
- public String getStackTrace() {
- return stackTrace;
+ public void clear() {
+ this.list.clear();
}
}
+
+
+
+
diff --git
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/IPipelineElementLogger.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/IExtensionsLogger.java
similarity index 72%
rename from
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/IPipelineElementLogger.java
rename to
streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/IExtensionsLogger.java
index e9418c41d..adade1251 100644
---
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/IPipelineElementLogger.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/IExtensionsLogger.java
@@ -18,7 +18,20 @@
package org.apache.streampipes.extensions.api.monitoring;
-import org.slf4j.Logger;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
-public interface IPipelineElementLogger extends Logger {
+public interface IExtensionsLogger {
+
+ void log(SpLogMessage logMessage);
+
+ void error(Exception e);
+
+ void error(String details,
+ Exception e);
+
+ void info(String title,
+ String details);
+
+ void warn(String title,
+ String details);
}
diff --git
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
index 8c15d8e8a..5f5e6fa6f 100644
---
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/monitoring/SpMonitoringManager.java
@@ -22,7 +22,6 @@ import
org.apache.streampipes.model.monitoring.SpEndpointMonitoringInfo;
import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.model.monitoring.SpMetricsEntry;
-import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -31,7 +30,7 @@ public enum SpMonitoringManager {
INSTANCE;
- private final Map<String, List<SpLogEntry>> logInfos;
+ private final Map<String, FixedSizeList<SpLogEntry>> logInfos;
private final Map<String, SpMetricsEntry> metricsInfos;
SpMonitoringManager() {
@@ -42,9 +41,9 @@ public enum SpMonitoringManager {
public void addErrorMessage(String resourceId,
SpLogEntry errorMessageEntry) {
if (!logInfos.containsKey(resourceId)) {
- logInfos.put(resourceId, new ArrayList<>());
+ logInfos.put(resourceId, new FixedSizeList<>(100));
}
- this.logInfos.get(resourceId).add(0, errorMessageEntry);
+ this.logInfos.get(resourceId).add(errorMessageEntry);
}
public void increaseInCounter(String resourceId,
@@ -87,18 +86,28 @@ public enum SpMonitoringManager {
return currentEntry;
}
- public Map<String, List<SpLogEntry>> getAllLogs() {
- return logInfos;
+ public SpEndpointMonitoringInfo getMonitoringInfo() {
+ var logInfos = makeLogInfos();
+ return new SpEndpointMonitoringInfo(logInfos, metricsInfos);
}
- public Map<String, SpMetricsEntry> getAllMetrics() {
- return this.metricsInfos;
+ public void clearAllLogs() {
+ this.logInfos.forEach((key, value) -> value.clear());
}
- public SpEndpointMonitoringInfo getMonitoringInfo() {
- return new SpEndpointMonitoringInfo(logInfos, metricsInfos);
+ private Map<String, List<SpLogEntry>> makeLogInfos() {
+ var logEntries = new HashMap<String, List<SpLogEntry>>();
+ this.logInfos.forEach((key, value) ->
+ logEntries.put(key, cloneList(value.getAllItems())));
+
+ return logEntries;
}
+ private List<SpLogEntry> cloneList(List<SpLogEntry> allItems) {
+ return allItems.stream().map(SpLogEntry::new).toList();
+ }
+
+
private void checkAndPrepareMetrics(String resourceId) {
if (!metricsInfos.containsKey(resourceId)) {
addMetricsObject(resourceId);
@@ -109,8 +118,4 @@ public enum SpMonitoringManager {
this.metricsInfos.put(resourceId, new SpMetricsEntry());
}
-
- public void clearAllLogs() {
- logInfos.forEach((key, value) -> value.clear());
- }
}
diff --git
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/pe/context/RuntimeContext.java
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/pe/context/RuntimeContext.java
index 5c6e1ae71..390d50318 100644
---
a/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/pe/context/RuntimeContext.java
+++
b/streampipes-extensions-api/src/main/java/org/apache/streampipes/extensions/api/pe/context/RuntimeContext.java
@@ -19,11 +19,11 @@ package org.apache.streampipes.extensions.api.pe.context;
import org.apache.streampipes.client.api.IStreamPipesClient;
import org.apache.streampipes.extensions.api.config.IConfigExtractor;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
public interface RuntimeContext {
- SpMonitoringManager getLogger();
+ IExtensionsLogger getLogger();
String getCorrespondingUser();
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
index 65a7c8ecf..958d413f1 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
@@ -23,6 +23,7 @@ import
org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
import
org.apache.streampipes.extensions.api.connect.context.IAdapterRuntimeContext;
import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import
org.apache.streampipes.extensions.management.connect.adapter.model.EventCollector;
+import
org.apache.streampipes.extensions.management.context.AdapterContextGenerator;
import org.apache.streampipes.extensions.management.init.IDeclarersSingleton;
import
org.apache.streampipes.extensions.management.init.RunningAdapterInstances;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
@@ -41,14 +42,10 @@ public class AdapterWorkerManagement {
private final RunningAdapterInstances runningAdapterInstances;
private final IDeclarersSingleton declarers;
- private final IAdapterRuntimeContext adapterRuntimeContext;
-
public AdapterWorkerManagement(RunningAdapterInstances
runningAdapterInstances,
- IDeclarersSingleton declarers,
- IAdapterRuntimeContext runtimeContext) {
+ IDeclarersSingleton declarers) {
this.runningAdapterInstances = runningAdapterInstances;
this.declarers = declarers;
- this.adapterRuntimeContext = runtimeContext;
}
public Collection<AdapterDescription> getAllRunningAdapterInstances() {
@@ -69,8 +66,9 @@ public class AdapterWorkerManagement {
var registeredParsers =
newAdapterInstance.declareConfig().getSupportedParsers();
var extractor = AdapterParameterExtractor.from(adapterDescription,
registeredParsers);
var eventCollector = EventCollector.from(adapterDescription);
+ var runtimeContext =
makeRuntimeContext(adapterDescription.getElementId());
- newAdapterInstance.onAdapterStarted(extractor, eventCollector,
adapterRuntimeContext);
+ newAdapterInstance.onAdapterStarted(extractor, eventCollector,
runtimeContext);
} else {
var errorMessage = "Adapter with id %s could not be
found".formatted(adapterDescription.getAppId());
LOG.error(errorMessage);
@@ -88,12 +86,17 @@ public class AdapterWorkerManagement {
var registeredParsers = adapter.declareConfig().getSupportedParsers();
var extractor = AdapterParameterExtractor.from(adapterDescription,
registeredParsers);
- adapter.onAdapterStopped(extractor, adapterRuntimeContext);
+ var runtimeContext = makeRuntimeContext(elementId);
+ adapter.onAdapterStopped(extractor, runtimeContext);
}
resetMonitoring(elementId);
}
+ private IAdapterRuntimeContext makeRuntimeContext(String adapterInstanceId) {
+ return new AdapterContextGenerator().makeRuntimeContext(adapterInstanceId);
+ }
+
private void resetMonitoring(String elementId) {
SpMonitoringManager.INSTANCE.reset(elementId);
}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/PullAdapterScheduler.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/PullAdapterScheduler.java
index eeb1aaf62..de9a1b414 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/PullAdapterScheduler.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/PullAdapterScheduler.java
@@ -20,8 +20,8 @@ package org.apache.streampipes.extensions.management.connect;
import org.apache.streampipes.extensions.api.connect.IPullAdapter;
import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.monitoring.SpLogEntry;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -45,7 +45,7 @@ public class PullAdapterScheduler {
} catch (ExecutionException | InterruptedException e) {
SpMonitoringManager.INSTANCE.addErrorMessage(
adapterElementId,
- SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e)));
+ SpLogEntry.from(System.currentTimeMillis(), SpLogMessage.from(e)));
} catch (TimeoutException e) {
LOGGER.warn("Timeout occurred", e);
}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
index 52f864db7..b9cebe258 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
@@ -24,12 +24,11 @@ import
org.apache.streampipes.dataformat.SpDataFormatDefinition;
import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import
org.apache.streampipes.extensions.management.connect.adapter.util.TransportFormatSelector;
+import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import org.apache.streampipes.messaging.EventProducer;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.grounding.TransportFormat;
import org.apache.streampipes.model.grounding.TransportProtocol;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
import java.util.Map;
@@ -78,9 +77,7 @@ public abstract class SendToBrokerAdapterSink<T extends
TransportProtocol> imple
System.currentTimeMillis());
}
} catch (RuntimeException e) {
- SpMonitoringManager.INSTANCE.addErrorMessage(
- adapterDescription.getElementId(),
- SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e)));
+ new ExtensionsLogger(adapterDescription.getElementId()).error(e);
}
return null;
}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/AdapterContextGenerator.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/AdapterContextGenerator.java
index 84a3aa346..77e24c234 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/AdapterContextGenerator.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/AdapterContextGenerator.java
@@ -20,15 +20,18 @@ package
org.apache.streampipes.extensions.management.context;
import
org.apache.streampipes.extensions.api.connect.context.IAdapterGuessSchemaContext;
import
org.apache.streampipes.extensions.api.connect.context.IAdapterRuntimeContext;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import static
org.apache.streampipes.extensions.management.util.RuntimeContextUtils.makeConfigExtractor;
import static
org.apache.streampipes.extensions.management.util.RuntimeContextUtils.makeStreamPipesClient;
public class AdapterContextGenerator {
- public IAdapterRuntimeContext makeRuntimeContext() {
- return new SpAdapterRuntimeContext(SpMonitoringManager.INSTANCE,
makeConfigExtractor(), makeStreamPipesClient());
+ public IAdapterRuntimeContext makeRuntimeContext(String adapterInstanceId) {
+ return new SpAdapterRuntimeContext(
+ new ExtensionsLogger(adapterInstanceId),
+ makeConfigExtractor(),
+ makeStreamPipesClient());
}
public IAdapterGuessSchemaContext makeGuessSchemaContext() {
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/SpAdapterRuntimeContext.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/SpAdapterRuntimeContext.java
index f0dd95b75..2fc6fd95e 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/SpAdapterRuntimeContext.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/context/SpAdapterRuntimeContext.java
@@ -20,7 +20,7 @@ package org.apache.streampipes.extensions.management.context;
import org.apache.streampipes.client.StreamPipesClient;
import
org.apache.streampipes.extensions.api.connect.context.IAdapterRuntimeContext;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
import org.apache.streampipes.extensions.management.config.ConfigExtractor;
import java.io.Serializable;
@@ -28,17 +28,17 @@ import java.io.Serializable;
public class SpAdapterRuntimeContext extends SpAdapterGuessSchemaContext
implements IAdapterRuntimeContext, Serializable {
- private SpMonitoringManager monitoringManager;
+ private final IExtensionsLogger extensionsLogger;
- public SpAdapterRuntimeContext(SpMonitoringManager monitoringManager,
+ public SpAdapterRuntimeContext(IExtensionsLogger extensionsLogger,
ConfigExtractor configExtractor,
StreamPipesClient streamPipesClient) {
super(configExtractor, streamPipesClient);
- this.monitoringManager = monitoringManager;
+ this.extensionsLogger = extensionsLogger;
}
@Override
- public SpMonitoringManager getLogger() {
- return monitoringManager;
+ public IExtensionsLogger getLogger() {
+ return extensionsLogger;
}
}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/ExtensionsLogger.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/ExtensionsLogger.java
new file mode 100644
index 000000000..bfe9b56be
--- /dev/null
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/monitoring/ExtensionsLogger.java
@@ -0,0 +1,63 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You 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.apache.streampipes.extensions.management.monitoring;
+
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
+import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.model.monitoring.SpLogEntry;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
+
+public class ExtensionsLogger implements IExtensionsLogger {
+
+ private final String resourceId;
+ private final SpMonitoringManager monitoringManager;
+
+ public ExtensionsLogger(String resourceId) {
+ this.resourceId = resourceId;
+ this.monitoringManager = SpMonitoringManager.INSTANCE;
+ }
+
+ @Override
+ public void log(SpLogMessage logMessage) {
+ monitoringManager.addErrorMessage(
+ resourceId,
+ SpLogEntry.from(System.currentTimeMillis(), logMessage));
+ }
+
+ @Override
+ public void error(Exception e) {
+ log(SpLogMessage.from(e));
+ }
+
+ @Override
+ public void error(String details, Exception e) {
+ log(SpLogMessage.from(e, details));
+ }
+
+ @Override
+ public void info(String title, String details) {
+ log(SpLogMessage.info(title, details));
+ }
+
+ @Override
+ public void warn(String title, String details) {
+ log(SpLogMessage.warn(title, details));
+ }
+
+}
diff --git
a/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
b/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
index de8645795..4df695db5 100644
---
a/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
+++
b/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
@@ -42,7 +42,7 @@ public class AdapterWorkerManagementTest {
when(declarerSingleton.getAdapter(any())).thenReturn(Optional.empty());
var adapterWorkerManagement = new AdapterWorkerManagement(
- null, declarerSingleton, null);
+ null, declarerSingleton);
adapterWorkerManagement.invokeAdapter(adapterDescription);
}
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/iolink/IfmAlMqttAdapter.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/iolink/IfmAlMqttAdapter.java
index c5db51d94..e2606996d 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/iolink/IfmAlMqttAdapter.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/iolink/IfmAlMqttAdapter.java
@@ -32,9 +32,7 @@ import
org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor;
import
org.apache.streampipes.extensions.management.connect.adapter.parser.JsonParsers;
import
org.apache.streampipes.extensions.management.connect.adapter.parser.json.JsonObjectParser;
import org.apache.streampipes.model.AdapterType;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.connect.guess.GuessSchema;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.pe.shared.config.mqtt.MqttConfig;
import org.apache.streampipes.pe.shared.config.mqtt.MqttConnectUtils;
import org.apache.streampipes.pe.shared.config.mqtt.MqttConsumer;
@@ -127,9 +125,7 @@ public class IfmAlMqttAdapter implements StreamPipesAdapter
{
} catch (Exception e) {
adapterRuntimeContext
.getLogger()
- .addErrorMessage(
- extractor.getAdapterDescription().getElementId(),
- SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e)));
+ .error(e);
LOG.error("Could not parse event", e);
}
});
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/simulator/machine/MachineDataSimulator.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/simulator/machine/MachineDataSimulator.java
index 020dfb3db..32b941a21 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/simulator/machine/MachineDataSimulator.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/simulator/machine/MachineDataSimulator.java
@@ -19,7 +19,6 @@ package
org.apache.streampipes.connect.iiot.adapters.simulator.machine;
import org.apache.streampipes.commons.exceptions.connect.AdapterException;
import org.apache.streampipes.extensions.api.connect.IEventCollector;
-import
org.apache.streampipes.extensions.management.connect.adapter.model.pipeline.AdapterPipeline;
import java.util.HashMap;
import java.util.Map;
@@ -33,20 +32,15 @@ public class MachineDataSimulator implements Runnable {
private Boolean running;
- public MachineDataSimulator(IEventCollector collector, Integer waitTimeMs,
String selectedSimulatorOption) {
+ public MachineDataSimulator(IEventCollector collector,
+ Integer waitTimeMs,
+ String selectedSimulatorOption) {
this.collector = collector;
this.waitTimeMs = waitTimeMs;
this.selectedSimulatorOption = selectedSimulatorOption;
this.running = true;
}
- @Deprecated
- public MachineDataSimulator(AdapterPipeline adapterPipeline, Integer
waitTimeMs, String selectedSimulatorOption) {
- this.waitTimeMs = waitTimeMs;
- this.selectedSimulatorOption = selectedSimulatorOption;
- this.running = true;
- }
-
@Override
public void run() {
this.running = true;
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
index c58254cbe..f9ad62a2c 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
@@ -31,9 +31,7 @@ import
org.apache.streampipes.extensions.management.connect.adapter.parser.Image
import
org.apache.streampipes.extensions.management.connect.adapter.parser.JsonParsers;
import
org.apache.streampipes.extensions.management.connect.adapter.parser.xml.XmlParser;
import org.apache.streampipes.model.AdapterType;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.connect.guess.GuessSchema;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.sdk.StaticProperties;
import org.apache.streampipes.sdk.builder.adapter.AdapterConfigurationBuilder;
import org.apache.streampipes.sdk.helpers.Alternatives;
@@ -199,8 +197,7 @@ public class FileReplayAdapter implements
StreamPipesAdapter {
} catch (AdapterException e) {
adapterRuntimeContext
.getLogger()
- .addErrorMessage(extractor.getAdapterDescription().getElementId(),
- SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e)));
+ .error(e);
}
}
diff --git
a/streampipes-extensions/streampipes-processors-geo-jvm/src/main/java/org/apache/streampipes/processors/geo/jvm/jts/processor/validation/complex/TopologyValidationProcessor.java
b/streampipes-extensions/streampipes-processors-geo-jvm/src/main/java/org/apache/streampipes/processors/geo/jvm/jts/processor/validation/complex/TopologyValidationProcessor.java
index f0910c893..03a2b3382 100644
---
a/streampipes-extensions/streampipes-processors-geo-jvm/src/main/java/org/apache/streampipes/processors/geo/jvm/jts/processor/validation/complex/TopologyValidationProcessor.java
+++
b/streampipes-extensions/streampipes-processors-geo-jvm/src/main/java/org/apache/streampipes/processors/geo/jvm/jts/processor/validation/complex/TopologyValidationProcessor.java
@@ -23,9 +23,9 @@ import
org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import
org.apache.streampipes.extensions.api.pe.context.EventProcessorRuntimeContext;
import org.apache.streampipes.extensions.api.pe.routing.SpOutputCollector;
import org.apache.streampipes.model.DataProcessorType;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.graph.DataProcessorDescription;
import org.apache.streampipes.model.monitoring.SpLogEntry;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.model.runtime.Event;
import org.apache.streampipes.model.schema.PropertyScope;
import
org.apache.streampipes.processors.geo.jvm.jts.exceptions.SpJtsGeoemtryException;
@@ -120,7 +120,7 @@ public class TopologyValidationProcessor extends
StreamPipesDataProcessor {
if (isLogOutput) {
SpMonitoringManager.INSTANCE.addErrorMessage(params.getGraph().getElementId(),
SpLogEntry.from(System.currentTimeMillis(),
- StreamPipesErrorMessage.from(new SpJtsGeoemtryException(
+ SpLogMessage.from(new SpJtsGeoemtryException(
validator.getValidationError().toString()))));
}
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogEntry.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogEntry.java
index 75fb385f6..959289ef3 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogEntry.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogEntry.java
@@ -18,27 +18,31 @@
package org.apache.streampipes.model.monitoring;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.shared.annotation.TsModel;
@TsModel
public class SpLogEntry {
private long timestamp;
- private StreamPipesErrorMessage errorMessage;
+ private SpLogMessage errorMessage;
public SpLogEntry() {
}
+ public SpLogEntry(SpLogEntry other) {
+ this.timestamp = other.getTimestamp();
+ this.errorMessage = new SpLogMessage(other.getErrorMessage());
+ }
+
private SpLogEntry(long timestamp,
- StreamPipesErrorMessage errorMessage) {
+ SpLogMessage errorMessage) {
this.timestamp = timestamp;
this.errorMessage = errorMessage;
}
public static SpLogEntry from(long timestamp,
- StreamPipesErrorMessage errorMessage) {
+ SpLogMessage errorMessage) {
return new SpLogEntry(timestamp, errorMessage);
}
@@ -50,11 +54,11 @@ public class SpLogEntry {
this.timestamp = timestamp;
}
- public StreamPipesErrorMessage getErrorMessage() {
+ public SpLogMessage getErrorMessage() {
return errorMessage;
}
- public void setErrorMessage(StreamPipesErrorMessage errorMessage) {
+ public void setErrorMessage(SpLogMessage errorMessage) {
this.errorMessage = errorMessage;
}
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/StreamPipesStatistics.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogLevel.java
similarity index 94%
rename from
streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/StreamPipesStatistics.java
rename to
streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogLevel.java
index f67f91d6e..ece9018a3 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/StreamPipesStatistics.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogLevel.java
@@ -18,5 +18,9 @@
package org.apache.streampipes.model.monitoring;
-public class StreamPipesStatistics {
+public enum SpLogLevel {
+
+ INFO,
+ WARN,
+ ERROR
}
diff --git
a/streampipes-model/src/main/java/org/apache/streampipes/model/StreamPipesErrorMessage.java
b/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogMessage.java
similarity index 59%
rename from
streampipes-model/src/main/java/org/apache/streampipes/model/StreamPipesErrorMessage.java
rename to
streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogMessage.java
index a52ee29e3..e1123a26c 100644
---
a/streampipes-model/src/main/java/org/apache/streampipes/model/StreamPipesErrorMessage.java
+++
b/streampipes-model/src/main/java/org/apache/streampipes/model/monitoring/SpLogMessage.java
@@ -16,39 +16,80 @@
*
*/
-package org.apache.streampipes.model;
+package org.apache.streampipes.model.monitoring;
import org.apache.streampipes.model.shared.annotation.TsModel;
import org.apache.commons.lang3.exception.ExceptionUtils;
@TsModel
-public class StreamPipesErrorMessage {
+public class SpLogMessage {
- private String level;
+ private SpLogLevel level;
private String title;
private String detail;
private String cause;
private String fullStackTrace;
- public StreamPipesErrorMessage() {
+ public static SpLogMessage from(Exception exception) {
+ return from(exception, "");
+ }
+
+ public static SpLogMessage from(Exception exception,
+ String detail) {
+ String cause = exception.getCause() != null ?
exception.getCause().getMessage() : exception.getMessage();
+ return new SpLogMessage(
+ SpLogLevel.ERROR,
+ exception.getMessage(),
+ detail,
+ ExceptionUtils.getStackTrace(exception),
+ cause);
+ }
+
+ public static SpLogMessage info(String title,
+ String details) {
+ return new SpLogMessage(
+ SpLogLevel.INFO,
+ title,
+ details
+ );
+ }
+
+ public static SpLogMessage warn(String title,
+ String details) {
+ return new SpLogMessage(
+ SpLogLevel.WARN,
+ title,
+ details
+ );
+ }
+
+ public SpLogMessage() {
}
- public StreamPipesErrorMessage(String level,
- String title,
- String detail) {
+ public SpLogMessage(SpLogMessage other) {
+ this.level = other.getLevel();
+ this.detail = other.getDetail();
+ this.title = other.getTitle();
+ this.cause = other.getCause();
+ this.fullStackTrace = other.getFullStackTrace();
+ }
+
+ public SpLogMessage(SpLogLevel level,
+ String title,
+ String detail) {
this.level = level;
this.title = title;
this.detail = detail;
}
- public StreamPipesErrorMessage(String level,
- String title,
- String detail,
- String fullStackTrace,
- String cause) {
+ public SpLogMessage(SpLogLevel level,
+ String title,
+ String detail,
+ String fullStackTrace,
+ String cause) {
this.level = level;
this.title = title;
this.detail = detail;
@@ -56,21 +97,11 @@ public class StreamPipesErrorMessage {
this.cause = cause;
}
- public static StreamPipesErrorMessage from(Exception exception) {
- String cause = exception.getCause() != null ?
exception.getCause().getMessage() : exception.getMessage();
- return new StreamPipesErrorMessage(
- "error",
- exception.getMessage(),
- "",
- ExceptionUtils.getStackTrace(exception),
- cause);
- }
-
- public String getLevel() {
+ public SpLogLevel getLevel() {
return level;
}
- public void setLevel(String level) {
+ public void setLevel(SpLogLevel level) {
this.level = level;
}
diff --git
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsLogProvider.java
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsLogProvider.java
index 013b55a6e..4517e30b1 100644
---
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsLogProvider.java
+++
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/monitoring/pipeline/ExtensionsLogProvider.java
@@ -53,7 +53,11 @@ public enum ExtensionsLogProvider {
infos.addAll(0, value);
if (infos.size() > MAX_ITEMS) {
- infos.subList(MAX_ITEMS, infos.size()).clear();
+ int numElementsToRemove = infos.size() - MAX_ITEMS;
+
+ for (int i = 0; i < numElementsToRemove; i++) {
+ infos.remove(infos.size() - 1);
+ }
}
});
}
diff --git
a/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
b/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
index b449e2dcf..a0d798c37 100644
---
a/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
+++
b/streampipes-platform-services/src/main/java/org/apache/streampipes/ps/DataLakeResourceV4.java
@@ -22,10 +22,10 @@ import
org.apache.streampipes.dataexplorer.DataExplorerQueryManagement;
import org.apache.streampipes.dataexplorer.DataExplorerSchemaManagement;
import org.apache.streampipes.dataexplorer.param.ProvidedRestQueryParams;
import org.apache.streampipes.dataexplorer.query.writer.OutputFormat;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.datalake.DataLakeMeasure;
import org.apache.streampipes.model.datalake.DataSeries;
import org.apache.streampipes.model.datalake.SpQueryResult;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.rest.core.base.impl.AbstractRestResource;
import io.swagger.v3.oas.annotations.Operation;
@@ -241,7 +241,7 @@ public class DataLakeResourceV4 extends
AbstractRestResource {
this.dataLakeManagement.getData(sanitizedParams,
isIgnoreMissingValues(missingValueBehaviour));
return ok(result);
} catch (RuntimeException e) {
- return badRequest(StreamPipesErrorMessage.from(e));
+ return badRequest(SpLogMessage.from(e));
}
}
}
diff --git
a/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/connect/AdapterWorkerResource.java
b/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/connect/AdapterWorkerResource.java
index d25ea6384..2877fc456 100644
---
a/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/connect/AdapterWorkerResource.java
+++
b/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/connect/AdapterWorkerResource.java
@@ -20,12 +20,11 @@ package org.apache.streampipes.rest.extensions.connect;
import org.apache.streampipes.commons.exceptions.connect.AdapterException;
import
org.apache.streampipes.extensions.management.connect.AdapterWorkerManagement;
-import
org.apache.streampipes.extensions.management.context.AdapterContextGenerator;
import org.apache.streampipes.extensions.management.init.DeclarersSingleton;
import
org.apache.streampipes.extensions.management.init.RunningAdapterInstances;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.message.Notifications;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.rest.shared.annotation.JacksonSerialized;
import org.apache.streampipes.rest.shared.impl.AbstractSharedRestInterface;
@@ -51,8 +50,7 @@ public class AdapterWorkerResource extends
AbstractSharedRestInterface {
public AdapterWorkerResource() {
adapterManagement = new AdapterWorkerManagement(
RunningAdapterInstances.INSTANCE,
- DeclarersSingleton.getInstance(),
- new AdapterContextGenerator().makeRuntimeContext()
+ DeclarersSingleton.getInstance()
);
}
@@ -83,7 +81,7 @@ public class AdapterWorkerResource extends
AbstractSharedRestInterface {
return ok(Notifications.success(responseMessage));
} catch (AdapterException e) {
logger.error("Error while starting adapter with id " +
adapterStreamDescription.getUri(), e);
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
}
}
@@ -107,7 +105,7 @@ public class AdapterWorkerResource extends
AbstractSharedRestInterface {
return ok(Notifications.success(responseMessage));
} catch (AdapterException e) {
logger.error("Error while stopping adapter with id " +
adapterStreamDescription.getElementId(), e);
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/monitoring/MonitoringResource.java
b/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/monitoring/MonitoringResource.java
index 0c2ee0aa4..8cfbea34d 100644
---
a/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/monitoring/MonitoringResource.java
+++
b/streampipes-rest-extensions/src/main/java/org/apache/streampipes/rest/extensions/monitoring/MonitoringResource.java
@@ -37,7 +37,7 @@ public class MonitoringResource extends
AbstractExtensionsResource {
try {
return ok(SpMonitoringManager.INSTANCE.getMonitoringInfo());
} finally {
- //SpLogManager.INSTANCE.clearAllLogs();
+ SpMonitoringManager.INSTANCE.clearAllLogs();
}
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
index 775e5c885..7de6ed08c 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/AdapterResource.java
@@ -20,9 +20,9 @@ package org.apache.streampipes.rest.impl.connect;
import org.apache.streampipes.commons.exceptions.connect.AdapterException;
import
org.apache.streampipes.connect.management.management.AdapterMasterManagement;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.message.Notifications;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.rest.security.AuthConstants;
import org.apache.streampipes.rest.shared.annotation.JacksonSerialized;
import org.apache.streampipes.storage.management.StorageDispatcher;
@@ -121,7 +121,7 @@ public class AdapterResource extends
AbstractAdapterResource<AdapterMasterManage
return ok(Notifications.success("Adapter started"));
} catch (AdapterException e) {
LOG.error("Could not stop adapter with id " + adapterId, e);
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
}
}
@@ -136,7 +136,7 @@ public class AdapterResource extends
AbstractAdapterResource<AdapterMasterManage
return ok(Notifications.success("Adapter stopped"));
} catch (AdapterException e) {
LOG.error("Could not start adapter with id " + adapterId, e);
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/GuessResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/GuessResource.java
index ea00214c1..50623e231 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/GuessResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/GuessResource.java
@@ -22,10 +22,10 @@ import
org.apache.streampipes.commons.exceptions.NoServiceEndpointsAvailableExce
import org.apache.streampipes.commons.exceptions.connect.ParseException;
import org.apache.streampipes.connect.management.management.GuessManagement;
import
org.apache.streampipes.extensions.api.connect.exception.WorkerAdapterException;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
import org.apache.streampipes.model.connect.guess.AdapterEventPreview;
import org.apache.streampipes.model.connect.guess.GuessSchema;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.rest.shared.annotation.JacksonSerialized;
import com.fasterxml.jackson.core.JsonProcessingException;
@@ -63,12 +63,12 @@ public class GuessResource extends
AbstractAdapterResource<GuessManagement> {
return ok(result);
} catch (ParseException e) {
LOG.error("Error while parsing events: ", e);
- return badRequest(StreamPipesErrorMessage.from(e));
+ return badRequest(SpLogMessage.from(e));
} catch (WorkerAdapterException e) {
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
} catch (NoServiceEndpointsAvailableException | IOException e) {
LOG.error(e.getMessage());
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/RuntimeResolvableResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/RuntimeResolvableResource.java
index e22a4d4fa..c8e5386b3 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/RuntimeResolvableResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/connect/RuntimeResolvableResource.java
@@ -24,7 +24,7 @@ import
org.apache.streampipes.commons.exceptions.connect.AdapterException;
import
org.apache.streampipes.connect.management.management.WorkerAdministrationManagement;
import org.apache.streampipes.connect.management.management.WorkerRestClient;
import org.apache.streampipes.connect.management.management.WorkerUrlProvider;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.model.runtime.RuntimeOptionsRequest;
import org.apache.streampipes.model.runtime.RuntimeOptionsResponse;
import org.apache.streampipes.rest.shared.annotation.JacksonSerialized;
@@ -68,13 +68,13 @@ public class RuntimeResolvableResource extends
AbstractAdapterResource<WorkerAdm
return ok(result);
} catch (AdapterException e) {
LOG.error("Adapter exception occurred", e);
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
} catch (NoServiceEndpointsAvailableException e) {
LOG.error("Could not find service endpoint for {} while fetching
configuration", appId);
- return serverError(StreamPipesErrorMessage.from(e));
+ return serverError(SpLogMessage.from(e));
} catch (SpConfigurationException e) {
LOG.error("Tried to fetch a runtime configuration with insufficient
settings");
- return badRequest(StreamPipesErrorMessage.from(e));
+ return badRequest(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataProcessorResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataProcessorResource.java
index edbdfa642..6f52faead 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataProcessorResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataProcessorResource.java
@@ -18,10 +18,10 @@
package org.apache.streampipes.rest.impl.pe;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.graph.DataProcessorDescription;
import org.apache.streampipes.model.graph.DataProcessorInvocation;
import org.apache.streampipes.model.message.NotificationType;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.resource.management.DataProcessorResourceManager;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import org.apache.streampipes.rest.security.AuthConstants;
@@ -83,7 +83,7 @@ public class DataProcessorResource extends
AbstractAuthGuardedRestResource {
try {
return ok(getDataProcessorResourceManager().findAsInvocation(elementId));
} catch (IllegalArgumentException e) {
- return notFound(StreamPipesErrorMessage.from(e));
+ return notFound(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataSinkResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataSinkResource.java
index 55c62f5ab..884b6b7aa 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataSinkResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataSinkResource.java
@@ -18,10 +18,10 @@
package org.apache.streampipes.rest.impl.pe;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.graph.DataSinkDescription;
import org.apache.streampipes.model.graph.DataSinkInvocation;
import org.apache.streampipes.model.message.NotificationType;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.resource.management.DataSinkResourceManager;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import org.apache.streampipes.rest.security.AuthConstants;
@@ -84,7 +84,7 @@ public class DataSinkResource extends
AbstractAuthGuardedRestResource {
try {
return ok(getDataSinkResourceManager().findAsInvocation(elementId));
} catch (IllegalArgumentException e) {
- return notFound(StreamPipesErrorMessage.from(e));
+ return notFound(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
index 9727f7ba9..5c582dace 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/pe/DataStreamResource.java
@@ -19,8 +19,8 @@
package org.apache.streampipes.rest.impl.pe;
import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.message.NotificationType;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.resource.management.DataStreamResourceManager;
import
org.apache.streampipes.rest.core.base.impl.AbstractAuthGuardedRestResource;
import org.apache.streampipes.rest.security.AuthConstants;
@@ -85,7 +85,7 @@ public class DataStreamResource extends
AbstractAuthGuardedRestResource {
try {
return ok(getDataStreamResourceManager().findAsInvocation(elementId));
} catch (IllegalArgumentException e) {
- return notFound(StreamPipesErrorMessage.from(e));
+ return notFound(SpLogMessage.from(e));
}
}
diff --git
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/FunctionContext.java
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/FunctionContext.java
index ed1f714fd..ff4326971 100644
---
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/FunctionContext.java
+++
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/FunctionContext.java
@@ -19,11 +19,11 @@
package org.apache.streampipes.wrapper.standalone.function;
import org.apache.streampipes.client.StreamPipesClient;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
import org.apache.streampipes.extensions.api.pe.routing.SpOutputCollector;
import org.apache.streampipes.extensions.management.config.ConfigExtractor;
+import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.model.schema.EventSchema;
import java.util.Collection;
@@ -40,6 +40,7 @@ public class FunctionContext {
private ConfigExtractor config;
private Map<String, SpOutputCollector> outputCollectors;
+ private IExtensionsLogger extensionsLogger;
public FunctionContext() {
this.streams = new HashMap<>();
@@ -56,6 +57,7 @@ public class FunctionContext {
this.functionId = functionId;
this.outputCollectors = outputCollectors;
this.client = client;
+ this.extensionsLogger = new ExtensionsLogger(functionId);
}
public Collection<SpDataStream> getStreams() {
@@ -78,8 +80,8 @@ public class FunctionContext {
return functionId;
}
- public void log(SpLogEntry logEntry) {
- SpMonitoringManager.INSTANCE.addErrorMessage(functionId, logEntry);
+ public IExtensionsLogger getLogger() {
+ return extensionsLogger;
}
public Map<String, SpOutputCollector> getOutputCollectors() {
diff --git
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/StreamPipesFunction.java
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/StreamPipesFunction.java
index f18ca9f9e..b4a09e20f 100644
---
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/StreamPipesFunction.java
+++
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/function/StreamPipesFunction.java
@@ -29,9 +29,9 @@ import
org.apache.streampipes.extensions.api.pe.routing.SpInputCollector;
import org.apache.streampipes.extensions.api.pe.routing.SpOutputCollector;
import org.apache.streampipes.extensions.management.util.GroundingDebugUtils;
import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.constants.PropertySelectorConstants;
import org.apache.streampipes.model.monitoring.SpLogEntry;
+import org.apache.streampipes.model.monitoring.SpLogMessage;
import org.apache.streampipes.model.runtime.Event;
import org.apache.streampipes.model.runtime.EventFactory;
import org.apache.streampipes.model.runtime.SchemaInfo;
@@ -136,7 +136,7 @@ public abstract class StreamPipesFunction implements
IStreamPipesFunctionDeclare
var functionId = this.getFunctionConfig().getFunctionId();
SpMonitoringManager.INSTANCE.addErrorMessage(
functionId.getId(),
- SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e)));
+ SpLogEntry.from(System.currentTimeMillis(), SpLogMessage.from(e)));
}
private void initializeProducers() {
diff --git
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
index c76df7cc6..671de2782 100644
---
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
+++
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
@@ -21,12 +21,11 @@ package org.apache.streampipes.wrapper.standalone.routing;
import org.apache.streampipes.commons.exceptions.SpRuntimeException;
import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import org.apache.streampipes.extensions.api.pe.routing.SpOutputCollector;
+import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import org.apache.streampipes.messaging.EventProducer;
import org.apache.streampipes.messaging.InternalEventProcessor;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.grounding.TransportFormat;
import org.apache.streampipes.model.grounding.TransportProtocol;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.model.runtime.Event;
import org.apache.streampipes.model.runtime.EventConverter;
import org.apache.streampipes.wrapper.standalone.manager.ProtocolManager;
@@ -44,6 +43,7 @@ public class StandaloneSpOutputCollector<T extends
TransportProtocol> extends
private final EventProducer producer;
private final String resourceId;
+ private final ExtensionsLogger extensionsLogger;
public StandaloneSpOutputCollector(T protocol,
TransportFormat format,
@@ -51,6 +51,7 @@ public class StandaloneSpOutputCollector<T extends
TransportProtocol> extends
super(protocol, format);
this.producer = protocolDefinition.getProducer(protocol);
this.resourceId = resourceId;
+ this.extensionsLogger = new ExtensionsLogger(resourceId);
}
public void collect(Event event) {
@@ -59,8 +60,7 @@ public class StandaloneSpOutputCollector<T extends
TransportProtocol> extends
producer.publish(dataFormatDefinition.fromMap(outEvent));
SpMonitoringManager.INSTANCE.increaseOutCounter(resourceId,
System.currentTimeMillis());
} catch (SpRuntimeException e) {
- var logEntry = SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e));
- SpMonitoringManager.INSTANCE.addErrorMessage(resourceId, logEntry);
+ extensionsLogger.error(e);
LOG.error("Could not publish event", e);
}
}
diff --git
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventProcessorRuntime.java
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventProcessorRuntime.java
index c3e250fe7..c2b70c540 100644
---
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventProcessorRuntime.java
+++
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventProcessorRuntime.java
@@ -70,13 +70,13 @@ public class StandaloneEventProcessorRuntime extends
StandalonePipelineElementRu
@Override
public void process(Map<String, Object> rawEvent, String sourceInfo) {
try {
- runtimeContext.getLogger().increaseInCounter(instanceId, sourceInfo,
System.currentTimeMillis());
+ monitoringManager.increaseInCounter(instanceId, sourceInfo,
System.currentTimeMillis());
var event = this.internalRuntimeParameters.makeEvent(runtimeParameters,
rawEvent, sourceInfo);
pipelineElement
.onEvent(event, outputCollector);
} catch (RuntimeException e) {
LOG.error("RuntimeException while processing event in {}",
pipelineElement.getClass().getCanonicalName(), e);
- addLogEntry(runtimeContext.getLogger(), instanceId, e);
+ addLogEntry(e);
}
}
diff --git
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventSinkRuntime.java
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventSinkRuntime.java
index 27be8ee8a..7c7f407c1 100644
---
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventSinkRuntime.java
+++
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandaloneEventSinkRuntime.java
@@ -51,11 +51,11 @@ public class StandaloneEventSinkRuntime extends
StandalonePipelineElementRuntime
@Override
public void process(Map<String, Object> rawEvent, String sourceInfo) {
try {
- runtimeContext.getLogger().increaseInCounter(instanceId, sourceInfo,
System.currentTimeMillis());
+ monitoringManager.increaseInCounter(instanceId, sourceInfo,
System.currentTimeMillis());
pipelineElement.onEvent(internalRuntimeParameters.makeEvent(runtimeParameters,
rawEvent, sourceInfo));
} catch (RuntimeException e) {
LOG.error("RuntimeException while processing event in {}",
pipelineElement.getClass().getCanonicalName(), e);
- addLogEntry(runtimeContext.getLogger(), instanceId, e);
+ addLogEntry(e);
}
}
diff --git
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandalonePipelineElementRuntime.java
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandalonePipelineElementRuntime.java
index 8005ca4d5..d35317cda 100644
---
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandalonePipelineElementRuntime.java
+++
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/runtime/StandalonePipelineElementRuntime.java
@@ -31,9 +31,7 @@ import
org.apache.streampipes.extensions.api.pe.routing.PipelineElementCollector
import org.apache.streampipes.extensions.api.pe.routing.RawDataProcessor;
import org.apache.streampipes.extensions.api.pe.routing.SpInputCollector;
import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.StreamPipesErrorMessage;
import org.apache.streampipes.model.base.InvocableStreamPipesEntity;
-import org.apache.streampipes.model.monitoring.SpLogEntry;
import org.apache.streampipes.wrapper.params.InternalRuntimeParameters;
import org.apache.streampipes.wrapper.runtime.PipelineElementRuntime;
import org.apache.streampipes.wrapper.standalone.manager.ProtocolManager;
@@ -58,10 +56,13 @@ public abstract class StandalonePipelineElementRuntime<
protected PeT pipelineElement;
protected IInternalRuntimeParameters internalRuntimeParameters;
+ protected final SpMonitoringManager monitoringManager;
+
public StandalonePipelineElementRuntime(IContextGenerator<RcT, IvT>
contextGenerator,
IParameterGenerator<IvT, ExT, PepT>
parameterGenerator) {
super(contextGenerator, parameterGenerator);
this.internalRuntimeParameters = new InternalRuntimeParameters();
+ this.monitoringManager = SpMonitoringManager.INSTANCE;
}
@Override
@@ -80,13 +81,12 @@ public abstract class StandalonePipelineElementRuntime<
@Override
public void stopRuntime() {
this.inputCollectors.forEach(is -> is.unregisterConsumer(instanceId));
- resetCounter(runtimeContext.getLogger(), instanceId);
+ resetCounter(instanceId);
afterStop();
}
- protected void resetCounter(SpMonitoringManager manager,
- String resourceId) throws SpRuntimeException {
- manager.resetCounter(resourceId);
+ protected void resetCounter(String resourceId) throws SpRuntimeException {
+ monitoringManager.resetCounter(resourceId);
}
protected List<SpInputCollector> getInputCollectors(List<SpDataStream>
inputStreams) throws SpRuntimeException {
@@ -99,12 +99,8 @@ public abstract class StandalonePipelineElementRuntime<
return inputCollectors;
}
- protected void addLogEntry(SpMonitoringManager manager,
- String elementId,
- RuntimeException e) {
- manager.addErrorMessage(
- elementId,
- SpLogEntry.from(System.currentTimeMillis(),
StreamPipesErrorMessage.from(e)));
+ protected void addLogEntry(RuntimeException e) {
+ runtimeContext.getLogger().error(e);
}
protected void connectInputCollectors() {
diff --git
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventProcessorRuntimeContext.java
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventProcessorRuntimeContext.java
index 6ad5ba0db..05fca370a 100644
---
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventProcessorRuntimeContext.java
+++
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventProcessorRuntimeContext.java
@@ -17,10 +17,10 @@
*/
package org.apache.streampipes.wrapper.context;
-import org.apache.streampipes.client.StreamPipesClient;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.extensions.api.config.IConfigExtractor;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
import
org.apache.streampipes.extensions.api.pe.context.EventProcessorRuntimeContext;
-import org.apache.streampipes.extensions.management.config.ConfigExtractor;
import java.io.Serializable;
@@ -29,14 +29,9 @@ public class SpEventProcessorRuntimeContext extends
SpRuntimeContext implements
public SpEventProcessorRuntimeContext(String correspondingUser,
- ConfigExtractor configExtractor,
- StreamPipesClient streamPipesClient,
- SpMonitoringManager logManager) {
- super(correspondingUser, configExtractor, streamPipesClient, logManager);
+ IConfigExtractor configExtractor,
+ IStreamPipesClient streamPipesClient,
+ IExtensionsLogger extensionsLogger) {
+ super(correspondingUser, configExtractor, streamPipesClient,
extensionsLogger);
}
-
- public SpEventProcessorRuntimeContext() {
- super();
- }
-
}
diff --git
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventSinkRuntimeContext.java
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventSinkRuntimeContext.java
index 6ef795a7e..6bfc03d03 100644
---
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventSinkRuntimeContext.java
+++
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpEventSinkRuntimeContext.java
@@ -17,10 +17,10 @@
*/
package org.apache.streampipes.wrapper.context;
-import org.apache.streampipes.client.StreamPipesClient;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.extensions.api.config.IConfigExtractor;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
import
org.apache.streampipes.extensions.api.pe.context.EventSinkRuntimeContext;
-import org.apache.streampipes.extensions.management.config.ConfigExtractor;
import java.io.Serializable;
@@ -28,9 +28,9 @@ public class SpEventSinkRuntimeContext extends
SpRuntimeContext implements
EventSinkRuntimeContext, Serializable {
public SpEventSinkRuntimeContext(String correspondingUser,
- ConfigExtractor configExtractor,
- StreamPipesClient streamPipesClient,
- SpMonitoringManager logManager) {
- super(correspondingUser, configExtractor, streamPipesClient, logManager);
+ IConfigExtractor configExtractor,
+ IStreamPipesClient streamPipesClient,
+ IExtensionsLogger extensionsLogger) {
+ super(correspondingUser, configExtractor, streamPipesClient,
extensionsLogger);
}
}
diff --git
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpRuntimeContext.java
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpRuntimeContext.java
index 081f7034e..aa4a73e76 100644
---
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpRuntimeContext.java
+++
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/SpRuntimeContext.java
@@ -17,36 +17,32 @@
*/
package org.apache.streampipes.wrapper.context;
-import org.apache.streampipes.client.StreamPipesClient;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.extensions.api.config.IConfigExtractor;
+import org.apache.streampipes.extensions.api.monitoring.IExtensionsLogger;
import org.apache.streampipes.extensions.api.pe.context.RuntimeContext;
-import org.apache.streampipes.extensions.management.config.ConfigExtractor;
public class SpRuntimeContext implements RuntimeContext {
- private String correspondingUser;
- private ConfigExtractor configExtractor;
- private StreamPipesClient streamPipesClient;
- private SpMonitoringManager spLogManager;
+ private final String correspondingUser;
+ private final IConfigExtractor configExtractor;
+ private final IStreamPipesClient streamPipesClient;
+ private final IExtensionsLogger extensionsLogger;
public SpRuntimeContext(String correspondingUser,
- ConfigExtractor configExtractor,
- StreamPipesClient streamPipesClient,
- SpMonitoringManager spLogManager) {
+ IConfigExtractor configExtractor,
+ IStreamPipesClient streamPipesClient,
+ IExtensionsLogger extensionsLogger) {
this.correspondingUser = correspondingUser;
this.configExtractor = configExtractor;
this.streamPipesClient = streamPipesClient;
- this.spLogManager = spLogManager;
- }
-
- public SpRuntimeContext() {
-
+ this.extensionsLogger = extensionsLogger;
}
@Override
- public SpMonitoringManager getLogger() {
- return spLogManager;
+ public IExtensionsLogger getLogger() {
+ return extensionsLogger;
}
@Override
@@ -55,12 +51,12 @@ public class SpRuntimeContext implements RuntimeContext {
}
@Override
- public ConfigExtractor getConfigStore() {
+ public IConfigExtractor getConfigStore() {
return this.configExtractor;
}
@Override
- public StreamPipesClient getStreamPipesClient() {
+ public IStreamPipesClient getStreamPipesClient() {
return streamPipesClient;
}
}
diff --git
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataProcessorContextGenerator.java
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataProcessorContextGenerator.java
index 60e04e549..2783e3e27 100644
---
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataProcessorContextGenerator.java
+++
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataProcessorContextGenerator.java
@@ -18,9 +18,9 @@
package org.apache.streampipes.wrapper.context.generator;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import
org.apache.streampipes.extensions.api.pe.context.EventProcessorRuntimeContext;
import org.apache.streampipes.extensions.api.pe.context.IContextGenerator;
+import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import org.apache.streampipes.extensions.management.util.RuntimeContextUtils;
import org.apache.streampipes.model.graph.DataProcessorInvocation;
import org.apache.streampipes.wrapper.context.SpEventProcessorRuntimeContext;
@@ -34,6 +34,6 @@ public class DataProcessorContextGenerator
invocation.getCorrespondingUser(),
RuntimeContextUtils.makeConfigExtractor(),
RuntimeContextUtils.makeStreamPipesClient(),
- SpMonitoringManager.INSTANCE);
+ new ExtensionsLogger(invocation.getElementId()));
}
}
diff --git
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataSinkContextGenerator.java
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataSinkContextGenerator.java
index 28166cd61..c2413a930 100644
---
a/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataSinkContextGenerator.java
+++
b/streampipes-wrapper/src/main/java/org/apache/streampipes/wrapper/context/generator/DataSinkContextGenerator.java
@@ -18,9 +18,9 @@
package org.apache.streampipes.wrapper.context.generator;
-import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import
org.apache.streampipes.extensions.api.pe.context.EventSinkRuntimeContext;
import org.apache.streampipes.extensions.api.pe.context.IContextGenerator;
+import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import org.apache.streampipes.extensions.management.util.RuntimeContextUtils;
import org.apache.streampipes.model.graph.DataSinkInvocation;
import org.apache.streampipes.wrapper.context.SpEventSinkRuntimeContext;
@@ -33,6 +33,6 @@ public class DataSinkContextGenerator implements
IContextGenerator<EventSinkRunt
invocation.getCorrespondingUser(),
RuntimeContextUtils.makeConfigExtractor(),
RuntimeContextUtils.makeStreamPipesClient(),
- SpMonitoringManager.INSTANCE);
+ new ExtensionsLogger(invocation.getElementId()));
}
}
diff --git
a/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
b/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
index 2698ec6f5..ca4deba8b 100644
---
a/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
+++
b/ui/projects/streampipes/platform-services/src/lib/model/gen/streampipes-model.ts
@@ -20,7 +20,7 @@
/* tslint:disable */
/* eslint-disable */
// @ts-nocheck
-// Generated using typescript-generator version 3.1.1185 on 2023-07-19
22:10:18.
+// Generated using typescript-generator version 3.1.1185 on 2023-07-23
10:14:46.
export class NamedStreamPipesEntity {
'@class':
@@ -3325,7 +3325,7 @@ export class SpDataStreamContainer {
}
export class SpLogEntry {
- errorMessage: StreamPipesErrorMessage;
+ errorMessage: SpLogMessage;
timestamp: number;
static fromData(data: SpLogEntry, target?: SpLogEntry): SpLogEntry {
@@ -3333,14 +3333,33 @@ export class SpLogEntry {
return data;
}
const instance = target || new SpLogEntry();
- instance.errorMessage = StreamPipesErrorMessage.fromData(
- data.errorMessage,
- );
+ instance.errorMessage = SpLogMessage.fromData(data.errorMessage);
instance.timestamp = data.timestamp;
return instance;
}
}
+export class SpLogMessage {
+ cause: string;
+ detail: string;
+ fullStackTrace: string;
+ level: SpLogLevel;
+ title: string;
+
+ static fromData(data: SpLogMessage, target?: SpLogMessage): SpLogMessage {
+ if (!data) {
+ return data;
+ }
+ const instance = target || new SpLogMessage();
+ instance.cause = data.cause;
+ instance.detail = data.detail;
+ instance.fullStackTrace = data.fullStackTrace;
+ instance.level = data.level;
+ instance.title = data.title;
+ return instance;
+ }
+}
+
export class SpMetricsEntry {
lastTimestamp: number;
messagesIn: { [index: string]: MessageCounter };
@@ -3521,30 +3540,6 @@ export class StreamPipesApplicationPackage {
}
}
-export class StreamPipesErrorMessage {
- cause: string;
- detail: string;
- fullStackTrace: string;
- level: string;
- title: string;
-
- static fromData(
- data: StreamPipesErrorMessage,
- target?: StreamPipesErrorMessage,
- ): StreamPipesErrorMessage {
- if (!data) {
- return data;
- }
- const instance = target || new StreamPipesErrorMessage();
- instance.cause = data.cause;
- instance.detail = data.detail;
- instance.fullStackTrace = data.fullStackTrace;
- instance.level = data.level;
- instance.title = data.title;
- return instance;
- }
-}
-
export class SuccessMessage extends Message {
static fromData(
data: SuccessMessage,
@@ -3852,6 +3847,8 @@ export type SelectionStaticPropertyUnion =
| AnyStaticProperty
| OneOfStaticProperty;
+export type SpLogLevel = 'INFO' | 'WARN' | 'ERROR';
+
export type SpQueryStatus = 'OK' | 'TOO_MUCH_DATA';
export type StaticPropertyType =
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/basic-nav-tabs/basic-nav-tabs.component.html
b/ui/projects/streampipes/shared-ui/src/lib/components/basic-nav-tabs/basic-nav-tabs.component.html
index 4759f541b..abe940f2a 100644
---
a/ui/projects/streampipes/shared-ui/src/lib/components/basic-nav-tabs/basic-nav-tabs.component.html
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/basic-nav-tabs/basic-nav-tabs.component.html
@@ -36,12 +36,7 @@
<mat-icon>arrow_back</mat-icon>
</button>
</div>
- <nav
- mat-tab-nav-bar
- mat-stretch-tabs="false"
- color="accent"
- class="w-100"
- >
+ <nav mat-tab-nav-bar mat-stretch-tabs="false" color="accent">
<a
mat-tab-link
*ngFor="let item of spNavigationItems"
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.html
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.html
index ed56aaf95..eff21738d 100644
---
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.html
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.html
@@ -31,9 +31,13 @@
<span class="color-warn">{{ title }}</span>
</div>
<div class="error-details-title">Probable cause</div>
- <div class="log-message" [innerText]="message.cause"></div>
+ <div
+ class="log-message"
+ [innerText]="message.cause || 'No more information'"
+ ></div>
<div class="mt-10">
<button
+ [disabled]="message.fullStackTrace === undefined"
mat-button
color="accent"
(click)="showDetails = !showDetails"
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.ts
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.ts
index dd5e42687..662c94d72 100644
---
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.ts
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/exception-details-dialog/exception-details-dialog.component.ts
@@ -17,7 +17,7 @@
*/
import { Component, Input, OnInit } from '@angular/core';
-import { StreamPipesErrorMessage } from '@streampipes/platform-services';
+import { SpLogMessage } from '@streampipes/platform-services';
import { DialogRef } from '../../../dialog/base-dialog/dialog-ref';
@Component({
@@ -30,7 +30,7 @@ import { DialogRef } from
'../../../dialog/base-dialog/dialog-ref';
})
export class SpExceptionDetailsDialogComponent implements OnInit {
@Input()
- message: StreamPipesErrorMessage;
+ message: SpLogMessage;
@Input()
title: string;
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.html
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.html
index 9a32302d6..8d59dcc17 100644
---
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.html
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.html
@@ -30,7 +30,9 @@
<h5 *ngIf="messageTimestamp" class="mr-15">
{{ messageTimestamp | date : 'short' }}
</h5>
- <h5 fxFlex class="color-warn">{{ message.title }}</h5>
+ <h5 fxFlex class="color-warn">
+ {{ message.title || 'Error' }}
+ </h5>
</div>
<span fxFlex></span>
<div fxLayoutAlign="end center" *ngIf="showDetails">
diff --git
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.ts
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.ts
index 431476999..59860d99a 100644
---
a/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.ts
+++
b/ui/projects/streampipes/shared-ui/src/lib/components/sp-exception-message/sp-exception-message.component.ts
@@ -17,7 +17,7 @@
*/
import { Component, Input } from '@angular/core';
-import { StreamPipesErrorMessage } from '@streampipes/platform-services';
+import { SpLogMessage } from '@streampipes/platform-services';
import { DialogService } from '../../dialog/base-dialog/base-dialog.service';
import { PanelType } from '../../dialog/base-dialog/base-dialog.model';
import { SpExceptionDetailsDialogComponent } from
'./exception-details-dialog/exception-details-dialog.component';
@@ -35,7 +35,7 @@ export class SpExceptionMessageComponent {
showDetails = true;
@Input()
- message: StreamPipesErrorMessage;
+ message: SpLogMessage;
@Input()
messageTimestamp: number;
diff --git
a/ui/src/app/connect/components/adapter-configuration/schema-editor/error-message/error-message.component.ts
b/ui/src/app/connect/components/adapter-configuration/schema-editor/error-message/error-message.component.ts
index 4ef638d05..5b1d87dc6 100644
---
a/ui/src/app/connect/components/adapter-configuration/schema-editor/error-message/error-message.component.ts
+++
b/ui/src/app/connect/components/adapter-configuration/schema-editor/error-message/error-message.component.ts
@@ -17,7 +17,7 @@
*/
import { Component, Input } from '@angular/core';
-import { StreamPipesErrorMessage } from
'../../../../../../../projects/streampipes/platform-services/src/lib/model/gen/streampipes-model';
+import { SpLogMessage } from '@streampipes/platform-services';
@Component({
selector: 'sp-error-message',
@@ -25,7 +25,7 @@ import { StreamPipesErrorMessage } from
'../../../../../../../projects/streampip
styleUrls: ['./error-message.component.scss'],
})
export class ErrorMessageComponent {
- @Input() errorMessage: StreamPipesErrorMessage;
+ @Input() errorMessage: SpLogMessage;
showErrorMessage = false;
diff --git
a/ui/src/app/connect/components/adapter-configuration/schema-editor/event-schema/event-schema.component.ts
b/ui/src/app/connect/components/adapter-configuration/schema-editor/event-schema/event-schema.component.ts
index 877395da6..ca0a4da55 100644
---
a/ui/src/app/connect/components/adapter-configuration/schema-editor/event-schema/event-schema.component.ts
+++
b/ui/src/app/connect/components/adapter-configuration/schema-editor/event-schema/event-schema.component.ts
@@ -37,7 +37,7 @@ import {
EventSchema,
FieldStatusInfo,
GuessSchema,
- StreamPipesErrorMessage,
+ SpLogMessage,
} from '@streampipes/platform-services';
import { MatStepper } from '@angular/material/stepper';
import { UserErrorMessage } from
'../../../../../core-model/base/UserErrorMessage';
@@ -102,7 +102,7 @@ export class EventSchemaComponent implements OnChanges {
isLoading = false;
isError = false;
isPreviewEnabled = false;
- errorMessage: StreamPipesErrorMessage;
+ errorMessage: SpLogMessage;
nodes: EventProperty[] = new Array<EventProperty>();
validEventSchema = false;
schemaErrorHints: UserErrorMessage[] = [];
diff --git
a/ui/src/app/connect/components/existing-adapters/existing-adapters.component.ts
b/ui/src/app/connect/components/existing-adapters/existing-adapters.component.ts
index 15fcb4cd9..4cfb1c7b8 100644
---
a/ui/src/app/connect/components/existing-adapters/existing-adapters.component.ts
+++
b/ui/src/app/connect/components/existing-adapters/existing-adapters.component.ts
@@ -23,8 +23,8 @@ import {
AdapterMonitoringService,
PipelineElementService,
SpMetricsEntry,
- StreamPipesErrorMessage,
PipelineService,
+ SpLogMessage,
} from '@streampipes/platform-services';
import { MatTableDataSource } from '@angular/material/table';
import { ConnectService } from '../../services/connect.service';
@@ -161,10 +161,7 @@ export class ExistingAdaptersComponent implements OnInit {
});
}
- openAdapterStatusErrorDialog(
- message: StreamPipesErrorMessage,
- title: string,
- ) {
+ openAdapterStatusErrorDialog(message: SpLogMessage, title: string) {
this.dialogService.open(SpExceptionDetailsDialogComponent, {
panelType: PanelType.STANDARD_PANEL,
title: 'Adapter Status',
diff --git
a/ui/src/app/core-ui/static-properties/static-runtime-resolvable-input/base-runtime-resolvable-input.ts
b/ui/src/app/core-ui/static-properties/static-runtime-resolvable-input/base-runtime-resolvable-input.ts
index e1075eb1c..925ffedc9 100644
---
a/ui/src/app/core-ui/static-properties/static-runtime-resolvable-input/base-runtime-resolvable-input.ts
+++
b/ui/src/app/core-ui/static-properties/static-runtime-resolvable-input/base-runtime-resolvable-input.ts
@@ -23,6 +23,7 @@ import {
RuntimeResolvableAnyStaticProperty,
RuntimeResolvableOneOfStaticProperty,
RuntimeResolvableTreeInputStaticProperty,
+ SpLogMessage,
StaticProperty,
StaticPropertyUnion,
TreeInputNode,
@@ -31,7 +32,6 @@ import { RuntimeResolvableService } from
'./runtime-resolvable.service';
import { Observable } from 'rxjs';
import { Directive, Input, OnChanges, SimpleChanges } from '@angular/core';
import { ConfigurationInfo } from '../../../connect/model/ConfigurationInfo';
-import { StreamPipesErrorMessage } from
'../../../../../projects/streampipes/platform-services/src/lib/model/gen/streampipes-model';
@Directive()
// eslint-disable-next-line @angular-eslint/directive-class-suffix
@@ -50,7 +50,7 @@ export abstract class BaseRuntimeResolvableInput<
showOptions = false;
loading = false;
error = false;
- errorMessage: StreamPipesErrorMessage;
+ errorMessage: SpLogMessage;
dependentStaticProperties: any = new Map();
constructor(private runtimeResolvableService: RuntimeResolvableService) {
@@ -111,8 +111,7 @@ export abstract class BaseRuntimeResolvableInput<
this.loading = false;
this.showOptions = true;
this.error = true;
- this.errorMessage =
- errorMessage.error as StreamPipesErrorMessage;
+ this.errorMessage = errorMessage.error as SpLogMessage;
this.afterErrorReceived();
},
);
diff --git
a/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.html
b/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.html
index bf0fd1298..82490189c 100644
---
a/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.html
+++
b/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.html
@@ -142,6 +142,7 @@
</div>
</div>
<div
+ fxLayout="column"
class="widget-content p-0 gridster-item-content ml-0 mr-0 h-100
mw-100"
>
<ng-template spWidgetHost class="h-100 p-0"></ng-template>
diff --git
a/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.ts
b/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.ts
index e3ecf6181..7c1b17b71 100644
---
a/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.ts
+++
b/ui/src/app/data-explorer/components/widget/data-explorer-dashboard-widget.component.ts
@@ -35,6 +35,7 @@ import {
DataLakeMeasure,
DataViewDataExplorerService,
DateRange,
+ SpLogMessage,
TimeSettings,
} from '@streampipes/platform-services';
import { DataDownloadDialogComponent } from
'../../../core-ui/data-download-dialog/data-download-dialog.component';
@@ -51,7 +52,6 @@ import {
DialogService,
PanelType,
} from '@streampipes/shared-ui';
-import { StreamPipesErrorMessage } from
'../../../../../projects/streampipes/platform-services/src/lib/model/gen/streampipes-model';
@Component({
selector: 'sp-data-explorer-dashboard-widget',
@@ -111,7 +111,7 @@ export class DataExplorerDashboardWidgetComponent
implements OnInit, OnDestroy {
widgetTypeChangedSubscription: Subscription;
intervalSubscription: Subscription;
- errorMessage: StreamPipesErrorMessage;
+ errorMessage: SpLogMessage;
componentRef: ComponentRef<BaseWidgetData<any>>;
diff --git
a/ui/src/app/data-explorer/components/widgets/base/base-data-explorer-widget.directive.ts
b/ui/src/app/data-explorer/components/widgets/base/base-data-explorer-widget.directive.ts
index 4ca1f9a85..cc197ac2c 100644
---
a/ui/src/app/data-explorer/components/widgets/base/base-data-explorer-widget.directive.ts
+++
b/ui/src/app/data-explorer/components/widgets/base/base-data-explorer-widget.directive.ts
@@ -33,8 +33,8 @@ import {
DataExplorerWidgetModel,
DatalakeRestService,
DataViewQueryGeneratorService,
+ SpLogMessage,
SpQueryResult,
- StreamPipesErrorMessage,
TimeSettings,
} from '@streampipes/platform-services';
import { ResizeService } from '../../../services/resize.service';
@@ -59,8 +59,8 @@ export abstract class BaseDataExplorerWidgetDirective<
timerCallback: EventEmitter<boolean> = new EventEmitter();
@Output()
- errorCallback: EventEmitter<StreamPipesErrorMessage> =
- new EventEmitter<StreamPipesErrorMessage>();
+ errorCallback: EventEmitter<SpLogMessage> =
+ new EventEmitter<SpLogMessage>();
@Input() gridsterItem: GridsterItem;
@Input() gridsterItemComponent: GridsterItemComponent;
diff --git
a/ui/src/app/data-explorer/components/widgets/base/data-explorer-widget-data.ts
b/ui/src/app/data-explorer/components/widgets/base/data-explorer-widget-data.ts
index dc234e9f6..9a43eb0d5 100644
---
a/ui/src/app/data-explorer/components/widgets/base/data-explorer-widget-data.ts
+++
b/ui/src/app/data-explorer/components/widgets/base/data-explorer-widget-data.ts
@@ -21,14 +21,14 @@ import { GridsterItem, GridsterItemComponent } from
'angular-gridster2';
import {
DashboardItem,
DataExplorerWidgetModel,
- StreamPipesErrorMessage,
+ SpLogMessage,
TimeSettings,
} from '@streampipes/platform-services';
export interface BaseWidgetData<T extends DataExplorerWidgetModel> {
removeWidgetCallback: EventEmitter<boolean>;
timerCallback: EventEmitter<boolean>;
- errorCallback: EventEmitter<StreamPipesErrorMessage>;
+ errorCallback: EventEmitter<SpLogMessage>;
gridsterItem: GridsterItem;
gridsterItemComponent: GridsterItemComponent;
diff --git a/ui/src/app/login/components/startup/startup.component.scss
b/ui/src/app/login/components/startup/startup.component.scss
index 9ebbfd0cc..eea22d50e 100644
--- a/ui/src/app/login/components/startup/startup.component.scss
+++ b/ui/src/app/login/components/startup/startup.component.scss
@@ -18,8 +18,8 @@
@import '../../../../scss/_variables.scss';
-::ng-deep .mat-progress-bar-fill::after {
- background: $sp-color-accent;
+::ng-deep .mdc-linear-progress__bar-inner {
+ border-color: $sp-color-accent !important;
}
.sp-progress {