Initial commit

This commit is contained in:
2024-05-04 13:00:06 +02:00
commit 80a33ab08e
21 changed files with 1172 additions and 0 deletions

5
src/main/kotlin/Main.kt Normal file
View File

@@ -0,0 +1,5 @@
package nl.astraeus
fun main() {
println("Hello World!")
}

View File

@@ -0,0 +1,229 @@
package nl.astraeus.nl.astraeus.persistence
import java.io.File
import java.io.ObjectInputStream
import java.io.ObjectOutputStream
import java.io.Serializable
import java.text.DecimalFormat
import java.util.*
import java.util.concurrent.ConcurrentHashMap
import kotlin.reflect.KClass
enum class ActionType {
STORE,
DELETE
}
class TypeData(
var nextId: Long = 1L,
val data: MutableMap<Any, Persistable> = ConcurrentHashMap(),
) : Serializable
class Action(
val type: ActionType,
val obj: Persistable
) : Serializable
class Datastore(
private val directory: File,
indexes: Array<PersistableIndex> = arrayOf(),
) {
private val fileManager = FileManager(directory)
private val transactionFormatter = DecimalFormat("#")
private var nextTransactionNumber = 1L
private val data: MutableMap<Class<*>, TypeData> = ConcurrentHashMap()
private val indexes: MutableMap<Class<*>, MutableMap<String, PersistableIndex>> = ConcurrentHashMap()
init {
if (!directory.exists()) {
directory.mkdirs()
}
for (index in indexes) {
this.indexes.getOrPut(index.cls) {
ConcurrentHashMap()
}[index.name] = index
}
loadTransactions()
}
private fun loadTransactions() {
synchronized(this) {
val snapshots: Array<File>? = directory.listFiles { _, name -> name.startsWith("transaction-") && name.endsWith(".snp") }
val files: Array<File>? = directory.listFiles { _, name -> name.startsWith("transaction-") && name.endsWith(".trn") }
var lastSnapshot: Long? = null
var lastSnapshotFile: File? = null
snapshots?.let {
it.forEach {
val trnx = getTrnx(it)
if (lastSnapshot == null || trnx > (lastSnapshot ?: 0L)) {
lastSnapshot = trnx
lastSnapshotFile = it
}
}
}
val lastSnapshotFile2 = fileManager.findLastSnapshot()
if (lastSnapshotFile != null) {
ObjectInputStream(lastSnapshotFile?.inputStream()).use { ois ->
readSnapshot(ois)
}
}
val trns = fileManager.findTransactionsAfter(lastSnapshot ?: 0L)
files?.also { snaphotFiles ->
Arrays.sort(snaphotFiles) { o1, o2 -> if (getTrnx(o1) > getTrnx(o2)) 1 else -1}
snaphotFiles.forEach { file ->
if (getTrnx(file) > (lastSnapshot ?: 0L)) {
ObjectInputStream(file.inputStream()).use { ois ->
val transactionNumber = ois.readLong()
nextTransactionNumber = transactionNumber + 1
val actions = ois.readObject() as MutableList<Action>
execute(actions)
}
}
}
}
}
}
private fun getTrnx(file: File): Long {
return file.name.substringAfterLast('/').substringAfter("transaction-").substringBefore(".").toLong()
}
fun execute(actions: MutableList<Action>) {
synchronized(this) {
for (action in actions) {
val typeData = data.getOrPut(action.obj::class.java) {
TypeData()
}
when (action.type) {
ActionType.STORE -> {
if (action.obj.id == 0L) {
action.obj.id = typeData.nextId++
}
typeData.data[action.obj.id] = action.obj
for (index in indexes[action.obj::class.java]?.values ?: listOf()) {
index.add(action.obj as Persistable)
}
}
ActionType.DELETE -> {
typeData.data.remove(action.obj.id)
for (index in indexes[action.obj::class.java]?.values ?: listOf()) {
index.remove(action.obj)
}
}
}
}
}
}
fun <T : Persistable> find(clazz: KClass<T>, id: Long): T? {
val typeData = data.getOrPut(clazz.java) {
TypeData()
}
val p: Persistable? = typeData.data[id]
return (p?.copy() as T?)
}
fun <T : Persistable> search(
clazz: KClass<T>,
search: (T) -> Boolean
): List<T> {
val typeData = data.getOrPut(clazz.java) {
TypeData()
}
return typeData.data.values
.filter { search(it as T) }
.map { o -> o.copy() as T }
}
fun findIndex(
kClass: KClass<*>,
indexName: String
): PersistableIndex? {
return indexes[kClass.java]?.get(indexName)
}
fun storeActions(actions: MutableList<Action>) {
if (actions.isNotEmpty()) {
synchronized(this) {
val number = transactionFormatter.format(nextTransactionNumber)
val file = File(directory, "transaction-$number.trn")
ObjectOutputStream(file.outputStream()).use { oos ->
oos.writeLong(nextTransactionNumber++)
oos.writeObject(actions)
}
}
}
}
fun snapshot() {
synchronized(this) {
val number = transactionFormatter.format(nextTransactionNumber)
val file = File(directory, "transaction-$number.snp")
ObjectOutputStream(file.outputStream()).use { oos ->
oos.writeLong(nextTransactionNumber++)
oos.writeObject(data)
oos.writeInt(indexes.size)
for ((cls, index) in indexes) {
oos.writeObject(cls)
oos.writeInt(index.size)
for ((name, idx) in index) {
oos.writeObject(name)
oos.writeObject(idx.index)
}
}
}
}
}
private fun readSnapshot(ois: ObjectInputStream) {
nextTransactionNumber = ois.readLong() + 1
data.clear()
data.putAll(ois.readObject() as MutableMap<Class<*>, TypeData>)
val foundIndexes = mutableMapOf<Class<*>, MutableList<String>>()
val numberOfClassesWithIndex = ois.readInt()
for (i in 0 until numberOfClassesWithIndex) {
val cls = ois.readObject() as Class<*>
val numberOfIndexesForClass = ois.readInt()
for (indexForClass in 0 until numberOfIndexesForClass) {
val name = ois.readObject() as String
val idx = ois.readObject() as MutableMap<Serializable, MutableSet<Long>>
foundIndexes.getOrPut(cls) { mutableListOf() }.add(name)
val index = indexes[cls]
if (index != null) {
index[name]?.index?.putAll(idx)
} // else ignore, index is removed
}
}
// any (new) index not serialized needs to be build now
for ((cls, indexes) in indexes) {
for ((name, index) in indexes) {
if (!foundIndexes.getOrDefault(cls, mutableListOf()).contains(name)) {
index.index.clear()
for (obj in data.getOrDefault(cls, TypeData()).data.values) {
index.add(obj)
}
}
}
}
}
}

