Skip to content

Commit b35f482

Browse files
committed
overload AbstractHollowProducer.restore to take in a hollow read state engine;
1 parent 2e99fd8 commit b35f482

3 files changed

Lines changed: 114 additions & 11 deletions

File tree

hollow/src/main/java/com/netflix/hollow/api/producer/AbstractHollowProducer.java

Lines changed: 24 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -291,26 +291,43 @@ public void initializeDataModel(HollowSchema... schemas) {
291291
* @see #initializeDataModel(Class[])
292292
*/
293293
public HollowProducer.ReadState restore(long versionDesired, HollowConsumer.BlobRetriever blobRetriever) {
294-
return restore(new HollowConsumer.VersionInfo(versionDesired), blobRetriever,
294+
return restore(versionDesired,
295+
() -> retrieveReadState(versionDesired, blobRetriever, this.updatePlanBlobVerifier),
295296
(restoreFrom, restoreTo) -> restoreTo.restoreFrom(restoreFrom));
296297
}
297298

298299
public HollowProducer.ReadState restore(HollowConsumer.VersionInfo versionDesired, HollowConsumer.BlobRetriever blobRetriever) {
299-
return restore(versionDesired, blobRetriever,
300+
return restore(versionDesired.getVersion(), blobRetriever);
301+
}
302+
303+
public HollowProducer.ReadState restore(long readStateEngineVersion, HollowReadStateEngine readStateEngine) {
304+
return restore(readStateEngineVersion,
305+
() -> ReadStateHelper.newReadState(readStateEngineVersion, readStateEngine),
300306
(restoreFrom, restoreTo) -> restoreTo.restoreFrom(restoreFrom));
301307
}
302308

309+
303310
HollowProducer.ReadState hardRestore(long versionDesired, HollowConsumer.BlobRetriever blobRetriever) {
304-
return restore(new HollowConsumer.VersionInfo(versionDesired), blobRetriever,
311+
return restore(versionDesired,
312+
() -> retrieveReadState(versionDesired, blobRetriever, this.updatePlanBlobVerifier),
305313
(restoreFrom, restoreTo) -> HollowWriteStateCreator.
306314
populateUsingReadEngine(restoreTo, restoreFrom, false));
307315
}
308316

317+
private static HollowProducer.ReadState retrieveReadState(long versionDesired, HollowConsumer.BlobRetriever blobRetriever,
318+
HollowConsumer.UpdatePlanBlobVerifier updatePlanBlobVerifier) {
319+
Objects.requireNonNull(blobRetriever);
320+
HollowConsumer client = HollowConsumer.withBlobRetriever(blobRetriever)
321+
.withUpdatePlanVerifier(updatePlanBlobVerifier)
322+
.build();
323+
client.triggerRefreshTo(new HollowConsumer.VersionInfo(versionDesired));
324+
return ReadStateHelper.newReadState(client.getCurrentVersionId(), client.getStateEngine());
325+
}
326+
309327
private HollowProducer.ReadState restore(
310-
HollowConsumer.VersionInfo versionInfoDesired, HollowConsumer.BlobRetriever blobRetriever,
328+
long versionDesired,Supplier<HollowProducer.ReadState> readStateSupplier,
311329
BiConsumer<HollowReadStateEngine, HollowWriteStateEngine> restoreAction) {
312-
long versionDesired = versionInfoDesired.getVersion();
313-
Objects.requireNonNull(blobRetriever);
330+
Objects.requireNonNull(readStateSupplier);
314331
Objects.requireNonNull(restoreAction);
315332

316333
if (!isInitialized) {
@@ -323,11 +340,7 @@ private HollowProducer.ReadState restore(
323340
Status.RestoreStageBuilder status = localListeners.fireProducerRestoreStart(versionDesired);
324341
try {
325342
if (versionDesired != HollowConstants.VERSION_NONE) {
326-
HollowConsumer client = HollowConsumer.withBlobRetriever(blobRetriever)
327-
.withUpdatePlanVerifier(updatePlanBlobVerifier)
328-
.build();
329-
client.triggerRefreshTo(versionInfoDesired);
330-
readState = ReadStateHelper.newReadState(client.getCurrentVersionId(), client.getStateEngine());
343+
readState = readStateSupplier.get();
331344
readStates = ReadStateHelper.restored(readState);
332345

333346
// Need to restore data to new ObjectMapper since can't restore to non empty Write State Engine

hollow/src/main/java/com/netflix/hollow/api/producer/HollowProducer.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,6 +220,11 @@ public HollowProducer.ReadState restore(HollowConsumer.VersionInfo versionDesire
220220
return super.restore(versionDesired, blobRetriever);
221221
}
222222

223+
@Override
224+
public HollowProducer.ReadState restore(long readStateEngineVersion, HollowReadStateEngine readStateEngine) {
225+
return super.restore(readStateEngineVersion, readStateEngine);
226+
}
227+
223228
@Override
224229
public HollowWriteStateEngine getWriteEngine() {
225230
return super.getWriteEngine();

hollow/src/test/java/com/netflix/hollow/api/producer/HollowProducerTest.java

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import static org.mockito.Mockito.doThrow;
2525
import static org.mockito.Mockito.mock;
2626
import static org.mockito.Mockito.spy;
27+
import static org.mockito.Mockito.verifyZeroInteractions;
2728
import static org.mockito.Mockito.when;
2829

2930
import com.netflix.hollow.api.consumer.HollowConsumer;
@@ -51,7 +52,9 @@
5152
import com.netflix.hollow.api.producer.model.CustomReferenceType;
5253
import com.netflix.hollow.api.producer.model.HasAllTypeStates;
5354
import com.netflix.hollow.core.HollowBlobHeader;
55+
import com.netflix.hollow.core.HollowConstants;
5456
import com.netflix.hollow.core.read.engine.HollowBlobHeaderReader;
57+
import com.netflix.hollow.core.read.engine.HollowReadStateEngine;
5558
import com.netflix.hollow.core.read.engine.HollowTypeReadState;
5659
import com.netflix.hollow.core.read.engine.object.HollowObjectTypeReadState;
5760
import com.netflix.hollow.core.schema.HollowMapSchema;
@@ -314,6 +317,88 @@ public void testPublishAndRestore() {
314317
assertEquals(producer.getCycleCountWithPrimaryStatus(), 1);
315318
}
316319

320+
@Test
321+
public void testRestoreFromReadStateEngine() {
322+
InMemoryBlobStore blobStore = new InMemoryBlobStore();
323+
HollowInMemoryBlobStager blobStager = new HollowInMemoryBlobStager();
324+
325+
// A producer publishes an initial version.
326+
HollowProducer producer1 = HollowProducer.withPublisher(blobStore)
327+
.withBlobStager(blobStager)
328+
.build();
329+
producer1.initializeDataModel(TestPojoV1.class);
330+
long version = producer1.runCycle(ws -> {
331+
for (int i = 1; i <= 3; i++) {
332+
ws.add(new TestPojoV1(i, i * 10));
333+
}
334+
});
335+
336+
// A caller builds their own HollowConsumer and refreshes it to that version.
337+
HollowConsumer consumer = HollowConsumer.withBlobRetriever(blobStore).build();
338+
consumer.triggerRefreshTo(version);
339+
assertEquals(version, consumer.getCurrentVersionId());
340+
341+
// A fresh producer restores directly from the consumer's read state engine
342+
// rather than downloading blobs itself.
343+
HollowProducer producer2 = HollowProducer.withPublisher(blobStore)
344+
.withBlobStager(blobStager)
345+
.build();
346+
producer2.initializeDataModel(TestPojoV1.class);
347+
HollowProducer.ReadState readState =
348+
producer2.restore(consumer.getCurrentVersionId(), consumer.getStateEngine());
349+
350+
// The producer is restored to the expected version with the expected data.
351+
Assert.assertNotNull(readState);
352+
assertEquals(version, readState.getVersion());
353+
HollowObjectTypeReadState typeState =
354+
(HollowObjectTypeReadState) readState.getStateEngine().getTypeState("TestPojo");
355+
BitSet populatedOrdinals = typeState.getPopulatedOrdinals();
356+
assertEquals(3, populatedOrdinals.cardinality());
357+
int ordinal = populatedOrdinals.nextSetBit(0);
358+
while (ordinal != -1) {
359+
GenericHollowObject obj = new GenericHollowObject(new HollowObjectGenericDelegate(typeState), ordinal);
360+
assertEquals("v1 should equal id * 10", obj.getInt("id") * 10, obj.getInt("v1"));
361+
ordinal = populatedOrdinals.nextSetBit(ordinal + 1);
362+
}
363+
364+
// Restore -> next delta: the cycle after restore must chain off the restored version,
365+
// so a consumer that loaded the restored version can move forward via a delta. If the
366+
// delta's "from" did not match the restored version the transition below would fail with
367+
// "Attempting to apply a delta to a state from which it was not originated!".
368+
long nextVersion = producer2.runCycle(ws -> {
369+
for (int i = 1; i <= 4; i++) {
370+
ws.add(new TestPojoV1(i, i * 10));
371+
}
372+
});
373+
Assert.assertNotEquals(version, nextVersion);
374+
375+
HollowConsumer follower = HollowConsumer.withBlobRetriever(blobStore).build();
376+
follower.triggerRefreshTo(version);
377+
assertEquals(version, follower.getCurrentVersionId());
378+
follower.triggerRefreshTo(nextVersion);
379+
assertEquals(nextVersion, follower.getCurrentVersionId());
380+
assertEquals(4, follower.getStateEngine().getTypeState("TestPojo")
381+
.getPopulatedOrdinals().cardinality());
382+
}
383+
384+
@Test
385+
public void testRestoreVersionNoneShortCircuits() {
386+
HollowProducer producer = createProducer(tmpFolder, schema);
387+
388+
// BlobRetriever overload: a VERSION_NONE restore must short-circuit, returning null
389+
// without ever building a consumer or touching the retriever.
390+
HollowConsumer.BlobRetriever retriever = mock(HollowConsumer.BlobRetriever.class);
391+
HollowProducer.ReadState readState = producer.restore(HollowConstants.VERSION_NONE, retriever);
392+
Assert.assertNull("VERSION_NONE restore should return null", readState);
393+
verifyZeroInteractions(retriever);
394+
395+
// ReadStateEngine overload: likewise returns null and never dereferences the engine,
396+
// so a null engine is tolerated when the version is VERSION_NONE.
397+
HollowProducer.ReadState readStateFromEngine =
398+
producer.restore(HollowConstants.VERSION_NONE, (HollowReadStateEngine) null);
399+
Assert.assertNull("VERSION_NONE restore should return null", readStateFromEngine);
400+
}
401+
317402
@Test
318403
public void testHeaderPublish() throws IOException {
319404
HollowProducer producer = createProducer(tmpFolder, schema);

0 commit comments

Comments
 (0)