0.1.1-SNAPSHOT: storages, sync tools, docs
This commit is contained in:
1 parent
3e6c487601
commit
babc3933eb
24 files changed
+925
-52
No files matched your search
@@ -27,9 +27,6 @@ interface CRC<T> {
|
||||
fun crc8(data: ByteArray, polynomial: UByte = 0xA7.toUByte()): UByte =
|
||||
CRC8(polynomial).also { it.update(data) }.value
|
||||
|
||||
fun crc8(data: UByteArray, polynomial: UByte = 0xA7.toUByte()): UByte =
|
||||
CRC8(polynomial).also { it.update(data) }.value
|
||||
|
||||
/**
|
||||
* Calculate CRC16 for a data array using a given polynomial (CRC16-CCITT polynomial (0x1021) by default)
|
||||
*/
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
package net.sergeych.bintools
|
||||
|
||||
import net.sergeych.bipack.BipackDecoder
|
||||
import net.sergeych.bipack.BipackEncoder
|
||||
import net.sergeych.synctools.WaitHandle
|
||||
import net.sergeych.tools.ProtectedOp
|
||||
import net.sergeych.tools.withLock
|
||||
|
||||
class DataKVStorage(private val provider: DataProvider) : KVStorage {
|
||||
|
||||
data class Lock(val name: String) {
|
||||
private val exclusive = ProtectedOp()
|
||||
|
||||
private var readerCount = 0
|
||||
|
||||
private var pulses = WaitHandle()
|
||||
|
||||
fun <T> lockExclusive(f: () -> T): T {
|
||||
while (true) {
|
||||
exclusive.withLock {
|
||||
if (readerCount == 0) {
|
||||
return f()
|
||||
} else {
|
||||
println("can't lock $this: count is $readerCount")
|
||||
}
|
||||
}
|
||||
pulses.await()
|
||||
}
|
||||
}
|
||||
|
||||
fun <T> lockRead(f: () -> T): T {
|
||||
try {
|
||||
exclusive.withLock { readerCount++ }
|
||||
return f()
|
||||
} finally {
|
||||
exclusive.withLock { readerCount-- }
|
||||
if (readerCount == 0) pulses.wakeUp()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private val locks = mutableMapOf<String, Lock>()
|
||||
private val access = ProtectedOp()
|
||||
private val keyIds = mutableMapOf<String, Int>()
|
||||
private var lastId: Int = 0
|
||||
|
||||
init {
|
||||
access.withLock {
|
||||
// TODO: read keys
|
||||
for (fn in provider.list()) {
|
||||
println("Scanning: $fn")
|
||||
if (fn.endsWith(".d")) {
|
||||
val id = fn.dropLast(2).toInt(16)
|
||||
println("found data record: $fn -> $id")
|
||||
|
||||
val name = provider.read(fn) { BipackDecoder.decode<String>(it) }
|
||||
println("Key=$name")
|
||||
keyIds[name] = id
|
||||
if (id > lastId) lastId = id
|
||||
} else println("ignoring record $fn")
|
||||
}
|
||||
}
|
||||
println("initialized, ${keyIds.size} records found, lastId=$lastId")
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Important: __it must be called with locked access op!__
|
||||
*/
|
||||
private fun lockFor(name: String): Lock = locks.getOrPut(name) { Lock(name) }
|
||||
|
||||
private fun recordName(id: Int) = "${id.toString(16)}.d"
|
||||
private fun <T> read(name: String, f: (DataSource) -> T): T {
|
||||
val lock: Lock
|
||||
val id: Int
|
||||
// global lock: fast
|
||||
access.withLock {
|
||||
id = keyIds[name] ?: throw DataProvider.NotFoundException()
|
||||
lock = lockFor(name)
|
||||
}
|
||||
// per-name lock: slow
|
||||
return lock.lockRead {
|
||||
provider.read(recordName(id), f)
|
||||
}
|
||||
}
|
||||
|
||||
private fun write(name: String, f: (DataSink) -> Unit) {
|
||||
val lock: Lock
|
||||
val id: Int
|
||||
// global lock: fast
|
||||
access.withLock {
|
||||
id = keyIds[name] ?: (++lastId).also { keyIds[name] = it }
|
||||
lock = lockFor(name)
|
||||
}
|
||||
// per-name lock: slow
|
||||
lock.lockExclusive { provider.write(recordName(id), f) }
|
||||
}
|
||||
|
||||
fun deleteEntry(name: String) {
|
||||
// fast pre-check:
|
||||
if (name !in keyIds) return
|
||||
// global lock: we can't now detect concurrent delete + write ops, so exclusive:
|
||||
access.withLock {
|
||||
val id = keyIds[name] ?: return
|
||||
provider.delete(recordName(id))
|
||||
locks.remove(name)
|
||||
keyIds.remove(name)
|
||||
}
|
||||
}
|
||||
|
||||
override fun get(key: String): ByteArray? = try {
|
||||
read(key) {
|
||||
BipackDecoder.decode<String>(it)
|
||||
// notice: not nullable byte array here!
|
||||
BipackDecoder.decode<ByteArray>(it)
|
||||
}
|
||||
} catch (_: DataProvider.NotFoundException) {
|
||||
null
|
||||
}
|
||||
|
||||
override fun set(key: String, value: ByteArray?) {
|
||||
if (value == null) {
|
||||
deleteEntry(key)
|
||||
} else write(key) {
|
||||
BipackEncoder.encode(key, it)
|
||||
BipackEncoder.encode(value, it)
|
||||
}
|
||||
}
|
||||
|
||||
override val keys: Set<String>
|
||||
get() = access.withLock { keyIds.keys }
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package net.sergeych.bintools
|
||||
|
||||
/**
|
||||
* Abstraction of some file- or named_record- storage. It could be a filesystem
|
||||
* on native and JVM targets and indexed DB or session storage in the browser.
|
||||
*/
|
||||
interface DataProvider {
|
||||
open class Error(msg: String,cause: Throwable?=null): Exception(msg,cause)
|
||||
class NotFoundException(msg: String="record not found",cause: Throwable? = null): Error(msg, cause)
|
||||
class WriteFailedException(msg: String="can't write", cause: Throwable?=null): Error(msg,cause)
|
||||
|
||||
/**
|
||||
* Read named record/file
|
||||
* @throws NotFoundException
|
||||
*/
|
||||
fun <T>read(name: String,f: (DataSource)->T): T
|
||||
|
||||
/**
|
||||
* Write to a named record / file.
|
||||
* @throws WriteFailedException
|
||||
*/
|
||||
fun write(name: String,f: (DataSink)->Unit)
|
||||
|
||||
/**
|
||||
* Delete if exists, or do nothing.
|
||||
*/
|
||||
fun delete(name: String)
|
||||
|
||||
/**
|
||||
* List all record names in this source
|
||||
*/
|
||||
fun list(): List<String>
|
||||
}
|
||||
@@ -0,0 +1,216 @@
|
||||
package net.sergeych.bintools
|
||||
|
||||
import kotlinx.serialization.serializer
|
||||
import net.sergeych.bipack.BipackDecoder
|
||||
import net.sergeych.bipack.BipackEncoder
|
||||
import kotlin.reflect.KProperty
|
||||
import kotlin.reflect.KType
|
||||
import kotlin.reflect.typeOf
|
||||
|
||||
|
||||
/**
|
||||
* Generic storage of binary content. PArsec uses boss encoding to store everything in it
|
||||
* in a convenient way. See [KVStorage.stored], [KVStorage.invoke] and
|
||||
* [KVStorage.optStored] delegates. The [MemoryKVStorage] is an implementation that stores
|
||||
* values in memory, allowing to connect some other (e.g. persistent storage) later in a
|
||||
* completely transparent way. It can also be used to cache values on the fly.
|
||||
*
|
||||
* Also, it is possible to use [read] and [write] where delegated properties
|
||||
* do not fit well.
|
||||
*/
|
||||
@Suppress("unused")
|
||||
interface KVStorage {
|
||||
operator fun get(key: String): ByteArray?
|
||||
operator fun set(key: String, value: ByteArray?)
|
||||
|
||||
/**
|
||||
* Check whether key is in storage.
|
||||
* Default implementation uses [keys]. You may override it for performance
|
||||
*/
|
||||
operator fun contains(key: String): Boolean = key in keys
|
||||
|
||||
val keys: Set<String>
|
||||
|
||||
|
||||
/**
|
||||
* Get number of object in the storage
|
||||
* Default implementation uses [keys]. You may override it for performance
|
||||
*/
|
||||
val size: Int get() = keys.size
|
||||
|
||||
/**
|
||||
* Clears all objects in the storage
|
||||
* Default implementation uses [keys]. You may override it for performance
|
||||
*/
|
||||
fun clear() {
|
||||
for (k in keys) this[k] = null
|
||||
}
|
||||
|
||||
/**
|
||||
* Default implementation uses [keys]. You may override it for performance
|
||||
*/
|
||||
fun isEmpty() = size == 0
|
||||
|
||||
/**
|
||||
* Default implementation uses [keys]. You may override it for performance
|
||||
*/
|
||||
fun isNotEmpty() = size != 0
|
||||
|
||||
/**
|
||||
* Add all elements from another storage, overwriting any existing
|
||||
* keys.
|
||||
*/
|
||||
fun addAll(other: KVStorage) {
|
||||
for (k in other.keys) {
|
||||
this[k] = other[k]
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Delete element by key
|
||||
*/
|
||||
fun delete(key: String) {
|
||||
set(key, null)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Write
|
||||
*/
|
||||
inline fun <reified T: Any>KVStorage.write(key: String,value: T) {
|
||||
this[key] = BipackEncoder.encode(value)
|
||||
}
|
||||
|
||||
inline fun <reified T:Any>KVStorage.read(key: String): T? =
|
||||
this[key]?.let { BipackDecoder.decode(it) }
|
||||
|
||||
|
||||
inline operator fun <reified T> KVStorage.invoke(defaultValue: T,overrideName: String? = null) =
|
||||
KVStorageDelegate<T>(this, typeOf<T>(), defaultValue, overrideName)
|
||||
|
||||
inline fun <reified T> KVStorage.stored(defaultValue: T, overrideName: String? = null) =
|
||||
KVStorageDelegate<T>(this, typeOf<T>(), defaultValue, overrideName)
|
||||
inline fun <reified T> KVStorage.optStored(overrideName: String? = null) =
|
||||
KVStorageDelegate<T?>(this, typeOf<T?>(), null, overrideName)
|
||||
|
||||
class KVStorageDelegate<T>(
|
||||
private val storage: KVStorage,
|
||||
type: KType,
|
||||
private val defaultValue: T,
|
||||
private val overrideName: String? = null,
|
||||
) {
|
||||
|
||||
private fun name(property: KProperty<*>): String = overrideName ?: property.name
|
||||
|
||||
private var cachedValue: T = defaultValue
|
||||
private var cacheReady = false
|
||||
private val serializer = serializer(type)
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
operator fun getValue(thisRef: Any?, property: KProperty<*>): T {
|
||||
if (cacheReady) return cachedValue
|
||||
val data = storage.get(name(property))
|
||||
println("Got data: ${data?.toDump()}")
|
||||
if (data == null)
|
||||
cachedValue = defaultValue
|
||||
else
|
||||
cachedValue = BipackDecoder.decode(data.toDataSource(), serializer) as T
|
||||
cacheReady = true
|
||||
return cachedValue
|
||||
}
|
||||
|
||||
operator fun setValue(thisRef: Any?, property: KProperty<*>, value: T) {
|
||||
// if (!cacheReady || value != cachedValue) {
|
||||
cachedValue = value
|
||||
cacheReady = true
|
||||
println("set ${name(property)} to ${BipackEncoder.encode(serializer, value).toDump()}")
|
||||
storage[name(property)] = BipackEncoder.encode(serializer, value)
|
||||
// }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Memory storage. Allows connecting to another storage (e.g. persistent one) at come
|
||||
* point later in a transparent way.
|
||||
*/
|
||||
class MemoryKVStorage(copyFrom: KVStorage? = null) : KVStorage {
|
||||
|
||||
// is used when connected:
|
||||
private var underlying: KVStorage? = null
|
||||
|
||||
// is used while underlying is null:
|
||||
private val data = mutableMapOf<String, ByteArray>()
|
||||
|
||||
/**
|
||||
* Connect some other storage. All existing data will be copied to the [other]
|
||||
* storage. After this call all data access will be routed to [other] storage.
|
||||
*/
|
||||
@Suppress("unused")
|
||||
fun connectToStorage(other: KVStorage) {
|
||||
other.addAll(this)
|
||||
underlying = other
|
||||
data.clear()
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Get data from either memory or a connected storage, see [connectToStorage]
|
||||
*/
|
||||
override fun get(key: String): ByteArray? {
|
||||
underlying?.let {
|
||||
return it[key]
|
||||
}
|
||||
return data[key]
|
||||
}
|
||||
|
||||
/**
|
||||
* Put data to memory storage or connected storage if [connectToStorage] was called
|
||||
*/
|
||||
override fun set(key: String, value: ByteArray?) {
|
||||
underlying?.let { it[key] = value } ?: run {
|
||||
if (value != null) data[key] = value
|
||||
else data.remove(key)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks the item exists in the memory storage or connected one, see[connectToStorage]
|
||||
*/
|
||||
override fun contains(key: String): Boolean {
|
||||
underlying?.let { return key in it }
|
||||
return key in data
|
||||
}
|
||||
|
||||
override val keys: Set<String>
|
||||
get() = underlying?.keys ?: data.keys
|
||||
|
||||
override fun clear() {
|
||||
underlying?.clear() ?: data.clear()
|
||||
}
|
||||
|
||||
init {
|
||||
copyFrom?.let { addAll(it) }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Create per-platform default named storage.
|
||||
*
|
||||
* - In the browser, it uses the `Window.localStorage` prefixing items
|
||||
* by a string containing the [name]
|
||||
*
|
||||
* - In the JVM environment it uses folder-based storage on the file system. The name
|
||||
* is considered to be a folder name (the whole path which will be automatically created)
|
||||
* using the following rules:
|
||||
* - when the name starts with slash (`/`) it is treated as an absolute path to a folder
|
||||
* - when the name contains slash, it is considered to be a relative folder to the
|
||||
* `User.home` directory, like "`~/`" on unix systems.
|
||||
* - otherwise, the folder will be created in "`~/.local_storage`" parent directory
|
||||
* (which also will be created if needed).
|
||||
*
|
||||
* - For the native platorms it is not yet implemented (but will be soon).
|
||||
*
|
||||
* See [DataKVStorage] and [DataProvider] to implement a KVStorage on filesystems and like,
|
||||
* and `FileDataProvider` class on JVM target.
|
||||
*/
|
||||
expect fun defaultNamedStorage(name: String): KVStorage
|
||||
File renamed without changes.
@@ -0,0 +1,15 @@
|
||||
package net.sergeych.tools
|
||||
|
||||
/**
|
||||
* Thread-safe multiplatform counter
|
||||
*/
|
||||
@Suppress("unused")
|
||||
class AtomicCounter(initialValue: Long = 0) : AtomicValue<Long>(initialValue) {
|
||||
|
||||
fun incrementAndGet(): Long = op { ++actualValue }
|
||||
fun getAndIncrement(): Long = op { actualValue++ }
|
||||
|
||||
fun decrementAndGet(): Long = op { --actualValue }
|
||||
|
||||
fun getAndDecrement(): Long = op { actualValue-- }
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package net.sergeych.tools
|
||||
|
||||
/**
|
||||
* Multiplatform (JS and battery included) atomically mutable value.
|
||||
* Actual value can be either changed in a block of [mutate] when
|
||||
* new value _depends on the current value_ or use a same [value]
|
||||
* property that is thread-safe where there are threads and just safe
|
||||
* otherwise ;)
|
||||
*/
|
||||
open class AtomicValue<T>(initialValue: T) {
|
||||
var actualValue = initialValue
|
||||
protected set
|
||||
|
||||
protected val op = ProtectedOp()
|
||||
|
||||
/**
|
||||
* Change the value: get the current and set to the returned, all in the
|
||||
* atomic operation. All other mutating requests including assigning to [value]
|
||||
* will be blocked and queued.
|
||||
* @return result of the mutation. Note that immediate call to property [value]
|
||||
* could already return modified bu some other thread value!
|
||||
*/
|
||||
fun mutate(mutator: (T) -> T): T = op {
|
||||
actualValue = mutator(actualValue)
|
||||
actualValue
|
||||
}
|
||||
|
||||
/**
|
||||
* Atomic get or set the value. Atomic get means if there is a [mutate] in progress
|
||||
* it will wait until the mutation finishes and then return the correct result.
|
||||
*/
|
||||
var value: T
|
||||
get() = op { actualValue }
|
||||
set(value) {
|
||||
mutate { value }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
package net.sergeych.tools
|
||||
|
||||
import kotlin.contracts.ExperimentalContracts
|
||||
import kotlin.contracts.InvocationKind
|
||||
import kotlin.contracts.contract
|
||||
|
||||
/**
|
||||
* Multiplatform interface to perform a regular (not suspend) operation
|
||||
* protected by a platform mutex (where necessary). Get real implementation
|
||||
* with [ProtectedOp] and use it with [ProtectedOpImplementation.withLock] and
|
||||
* [ProtectedOpImplementation.invoke]
|
||||
*/
|
||||
interface ProtectedOpImplementation {
|
||||
/**
|
||||
* Get a lock. Be sure to release it.
|
||||
* The recommended way is using [ProtectedOpImplementation.withLock] and
|
||||
* [ProtectedOpImplementation.invoke]
|
||||
*/
|
||||
fun lock()
|
||||
|
||||
/**
|
||||
* Release a lock.
|
||||
* The recommended way is using [ProtectedOpImplementation.withLock] and
|
||||
* [ProtectedOpImplementation.invoke]
|
||||
*/
|
||||
fun unlock()
|
||||
}
|
||||
|
||||
@ExperimentalContracts
|
||||
inline fun <T> ProtectedOpImplementation.withLock(f: () -> T): T {
|
||||
contract {
|
||||
callsInPlace(f, InvocationKind.EXACTLY_ONCE)
|
||||
}
|
||||
lock()
|
||||
return try {
|
||||
f()
|
||||
} finally {
|
||||
unlock()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Run a block mutualy-exclusively, see [ProtectedOp]
|
||||
*/
|
||||
operator fun <T> ProtectedOpImplementation.invoke(f: () -> T): T = withLock(f)
|
||||
|
||||
/**
|
||||
* Get the platform-depended implementation of a mutex. It does nothing in the
|
||||
* browser and use appropriate mechanics on JVM and native targets. See
|
||||
* [ProtectedOpImplementation.invoke], [ProtectedOpImplementation.withLock]
|
||||
* ```kotlin
|
||||
* val op = ProtectedOp()
|
||||
* //...
|
||||
* op {
|
||||
* // mutually exclusive execution
|
||||
* println("sequential execution here")
|
||||
* }
|
||||
* ~~~
|
||||
*/
|
||||
expect fun ProtectedOp(): ProtectedOpImplementation
|
||||
@@ -0,0 +1,23 @@
|
||||
package net.sergeych.synctools
|
||||
|
||||
/**
|
||||
* Platform-independent interface to thread wait/notify. Does nothing in JS/browser,
|
||||
* and uses appropriate mechanics on other platforms.
|
||||
*/
|
||||
@Suppress("EXPECT_ACTUAL_CLASSIFIERS_ARE_IN_BETA_WARNING")
|
||||
expect class WaitHandle() {
|
||||
|
||||
/**
|
||||
* Wait for [wakeUp] as long as [milliseconds] milliseconds, or forever.
|
||||
* Notice it returns immediately on the single-theaded platforms like JS.
|
||||
* @param milliseconds to wait, use 0 to wait indefinitely
|
||||
* @return true if [wakeUp] was called before the timeout, false on timeout.
|
||||
*/
|
||||
fun await(milliseconds: Long=0): Boolean
|
||||
|
||||
/**
|
||||
* Awake all [await]'ing threads. Does nothing in the single-threaded JS
|
||||
* environment
|
||||
*/
|
||||
fun wakeUp()
|
||||
}
|
||||
Reference in new issue
Block a user