diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index babadd1..02f3007 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -9,7 +9,7 @@ permissions: contents: read concurrency: - group: ${{ github.workflow }}-${{ github.ref }} + group: ${{ github.workflow }}-${{ github.event.pull_request.number }} cancel-in-progress: true jobs: @@ -73,7 +73,7 @@ jobs: load: true tags: hstore:ci cache-from: type=gha - cache-to: type=gha,mode=max + cache-to: type=gha,mode=max,ignore-error=true - name: Start container and wait for health run: | docker run -d --name hstore -e HSTORE_PASSWORD=ci -p 7432:7432 -p 7480:7480 hstore:ci diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index a20f3b5..14f696d 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -85,7 +85,7 @@ jobs: org.opencontainers.image.revision=${{ github.sha }} org.opencontainers.image.licenses=PolyForm-Noncommercial-1.0.0 cache-from: type=gha - cache-to: type=gha,mode=max + cache-to: type=gha,mode=max,ignore-error=true - name: Start the container and wait for health run: | docker run -d --name hstore -e HSTORE_PASSWORD=ci -p 7480:7480 "$IMAGE:$VERSION-amd64" diff --git a/README.md b/README.md index bfee9a2..649ef14 100644 --- a/README.md +++ b/README.md @@ -5,7 +5,7 @@ number of atoms, carries roles, weights, validity intervals and properties, and other hyperedges. HStore stores this structure directly instead of reifying it into pairwise edges or relational join tables. -The engine implements the design of *Native Persistent Hypergraph Storage Engine* (VLDB) together with its +The engine implements the design of *Native Persistent Hypergraph Storage Engine* together with its higher-order extension. Membership and incidence are copy-on-write counted B+trees with monoid summaries. Commits are atomic root swaps. The database layer adds a query language (HQL), hypergraph relational operators, a semantic plane, multi-tenancy and a web console. diff --git a/database/src/main/java/io/hstore/db/view/MaterializedViews.java b/database/src/main/java/io/hstore/db/view/MaterializedViews.java index 5bc1a6d..b31769d 100644 --- a/database/src/main/java/io/hstore/db/view/MaterializedViews.java +++ b/database/src/main/java/io/hstore/db/view/MaterializedViews.java @@ -118,7 +118,7 @@ public Descriptor refresh(Principal principal, String name) { } private Descriptor refresh(Writer writer, Descriptor descriptor) { - long target = database.engine().feed().lastGeneration(); + long target = Math.min(database.engine().feed().lastGeneration(), writer.generation()); if (target <= descriptor.lastGeneration()) { return descriptor; } diff --git a/docs/database/views-signals-stats.md b/docs/database/views-signals-stats.md index 209ee65..212c2db 100644 --- a/docs/database/views-signals-stats.md +++ b/docs/database/views-signals-stats.md @@ -22,7 +22,7 @@ sequenceDiagram TM->>CF: append CommitEvent(generation g, members, slots) CF-->>MV: onCommit(event) on main with member changes MV->>MV: for each CONTINUOUS view: refresh(name) - MV->>CF: replay(lastGeneration) up to feed.lastGeneration + MV->>CF: replay(lastGeneration) up to min(feed.lastGeneration, own snapshot) MV->>TM: write txn: apply deltas to view-data, lastGeneration := g TM->>CF: append slot-only event (no member changes, ignored by onCommit) ``` diff --git a/docs/storage/change-feed.md b/docs/storage/change-feed.md index 5cb3513..0cb2ac2 100644 --- a/docs/storage/change-feed.md +++ b/docs/storage/change-feed.md @@ -43,18 +43,19 @@ sequenceDiagram TM->>WAL: Commit(txn, g) TM->>G: submit(Pending g) G->>G: makeDurable (fsync data and WAL under SYNC) + G->>G: current = g, history G->>CF: append(event) CF-->>S: signal "appended" - G->>G: current = g, history (CATALOG_PUBLISH) + G->>G: CATALOG_PUBLISH S->>CF: replay(acknowledged) and consume ``` Two copies exist: 1. The **WAL `Feed` record** (type 5), written inside the commit before the `Commit` record. It is the durable copy and is covered by the commit's atomicity. -2. The **feed segments**, appended by the group-commit flusher during publication: after the commit is durable (under `SYNC`) and **immediately before** the generation becomes `current`. Appends happen on the single flusher thread in queue order, so the feed is in generation order. +2. The **feed segments**, appended by the group-commit flusher during publication: after the commit is durable (under `SYNC`) and **immediately after** the generation becomes `current`. Appends happen on the single flusher thread in queue order, so the feed is in generation order. -Appending before `current.set(next)` establishes the invariant that any reader which observes generation `g` as current also observes `feed.lastGeneration()` covering every non-empty commit `<= g`. A subscriber may therefore briefly see an event for generation `g` while `current` is still `g - 1`; the commit is already durable at that point, so the event will never be rolled back by a process crash. +Appending after `current.set(next)` guarantees that a subscriber told about generation `g` can read it: any transaction it starts sees `g` or later. Subscribers such as `MaterializedViews` look up the atoms an event mentions, so an event that arrived before its commit was visible would make them skip changes they can't yet see. The opposite direction is not guaranteed: a reader whose snapshot is `g` may find `feed.lastGeneration()` still at `g - 1` while `g` is being published. Code that combines a snapshot with the feed replays only up to `min(snapshot, feed.lastGeneration())` (`SemanticPlane`, `MaterializedViews.refresh`) and picks up the rest from the next event. A committing client is acknowledged only after both steps, so its own commit is always in the feed by then. The commit is durable before either step, so an event is never rolled back by a process crash. The feed is *not* forced on every append. The active segment is forced by `ChangeFeed.sync()` during every checkpoint (before the catalog image is published) and on `close()`; a segment is also forced when the feed rolls past it. After a crash, recovery re-appends the `Feed` records of every replayed commit (`ChangeFeed.append` ignores any generation `<= lastGeneration`, so events already present are not duplicated), then `truncateAfter(recovered generation)` removes events for commits that did not survive, and forces the active segment. The checkpoint ordering (`feed.sync()` before the catalog image that moves the WAL start) guarantees that events older than the checkpoint LSN are already durable in the feed when their WAL records are truncated. diff --git a/engine/src/main/java/io/hstore/engine/txn/TransactionManager.java b/engine/src/main/java/io/hstore/engine/txn/TransactionManager.java index f56fc91..443f9bc 100644 --- a/engine/src/main/java/io/hstore/engine/txn/TransactionManager.java +++ b/engine/src/main/java/io/hstore/engine/txn/TransactionManager.java @@ -416,11 +416,11 @@ private Appended append(Generation latest, Branch changed, long txnId, List { + current.set(next); + history.put(next.id(), next); if (!event.isEmpty()) { storage.feed().append(event); } - current.set(next); - history.put(next.id(), next); trimHistory(); storage.faults().reach(CrashPoint.CATALOG_PUBLISH); }); diff --git a/engine/src/test/java/io/hstore/engine/EngineTest.java b/engine/src/test/java/io/hstore/engine/EngineTest.java index f050551..761ea53 100644 --- a/engine/src/test/java/io/hstore/engine/EngineTest.java +++ b/engine/src/test/java/io/hstore/engine/EngineTest.java @@ -17,6 +17,7 @@ import io.hstore.engine.tree.TreeSchema; import io.hstore.engine.tree.ValueCodec; import io.hstore.engine.txn.CommitResult; +import io.hstore.engine.txn.Durability; import io.hstore.engine.txn.IncidentEdge; import io.hstore.engine.txn.Isolation; import io.hstore.engine.txn.Snapshot; @@ -37,6 +38,7 @@ import java.util.Set; import java.util.TreeSet; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.atomic.AtomicLong; import java.util.random.RandomGenerator; import java.util.random.RandomGeneratorFactory; import java.util.stream.LongStream; @@ -146,6 +148,45 @@ void feedRetentionReleasesWholeSegmentsBehindHolds() { } } + @Test + void subscribersSeeTheCommitTheyAreToldAbout() throws InterruptedException { + try (StorageEngine engine = StorageEngine.open(directory, small().withDurability(Durability.ASYNC))) { + long edge = engine.write(txn -> txn.createEdge(EDGE, EdgeKind.SET)); + long start = engine.transactions().current().id(); + List early = new CopyOnWriteArrayList<>(); + List delivered = new ArrayList<>(); + List subscriptions = new ArrayList<>(); + for (int s = 0; s < 4; s++) { + AtomicLong seen = new AtomicLong(); + delivered.add(seen); + subscriptions.add(engine.feed().subscribe(start, event -> { + if (engine.transactions().current().id() < event.generation()) { + early.add(event.generation()); + } + seen.set(event.generation()); + })); + } + try { + for (int i = 0; i < 500; i++) { + int index = i; + engine.write(txn -> { + txn.insert(edge, txn.createNode(NODE, "v" + index)); + return null; + }); + } + long last = engine.transactions().current().id(); + long deadline = System.currentTimeMillis() + 5000; + while (delivered.stream().anyMatch(seen -> seen.get() < last) && System.currentTimeMillis() < deadline) { + Thread.sleep(10); + } + delivered.forEach(seen -> assertEquals(last, seen.get())); + assertEquals(List.of(), early); + } finally { + subscriptions.forEach(ChangeFeed.Subscription::close); + } + } + } + @Test void subscribersBehindRetentionAreToldAboutTheGap() throws InterruptedException { EngineOptions options = small().withWalSegmentBytes(4096).withFeedRetention(10);