From ddc34d04f738c14b244d321a4b8f1086af48deb2 Mon Sep 17 00:00:00 2001 From: npostnikova Date: Thu, 9 Jun 2022 00:33:24 +0300 Subject: [PATCH 1/5] Send logs one by one --- .../ganttproject/GanttProject.java | 22 ++++++++++++++----- .../ganttproject/GanttProjectBase.java | 11 +++++++++- .../storage/SqlProjectDatabaseImpl.kt | 22 +++++++++++++++++++ 3 files changed, 49 insertions(+), 6 deletions(-) diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java index 2ff84c2f79..8505d38c71 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java @@ -86,6 +86,7 @@ import java.util.Arrays; import java.util.HashMap; import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; @@ -208,10 +209,6 @@ public GanttProject(boolean isOnlyViewer) { getWebSocket().register(null); getWebSocket().onCommitResponseReceived(this::fireXlogReceived); getWebSocket().onBaseTxnIdReceived(this::onBaseTxnIdReceived); - var taskListenerAdapter = new TaskListenerAdapter(); - // TODO: add listeners sensibly. - taskListenerAdapter.setTaskAddedHandler(event -> this.sendProjectStateLogs()); - getTaskManager().addTaskListener(taskListenerAdapter); } area = new GanttGraphicArea(this, getTaskManager(), getZoomManager(), getUndoManager(), @@ -956,7 +953,8 @@ public void refresh() { super.repaint(); } - // TODO: Accumulate changes instead of sending it every time. + private final AtomicBoolean isSendingInProgress = new AtomicBoolean(); + private Unit sendProjectStateLogs() { gpLogger.debug("Sending project state logs"); try { @@ -969,20 +967,34 @@ private Unit sendProjectStateLogs() { "refid", txns )); + isSendingInProgress.set(true); } } catch (ProjectDatabaseException e) { gpLogger.error("Failed to send logs", new Object[]{}, ImmutableMap.of(), e); + isSendingInProgress.set(false); + } + return Unit.INSTANCE; + } + + @Override + protected Unit onProjectLogUpdate() { + if (isColloboqueLocalTest()) { + super.onProjectLogUpdate(); + if (!isSendingInProgress.get()) sendProjectStateLogs(); } return Unit.INSTANCE; } private Unit fireXlogReceived(ServerCommitResponse response) { myBaseTxnCommitInfo.update(response.getBaseTxnId(), response.getNewBaseTxnId(), 1); + isSendingInProgress.set(false); + sendProjectStateLogs(); return Unit.INSTANCE; } private Unit onBaseTxnIdReceived(String baseTxnId) { myBaseTxnCommitInfo.update("", baseTxnId, 0); + sendProjectStateLogs(); return Unit.INSTANCE; } } diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java index e02472d5d7..a0da77ad76 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java @@ -40,6 +40,7 @@ of the License, or (at your option) any later version. import javafx.beans.property.SimpleIntegerProperty; import javafx.beans.property.SimpleObjectProperty; import javafx.collections.FXCollections; +import kotlin.Unit; import kotlin.jvm.functions.Function0; import net.sourceforge.ganttproject.chart.Chart; import net.sourceforge.ganttproject.chart.ChartModelBase; @@ -79,6 +80,8 @@ of the License, or (at your option) any later version. import java.util.Map; import java.util.function.Supplier; +import static biz.ganttproject.storage.cloud.GPCloudHttpImplKt.isColloboqueLocalTest; + /** * This class is designed to be a GanttProject-after-refactorings. I am going to * refactor GanttProject in order to make true view communicating with other @@ -192,7 +195,10 @@ GPOptionGroup getTaskOptions() { protected GanttProjectBase() { super("GanttProject"); - var databaseProxy = new LazyProjectDatabaseProxy(SqlProjectDatabaseImpl.Factory::createInMemoryDatabase, this::getTaskManager); + var databaseProxy = new LazyProjectDatabaseProxy( + () -> SqlProjectDatabaseImpl.Factory.createInMemoryDatabase(this::onProjectLogUpdate), + this::getTaskManager + ); myProjectDatabase = databaseProxy; myTaskManagerConfig = new TaskManagerConfigImpl(); @@ -258,6 +264,9 @@ protected ParserFactory getParserFactory() { protected GanttProjectImpl getProjectImpl() { return myProjectImpl; } + + protected Unit onProjectLogUpdate() { return Unit.INSTANCE; } + @Override public void restore(@NotNull Document fromDocument) throws Document.DocumentException, IOException { GanttProjectImplKt.restoreProject(this, fromDocument, myProjectImpl.getListeners()); diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt b/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt index 4bcb66583c..add16d8f1f 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt @@ -47,12 +47,33 @@ class SqlProjectDatabaseImpl(private val dataSource: DataSource) : ProjectDataba dataSource.setURL(H2_IN_MEMORY_URL) return SqlProjectDatabaseImpl(dataSource) } + + fun createInMemoryDatabase(logUpdateCallback: () -> Unit): ProjectDatabase { + val dataSource = JdbcDataSource() + dataSource.setURL(H2_IN_MEMORY_URL) + val database = SqlProjectDatabaseImpl(dataSource) + database.addLogUpdateCallback(logUpdateCallback) + return database + } } /** Queries which belong to the current transaction. Null if each statement should be committed separately. */ private var currentTxn: TransactionImpl? = null private var localTxnId: Int = 1 + private val logUpdateCallbacks: MutableList<() -> Unit> = mutableListOf() + + /** Log update callbacks are invoked when a new log record is added. */ + fun addLogUpdateCallback(listener: () -> Unit) = logUpdateCallbacks.add(listener) + + private fun onLogUpdate() = logUpdateCallbacks.forEach { + try { + it.invoke() + } catch (e: Exception) { + LOG.error("Failed to execute update callback", e) + } + } + private fun withDSL( errorMessage: () -> String = { "Failed to execute query" }, body: (dsl: DSLContext) -> T @@ -89,6 +110,7 @@ class SqlProjectDatabaseImpl(private val dataSource: DataSource) : ProjectDataba throw ProjectDatabaseException(it.errorMessage(), e) } } + onLogUpdate() } } From c4b36bc33d5f4c24112cd94369fa59ca001fb82b Mon Sep 17 00:00:00 2001 From: npostnikova Date: Mon, 13 Jun 2022 21:19:38 +0300 Subject: [PATCH 2/5] Avoid committing empty txns --- .../sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt b/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt index add16d8f1f..34b31ec644 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/SqlProjectDatabaseImpl.kt @@ -110,8 +110,8 @@ class SqlProjectDatabaseImpl(private val dataSource: DataSource) : ProjectDataba throw ProjectDatabaseException(it.errorMessage(), e) } } - onLogUpdate() } + onLogUpdate() } /** Add a query to the current txn. Executes immediately if no transaction started. */ @@ -177,6 +177,7 @@ class SqlProjectDatabaseImpl(private val dataSource: DataSource) : ProjectDataba @Throws(ProjectDatabaseException::class) internal fun commitTransaction(txn: TransactionImpl) { try { + if (txn.statements.isEmpty()) return executeAndLog(txn.statements, localTxnId) localTxnId++ // Increment only on success. } finally { From 86c319fdaf2d15956efaf1b9d283ca12834e0947 Mon Sep 17 00:00:00 2001 From: npostnikova Date: Wed, 15 Jun 2022 18:14:23 +0300 Subject: [PATCH 3/5] Short-living listener for sending txns --- .../ganttproject/GanttProject.java | 38 ++++++++++++------- 1 file changed, 25 insertions(+), 13 deletions(-) diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java index 8505d38c71..f47f4ec34e 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java @@ -71,7 +71,6 @@ import net.sourceforge.ganttproject.storage.ServerCommitResponse; import net.sourceforge.ganttproject.task.CustomColumnsStorage; import net.sourceforge.ganttproject.task.Task; -import net.sourceforge.ganttproject.task.event.TaskListenerAdapter; import net.sourceforge.ganttproject.undo.GPUndoListener; import org.apache.commons.lang3.tuple.ImmutablePair; import org.jetbrains.annotations.NotNull; @@ -86,7 +85,6 @@ import java.util.Arrays; import java.util.HashMap; import java.util.List; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; @@ -953,25 +951,36 @@ public void refresh() { super.repaint(); } - private final AtomicBoolean isSendingInProgress = new AtomicBoolean(); + private interface TxnSendListener { + void onSendCompleted(); + } + + private final AtomicReference txnSendingListener = new AtomicReference<>(); private Unit sendProjectStateLogs() { gpLogger.debug("Sending project state logs"); + if (txnSendingListener.get() != null) return Unit.INSTANCE; try { var baseTxnCommitInfo = myBaseTxnCommitInfo.get(); var txns = myProjectDatabase.fetchTransactions(baseTxnCommitInfo.right + 1, 1); if (!txns.isEmpty()) { - getWebSocket().sendLogs(new InputXlog( - baseTxnCommitInfo.left, - "userId", - "refid", - txns - )); - isSendingInProgress.set(true); + var listener = new TxnSendListener() { + @Override + public void onSendCompleted() { + txnSendingListener.compareAndSet(this, null); + } + }; + if (txnSendingListener.compareAndSet(null, listener)) { + getWebSocket().sendLogs(new InputXlog( + baseTxnCommitInfo.left, + "userId", + "refid", + txns + )); + } } } catch (ProjectDatabaseException e) { gpLogger.error("Failed to send logs", new Object[]{}, ImmutableMap.of(), e); - isSendingInProgress.set(false); } return Unit.INSTANCE; } @@ -980,20 +989,23 @@ private Unit sendProjectStateLogs() { protected Unit onProjectLogUpdate() { if (isColloboqueLocalTest()) { super.onProjectLogUpdate(); - if (!isSendingInProgress.get()) sendProjectStateLogs(); + sendProjectStateLogs(); } return Unit.INSTANCE; } private Unit fireXlogReceived(ServerCommitResponse response) { myBaseTxnCommitInfo.update(response.getBaseTxnId(), response.getNewBaseTxnId(), 1); - isSendingInProgress.set(false); + txnSendingListener.get().onSendCompleted(); sendProjectStateLogs(); return Unit.INSTANCE; } private Unit onBaseTxnIdReceived(String baseTxnId) { myBaseTxnCommitInfo.update("", baseTxnId, 0); + var listener = txnSendingListener.get(); + // Websocket is [re-]started. Previous messages are discarded. + if (listener != null) listener.onSendCompleted(); sendProjectStateLogs(); return Unit.INSTANCE; } From df07f915cff0288ed3114ee9ee87e3d954f128fc Mon Sep 17 00:00:00 2001 From: npostnikova Date: Thu, 16 Jun 2022 12:20:31 +0300 Subject: [PATCH 4/5] Fixes --- .../colloboque/ColloboqueServer.kt | 2 +- .../ganttproject/colloboque/DevServer.kt | 4 +- .../ganttproject/GanttProject.java | 69 +++++++++++-------- .../ganttproject/GanttProjectBase.java | 1 - 4 files changed, 45 insertions(+), 31 deletions(-) diff --git a/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/ColloboqueServer.kt b/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/ColloboqueServer.kt index a85bb1f05c..d52d31c414 100644 --- a/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/ColloboqueServer.kt +++ b/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/ColloboqueServer.kt @@ -91,7 +91,7 @@ class ColloboqueServer( ) } } - refidToBaseTxnId[projectRefid] = "abacaba" // TODO: get from the database + refidToBaseTxnId[projectRefid] = EMPTY_LOG_BASE_TXN_ID // TODO: get from the database } } catch (e: Exception) { throw ColloboqueServerException("Failed to init project $projectRefid", e) diff --git a/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/DevServer.kt b/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/DevServer.kt index ee8f64a53d..9f6ff22bc4 100644 --- a/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/DevServer.kt +++ b/cloud.ganttproject.colloboque/src/main/kotlin/cloud/ganttproject/colloboque/DevServer.kt @@ -123,8 +123,8 @@ class ColloboqueWebSocketServer(port: Int, private val colloboqueServer: Collobo } override fun onMessage(message: WebSocketFrame) { - LOG.debug("Message received\n {}", message.textPayload) val inputXlog = parseInputXlog(message.textPayload) ?: return + LOG.debug("Xlog received\n {}", inputXlog) if (inputXlog.transactions.size != 1) { // TODO: add multiple transactions support. LOG.error("Only single transaction commit supported") @@ -138,7 +138,7 @@ class ColloboqueWebSocketServer(port: Int, private val colloboqueServer: Collobo override fun onPong(pong: WebSocketFrame?) {} override fun onException(exception: IOException) { - LOG.error("WebSocket exception", exception) + LOG.error("WebSocket exception", exception = exception) } } } diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java index f47f4ec34e..00a5efb449 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProject.java @@ -69,6 +69,7 @@ import net.sourceforge.ganttproject.storage.InputXlog; import net.sourceforge.ganttproject.storage.ProjectDatabaseException; import net.sourceforge.ganttproject.storage.ServerCommitResponse; +import net.sourceforge.ganttproject.storage.XlogRecord; import net.sourceforge.ganttproject.task.CustomColumnsStorage; import net.sourceforge.ganttproject.task.Task; import net.sourceforge.ganttproject.undo.GPUndoListener; @@ -81,15 +82,14 @@ import java.awt.*; import java.awt.event.*; import java.io.IOException; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.HashMap; +import java.util.*; import java.util.List; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; import static biz.ganttproject.storage.cloud.GPCloudHttpImplKt.getWebSocket; import static biz.ganttproject.storage.cloud.GPCloudHttpImplKt.isColloboqueLocalTest; +import static net.sourceforge.ganttproject.storage.XlogKt.EMPTY_LOG_BASE_TXN_ID; /** * Main frame of the project @@ -166,7 +166,7 @@ private static class TxnCommitInfo { /** If `oldTxnId` is currently being hold, sets the txn ID to `newTxnId` and moves the local ID ahead by `committedNum`. */ void update(String oldTxnId, String newTxnId, int committedNum) { myTxnId.updateAndGet(oldValue -> { - if (oldValue.left.equals(oldTxnId)) { + if (Objects.equals(oldValue.left, oldTxnId)) { return new ImmutablePair<>(newTxnId, oldValue.right + committedNum); } else { return oldValue; @@ -177,9 +177,13 @@ void update(String oldTxnId, String newTxnId, int committedNum) { ImmutablePair get() { return myTxnId.get(); } + + void reset() { + myTxnId.set(new ImmutablePair<>(null, 0)); + } } - private final TxnCommitInfo myBaseTxnCommitInfo = new TxnCommitInfo("", 0); + private final TxnCommitInfo myBaseTxnCommitInfo = new TxnCommitInfo(null, 0); public GanttProject(boolean isOnlyViewer) { @@ -960,27 +964,35 @@ private interface TxnSendListener { private Unit sendProjectStateLogs() { gpLogger.debug("Sending project state logs"); if (txnSendingListener.get() != null) return Unit.INSTANCE; + var baseTxnCommitInfo = myBaseTxnCommitInfo.get(); + if (baseTxnCommitInfo.left == null) { + // Connection with the server was not established. + return Unit.INSTANCE; + } + List transactions; try { - var baseTxnCommitInfo = myBaseTxnCommitInfo.get(); - var txns = myProjectDatabase.fetchTransactions(baseTxnCommitInfo.right + 1, 1); - if (!txns.isEmpty()) { - var listener = new TxnSendListener() { - @Override - public void onSendCompleted() { - txnSendingListener.compareAndSet(this, null); - } - }; - if (txnSendingListener.compareAndSet(null, listener)) { - getWebSocket().sendLogs(new InputXlog( - baseTxnCommitInfo.left, - "userId", - "refid", - txns - )); - } - } + transactions = myProjectDatabase.fetchTransactions(baseTxnCommitInfo.right + 1, 1); } catch (ProjectDatabaseException e) { gpLogger.error("Failed to send logs", new Object[]{}, ImmutableMap.of(), e); + return Unit.INSTANCE; + } + if (!transactions.isEmpty()) { + var listener = new TxnSendListener() { + @Override + public void onSendCompleted() { + txnSendingListener.compareAndSet(this, null); + } + }; + if (txnSendingListener.compareAndSet(null, listener)) { + getWebSocket().sendLogs(new InputXlog( + baseTxnCommitInfo.left, + "userId", + "refid", + transactions + )); + } else { + // Logs were sent by another thread, no action required. + } } return Unit.INSTANCE; } @@ -1001,12 +1013,15 @@ private Unit fireXlogReceived(ServerCommitResponse response) { return Unit.INSTANCE; } + // TODO: sync logs with the server. private Unit onBaseTxnIdReceived(String baseTxnId) { - myBaseTxnCommitInfo.update("", baseTxnId, 0); - var listener = txnSendingListener.get(); // Websocket is [re-]started. Previous messages are discarded. - if (listener != null) listener.onSendCompleted(); - sendProjectStateLogs(); + txnSendingListener.set(null); + if (EMPTY_LOG_BASE_TXN_ID.equals(baseTxnId)) { + myBaseTxnCommitInfo.reset(); + myBaseTxnCommitInfo.update(null, baseTxnId, 0); + sendProjectStateLogs(); + } return Unit.INSTANCE; } } diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java index a0da77ad76..2a54028d97 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/GanttProjectBase.java @@ -80,7 +80,6 @@ of the License, or (at your option) any later version. import java.util.Map; import java.util.function.Supplier; -import static biz.ganttproject.storage.cloud.GPCloudHttpImplKt.isColloboqueLocalTest; /** * This class is designed to be a GanttProject-after-refactorings. I am going to From fa03788bff6c2896b31d43607522f597e83c6f8a Mon Sep 17 00:00:00 2001 From: npostnikova Date: Thu, 16 Jun 2022 12:25:16 +0300 Subject: [PATCH 5/5] Add empty log txn id constant --- .../main/java/net/sourceforge/ganttproject/storage/Xlog.kt | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/Xlog.kt b/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/Xlog.kt index f9ba20911e..2b0b2b0446 100644 --- a/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/Xlog.kt +++ b/ganttproject/src/main/java/net/sourceforge/ganttproject/storage/Xlog.kt @@ -81,4 +81,7 @@ data class ServerCommitError( ) const val SERVER_COMMIT_RESPONSE_TYPE = "ServerCommitResponse" -const val SERVER_COMMIT_ERROR_TYPE = "ServerCommitError" \ No newline at end of file +const val SERVER_COMMIT_ERROR_TYPE = "ServerCommitError" + +/** Base txn ID for the empty log state. */ +const val EMPTY_LOG_BASE_TXN_ID = "abacaba" \ No newline at end of file