Skip to main content

graphforge_storage/
lib.rs

1//! GraphForge `StorageProvider` trait, canonical Arrow schemas, and Parquet backend.
2//!
3//! - [`schemas`] — Arrow schema constants for every Parquet file
4//! - [`catalog`] — DataFusion `TableProvider` / `CatalogProvider` implementations (#572)
5//! - [`writer`] — buffered Parquet write path ([`GraphWriter`]) (#579)
6//! - [`mutator`] — in-place rewrite primitives for `DELETE`/`DETACH DELETE` (#740)
7//! - [`staging`] — temp-file + atomic-rename Parquet commit ([`RewriteBatch`]) (#790)
8//! - [`adjacency`] — on-disk CSR format for the derived adjacency index (#758, ADR 0005)
9//! - [`generation`] — `topology_generation` counter, the staleness signal for derived indexes (#759)
10//! - [`search_manifest`] / [`search_publication`] — shared M19 search freshness and atomic publication
11#![forbid(unsafe_code)]
12
13pub mod adjacency;
14pub mod adjacency_delta;
15
16pub mod generation;
17pub use generation::{
18    commit_topology_aware, read_search_generation, read_topology_generation, touches_search_source,
19};
20
21pub mod graph_projection;
22pub use graph_projection::{
23    GraphProjectionSelection, GraphProjectionSummary, materialize_graph_projection,
24};
25
26pub mod project_generation;
27pub use project_generation::{
28    CURRENT_FILE, FORMAT_FILE, PROJECT_FORMAT_BYTES, ProjectCapabilityDescriptor,
29    ProjectParticipantDescriptor, ProjectParticipantSnapshot, ResolvedProjectGeneration,
30    open_or_initialize_project, resolve_project_generation,
31};
32
33mod project_failpoint;
34
35pub mod project_checkpoints;
36pub use project_checkpoints::{
37    CheckpointCreateRequest, CheckpointDeleteRequest, CheckpointReceipt, CheckpointRecord,
38    CheckpointRevertRequest, create_checkpoint, delete_checkpoint, list_checkpoints,
39    open_checkpoint_generation, revert_checkpoint,
40};
41
42pub mod project_publication;
43pub use project_publication::{
44    ProjectCapability, ProjectGenerationRequest, ProjectParticipant, ProjectParticipantEncoding,
45    ProjectPublicationReceipt, ProjectStageOutcome, StagedParticipant, StagedProjectGeneration,
46    ValidatedProjectGeneration, published_project_transaction, stage_project_generation,
47    stage_project_generation_optimistic,
48};
49
50pub mod project_recovery;
51pub use project_recovery::{ProjectRecoveryReport, recover_project_transactions};
52
53pub mod project_portable;
54pub use project_portable::{
55    PortableExportReceipt, PortableImportReceipt, PortableProjectLimits, encode_portable_project,
56    export_portable_project, import_portable_project, import_portable_project_file,
57};
58
59pub mod workspace_participants;
60pub use workspace_participants::{
61    MAX_WORKSPACE_REPOSITORY_SNAPSHOT_BYTES, MAX_WORKSPACE_REPOSITORY_SNAPSHOT_ENTRIES,
62    MAX_WORKSPACE_REPOSITORY_SNAPSHOT_ID_BYTES, WORKSPACE_CAPABILITY_ID,
63    WORKSPACE_CAPABILITY_VERSION, WORKSPACE_CONFIGURATION_FAMILY, WORKSPACE_ONTOLOGY_FAMILY,
64    WORKSPACE_REPOSITORY_SNAPSHOT_FAMILY, WORKSPACE_REPOSITORY_SNAPSHOT_VERSION,
65    WorkspaceConfiguration, WorkspaceOntology, WorkspaceOntologyMode,
66    WorkspaceOntologySourceFormat, WorkspaceRepositoryDefinitionDigest,
67    WorkspaceRepositoryGitProvenance, WorkspaceRepositorySnapshot, WorkspaceRepositorySourceDigest,
68    empty_workspace_participants,
69};
70
71pub mod embedding_identity;
72pub use embedding_identity::{
73    ChunkingIdentity, EmbeddingCompatibilityDescriptor, EmbeddingCompatibilityId,
74    EmbeddingCompatibilityInput, EmbeddingContentDigest, EmbeddingDisplayName, EmbeddingDistance,
75    EmbeddingGenerationId, EmbeddingNormalization, EmbeddingProducerIdentity,
76    EmbeddingSourceFingerprint, EmbeddingValueType, TokenCountClass, TokenizerIdentity,
77};
78
79pub mod embedding_catalog;
80pub use embedding_catalog::{
81    EMBEDDING_SPACE_CATALOG_VERSION, EmbeddingSpaceCatalog, EmbeddingSpaceCatalogEntry,
82    EmbeddingSpaceCatalogLimits, EmbeddingSpaceCatalogUpdate,
83    bind_existing_embedding_space_catalog_entry, read_embedding_space_catalog,
84    remove_embedding_space_catalog_identity, update_embedding_space_catalog,
85};
86
87pub mod embedding_discovery;
88pub use embedding_discovery::{
89    DiscoveredEmbeddingSpace, EmbeddingSpaceDiscoveryLimits,
90    MAX_DISCOVERED_EMBEDDING_DESCRIPTOR_BYTES, MAX_DISCOVERED_EMBEDDING_SPACES,
91    MAX_EMBEDDING_SPACE_DIRECTORY_ENTRIES, discover_embedding_spaces,
92};
93
94pub mod embedding_batch;
95pub use embedding_batch::{EmbeddingBatchRow, ValidatedEmbeddingBatch, validate_embedding_batch};
96
97pub mod embedding_manifest;
98pub use embedding_manifest::{
99    EMBEDDING_GENERATION_MANIFEST_VERSION, EmbeddingGenerationManifest,
100    EmbeddingGenerationManifestInput, EmbeddingPublicationFingerprint, EmbeddingSourceState,
101    MAX_EMBEDDING_GENERATION_MANIFEST_BYTES,
102};
103
104pub mod embedding_publication;
105pub use embedding_publication::{
106    EmbeddingGenerationPublication, EmbeddingPublicationOutcome, EmbeddingPublicationRequest,
107    current_embedding_generation, delete_embedding_space_lineage, publish_embedding_generation,
108};
109
110pub mod embedding_freshness;
111pub use embedding_freshness::{
112    EMBEDDING_SUBSTANTIAL_CHANGED_PERCENT, EMBEDDING_SUBSTANTIAL_MUTATION_BATCHES,
113    EmbeddingForcedStaleDiagnostic, EmbeddingFreshness, EmbeddingFreshnessReason,
114    EmbeddingFreshnessState, EmbeddingMutationObservation, EmbeddingReadDecision,
115    classify_embedding_freshness, decide_embedding_read,
116};
117
118pub mod embedding_refresh_config;
119pub use embedding_refresh_config::{
120    DEFAULT_EMBEDDING_REFRESH_DEBOUNCE, EMBEDDING_REFRESH_CONFIG_VERSION, EmbeddingRefreshConfig,
121    EmbeddingRefreshConfigLimits, EmbeddingRefreshConfigUpdate, EmbeddingRefreshFailureClass,
122    EmbeddingRefreshOutcomeRecord, EmbeddingRefreshOutcomeStatus, EmbeddingRefreshProjectPolicy,
123    EmbeddingRefreshSpacePolicy, EmbeddingRefreshSpaceState, MAX_EMBEDDING_REFRESH_CONFIG_BYTES,
124    MAX_EMBEDDING_REFRESH_CONFIG_ENTRIES, MAX_EMBEDDING_REFRESH_JOBS,
125    ResolvedEmbeddingRefreshPolicy, read_embedding_refresh_config, update_embedding_refresh_config,
126};
127
128pub mod embedding_mutations;
129pub use embedding_mutations::{
130    EMBEDDING_MUTATION_JOURNAL_VERSION, EmbeddingMutationBatch, EmbeddingMutationJournal,
131    EmbeddingMutationJournalLimits, merge_embedding_mutation_batch,
132    read_embedding_mutation_journal, reset_embedding_mutation_journal,
133};
134
135pub mod search_manifest;
136pub use search_manifest::{
137    MAX_SEARCH_ARTIFACT_KEY_BYTES, MAX_SEARCH_MANIFEST_BYTES, MAX_SEARCH_SELECTOR_BYTES,
138    SEARCH_MANIFEST_VERSION, SearchArtifactError, SearchArtifactKey, SearchIndexKind,
139    SearchManifest, SearchSourcePart, SearchSourceSnapshot, canonical_source_fingerprint,
140};
141
142pub mod search_publication;
143pub use search_publication::{
144    PublishedSearchArtifact, SearchCoordinationLimits, SearchPublicationMode,
145    SearchPublicationOutcome, SearchPublicationPlan, SearchUpdateBuild,
146    cleanup_abandoned_search_builds, coordinate_search_publication, coordinate_search_update,
147    current_search_artifact,
148};
149
150pub mod vector_store;
151pub use vector_store::{
152    StoredVector, VECTOR_BACKEND_VERSION, VECTOR_CONTRACT_VERSION, VECTOR_DATA_FILE,
153    VectorSearchHit, VectorStoreLimits, VectorUpsertChange, apply_vector_upsert,
154    exact_cosine_search, read_vector_snapshot, search_published_vectors, upsert_published_vector,
155    validate_published_vectors, validate_vector, vector_schema, write_vector_snapshot,
156};
157
158pub mod io_stats;
159pub use io_stats::{IoSnapshot, snapshot as io_snapshot};
160
161pub mod catalog;
162pub use catalog::{
163    EdgePropertyTable, GraphCatalog, PropertyTable, TopologyNodeTable, TypedEdgeTable,
164    UnionEdgeTable, list_edge_property_stems, list_property_stems, read_edge_properties,
165    read_edges, read_edges_filtered, read_edges_filtered_observed, read_nodes, read_nodes_filtered,
166    read_nodes_filtered_observed, read_properties,
167};
168
169pub mod schemas;
170pub use schemas::{
171    ADJACENCY_CSR_SCHEMA, ADJACENCY_MANIFEST_SCHEMA, EDGE_PROPERTY_BASE_SCHEMA,
172    EXPLORATORY_EDGE_SCHEMA, PROPERTY_BASE_SCHEMA, TOPOLOGY_NODES_SCHEMA, TYPED_EDGE_SCHEMA,
173    property_schema, property_type_to_arrow, result_schema,
174};
175
176pub mod writer;
177pub use writer::{
178    GraphWriter, count_entity_properties, read_entity_properties, read_entity_property_keys,
179    read_node_property_rows, remove_edge_properties, remove_node_properties,
180    set_edge_properties_rewrite, set_node_properties, stage_remove_edge_properties,
181    stage_remove_node_properties, stage_set_edge_properties, stage_set_node_properties,
182};
183
184pub mod mutator;
185pub use mutator::{
186    delete_edges, delete_nodes, delete_nodes_and_edges, incident_edge_uuids, stage_add_node_labels,
187    stage_delete_edges, stage_delete_nodes, stage_mutate_node_labels,
188};
189
190pub mod staging;
191pub use staging::{RewriteBatch, remove_stale_temps};
192
193pub use graphforge_core::GfError;
194
195/// Minimal row type exchanged between the storage layer and the executor.
196/// Will be replaced by Arrow `RecordBatch` in Milestone 13.
197#[derive(Debug, Clone, Default)]
198pub struct StorageRow {
199    /// Column name → string value pairs (all types as strings at stub stage).
200    pub columns: Vec<(String, String)>,
201}
202
203/// Abstraction over different storage backends.
204pub trait StorageProvider: Send + Sync {
205    /// Scan all rows for the given node label.
206    ///
207    /// # Errors
208    /// Returns [`GfError`] on I/O failure or if the label is unknown.
209    fn scan_nodes(&self, label: &str) -> Result<Vec<StorageRow>, GfError>;
210}
211
212/// Parquet-backed storage provider stub.
213#[derive(Debug, Default)]
214pub struct ParquetProvider {
215    /// Optional path to the Parquet directory.
216    pub path: Option<std::path::PathBuf>,
217}
218
219impl StorageProvider for ParquetProvider {
220    fn scan_nodes(&self, _label: &str) -> Result<Vec<StorageRow>, GfError> {
221        Err(GfError::NotImplemented("scan_nodes"))
222    }
223}