Optimistic locking option
This commit is contained in:
@@ -17,16 +17,29 @@ enum class ActionType {
|
|||||||
class TypeData(
|
class TypeData(
|
||||||
var nextId: AtomicLong = AtomicLong(1L),
|
var nextId: AtomicLong = AtomicLong(1L),
|
||||||
val data: MutableMap<Serializable, Persistable> = ConcurrentHashMap(),
|
val data: MutableMap<Serializable, Persistable> = ConcurrentHashMap(),
|
||||||
) : Serializable
|
) : Serializable {
|
||||||
|
companion object {
|
||||||
|
private const val serialVersionUID: Long = 1L
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
class Action(
|
class Action(
|
||||||
val type: ActionType,
|
val type: ActionType,
|
||||||
val obj: Persistable
|
val obj: Persistable
|
||||||
) : Serializable
|
) : Serializable {
|
||||||
|
override fun toString(): String {
|
||||||
|
return "Action(type=$type, obj=$obj)"
|
||||||
|
}
|
||||||
|
|
||||||
|
companion object {
|
||||||
|
private const val serialVersionUID: Long = 1L
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
class Datastore(
|
class Datastore(
|
||||||
private val directory: File,
|
private val directory: File,
|
||||||
indexes: Array<PersistableIndex> = arrayOf(),
|
indexes: Array<PersistableIndex> = arrayOf(),
|
||||||
|
val enableOptimisticLocking: Boolean = false,
|
||||||
) {
|
) {
|
||||||
private val fileManager = FileManager(directory)
|
private val fileManager = FileManager(directory)
|
||||||
private val transactionFormatter = DecimalFormat("#")
|
private val transactionFormatter = DecimalFormat("#")
|
||||||
@@ -58,7 +71,8 @@ class Datastore(
|
|||||||
|
|
||||||
fun setMaxId(javaClass: Class<Persistable>, id: Long) {
|
fun setMaxId(javaClass: Class<Persistable>, id: Long) {
|
||||||
val nextId = data.getOrPut(javaClass) { TypeData() }.nextId
|
val nextId = data.getOrPut(javaClass) { TypeData() }.nextId
|
||||||
if (nextId.get() <= id) nextId.set(id + 1)
|
val current = nextId.get()
|
||||||
|
if (current <= id) nextId.addAndGet(id - current)
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun toString(): String {
|
override fun toString(): String {
|
||||||
@@ -97,15 +111,24 @@ class Datastore(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
Logger.debug("Loaded transactions in ${(System.nanoTime() - start) / 1_000_000}ms")
|
Logger.debug("Loaded transactions in %6.3fms", ((System.nanoTime() - start) / 1_000_000f))
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun getTrnx(file: File): Long {
|
private fun execute(actions: Set<Action>) {
|
||||||
return file.name.substringAfterLast('/').substringAfter("transaction-").substringBefore(".").toLong()
|
|
||||||
}
|
|
||||||
|
|
||||||
fun execute(actions: MutableList<Action>) {
|
|
||||||
synchronized(this) {
|
synchronized(this) {
|
||||||
|
if (enableOptimisticLocking) {
|
||||||
|
for (action in actions) {
|
||||||
|
val typeData = data.getOrPut(action.obj::class.java) {
|
||||||
|
TypeData()
|
||||||
|
}
|
||||||
|
if (action.type == ActionType.STORE) {
|
||||||
|
if ((typeData.data[action.obj.id]?.version ?: -1L) >= action.obj.version) {
|
||||||
|
throw OptimisticLockingException(action.obj)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
for (action in actions) {
|
for (action in actions) {
|
||||||
val typeData = data.getOrPut(action.obj::class.java) {
|
val typeData = data.getOrPut(action.obj::class.java) {
|
||||||
TypeData()
|
TypeData()
|
||||||
@@ -137,7 +160,6 @@ class Datastore(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
fun <T : Persistable> count(clazz: KClass<T>): Int {
|
fun <T : Persistable> count(clazz: KClass<T>): Int {
|
||||||
val typeData = data.getOrPut(clazz.java) {
|
val typeData = data.getOrPut(clazz.java) {
|
||||||
TypeData()
|
TypeData()
|
||||||
@@ -176,25 +198,25 @@ class Datastore(
|
|||||||
return indexes[kClass.java]?.get(indexName)
|
return indexes[kClass.java]?.get(indexName)
|
||||||
}
|
}
|
||||||
|
|
||||||
fun storeAndExecute(actions: MutableList<Action>) {
|
fun storeAndExecute(actions: Set<Action>) {
|
||||||
if (actions.isNotEmpty()) {
|
if (actions.isNotEmpty()) {
|
||||||
synchronized(this) {
|
synchronized(this) {
|
||||||
writeTransaction(actions)
|
|
||||||
execute(actions)
|
execute(actions)
|
||||||
|
writeTransaction(actions)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun readTransaction(ois: ObjectInputStream) {
|
private fun readTransaction(ois: ObjectInputStream) {
|
||||||
val versionNumber = ois.readInt()
|
val versionNumber = ois.readInt()
|
||||||
check (versionNumber == 1) { "Unsupported version number: $versionNumber" }
|
check(versionNumber == 1) { "Unsupported version number: $versionNumber" }
|
||||||
val transactionNumber = ois.readLong()
|
val transactionNumber = ois.readLong()
|
||||||
nextTransactionNumber = transactionNumber + 1
|
nextTransactionNumber = transactionNumber + 1
|
||||||
val actions = ois.readObject() as MutableList<Action>
|
val actions = ois.readObject() as Set<Action>
|
||||||
execute(actions)
|
execute(actions)
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun writeTransaction(actions: MutableList<Action>) {
|
private fun writeTransaction(actions: Set<Action>) {
|
||||||
val number = transactionFormatter.format(nextTransactionNumber)
|
val number = transactionFormatter.format(nextTransactionNumber)
|
||||||
val file = File(directory, "transaction-$number.trn")
|
val file = File(directory, "transaction-$number.trn")
|
||||||
ObjectOutputStream(file.outputStream()).use { oos ->
|
ObjectOutputStream(file.outputStream()).use { oos ->
|
||||||
@@ -226,12 +248,12 @@ class Datastore(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Logger.debug("Snapshot in ${(System.nanoTime() - start) / 1_000_000}ms")
|
Logger.debug("Snapshot in %6.3fms", ((System.nanoTime() - start) / 1_000_000f))
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun readSnapshot(ois: ObjectInputStream) {
|
private fun readSnapshot(ois: ObjectInputStream) {
|
||||||
val versionNumber = ois.readInt()
|
val versionNumber = ois.readInt()
|
||||||
check (versionNumber == 1) { "Unsupported version number: $versionNumber" }
|
check(versionNumber == 1) { "Unsupported version number: $versionNumber" }
|
||||||
nextTransactionNumber = ois.readLong() + 1
|
nextTransactionNumber = ois.readLong() + 1
|
||||||
data.clear()
|
data.clear()
|
||||||
data.putAll(ois.readObject() as MutableMap<Class<*>, TypeData>)
|
data.putAll(ois.readObject() as MutableMap<Class<*>, TypeData>)
|
||||||
|
|||||||
@@ -1,8 +1,62 @@
|
|||||||
package nl.astraeus.nl.astraeus.persistence
|
package nl.astraeus.nl.astraeus.persistence
|
||||||
|
|
||||||
|
enum class LogLevel {
|
||||||
|
TRACE,
|
||||||
|
DEBUG,
|
||||||
|
INFO,
|
||||||
|
WARN,
|
||||||
|
ERROR
|
||||||
|
}
|
||||||
|
|
||||||
object Logger {
|
object Logger {
|
||||||
var debug: (String) -> Unit = { println("DEBUG: $it") }
|
var level: LogLevel = LogLevel.DEBUG
|
||||||
var info: (String) -> Unit = { println("INFO: $it") }
|
|
||||||
var warn: (String) -> Unit = { println("WARN: $it") }
|
var tracePrinter: (String) -> Unit = { println(it) }
|
||||||
var error: (String) -> Unit = { println("ERROR: $it") }
|
var debugPrinter: (String) -> Unit = { println(it) }
|
||||||
}
|
var infoPrinter: (String) -> Unit = { println(it) }
|
||||||
|
var warnPrinter: (String) -> Unit = { println(it) }
|
||||||
|
var errorPrinter: (String) -> Unit = { System.err.println(it) }
|
||||||
|
|
||||||
|
fun trace(message: String, vararg parameters: Any?) {
|
||||||
|
if (level <= LogLevel.TRACE) {
|
||||||
|
writeLogMessage(LogLevel.TRACE, message, *parameters)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun debug(message: String, vararg parameters: Any?) {
|
||||||
|
if (level <= LogLevel.DEBUG) {
|
||||||
|
writeLogMessage(LogLevel.DEBUG, message, *parameters)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun info(message: String, vararg parameters: Any?) {
|
||||||
|
if (level <= LogLevel.INFO) {
|
||||||
|
writeLogMessage(LogLevel.INFO, message, *parameters)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun warn(message: String, vararg parameters: Any?) {
|
||||||
|
if (level <= LogLevel.DEBUG) {
|
||||||
|
writeLogMessage(LogLevel.DEBUG, message, *parameters)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun error(message: String, vararg parameters: Any?) {
|
||||||
|
if (level <= LogLevel.ERROR) {
|
||||||
|
writeLogMessage(LogLevel.ERROR, message, *parameters)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun writeLogMessage(level: LogLevel, message: String, vararg parameters: Any?) {
|
||||||
|
val formattedMessage = "[${level}] - ${message.format(*parameters)}"
|
||||||
|
|
||||||
|
when (level) {
|
||||||
|
LogLevel.TRACE -> tracePrinter(formattedMessage)
|
||||||
|
LogLevel.DEBUG -> debugPrinter(formattedMessage)
|
||||||
|
LogLevel.INFO -> infoPrinter(formattedMessage)
|
||||||
|
LogLevel.WARN -> warnPrinter(formattedMessage)
|
||||||
|
LogLevel.ERROR -> errorPrinter(formattedMessage)
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,18 @@
|
|||||||
|
package nl.astraeus.nl.astraeus.persistence
|
||||||
|
|
||||||
|
class OptimisticLockingException : Exception {
|
||||||
|
constructor(
|
||||||
|
obj: Persistable
|
||||||
|
) : this("Optimistic locking failed for ${obj.javaClass.simpleName} with id ${obj.id}, version ${obj.version}")
|
||||||
|
|
||||||
|
constructor() : super()
|
||||||
|
constructor(message: String?) : super(message)
|
||||||
|
constructor(message: String?, cause: Throwable?) : super(message, cause)
|
||||||
|
constructor(cause: Throwable?) : super(cause)
|
||||||
|
constructor(message: String?, cause: Throwable?, enableSuppression: Boolean, writableStackTrace: Boolean) : super(
|
||||||
|
message,
|
||||||
|
cause,
|
||||||
|
enableSuppression,
|
||||||
|
writableStackTrace
|
||||||
|
)
|
||||||
|
}
|
||||||
@@ -17,7 +17,9 @@ interface Persistable : Serializable {
|
|||||||
}
|
}
|
||||||
ByteArrayInputStream(baos.toByteArray()).use { bais ->
|
ByteArrayInputStream(baos.toByteArray()).use { bais ->
|
||||||
ObjectInputStream(bais).use { ois ->
|
ObjectInputStream(bais).use { ois ->
|
||||||
return ois.readObject() as Persistable
|
val result = ois.readObject() as Persistable
|
||||||
|
result.version++
|
||||||
|
return result
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,8 +11,9 @@ fun currentTransaction(): Transaction? {
|
|||||||
class Persistent(
|
class Persistent(
|
||||||
directory: File,
|
directory: File,
|
||||||
indexes: Array<PersistableIndex> = arrayOf(),
|
indexes: Array<PersistableIndex> = arrayOf(),
|
||||||
|
enableOptimisticLocking: Boolean = false,
|
||||||
) {
|
) {
|
||||||
val datastore: Datastore = Datastore(directory, indexes)
|
val datastore: Datastore = Datastore(directory, indexes, enableOptimisticLocking)
|
||||||
|
|
||||||
fun <T> query(block: Query.() -> T): T {
|
fun <T> query(block: Query.() -> T): T {
|
||||||
return block(Query(this))
|
return block(Query(this))
|
||||||
|
|||||||
@@ -95,11 +95,13 @@ open class Query(
|
|||||||
class Transaction(
|
class Transaction(
|
||||||
persistent: Persistent,
|
persistent: Persistent,
|
||||||
) : Query(persistent), Serializable {
|
) : Query(persistent), Serializable {
|
||||||
private val actions = mutableListOf<Action>()
|
private val actions = mutableSetOf<Action>()
|
||||||
|
|
||||||
fun store(obj: Persistable) {
|
fun store(obj: Persistable) {
|
||||||
if (obj.id == 0L) {
|
if (obj.id == 0L) {
|
||||||
obj.id = persistent.datastore.getNextId(obj.javaClass)
|
obj.id = persistent.datastore.getNextId(obj.javaClass)
|
||||||
|
} else if (obj.id > persistent.datastore.getNextId(obj.javaClass)) {
|
||||||
|
persistent.datastore.setMaxId(obj.javaClass, obj.id + 1)
|
||||||
}
|
}
|
||||||
|
|
||||||
actions.add(Action(ActionType.STORE, obj))
|
actions.add(Action(ActionType.STORE, obj))
|
||||||
|
|||||||
33
src/main/kotlin/nl/astraeus/persistence/TransactionLog.kt
Normal file
33
src/main/kotlin/nl/astraeus/persistence/TransactionLog.kt
Normal file
@@ -0,0 +1,33 @@
|
|||||||
|
package nl.astraeus.nl.astraeus.persistence
|
||||||
|
|
||||||
|
import java.io.File
|
||||||
|
import java.io.ObjectInputStream
|
||||||
|
|
||||||
|
class TransactionLog(
|
||||||
|
val directory: File,
|
||||||
|
) {
|
||||||
|
val fileManager = FileManager(directory)
|
||||||
|
|
||||||
|
fun showTransactions() {
|
||||||
|
fileManager.findLastSnapshot().let { (after, snapshot) ->
|
||||||
|
println("Last snapshot: $snapshot")
|
||||||
|
|
||||||
|
val transactions = fileManager.findTransactionsAfter(after ?: 0L)
|
||||||
|
|
||||||
|
println("Transactions:")
|
||||||
|
transactions?.forEach { transaction ->
|
||||||
|
transaction.inputStream().use { input ->
|
||||||
|
ObjectInputStream(input).use { ois ->
|
||||||
|
val versionNumber = ois.readInt()
|
||||||
|
check(versionNumber == 1) { "Unsupported version number: $versionNumber" }
|
||||||
|
val transactionNumber = ois.readLong()
|
||||||
|
val actions = ois.readObject() as Set<Action>
|
||||||
|
println("[$versionNumber] $transactionNumber - ${actions.joinToString(",")}")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
@@ -74,7 +74,8 @@ public class TestPersistenceJava {
|
|||||||
"name",
|
"name",
|
||||||
(p) -> ((Person)p).getName()
|
(p) -> ((Person)p).getName()
|
||||||
)
|
)
|
||||||
}
|
},
|
||||||
|
false
|
||||||
);
|
);
|
||||||
|
|
||||||
persistent.transaction((t) -> {
|
persistent.transaction((t) -> {
|
||||||
|
|||||||
114
src/test/kotlin/nl/astraeus/persistence/TestOptimisticLocking.kt
Normal file
114
src/test/kotlin/nl/astraeus/persistence/TestOptimisticLocking.kt
Normal file
@@ -0,0 +1,114 @@
|
|||||||
|
package nl.astraeus.persistence
|
||||||
|
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.OptimisticLockingException
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.Persistable
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.Persistent
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.TransactionLog
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.find
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.findByIndex
|
||||||
|
import nl.astraeus.nl.astraeus.persistence.index
|
||||||
|
import org.junit.jupiter.api.Assertions.assertNotNull
|
||||||
|
import org.junit.jupiter.api.assertThrows
|
||||||
|
import java.io.File
|
||||||
|
import kotlin.test.Test
|
||||||
|
|
||||||
|
class TestOptimisticLocking {
|
||||||
|
|
||||||
|
class Person(
|
||||||
|
override var id: Long = 0,
|
||||||
|
override var version: Long = 0,
|
||||||
|
val name: String,
|
||||||
|
val age: Int,
|
||||||
|
) : Persistable, Cloneable {
|
||||||
|
companion object {
|
||||||
|
private const val serialVersionUID: Long = 1L
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun toString(): String {
|
||||||
|
return "Person(id=$id, version=$version, name='$name', age=$age)"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun showTransactions() {
|
||||||
|
val log = TransactionLog(File("data", "test-locking"))
|
||||||
|
|
||||||
|
log.showTransactions()
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun testOptimisticLocking() {
|
||||||
|
println("Test locking")
|
||||||
|
|
||||||
|
val pst = Persistent(
|
||||||
|
directory = File("data", "test-locking"),
|
||||||
|
arrayOf(
|
||||||
|
index<Person>("name") { p -> (p as? Person)?.name ?: "" },
|
||||||
|
),
|
||||||
|
true
|
||||||
|
)
|
||||||
|
|
||||||
|
pst.transaction {
|
||||||
|
val person = find<Person>(1L) ?: Person(
|
||||||
|
id = 1L,
|
||||||
|
name = "John Doe",
|
||||||
|
age = 25
|
||||||
|
)
|
||||||
|
|
||||||
|
store(person)
|
||||||
|
|
||||||
|
findByIndex<Person>("name", "John Doe").forEach { p ->
|
||||||
|
println("Found person by name: ${p.name} - ${p.age}")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pst.query {
|
||||||
|
val person = find<Person>(1L)
|
||||||
|
|
||||||
|
assertNotNull(person)
|
||||||
|
}
|
||||||
|
|
||||||
|
val threads = Array(2) { index ->
|
||||||
|
Thread {
|
||||||
|
println("Start thread $index")
|
||||||
|
var person: Person? = null
|
||||||
|
|
||||||
|
pst.transaction {
|
||||||
|
Thread.sleep(10L)
|
||||||
|
person = find<Person>(1L)
|
||||||
|
println("Thread $index -> ${person?.version}")
|
||||||
|
}
|
||||||
|
|
||||||
|
if (person != null) {
|
||||||
|
Thread.sleep((index + 1) * 10L)
|
||||||
|
|
||||||
|
if (index == 1) {
|
||||||
|
assertThrows<OptimisticLockingException> {
|
||||||
|
println("Store thread $index -> ${person!!.version}")
|
||||||
|
pst.transaction {
|
||||||
|
store(person!!)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
println("Store thread $index -> ${person!!.version}")
|
||||||
|
pst.transaction {
|
||||||
|
store(person!!)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for (thread in threads) {
|
||||||
|
thread.start()
|
||||||
|
}
|
||||||
|
|
||||||
|
for (thread in threads) {
|
||||||
|
thread.join()
|
||||||
|
}
|
||||||
|
|
||||||
|
pst.datastore.printStatus()
|
||||||
|
//pst.snapshot()
|
||||||
|
pst.removeOldFiles()
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -127,6 +127,14 @@ class TestPersistence {
|
|||||||
assertNotNull(person)
|
assertNotNull(person)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pst.transaction {
|
||||||
|
val p1 = find<Person>(1L)
|
||||||
|
val p2 = find<Person>(1L)
|
||||||
|
|
||||||
|
store(p2!!)
|
||||||
|
store(p1!!)
|
||||||
|
}
|
||||||
|
|
||||||
pst.transaction {
|
pst.transaction {
|
||||||
val person = find(Person::class, 1L)
|
val person = find(Person::class, 1L)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user