Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions plc4j/tools/merlot/org.apache.plc4x.logrecorder/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -199,10 +199,13 @@
<configuration>
<artifacts>
<artifact>
<file
>src/main/cfg/org.apache.plc4x.merlot.logrecorder.cfg</file>
<file>src/main/cfg/org.apache.plc4x.merlot.logrecorder.cfg</file>
<type>cfg</type>
</artifact>
<artifact>
<file >src/main/cfg/org.apache.plc4x.merlot.logrecorder.connection.database.cfg</file>
<type>cfg2</type>
</artifact>
</artifacts>
</configuration>
</execution>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,12 +50,11 @@ public class MerlotLogRecorderJDBCAppender implements EventHandler, ManagedServi

private Map<String, String> connectionProperties = new HashMap();

//OPS4J support: Derby, H2, MariaDB, MySQL, PostgreSQL, SQLite
// //OPS4J support: Derby, H2, MariaDB, MySQL, PostgreSQL, SQLite
private final static String createTableQueryGenericTemplate
= "CREATE TABLE IF NOT EXISTS TABLENAME(id BIGINT NOT NULL PRIMARY KEY, owner VARCHAR(255),"
+ " level VARCHAR(100), description VARCHAR(4000), title VARCHAR(255), createdDate BIGINT, tags VARCHAR(255), logbooks VARCHAR(255), attachments_path VARCHAR(1500))";

