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
13 changes: 13 additions & 0 deletions hugegraph-hubble/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,20 @@ RUN set -x \

FROM eclipse-temurin:11-jre-jammy

RUN apt-get -q update \
&& apt-get install --no-install-recommends curl -yq \
&& apt-get clean \
&& rm -rf /var/lib/apt/lists/*

COPY --from=build /pkg/hugegraph-hubble/apache-hugegraph-hubble-*/ /hubble
RUN sed -i \
-e 's/^server\.host=.*/server.host=0.0.0.0/' \
-e 's/^dashboard\.address=.*/dashboard.address=/' \
/hubble/conf/hugegraph-hubble.properties \
&& grep -Fqx 'server.host=0.0.0.0' \
/hubble/conf/hugegraph-hubble.properties \
&& grep -Fqx 'dashboard.address=' \
/hubble/conf/hugegraph-hubble.properties
WORKDIR /hubble/

# SECURITY: This is a plain HTTP port. Do not publish it to an untrusted network;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.apache.hugegraph.common.Mergeable;
import org.apache.hugegraph.util.EntityUtil;
import org.apache.hugegraph.util.Ex;
import org.apache.hugegraph.util.HugeClientUtil;

@Component
public abstract class BaseController {
Expand Down Expand Up @@ -196,7 +197,11 @@ protected HugeClient tempTokenClient() {
HttpServletRequest request = getRequest();
if (request.getAttribute("hugeClient") != null) {
HugeClient client = (HugeClient) request.getAttribute("hugeClient");
client.setAuthContext("Basic " + this.getToken());
String token = this.getToken();
if (org.apache.commons.lang3.StringUtils.isNotBlank(token)) {
client.setAuthContext(
HugeClientUtil.bearerAuthContext(token));
}
return client;
}
HugeClient client = this.hugeClientPoolService.createTempTokenClient(this.getToken());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,9 @@

import com.baomidou.mybatisplus.core.metadata.IPage;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.Data;
import lombok.extern.log4j.Log4j2;
import org.apache.commons.io.FileUtils;
Expand All @@ -45,6 +48,7 @@
import org.apache.hugegraph.entity.load.ValueMappingItem;
import org.apache.hugegraph.entity.load.VertexMapping;
import org.apache.hugegraph.exception.InternalException;
import org.apache.hugegraph.loader.source.file.FileFormat;
import org.apache.hugegraph.options.HubbleOptions;
import org.apache.hugegraph.service.load.DatasourceService;
import org.apache.hugegraph.service.load.FileMappingService;
Expand Down Expand Up @@ -76,6 +80,7 @@
import java.util.Collections;
import java.util.Date;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
Expand All @@ -88,6 +93,10 @@
@RequestMapping(Constant.API_VERSION + "ingest")
public class IngestController extends BaseController {

private static final ObjectMapper JSON_MAPPER =
new ObjectMapper().disable(
JsonParser.Feature.INCLUDE_SOURCE_IN_LOCATION);

@Autowired
private JobManagerService jobManagerService;
@Autowired
Expand Down Expand Up @@ -193,10 +202,11 @@ public Response datasourceSchema(@RequestParam("datasource") int id) {
.data(header).build();
}

boolean hasPhysicalHeader = this.hasPhysicalHeader(config);
FileSetting setting = this.buildFileSetting(config, Collections.emptyList(),
false);
hasPhysicalHeader);
ColumnInfo columns = this.readColumns(this.requireUploadFile(config),
setting, false);
setting, hasPhysicalHeader);
return Response.builder().status(Constant.STATUS_OK)
.data(columns.names).build();
}
Expand Down Expand Up @@ -234,8 +244,7 @@ public Response createTask(@RequestBody IngestTaskRequest request) {
File sourceFile = this.requireUploadFile(input);
long totalSize = sourceFile.length();
List<String> header = this.stringList(input.get("header"));
boolean hasPhysicalHeader = !this.stringList(dsConfig.get("header"))
.isEmpty();
boolean hasPhysicalHeader = this.hasPhysicalHeader(dsConfig);
FileSetting setting = this.buildFileSetting(input, header,
hasPhysicalHeader);
Comment thread
imbajin marked this conversation as resolved.
if (setting.getColumnNames() == null ||
Expand Down Expand Up @@ -548,6 +557,13 @@ private FileSetting buildFileSetting(Map<String, Object> input,
return setting;
}

private boolean hasPhysicalHeader(Map<String, Object> config) {
FileFormat format = FileFormat.valueOf(
this.stringOrDefault(config.get("format"), "CSV"));
return format.needHeader() &&
this.stringList(config.get("header")).isEmpty();
}

private ColumnInfo readColumns(File file, FileSetting setting,
boolean hasPhysicalHeader) {
try (BufferedReader reader = Files.newBufferedReader(file.toPath())) {
Expand All @@ -559,6 +575,10 @@ private ColumnInfo readColumns(File file, FileSetting setting,
}
Ex.check(line != null, "The file has no data line can treat as header");

if (FileFormat.JSON.name().equals(setting.getFormat())) {
return this.readJsonColumns(line, file);
}

List<String> firstLine = this.splitLine(line, setting.getDelimiter());
if (hasPhysicalHeader) {
String sample = reader.readLine();
Expand All @@ -579,6 +599,23 @@ private ColumnInfo readColumns(File file, FileSetting setting,
}
}

private ColumnInfo readJsonColumns(String line, File file) {
try {
Map<String, Object> fields = JSON_MAPPER.readValue(
line, new TypeReference<LinkedHashMap<String, Object>>() {
});
List<String> names = new ArrayList<>(fields.keySet());
List<String> values = names.stream()
.map(fields::get)
.map(String::valueOf)
.collect(Collectors.toList());
return new ColumnInfo(names, values);
} catch (IOException ignored) {
throw new InternalException(
"Failed to read JSON fields from file %s", file);
}
}

private List<String> splitLine(String line, String delimiter) {
if (line == null) {
return Collections.emptyList();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -285,13 +285,13 @@ private static class SaaSMetrics {
private long vertexLabelCount;

@JsonProperty("vertex-count")
private long vertexCount;
private Long vertexCount;

@JsonProperty("edge-label-count")
private long edgeLabelCount;

@JsonProperty("edge-count")
private long edgeCount;
private Long edgeCount;

@JsonProperty("graph-space-count")
private long graphSpaceCount;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.apache.hugegraph.entity.space.BuiltInEntity;
import org.apache.hugegraph.exception.ServerException;
import org.apache.hugegraph.loader.util.JsonUtil;
import org.apache.hugegraph.options.HubbleOptions;
import org.apache.hugegraph.service.algorithm.AsyncTaskService;
import org.apache.hugegraph.service.auth.UserService;
import org.apache.hugegraph.service.load.LoadTaskService;
Expand Down Expand Up @@ -784,11 +785,15 @@ public Map<String, Object> evCount(HugeClient client,
Long vertexCount = null;
String statisticDate = HubbleUtil.dateFormatDay(HubbleUtil.nowDate());
client.assignGraph(graphSpace, graph);
GraphMetricsAPI.ElementCount statistic =
client.graph().getEVCount(statisticDate);
if (statistic == null) {
statisticDate = HubbleUtil.dateFormatLastDay();
GraphMetricsAPI.ElementCount statistic = null;
boolean pdEnabled = this.config != null &&
this.config.get(HubbleOptions.PD_ENABLED);
if (!pdEnabled) {
statistic = client.graph().getEVCount(statisticDate);
Comment thread
imbajin marked this conversation as resolved.
if (statistic == null) {
statisticDate = HubbleUtil.dateFormatLastDay();
statistic = client.graph().getEVCount(statisticDate);
}
}

if (statistic != null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -557,6 +557,7 @@ private FileSource buildFileSource(FileMapping fileMapping) {
Ex.check(setting.getColumnNames() != null,
"Must do file setting firstly");
source.header(setting.getColumnNames().toArray(new String[]{}));
source.hasHeader(setting.isHasHeader());
// NOTE: format and delimiter must be CSV and "," temporarily
source.format(FileFormat.valueOf(setting.getFormat()));
source.delimiter(setting.getDelimiter());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,8 +72,8 @@ public class GraphSpaceService {
public Map<String, Long> metrics(HugeClient client) {
long gsCount = 0L;
long gCount = 0L;
long vCount = 0L;
long eCount = 0L;
Long vCount = 0L;
Long eCount = 0L;
long vlCount = 0L;
long elCount = 0L;
long preDayTaskCount = 0L;
Expand All @@ -86,8 +86,10 @@ public Map<String, Long> metrics(HugeClient client) {
Map<String, Object> elVl = elAndVlCount(client, gs);
Map<String, Object> task = preDayTaskCount(client, gs);

vCount += ((Number) ev.get("vertex")).longValue();
eCount += ((Number) ev.get("edge")).longValue();
vCount = addAvailableCount(vCount,
(Number) ev.get("vertex"));
eCount = addAvailableCount(eCount,
(Number) ev.get("edge"));
vlCount += ((Number) elVl.get("vertexlabel")).longValue();
elCount += ((Number) elVl.get("edgelabel")).longValue();
preDayTaskCount += ((Number) task.get("task")).longValue();
Expand Down Expand Up @@ -199,8 +201,8 @@ private static void removeSensitiveFields(Map<String, Object> info) {
* @return
*/
Map<String, Object> evCount(HugeClient client, String graphSpace) {
long vertexTotal = 0L;
long edgeTotal = 0L;
Long vertexTotal = 0L;
Long edgeTotal = 0L;
Map<String, Object> statisticTotal = new HashMap<>();
client.assignGraph(graphSpace, "");
Set<String> graphs = graphsService.listGraphNames(client, graphSpace, "");
Expand All @@ -219,8 +221,10 @@ Map<String, Object> evCount(HugeClient client, String graphSpace) {
statisticDate = null;
}

vertexTotal += ((Number) graphEvCount.get("vertex")).longValue();
edgeTotal += ((Number) graphEvCount.get("edge")).longValue();
Number vertexCount = (Number) graphEvCount.get("vertex");
Number edgeCount = (Number) graphEvCount.get("edge");
vertexTotal = addAvailableCount(vertexTotal, vertexCount);
edgeTotal = addAvailableCount(edgeTotal, edgeCount);
}
if (graphs.isEmpty()) {
statisticDate = HubbleUtil.dateFormatDay(HubbleUtil.nowDate());
Expand All @@ -232,6 +236,13 @@ Map<String, Object> evCount(HugeClient client, String graphSpace) {
return statisticTotal;
}

private static Long addAvailableCount(Long total, Number count) {
if (total == null || count == null) {
return null;
}
return Math.addExact(total, count.longValue());
}

/**
* 统计指定图空间下的edgeLabel总数和vertexLabel边总数
* @param client
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
public final class HugeClientUtil {

private static final String DEFAULT_PROTOCOL = "http";
private static final String BEARER_SCHEME = "Bearer";

private static final Set<String> ACCEPTABLE_EXCEPTIONS = ImmutableSet.of(
"Permission denied: execute Resource"
Expand All @@ -44,7 +45,12 @@ public static HugeClient tryConnect(GraphConnection connection) {
String graph = connection.getGraph();
String host = connection.getHost();
Integer port = connection.getPort();
String token = normalizeToken(connection.getToken());
String tokenInput = connection.getToken();
String token = normalizeTokenPayload(tokenInput);
String authContext = null;
if (StringUtils.isNotBlank(tokenInput)) {
authContext = bearerAuthContext(tokenInput);
}
String username = connection.getUsername();
String password = connection.getPassword();
int timeout = connection.getTimeout();
Expand Down Expand Up @@ -112,14 +118,36 @@ public static HugeClient tryConnect(GraphConnection connection) {
throw e;
}

if (authContext != null) {
client.setAuthContext(authContext);
}

return client;
}

private static String normalizeToken(String token) {
if (StringUtils.isBlank(token) || token.startsWith(" ")) {
return token;
public static String bearerAuthContext(String token) {
String normalized = normalizeTokenPayload(token);
if (StringUtils.isBlank(normalized)) {
throw new IllegalArgumentException(
"Bearer token must contain a non-blank payload");
}
return BEARER_SCHEME + " " + normalized;
}

private static String normalizeTokenPayload(String token) {
String normalized = StringUtils.stripToNull(token);
if (normalized == null) {
return null;
}
int schemeLength = BEARER_SCHEME.length();
if (normalized.length() >= schemeLength &&
normalized.regionMatches(true, 0, BEARER_SCHEME, 0, schemeLength) &&
(normalized.length() == schemeLength ||
Character.isWhitespace(normalized.charAt(schemeLength)))) {
normalized = StringUtils.stripToNull(
normalized.substring(schemeLength));
}
return " " + token;
return normalized;
}

private static boolean isAcceptable(String message) {
Expand Down
Loading
Loading