diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/ingest/ImportStateMachine.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/ingest/ImportStateMachine.java new file mode 100644 index 000000000..a7aad805a --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/ingest/ImportStateMachine.java @@ -0,0 +1,55 @@ +/* + * 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.retrieval.ingest; + +import org.apache.geaflow.ai.retrieval.metadata.ImportState; +import org.apache.geaflow.ai.retrieval.metadata.MetadataException; + +/** Monotonic transitions within one attempt; retries create a new graph version. */ +public final class ImportStateMachine { + + private ImportStateMachine() { + } + + public static boolean canTransition(ImportState current, ImportState target) { + return current == ImportState.IMPORTING + && (target == ImportState.INDEXING || target == ImportState.FAILED) + || current == ImportState.INDEXING + && (target == ImportState.READY || target == ImportState.FAILED); + } + + public static void validate(ImportState current, ImportState target) { + if (!canTransition(current, target)) { + throw new MetadataException(MetadataException.Code.INVALID_TRANSITION, + "illegal import transition: " + current + " -> " + target); + } + } + + /** Validates and returns the next state for callers driving an import explicitly. */ + public static ImportState transition(ImportState current, ImportState target) { + validate(current, target); + return target; + } + + /** Alias used by adapters that name the operation as a transition validation. */ + public static void validateTransition(ImportState current, ImportState target) { + validate(current, target); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ChunkingConfiguration.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ChunkingConfiguration.java new file mode 100644 index 000000000..5f9f09b84 --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ChunkingConfiguration.java @@ -0,0 +1,82 @@ +/* + * 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.retrieval.metadata; + +import java.util.Objects; + +/** Reproducible character-based chunking policy. */ +public final class ChunkingConfiguration { + + private final String policyVersion; + private final int chunkSize; + private final int overlap; + + public ChunkingConfiguration( + String policyVersion, + int chunkSize, + int overlap) { + MetadataValidation.required(policyVersion, "policyVersion"); + if (chunkSize <= 0 || overlap < 0 || overlap >= chunkSize) { + throw MetadataValidation.invalid("chunkSize must be positive and 0 <= overlap < chunkSize"); + } + this.policyVersion = policyVersion; + this.chunkSize = chunkSize; + this.overlap = overlap; + } + + public String getPolicyVersion() { + return policyVersion; + } + + public String getStrategyVersion() { + return policyVersion; + } + + public int getChunkSize() { + return chunkSize; + } + + public int getOverlap() { + return overlap; + } + + public int getOverlapSize() { + return overlap; + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof ChunkingConfiguration)) { + return false; + } + ChunkingConfiguration that = (ChunkingConfiguration) object; + return chunkSize == that.chunkSize + && overlap == that.overlap + && Objects.equals(policyVersion, that.policyVersion); + } + + @Override + public int hashCode() { + return Objects.hash(policyVersion, chunkSize, overlap); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/DatasetManifest.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/DatasetManifest.java new file mode 100644 index 000000000..31a245e19 --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/DatasetManifest.java @@ -0,0 +1,202 @@ +/* + * 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.retrieval.metadata; + +import java.util.Objects; + +/** Immutable inputs needed to reproduce a dataset build. */ +public final class DatasetManifest { + + public static final String CURRENT_VERSION = "v1"; + + private final String manifestVersion; + private final String dataset; + private final String datasetRelease; + private final String split; + private final String sourceUri; + private final String cachePath; + private final String sha256; + private final String preprocessingVersion; + private final ChunkingConfiguration chunking; + private final String graphSchemaVersion; + private final String vectorSource; + private final String vectorVersion; + private final long randomSeed; + + public DatasetManifest( + String manifestVersion, + String dataset, + String datasetRelease, + String split, + String sourceUri, + String cachePath, + String sha256, + String preprocessingVersion, + ChunkingConfiguration chunking, + String graphSchemaVersion, + String vectorSource, + String vectorVersion, + long randomSeed) { + if (!CURRENT_VERSION.equals(manifestVersion)) { + throw MetadataValidation.invalid("unsupported manifestVersion"); + } + MetadataValidation.required(dataset, "dataset"); + MetadataValidation.required(datasetRelease, "datasetRelease"); + MetadataValidation.required(split, "split"); + if ((sourceUri == null || sourceUri.trim().isEmpty()) && (cachePath == null || cachePath.trim().isEmpty())) { + throw MetadataValidation.invalid("sourceUri or cachePath is required"); + } + sha256 = MetadataValidation.checksum(sha256); + MetadataValidation.required(preprocessingVersion, "preprocessingVersion"); + if (chunking == null) { + throw MetadataValidation.invalid("chunking is required"); + } + MetadataValidation.required(graphSchemaVersion, "graphSchemaVersion"); + if ((vectorSource == null) != (vectorVersion == null)) { + throw MetadataValidation.invalid("vectorSource and vectorVersion must be supplied together"); + } + if (vectorSource != null) { + MetadataValidation.required(vectorSource, "vectorSource"); + MetadataValidation.required(vectorVersion, "vectorVersion"); + } + this.manifestVersion = manifestVersion; + this.dataset = dataset; + this.datasetRelease = datasetRelease; + this.split = split; + this.sourceUri = sourceUri; + this.cachePath = cachePath; + this.sha256 = sha256; + this.preprocessingVersion = preprocessingVersion; + this.chunking = chunking; + this.graphSchemaVersion = graphSchemaVersion; + this.vectorSource = vectorSource; + this.vectorVersion = vectorVersion; + this.randomSeed = randomSeed; + } + + public String getManifestVersion() { + return manifestVersion; + } + + public String getDataset() { + return dataset; + } + + public String getDatasetRelease() { + return datasetRelease; + } + + public String getDatasetVersion() { + return datasetRelease; + } + + public String getRelease() { + return datasetRelease; + } + + public String getSplit() { + return split; + } + + public String getSourceUri() { + return sourceUri; + } + + public String getSourceUrl() { + return sourceUri; + } + + public String getCachePath() { + return cachePath; + } + + public String getSha256() { + return sha256; + } + + public String getChecksum() { + return sha256; + } + + public String getPreprocessingVersion() { + return preprocessingVersion; + } + + public ChunkingConfiguration getChunking() { + return chunking; + } + + public String getGraphSchemaVersion() { + return graphSchemaVersion; + } + + public String getVectorSource() { + return vectorSource; + } + + public String getVectorVersion() { + return vectorVersion; + } + + public long getRandomSeed() { + return randomSeed; + } + + /** Returns the dataset identifier using the terminology used by ingestion clients. */ + public String getDatasetId() { + return dataset; + } + + /** Returns the dataset name alias retained for manifest consumers. */ + public String getDatasetName() { + return dataset; + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof DatasetManifest)) { + return false; + } + DatasetManifest that = (DatasetManifest) object; + return randomSeed == that.randomSeed + && Objects.equals(manifestVersion, that.manifestVersion) + && Objects.equals(dataset, that.dataset) + && Objects.equals(datasetRelease, that.datasetRelease) + && Objects.equals(split, that.split) + && Objects.equals(sourceUri, that.sourceUri) + && Objects.equals(cachePath, that.cachePath) + && Objects.equals(sha256, that.sha256) + && Objects.equals(preprocessingVersion, that.preprocessingVersion) + && Objects.equals(chunking, that.chunking) + && Objects.equals(graphSchemaVersion, that.graphSchemaVersion) + && Objects.equals(vectorSource, that.vectorSource) + && Objects.equals(vectorVersion, that.vectorVersion); + } + + @Override + public int hashCode() { + return Objects.hash(manifestVersion, dataset, datasetRelease, split, sourceUri, cachePath, + sha256, preprocessingVersion, chunking, graphSchemaVersion, vectorSource, vectorVersion, + randomSeed); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/GraphBuildMetadata.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/GraphBuildMetadata.java new file mode 100644 index 000000000..38c7f40f4 --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/GraphBuildMetadata.java @@ -0,0 +1,157 @@ +/* + * 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.retrieval.metadata; + +import java.util.List; +import java.util.Objects; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; +import org.apache.geaflow.ai.retrieval.validation.ModelValidation; + +/** Graph artifact and the exact required index names for its publication. */ +public final class GraphBuildMetadata { + + private final GraphVersion graphVersion; + private final String artifactUri; + private final String graphType; + private final String builderVersion; + private final boolean ready; + private final List requiredIndexes; + + /** Creates graph metadata with the artifact details used by the current build. */ + public GraphBuildMetadata( + GraphVersion graphVersion, + String artifactUri, + String graphType, + String builderVersion, + boolean ready, + List requiredIndexes) { + this.graphVersion = Objects.requireNonNull(graphVersion, "graphVersion"); + this.artifactUri = artifactUri; + this.graphType = MetadataValidation.required(graphType, "graphType"); + this.builderVersion = MetadataValidation.required(builderVersion, "builderVersion"); + if (ready) { + MetadataValidation.required(artifactUri, "artifactUri"); + } + List names = ModelValidation.sortedStrings(requiredIndexes, "requiredIndexes"); + if (names.stream().distinct().count() != names.size()) { + throw MetadataValidation.invalid("requiredIndexes must be unique"); + } + this.ready = ready; + this.requiredIndexes = names; + } + + /** Creates graph metadata while retaining the original four-argument API. */ + public GraphBuildMetadata( + GraphVersion graphVersion, + String artifactUri, + boolean ready, + List requiredIndexes) { + this(graphVersion, artifactUri, "graph", "unknown", ready, requiredIndexes); + } + + /** Creates graph metadata without required indexes for a graph-only build. */ + public GraphBuildMetadata( + GraphVersion graphVersion, + String artifactUri, + String graphType, + String builderVersion, + boolean ready) { + this(graphVersion, artifactUri, graphType, builderVersion, ready, java.util.Collections.emptyList()); + } + + /** Creates graph metadata with the list before the readiness flag. */ + public GraphBuildMetadata( + GraphVersion graphVersion, + String artifactUri, + String graphType, + String builderVersion, + List requiredIndexes, + boolean ready) { + this(graphVersion, artifactUri, graphType, builderVersion, ready, requiredIndexes); + } + + public GraphVersion getGraphVersion() { + return graphVersion; + } + + public GraphVersion getVersion() { + return graphVersion; + } + + public String getArtifactUri() { + return artifactUri; + } + + public String getGraphType() { + return graphType; + } + + /** Alias for consumers that use a generic artifact type name. */ + public String getType() { + return graphType; + } + + public String getBuilderVersion() { + return builderVersion; + } + + public String getBuildVersion() { + return builderVersion; + } + + public boolean isReady() { + return ready; + } + + public boolean getReady() { + return ready; + } + + public List getRequiredIndexes() { + return requiredIndexes; + } + + /** Alias that makes the publication role explicit at call sites. */ + public List getRequiredIndexNames() { + return requiredIndexes; + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof GraphBuildMetadata)) { + return false; + } + GraphBuildMetadata that = (GraphBuildMetadata) object; + return ready == that.ready + && Objects.equals(graphVersion, that.graphVersion) + && Objects.equals(artifactUri, that.artifactUri) + && Objects.equals(graphType, that.graphType) + && Objects.equals(builderVersion, that.builderVersion) + && Objects.equals(requiredIndexes, that.requiredIndexes); + } + + @Override + public int hashCode() { + return Objects.hash(graphVersion, artifactUri, graphType, builderVersion, ready, requiredIndexes); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ImportMetadata.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ImportMetadata.java new file mode 100644 index 000000000..5ae5325d1 --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ImportMetadata.java @@ -0,0 +1,223 @@ +/* + * 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.retrieval.metadata; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Objects; +import org.apache.geaflow.ai.retrieval.validation.ModelValidation; + +/** Immutable attempt snapshot; timestamps are UTC epoch milliseconds. */ +public final class ImportMetadata { + + private final DatasetManifest manifest; + private final String importerVersion; + private final GraphBuildMetadata graph; + private final List indexes; + private final QualityCounters counters; + private final ImportState state; + private final String verifiedSha256; + private final MetadataException.Code failureCode; + private final String failureReason; + private final long startedAt; + private final long updatedAt; + + public ImportMetadata(DatasetManifest manifest, String importerVersion, GraphBuildMetadata graph, + List indexes, QualityCounters counters, ImportState state, + String verifiedSha256, MetadataException.Code failureCode, String failureReason, + long startedAt, long updatedAt) { + this.manifest = Objects.requireNonNull(manifest, "manifest"); + this.importerVersion = MetadataValidation.required(importerVersion, "importerVersion"); + this.graph = Objects.requireNonNull(graph, "graph"); + this.indexes = ModelValidation.immutableList(indexes, "indexes"); + this.counters = Objects.requireNonNull(counters, "counters"); + this.state = Objects.requireNonNull(state, "state"); + this.verifiedSha256 = verifiedSha256 == null ? null : MetadataValidation.checksum(verifiedSha256); + this.failureCode = failureCode; + this.failureReason = failureReason; + this.startedAt = MetadataValidation.nonNegative(startedAt, "startedAt"); + this.updatedAt = MetadataValidation.nonNegative(updatedAt, "updatedAt"); + if (updatedAt < startedAt) { + throw MetadataValidation.invalid("updatedAt precedes startedAt"); + } + if (this.verifiedSha256 != null && !manifest.getSha256().equals(this.verifiedSha256)) { + throw MetadataValidation.invalid("verifiedSha256 does not match manifest sha256"); + } + if (state == ImportState.FAILED) { + Objects.requireNonNull(failureCode, "failureCode"); + MetadataValidation.required(failureReason, "failureReason"); + } else if (failureCode != null || failureReason != null) { + throw MetadataValidation.invalid("failure is only allowed in FAILED"); + } + if (this.indexes.stream().map(index -> index.getIndexVersion().getIndexName()).distinct().count() + != this.indexes.size()) { + throw MetadataValidation.invalid("duplicate index names"); + } + for (IndexBuildMetadata index : this.indexes) { + if (!graph.getGraphVersion().equals(index.getGraphVersion())) { + throw MetadataValidation.invalid("index belongs to another graph"); + } + } + if (state == ImportState.IMPORTING && graph.isReady()) { + throw MetadataValidation.invalid("IMPORTING graph must not be ready"); + } + if (state == ImportState.INDEXING || state == ImportState.READY) { + if (!graph.isReady() || !manifest.getSha256().equals(this.verifiedSha256)) { + throw new MetadataException(MetadataException.Code.NOT_READY, + "verified source and completed graph required"); + } + } + if (state == ImportState.READY && !readinessReasons().isEmpty()) { + throw new MetadataException(MetadataException.Code.NOT_READY, readinessReasons().toString()); + } + } + + public List readinessReasons() { + List reasons = new ArrayList<>(); + if (state == ImportState.FAILED) { + reasons.add("FAILED: " + failureReason); + } + if (state == ImportState.IMPORTING) { + reasons.add("STATE_NOT_READY: IMPORTING"); + } + if (!manifest.getSha256().equals(verifiedSha256)) { + reasons.add("CHECKSUM_NOT_VERIFIED"); + } + if (!graph.isReady()) { + reasons.add("GRAPH_NOT_READY"); + } + for (String name : graph.getRequiredIndexes()) { + if (indexes.stream().noneMatch(index -> name.equals(index.getIndexVersion().getIndexName()) && index.isReady())) { + reasons.add("INDEX_NOT_READY: " + name); + } + } + return Collections.unmodifiableList(reasons); + } + + public DatasetManifest getManifest() { + return manifest; + } + + public String getImporterVersion() { + return importerVersion; + } + + public GraphBuildMetadata getGraph() { + return graph; + } + + public List getIndexes() { + return indexes; + } + + public QualityCounters getCounters() { + return counters; + } + + public ImportState getState() { + return state; + } + + public String getVerifiedSha256() { + return verifiedSha256; + } + + public MetadataException.Code getFailureCode() { + return failureCode; + } + + public String getFailureReason() { + return failureReason; + } + + public long getStartedAt() { + return startedAt; + } + + public long getUpdatedAt() { + return updatedAt; + } + + public GraphBuildMetadata getGraphMetadata() { + return graph; + } + + public GraphBuildMetadata getGraphBuildMetadata() { + return graph; + } + + public org.apache.geaflow.ai.retrieval.model.version.GraphVersion getGraphVersion() { + return graph.getGraphVersion(); + } + + public List getReadinessReasons() { + return readinessReasons(); + } + + public String getReadinessReason() { + List reasons = readinessReasons(); + return reasons.isEmpty() ? null : String.join(", ", reasons); + } + + public String getChecksum() { + return verifiedSha256; + } + + public long getCreatedAt() { + return startedAt; + } + + public long getLastUpdatedAt() { + return updatedAt; + } + + public boolean isReady() { + return state == ImportState.READY && readinessReasons().isEmpty(); + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof ImportMetadata)) { + return false; + } + ImportMetadata that = (ImportMetadata) object; + return startedAt == that.startedAt + && updatedAt == that.updatedAt + && Objects.equals(manifest, that.manifest) + && Objects.equals(importerVersion, that.importerVersion) + && Objects.equals(graph, that.graph) + && Objects.equals(indexes, that.indexes) + && Objects.equals(counters, that.counters) + && state == that.state + && Objects.equals(verifiedSha256, that.verifiedSha256) + && failureCode == that.failureCode + && Objects.equals(failureReason, that.failureReason); + } + + @Override + public int hashCode() { + return Objects.hash(manifest, importerVersion, graph, indexes, counters, state, + verifiedSha256, failureCode, failureReason, startedAt, updatedAt); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ImportState.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ImportState.java new file mode 100644 index 000000000..c91d2e080 --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/ImportState.java @@ -0,0 +1,28 @@ +/* + * 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.retrieval.metadata; + +/** Lifecycle of a single import attempt. */ +public enum ImportState { + IMPORTING, + INDEXING, + READY, + FAILED +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/InMemoryMetadataStore.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/InMemoryMetadataStore.java new file mode 100644 index 000000000..58492ebaa --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/InMemoryMetadataStore.java @@ -0,0 +1,245 @@ +/* + * 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.retrieval.metadata; + +import java.io.IOException; +import java.io.InputStream; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.Clock; +import java.util.ArrayList; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.stream.Collectors; +import org.apache.geaflow.ai.retrieval.ingest.ImportStateMachine; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; + +/** Single-process reference implementation with atomic snapshot publication. */ +public final class InMemoryMetadataStore implements MetadataStore { + + private final Map attempts = new LinkedHashMap<>(); + private final Map published = new LinkedHashMap<>(); + private final Clock clock; + + public InMemoryMetadataStore() { + this(Clock.systemUTC()); + } + + public InMemoryMetadataStore(Clock clock) { + this.clock = Objects.requireNonNull(clock, "clock"); + } + + @Override + public synchronized ImportMetadata begin(DatasetManifest manifest, String importerVersion, GraphBuildMetadata graph) { + Objects.requireNonNull(manifest, "manifest"); + MetadataValidation.required(importerVersion, "importerVersion"); + Objects.requireNonNull(graph, "graph"); + GraphVersion version = graph.getGraphVersion(); + if (attempts.containsKey(version)) { + throw new MetadataException(MetadataException.Code.VERSION_CONFLICT, "attempt already exists"); + } + if (graph.isReady()) { + throw MetadataValidation.invalid("new attempt cannot have a ready graph"); + } + long now = clock.millis(); + ImportMetadata value = new ImportMetadata(manifest, importerVersion, graph, Collections.emptyList(), + new QualityCounters(0, 0, 0, 0, 0, 0), ImportState.IMPORTING, null, null, null, now, now); + attempts.put(version, value); + return value; + } + + @Override + public synchronized ImportMetadata retry(GraphVersion failedVersion, GraphBuildMetadata newGraph) { + Objects.requireNonNull(failedVersion, "failedVersion"); + Objects.requireNonNull(newGraph, "newGraph"); + ImportMetadata old = require(failedVersion); + if (old.getState() != ImportState.FAILED + || !failedVersion.getGraphName().equals(newGraph.getGraphVersion().getGraphName())) { + throw new MetadataException(MetadataException.Code.INVALID_TRANSITION, + "retry requires a failed attempt of the same graph"); + } + return begin(old.getManifest(), old.getImporterVersion(), newGraph); + } + + @Override + public synchronized ImportMetadata verifySource(GraphVersion version, InputStream source) throws IOException { + Objects.requireNonNull(source, "source"); + final ImportMetadata old = requireState(version, ImportState.IMPORTING); + String hash = digest(source); + if (!old.getManifest().getSha256().equals(hash)) { + fail(version, MetadataException.Code.CHECKSUM_MISMATCH, "source SHA-256 differs from manifest"); + throw new MetadataException(MetadataException.Code.CHECKSUM_MISMATCH, "source SHA-256 differs from manifest"); + } + return save(old, old.getGraph(), old.getIndexes(), old.getCounters(), old.getState(), hash, null, null); + } + + @Override + public synchronized ImportMetadata finishImport(GraphVersion version, GraphBuildMetadata graph, QualityCounters counters) { + Objects.requireNonNull(graph, "graph"); + Objects.requireNonNull(counters, "counters"); + ImportMetadata old = requireState(version, ImportState.IMPORTING); + if (!version.equals(graph.getGraphVersion()) + || !old.getGraph().getRequiredIndexes().equals(graph.getRequiredIndexes())) { + throw MetadataValidation.invalid("graph identity and required indexes cannot change within an attempt"); + } + if (!graph.isReady() || old.getVerifiedSha256() == null) { + throw new MetadataException(MetadataException.Code.NOT_READY, "verified source and completed graph required"); + } + ImportStateMachine.validate(old.getState(), ImportState.INDEXING); + return save(old, graph, old.getIndexes(), counters, ImportState.INDEXING, old.getVerifiedSha256(), null, null); + } + + @Override + public synchronized ImportMetadata putIndex(GraphVersion version, IndexBuildMetadata index) { + Objects.requireNonNull(index, "index"); + ImportMetadata old = requireState(version, ImportState.INDEXING); + if (!version.equals(index.getGraphVersion())) { + throw MetadataValidation.invalid("index belongs to another graph version"); + } + List indexes = new ArrayList<>(old.getIndexes()); + for (IndexBuildMetadata existing : indexes) { + if (existing.getIndexVersion().getIndexName().equals(index.getIndexVersion().getIndexName()) + && existing.isReady()) { + throw new MetadataException(MetadataException.Code.VERSION_CONFLICT, + "completed index is immutable"); + } + } + indexes.removeIf(existing -> existing.getIndexVersion().getIndexName() + .equals(index.getIndexVersion().getIndexName())); + indexes.add(index); + return save(old, old.getGraph(), indexes, old.getCounters(), old.getState(), + old.getVerifiedSha256(), null, null); + } + + @Override + public synchronized ImportMetadata fail(GraphVersion version, MetadataException.Code code, String reason) { + Objects.requireNonNull(code, "code"); + MetadataValidation.required(reason, "failureReason"); + ImportMetadata old = require(version); + ImportStateMachine.validate(old.getState(), ImportState.FAILED); + return save(old, old.getGraph(), old.getIndexes(), old.getCounters(), ImportState.FAILED, + old.getVerifiedSha256(), code, reason); + } + + @Override + public synchronized ImportMetadata publish(GraphVersion version, GraphVersion expectedPublishedVersion) { + ImportMetadata old = requireState(version, ImportState.INDEXING); + ImportMetadata current = published.get(version.getGraphName()); + GraphVersion actual = current == null ? null : current.getGraph().getGraphVersion(); + if (!Objects.equals(expectedPublishedVersion, actual)) { + throw new MetadataException(MetadataException.Code.VERSION_CONFLICT, "published version changed"); + } + if (!old.readinessReasons().isEmpty()) { + throw new MetadataException(MetadataException.Code.NOT_READY, + old.readinessReasons().toString()); + } + ImportStateMachine.validate(old.getState(), ImportState.READY); + ImportMetadata ready = save(old, old.getGraph(), old.getIndexes(), old.getCounters(), + ImportState.READY, old.getVerifiedSha256(), null, null); + published.put(version.getGraphName(), ready); + return ready; + } + + @Override + public synchronized Optional find(GraphVersion version) { + return Optional.ofNullable(attempts.get(version)); + } + + @Override + public synchronized Optional published(String graphName) { + return Optional.ofNullable(published.get(graphName)); + } + + @Override + public synchronized List findDataset(String dataset, String release, String split) { + return Collections.unmodifiableList(attempts.values().stream().filter(value -> + value.getManifest().getDataset().equals(dataset) + && value.getManifest().getDatasetRelease().equals(release) + && value.getManifest().getSplit().equals(split)).collect(Collectors.toList())); + } + + /** Returns the published version pointer without exposing mutable storage state. */ + public synchronized Optional publishedVersion(String graphName) { + ImportMetadata value = published.get(graphName); + return value == null ? Optional.empty() : Optional.of(value.getGraphVersion()); + } + + /** Alias for callers that use a get-style publication query. */ + public synchronized Optional getPublished(String graphName) { + return published(graphName); + } + + /** Alias for callers that use a find-style publication query. */ + public synchronized Optional findPublished(String graphName) { + return published(graphName); + } + + private ImportMetadata require(GraphVersion version) { + ImportMetadata value = attempts.get(version); + if (value == null) { + throw new MetadataException(MetadataException.Code.NOT_FOUND, "unknown graph version"); + } + return value; + } + + private ImportMetadata requireState(GraphVersion version, ImportState state) { + ImportMetadata value = require(version); + if (value.getState() != state) { + throw new MetadataException(MetadataException.Code.INVALID_TRANSITION, + "expected " + state + " but was " + value.getState()); + } + return value; + } + + private ImportMetadata save(ImportMetadata old, GraphBuildMetadata graph, + List indexes, QualityCounters counters, + ImportState state, String checksum, MetadataException.Code code, + String reason) { + ImportMetadata value = new ImportMetadata(old.getManifest(), old.getImporterVersion(), graph, indexes, + counters, state, checksum, code, reason, old.getStartedAt(), + Math.max(old.getUpdatedAt(), clock.millis())); + attempts.put(graph.getGraphVersion(), value); + return value; + } + + private static String digest(InputStream source) throws IOException { + MessageDigest digest; + try { + digest = MessageDigest.getInstance("SHA-256"); + } catch (NoSuchAlgorithmException e) { + throw new IllegalStateException("SHA-256 unavailable", e); + } + byte[] buffer = new byte[8192]; + int length; + while ((length = source.read(buffer)) != -1) { + digest.update(buffer, 0, length); + } + StringBuilder hash = new StringBuilder(64); + for (byte value : digest.digest()) { + hash.append(Character.forDigit((value & 0xff) >>> 4, 16)); + hash.append(Character.forDigit(value & 0xf, 16)); + } + return hash.toString(); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/IndexBuildMetadata.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/IndexBuildMetadata.java new file mode 100644 index 000000000..feef524a2 --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/IndexBuildMetadata.java @@ -0,0 +1,119 @@ +/* + * 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.retrieval.metadata; + +import java.util.Objects; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; +import org.apache.geaflow.ai.retrieval.model.version.IndexVersion; + +/** Index artifact bound to a complete graph identity. */ +public final class IndexBuildMetadata { + + private final GraphVersion graphVersion; + private final IndexVersion indexVersion; + private final String indexType; + private final String builderVersion; + private final String artifactUri; + private final boolean ready; + + public IndexBuildMetadata( + GraphVersion graphVersion, + IndexVersion indexVersion, + String indexType, + String builderVersion, + String artifactUri, + boolean ready) { + Objects.requireNonNull(graphVersion, "graphVersion"); + Objects.requireNonNull(indexVersion, "indexVersion"); + if (!graphVersion.getVersion().equals(indexVersion.getGraphVersion())) { + throw MetadataValidation.invalid("index source graph version mismatch"); + } + MetadataValidation.required(indexType, "indexType"); + MetadataValidation.required(builderVersion, "builderVersion"); + if (ready) { + MetadataValidation.required(artifactUri, "artifactUri"); + } + this.graphVersion = graphVersion; + this.indexVersion = indexVersion; + this.indexType = indexType; + this.builderVersion = builderVersion; + this.artifactUri = artifactUri; + this.ready = ready; + } + + public GraphVersion getGraphVersion() { + return graphVersion; + } + + public IndexVersion getIndexVersion() { + return indexVersion; + } + + public String getIndexType() { + return indexType; + } + + public String getBuilderVersion() { + return builderVersion; + } + + public String getBuildVersion() { + return builderVersion; + } + + public String getArtifactUri() { + return artifactUri; + } + + /** Alias for callers that use the shorter artifact terminology. */ + public String getType() { + return indexType; + } + + public boolean isReady() { + return ready; + } + + public boolean getReady() { + return ready; + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof IndexBuildMetadata)) { + return false; + } + IndexBuildMetadata that = (IndexBuildMetadata) object; + return ready == that.ready + && Objects.equals(graphVersion, that.graphVersion) + && Objects.equals(indexVersion, that.indexVersion) + && Objects.equals(indexType, that.indexType) + && Objects.equals(builderVersion, that.builderVersion) + && Objects.equals(artifactUri, that.artifactUri); + } + + @Override + public int hashCode() { + return Objects.hash(graphVersion, indexVersion, indexType, builderVersion, artifactUri, ready); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataException.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataException.java new file mode 100644 index 000000000..0c808437e --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataException.java @@ -0,0 +1,48 @@ +/* + * 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.retrieval.metadata; + +import java.util.Objects; + +/** Storage-neutral failure with a stable machine-readable code. */ +public final class MetadataException extends IllegalArgumentException { + + private static final long serialVersionUID = 1L; + + public enum Code { + INVALID_METADATA, + CHECKSUM_MISMATCH, + INVALID_TRANSITION, + NOT_READY, + VERSION_CONFLICT, + NOT_FOUND + } + + private final Code code; + + public MetadataException(Code code, String message) { + super(message); + this.code = Objects.requireNonNull(code, "code"); + } + + public Code getCode() { + return code; + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataJson.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataJson.java new file mode 100644 index 000000000..3c50ad30b --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataJson.java @@ -0,0 +1,204 @@ +/* + * 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.retrieval.metadata; + +import com.google.gson.Gson; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import com.google.gson.JsonParseException; +import com.google.gson.JsonParser; +import java.math.BigDecimal; +import java.util.ArrayList; +import java.util.List; +import org.apache.geaflow.ai.retrieval.codec.RetrievalModelJson; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; +import org.apache.geaflow.ai.retrieval.model.version.IndexVersion; + +/** Constructor-validated JSON boundary; unknown optional fields are ignored. */ +public final class MetadataJson { + + private static final Gson GSON = new Gson(); + + private MetadataJson() { + } + + public static String toJson(Object value) { + return GSON.toJson(value); + } + + /** Alias for serializers that use a serialize-style name. */ + public static String serialize(Object value) { + return toJson(value); + } + + public static T fromJson(String json, Class type) { + if (json == null) { + throw new JsonParseException("metadata JSON is required"); + } + if (type == null) { + throw new JsonParseException("metadata type is required"); + } + try { + return type.cast(decode(new JsonParser().parse(json).getAsJsonObject(), type)); + } catch (RuntimeException e) { + throw new JsonParseException("invalid metadata: " + e.getMessage(), e); + } + } + + /** Alias for callers that use a parse-style name. */ + public static T parse(String json, Class type) { + return fromJson(json, type); + } + + private static Object decode(JsonObject object, Class type) { + if (type == ChunkingConfiguration.class) { + return new ChunkingConfiguration( + string(object, "policyVersion"), + Math.toIntExact(number(object, "chunkSize")), + Math.toIntExact(number(object, "overlap"))); + } + if (type == QualityCounters.class) { + return new QualityCounters( + number(object, "documents"), + number(object, "chunks"), + number(object, "entities"), + number(object, "edges"), + number(object, "skippedRecords"), + number(object, "validationErrors")); + } + if (type == DatasetManifest.class) { + return new DatasetManifest( + string(object, "manifestVersion"), + string(object, "dataset"), + string(object, "datasetRelease"), + string(object, "split"), + string(object, "sourceUri"), + string(object, "cachePath"), + string(object, "sha256"), + string(object, "preprocessingVersion"), + (ChunkingConfiguration) decode(object.getAsJsonObject("chunking"), ChunkingConfiguration.class), + string(object, "graphSchemaVersion"), + string(object, "vectorSource"), + string(object, "vectorVersion"), + number(object, "randomSeed")); + } + if (type == GraphBuildMetadata.class) { + return new GraphBuildMetadata( + RetrievalModelJson.fromJson(object.get("graphVersion").toString(), GraphVersion.class), + string(object, "artifactUri"), + stringOrDefault(object, "graphType", "graph"), + stringOrDefault(object, "builderVersion", "unknown"), + bool(object, "ready"), + strings(object, "requiredIndexes")); + } + if (type == IndexBuildMetadata.class) { + return new IndexBuildMetadata( + RetrievalModelJson.fromJson(object.get("graphVersion").toString(), GraphVersion.class), + RetrievalModelJson.fromJson(object.get("indexVersion").toString(), IndexVersion.class), + string(object, "indexType"), + string(object, "builderVersion"), + string(object, "artifactUri"), + bool(object, "ready")); + } + if (type == ImportMetadata.class) { + return new ImportMetadata( + (DatasetManifest) decode(object.getAsJsonObject("manifest"), DatasetManifest.class), + string(object, "importerVersion"), + (GraphBuildMetadata) decode(object.getAsJsonObject("graph"), GraphBuildMetadata.class), + indexes(object), + (QualityCounters) decode(object.getAsJsonObject("counters"), QualityCounters.class), + ImportState.valueOf(string(object, "state")), + string(object, "verifiedSha256"), + string(object, "failureCode") == null ? null : MetadataException.Code.valueOf(string(object, "failureCode")), + string(object, "failureReason"), + number(object, "startedAt"), + number(object, "updatedAt")); + } + throw new JsonParseException("unsupported metadata type: " + type); + } + + private static String string(JsonObject object, String name) { + JsonElement value = object.get(name); + if (value == null || value.isJsonNull()) { + return null; + } + if (!value.isJsonPrimitive() || !value.getAsJsonPrimitive().isString()) { + throw new JsonParseException(name + " must be a string"); + } + return value.getAsString(); + } + + private static long number(JsonObject object, String name) { + JsonElement value = object.get(name); + if (value == null || !value.isJsonPrimitive() || !value.getAsJsonPrimitive().isNumber()) { + throw new JsonParseException(name + " must be a number"); + } + return new BigDecimal(value.getAsString()).longValueExact(); + } + + private static boolean bool(JsonObject object, String name) { + JsonElement value = object.get(name); + if (value == null || !value.isJsonPrimitive() || !value.getAsJsonPrimitive().isBoolean()) { + throw new JsonParseException(name + " must be boolean"); + } + return value.getAsBoolean(); + } + + private static List strings(JsonObject object, String name) { + List result = new ArrayList<>(); + JsonElement element = object.get(name); + if (element == null || element.isJsonNull()) { + return result; + } + if (!element.isJsonArray()) { + throw new JsonParseException(name + " must be an array"); + } + for (JsonElement value : element.getAsJsonArray()) { + if (!value.isJsonPrimitive() || !value.getAsJsonPrimitive().isString()) { + throw new JsonParseException(name + " must contain strings"); + } + result.add(value.getAsString()); + } + return result; + } + + private static List indexes(JsonObject object) { + List result = new ArrayList<>(); + JsonElement element = object.get("indexes"); + if (element == null || element.isJsonNull()) { + return result; + } + if (!element.isJsonArray()) { + throw new JsonParseException("indexes must be an array"); + } + for (JsonElement value : element.getAsJsonArray()) { + if (!value.isJsonObject()) { + throw new JsonParseException("indexes must contain objects"); + } + result.add((IndexBuildMetadata) decode(value.getAsJsonObject(), IndexBuildMetadata.class)); + } + return result; + } + + private static String stringOrDefault(JsonObject object, String name, String defaultValue) { + String value = string(object, name); + return value == null ? defaultValue : value; + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataStore.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataStore.java new file mode 100644 index 000000000..bb0febf8c --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataStore.java @@ -0,0 +1,89 @@ +/* + * 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.retrieval.metadata; + +import java.io.IOException; +import java.io.InputStream; +import java.util.List; +import java.util.Optional; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; + +/** + * Storage-neutral build lifecycle. Published snapshots are immutable. + * + *

Implementations must serialize mutations per attempt and atomically compare and replace the + * published graph pointer. Source streams remain owned by the caller. Retry uses a new graph + * version, preserving the failed attempt and any previously published version.

+ */ +public interface MetadataStore { + + ImportMetadata begin(DatasetManifest manifest, String importerVersion, GraphBuildMetadata graph); + + ImportMetadata retry(GraphVersion failedVersion, GraphBuildMetadata newGraph); + + ImportMetadata verifySource(GraphVersion version, InputStream source) throws IOException; + + ImportMetadata finishImport(GraphVersion version, GraphBuildMetadata graph, QualityCounters counters); + + ImportMetadata putIndex(GraphVersion version, IndexBuildMetadata index); + + ImportMetadata fail(GraphVersion version, MetadataException.Code code, String reason); + + ImportMetadata publish(GraphVersion version, GraphVersion expectedPublishedVersion); + + Optional find(GraphVersion version); + + Optional published(String graphName); + + List findDataset(String dataset, String release, String split); + + default List findByDataset(String dataset, String release, String split) { + return findDataset(dataset, release, split); + } + + /** Alias for implementations and callers that use get terminology. */ + default Optional get(GraphVersion version) { + return find(version); + } + + /** Returns the published snapshot for a graph name, if one exists. */ + default Optional getPublished(String graphName) { + return published(graphName); + } + + /** Reads only the published graph version while retaining snapshot immutability. */ + default Optional getPublishedVersion(String graphName) { + Optional snapshot = published(graphName); + return snapshot.map(ImportMetadata::getGraphVersion); + } + + default Optional publishedVersion(String graphName) { + return getPublishedVersion(graphName); + } + + default Optional readPublishedVersion(String graphName) { + return getPublishedVersion(graphName); + } + + /** Alias for stores that expose a find-style publication query. */ + default Optional findPublished(String graphName) { + return published(graphName); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataValidation.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataValidation.java new file mode 100644 index 000000000..c08ce337c --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/MetadataValidation.java @@ -0,0 +1,54 @@ +/* + * 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.retrieval.metadata; + +import java.util.Locale; + +/** Validation shared by metadata constructors and their JSON boundary. */ +final class MetadataValidation { + + private MetadataValidation() { + } + + static String required(String value, String field) { + if (value == null || value.trim().isEmpty()) { + throw invalid(field + " is required"); + } + return value; + } + + static String checksum(String value) { + if (value == null || !value.matches("[a-fA-F0-9]{64}")) { + throw invalid("sha256 must contain 64 hexadecimal characters"); + } + return value.toLowerCase(Locale.ROOT); + } + + static long nonNegative(long value, String field) { + if (value < 0) { + throw invalid(field + " must be non-negative"); + } + return value; + } + + static MetadataException invalid(String message) { + return new MetadataException(MetadataException.Code.INVALID_METADATA, message); + } +} diff --git a/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/QualityCounters.java b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/QualityCounters.java new file mode 100644 index 000000000..79e013b7b --- /dev/null +++ b/geaflow-ai/src/main/java/org/apache/geaflow/ai/retrieval/metadata/QualityCounters.java @@ -0,0 +1,124 @@ +/* + * 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.retrieval.metadata; + +import java.util.Objects; + +/** Non-negative quality measurements for a build. */ +public final class QualityCounters { + + private final long documents; + private final long chunks; + private final long entities; + private final long edges; + private final long skippedRecords; + private final long validationErrors; + + public QualityCounters( + long documents, + long chunks, + long entities, + long edges, + long skippedRecords, + long validationErrors) { + MetadataValidation.nonNegative(documents, "documents"); + MetadataValidation.nonNegative(chunks, "chunks"); + MetadataValidation.nonNegative(entities, "entities"); + MetadataValidation.nonNegative(edges, "edges"); + MetadataValidation.nonNegative(skippedRecords, "skippedRecords"); + MetadataValidation.nonNegative(validationErrors, "validationErrors"); + this.documents = documents; + this.chunks = chunks; + this.entities = entities; + this.edges = edges; + this.skippedRecords = skippedRecords; + this.validationErrors = validationErrors; + } + + public long getDocuments() { + return documents; + } + + public long getChunks() { + return chunks; + } + + public long getEntities() { + return entities; + } + + public long getEdges() { + return edges; + } + + public long getSkippedRecords() { + return skippedRecords; + } + + public long getValidationErrors() { + return validationErrors; + } + + public long getDocumentCount() { + return documents; + } + + public long getChunkCount() { + return chunks; + } + + public long getEntityCount() { + return entities; + } + + public long getEdgeCount() { + return edges; + } + + public long getSkippedRecordCount() { + return skippedRecords; + } + + public long getValidationErrorCount() { + return validationErrors; + } + + @Override + public boolean equals(Object object) { + if (this == object) { + return true; + } + if (!(object instanceof QualityCounters)) { + return false; + } + QualityCounters that = (QualityCounters) object; + return documents == that.documents + && chunks == that.chunks + && entities == that.entities + && edges == that.edges + && skippedRecords == that.skippedRecords + && validationErrors == that.validationErrors; + } + + @Override + public int hashCode() { + return Objects.hash(documents, chunks, entities, edges, skippedRecords, validationErrors); + } +} diff --git a/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/ingest/ImportStateMachineTest.java b/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/ingest/ImportStateMachineTest.java new file mode 100644 index 000000000..7d219e3d7 --- /dev/null +++ b/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/ingest/ImportStateMachineTest.java @@ -0,0 +1,58 @@ +/* + * 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.retrieval.ingest; + +import org.apache.geaflow.ai.retrieval.metadata.ImportState; +import org.apache.geaflow.ai.retrieval.metadata.MetadataException; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +/** State transition coverage for one immutable import attempt. */ +public class ImportStateMachineTest { + + @Test + public void allowsOnlyForwardBuildTransitionsAndFailure() { + Assertions.assertTrue(ImportStateMachine.canTransition(ImportState.IMPORTING, + ImportState.INDEXING)); + Assertions.assertTrue(ImportStateMachine.canTransition(ImportState.IMPORTING, + ImportState.FAILED)); + Assertions.assertTrue(ImportStateMachine.canTransition(ImportState.INDEXING, + ImportState.READY)); + Assertions.assertTrue(ImportStateMachine.canTransition(ImportState.INDEXING, + ImportState.FAILED)); + Assertions.assertEquals(ImportState.INDEXING, ImportStateMachine.transition( + ImportState.IMPORTING, ImportState.INDEXING)); + } + + @Test + public void rejectsTerminalAndBackwardTransitionsWithTypedError() { + ImportState[] states = ImportState.values(); + for (ImportState current : states) { + for (ImportState target : states) { + if (!ImportStateMachine.canTransition(current, target)) { + MetadataException exception = Assertions.assertThrows(MetadataException.class, + () -> ImportStateMachine.validate(current, target)); + Assertions.assertEquals(MetadataException.Code.INVALID_TRANSITION, + exception.getCode()); + } + } + } + } +} diff --git a/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/metadata/InMemoryMetadataStoreTest.java b/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/metadata/InMemoryMetadataStoreTest.java new file mode 100644 index 000000000..ce8a62e2c --- /dev/null +++ b/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/metadata/InMemoryMetadataStoreTest.java @@ -0,0 +1,125 @@ +/* + * 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.retrieval.metadata; + +import java.io.ByteArrayInputStream; +import java.nio.charset.StandardCharsets; +import java.util.Arrays; +import java.util.Collections; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; +import org.apache.geaflow.ai.retrieval.model.version.IndexVersion; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +/** Lifecycle, readiness, CAS, and retry tests for the in-memory metadata store. */ +public class InMemoryMetadataStoreTest { + + private static final String SOURCE_SHA256 = "41cf6794ba4200b839c53531555f0f3998df4cbb01a4d5cb0b94e3ca5e23947d"; + + @Test + public void mismatchFailsAttemptAndRetryKeepsTheFailureRecord() throws Exception { + InMemoryMetadataStore store = new InMemoryMetadataStore(); + GraphVersion failedVersion = new GraphVersion("graph", "g1"); + store.begin(manifest(), "importer-v1", graph(failedVersion, false, Collections.emptyList())); + + MetadataException mismatch = Assertions.assertThrows(MetadataException.class, + () -> store.verifySource(failedVersion, input("wrong"))); + Assertions.assertEquals(MetadataException.Code.CHECKSUM_MISMATCH, mismatch.getCode()); + Assertions.assertEquals(ImportState.FAILED, store.find(failedVersion).get().getState()); + + GraphVersion retryVersion = new GraphVersion("graph", "g2"); + ImportMetadata retry = store.retry(failedVersion, + graph(retryVersion, false, Collections.emptyList())); + Assertions.assertEquals(ImportState.IMPORTING, retry.getState()); + Assertions.assertEquals(ImportState.FAILED, store.find(failedVersion).get().getState()); + Assertions.assertEquals(2, store.findDataset("dataset", "release", "dev").size()); + } + + @Test + public void publicationRequiresEveryRequiredIndexAndUsesCompareAndSet() throws Exception { + InMemoryMetadataStore store = new InMemoryMetadataStore(); + GraphVersion version = new GraphVersion("graph", "g1"); + store.begin(manifest(), "importer-v1", graph(version, false, + Collections.singletonList("keyword"))); + store.verifySource(version, input("source")); + store.finishImport(version, graph(version, true, Collections.singletonList("keyword")), + new QualityCounters(1, 1, 1, 0, 0, 0)); + + MetadataException notReady = Assertions.assertThrows(MetadataException.class, + () -> store.publish(version, null)); + Assertions.assertEquals(MetadataException.Code.NOT_READY, notReady.getCode()); + Assertions.assertFalse(store.published("graph").isPresent()); + + IndexVersion indexVersion = new IndexVersion("keyword", "i1", "g1"); + store.putIndex(version, new IndexBuildMetadata(version, indexVersion, "bm25", + "builder-v1", null, false)); + Assertions.assertThrows(MetadataException.class, () -> store.publish(version, null)); + store.putIndex(version, new IndexBuildMetadata(version, indexVersion, "bm25", + "builder-v1", "memory:i1", true)); + ImportMetadata published = store.publish(version, null); + Assertions.assertEquals(ImportState.READY, published.getState()); + Assertions.assertEquals(version, store.getPublishedVersion("graph").get()); + + GraphVersion nextVersion = new GraphVersion("graph", "g2"); + store.begin(manifest(), "importer-v1", graph(nextVersion, false, Collections.emptyList())); + store.verifySource(nextVersion, input("source")); + store.finishImport(nextVersion, graph(nextVersion, true, Collections.emptyList()), + new QualityCounters(1, 1, 1, 0, 0, 0)); + MetadataException conflict = Assertions.assertThrows(MetadataException.class, + () -> store.publish(nextVersion, new GraphVersion("graph", "stale"))); + Assertions.assertEquals(MetadataException.Code.VERSION_CONFLICT, conflict.getCode()); + Assertions.assertEquals(version, store.getPublishedVersion("graph").get()); + } + + @Test + public void completedIndexCannotBeReplaced() throws Exception { + InMemoryMetadataStore store = new InMemoryMetadataStore(); + GraphVersion version = new GraphVersion("graph", "g1"); + store.begin(manifest(), "importer-v1", graph(version, false, + Collections.singletonList("keyword"))); + store.verifySource(version, input("source")); + store.finishImport(version, graph(version, true, Collections.singletonList("keyword")), + new QualityCounters(1, 1, 1, 0, 0, 0)); + IndexVersion indexVersion = new IndexVersion("keyword", "i1", "g1"); + store.putIndex(version, new IndexBuildMetadata(version, indexVersion, "bm25", + "builder-v1", "memory:i1", true)); + + MetadataException conflict = Assertions.assertThrows(MetadataException.class, + () -> store.putIndex(version, new IndexBuildMetadata(version, + new IndexVersion("keyword", "i2", "g1"), "bm25", "builder-v2", "memory:i2", true))); + Assertions.assertEquals(MetadataException.Code.VERSION_CONFLICT, conflict.getCode()); + } + + private static DatasetManifest manifest() { + return new DatasetManifest("v1", "dataset", "release", "dev", "uri", null, + SOURCE_SHA256, "preprocess-v1", new ChunkingConfiguration("chunk-v1", 100, 10), + "schema-v1", null, null, 7); + } + + private static GraphBuildMetadata graph(GraphVersion version, boolean ready, + java.util.List requiredIndexes) { + return new GraphBuildMetadata(version, ready ? "memory:" + version.getVersion() : null, + "memory", "builder-v1", ready, requiredIndexes); + } + + private static ByteArrayInputStream input(String value) { + return new ByteArrayInputStream(value.getBytes(StandardCharsets.UTF_8)); + } +} diff --git a/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/metadata/MetadataModelTest.java b/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/metadata/MetadataModelTest.java new file mode 100644 index 000000000..b371b98d6 --- /dev/null +++ b/geaflow-ai/src/test/java/org/apache/geaflow/ai/retrieval/metadata/MetadataModelTest.java @@ -0,0 +1,114 @@ +/* + * 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.retrieval.metadata; + +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import org.apache.geaflow.ai.retrieval.model.version.GraphVersion; +import org.apache.geaflow.ai.retrieval.model.version.IndexVersion; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +/** Contract tests for immutable import metadata models and their JSON boundary. */ +public class MetadataModelTest { + + private static final String SHA256 = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"; + + @Test + public void manifestValidatesRequiredFieldsAndPreservesLongSeed() { + ChunkingConfiguration chunking = new ChunkingConfiguration("chunk-v1", 512, 64); + DatasetManifest manifest = new DatasetManifest("v1", "hotpotqa", "2025", "dev", + "https://example.test/data", null, SHA256.toUpperCase(), "clean-v1", chunking, + "schema-v1", "offline-v1", "embeddings-v2", Long.MIN_VALUE); + + Assertions.assertEquals("hotpotqa", manifest.getDataset()); + Assertions.assertEquals(-9223372036854775808L, manifest.getRandomSeed()); + Assertions.assertEquals(SHA256, manifest.getSha256()); + Assertions.assertEquals(manifest, + MetadataJson.fromJson(MetadataJson.toJson(manifest), DatasetManifest.class)); + } + + @Test + public void manifestRejectsInvalidRequiredFieldsAndPairs() { + ChunkingConfiguration chunking = new ChunkingConfiguration("chunk-v1", 10, 2); + Assertions.assertThrows(MetadataException.class, () -> new DatasetManifest("v2", "d", "r", + "dev", "uri", null, SHA256, "p", chunking, "schema", null, null, 1)); + Assertions.assertThrows(MetadataException.class, () -> new DatasetManifest("v1", "d", "r", + "dev", null, " ", SHA256, "p", chunking, "schema", null, null, 1)); + Assertions.assertThrows(MetadataException.class, () -> new DatasetManifest("v1", "d", "r", + "dev", "uri", null, "bad", "p", chunking, "schema", null, null, 1)); + Assertions.assertThrows(MetadataException.class, () -> new DatasetManifest("v1", "d", "r", + "dev", "uri", null, SHA256, "p", chunking, "schema", "vectors", null, 1)); + Assertions.assertThrows(MetadataException.class, () -> new ChunkingConfiguration("v1", 2, 2)); + } + + @Test + public void graphAndIndexMetadataValidateBindingAndGraphOnlyBuilds() { + GraphVersion graphVersion = new GraphVersion("graph", "g1"); + GraphBuildMetadata graph = new GraphBuildMetadata(graphVersion, "memory:g1", "memory", + "builder-v1", true, Collections.emptyList()); + Assertions.assertTrue(graph.getRequiredIndexes().isEmpty()); + Assertions.assertEquals(graph, MetadataJson.fromJson(MetadataJson.toJson(graph), + GraphBuildMetadata.class)); + + IndexVersion indexVersion = new IndexVersion("keyword", "i1", "g1"); + IndexBuildMetadata index = new IndexBuildMetadata(graphVersion, indexVersion, "bm25", + "builder-v1", "memory:i1", true); + Assertions.assertEquals(index, MetadataJson.fromJson(MetadataJson.toJson(index), + IndexBuildMetadata.class)); + Assertions.assertThrows(MetadataException.class, () -> new IndexBuildMetadata(graphVersion, + new IndexVersion("keyword", "i1", "other"), "bm25", "builder-v1", "uri", true)); + Assertions.assertThrows(MetadataException.class, () -> new GraphBuildMetadata(graphVersion, + "uri", "memory", "builder-v1", false, Arrays.asList("keyword", "keyword"))); + } + + @Test + public void importJsonIgnoresUnknownFieldsAndDefaultsOldIndexes() { + DatasetManifest manifest = manifest(); + GraphBuildMetadata graph = new GraphBuildMetadata(new GraphVersion("graph", "g1"), null, + false, Collections.emptyList()); + ImportMetadata metadata = new ImportMetadata(manifest, "importer-v1", graph, + Collections.emptyList(), new QualityCounters(1, 2, 3, 4, 5, 6), ImportState.IMPORTING, + null, null, null, 10, 10); + JsonObject json = new JsonParser().parse(MetadataJson.toJson(metadata)).getAsJsonObject(); + json.remove("indexes"); + json.addProperty("futureField", "ignored"); + + ImportMetadata restored = MetadataJson.fromJson(json.toString(), ImportMetadata.class); + Assertions.assertTrue(restored.getIndexes().isEmpty()); + Assertions.assertEquals(metadata.getCounters(), restored.getCounters()); + Assertions.assertEquals(ImportState.IMPORTING, restored.getState()); + } + + @Test + public void countersRejectNegativeValues() { + Assertions.assertThrows(MetadataException.class, + () -> new QualityCounters(0, 0, -1, 0, 0, 0)); + } + + private static DatasetManifest manifest() { + return new DatasetManifest("v1", "dataset", "release", "dev", "uri", null, SHA256, + "preprocess-v1", new ChunkingConfiguration("chunk-v1", 100, 10), "schema-v1", + null, null, 7); + } +}