View File

@@ -0,0 +1,37 @@
package nl.astraeus.nl.astraeus.persistence
import java.io.Serializable
import kotlin.reflect.KClass
typealias PersistableIndex = Index<out Persistable>
class Index<T : Persistable>(
kcls: KClass<T>,
val name: String,
val value: (Persistable) -> Serializable?,
) : Serializable {
val cls: Class<T> = kcls.java
val index = mutableMapOf<Serializable, MutableSet<Long>>()
fun add(obj: Persistable) {
val key = value(obj)
key?.also {
index.getOrPut(it) { mutableSetOf() }.add(obj.id)
}
}
fun remove(obj: Persistable) {
val key = value(obj)
index[key]?.remove(obj.id)
}
fun find(key: Any): List<T> {
return index[key]?.mapNotNull { currentTransaction()?.find(cls.kotlin, it) } ?: emptyList()
}
fun matches(obj: Persistable, value: Any): Boolean {
return value(obj) == value
}
}

View File

@@ -0,0 +1,25 @@
package nl.astraeus.nl.astraeus.persistence
import java.io.ByteArrayInputStream
import java.io.ByteArrayOutputStream
import java.io.ObjectInputStream
import java.io.ObjectOutputStream
import java.io.Serializable
interface Persistable : Serializable, Cloneable {
var id: Long
var version: Long
fun copy(): Persistable {
ByteArrayOutputStream().use { baos ->
ObjectOutputStream(baos).use { oos ->
oos.writeObject(this)
}
ByteArrayInputStream(baos.toByteArray()).use { bais ->
ObjectInputStream(bais).use { ois ->
return ois.readObject() as Persistable
}
}
}
}
}

View File

@@ -0,0 +1,38 @@
package nl.astraeus.nl.astraeus.persistence
import java.io.File
private val transactions: ThreadLocal<Transaction> = ThreadLocal<Transaction>()
fun currentTransaction(): Transaction? {
return transactions.get()
}
class Persistent(
directory: File,
indexes: Array<PersistableIndex> = arrayOf(),
) {
val datastore: Datastore = Datastore(directory, indexes)
fun transaction(block: Transaction.() -> Unit) {
var cleanup = false
if (transactions.get() == null) {
transactions.set(Transaction(this))
cleanup = true
}
try {
block(transactions.get())
transactions.get().commit()
} finally {
if (cleanup) {
transactions.remove()
}
}
}
fun snapshot() {
datastore.snapshot()
}
}

