diff options
| author | ketsuban <ketsuban@192.168.3.110> | 2026-08-26 06:54:51 +0000 |
|---|---|---|
| committer | ketsuban <ketsuban@192.168.3.110> | 2026-08-26 06:54:51 +0000 |
| commit | 4468ed06508414e505dec106ff6c0a32a126d1e2 (patch) | |
| tree | c489806bbb9104e80e530b0c728aeb962bb18171 | |
| parent | 6d14259175d07eeee7e5e1bc9eecf838ff442fe4 (diff) | |
| download | kukuri-4468ed06508414e505dec106ff6c0a32a126d1e2.tar.gz kukuri-4468ed06508414e505dec106ff6c0a32a126d1e2.tar.bz2 kukuri-4468ed06508414e505dec106ff6c0a32a126d1e2.zip | |
Vendor nabu adapter (jvm) + variable CDC splitter tests
8 files changed, 118 insertions, 24 deletions
diff --git a/modules/orgflow-content-store/src/commonTest/kotlin/jp/orgflow/contentstore/ContentSplitterTest.kt b/modules/orgflow-content-store/src/commonTest/kotlin/jp/orgflow/contentstore/ContentSplitterTest.kt new file mode 100644 index 0000000..e227917 --- /dev/null +++ b/modules/orgflow-content-store/src/commonTest/kotlin/jp/orgflow/contentstore/ContentSplitterTest.kt @@ -0,0 +1,50 @@ +package jp.orgflow.contentstore + +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class ContentSplitterTest { + private val splitter = ContentSplitter(min = 64, max = 512) + + private fun pseudoRandom(size: Int, seed: Long = 42): ByteArray { + var s = seed + return ByteArray(size) { s = s * 6364136223846793005L + 1442695040888963407L; (s shr 33).toByte() } + } + + @Test + fun emptyInputYieldsNoChunks() { + assertTrue(splitter.split(ByteArray(0)).isEmpty()) + } + + @Test + fun smallInputSingleChunk() { + val chunks = splitter.split("hello".encodeToByteArray()) + assertEquals(1, chunks.size) + } + + @Test + fun deterministicBoundaries() { + val data = pseudoRandom(10_000) + assertEquals(splitter.split(data), splitter.split(data)) + } + + @Test + fun chunksWithinSizeBounds() { + val chunks = splitter.split(pseudoRandom(50_000)) + assertTrue(chunks.size > 1) + chunks.forEach { assertTrue(it.size <= 512, "chunk size ${it.size}") } + assertEquals(50_000, chunks.sumOf { it.size }) + } + + @Test + fun reassemblyRoundTripViaCids() { + val data = pseudoRandom(20_000) + val chunks = splitter.split(data) + val cids = chunks.map { ContentAddresser.cidOf(it) } + val restored = cids.flatMapIndexed { i, cid -> + chunks[i].let { ch -> check(cid == ContentAddresser.cidOf(ch)); ch.toList() } + }.toByteArray() + assertTrue(data.contentEquals(restored)) + } +} diff --git a/modules/orgflow-content-store/src/jvmMain/kotlin/jp/orgflow/contentstore/jvm/NabuBlockStoreAdapter.kt b/modules/orgflow-content-store/src/jvmMain/kotlin/jp/orgflow/contentstore/jvm/NabuBlockStoreAdapter.kt new file mode 100644 index 0000000..2b9bdb5 --- /dev/null +++ b/modules/orgflow-content-store/src/jvmMain/kotlin/jp/orgflow/contentstore/jvm/NabuBlockStoreAdapter.kt @@ -0,0 +1,35 @@ +package jp.orgflow.contentstore.jvm + +import io.ipfs.multihash.Multihash +import jp.orgflow.contentstore.BlockStore +import jp.orgflow.contentstore.ContentBlock +import org.peergos.blockstore.RamBlockstore +import java.util.concurrent.TimeUnit + +// ch.20a: adapter exposing the vendored nabu blockstore as our BlockStore. +// Nabu is JVM-only, hence this adapter lives in jvmMain. + +class NabuBlockStoreAdapter( + private val backing: RamBlockstore = RamBlockstore(), +) : BlockStore { + + private fun nabuMultihash(hex: String): Multihash = + Multihash(Multihash.Type.sha2_256, hexToBytes(hex)) + + private fun nabuCid(hex: String): org.ipfs.cid.Cid = + org.ipfs.cid.Cid.build(1, org.ipfs.cid.Cid.Codec.Raw, nabuMultihash(hex)) + + override suspend fun put(block: ContentBlock) { + val stored = backing.put(block.bytes, org.ipfs.cid.Cid.Codec.Raw).get(30, TimeUnit.SECONDS) + check(stored.bareMultihash() == nabuMultihash(block.cid.hashHex)) { "cid mismatch after put" } + } + + override suspend fun get(cid: jp.orgflow.contentstore.Cid): ContentBlock? { + val bytes = backing.get(nabuCid(cid.hashHex)).get(30, TimeUnit.SECONDS) ?: return null + return ContentBlock(cid, bytes) + } +} + +private fun hexToBytes(hex: String): ByteArray = ByteArray(hex.length / 2) { i -> + ((Character.digit(hex[i * 2], 16) shl 4) + Character.digit(hex[i * 2 + 1], 16)).toByte() +} diff --git a/modules/orgflow-content-store/src/jvmTest/kotlin/jp/orgflow/contentstore/NabuAdapterTest.kt b/modules/orgflow-content-store/src/jvmTest/kotlin/jp/orgflow/contentstore/NabuAdapterTest.kt new file mode 100644 index 0000000..32befb3 --- /dev/null +++ b/modules/orgflow-content-store/src/jvmTest/kotlin/jp/orgflow/contentstore/NabuAdapterTest.kt @@ -0,0 +1,18 @@ +package jp.orgflow.contentstore + +import jp.orgflow.contentstore.jvm.NabuBlockStoreAdapter +import kotlin.test.Test +import kotlin.test.assertEquals + +class NabuAdapterTest { + @Test + fun putGetRoundTripThroughNabu() { + val adapter = NabuBlockStoreAdapter() + val bytes = "nabu vendored core works".encodeToByteArray() + val cid = ContentAddresser.cidOf(bytes) + val block = ContentBlock(cid, bytes) + + adapter.put(block) + assertEquals(block, adapter.get(cid)) + } +} diff --git a/modules/vendor-nabu/src/main/java/org/peergos/Hash.java b/modules/vendor-nabu/src/main/java/org/peergos/Hash.java new file mode 100644 index 0000000..6c88628 --- /dev/null +++ b/modules/vendor-nabu/src/main/java/org/peergos/Hash.java @@ -0,0 +1,15 @@ +package org.peergos; + +import java.security.*; + +public class Hash { + + public static byte[] sha256(byte[] in) { + try { + MessageDigest hasher = MessageDigest.getInstance("SHA-256"); + return hasher.digest(in); + } catch (NoSuchAlgorithmException e) { + throw new RuntimeException(e); + } + } +} diff --git a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/Blockstore.java b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/Blockstore.java index f845b3a..38bec0c 100644 --- a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/Blockstore.java +++ b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/Blockstore.java @@ -3,8 +3,6 @@ package org.peergos.blockstore; import io.ipfs.cid.*; import io.ipfs.multibase.binary.Base32; import io.ipfs.multihash.Multihash; -import org.peergos.blockstore.metadatadb.BlockMetadata; -import org.peergos.blockstore.metadatadb.BlockMetadataStore; import java.util.*; import java.util.concurrent.*; @@ -40,6 +38,4 @@ public interface Blockstore { CompletableFuture<Boolean> applyToAll(Consumer<Cid> action, boolean useBlockstore); CompletableFuture<Boolean> bloomAdd(Cid cid); - - CompletableFuture<BlockMetadata> getBlockMetadata(Cid h); }
\ No newline at end of file diff --git a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/FileBlockstore.java b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/FileBlockstore.java index 7bcd40c..abb6c3d 100644 --- a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/FileBlockstore.java +++ b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/FileBlockstore.java @@ -3,7 +3,6 @@ package org.peergos.blockstore; import io.ipfs.cid.Cid; import io.ipfs.multihash.Multihash; import org.peergos.Hash; -import org.peergos.blockstore.metadatadb.BlockMetadata; import org.peergos.util.*; import java.io.*; @@ -181,9 +180,4 @@ public class FileBlockstore implements Blockstore { throw new RuntimeException(e); } } - - @Override - public CompletableFuture<BlockMetadata> getBlockMetadata(Cid h) { - throw new IllegalStateException("Unsupported operation!"); - } } diff --git a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/RamBlockstore.java b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/RamBlockstore.java index 5b67a63..4dc7438 100644 --- a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/RamBlockstore.java +++ b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/RamBlockstore.java @@ -3,7 +3,6 @@ package org.peergos.blockstore; import io.ipfs.cid.*; import io.ipfs.multihash.*; import org.peergos.*; -import org.peergos.blockstore.metadatadb.BlockMetadata; import org.peergos.cbor.*; import org.peergos.util.*; @@ -70,10 +69,4 @@ public class RamBlockstore implements Blockstore { blocks.keySet().stream().forEach(action); return Futures.of(true); } - - @Override - public CompletableFuture<BlockMetadata> getBlockMetadata(Cid h) { - byte[] block = get(h).join().get(); - return Futures.of(new BlockMetadata(block.length, CborObject.getLinks(h, block))); - } } diff --git a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/TypeLimitedBlockstore.java b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/TypeLimitedBlockstore.java index 9e3ed29..49dc554 100644 --- a/modules/vendor-nabu/src/main/java/org/peergos/blockstore/TypeLimitedBlockstore.java +++ b/modules/vendor-nabu/src/main/java/org/peergos/blockstore/TypeLimitedBlockstore.java @@ -2,7 +2,6 @@ package org.peergos.blockstore; import io.ipfs.cid.Cid; import io.ipfs.multihash.*; -import org.peergos.blockstore.metadatadb.BlockMetadata; import org.peergos.util.*; import java.util.List; @@ -86,10 +85,4 @@ public class TypeLimitedBlockstore implements Blockstore { public CompletableFuture<Boolean> applyToAll(Consumer<Cid> action, boolean useBlockstore) { return blocks.applyToAll(action, useBlockstore); } - - @Override - public CompletableFuture<BlockMetadata> getBlockMetadata(Cid h) { - return blocks.getBlockMetadata(h); - } - } |
