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
4 changes: 2 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
2 changes: 1 addition & 1 deletion docs/database/views-signals-stats.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
```
Expand Down
7 changes: 4 additions & 3 deletions docs/storage/change-feed.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -416,11 +416,11 @@ private Appended append(Generation latest, Branch changed, long txnId, List<Memb
pagesWritten.add(materializer.pagesWritten());
walBytes.add(written);
GroupCommitter.Pending pending = new GroupCommitter.Pending(next, () -> {
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);
});
Expand Down
41 changes: 41 additions & 0 deletions engine/src/test/java/io/hstore/engine/EngineTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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<Long> early = new CopyOnWriteArrayList<>();
List<AtomicLong> delivered = new ArrayList<>();
List<ChangeFeed.Subscription> 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);
Expand Down
Loading