View File

@@ -0,0 +1,133 @@
package nl.astraeus.nl.astraeus.persistence
import java.io.Serializable
import kotlin.reflect.KProperty
class Reference<S : Persistable, H : Persistable>(
val cls: Class<S>,
) : Serializable {
companion object {
private const val serialVersionUID: Long = 1L
}
var id: Long = 0
operator fun getValue(thisRef: H, property: KProperty<*>): S {
return currentTransaction()?.find(cls.kotlin, id) ?: throw IllegalStateException("Reference not found")
}
operator fun setValue(thisRef: H, property: KProperty<*>, value: S) {
id = value.id
currentTransaction()?.store(value)
}
}
class ListReference<S : Persistable, H : Persistable>(
val cls: Class<S>,
) : Serializable {
companion object {
private const val serialVersionUID: Long = 1L
}
var ids: ReferenceList<S> = ReferenceList(cls)
operator fun getValue(thisRef: H, property: KProperty<*>): ReferenceList<S> {
return ids
}
operator fun setValue(thisRef: H, property: KProperty<*>, value: List<S>) {
this.ids.clear()
this.ids.addAll(value)
}
}
class ReferenceList<T : Persistable>(
val cls: Class<T>,
) : MutableList<T> {
val ids = ArrayList<Long>()
private fun checkElementIsPersisted(element: T) {
if (currentTransaction()?.find(cls.kotlin, element.id) == null) {
currentTransaction()?.store(element)
}
}
override val size: Int = ids.size
override fun clear() = ids.clear()
override fun addAll(elements: Collection<T>): Boolean {
TODO("Not yet implemented")
}
override fun addAll(index: Int, elements: Collection<T>): Boolean {
TODO("Not yet implemented")
}
override fun add(index: Int, element: T) {
ids.add(index, element.id)
}
override fun add(element: T): Boolean {
return ids.add(element.id)
}
override fun get(index: Int): T = currentTransaction()?.find(cls.kotlin, ids[index]) ?: throw IllegalStateException("Reference not found")
override fun isEmpty(): Boolean = ids.isEmpty()
override fun iterator(): MutableIterator<T> {
TODO("Not yet implemented")
}
override fun listIterator(): MutableListIterator<T> {
TODO("Not yet implemented")
}
override fun listIterator(index: Int): MutableListIterator<T> {
TODO("Not yet implemented")
}
override fun removeAt(index: Int): T {
val id = ids.removeAt(index)
return currentTransaction()?.find(cls.kotlin, id) ?: throw IllegalStateException("Reference not found")
}
override fun set(index: Int, element: T): T {
TODO("Not yet implemented")
}
override fun retainAll(elements: Collection<T>): Boolean {
TODO("Not yet implemented")
}
override fun removeAll(elements: Collection<T>): Boolean {
TODO("Not yet implemented")
}
override fun remove(element: T): Boolean {
TODO("Not yet implemented")
}
override fun subList(fromIndex: Int, toIndex: Int): MutableList<T> {
TODO("Not yet implemented")
}
override fun lastIndexOf(element: T): Int {
TODO("Not yet implemented")
}
override fun indexOf(element: T): Int {
TODO("Not yet implemented")
}
override fun containsAll(elements: Collection<T>): Boolean {
TODO("Not yet implemented")
}
override fun contains(element: T): Boolean {
TODO("Not yet implemented")
}
}

View File

