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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 106 additions & 6 deletions go/cstx_ffi.h
Original file line number Diff line number Diff line change
Expand Up @@ -131,11 +131,34 @@ CstxStatusCode cstx_graph_add_nodes(struct CstxHandle *handle,
uint64_t *affected,
struct CstxBuffer *error);

/**
* Write each node as its current state, replacing the stored record.
*
* The merge path (`cstx_graph_add_nodes`) owns bulk ingest and keeps its JSON
* fast path. A replace batch is a caller restating records it already holds —
* a task's oracles, a document's current revision — so it goes through the
* shared `Value` path rather than earning a second parser.
*/
CstxStatusCode cstx_graph_replace_nodes(struct CstxHandle *handle,
struct CstxSlice data,
uint64_t *affected,
struct CstxBuffer *error);

CstxStatusCode cstx_graph_add_edges(struct CstxHandle *handle,
struct CstxSlice data,
uint64_t *affected,
struct CstxBuffer *error);

CstxStatusCode cstx_graph_delete_nodes(struct CstxHandle *handle,
struct CstxSlice node_ids_json,
uint64_t *output,
struct CstxBuffer *error);

CstxStatusCode cstx_graph_delete_edges(struct CstxHandle *handle,
struct CstxSlice edge_ids_json,
uint64_t *output,
struct CstxBuffer *error);

