Skip to content

Commit 196da44

Browse files
committed
align and deduplicate item state restore logic
Signed-off-by: Mark Herwege <mark.herwege@telenet.be>
1 parent 5869383 commit 196da44

2 files changed

Lines changed: 117 additions & 72 deletions

File tree

bundles/org.openhab.core.persistence/src/main/java/org/openhab/core/persistence/internal/PersistenceManagerImpl.java

Lines changed: 101 additions & 72 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,7 @@
9494
* @author Jan N. Klug - Added time series support
9595
* @author Mark Herwege - Added restoring lastState, lastStateChange and lastStateUpdate
9696
* @author Mark Herwege - Make default strategy to be only a configuration suggestion
97+
* @author Mark Herwege - Fix and enhance handling of time series and external persistence updates
9798
*/
9899
@Component(immediate = true, service = PersistenceManager.class)
99100
@NonNullByDefault
@@ -402,7 +403,7 @@ public void timeSeriesUpdated(Item item, TimeSeries timeSeries) {
402403
ZonedDateTime lastStateUpdate = item.getLastStateUpdate();
403404
ZonedDateTime timestamp = s.timestamp().atZone(ZoneId.systemDefault());
404405
if (lastStateUpdate == null || timestamp.isAfter(lastStateUpdate)) {
405-
container.restoreItemState(item, timestamp, s.state());
406+
container.restoreItemStateFromTimeSeriesEntry(item, timestamp, s.state());
406407
}
407408
});
408409
}));
@@ -462,22 +463,22 @@ public void handleExternalPersistenceDataChange(PersistenceService persistenceSe
462463
persistenceServiceContainers.values().stream()
463464
.filter(container -> container.persistenceService.equals(persistenceService) && Stream
464465
.concat(container.getMatchingConfigurations(UPDATE),
465-
container.getMatchingConfigurations(FORECAST))
466+
Stream.concat(container.getMatchingConfigurations(CHANGE),
467+
container.getMatchingConfigurations(FORECAST)))
466468
.distinct().anyMatch(itemConf -> appliesToItem(itemConf, item)))
467469
.forEach(container -> {
468470
container.scheduleNextPersistedForecastForItem(item.getName());
469471
PersistedItem persistedItem = container.getPersistedItem(item);
470472
if (persistedItem != null) {
471-
container.restoreItemState(item, persistedItem);
473+
container.restoreItemStateFromPersistenceUpdate(item, persistedItem);
472474
}
473475
});
474476
}
475477