@@ -0,0 +1,100 @@
package nl.astraeus.nl.astraeus.persistence
import java.io.Serializable
import kotlin.reflect.KClass
class Transaction(
val persistent: Persistent,
) : Serializable {
private val actions = mutableListOf<Action>()
fun store(obj: Persistable) {
actions.add(Action(ActionType.STORE, obj))
}
fun delete(obj: Persistable) {
actions.add(Action(ActionType.DELETE, obj))
}
fun <T : Persistable> find(clazz: KClass<T>, id: Long): T? {
var result: T? = persistent.datastore.find(clazz, id)
for (action in actions) {
if (action.obj::class == clazz && action.obj.id == id) {
result = when {
action.type == ActionType.DELETE -> {
null
}
action.type == ActionType.STORE -> {
action.obj as? T
}
else -> {
result
}
}
}
}
return result
}
fun <T : Persistable> search(clazz: KClass<T>, search: (T) -> Boolean): List<T> {
val fromDatastore: List<T> = persistent.datastore.search(clazz, search)
val result = mutableListOf<T>()
result.addAll(fromDatastore)
for (obj in result) {
for (action in actions) {
if (action.obj.id == obj.id) {
if (action.type == ActionType.DELETE) {
result.remove(obj)
} else if (action.type == ActionType.STORE) {
result.remove(obj)
result.add(action.obj as T)
}
}
}
}
return result
}
fun commit() {
persistent.datastore.storeActions(actions)
persistent.datastore.execute(actions)
actions.clear()
}
fun <T : Persistable> findByIndex(
kClass: KClass<T>,
indexName: String,
search: Any
): List<T> {
val result = mutableListOf<T>()
val index = persistent.datastore.findIndex(kClass, indexName) ?: throw IllegalArgumentException("Index not found")
index.find(search).forEach { id ->
result.add(id as T)
}
for (action in actions) {
if (action.obj::class == kClass) {
if (action.type == ActionType.DELETE) {
if (index.matches(action.obj, search)) {
result.remove(action.obj as T)
}
} else if (action.type == ActionType.STORE) {
if (index.matches(action.obj, search)) {
result.remove(action.obj)
result.add(action.obj as T)
}
}
}
}
return result
}
}

View File

@@ -0,0 +1,152 @@
package nl.astraeus.persistence
import nl.astraeus.nl.astraeus.persistence.Index
import nl.astraeus.nl.astraeus.persistence.Persistable
import nl.astraeus.nl.astraeus.persistence.Persistent
import nl.astraeus.nl.astraeus.persistence.Reference
import java.io.File
import kotlin.test.Test
class TestPersistence {
class Company(
override var id: Long = 0,
override var version: Long = 0,
val name: String
) : Persistable, Cloneable {
//var persons: MutableList<Person> by ListReference<Person, Company>(Person::class.java)
companion object {
private const val serialVersionUID: Long = 1L
}
}
class Person(
override var id: Long = 0,
override var version: Long = 0,
val name: String,
val age: Int,
) : Persistable, Cloneable {
var company: Company by Reference<Company, Person>(Company::class.java)
companion object {
private const val serialVersionUID: Long = 1L
}
}
@Test
fun testPersistence() {
println("Test persistence")
val pst = Persistent(
directory = File("data"),
arrayOf(
Index(Person::class, "name") { p -> (p as? Person)?.name ?: "" },
Index(Person::class, "age") { p -> (p as? Person)?.age ?: -1 },
Index(Person::class, "ageGt20") { p -> ((p as? Person)?.age ?: 0) > 20 },
Index(Person::class, "ageGt23") { p -> ((p as? Person)?.age ?: 0) > 23 },
Index(Person::class, "ageOnlyGt20") { p ->
if (((p as? Person)?.age ?: 0) > 20) {
true
} else {
null
}
},
Index(Company::class, "name") { p -> (p as? Company)?.name ?: "" },
)
)
pst.transaction {
val person = find(Person::class, 1L) ?: Person(
id = 1L,
name = "John Doe",
age = 25
)
val company = find(Company::class, 1L) ?: Company(
id = 1L,
name = "ACME"
)
person.company = company
//company.persons.add(person)
store(person)
store(Person(
id = 2L,
name = "John Doe",
age = 23
))
store(Person(
id = 3L,
name = "John Doe",
age = 18
))
findByIndex(Person::class, "name", "John Doe").forEach { p ->
println("Found person by name: ${p.name} - ${p.age}")
}
findByIndex(Person::class, "age", 23).forEach { p ->
println("Found person by age: ${p.name} - ${p.age}")
}
findByIndex(Person::class, "ageGt20", true).forEach { p ->
println("Found person by age > 20: ${p.name} - ${p.age}")
}
findByIndex(Person::class, "ageGt23", true).forEach { p ->
println("Found person by age > 23: ${p.name} - ${p.age}")
}
findByIndex(Person::class, "ageGt20", false).forEach { p ->
println("Found person by age <= 20: ${p.name} - ${p.age}")
}
val p2 = find(Person::class, 1L)
assert(p2 != null)
val c2 = find(Company::class, 1L)
assert(c2 != null)
}
pst.transaction {
val person = find(Person::class, 1L)
delete(person!!)
val p2 = find(Person::class, 1L)
assert(p2 == null)
}
pst.transaction {
val persons = search(Person::class) { p -> p.name == "John Doe" }
if (persons.isNotEmpty()) {
delete(persons[0])
}
}
pst.snapshot()
pst.transaction {
store(
Person(
id = 10L,
name = "Pipo",
age = 23
)
)
store(
Person(
id = 11L,
name = "Clown",
age = 18
)
)
}
}
}