diff --git a/geaflow-ai/pom.xml b/geaflow-ai/pom.xml
index 006215b74..9d450f890 100644
--- a/geaflow-ai/pom.xml
+++ b/geaflow-ai/pom.xml
@@ -121,8 +121,31 @@
geaflow-api
${project.version}
+
+ org.apache.geaflow
+ geaflow-dsl-common
+ ${project.version}
+
+
+ org.apache.geaflow
+ geaflow-pipeline
+ ${project.version}
+ test
+
+
+ org.apache.geaflow
+ geaflow-on-local
+ ${project.version}
+ test
+
+
+ org.apache.geaflow
+ geaflow-dsl-runtime
+ ${project.version}
+ test
+
org.junit.jupiter
junit-jupiter
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/adapter/MemoryGraphAdapter.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/adapter/MemoryGraphAdapter.java
new file mode 100644
index 000000000..8edb2fb0f
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/adapter/MemoryGraphAdapter.java
@@ -0,0 +1,1247 @@
+/*
+ * 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.geaflow.ai.temporal.adapter;
+
+import java.time.Instant;
+import java.time.format.DateTimeParseException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.TreeMap;
+import org.apache.geaflow.ai.graph.io.Edge;
+import org.apache.geaflow.ai.graph.io.EdgeGroup;
+import org.apache.geaflow.ai.graph.io.EdgeSchema;
+import org.apache.geaflow.ai.graph.io.EntityGroup;
+import org.apache.geaflow.ai.graph.io.GraphSchema;
+import org.apache.geaflow.ai.graph.io.MemoryGraph;
+import org.apache.geaflow.ai.graph.io.Vertex;
+import org.apache.geaflow.ai.graph.io.VertexGroup;
+import org.apache.geaflow.ai.graph.io.VertexSchema;
+import org.apache.geaflow.ai.temporal.model.Evidence;
+import org.apache.geaflow.ai.temporal.model.FactKey;
+import org.apache.geaflow.ai.temporal.model.FactValue;
+import org.apache.geaflow.ai.temporal.model.MemoryEntity;
+import org.apache.geaflow.ai.temporal.model.MemoryEvent;
+import org.apache.geaflow.ai.temporal.model.MemoryEventOperation;
+import org.apache.geaflow.ai.temporal.model.MemoryFact;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersion;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersionStatus;
+import org.apache.geaflow.ai.temporal.model.Source;
+import org.apache.geaflow.ai.temporal.model.TimeInterval;
+import org.apache.geaflow.ai.temporal.model.VersionRelation;
+import org.apache.geaflow.ai.temporal.model.VersionRelationType;
+import org.apache.geaflow.ai.temporal.semantics.CanonicalSnapshot;
+import org.apache.geaflow.ai.temporal.semantics.EventNormalizer;
+import org.apache.geaflow.ai.temporal.semantics.NormalizedMemoryEvent;
+import org.apache.geaflow.ai.temporal.semantics.TemporalState;
+
+/**
+ * Projects canonical temporal snapshots to the existing in-memory graph model.
+ */
+public final class MemoryGraphAdapter {
+
+ private static final String ENTITY = "entity";
+ private static final String FACT_VERSION = "fact_version";
+ private static final String MEMORY_EVENT = "memory_event";
+ private static final String EVIDENCE = "evidence";
+ private static final String SOURCE = "source";
+
+ private static final String SUBJECT = "subject";
+ private static final String OBJECT = "object";
+ private static final String GENERATES = "generates";
+ private static final String SUPPORTED_BY = "supported_by";
+ private static final String FROM_SOURCE = "from_source";
+ private static final String SUPERSEDES = "supersedes";
+ private static final String DUPLICATE_OF = "duplicate_of";
+ private static final String CONFLICTS_WITH = "conflicts_with";
+
+ private static final String ENTITY_PREFIX = "entity:";
+ private static final String VERSION_PREFIX = "version:";
+ private static final String EVENT_PREFIX = "event:";
+ private static final String EVIDENCE_PREFIX = "evidence:";
+ private static final String SOURCE_PREFIX = "source:";
+ private static final String EMPTY = "";
+
+ private static final List ENTITY_FIELDS = fields("label");
+ private static final List VERSION_FIELDS = fields(
+ "factId",
+ "predicate",
+ "scope",
+ "valueKind",
+ "literalValue",
+ "status",
+ "validStart",
+ "validEnd",
+ "transactionStart",
+ "transactionEnd");
+ private static final List EVENT_FIELDS = fields(
+ "operation",
+ "factId",
+ "subjectId",
+ "predicate",
+ "scope",
+ "valueKind",
+ "value",
+ "validStart",
+ "validEnd",
+ "recordedAt",
+ "payloadHash");
+ private static final List EVIDENCE_FIELDS = fields("content");
+ private static final List SOURCE_FIELDS = fields("name");
+ private static final List NO_FIELDS = Collections.emptyList();
+ private static final List COUNT_FIELDS = fields(
+ "occurrenceCount");
+ private static final List RELATION_FIELDS = fields(
+ "relationId", "occurrenceCount");
+
+ private static final List VERTEX_LABELS = fields(
+ ENTITY, FACT_VERSION, MEMORY_EVENT, EVIDENCE, SOURCE);
+ private static final List EDGE_LABELS = fields(
+ SUBJECT,
+ OBJECT,
+ GENERATES,
+ SUPPORTED_BY,
+ FROM_SOURCE,
+ SUPERSEDES,
+ DUPLICATE_OF,
+ CONFLICTS_WITH);
+
+ private static final Comparator VERTEX_ORDER =
+ Comparator.comparing(Vertex::getId);
+ private static final Comparator EDGE_ORDER =
+ Comparator.comparing(Edge::getSrcId)
+ .thenComparing(Edge::getDstId)
+ .thenComparing(edge -> edge.getValues().toString());
+
+ public MemoryGraph toGraph(CanonicalSnapshot snapshot) {
+ Objects.requireNonNull(snapshot, "snapshot");
+
+ Map memoryEntities = new TreeMap<>();
+ Map evidenceById = new TreeMap<>();
+ Map sources = new TreeMap<>();
+ Map versionVertices = new TreeMap<>();
+ Map eventVertices = new TreeMap<>();
+ Map> edges = emptyEdgeLists();
+ Map supported = new TreeMap<>();
+ Map> relationEdges =
+ relationEdgeMaps();
+
+ for (NormalizedMemoryEvent event : snapshot.getEvents()) {
+ collectEventEntities(event, memoryEntities);
+ collectEvidence(event.getEvidence(), evidenceById, sources);
+ String eventGraphId = eventId(event.getEventId());
+ putUniqueVertex(
+ eventVertices,
+ new Vertex(MEMORY_EVENT, eventGraphId, eventValues(event)));
+ for (Evidence evidence : event.getEvidence()) {
+ addCountedEdge(
+ supported,
+ eventGraphId,
+ evidenceId(evidence.getId()));
+ }
+ }
+
+ for (Map.Entry> entry
+ : snapshot.getState().getVersionsByFactKey().entrySet()) {
+ FactKey key = entry.getKey();
+ for (MemoryFactVersion version : entry.getValue()) {
+ MemoryFact fact = version.getFact();
+ collectEntity(memoryEntities, fact.getSubject());
+ if (fact.isRelationship()) {
+ collectEntity(memoryEntities, fact.getTarget().get());
+ }
+ collectEvidence(version.getEvidence(), evidenceById, sources);
+
+ String versionGraphId = versionId(version.getId());
+ putUniqueVertex(
+ versionVertices,
+ new Vertex(
+ FACT_VERSION,
+ versionGraphId,
+ versionValues(version, key)));
+ edges.get(SUBJECT).add(new Edge(
+ SUBJECT,
+ versionGraphId,
+ entityId(fact.getSubject().getId()),
+ NO_FIELDS));
+ if (fact.isRelationship()) {
+ edges.get(OBJECT).add(new Edge(
+ OBJECT,
+ versionGraphId,
+ entityId(fact.getTarget().get().getId()),
+ NO_FIELDS));
+ }
+ String generatingEventId =
+ snapshot.getGeneratingEventIds().get(version.getId());
+ edges.get(GENERATES).add(new Edge(
+ GENERATES,
+ eventId(generatingEventId),
+ versionGraphId,
+ NO_FIELDS));
+ for (Evidence evidence : version.getEvidence()) {
+ addCountedEdge(
+ supported,
+ versionGraphId,
+ evidenceId(evidence.getId()));
+ }
+ }
+ }
+
+ for (VersionRelation relation
+ : snapshot.getState().getRelations()) {
+ String label = relationLabel(relation.getType());
+ addCountedEdge(
+ relationEdges.get(label),
+ versionId(relation.getFromVersionId()),
+ versionId(relation.getToVersionId()));
+ }
+
+ List entityVertices = new ArrayList<>();
+ for (MemoryEntity entity : memoryEntities.values()) {
+ entityVertices.add(new Vertex(
+ ENTITY,
+ entityId(entity.getId()),
+ fields(entity.getLabel())));
+ }
+ List evidenceVertices = new ArrayList<>();
+ for (Evidence evidence : evidenceById.values()) {
+ evidenceVertices.add(new Vertex(
+ EVIDENCE,
+ evidenceId(evidence.getId()),
+ fields(evidence.getContent())));
+ edges.get(FROM_SOURCE).add(new Edge(
+ FROM_SOURCE,
+ evidenceId(evidence.getId()),
+ sourceId(evidence.getSource().getId()),
+ NO_FIELDS));
+ }
+ List sourceVertices = new ArrayList<>();
+ for (Source source : sources.values()) {
+ sourceVertices.add(new Vertex(
+ SOURCE,
+ sourceId(source.getId()),
+ fields(source.getName())));
+ }
+
+ edges.put(
+ SUPPORTED_BY,
+ countedEdges(SUPPORTED_BY, supported, false));
+ for (String label : Arrays.asList(
+ SUPERSEDES, DUPLICATE_OF, CONFLICTS_WITH)) {
+ edges.put(
+ label,
+ countedEdges(label, relationEdges.get(label), true));
+ }
+
+ Map> vertices = new LinkedHashMap<>();
+ vertices.put(ENTITY, entityVertices);
+ vertices.put(
+ FACT_VERSION,
+ new ArrayList<>(versionVertices.values()));
+ vertices.put(
+ MEMORY_EVENT,
+ new ArrayList<>(eventVertices.values()));
+ vertices.put(EVIDENCE, evidenceVertices);
+ vertices.put(SOURCE, sourceVertices);
+ return createGraph(vertices, edges);
+ }
+
+ public CanonicalSnapshot fromGraph(MemoryGraph graph) {
+ Objects.requireNonNull(graph, "graph");
+ validateGraphShape(graph);
+
+ Map entityVertices =
+ vertexIndex(graph, ENTITY, ENTITY_PREFIX, ENTITY_FIELDS);
+ Map versionVertices =
+ vertexIndex(
+ graph,
+ FACT_VERSION,
+ VERSION_PREFIX,
+ VERSION_FIELDS);
+ Map eventVertices =
+ vertexIndex(
+ graph,
+ MEMORY_EVENT,
+ EVENT_PREFIX,
+ EVENT_FIELDS);
+ Map evidenceVertices =
+ vertexIndex(
+ graph,
+ EVIDENCE,
+ EVIDENCE_PREFIX,
+ EVIDENCE_FIELDS);
+ Map sourceVertices =
+ vertexIndex(graph, SOURCE, SOURCE_PREFIX, SOURCE_FIELDS);
+ Map> edges = edgeIndex(graph);
+
+ Map memoryEntities =
+ readEntities(entityVertices);
+ Map sources = readSources(sourceVertices);
+ Map evidence = readEvidence(
+ evidenceVertices,
+ sources,
+ edges.get(FROM_SOURCE));
+ Map> evidenceByOwner =
+ readSupportedEvidence(
+ edges.get(SUPPORTED_BY),
+ eventVertices,
+ versionVertices,
+ evidence);
+
+ Map events = readEvents(
+ eventVertices,
+ memoryEntities,
+ evidenceByOwner);
+ Map> versions = readVersions(
+ versionVertices,
+ memoryEntities,
+ evidenceByOwner,
+ edges.get(SUBJECT),
+ edges.get(OBJECT));
+ validateNoOrphanVertices(
+ memoryEntities,
+ evidence,
+ sources,
+ events,
+ versions);
+ Map generating = readGeneratingEvents(
+ edges.get(GENERATES),
+ events,
+ versionVertices);
+ List relations = readRelations(
+ edges,
+ versionVertices);
+
+ return new CanonicalSnapshot(
+ new TemporalState(versions, relations),
+ new ArrayList<>(events.values()),
+ generating);
+ }
+
+ private static MemoryGraph createGraph(
+ Map> vertices,
+ Map> edges) {
+ GraphSchema schema = new GraphSchema();
+ Map vertexSchemas = new LinkedHashMap<>();
+ addVertexSchema(schema, vertexSchemas, ENTITY, ENTITY_FIELDS);
+ addVertexSchema(
+ schema, vertexSchemas, FACT_VERSION, VERSION_FIELDS);
+ addVertexSchema(
+ schema, vertexSchemas, MEMORY_EVENT, EVENT_FIELDS);
+ addVertexSchema(schema, vertexSchemas, EVIDENCE, EVIDENCE_FIELDS);
+ addVertexSchema(schema, vertexSchemas, SOURCE, SOURCE_FIELDS);
+
+ Map edgeSchemas = new LinkedHashMap<>();
+ addEdgeSchema(schema, edgeSchemas, SUBJECT, NO_FIELDS);
+ addEdgeSchema(schema, edgeSchemas, OBJECT, NO_FIELDS);
+ addEdgeSchema(schema, edgeSchemas, GENERATES, NO_FIELDS);
+ addEdgeSchema(schema, edgeSchemas, SUPPORTED_BY, COUNT_FIELDS);
+ addEdgeSchema(schema, edgeSchemas, FROM_SOURCE, NO_FIELDS);
+ addEdgeSchema(schema, edgeSchemas, SUPERSEDES, RELATION_FIELDS);
+ addEdgeSchema(schema, edgeSchemas, DUPLICATE_OF, RELATION_FIELDS);
+ addEdgeSchema(
+ schema, edgeSchemas, CONFLICTS_WITH, RELATION_FIELDS);
+
+ Map groups = new LinkedHashMap<>();
+ for (String label : VERTEX_LABELS) {
+ List rows = vertices.get(label);
+ Collections.sort(rows, VERTEX_ORDER);
+ groups.put(
+ label,
+ new VertexGroup(vertexSchemas.get(label), rows));
+ }
+ for (String label : EDGE_LABELS) {
+ List rows = edges.get(label);
+ Collections.sort(rows, EDGE_ORDER);
+ groups.put(
+ label,
+ new EdgeGroup(edgeSchemas.get(label), rows));
+ }
+ return new MemoryGraph(schema, groups);
+ }
+
+ private static void addVertexSchema(
+ GraphSchema graphSchema,
+ Map schemas,
+ String label,
+ List fields) {
+ VertexSchema schema = new VertexSchema(label, "id", fields);
+ graphSchema.addVertex(schema);
+ schemas.put(label, schema);
+ }
+
+ private static void addEdgeSchema(
+ GraphSchema graphSchema,
+ Map schemas,
+ String label,
+ List fields) {
+ EdgeSchema schema = new EdgeSchema(
+ label, "srcId", "dstId", fields);
+ graphSchema.addEdge(schema);
+ schemas.put(label, schema);
+ }
+
+ private static Map> emptyEdgeLists() {
+ Map> edges = new LinkedHashMap<>();
+ for (String label : EDGE_LABELS) {
+ edges.put(label, new ArrayList<>());
+ }
+ return edges;
+ }
+
+ private static Map>
+ relationEdgeMaps() {
+ Map> relations =
+ new LinkedHashMap<>();
+ relations.put(SUPERSEDES, new TreeMap<>());
+ relations.put(DUPLICATE_OF, new TreeMap<>());
+ relations.put(CONFLICTS_WITH, new TreeMap<>());
+ return relations;
+ }
+
+ private static List countedEdges(
+ String label,
+ Map counted,
+ boolean relation) {
+ List edges = new ArrayList<>();
+ for (CountedEdge item : counted.values()) {
+ List values = relation
+ ? fields(
+ relationId(label, item.sourceId, item.targetId),
+ Integer.toString(item.count))
+ : fields(Integer.toString(item.count));
+ edges.add(new Edge(
+ label, item.sourceId, item.targetId, values));
+ }
+ return edges;
+ }
+
+ private static void addCountedEdge(
+ Map edges,
+ String sourceId,
+ String targetId) {
+ String key = tuple(sourceId, targetId);
+ CountedEdge current = edges.get(key);
+ if (current == null) {
+ edges.put(key, new CountedEdge(sourceId, targetId));
+ } else {
+ current.count++;
+ }
+ }
+
+ private static List eventValues(NormalizedMemoryEvent event) {
+ String kind = EMPTY;
+ String value = EMPTY;
+ if (event.getFactValue().isPresent()) {
+ kind = event.getFactValue().get().getKind().name();
+ value = event.getFactValue().get().getValue();
+ }
+ return fields(
+ event.getOperation().name(),
+ event.getFactId(),
+ event.getFactKey().getSubjectId(),
+ event.getFactKey().getPredicate(),
+ event.getFactKey().getScope(),
+ kind,
+ value,
+ instant(event.getValidTime().getStart()),
+ optionalInstant(event.getValidTime()),
+ instant(event.getRecordedAt()),
+ event.getPayloadHash());
+ }
+
+ private static List versionValues(
+ MemoryFactVersion version,
+ FactKey key) {
+ MemoryFact fact = version.getFact();
+ String kind = fact.isRelationship()
+ ? FactValue.Kind.ENTITY_REF.name()
+ : FactValue.Kind.LITERAL.name();
+ String literalValue = fact.getLiteralValue().orElse(EMPTY);
+ return fields(
+ fact.getId(),
+ fact.getPredicate(),
+ key.getScope(),
+ kind,
+ literalValue,
+ version.getStatus().name(),
+ instant(version.getValidTime().getStart()),
+ optionalInstant(version.getValidTime()),
+ instant(version.getTransactionTime().getStart()),
+ optionalInstant(version.getTransactionTime()));
+ }
+
+ private static void collectEventEntities(
+ NormalizedMemoryEvent event,
+ Map entities) {
+ if (!event.getEvent().getFact().isPresent()) {
+ return;
+ }
+ MemoryFact fact = event.getEvent().getFact().get();
+ collectEntity(entities, fact.getSubject());
+ if (fact.isRelationship()) {
+ collectEntity(entities, fact.getTarget().get());
+ }
+ }
+
+ private static void collectEntity(
+ Map entities,
+ MemoryEntity entity) {
+ putUnique(entities, entity.getId(), entity, "entity");
+ }
+
+ private static void collectEvidence(
+ List evidence,
+ Map evidenceById,
+ Map sources) {
+ for (Evidence item : evidence) {
+ putUnique(evidenceById, item.getId(), item, "evidence");
+ putUnique(
+ sources,
+ item.getSource().getId(),
+ item.getSource(),
+ "source");
+ }
+ }
+
+ private static void putUnique(
+ Map values,
+ String id,
+ T value,
+ String type) {
+ T previous = values.get(id);
+ if (previous != null && !previous.equals(value)) {
+ throw new IllegalArgumentException(
+ "Conflicting " + type + " id: " + id);
+ }
+ values.put(id, value);
+ }
+
+ private static void putUniqueVertex(
+ Map vertices,
+ Vertex vertex) {
+ Vertex previous = vertices.put(vertex.getId(), vertex);
+ if (previous != null
+ && !previous.getValues().equals(vertex.getValues())) {
+ throw new IllegalArgumentException(
+ "Conflicting vertex id: " + vertex.getId());
+ }
+ }
+
+ private static void validateGraphShape(MemoryGraph graph) {
+ GraphSchema schema = Objects.requireNonNull(
+ graph.getGraphSchema(), "graphSchema");
+ require(
+ schema.getVertexSchemaList().size() == VERTEX_LABELS.size(),
+ "Unexpected vertex schema count");
+ require(
+ schema.getEdgeSchemaList().size() == EDGE_LABELS.size(),
+ "Unexpected edge schema count");
+
+ for (int index = 0; index < VERTEX_LABELS.size(); index++) {
+ String label = VERTEX_LABELS.get(index);
+ VertexSchema actual = schema.getVertexSchemaList().get(index);
+ require(label.equals(actual.getLabel()),
+ "Unexpected vertex schema: " + actual.getLabel());
+ require("id".equals(actual.getIdField()),
+ "Unexpected vertex id field: " + label);
+ require(vertexFields(label).equals(actual.getFields()),
+ "Unexpected vertex fields: " + label);
+ }
+ for (int index = 0; index < EDGE_LABELS.size(); index++) {
+ String label = EDGE_LABELS.get(index);
+ EdgeSchema actual = schema.getEdgeSchemaList().get(index);
+ require(label.equals(actual.getLabel()),
+ "Unexpected edge schema: " + actual.getLabel());
+ require("srcId".equals(actual.getSrcIdField()),
+ "Unexpected edge source field: " + label);
+ require("dstId".equals(actual.getDstIdField()),
+ "Unexpected edge target field: " + label);
+ require(edgeFields(label).equals(actual.getFields()),
+ "Unexpected edge fields: " + label);
+ }
+
+ Map groups = Objects.requireNonNull(
+ graph.entities, "graph.entities");
+ List expectedGroups = new ArrayList<>(VERTEX_LABELS);
+ expectedGroups.addAll(EDGE_LABELS);
+ require(
+ expectedGroups.equals(new ArrayList<>(groups.keySet())),
+ "Unexpected graph entity groups");
+ for (String label : VERTEX_LABELS) {
+ require(groups.get(label) instanceof VertexGroup,
+ "Expected vertex group: " + label);
+ }
+ for (String label : EDGE_LABELS) {
+ require(groups.get(label) instanceof EdgeGroup,
+ "Expected edge group: " + label);
+ }
+ }
+
+ private static Map vertexIndex(
+ MemoryGraph graph,
+ String label,
+ String prefix,
+ List expectedFields) {
+ Map result = new TreeMap<>();
+ for (Vertex vertex
+ : ((VertexGroup) graph.entities.get(label)).getVertices()) {
+ require(vertex != null, "Null vertex in group: " + label);
+ require(label.equals(vertex.getLabel()),
+ "Vertex label does not match group: " + vertex.getId());
+ rawId(vertex.getId(), prefix);
+ require(vertex.getValues() != null,
+ "Null vertex values: " + vertex.getId());
+ require(vertex.getValues().size() == expectedFields.size(),
+ "Unexpected vertex value count: " + vertex.getId());
+ for (String fieldValue : vertex.getValues()) {
+ require(fieldValue != null,
+ "Null vertex field: " + vertex.getId());
+ }
+ require(result.put(vertex.getId(), vertex) == null,
+ "Duplicate vertex id: " + vertex.getId());
+ }
+ return result;
+ }
+
+ private static Map> edgeIndex(MemoryGraph graph) {
+ Map> result = new LinkedHashMap<>();
+ for (String label : EDGE_LABELS) {
+ List rows = new ArrayList<>();
+ Set identities = new HashSet<>();
+ for (Edge edge
+ : ((EdgeGroup) graph.entities.get(label)).getOutEdges()) {
+ require(edge != null, "Null edge in group: " + label);
+ require(label.equals(edge.getLabel()),
+ "Edge label does not match group: " + label);
+ require(edge.getSrcId() != null && edge.getDstId() != null,
+ "Null edge endpoint: " + label);
+ require(edge.getValues() != null,
+ "Null edge values: " + label);
+ require(edge.getValues().size() == edgeFields(label).size(),
+ "Unexpected edge value count: " + label);
+ for (String fieldValue : edge.getValues()) {
+ require(fieldValue != null,
+ "Null edge field: " + label);
+ }
+ require(identities.add(edge),
+ "Duplicate edge identity: " + edge);
+ rows.add(edge);
+ }
+ Collections.sort(rows, EDGE_ORDER);
+ result.put(label, rows);
+ }
+ return result;
+ }
+
+ private static Map readEntities(
+ Map vertices) {
+ Map result = new TreeMap<>();
+ for (Vertex vertex : vertices.values()) {
+ String rawId = rawId(vertex.getId(), ENTITY_PREFIX);
+ result.put(
+ vertex.getId(),
+ new MemoryEntity(
+ rawId,
+ value(vertex, ENTITY_FIELDS, "label")));
+ }
+ return result;
+ }
+
+ private static Map readSources(
+ Map vertices) {
+ Map result = new TreeMap<>();
+ for (Vertex vertex : vertices.values()) {
+ String rawId = rawId(vertex.getId(), SOURCE_PREFIX);
+ result.put(
+ vertex.getId(),
+ new Source(
+ rawId,
+ value(vertex, SOURCE_FIELDS, "name")));
+ }
+ return result;
+ }
+
+ private static Map readEvidence(
+ Map vertices,
+ Map sources,
+ List fromSourceEdges) {
+ Map sourceByEvidence = new HashMap<>();
+ for (Edge edge : fromSourceEdges) {
+ require(vertices.containsKey(edge.getSrcId()),
+ "Unknown evidence in from_source edge");
+ require(sources.containsKey(edge.getDstId()),
+ "Unknown source in from_source edge");
+ require(sourceByEvidence.put(
+ edge.getSrcId(), edge.getDstId()) == null,
+ "Evidence has multiple sources: " + edge.getSrcId());
+ }
+ require(sourceByEvidence.keySet().equals(vertices.keySet()),
+ "Every evidence must have exactly one source");
+
+ Map result = new TreeMap<>();
+ for (Vertex vertex : vertices.values()) {
+ String rawId = rawId(vertex.getId(), EVIDENCE_PREFIX);
+ result.put(
+ vertex.getId(),
+ new Evidence(
+ rawId,
+ sources.get(sourceByEvidence.get(vertex.getId())),
+ value(vertex, EVIDENCE_FIELDS, "content")));
+ }
+ return result;
+ }
+
+ private static Map> readSupportedEvidence(
+ List supportedEdges,
+ Map eventVertices,
+ Map versionVertices,
+ Map evidence) {
+ Map> result = new HashMap<>();
+ for (Edge edge : supportedEdges) {
+ require(
+ eventVertices.containsKey(edge.getSrcId())
+ || versionVertices.containsKey(edge.getSrcId()),
+ "Unknown supported_by owner: " + edge.getSrcId());
+ Evidence item = evidence.get(edge.getDstId());
+ require(item != null,
+ "Unknown supported_by evidence: " + edge.getDstId());
+ int count = positiveCount(edge.getValues().get(0));
+ List ownerEvidence = result.computeIfAbsent(
+ edge.getSrcId(), ignored -> new ArrayList<>());
+ for (int occurrence = 0; occurrence < count; occurrence++) {
+ ownerEvidence.add(item);
+ }
+ }
+ return result;
+ }
+
+ private static Map readEvents(
+ Map vertices,
+ Map entities,
+ Map> evidenceByOwner) {
+ Map result = new TreeMap<>();
+ EventNormalizer normalizer = new EventNormalizer();
+ for (Vertex vertex : vertices.values()) {
+ String rawEventId = rawId(vertex.getId(), EVENT_PREFIX);
+ MemoryEventOperation operation = enumValue(
+ MemoryEventOperation.class,
+ value(vertex, EVENT_FIELDS, "operation"),
+ "event operation");
+ String factId = value(vertex, EVENT_FIELDS, "factId");
+ String subjectId = value(vertex, EVENT_FIELDS, "subjectId");
+ String predicate = value(vertex, EVENT_FIELDS, "predicate");
+ FactKey key = new FactKey(
+ subjectId,
+ predicate,
+ value(vertex, EVENT_FIELDS, "scope"));
+ TimeInterval validTime = interval(
+ value(vertex, EVENT_FIELDS, "validStart"),
+ value(vertex, EVENT_FIELDS, "validEnd"));
+ Instant recordedAt = parseInstant(
+ value(vertex, EVENT_FIELDS, "recordedAt"));
+ List eventEvidence = evidenceByOwner.get(
+ vertex.getId());
+ require(eventEvidence != null && !eventEvidence.isEmpty(),
+ "Memory event must have evidence: " + rawEventId);
+
+ MemoryEvent event;
+ if (operation == MemoryEventOperation.RETRACT) {
+ require(value(vertex, EVENT_FIELDS, "valueKind").isEmpty(),
+ "Retract event must not have a value kind");
+ require(value(vertex, EVENT_FIELDS, "value").isEmpty(),
+ "Retract event must not have a value");
+ event = MemoryEvent.retract(
+ rawEventId,
+ factId,
+ validTime,
+ recordedAt,
+ eventEvidence);
+ } else {
+ MemoryFact fact = eventFact(
+ vertex,
+ factId,
+ subjectId,
+ predicate,
+ entities);
+ if (operation == MemoryEventOperation.ADD) {
+ event = MemoryEvent.add(
+ rawEventId,
+ fact,
+ validTime,
+ recordedAt,
+ eventEvidence);
+ } else {
+ require(operation == MemoryEventOperation.CORRECT,
+ "Unsupported memory event operation");
+ event = MemoryEvent.correct(
+ rawEventId,
+ fact,
+ validTime,
+ recordedAt,
+ eventEvidence);
+ }
+ }
+
+ NormalizedMemoryEvent normalized = normalizer.normalize(event, key);
+ require(eventId(normalized.getEventId()).equals(vertex.getId()),
+ "Event id is not canonical: " + rawEventId);
+ require(
+ normalized.getPayloadHash().equals(
+ value(vertex, EVENT_FIELDS, "payloadHash")),
+ "Event payload hash does not match: " + rawEventId);
+ result.put(vertex.getId(), normalized);
+ }
+ return result;
+ }
+
+ private static MemoryFact eventFact(
+ Vertex vertex,
+ String factId,
+ String subjectId,
+ String predicate,
+ Map entities) {
+ MemoryEntity subject = entities.get(entityId(subjectId));
+ require(subject != null,
+ "Unknown event subject: " + subjectId);
+ FactValue.Kind kind = enumValue(
+ FactValue.Kind.class,
+ value(vertex, EVENT_FIELDS, "valueKind"),
+ "event value kind");
+ String factValue = value(vertex, EVENT_FIELDS, "value");
+ if (kind == FactValue.Kind.LITERAL) {
+ return MemoryFact.attribute(
+ factId, subject, predicate, factValue);
+ }
+ MemoryEntity target = entities.get(entityId(factValue));
+ require(target != null,
+ "Unknown event object: " + factValue);
+ return MemoryFact.relationship(
+ factId, subject, predicate, target);
+ }
+
+ private static Map> readVersions(
+ Map vertices,
+ Map entities,
+ Map> evidenceByOwner,
+ List subjectEdges,
+ List objectEdges) {
+ Map subjects = uniqueTargets(
+ subjectEdges,
+ vertices,
+ entities,
+ "subject");
+ require(subjects.keySet().equals(vertices.keySet()),
+ "Every fact version must have exactly one subject");
+ Map objects = uniqueTargets(
+ objectEdges,
+ vertices,
+ entities,
+ "object");
+
+ Map> result = new TreeMap<>();
+ for (Vertex vertex : vertices.values()) {
+ String rawVersionId = rawId(vertex.getId(), VERSION_PREFIX);
+ MemoryEntity subject = entities.get(subjects.get(vertex.getId()));
+ String predicate = value(vertex, VERSION_FIELDS, "predicate");
+ FactKey key = new FactKey(
+ subject.getId(),
+ predicate,
+ value(vertex, VERSION_FIELDS, "scope"));
+
+ FactValue.Kind kind = enumValue(
+ FactValue.Kind.class,
+ value(vertex, VERSION_FIELDS, "valueKind"),
+ "version value kind");
+ String factId = value(vertex, VERSION_FIELDS, "factId");
+ String literalValue = value(
+ vertex, VERSION_FIELDS, "literalValue");
+ MemoryFact fact;
+ if (kind == FactValue.Kind.LITERAL) {
+ require(!objects.containsKey(vertex.getId()),
+ "Literal version must not have an object");
+ fact = MemoryFact.attribute(
+ factId, subject, predicate, literalValue);
+ } else {
+ require(literalValue.isEmpty(),
+ "Entity-reference version must not have a literal value");
+ String objectId = objects.get(vertex.getId());
+ require(objectId != null,
+ "Entity-reference version must have an object");
+ fact = MemoryFact.relationship(
+ factId, subject, predicate, entities.get(objectId));
+ }
+
+ List versionEvidence = evidenceByOwner.get(
+ vertex.getId());
+ require(versionEvidence != null && !versionEvidence.isEmpty(),
+ "Fact version must have evidence: " + rawVersionId);
+ MemoryFactVersion version = new MemoryFactVersion(
+ rawVersionId,
+ fact,
+ enumValue(
+ MemoryFactVersionStatus.class,
+ value(vertex, VERSION_FIELDS, "status"),
+ "version status"),
+ interval(
+ value(vertex, VERSION_FIELDS, "validStart"),
+ value(vertex, VERSION_FIELDS, "validEnd")),
+ interval(
+ value(vertex, VERSION_FIELDS, "transactionStart"),
+ value(vertex, VERSION_FIELDS, "transactionEnd")),
+ versionEvidence);
+ result.computeIfAbsent(
+ key, ignored -> new ArrayList<>()).add(version);
+ }
+ return result;
+ }
+
+ private static Map uniqueTargets(
+ List edges,
+ Map sourceVertices,
+ Map targets,
+ String relationName) {
+ Map result = new HashMap<>();
+ for (Edge edge : edges) {
+ require(sourceVertices.containsKey(edge.getSrcId()),
+ "Unknown " + relationName + " source");
+ require(targets.containsKey(edge.getDstId()),
+ "Unknown " + relationName + " target");
+ require(result.put(edge.getSrcId(), edge.getDstId()) == null,
+ "Multiple " + relationName + " targets");
+ }
+ return result;
+ }
+
+ private static void validateNoOrphanVertices(
+ Map entities,
+ Map evidence,
+ Map sources,
+ Map events,
+ Map> versions) {
+ Map referencedEntities = new TreeMap<>();
+ Map referencedEvidence = new TreeMap<>();
+ Map referencedSources = new TreeMap<>();
+ for (NormalizedMemoryEvent event : events.values()) {
+ collectEventEntities(event, referencedEntities);
+ collectEvidence(
+ event.getEvidence(),
+ referencedEvidence,
+ referencedSources);
+ }
+ for (List factVersions : versions.values()) {
+ for (MemoryFactVersion version : factVersions) {
+ MemoryFact fact = version.getFact();
+ collectEntity(referencedEntities, fact.getSubject());
+ if (fact.isRelationship()) {
+ collectEntity(
+ referencedEntities,
+ fact.getTarget().get());
+ }
+ collectEvidence(
+ version.getEvidence(),
+ referencedEvidence,
+ referencedSources);
+ }
+ }
+ requireAllReferenced(
+ entities, referencedEntities, ENTITY_PREFIX, "entity");
+ requireAllReferenced(
+ evidence, referencedEvidence, EVIDENCE_PREFIX, "evidence");
+ requireAllReferenced(
+ sources, referencedSources, SOURCE_PREFIX, "source");
+ }
+
+ private static void requireAllReferenced(
+ Map graphValues,
+ Map referencedValues,
+ String prefix,
+ String type) {
+ for (String graphId : graphValues.keySet()) {
+ require(referencedValues.containsKey(rawId(graphId, prefix)),
+ "Orphan " + type + " vertex: " + graphId);
+ }
+ }
+
+ private static Map readGeneratingEvents(
+ List generateEdges,
+ Map events,
+ Map versions) {
+ Map byVersion = new HashMap<>();
+ for (Edge edge : generateEdges) {
+ require(events.containsKey(edge.getSrcId()),
+ "Unknown generating event: " + edge.getSrcId());
+ require(versions.containsKey(edge.getDstId()),
+ "Unknown generated version: " + edge.getDstId());
+ require(byVersion.put(
+ edge.getDstId(), edge.getSrcId()) == null,
+ "Version has multiple generating events: "
+ + edge.getDstId());
+ }
+ require(byVersion.keySet().equals(versions.keySet()),
+ "Every fact version must have one generating event");
+
+ Map result = new TreeMap<>();
+ for (Map.Entry entry : byVersion.entrySet()) {
+ result.put(
+ rawId(entry.getKey(), VERSION_PREFIX),
+ rawId(entry.getValue(), EVENT_PREFIX));
+ }
+ return result;
+ }
+
+ private static List readRelations(
+ Map> edges,
+ Map versions) {
+ List result = new ArrayList<>();
+ for (String label : Arrays.asList(
+ SUPERSEDES, DUPLICATE_OF, CONFLICTS_WITH)) {
+ for (Edge edge : edges.get(label)) {
+ require(versions.containsKey(edge.getSrcId()),
+ "Unknown relation source: " + edge.getSrcId());
+ require(versions.containsKey(edge.getDstId()),
+ "Unknown relation target: " + edge.getDstId());
+ require(
+ relationId(label, edge.getSrcId(), edge.getDstId())
+ .equals(edge.getValues().get(0)),
+ "Version relation id does not match endpoints");
+ int count = positiveCount(edge.getValues().get(1));
+ for (int occurrence = 0; occurrence < count; occurrence++) {
+ VersionRelation relation = new VersionRelation(
+ relationType(label),
+ rawId(edge.getSrcId(), VERSION_PREFIX),
+ rawId(edge.getDstId(), VERSION_PREFIX));
+ require(
+ versionId(relation.getFromVersionId()).equals(
+ edge.getSrcId())
+ && versionId(relation.getToVersionId()).equals(
+ edge.getDstId()),
+ "Version relation endpoints are not canonical");
+ result.add(relation);
+ }
+ }
+ }
+ return result;
+ }
+
+ private static List vertexFields(String label) {
+ if (ENTITY.equals(label)) {
+ return ENTITY_FIELDS;
+ }
+ if (FACT_VERSION.equals(label)) {
+ return VERSION_FIELDS;
+ }
+ if (MEMORY_EVENT.equals(label)) {
+ return EVENT_FIELDS;
+ }
+ if (EVIDENCE.equals(label)) {
+ return EVIDENCE_FIELDS;
+ }
+ if (SOURCE.equals(label)) {
+ return SOURCE_FIELDS;
+ }
+ throw new IllegalArgumentException(
+ "Unknown vertex label: " + label);
+ }
+
+ private static List edgeFields(String label) {
+ if (SUPPORTED_BY.equals(label)) {
+ return COUNT_FIELDS;
+ }
+ if (SUPERSEDES.equals(label)
+ || DUPLICATE_OF.equals(label)
+ || CONFLICTS_WITH.equals(label)) {
+ return RELATION_FIELDS;
+ }
+ if (SUBJECT.equals(label)
+ || OBJECT.equals(label)
+ || GENERATES.equals(label)
+ || FROM_SOURCE.equals(label)) {
+ return NO_FIELDS;
+ }
+ throw new IllegalArgumentException(
+ "Unknown edge label: " + label);
+ }
+
+ private static String value(
+ Vertex vertex,
+ List fields,
+ String field) {
+ int index = fields.indexOf(field);
+ if (index < 0) {
+ throw new IllegalArgumentException("Unknown field: " + field);
+ }
+ return vertex.getValues().get(index);
+ }
+
+ private static TimeInterval interval(String start, String end) {
+ Instant parsedStart = parseInstant(start);
+ return end.isEmpty()
+ ? TimeInterval.unboundedFrom(parsedStart)
+ : new TimeInterval(parsedStart, parseInstant(end));
+ }
+
+ private static Instant parseInstant(String value) {
+ try {
+ return Instant.parse(value);
+ } catch (DateTimeParseException exception) {
+ throw new IllegalArgumentException(
+ "Invalid instant: " + value,
+ exception);
+ }
+ }
+
+ private static int positiveCount(String value) {
+ try {
+ int count = Integer.parseInt(value);
+ require(count > 0, "Occurrence count must be positive");
+ return count;
+ } catch (NumberFormatException exception) {
+ throw new IllegalArgumentException(
+ "Invalid occurrence count: " + value,
+ exception);
+ }
+ }
+
+ private static > T enumValue(
+ Class type,
+ String value,
+ String fieldName) {
+ try {
+ return Enum.valueOf(type, value);
+ } catch (IllegalArgumentException exception) {
+ throw new IllegalArgumentException(
+ "Invalid " + fieldName + ": " + value,
+ exception);
+ }
+ }
+
+ private static String relationLabel(VersionRelationType type) {
+ if (type == VersionRelationType.SUPERSEDES) {
+ return SUPERSEDES;
+ }
+ if (type == VersionRelationType.DUPLICATE_OF) {
+ return DUPLICATE_OF;
+ }
+ if (type == VersionRelationType.CONFLICTS_WITH) {
+ return CONFLICTS_WITH;
+ }
+ throw new IllegalArgumentException(
+ "Unsupported version relation type: " + type);
+ }
+
+ private static VersionRelationType relationType(String label) {
+ if (SUPERSEDES.equals(label)) {
+ return VersionRelationType.SUPERSEDES;
+ }
+ if (DUPLICATE_OF.equals(label)) {
+ return VersionRelationType.DUPLICATE_OF;
+ }
+ if (CONFLICTS_WITH.equals(label)) {
+ return VersionRelationType.CONFLICTS_WITH;
+ }
+ throw new IllegalArgumentException(
+ "Unsupported version relation label: " + label);
+ }
+
+ private static String relationId(
+ String label,
+ String sourceId,
+ String targetId) {
+ return "relation:" + tuple(label, sourceId, targetId);
+ }
+
+ private static String tuple(String... values) {
+ StringBuilder result = new StringBuilder();
+ for (String value : values) {
+ result.append(value.length()).append(':').append(value);
+ }
+ return result.toString();
+ }
+
+ private static String instant(Instant instant) {
+ return instant.toString();
+ }
+
+ private static String optionalInstant(TimeInterval interval) {
+ return interval.getEnd().isPresent()
+ ? instant(interval.getEnd().get()) : EMPTY;
+ }
+
+ private static String entityId(String rawId) {
+ return ENTITY_PREFIX + rawId;
+ }
+
+ private static String versionId(String rawId) {
+ return VERSION_PREFIX + rawId;
+ }
+
+ private static String eventId(String rawId) {
+ return EVENT_PREFIX + rawId;
+ }
+
+ private static String evidenceId(String rawId) {
+ return EVIDENCE_PREFIX + rawId;
+ }
+
+ private static String sourceId(String rawId) {
+ return SOURCE_PREFIX + rawId;
+ }
+
+ private static String rawId(String graphId, String prefix) {
+ require(graphId != null && graphId.startsWith(prefix),
+ "Graph id has the wrong prefix: " + graphId);
+ String rawId = graphId.substring(prefix.length());
+ require(!rawId.trim().isEmpty(),
+ "Graph id has an empty raw id: " + graphId);
+ return rawId;
+ }
+
+ private static List fields(String... values) {
+ return Collections.unmodifiableList(Arrays.asList(values));
+ }
+
+ private static void require(boolean condition, String message) {
+ if (!condition) {
+ throw new IllegalArgumentException(message);
+ }
+ }
+
+ private static final class CountedEdge {
+
+ private final String sourceId;
+ private final String targetId;
+ private int count;
+
+ private CountedEdge(String sourceId, String targetId) {
+ this.sourceId = sourceId;
+ this.targetId = targetId;
+ this.count = 1;
+ }
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/baseline/LwwBaseline.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/baseline/LwwBaseline.java
new file mode 100644
index 000000000..cdcfd7ef3
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/baseline/LwwBaseline.java
@@ -0,0 +1,140 @@
+/*
+ * 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.geaflow.ai.temporal.baseline;
+
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.TreeMap;
+import org.apache.geaflow.ai.temporal.model.FactKey;
+import org.apache.geaflow.ai.temporal.model.MemoryEventOperation;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersion;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersionStatus;
+import org.apache.geaflow.ai.temporal.model.TimeInterval;
+import org.apache.geaflow.ai.temporal.oracle.ReplayMethod;
+import org.apache.geaflow.ai.temporal.semantics.CanonicalSnapshot;
+import org.apache.geaflow.ai.temporal.semantics.EventLedger;
+import org.apache.geaflow.ai.temporal.semantics.EventLedgerDecision;
+import org.apache.geaflow.ai.temporal.semantics.NormalizedMemoryEvent;
+import org.apache.geaflow.ai.temporal.semantics.TemporalState;
+
+/**
+ * Replays events using deterministic last-write-wins current-state semantics.
+ */
+public final class LwwBaseline implements ReplayMethod {
+
+ private static final Comparator EVENT_ORDER =
+ Comparator.comparing(NormalizedMemoryEvent::getRecordedAt)
+ .thenComparing(NormalizedMemoryEvent::getEventId);
+
+ private static final TimeInterval ALL_VALID_TIME =
+ TimeInterval.unboundedFrom(
+ Instant.ofEpochMilli(Long.MIN_VALUE));
+
+ @Override
+ public CanonicalSnapshot replayToSnapshot(
+ List events) {
+ List ordered = new ArrayList<>(
+ Objects.requireNonNull(events, "events"));
+ for (NormalizedMemoryEvent event : ordered) {
+ Objects.requireNonNull(event, "event");
+ }
+ Collections.sort(ordered, EVENT_ORDER);
+
+ EventLedger ledger = new EventLedger();
+ List accepted = new ArrayList<>();
+ Map current = new TreeMap<>();
+ Map generatingEventIds = new TreeMap<>();
+ for (NormalizedMemoryEvent event : ordered) {
+ EventLedgerDecision decision = ledger.check(event);
+ if (decision == EventLedgerDecision.DUPLICATE_NOOP) {
+ continue;
+ }
+ if (decision == EventLedgerDecision.REJECT_EVENT_ID_REUSE) {
+ throw new IllegalArgumentException(
+ "Event id reused with a different payload: "
+ + event.getEventId());
+ }
+
+ apply(event, current, generatingEventIds);
+ ledger.commit(event);
+ accepted.add(event);
+ }
+
+ Map> versions =
+ new TreeMap<>();
+ for (Map.Entry entry
+ : current.entrySet()) {
+ versions.put(
+ entry.getKey(),
+ Collections.singletonList(entry.getValue()));
+ }
+ return new CanonicalSnapshot(
+ new TemporalState(versions, Collections.emptyList()),
+ accepted,
+ generatingEventIds);
+ }
+
+ private static void apply(
+ NormalizedMemoryEvent event,
+ Map current,
+ Map generatingEventIds) {
+ FactKey key = event.getFactKey();
+ MemoryEventOperation operation = event.getOperation();
+ if (operation == MemoryEventOperation.RETRACT) {
+ MemoryFactVersion removed = current.remove(key);
+ if (removed == null) {
+ throw new IllegalArgumentException(
+ "Cannot retract a fact without current state: "
+ + event.getFactId());
+ }
+ generatingEventIds.remove(removed.getId());
+ return;
+ }
+ if (operation == MemoryEventOperation.CORRECT
+ && !current.containsKey(key)) {
+ throw new IllegalArgumentException(
+ "Cannot correct a fact without current state: "
+ + event.getFactId());
+ }
+ if (operation != MemoryEventOperation.ADD
+ && operation != MemoryEventOperation.CORRECT) {
+ throw new UnsupportedOperationException(
+ "Unsupported memory event operation: " + operation);
+ }
+
+ MemoryFactVersion version = new MemoryFactVersion(
+ event.getEventId() + ":version:0",
+ event.getEvent().getFact().get(),
+ MemoryFactVersionStatus.ACTIVE,
+ ALL_VALID_TIME,
+ TimeInterval.unboundedFrom(event.getRecordedAt()),
+ event.getEvidence());
+ MemoryFactVersion replaced = current.put(key, version);
+ if (replaced != null) {
+ generatingEventIds.remove(replaced.getId());
+ }
+ generatingEventIds.put(version.getId(), event.getEventId());
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/baseline/SingleTimestampBaseline.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/baseline/SingleTimestampBaseline.java
new file mode 100644
index 000000000..c58d23efc
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/baseline/SingleTimestampBaseline.java
@@ -0,0 +1,169 @@
+/*
+ * 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.geaflow.ai.temporal.baseline;
+
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.TreeMap;
+import org.apache.geaflow.ai.temporal.model.FactKey;
+import org.apache.geaflow.ai.temporal.model.FactValue;
+import org.apache.geaflow.ai.temporal.model.MemoryEventOperation;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersion;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersionStatus;
+import org.apache.geaflow.ai.temporal.model.TimeInterval;
+import org.apache.geaflow.ai.temporal.model.VersionRelation;
+import org.apache.geaflow.ai.temporal.model.VersionRelationType;
+import org.apache.geaflow.ai.temporal.oracle.ReplayMethod;
+import org.apache.geaflow.ai.temporal.semantics.CanonicalSnapshot;
+import org.apache.geaflow.ai.temporal.semantics.EventLedger;
+import org.apache.geaflow.ai.temporal.semantics.EventLedgerDecision;
+import org.apache.geaflow.ai.temporal.semantics.NormalizedMemoryEvent;
+import org.apache.geaflow.ai.temporal.semantics.TemporalState;
+
+/**
+ * Keeps the current distinct values while ignoring valid-time intervals.
+ */
+public final class SingleTimestampBaseline implements ReplayMethod {
+
+ private static final Comparator EVENT_ORDER =
+ Comparator.comparing(NormalizedMemoryEvent::getRecordedAt)
+ .thenComparing(NormalizedMemoryEvent::getEventId);
+
+ private static final TimeInterval ALL_VALID_TIME =
+ TimeInterval.unboundedFrom(Instant.ofEpochMilli(Long.MIN_VALUE));
+
+ @Override
+ public CanonicalSnapshot replayToSnapshot(
+ List events) {
+ List ordered = new ArrayList<>(
+ Objects.requireNonNull(events, "events"));
+ for (NormalizedMemoryEvent event : ordered) {
+ Objects.requireNonNull(event, "event");
+ }
+ Collections.sort(ordered, EVENT_ORDER);
+
+ EventLedger ledger = new EventLedger();
+ List accepted = new ArrayList<>();
+ Map> current =
+ new TreeMap<>();
+ for (NormalizedMemoryEvent event : ordered) {
+ EventLedgerDecision decision = ledger.check(event);
+ if (decision == EventLedgerDecision.DUPLICATE_NOOP) {
+ continue;
+ }
+ if (decision == EventLedgerDecision.REJECT_EVENT_ID_REUSE) {
+ throw new IllegalArgumentException(
+ "Event id reused with a different payload: "
+ + event.getEventId());
+ }
+
+ Map next = new TreeMap<>();
+ Map existing =
+ current.get(event.getFactKey());
+ if (existing != null) {
+ next.putAll(existing);
+ }
+ apply(event, next);
+ ledger.commit(event);
+ if (next.isEmpty()) {
+ current.remove(event.getFactKey());
+ } else {
+ current.put(event.getFactKey(), next);
+ }
+ accepted.add(event);
+ }
+
+ Map> versions =
+ new TreeMap<>();
+ Map generatingEventIds = new TreeMap<>();
+ List relations = new ArrayList<>();
+ for (Map.Entry> entry
+ : current.entrySet()) {
+ List versionsForKey = new ArrayList<>();
+ for (NormalizedMemoryEvent event : entry.getValue().values()) {
+ MemoryFactVersion version = version(event);
+ versionsForKey.add(version);
+ generatingEventIds.put(
+ version.getId(), event.getEventId());
+ }
+ versions.put(entry.getKey(), versionsForKey);
+ addConflicts(versionsForKey, relations);
+ }
+ return new CanonicalSnapshot(
+ new TemporalState(versions, relations),
+ accepted,
+ generatingEventIds);
+ }
+
+ private static void apply(
+ NormalizedMemoryEvent event,
+ Map current) {
+ MemoryEventOperation operation = event.getOperation();
+ if (operation == MemoryEventOperation.ADD) {
+ current.put(event.getFactValue().get(), event);
+ return;
+ }
+ if (current.isEmpty()) {
+ throw new IllegalArgumentException(
+ "Operation requires current values for fact key: "
+ + event.getFactKey());
+ }
+ current.clear();
+ if (operation == MemoryEventOperation.CORRECT) {
+ current.put(event.getFactValue().get(), event);
+ } else if (operation != MemoryEventOperation.RETRACT) {
+ throw new UnsupportedOperationException(
+ "Unsupported memory event operation: " + operation);
+ }
+ }
+
+ private static MemoryFactVersion version(
+ NormalizedMemoryEvent event) {
+ return new MemoryFactVersion(
+ versionId(event),
+ event.getEvent().getFact().get(),
+ MemoryFactVersionStatus.ACTIVE,
+ ALL_VALID_TIME,
+ TimeInterval.unboundedFrom(event.getRecordedAt()),
+ event.getEvidence());
+ }
+
+ private static String versionId(NormalizedMemoryEvent event) {
+ return event.getEventId() + ":version:0";
+ }
+
+ private static void addConflicts(
+ List versions,
+ List relations) {
+ for (int left = 0; left < versions.size(); left++) {
+ for (int right = left + 1; right < versions.size(); right++) {
+ relations.add(new VersionRelation(
+ VersionRelationType.CONFLICTS_WITH,
+ versions.get(left).getId(),
+ versions.get(right).getId()));
+ }
+ }
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/integration/IncrementalTemporalIntegrator.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/integration/IncrementalTemporalIntegrator.java
new file mode 100644
index 000000000..5a62af740
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/integration/IncrementalTemporalIntegrator.java
@@ -0,0 +1,716 @@
+/*
+ * 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.geaflow.ai.temporal.integration;
+
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.TreeMap;
+import org.apache.geaflow.ai.temporal.model.Evidence;
+import org.apache.geaflow.ai.temporal.model.FactKey;
+import org.apache.geaflow.ai.temporal.model.FactValue;
+import org.apache.geaflow.ai.temporal.model.MemoryEvent;
+import org.apache.geaflow.ai.temporal.model.MemoryEventOperation;
+import org.apache.geaflow.ai.temporal.model.MemoryFact;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersion;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersionStatus;
+import org.apache.geaflow.ai.temporal.model.TimeInterval;
+import org.apache.geaflow.ai.temporal.model.VersionRelation;
+import org.apache.geaflow.ai.temporal.model.VersionRelationType;
+import org.apache.geaflow.ai.temporal.semantics.EventLedger;
+import org.apache.geaflow.ai.temporal.semantics.EventLedgerDecision;
+import org.apache.geaflow.ai.temporal.semantics.NormalizedMemoryEvent;
+import org.apache.geaflow.ai.temporal.semantics.TemporalState;
+
+/**
+ * Incrementally integrates legacy events by fact id and normalized events by
+ * fact key.
+ */
+public final class IncrementalTemporalIntegrator {
+
+ private static final Comparator EVENT_ORDER =
+ Comparator.comparing(MemoryEvent::getTransactionTime)
+ .thenComparing(MemoryEvent::getId);
+
+ private static final Comparator
+ NORMALIZED_EVENT_ORDER =
+ Comparator.comparing(NormalizedMemoryEvent::getRecordedAt)
+ .thenComparing(NormalizedMemoryEvent::getEventId);
+
+ private static final Comparator VERSION_ORDER =
+ Comparator.comparing(
+ (MemoryFactVersion version) ->
+ version.getTransactionTime().getStart())
+ .thenComparing(
+ version -> version.getValidTime().getStart())
+ .thenComparing(MemoryFactVersion::getId);
+
+ private static final Comparator VALID_TIME_ORDER =
+ Comparator.comparing(
+ (MemoryFactVersion version) ->
+ version.getValidTime().getStart())
+ .thenComparing(MemoryFactVersion::getId);
+
+ private static final Comparator EVIDENCE_ORDER =
+ Comparator.comparing(Evidence::getId)
+ .thenComparing(evidence -> evidence.getSource().getId())
+ .thenComparing(evidence -> evidence.getSource().getName())
+ .thenComparing(Evidence::getContent);
+
+ private final Map eventsById = new HashMap<>();
+ private final Map> eventsByFactId =
+ new HashMap<>();
+ private final Map> versionsByFactId =
+ new HashMap<>();
+
+ private final EventLedger normalizedEventLedger = new EventLedger();
+ private final Map>
+ normalizedEventsByFactKey = new HashMap<>();
+ private final Map>
+ normalizedVersionsByFactKey = new HashMap<>();
+ private final Map>
+ normalizedRelationsByFactKey = new HashMap<>();
+
+ public void apply(MemoryEvent event) {
+ Objects.requireNonNull(event, "event");
+
+ MemoryEvent existing = eventsById.get(event.getId());
+ if (existing != null) {
+ if (!existing.equals(event)) {
+ throw new IllegalArgumentException(
+ "Conflicting event id: " + event.getId());
+ }
+ return;
+ }
+
+ String factId = event.getFactId();
+ List existingEvents = eventsByFactId.get(factId);
+ List updatedEvents = existingEvents == null
+ ? new ArrayList<>() : new ArrayList<>(existingEvents);
+ boolean appended = existingEvents == null
+ || EVENT_ORDER.compare(
+ existingEvents.get(existingEvents.size() - 1),
+ event) < 0;
+
+ updatedEvents.add(event);
+ Collections.sort(updatedEvents, EVENT_ORDER);
+
+ List updatedVersions;
+ if (appended) {
+ List existingVersions =
+ versionsByFactId.get(factId);
+ updatedVersions = existingVersions == null
+ ? new ArrayList<>()
+ : new ArrayList<>(existingVersions);
+ applyOrderedEvent(event, updatedVersions);
+ } else {
+ updatedVersions = replayFact(updatedEvents);
+ }
+
+ Collections.sort(updatedVersions, VERSION_ORDER);
+ eventsById.put(event.getId(), event);
+ eventsByFactId.put(
+ factId,
+ Collections.unmodifiableList(updatedEvents));
+ versionsByFactId.put(
+ factId,
+ Collections.unmodifiableList(updatedVersions));
+ }
+
+ /**
+ * Applies an already normalized event to the state for its fact key.
+ */
+ public void apply(NormalizedMemoryEvent event) {
+ Objects.requireNonNull(event, "event");
+
+ EventLedgerDecision decision = normalizedEventLedger.check(event);
+ if (decision == EventLedgerDecision.DUPLICATE_NOOP) {
+ return;
+ }
+ if (decision == EventLedgerDecision.REJECT_EVENT_ID_REUSE) {
+ throw new IllegalArgumentException(
+ "Event id reused with a different payload: "
+ + event.getEventId());
+ }
+
+ FactKey factKey = event.getFactKey();
+ List existingEvents =
+ normalizedEventsByFactKey.get(factKey);
+ List updatedEvents =
+ existingEvents == null
+ ? new ArrayList<>()
+ : new ArrayList<>(existingEvents);
+ boolean appended = existingEvents == null
+ || NORMALIZED_EVENT_ORDER.compare(
+ existingEvents.get(existingEvents.size() - 1),
+ event) < 0;
+ updatedEvents.add(event);
+ Collections.sort(updatedEvents, NORMALIZED_EVENT_ORDER);
+
+ List updatedVersions;
+ List updatedRelations;
+ if (appended) {
+ List existingVersions =
+ normalizedVersionsByFactKey.get(factKey);
+ List existingRelations =
+ normalizedRelationsByFactKey.get(factKey);
+ updatedVersions = existingVersions == null
+ ? new ArrayList<>()
+ : new ArrayList<>(existingVersions);
+ updatedRelations = existingRelations == null
+ ? new ArrayList<>()
+ : new ArrayList<>(existingRelations);
+ applyNormalizedOrderedEvent(
+ event,
+ updatedVersions,
+ updatedRelations);
+ } else {
+ updatedVersions = new ArrayList<>();
+ updatedRelations = new ArrayList<>();
+ for (NormalizedMemoryEvent orderedEvent : updatedEvents) {
+ applyNormalizedOrderedEvent(
+ orderedEvent,
+ updatedVersions,
+ updatedRelations);
+ }
+ }
+
+ Collections.sort(updatedVersions, VERSION_ORDER);
+ Collections.sort(updatedRelations);
+ normalizedEventLedger.commit(event);
+ normalizedEventsByFactKey.put(
+ factKey,
+ Collections.unmodifiableList(updatedEvents));
+ normalizedVersionsByFactKey.put(
+ factKey,
+ Collections.unmodifiableList(updatedVersions));
+ normalizedRelationsByFactKey.put(
+ factKey,
+ Collections.unmodifiableList(updatedRelations));
+ }
+
+ public List snapshot() {
+ List snapshot = new ArrayList<>();
+ for (List versions :
+ versionsByFactId.values()) {
+ snapshot.addAll(versions);
+ }
+
+ Collections.sort(snapshot, VERSION_ORDER);
+ return Collections.unmodifiableList(snapshot);
+ }
+
+ /**
+ * Returns an immutable, deterministic snapshot of normalized state.
+ */
+ public TemporalState stateSnapshot() {
+ Map> versionsByFactKey =
+ new TreeMap<>();
+ versionsByFactKey.putAll(normalizedVersionsByFactKey);
+
+ List relations = new ArrayList<>();
+ for (List keyRelations :
+ normalizedRelationsByFactKey.values()) {
+ relations.addAll(keyRelations);
+ }
+ Collections.sort(relations);
+ return new TemporalState(versionsByFactKey, relations);
+ }
+
+ List eventSnapshot() {
+ List snapshot =
+ new ArrayList<>(eventsById.values());
+ Collections.sort(snapshot, EVENT_ORDER);
+ return Collections.unmodifiableList(snapshot);
+ }
+
+ private static void applyNormalizedOrderedEvent(
+ NormalizedMemoryEvent event,
+ List versions,
+ List relations) {
+ MemoryEventOperation operation = event.getOperation();
+ if (operation == MemoryEventOperation.ADD) {
+ applyNormalizedAdd(event, versions, relations);
+ } else if (operation == MemoryEventOperation.CORRECT
+ || operation == MemoryEventOperation.RETRACT) {
+ applyNormalizedChange(event, versions, relations);
+ } else {
+ throw new UnsupportedOperationException(
+ "Unsupported memory event operation: "
+ + operation);
+ }
+ }
+
+ private static void applyNormalizedAdd(
+ NormalizedMemoryEvent event,
+ List versions,
+ List relations) {
+ List duplicates = new ArrayList<>();
+ for (MemoryFactVersion version : versions) {
+ if (isCurrentActive(version)
+ && version.getValidTime().overlaps(
+ event.getValidTime())
+ && hasSameValue(event, version)) {
+ duplicates.add(version);
+ }
+ }
+ Collections.sort(duplicates, VALID_TIME_ORDER);
+
+ TimeInterval validTime = event.getValidTime();
+ List evidence = new ArrayList<>();
+ mergeEvidence(evidence, event.getEvidence());
+ Map materializedDuplicates = new HashMap<>();
+ for (MemoryFactVersion duplicate : duplicates) {
+ versions.remove(duplicate);
+ boolean materialized = materializeClosedVersion(
+ duplicate,
+ event.getRecordedAt(),
+ versions,
+ relations);
+ materializedDuplicates.put(duplicate.getId(), materialized);
+ validTime = span(validTime, duplicate.getValidTime());
+ mergeEvidence(evidence, duplicate.getEvidence());
+ }
+
+ MemoryFactVersion added = new MemoryFactVersion(
+ event.getEventId() + ":version:0",
+ event.getEvent().getFact().get(),
+ MemoryFactVersionStatus.ACTIVE,
+ validTime,
+ TimeInterval.unboundedFrom(event.getRecordedAt()),
+ evidence);
+ versions.add(added);
+
+ for (MemoryFactVersion duplicate : duplicates) {
+ if (Boolean.TRUE.equals(
+ materializedDuplicates.get(duplicate.getId()))) {
+ addRelation(
+ relations,
+ new VersionRelation(
+ VersionRelationType.DUPLICATE_OF,
+ added.getId(),
+ duplicate.getId()));
+ }
+ }
+
+ for (MemoryFactVersion version : versions) {
+ if (version != added
+ && isCurrentActive(version)
+ && version.getValidTime().overlaps(validTime)
+ && !hasSameValue(event, version)) {
+ addRelation(
+ relations,
+ new VersionRelation(
+ VersionRelationType.CONFLICTS_WITH,
+ added.getId(),
+ version.getId()));
+ }
+ }
+ }
+
+ private static void applyNormalizedChange(
+ NormalizedMemoryEvent event,
+ List versions,
+ List relations) {
+ List affected = new ArrayList<>();
+ for (MemoryFactVersion version : versions) {
+ if (isCurrentActive(version)
+ && version.getValidTime().overlaps(
+ event.getValidTime())) {
+ affected.add(version);
+ }
+ }
+ Collections.sort(affected, VALID_TIME_ORDER);
+ if (!isFullyCovered(event.getValidTime(), affected)) {
+ throw new IllegalArgumentException(
+ "Event interval is not fully covered for fact key: "
+ + event.getFactKey());
+ }
+
+ Map materializedSources = new HashMap<>();
+ for (MemoryFactVersion version : affected) {
+ versions.remove(version);
+ materializedSources.put(
+ version.getId(),
+ materializeClosedVersion(
+ version,
+ event.getRecordedAt(),
+ versions,
+ relations));
+ }
+
+ int versionIndex;
+ if (event.getOperation() == MemoryEventOperation.CORRECT) {
+ MemoryFactVersion corrected = new MemoryFactVersion(
+ event.getEventId() + ":version:0",
+ event.getEvent().getFact().get(),
+ MemoryFactVersionStatus.ACTIVE,
+ event.getValidTime(),
+ TimeInterval.unboundedFrom(event.getRecordedAt()),
+ event.getEvidence());
+ versions.add(corrected);
+ for (MemoryFactVersion source : affected) {
+ addSupersedesIfMaterialized(
+ corrected,
+ source,
+ materializedSources,
+ relations);
+ }
+ versionIndex = 1;
+ } else {
+ versionIndex = addTombstones(
+ event,
+ affected,
+ materializedSources,
+ versions,
+ relations);
+ }
+
+ for (MemoryFactVersion source : affected) {
+ for (TimeInterval remaining : source.getValidTime()
+ .subtract(event.getValidTime())) {
+ MemoryFactVersion residue = new MemoryFactVersion(
+ event.getEventId() + ":version:"
+ + versionIndex++,
+ source.getFact(),
+ MemoryFactVersionStatus.ACTIVE,
+ remaining,
+ TimeInterval.unboundedFrom(
+ event.getRecordedAt()),
+ source.getEvidence());
+ versions.add(residue);
+ addSupersedesIfMaterialized(
+ residue,
+ source,
+ materializedSources,
+ relations);
+ }
+ }
+
+ addCurrentConflicts(versions, relations);
+ }
+
+ private static void addCurrentConflicts(
+ List versions,
+ List relations) {
+ for (int leftIndex = 0;
+ leftIndex < versions.size(); leftIndex++) {
+ MemoryFactVersion left = versions.get(leftIndex);
+ if (!isCurrentActive(left)) {
+ continue;
+ }
+ for (int rightIndex = leftIndex + 1;
+ rightIndex < versions.size(); rightIndex++) {
+ MemoryFactVersion right = versions.get(rightIndex);
+ if (isCurrentActive(right)
+ && left.getValidTime().overlaps(
+ right.getValidTime())
+ && !factValue(left.getFact()).equals(
+ factValue(right.getFact()))) {
+ addRelation(
+ relations,
+ new VersionRelation(
+ VersionRelationType.CONFLICTS_WITH,
+ left.getId(),
+ right.getId()));
+ }
+ }
+ }
+ }
+
+ private static int addTombstones(
+ NormalizedMemoryEvent event,
+ List affected,
+ Map materializedSources,
+ List versions,
+ List relations) {
+ int versionIndex = 0;
+ for (MemoryFactVersion source : affected) {
+ TimeInterval overlap = source.getValidTime()
+ .intersection(event.getValidTime()).get();
+ MemoryFactVersion tombstone = new MemoryFactVersion(
+ event.getEventId() + ":version:" + versionIndex++,
+ source.getFact(),
+ MemoryFactVersionStatus.RETRACTED,
+ overlap,
+ TimeInterval.unboundedFrom(event.getRecordedAt()),
+ event.getEvidence());
+ versions.add(tombstone);
+ addSupersedesIfMaterialized(
+ tombstone,
+ source,
+ materializedSources,
+ relations);
+ }
+ return versionIndex;
+ }
+
+ private static boolean materializeClosedVersion(
+ MemoryFactVersion version,
+ Instant recordedAt,
+ List versions,
+ List relations) {
+ if (!version.getTransactionTime().getStart()
+ .isBefore(recordedAt)) {
+ removeRelationsFor(version.getId(), relations);
+ return false;
+ }
+
+ versions.add(new MemoryFactVersion(
+ version.getId(),
+ version.getFact(),
+ version.getStatus(),
+ version.getValidTime(),
+ new TimeInterval(
+ version.getTransactionTime().getStart(),
+ recordedAt),
+ version.getEvidence()));
+ return true;
+ }
+
+ private static void addSupersedesIfMaterialized(
+ MemoryFactVersion replacement,
+ MemoryFactVersion source,
+ Map materializedSources,
+ List relations) {
+ if (Boolean.TRUE.equals(
+ materializedSources.get(source.getId()))) {
+ addRelation(
+ relations,
+ new VersionRelation(
+ VersionRelationType.SUPERSEDES,
+ replacement.getId(),
+ source.getId()));
+ }
+ }
+
+ private static void addRelation(
+ List relations,
+ VersionRelation relation) {
+ if (!relations.contains(relation)) {
+ relations.add(relation);
+ }
+ }
+
+ private static void removeRelationsFor(
+ String versionId,
+ List relations) {
+ for (int index = relations.size() - 1; index >= 0; index--) {
+ VersionRelation relation = relations.get(index);
+ if (relation.getFromVersionId().equals(versionId)
+ || relation.getToVersionId().equals(versionId)) {
+ relations.remove(index);
+ }
+ }
+ }
+
+ private static boolean hasSameValue(
+ NormalizedMemoryEvent event,
+ MemoryFactVersion version) {
+ Optional eventValue = event.getFactValue();
+ return eventValue.isPresent()
+ && eventValue.get().equals(factValue(version.getFact()));
+ }
+
+ private static FactValue factValue(MemoryFact fact) {
+ if (fact.isRelationship()) {
+ return FactValue.entityReference(
+ fact.getTarget().get().getId());
+ }
+ return FactValue.literal(fact.getLiteralValue().get());
+ }
+
+ private static TimeInterval span(
+ TimeInterval left,
+ TimeInterval right) {
+ Instant start = left.getStart().isBefore(right.getStart())
+ ? left.getStart() : right.getStart();
+ Instant end;
+ if (!left.getEnd().isPresent()
+ || !right.getEnd().isPresent()) {
+ end = null;
+ } else {
+ Instant leftEnd = left.getEnd().get();
+ Instant rightEnd = right.getEnd().get();
+ end = leftEnd.isAfter(rightEnd) ? leftEnd : rightEnd;
+ }
+ return new TimeInterval(start, end);
+ }
+
+ private static void mergeEvidence(
+ List target,
+ List additions) {
+ for (Evidence evidence : additions) {
+ if (!target.contains(evidence)) {
+ target.add(evidence);
+ }
+ }
+ Collections.sort(target, EVIDENCE_ORDER);
+ }
+
+ private static boolean isCurrentActive(
+ MemoryFactVersion version) {
+ return version.getStatus() == MemoryFactVersionStatus.ACTIVE
+ && isCurrent(version);
+ }
+
+ private static List replayFact(
+ List events) {
+ List versions = new ArrayList<>();
+ for (MemoryEvent event : events) {
+ applyOrderedEvent(event, versions);
+ }
+ return versions;
+ }
+
+ private static void applyOrderedEvent(
+ MemoryEvent event,
+ List versions) {
+ MemoryEventOperation operation = event.getOperation();
+ if (operation == MemoryEventOperation.ADD) {
+ applyAdd(event, versions);
+ } else if (operation == MemoryEventOperation.CORRECT
+ || operation == MemoryEventOperation.RETRACT) {
+ applyChange(event, versions);
+ } else {
+ throw new UnsupportedOperationException(
+ "Unsupported memory event operation: "
+ + operation);
+ }
+ }
+
+ private static void applyAdd(
+ MemoryEvent event,
+ List versions) {
+ for (MemoryFactVersion version : versions) {
+ if (isCurrent(version)
+ && version.getFact().getId().equals(event.getFactId())
+ && version.getValidTime().overlaps(
+ event.getValidTime())) {
+ throw new IllegalArgumentException(
+ "Overlapping add for fact id: "
+ + event.getFactId());
+ }
+ }
+
+ versions.add(new MemoryFactVersion(
+ event.getId() + ":version:0",
+ event.getFact().get(),
+ event.getValidTime(),
+ TimeInterval.unboundedFrom(
+ event.getTransactionTime()),
+ event.getEvidence()));
+ }
+
+ private static void applyChange(
+ MemoryEvent event,
+ List versions) {
+ List affected = new ArrayList<>();
+ for (MemoryFactVersion version : versions) {
+ if (isCurrent(version)
+ && version.getFact().getId().equals(event.getFactId())
+ && version.getValidTime().overlaps(
+ event.getValidTime())) {
+ affected.add(version);
+ }
+ }
+
+ Collections.sort(affected, VALID_TIME_ORDER);
+ if (!isFullyCovered(event.getValidTime(), affected)) {
+ throw new IllegalArgumentException(
+ "Event interval is not fully covered for fact id: "
+ + event.getFactId());
+ }
+
+ versions.removeAll(affected);
+
+ int fragmentIndex = 1;
+ for (MemoryFactVersion version : affected) {
+ if (version.getTransactionTime().getStart()
+ .isBefore(event.getTransactionTime())) {
+ versions.add(new MemoryFactVersion(
+ version.getId(),
+ version.getFact(),
+ version.getValidTime(),
+ new TimeInterval(
+ version.getTransactionTime().getStart(),
+ event.getTransactionTime()),
+ version.getEvidence()));
+ }
+
+ for (TimeInterval remaining :
+ version.getValidTime().subtract(
+ event.getValidTime())) {
+ versions.add(new MemoryFactVersion(
+ event.getId() + ":version:"
+ + fragmentIndex++,
+ version.getFact(),
+ remaining,
+ TimeInterval.unboundedFrom(
+ event.getTransactionTime()),
+ version.getEvidence()));
+ }
+ }
+
+ if (event.getOperation() == MemoryEventOperation.CORRECT) {
+ versions.add(new MemoryFactVersion(
+ event.getId() + ":version:0",
+ event.getFact().get(),
+ event.getValidTime(),
+ TimeInterval.unboundedFrom(
+ event.getTransactionTime()),
+ event.getEvidence()));
+ }
+ }
+
+ private static boolean isFullyCovered(
+ TimeInterval target,
+ List coveringVersions) {
+ List uncovered = new ArrayList<>();
+ uncovered.add(target);
+
+ for (MemoryFactVersion version : coveringVersions) {
+ List remaining = new ArrayList<>();
+ for (TimeInterval interval : uncovered) {
+ remaining.addAll(
+ interval.subtract(version.getValidTime()));
+ }
+
+ uncovered = remaining;
+ if (uncovered.isEmpty()) {
+ return true;
+ }
+ }
+
+ return false;
+ }
+
+ private static boolean isCurrent(
+ MemoryFactVersion version) {
+ return !version.getTransactionTime()
+ .getEnd().isPresent();
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/integration/TemporalEventAggregateFunction.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/integration/TemporalEventAggregateFunction.java
new file mode 100644
index 000000000..dc37bda5d
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/integration/TemporalEventAggregateFunction.java
@@ -0,0 +1,73 @@
+/*
+ * 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.geaflow.ai.temporal.integration;
+
+import java.util.List;
+import java.util.Objects;
+import org.apache.geaflow.ai.temporal.model.MemoryEvent;
+import org.apache.geaflow.ai.temporal.model.MemoryFactVersion;
+import org.apache.geaflow.api.function.base.AggregateFunction;
+
+/**
+ * Adapts temporal event integration to GeaFlow keyed aggregation.
+ */
+public final class TemporalEventAggregateFunction implements
+ AggregateFunction> {
+
+ @Override
+ public IncrementalTemporalIntegrator createAccumulator() {
+ return new IncrementalTemporalIntegrator();
+ }
+
+ @Override
+ public void add(
+ MemoryEvent value,
+ IncrementalTemporalIntegrator accumulator) {
+ Objects.requireNonNull(accumulator, "accumulator")
+ .apply(value);
+ }
+
+ @Override
+ public List getResult(
+ IncrementalTemporalIntegrator accumulator) {
+ return Objects.requireNonNull(
+ accumulator,
+ "accumulator").snapshot();
+ }
+
+ @Override
+ public IncrementalTemporalIntegrator merge(
+ IncrementalTemporalIntegrator left,
+ IncrementalTemporalIntegrator right) {
+ Objects.requireNonNull(left, "left");
+ Objects.requireNonNull(right, "right");
+
+ IncrementalTemporalIntegrator merged =
+ new IncrementalTemporalIntegrator();
+ for (MemoryEvent event : left.eventSnapshot()) {
+ merged.apply(event);
+ }
+ for (MemoryEvent event : right.eventSnapshot()) {
+ merged.apply(event);
+ }
+ return merged;
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/Evidence.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/Evidence.java
new file mode 100644
index 000000000..998137fda
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/Evidence.java
@@ -0,0 +1,78 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+import java.util.Objects;
+
+/**
+ * An immutable piece of evidence and its source.
+ */
+public final class Evidence {
+
+ private final String id;
+ private final Source source;
+ private final String content;
+
+ public Evidence(String id, Source source, String content) {
+ this.id = requireText(id, "id");
+ this.source = Objects.requireNonNull(source, "source");
+ this.content = requireText(content, "content");
+ }
+
+ public String getId() {
+ return id;
+ }
+
+ public Source getSource() {
+ return source;
+ }
+
+ public String getContent() {
+ return content;
+ }
+
+ private static String requireText(String value, String fieldName) {
+ Objects.requireNonNull(value, fieldName);
+ if (value.trim().isEmpty()) {
+ throw new IllegalArgumentException(
+ "Evidence " + fieldName + " must not be blank");
+ }
+ return value;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (!(object instanceof Evidence)) {
+ return false;
+ }
+ Evidence that = (Evidence) object;
+ return id.equals(that.id)
+ && source.equals(that.source)
+ && content.equals(that.content);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(id, source, content);
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/FactKey.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/FactKey.java
new file mode 100644
index 000000000..b74621b37
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/FactKey.java
@@ -0,0 +1,96 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+import java.util.Objects;
+
+/**
+ * An immutable identity for a fact independent of its value.
+ */
+public final class FactKey implements Comparable {
+
+ private final String subjectId;
+ private final String predicate;
+ private final String scope;
+
+ public FactKey(
+ String subjectId,
+ String predicate,
+ String scope) {
+ this.subjectId = requireText(subjectId, "subjectId");
+ this.predicate = requireText(predicate, "predicate");
+ this.scope = requireText(scope, "scope");
+ }
+
+ public String getSubjectId() {
+ return subjectId;
+ }
+
+ public String getPredicate() {
+ return predicate;
+ }
+
+ public String getScope() {
+ return scope;
+ }
+
+ @Override
+ public int compareTo(FactKey that) {
+ int comparison = subjectId.compareTo(that.subjectId);
+ if (comparison != 0) {
+ return comparison;
+ }
+ comparison = predicate.compareTo(that.predicate);
+ if (comparison != 0) {
+ return comparison;
+ }
+ return scope.compareTo(that.scope);
+ }
+
+ private static String requireText(
+ String value,
+ String fieldName) {
+ Objects.requireNonNull(value, fieldName);
+ if (value.trim().isEmpty()) {
+ throw new IllegalArgumentException(
+ "Fact key " + fieldName + " must not be blank");
+ }
+ return value;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (!(object instanceof FactKey)) {
+ return false;
+ }
+ FactKey that = (FactKey) object;
+ return subjectId.equals(that.subjectId)
+ && predicate.equals(that.predicate)
+ && scope.equals(that.scope);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(subjectId, predicate, scope);
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/FactValue.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/FactValue.java
new file mode 100644
index 000000000..5905f82ea
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/FactValue.java
@@ -0,0 +1,101 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+import java.util.Objects;
+import java.util.Optional;
+
+/**
+ * An immutable literal or entity-reference fact value.
+ */
+public final class FactValue implements Comparable {
+
+ public enum Kind {
+ LITERAL,
+ ENTITY_REF
+ }
+
+ private final Kind kind;
+ private final String value;
+
+ private FactValue(Kind kind, String value) {
+ this.kind = Objects.requireNonNull(kind, "kind");
+ this.value = requireText(value);
+ }
+
+ public static FactValue literal(String value) {
+ return new FactValue(Kind.LITERAL, value);
+ }
+
+ public static FactValue entityReference(String entityId) {
+ return new FactValue(Kind.ENTITY_REF, entityId);
+ }
+
+ public Kind getKind() {
+ return kind;
+ }
+
+ public String getValue() {
+ return value;
+ }
+
+ public Optional getLiteralValue() {
+ return kind == Kind.LITERAL
+ ? Optional.of(value) : Optional.empty();
+ }
+
+ public Optional getEntityId() {
+ return kind == Kind.ENTITY_REF
+ ? Optional.of(value) : Optional.empty();
+ }
+
+ @Override
+ public int compareTo(FactValue that) {
+ int comparison = kind.compareTo(that.kind);
+ return comparison != 0
+ ? comparison : value.compareTo(that.value);
+ }
+
+ private static String requireText(String value) {
+ Objects.requireNonNull(value, "value");
+ if (value.trim().isEmpty()) {
+ throw new IllegalArgumentException(
+ "Fact value must not be blank");
+ }
+ return value;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (!(object instanceof FactValue)) {
+ return false;
+ }
+ FactValue that = (FactValue) object;
+ return kind == that.kind && value.equals(that.value);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(kind, value);
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEntity.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEntity.java
new file mode 100644
index 000000000..87ce52a8c
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEntity.java
@@ -0,0 +1,70 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+import java.util.Objects;
+
+/**
+ * An immutable identity anchor for temporal memory facts.
+ */
+public final class MemoryEntity {
+
+ private final String id;
+ private final String label;
+
+ public MemoryEntity(String id, String label) {
+ this.id = requireText(id, "id");
+ this.label = requireText(label, "label");
+ }
+
+ public String getId() {
+ return id;
+ }
+
+ public String getLabel() {
+ return label;
+ }
+
+ private static String requireText(String value, String fieldName) {
+ Objects.requireNonNull(value, fieldName);
+ if (value.trim().isEmpty()) {
+ throw new IllegalArgumentException(
+ "Memory entity " + fieldName + " must not be blank");
+ }
+ return value;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (!(object instanceof MemoryEntity)) {
+ return false;
+ }
+ MemoryEntity that = (MemoryEntity) object;
+ return id.equals(that.id) && label.equals(that.label);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(id, label);
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEvent.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEvent.java
new file mode 100644
index 000000000..122cee8a3
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEvent.java
@@ -0,0 +1,194 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+
+/**
+ * An immutable input event for temporal memory replay.
+ */
+public final class MemoryEvent {
+
+ private final String id;
+ private final MemoryEventOperation operation;
+ private final String factId;
+ private final MemoryFact fact;
+ private final TimeInterval validTime;
+ private final Instant transactionTime;
+ private final List evidence;
+
+ private MemoryEvent(
+ String id,
+ MemoryEventOperation operation,
+ String factId,
+ MemoryFact fact,
+ TimeInterval validTime,
+ Instant transactionTime,
+ List evidence) {
+ this.id = requireText(id, "id");
+ this.operation =
+ Objects.requireNonNull(operation, "operation");
+ this.factId = requireText(factId, "fact id");
+ this.fact = fact;
+ this.validTime =
+ Objects.requireNonNull(validTime, "validTime");
+ this.transactionTime =
+ Objects.requireNonNull(transactionTime, "transactionTime");
+ this.evidence = copyEvidence(evidence);
+ }
+
+ public static MemoryEvent add(
+ String id,
+ MemoryFact fact,
+ TimeInterval validTime,
+ Instant transactionTime,
+ List evidence) {
+ Objects.requireNonNull(fact, "fact");
+ return new MemoryEvent(
+ id,
+ MemoryEventOperation.ADD,
+ fact.getId(),
+ fact,
+ validTime,
+ transactionTime,
+ evidence);
+ }
+
+ public static MemoryEvent correct(
+ String id,
+ MemoryFact fact,
+ TimeInterval validTime,
+ Instant transactionTime,
+ List evidence) {
+ Objects.requireNonNull(fact, "fact");
+ return new MemoryEvent(
+ id,
+ MemoryEventOperation.CORRECT,
+ fact.getId(),
+ fact,
+ validTime,
+ transactionTime,
+ evidence);
+ }
+
+ public static MemoryEvent retract(
+ String id,
+ String factId,
+ TimeInterval validTime,
+ Instant transactionTime,
+ List evidence) {
+ return new MemoryEvent(
+ id,
+ MemoryEventOperation.RETRACT,
+ factId,
+ null,
+ validTime,
+ transactionTime,
+ evidence);
+ }
+
+ public String getId() {
+ return id;
+ }
+
+ public MemoryEventOperation getOperation() {
+ return operation;
+ }
+
+ public String getFactId() {
+ return factId;
+ }
+
+ public Optional getFact() {
+ return Optional.ofNullable(fact);
+ }
+
+ public TimeInterval getValidTime() {
+ return validTime;
+ }
+
+ public Instant getTransactionTime() {
+ return transactionTime;
+ }
+
+ public List getEvidence() {
+ return evidence;
+ }
+
+ private static List copyEvidence(
+ List evidence) {
+ List copy =
+ new ArrayList<>(Objects.requireNonNull(evidence, "evidence"));
+ if (copy.isEmpty()) {
+ throw new IllegalArgumentException(
+ "Memory event evidence must not be empty");
+ }
+ for (Evidence item : copy) {
+ Objects.requireNonNull(item, "evidence item");
+ }
+ return Collections.unmodifiableList(copy);
+ }
+
+ private static String requireText(
+ String value,
+ String fieldName) {
+ Objects.requireNonNull(value, fieldName);
+ if (value.trim().isEmpty()) {
+ throw new IllegalArgumentException(
+ "Memory event " + fieldName + " must not be blank");
+ }
+ return value;
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ if (this == object) {
+ return true;
+ }
+ if (!(object instanceof MemoryEvent)) {
+ return false;
+ }
+ MemoryEvent that = (MemoryEvent) object;
+ return id.equals(that.id)
+ && operation == that.operation
+ && factId.equals(that.factId)
+ && Objects.equals(fact, that.fact)
+ && validTime.equals(that.validTime)
+ && transactionTime.equals(that.transactionTime)
+ && evidence.equals(that.evidence);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(
+ id,
+ operation,
+ factId,
+ fact,
+ validTime,
+ transactionTime,
+ evidence);
+ }
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEventOperation.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEventOperation.java
new file mode 100644
index 000000000..7e5d63460
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryEventOperation.java
@@ -0,0 +1,30 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+/**
+ * Supported operations for temporal memory events.
+ */
+public enum MemoryEventOperation {
+
+ ADD,
+ CORRECT,
+ RETRACT
+}
diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryFact.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryFact.java
new file mode 100644
index 000000000..23ead3e59
--- /dev/null
+++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/temporal/model/MemoryFact.java
@@ -0,0 +1,133 @@
+/*
+ * 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.geaflow.ai.temporal.model;
+
+import java.util.Objects;
+import java.util.Optional;
+
+/**
+ * An immutable structured assertion about a memory entity.
+ */
+public final class MemoryFact {
+
+ private final String id;
+ private final MemoryEntity subject;
+ private final String predicate;
+ private final String literalValue;
+ private final MemoryEntity target;
+
+ private MemoryFact(
+ String id,
+ MemoryEntity subject,
+ String predicate,
+ String literalValue,
+ MemoryEntity target) {
+ this.id = requireText(id, "id");
+ this.subject = Objects.requireNonNull(subject, "subject");
+ this.predicate = requireText(predicate, "predicate");
+ this.literalValue = literalValue;
+ this.target = target;
+ }
+
+ public static MemoryFact attribute(
+ String id,
+ MemoryEntity subject,
+ String predicate,
+ String literalValue) {
+ return new MemoryFact(
+ id,
+ subject,
+ predicate,
+ requireText(literalValue, "literal value"),
+ null);
+ }
+
+ public static MemoryFact relationship(
+ String id,
+ MemoryEntity subject,
+ String predicate,
+ MemoryEntity target) {
+ return new MemoryFact(
+ id,
+ subject,
+ predicate,
+ null,
+ Objects.requireNonNull(target, "target"));
+ }
+
+ public String getId() {
+ return id;
+ }
+
+ public MemoryEntity getSubject() {
+ return subject;
+ }
+
+ public String getPredicate() {
+ return predicate;
+ }
+
+ public Optional