476-
private void storeInOtherServices(String excludeContainerId, Item item, State oldState) {
478+
private void storeInOtherServices(PersistenceService persistenceService, Item item, State oldState) {
477479
boolean changed = !item.getState().equals(oldState);
478480
persistenceServiceContainers.values().stream()
479-
.filter(container -> !container.getPersistenceService().getId().equals(excludeContainerId))
480-
.forEach(container -> {
481+
.filter(container -> !container.persistenceService.equals(persistenceService)).forEach(container -> {
481482
if (changed) {
482483
storeItem(container, item, PersistenceStrategy.Globals.CHANGE);
483484
}
@@ -585,7 +586,7 @@ public void addItem(Item item) {
585586
.anyMatch(configuration -> appliesToItem(configuration, item)))
586587
|| getMatchingConfigurations(FORECAST)
587588
.anyMatch(configuration -> appliesToItem(configuration, item))) {
588-
restoreItemStateIfPossible(item);
589+
restoreItemStateOnStartup(item);
589590
}
590591
if (getMatchingConfigurations(FORECAST).anyMatch(configuration -> appliesToItem(configuration, item))) {
591592
scheduleNextPersistedForecastForItem(item.getName());
@@ -600,51 +601,6 @@ public void removeItem(String itemName) {
600601
}
601602
}
602603

603-
private void restoreItemStateIfPossible(Item item) {
604-
PersistedItem persistedItem = getPersistedItem(item);
605-
if (persistedItem == null) {
606-
// in case of an exception or timeout, the safe caller returns null
607-
return;
608-
}
609-
GenericItem genericItem = (GenericItem) item;
610-
State state = item.getState();
611-
State lastState = null;
612-
ZonedDateTime lastStateUpdate = null;
613-
ZonedDateTime lastStateChange = null;
614-
if (UnDefType.NULL.equals(state)) {
615-
state = persistedItem.getState();
616-
lastState = persistedItem.getLastState();
617-
lastStateUpdate = persistedItem.getTimestamp();
618-
lastStateChange = persistedItem.getLastStateChange();
619-
} else {
620-
// someone else already restored the state or a new state was set
621-
// try restoring the previous state if not yet set
622-
if (item.getLastState() != null && item.getLastState() != UnDefType.NULL) {
623-
// there is already a previous state, nothing to restore
624-
return;
625-
}
626-
lastStateUpdate = item.getLastStateUpdate();
627-
if (state.equals(persistedItem.getState())) {
628-
lastState = persistedItem.getLastState();
629-
lastStateChange = persistedItem.getLastStateChange();
630-
} else {
631-
lastState = persistedItem.getState();
632-
lastStateChange = item.getLastStateChange();
633-
}
634-
}
635-
genericItem.removeStateChangeListener(PersistenceManagerImpl.this);
636-
try {
637-
genericItem.setState(state, lastState, lastStateUpdate, lastStateChange, PERSISTENCE_SOURCE);
638-
} finally {
639-
genericItem.addStateChangeListener(PersistenceManagerImpl.this);
640-
}
641-
if (logger.isDebugEnabled()) {
642-
logger.debug("Restored item state from '{}' for item '{}' -> '{}'",
643-
DateTimeFormatter.ISO_ZONED_DATE_TIME.format(persistedItem.getTimestamp()), item.getName(),
644-
persistedItem.getState());
645-
}
646-
}
647-
648604
private @Nullable PersistedItem getPersistedItem(Item item) {
649605
QueryablePersistenceService queryService = (QueryablePersistenceService) persistenceService;
650606
String alias = getAlias(item);
@@ -666,8 +622,8 @@ public void scheduleNextForecastForItem(Item item, Instant time, State state) {
666622
if (oldJob != null) {
667623
oldJob.cancel(true);
668624
}
669-
forecastJobs.put(itemName,
670-
scheduler.at(() -> restoreItemState(item, time.atZone(ZoneId.systemDefault()), state), time));
625+
forecastJobs.put(itemName, scheduler.at(
626+
() -> restoreItemStateFromTimeSeriesEntry(item, time.atZone(ZoneId.systemDefault()), state), time));
671627
logger.trace("Scheduled forecasted value for {} at {}", item.getName(), time);
672628
}
673629

@@ -695,7 +651,34 @@ public void scheduleNextPersistedForecastForItem(String itemName) {
695651
}
696652
}
697653

698-
private void restoreItemState(Item item, ZonedDateTime timestamp, State state) {
654+
private void restoreItemStateOnStartup(Item item) {
655+
PersistedItem persistedItem = getPersistedItem(item);
656+
if (persistedItem == null) {
657+
// in case of an exception or timeout, the safe caller returns null
658+
return;
659+
}
660+
661+
PersistedItem newItemState = itemState(item, persistedItem);
662+
if (newItemState == null) {
663+
return;
664+
}
665+
666+
GenericItem genericItem = (GenericItem) item;
667+
genericItem.removeStateChangeListener(PersistenceManagerImpl.this);
668+
try {
669+
genericItem.setState(newItemState.getState(), newItemState.getLastState(), newItemState.getTimestamp(),
670+
newItemState.getLastStateChange(), PERSISTENCE_SOURCE);
671+
} finally {
672+
genericItem.addStateChangeListener(PersistenceManagerImpl.this);
673+
}
674+
if (logger.isDebugEnabled()) {
675+
logger.debug("Restored item state from '{}' for item '{}' -> '{}'",
676+
DateTimeFormatter.ISO_ZONED_DATE_TIME.format(persistedItem.getTimestamp()), item.getName(),
677+
persistedItem.getState());
678+
}
679+
}
680+
681+
private void restoreItemStateFromTimeSeriesEntry(Item item, ZonedDateTime timestamp, State state) {
699682
PersistedItem persistedItem = new PersistedItem() {
700683

701684
@Override
@@ -723,23 +706,46 @@ public String getName() {
723706
return null;
724707
}
725708
};
726-
restoreItemState(item, persistedItem);
709+
restoreItemStateFromPersistenceUpdate(item, persistedItem);
727710
scheduleNextPersistedForecastForItem(item.getName());
728711
}
729712

730-
private void restoreItemState(Item item, PersistedItem persistedItem) {
713+
private void restoreItemStateFromPersistenceUpdate(Item item, PersistedItem persistedItem) {
714+
PersistedItem newItemState = itemState(item, persistedItem);
715+
if (newItemState == null) {
716+
return;
717+
}
718+
731719
GenericItem genericItem = (GenericItem) item;
720+
genericItem.removeStateChangeListener(PersistenceManagerImpl.this);
721+
try {
722+
genericItem.setState(newItemState.getState(), newItemState.getLastState(), newItemState.getTimestamp(),
723+
newItemState.getLastStateChange(), PERSISTENCE_SOURCE);
724+
// other services with update or change strategy should persist new state
725+
storeInOtherServices(persistenceService, item, newItemState.getState());
726+
} finally {
727+
genericItem.addStateChangeListener(PersistenceManagerImpl.this);
728+
}
729+
if (logger.isDebugEnabled()) {
730+
logger.debug("Reset item state from '{}' for item '{}' -> '{}'",
731+
DateTimeFormatter.ISO_ZONED_DATE_TIME.format(persistedItem.getTimestamp()), item.getName(),
732+
persistedItem.getState());
733+
}
734+
}
735+
736+
private @Nullable PersistedItem itemState(Item item, PersistedItem persistedItem) {
732737
ZonedDateTime itemLastStateUpdate = item.getLastStateUpdate();
733738
ZonedDateTime persistedItemTimestamp = persistedItem.getTimestamp();
739+
State itemState = item.getState();
734740

735-
if (itemLastStateUpdate != null && persistedItemTimestamp.isBefore(itemLastStateUpdate)) {
736-
return;
741+
if (itemState != UnDefType.NULL && itemLastStateUpdate != null
742+
&& persistedItemTimestamp.isBefore(itemLastStateUpdate)) {
743+
return null;
737744
}
738745

739746
State persistedItemState = persistedItem.getState();
740747
State persistedItemLastState = persistedItem.getLastState();
741748
ZonedDateTime persistedItemLastStateChange = persistedItem.getLastStateChange();
742-
State itemState = item.getState();
743749
State itemLastState = item.getLastState();
744750
ZonedDateTime itemLastStateChange = item.getLastStateChange();
745751

@@ -756,8 +762,9 @@ private void restoreItemState(Item item, PersistedItem persistedItem) {
756762
: itemLastStateChange;
757763
} else {
758764
state = persistedItemState;
759-
if (persistedItemLastStateChange != null && persistedItemLastState != null
760-
&& itemLastStateUpdate != null && persistedItemLastStateChange.isAfter(itemLastStateUpdate)) {
765+
if (itemState == UnDefType.NULL || (persistedItemLastStateChange != null
766+
&& persistedItemLastState != null && itemLastStateUpdate != null
767+
&& persistedItemLastStateChange.isAfter(itemLastStateUpdate))) {
761768
lastState = persistedItemLastState;
762769
lastStateChange = persistedItemLastStateChange;
763770
} else {
@@ -768,17 +775,39 @@ private void restoreItemState(Item item, PersistedItem persistedItem) {
768775

769776
// Check again if item has not been updated in the mean time before commit
770777
itemLastStateUpdate = item.getLastStateUpdate();
771-
if (itemLastStateUpdate == null
772-
|| itemLastStateUpdate.toInstant().compareTo(lastStateUpdate.toInstant()) <= 0) {
773-
genericItem.removeStateChangeListener(PersistenceManagerImpl.this);
774-
try {
775-
genericItem.setState(state, lastState, lastStateUpdate, lastStateChange, PERSISTENCE_SOURCE);
776-
// other services with update or change strategy should persist new state
777-
storeInOtherServices(getPersistenceService().getId(), item, itemState);
778-
} finally {
779-
genericItem.addStateChangeListener(PersistenceManagerImpl.this);
780-
}
778+
itemState = item.getState();
779+
if (itemState != UnDefType.NULL && itemLastStateUpdate != null
780+
&& persistedItemTimestamp.isBefore(itemLastStateUpdate)) {
781+
return null;
781782
}
783+
784+
return new PersistedItem() {
785+
786+
@Override
787+
public ZonedDateTime getTimestamp() {
788+
return lastStateUpdate;
789+
}
790+
791+
@Override
792+
public State getState() {
793+
return state;
794+
}
795+
796+
@Override
797+
public String getName() {
798+
return item.getName();
799+
}
800+
801+
@Override
802+
public @Nullable ZonedDateTime getLastStateChange() {
803+
return lastStateChange;
804+
}
805+
806+
@Override
807+
public @Nullable State getLastState() {
808+
return lastState;
809+
}
810+
};
782811
}
783812

784813
private void persistJob(List<PersistenceItemConfiguration> itemConfigs) {

bundles/org.openhab.core.persistence/src/test/java/org/openhab/core/persistence/internal/PersistenceManagerTest.java

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -498,6 +498,22 @@ public void storeTimeSeriesAndForecastsScheduled() {
498498
inOrder.verify(schedulerMock).at(any(SchedulerRunnable.class), eq(time5));
499499
}
500500

501+
@Test
502+
public void externalPersistenceDataChangeIsHandled() {
503+
setupPersistence(new PersistenceAllConfig());
504+
addConfiguration(TEST_PERSISTENCE_SERVICE_ID, List.of(new PersistenceAllConfig()),
505+
PersistenceStrategy.Globals.UPDATE, null);
506+
addConfiguration(TEST_QUERYABLE_PERSISTENCE_SERVICE_ID, List.of(new PersistenceAllConfig()),
507+
PersistenceStrategy.Globals.UPDATE, null);
508+
manager.handleExternalPersistenceDataChange(persistenceServiceMock, TEST_ITEM);
509+
assertNotEquals(TEST_STATE, TEST_ITEM.getState());
510+
511+
manager.handleExternalPersistenceDataChange(queryablePersistenceServiceMock, TEST_ITEM);
512+
verify(queryablePersistenceServiceMock).persistedItem(eq(TEST_ITEM_NAME), any());
513+
assertEquals(TEST_STATE, TEST_ITEM.getState());
514+
verify(persistenceServiceMock).store(TEST_ITEM, null);
515+
}
516+
501517
@Test
502518
public void cronStrategyIsScheduledAndCancelledAndPersistsValue() throws Exception {
503519
ArgumentCaptor<SchedulerRunnable> runnableCaptor = ArgumentCaptor.forClass(SchedulerRunnable.class);

0 commit comments

Comments
 (0)