CstxStatusCode cstx_graph_ingest(struct CstxHandle *handle,
struct CstxSlice source,
struct CstxSlice data,
Expand Down Expand Up @@ -362,20 +385,97 @@ CstxStatusCode cstx_repo_commit(struct CstxHandle *handle,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_prepare(struct CstxHandle *handle,
struct CstxSlice message,
struct CstxSlice ref_name,
struct CstxSlice expected_head,
struct CstxSlice metadata_json,
int64_t timestamp,
uint8_t has_timestamp,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_accept(struct CstxHandle *handle,
struct CstxSlice commit,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_discard(struct CstxHandle *handle, struct CstxBuffer *error);

CstxStatusCode cstx_repo_synchronize(struct CstxHandle *handle,
struct CstxSlice payload_json,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_contains(struct CstxHandle *handle,
struct CstxSlice object,
uint8_t *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_tree(struct CstxHandle *handle,
struct CstxSlice commit,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_object_closure(struct CstxHandle *handle,
struct CstxSlice commit,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_prepare(struct CstxHandle *handle,
struct CstxSlice commit,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_history(struct CstxHandle *handle,
struct CstxSlice commit,
struct CstxSlice entity_id,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_stat(struct CstxHandle *handle,
struct CstxSlice commit,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_commits(struct CstxHandle *handle,
struct CstxSlice commit,
size_t limit,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_diff(struct CstxHandle *handle,
struct CstxSlice base,
struct CstxSlice head,
struct CstxSlice detail,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_delta(struct CstxHandle *handle,
struct CstxSlice commit,
int64_t start_timestamp,
uint8_t has_start,
int64_t end_timestamp,
uint8_t has_end,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_missing_merge(struct CstxHandle *handle,
struct CstxSlice source,
struct CstxSlice target,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_release_transient_objects(struct CstxHandle *handle,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_diff(struct CstxHandle *handle,
struct CstxSlice base_ref,
struct CstxSlice head_ref,
size_t limit,
uint8_t has_limit,
struct CstxSlice detail,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_diff_stat(struct CstxHandle *handle,
struct CstxSlice base_ref,
struct CstxSlice head_ref,
struct CstxBuffer *output,
struct CstxBuffer *error);

CstxStatusCode cstx_repo_head(struct CstxHandle *handle,
struct CstxSlice ref_name,
struct CstxBuffer *output,
Expand Down
155 changes: 150 additions & 5 deletions go/cstx_native_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,151 @@ func addDomain(t *testing.T, rt *CSTX, value string) uint64 {
return affected
}

func TestRepositoryExternalPersistenceRoundTrip(t *testing.T) {
writer := openRuntime(t)
addDomain(t, writer, "persisted.example")

prepared, err := writer.Repo.Prepare(
testContext,
"external persistence",
"main",
nil,
map[string]any{"source": "go-test"},
nil,
)
if err != nil {
t.Fatalf("prepare: %v", err)
}
if prepared.Commit.ID == "" || prepared.IndexRoot == "" || len(prepared.Objects) == 0 {
t.Fatalf("incomplete prepared payload: %+v", prepared)
}

objects := make(map[string]RepositoryObject, len(prepared.Objects))
var commitObject, indexObject RepositoryObject
for _, object := range prepared.Objects {
stored := RepositoryObject{ID: object.ID, Envelope: append([]byte(nil), object.Envelope...)}
objects[object.ID] = stored
if object.Kind == "commit" && object.ID == prepared.Commit.ID {
commitObject = stored
}
if object.ID == prepared.IndexRoot {
indexObject = stored
}
}
if commitObject.ID == "" || indexObject.ID == "" {
t.Fatal("prepared payload does not contain commit and index-root envelopes")
}
if err := writer.Repo.Accept(testContext, prepared.Commit.ID); err != nil {
t.Fatalf("accept: %v", err)
}

reader := openRuntime(t)
head := prepared.Commit.ID
if err := reader.Repo.Synchronize(testContext, RepositorySync{
Objects: []RepositoryObject{commitObject, indexObject},
}); err != nil {
t.Fatalf("synchronize commit objects: %v", err)
}
if err := reader.Repo.Synchronize(testContext, RepositorySync{
Refs: []RepositoryRef{{Name: "main", Commit: &head}},
Indexes: []RepositoryIndex{{Commit: head, IndexRoot: prepared.IndexRoot}},
}); err != nil {
t.Fatalf("synchronize commit frontier: %v", err)
}

for {
missing, err := reader.Repo.MissingTree(testContext, head)
if err != nil {
t.Fatalf("plan missing tree: %v", err)
}
if len(missing) == 0 {
break
}
batch := make([]RepositoryObject, 0, len(missing))
for _, id := range missing {
object, ok := objects[id]
if !ok {
t.Fatalf("planner requested unknown object %s", id)
}
batch = append(batch, object)
}
if err := reader.Repo.Synchronize(testContext, RepositorySync{Objects: batch}); err != nil {
t.Fatalf("hydrate tree: %v", err)
}
}
if _, err := reader.Repo.Checkout(testContext, "main", true); err != nil {
t.Fatalf("checkout hydrated main: %v", err)
}
node, err := reader.Graph.Node(testContext, "domain:persisted.example")
if err != nil || node.Value != "persisted.example" {
t.Fatalf("restored node: %+v err=%v", node, err)
}
if err := reader.Repo.ReleaseTransientObjects(testContext); err != nil {
t.Fatalf("release transient objects: %v", err)
}
}

func TestGraphDeleteNodesCascadesAndCommits(t *testing.T) {
rt := openRuntime(t)
addDomain(t, rt, "delete-a.example")
addDomain(t, rt, "delete-b.example")
addDomain(t, rt, "keep.example")
edges := []Edge{
relatedEdge("domain:delete-a.example", "domain:delete-b.example"),
relatedEdge("domain:delete-b.example", "domain:keep.example"),
}
if _, err := rt.Graph.AddEdges(testContext, edges); err != nil {
t.Fatalf("add edges: %v", err)
}
base, err := rt.Repo.Commit(testContext, "base", "main", nil, nil)
if err != nil {
t.Fatalf("commit base: %v", err)
}
cursor, err := rt.Graph.Nodes(testContext, NodeFilter{}, CollectionOptions{})
if err != nil {
t.Fatalf("open cursor: %v", err)
}
defer cursor.Close()

affected, err := rt.Graph.DeleteNodes(testContext, []string{"domain:delete-b.example"})
if err != nil || affected != 3 {
t.Fatalf("delete node: affected=%d err=%v", affected, err)
}
if _, err := cursor.Page(testContext, 10, 1); !IsCode(err, CodeCursorInvalidated) {
t.Fatalf("expected cursor invalidation, got %v", err)
}
if count, _ := rt.Graph.NodeCount(testContext); count != 2 {
t.Fatalf("node count after delete=%d", count)
}
if count, _ := rt.Graph.EdgeCount(testContext); count != 0 {
t.Fatalf("edge count after cascade=%d", count)
}
change, err := rt.LastChange(testContext)
if err != nil || !reflect.DeepEqual(change.RemovedNodeIDs, []string{"domain:delete-b.example"}) || len(change.RemovedEdgeIDs) != 2 {
t.Fatalf("delete change=%+v err=%v", change, err)
}
head, err := rt.Repo.Commit(testContext, "delete", "main", &base.ID, nil)
if err != nil {
t.Fatalf("commit delete: %v", err)
}
diff, err := rt.Repo.Diff(testContext, base.ID, head.ID, DiffOptions{})
if err != nil || !reflect.DeepEqual(diff.Removed["domain"], []string{"domain:delete-b.example"}) || len(diff.Removed["edge:related"]) != 2 {
t.Fatalf("delete diff=%+v err=%v", diff, err)
}
}

func TestGraphDeleteIsAtomicOnMissingID(t *testing.T) {
rt := openRuntime(t)
addDomain(t, rt, "present.example")
affected, err := rt.Graph.DeleteNodes(testContext, []string{"domain:present.example", "domain:missing.example"})
if err == nil || affected != 0 {
t.Fatalf("expected atomic validation failure: affected=%d err=%v", affected, err)
}
if count, _ := rt.Graph.NodeCount(testContext); count != 1 {
t.Fatalf("failed delete changed graph: count=%d", count)
}
}

func TestSchemas(t *testing.T) {
rt := openRuntime(t)
valueField := "domain"
Expand Down Expand Up @@ -500,13 +645,13 @@ func TestRepositoryRoundTrip(t *testing.T) {
if err != nil {
t.Fatalf("second commit: %v", err)
}
diff, err := rt.Repo.Diff(testContext, commit.ID, second.ID, nil)
if err != nil || len(diff.Added["domain"]) != 1 {
diff, err := rt.Repo.Diff(testContext, commit.ID, second.ID, DiffOptions{})
if err != nil || len(diff.Added["domain"]) != 1 || diff.Stats.AddedNodes != 1 {
t.Fatalf("diff: %+v %v", diff, err)
}
diffStat, err := rt.Repo.DiffStat(testContext, commit.ID, second.ID)
if err != nil || diffStat.AddedNodes != 1 {
t.Fatalf("diff stat: %+v %v", diffStat, err)
counted, err := rt.Repo.Diff(testContext, commit.ID, second.ID, DiffOptions{Detail: DiffCounts})
if err != nil || counted.Stats.AddedNodes != 1 || len(counted.Added) != 0 {
t.Fatalf("counted diff: %+v %v", counted, err)
}
log, err := rt.Repo.Log(testContext, "main", 10)
if err != nil || len(log) != 2 {
Expand Down
21 changes: 19 additions & 2 deletions go/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,10 @@ type engine interface {
schemaAnchorConcepts(context.Context) ([]AnchorConcept, error)

graphAddNodes(context.Context, []Node) (uint64, error)
graphReplaceNodes(context.Context, []Node) (uint64, error)
graphAddEdges(context.Context, []Edge) (uint64, error)
graphDeleteNodes(context.Context, []string) (uint64, error)
graphDeleteEdges(context.Context, []string) (uint64, error)
graphIngest(context.Context, string, []byte) (uint64, error)
graphNode(context.Context, string) (Node, error)
graphContains(context.Context, string) (bool, error)
Expand All @@ -43,8 +46,22 @@ type engine interface {
repoHead(context.Context, string) (*string, error)
repoCheckout(context.Context, string, bool) (Commit, error)
repoCommit(context.Context, string, string, *string, any) (Commit, error)
repoDiff(context.Context, string, string, *int) (GraphDiff, error)
repoDiffStat(context.Context, string, string) (Delta, error)
repoPrepare(context.Context, string, string, *string, any, *int64) (PreparedCommit, error)
repoAccept(context.Context, string) error
repoDiscard(context.Context) error
repoSynchronize(context.Context, RepositorySync) error
repoContains(context.Context, string) (bool, error)
repoMissingTree(context.Context, string) ([]string, error)
repoObjectClosure(context.Context, string) ([]string, error)
repoMissingPrepare(context.Context, string) ([]string, error)
repoMissingHistory(context.Context, string, string) ([]string, error)
repoMissingStat(context.Context, string) ([]string, error)
repoMissingCommits(context.Context, string, int) ([]string, error)
repoMissingDiff(context.Context, string, string, DiffDetail) ([]string, error)
repoMissingDelta(context.Context, string, *int64, *int64) ([]string, error)
repoMissingMerge(context.Context, string, string) ([]string, error)
repoReleaseTransientObjects(context.Context) error
repoDiff(context.Context, string, string, DiffOptions) (GraphDiff, error)
repoLog(context.Context, string, int) ([]map[string]any, error)
repoHistory(context.Context, string, string, *int) ([]map[string]any, error)
repoBranch(context.Context, string, string) (string, error)
Expand Down
Loading
Loading