//OPS4J support: Oracle
//OPS4J support: Oracle
private final static String createTableQueryOracleTemplate
= "CREATE TABLE IF NOT EXISTS TABLENAME(id NUMBER(19) NOT NULL PRIMARY KEY, owner VARCHAR2(255),"
+ " level VARCHAR2(100), description VARCHAR2(4000), title VARCHAR2(255), createdDate NUMBER(19), tags VARCHAR2(255), logbooks VARCHAR2(255), attachments_path VARCHAR(1500))";
Expand Down Expand Up @@ -96,6 +95,7 @@ public void handleEvent(Event event) {

try (Connection connection = dataSource.getConnection()) {
String insertQuery = insertQueryTemplate.replaceAll("TABLENAME", this.connectionProperties.get(TABLE_NAME_PROPERTY));

try (PreparedStatement insertStatement = connection.prepareStatement(insertQuery)) {
insertStatement.setLong(1, id);
insertStatement.setString(2, owner);
Expand Down Expand Up @@ -183,6 +183,7 @@ private void createTable(Connection connection) {

String createTableQuery = createTemplate.replaceAll("TABLENAME", this.connectionProperties.get(TABLE_NAME_PROPERTY));


try (Statement createStatement = connection.createStatement()) {
createStatement.executeUpdate(createTableQuery);
LOGGER.info("Table {} has been created", this.connectionProperties.get(TABLE_NAME_PROPERTY));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,20 +27,25 @@

public class MerlotLogRecorderFileExplorer {

private final static String MERLOT_DATA_DIRECTORY = "karaf.data";

private MerlotLogRecorderFileExplorer() {

}

public static File findFileByFilename(String searchTerm, BundleContext ctx) {

String fullDirPath = System.getenv("MERLOT_OLOG_ATTACHMENT");

if (fullDirPath == null || fullDirPath.trim().isEmpty()) {
return null;
}


String karafDataDir = ctx.getProperty(MERLOT_DATA_DIRECTORY);
//Search the data/tmp directory in Karaf
Path fileTarget = Paths.get(karafDataDir, "tmp", searchTerm);
Path fileTarget = Paths.get(fullDirPath, searchTerm);

if (Files.exists(fileTarget) && Files.isRegularFile(fileTarget)) {
return fileTarget.toFile();
}

return null;

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.util.Map;
import java.util.Random;
import javax.security.auth.login.LoginException;
import javax.servlet.ServletException;
Expand Down Expand Up @@ -75,11 +76,10 @@ private void processMultipart(HttpServletRequest req, HttpServletResponse resp

//Getting the application sections
for (Part part : req.getParts()) {
String directoryPath = "data/tmp";
String directoryPath = System.getenv("MERLOT_OLOG_ATTACHMENT");
String fileName = part.getSubmittedFileName();

if (part.getContentType().equals("application/json")) {
// String json = new String(part.getInputStream().readAllBytes());
try (InputStream is = part.getInputStream()) {
JsonNode node = mapper.readTree(is);

Expand Down Expand Up @@ -220,4 +220,6 @@ public String extractJSONData(JsonNode n, String nodeName) {
return dataResult;

}


}
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws Se
JSONObject log = new JSONObject();
log.put("id", dataLog.getId());
log.put("owner", dataLog.getOwner());
log.put("source", "source");//TODO: This is where the machine's parameters should be listed—the machine that generated the log—but Phoebus doesn't send them.
log.put("source", dataLog.getDescription());
log.put("level", dataLog.getLevel());
log.put("title", dataLog.getTitle());
log.put("createdDate", dataLog.getCreatedDate());
Expand Down Expand Up @@ -330,8 +330,7 @@ private void addBetweenFilter(StringBuilder sql, List<Object> params,
}
}


public String mappingTypeToMIMEHTTP(String paramType) {
public String mappingTypeToMIMEHTTP(String paramType) {
String extension = paramType.substring(paramType.lastIndexOf(".") + 1);
String extensionFinal = null;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,10 @@
*/
package org.apache.plc4x.merlot.archiver.impl;

import io.netty.buffer.ByteBuf;
import io.netty.buffer.ByteBufUtil;
import io.netty.buffer.Unpooled;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.OutputStream;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
Expand All @@ -39,9 +37,11 @@ public class MerlotDataBrowserIoTDBGetDataPVImpl extends HttpServlet {

private static final Logger LOGGER = LoggerFactory.getLogger(MerlotDataBrowserRTGetDataPVImpl.class);

private final Pattern opti_pattern = Pattern.compile("optimized_11520\\(([^)]+)\\)");
private final Pattern opti_pattern = Pattern.compile("optimized_\\w+\\(([^)]+)\\)");

private final Pattern ncount_pattern = Pattern.compile("ncount\\(([^)]+)\\)");
private final Pattern count_pattern = Pattern.compile("count_3600\\(([^)]+)\\)");

private final Pattern count_pattern = Pattern.compile("count_\\w+\\(([^)]+)\\)");

private Matcher opti_matcher = null;
private Matcher ncount_matcher = null;
Expand All @@ -55,58 +55,70 @@ public MerlotDataBrowserIoTDBGetDataPVImpl(MerlotHtc mhtc) {

@Override
protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {


String from = req.getParameter("from");
String to = req.getParameter("to");
String[] pvs = req.getParameterValues("pv");




LOGGER.info("Inicio Servlet.");
if ((null == from) || (null == to)) {
if (from == null || to == null) {
resp.sendError(HttpServletResponse.SC_BAD_REQUEST, "Missing 'from' or 'to'");
return;
}
if ((null == pvs) || (pvs.length == 0)) {
if (pvs == null || pvs.length == 0) {
resp.sendError(HttpServletResponse.SC_BAD_REQUEST, "Missing 'pv' parameter");
return;
}

resp.setContentType("application/octet-stream");
for (String pv : pvs) {
opti_matcher = opti_pattern.matcher(pv);
ncount_matcher = ncount_pattern.matcher(pv);
count_matcher = count_pattern.matcher(pv);
if (pvs.length > 1) {
resp.sendError(HttpServletResponse.SC_BAD_REQUEST, "Multiple PVs not supported");
return;
}

String pv = pvs[0];

opti_matcher = opti_pattern.matcher(pv);
ncount_matcher = ncount_pattern.matcher(pv);
count_matcher = count_pattern.matcher(pv);

resp.setContentType("application/octet-stream");
try (OutputStream out = resp.getOutputStream()) {
if (opti_matcher.matches()) {
LOGGER.info("optimized_11520(pv) not supported.");
resp.sendError(HttpServletResponse.SC_NOT_IMPLEMENTED, "optimized_* not supported");
return;
} else if (ncount_matcher.matches()) {
String strpv = ncount_matcher.group(1);
int countpv = mhtc.countPVs(strpv, from, to);
LOGGER.info("Number of events: " + countpv);
resp.getWriter().print(countpv);
resp.getWriter().close();
byte[] bytes = Integer.toString(countpv).getBytes(StandardCharsets.UTF_8);
resp.setContentType("text/plain; charset=utf-8");
resp.setContentLength(bytes.length);
out.write(bytes);
out.flush();
return;
} else if (count_matcher.matches()) {
LOGGER.info("count_3600(pv) not supported.");
resp.getWriter().close();
resp.sendError(HttpServletResponse.SC_NOT_IMPLEMENTED, "count_* not supported");
return;
} else {
createRawResponse(pv, from, to, resp.getOutputStream(), resp);
resp.getOutputStream().close();
boolean ok = createRawResponse(pv, from, to, out, resp);
if (!ok) {
return;
}
out.flush();
}
}

}

private void createRawResponse(String pv, String from, String to, OutputStream out, HttpServletResponse resp) throws IOException {

private boolean createRawResponse(String pv, String from, String to, OutputStream out, HttpServletResponse resp) throws IOException {
List<VType> values = mhtc.getPVs(pv, from, to);
if (values == null) {
out.close();
resp.sendError(HttpServletResponse.SC_NOT_FOUND, "No values for pv: " + pv);
return false;
}
ByteArrayOutputStream bout = new ByteArrayOutputStream();
MerlotPBRawSerializer.serializeIoTDBToPBRaw(values, pv, bout);
ByteBuf buf = Unpooled.wrappedBuffer(bout.toByteArray());
// System.out.println(ByteBufUtil.prettyHexDump(buf));
out.write(bout.toByteArray());
byte[] payload = bout.toByteArray();
resp.setContentType("application/octet-stream");
out.write(payload);
return true;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,6 @@
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import org.apache.iotdb.isession.pool.SessionDataSetWrapper;
import org.apache.iotdb.pipe.api.type.Type;
import org.apache.iotdb.rpc.IoTDBConnectionException;
Expand All @@ -46,10 +44,6 @@
import org.osgi.service.cm.ManagedService;
import org.slf4j.LoggerFactory;

/**
*
* @author cgarcia
*/
public class MerlotHtcIoTDBImpl implements MerlotHtc, ManagedService {

private static final org.slf4j.Logger LOGGER = LoggerFactory.getLogger(MerlotHtcIoTDBImpl.class);
Expand Down Expand Up @@ -77,15 +71,15 @@ public synchronized SessionPool getIoTDBConnection() throws IoTDBConnectionExcep
} else {
try {
sp = new SessionPool.Builder()
.nodeUrls((new ArrayList<>(this.urls)))
.user(this.username)
.password(this.password)
.maxSize(this.maxThreadPool)
.enableAutoFetch(this.enableAutoFetch)
.build();
.nodeUrls((new ArrayList<>(this.urls)))
.user(this.username)
.password(this.password)
.maxSize(this.maxThreadPool)
.enableAutoFetch(this.enableAutoFetch)
.build();

} catch (Exception e) {
System.out.println("Error: " + e.getMessage());
LOGGER.info("Error: {}", e.getMessage());
}
}
return null;
Expand Down Expand Up @@ -124,8 +118,7 @@ public Set<String> getPVs() {
while (ds.hasNext()) {

String fullPath = ds.next().getFields().get(0).getStringValue();

pvs.add(fullPath.replaceFirst("^[^.]+\\.(.*)", "$1"));
pvs.add(String.format("pva://%s", fullPath.replaceFirst("^[^.]+\\.(.*)", "$1").replaceAll("\\.(?![^.]*$)", "/")));
}
} catch (Exception ex) {
LOGGER.error("Error retrieving PVs from IoTDB: {}", ex.getMessage());
Expand All @@ -137,15 +130,21 @@ public Set<String> getPVs() {
@Override
public List<VType> getPVs(String strPV, String init, String end) {
List<VType> listResult = new ArrayList<>(5000);

try {
String device = getBasePath(strPV);
String measurement = getTimeserieNameSimple(strPV);
String adjusmentVariable = strPV.replaceFirst("^pva://", "").replaceAll("/", ".");
String device = getBasePath(adjusmentVariable);
String measurement = getTimeserieNameSimple(adjusmentVariable);

long startT = Instant.parse(init).toEpochMilli();
long endT = Instant.parse(end).toEpochMilli();

if (startT == endT) {
startT = endT - (60 * 1000);
}

String sql = String.format("SELECT %s FROM root.%s WHERE time >= %d AND time <= %d",
measurement, device, startT, endT);
measurement, device, startT, endT);

SessionPool pool = getIoTDBConnection();
if (pool == null) {
Expand All @@ -166,7 +165,7 @@ public List<VType> getPVs(String strPV, String init, String end) {
Instant inst = Instant.ofEpochMilli(timestamp);

MerlotIoTDBMapping.EpicsMetadata meta
= new MerlotIoTDBMapping.EpicsMetadata((int) inst.getEpochSecond(), inst.getNano(), 0, 0);
= new MerlotIoTDBMapping.EpicsMetadata((int) inst.getEpochSecond(), inst.getNano(), 0, 0);

Object epicsEvent = mapper.convert(field.getObjectValue(field.getDataType()), meta);
translateToScalarType(inst, epicsEvent, listResult);
Expand All @@ -180,6 +179,7 @@ public List<VType> getPVs(String strPV, String init, String end) {
return listResult;
}

//TODO: Mover esto a MerlotIoTDBMapping.java
private static List<VType> translateToScalarType(Instant inst, Object epicsEvent, List<VType> listEvents) {

// --- CONVERSION BASED STRICTLY ON MerlotIoTDBMapping ---
Expand Down Expand Up @@ -211,7 +211,7 @@ private static List<VType> translateToScalarType(Instant inst, Object epicsEvent

@Override
public int countPVs(String strPV, String init,
String end
String end
) {

return 0;
Expand All @@ -221,14 +221,10 @@ public int countPVs(String strPV, String init,
Returns the variable stored in the IoTDB device (pvName)
*/
private static String getTimeserieNameSimple(String pvName) {
Pattern pattern = Pattern.compile("[^.]+$");
Matcher matcher = pattern.matcher(pvName);

if (matcher.find()) {
return matcher.group();
if (pvName == null || pvName.isEmpty()) {
return "";
}

return "";
return pvName.replaceFirst("^.*\\.", "");
}

/*
Expand All @@ -238,7 +234,7 @@ public static String getBasePath(String pvName) {
if (pvName == null || pvName.isEmpty()) {
return "";
}
return pvName.replaceFirst("^(?:root\\.)?(.*)\\.[^.]+$", "$1");
return pvName.replaceFirst("\\.[^.]*$", "");
}

public void updated(Dictionary<String, ?> properties) throws ConfigurationException {
Expand Down
Loading
Loading