get_roots
create_wallet
rename
This commit is contained in:
almog
2022-01-14 17:58:26 +02:00
parent 831f18cd89
commit b0cdb93e88
8 changed files with 463 additions and 191 deletions
+19 -7
View File
@@ -90,16 +90,16 @@ def create_rpc_port_option() -> "IdentityFunction":
)
@data_cmd.command("create_kv_store", short_help="Get a data row by its hash")
@data_cmd.command("create_data_store", short_help="Get a data row by its hash")
@create_kv_store_id_option()
@create_rpc_port_option()
def create_kv_store(
def create_data_store(
table_string: str,
data_rpc_port: int,
) -> None:
from chia.cmds.data_funcs import create_kv_store_cmd
from chia.cmds.data_funcs import create_data_store_cmd
run(create_kv_store_cmd(rpc_port=data_rpc_port, table_string=table_string))
run(create_data_store_cmd(rpc_port=data_rpc_port, table_string=table_string))
@data_cmd.command("get_value", short_help="Get a data row by its hash")
@@ -116,7 +116,7 @@ def get_value(
run(get_value_cmd(rpc_port=data_rpc_port, tree_id=tree_id, key=key))
@data_cmd.command("update_kv_store", short_help="Update a table.")
@data_cmd.command("update_data_store", short_help="Update a table.")
@create_kv_store_id_option()
@create_rpc_port_option()
@create_changelist_option()
@@ -125,8 +125,20 @@ def update_kv_store(
changelist_string: str,
data_rpc_port: int,
) -> None:
from chia.cmds.data_funcs import update_kv_store_cmd
from chia.cmds.data_funcs import update_data_store_cmd
changelist = json.loads(changelist_string)
run(update_kv_store_cmd(rpc_port=data_rpc_port, tree_id=tree_id, changelist=changelist))
run(update_data_store_cmd(rpc_port=data_rpc_port, tree_id=tree_id, changelist=changelist))
@data_cmd.command("get_root", short_help="Get a data row by its hash")
@create_kv_store_id_option()
@create_rpc_port_option()
def get_root(
tree_id: str,
data_rpc_port: int,
) -> None:
from chia.cmds.data_funcs import get_root_cmd
run(get_root_cmd(rpc_port=data_rpc_port, tree_id=tree_id))
+45 -10
View File
@@ -24,12 +24,29 @@ async def get_client(rpc_port: Optional[int]) -> Tuple[DataLayerRpcClient, int]:
return client, rpc_port
async def create_kv_store_cmd(rpc_port: Optional[int], table_string: str) -> Optional[Dict[str, Any]]:
async def create_wallet(rpc_port: Optional[int]) -> Optional[Dict[str, Any]]:
# TODO: nice cli error handling
try:
client, rpc_port = await get_client(rpc_port)
response = await client.create_wallet()
except aiohttp.ClientConnectorError:
print(f"Connection error. Check if data is running at {rpc_port}")
return None
except Exception as e:
print(f"Exception from 'data': {e}")
return None
client.close()
await client.await_closed()
return response
async def create_data_store_cmd(rpc_port: Optional[int], table_string: str) -> Optional[Dict[str, Any]]:
# TODO: nice cli error handling
try:
client, rpc_port = await get_client(rpc_port)
response = await client.create_kv_store()
response = await client.create_data_store()
except aiohttp.ClientConnectorError:
print(f"Connection error. Check if data is running at {rpc_port}")
return None
@@ -45,12 +62,11 @@ async def create_kv_store_cmd(rpc_port: Optional[int], table_string: str) -> Opt
async def get_value_cmd(rpc_port: Optional[int], tree_id: str, key: str) -> Optional[Dict[str, Any]]:
# TODO: nice cli error handling
tree_id_bytes = bytes32(hexstr_to_bytes(tree_id))
store_id_bytes = bytes32(hexstr_to_bytes(tree_id))
key_bytes = hexstr_to_bytes(key)
try:
client, rpc_port = await get_client(rpc_port)
response = await client.get_value(tree_id=tree_id_bytes, key=key_bytes)
print(json.dumps(response, indent=4))
response = await client.get_value(store_id=store_id_bytes, key=key_bytes)
except aiohttp.ClientConnectorError:
print(f"Connection error. Check if data is running at {rpc_port}")
return None
@@ -63,18 +79,37 @@ async def get_value_cmd(rpc_port: Optional[int], tree_id: str, key: str) -> Opti
return response
async def update_kv_store_cmd(
async def update_data_store_cmd(
rpc_port: Optional[int],
tree_id: str,
changelist: Dict[str, str],
) -> Optional[Dict[str, Any]]:
# TODO: nice cli error handling
tree_id_bytes = bytes32(hexstr_to_bytes(tree_id))
store_id_bytes = bytes32(hexstr_to_bytes(tree_id))
try:
client, rpc_port = await get_client(rpc_port)
response = await client.update_kv_store(tree_id=tree_id_bytes, changelist=changelist)
print(json.dumps(response, indent=4))
response = await client.update_data_store(store_id=store_id_bytes, changelist=changelist)
except aiohttp.ClientConnectorError:
print(f"Connection error. Check if data is running at {rpc_port}")
return None
except Exception as e:
print(f"Exception from 'data': {e}")
return None
client.close()
await client.await_closed()
return response
async def get_root_cmd(
rpc_port: Optional[int],
tree_id: str,
) -> Optional[Dict[str, Any]]:
# TODO: nice cli error handling
store_id_bytes = bytes32(hexstr_to_bytes(tree_id))
try:
client, rpc_port = await get_client(rpc_port)
response = await client.get_root(store_id=store_id_bytes)
except aiohttp.ClientConnectorError:
print(f"Connection error. Check if data is running at {rpc_port}")
return None
+44 -14
View File
@@ -9,11 +9,13 @@ from chia.data_layer.data_layer_wallet import DataLayerWallet
from chia.data_layer.data_store import DataStore
from chia.server.server import ChiaServer
from chia.types.blockchain_format.sized_bytes import bytes32
from chia.util.byte_types import hexstr_to_bytes
from chia.util.config import load_config
from chia.util.db_wrapper import DBWrapper
from chia.util.ints import uint64
from chia.util.path import mkdir, path_from_root
from chia.wallet.wallet_node import WalletNode
from chia.wallet.transaction_record import TransactionRecord
class DataLayer:
@@ -69,16 +71,21 @@ class DataLayer:
self.connection = await aiosqlite.connect(self.db_path)
self.db_wrapper = DBWrapper(self.connection)
self.data_store = await DataStore.create(self.db_wrapper)
return True
async def create_wallet(self) -> Optional[List[TransactionRecord]]:
# TODO: review for anything else we need to do here
assert self.wallet_node
assert self.wallet_node.wallet_state_manager
main_wallet = self.wallet_node.wallet_state_manager.main_wallet
amount = uint64(1) # todo what should amount be ?
async with self.wallet_node.wallet_state_manager.lock:
creation_record = await self.wallet.create_new_dl_wallet(
creation_record = await DataLayerWallet.create_new_dl_wallet(
self.wallet_node.wallet_state_manager, main_wallet, amount, None
)
self.wallet = creation_record.item
self.initialized = True
return True
return creation_record.transaction_records
def _close(self) -> None:
# TODO: review for anything else we need to do here
@@ -97,11 +104,11 @@ class DataLayer:
self.log.fatal("failed creating store")
return tree_id
async def insert(
async def batch_update(
self,
tree_id: bytes32,
changelist: List[Dict[str, Any]],
) -> bool:
) -> Optional[TransactionRecord]:
for change in changelist:
if change["action"] == "insert":
key = change["key"]
@@ -116,27 +123,50 @@ class DataLayer:
key = change["key"]
await self.data_store.delete(key, tree_id)
await self.data_store.get_tree_root(tree_id)
root = await self.data_store.get_tree_root(tree_id)
assert root.node_hash
res = await self.wallet.create_update_state_spend(root.node_hash)
# todo return empty node hash from get_tree_root
if root.node_hash is not None:
node_hash = root.node_hash
else:
node_hash = bytes32([0] * 32) # todo change
res = await self.wallet.create_update_state_spend(node_hash)
assert res
# todo need to mark data as pending and change once tx is confirmed
return True
# todo register callback to change status in data store
# await self.data_store.change_root_status(root, Status.COMMITTED)
return None
async def get_value(self, store_id: bytes32, key: bytes) -> bytes:
async def get_value(self, store_id: bytes32, key: bytes) -> Optional[bytes]:
res = await self.data_store.get_node_by_key(tree_id=store_id, key=key)
if res is None:
self.log.error("Failed to create tree")
self.log.error("Failed to fetch key")
return None
return res.value
async def get_pairs(self, store_id: bytes32) -> List[TerminalNode]:
res = await self.data_store.get_pairs(store_id)
async def get_keys_values(self, store_id: bytes32) -> List[TerminalNode]:
res = await self.data_store.get_keys_values(store_id)
if res is None:
self.log.error("Failed to create tree")
self.log.error("Failed to fetch keys values")
return res
async def get_ancestors(self, node_hash: bytes32, store_id: bytes32) -> List[InternalNode]:
res = await self.data_store.get_ancestors(store_id, node_hash)
if res is None:
self.log.error("Failed to create tree")
self.log.error("Failed to get ancestors")
return res
async def get_root(self, store_id: bytes32) -> Optional[bytes32]:
res = await self.data_store.get_tree_root(tree_id=store_id)
if res is None:
self.log.error(f"Failed to get root for {store_id.hex()}")
return res.node_hash
async def get_roots(self, store_ids: List[str]) -> List[Optional[bytes32]]:
roots = []
for id in store_ids:
res = await self.data_store.get_tree_root(tree_id=bytes32(hexstr_to_bytes(id)))
if res is None:
self.log.error(f"Failed to get root for {id}")
continue
roots.append(res.node_hash)
return roots
+4 -4
View File
@@ -350,7 +350,7 @@ class DataStore:
return ancestors
async def get_pairs(self, tree_id: bytes32, *, lock: bool = True) -> List[TerminalNode]:
async def get_keys_values(self, tree_id: bytes32, *, lock: bool = True) -> List[TerminalNode]:
async with self.db_wrapper.locked_transaction(lock=lock):
root = await self.get_tree_root(tree_id=tree_id, lock=False)
@@ -421,7 +421,7 @@ class DataStore:
lock: bool = True,
) -> bytes32:
async with self.db_wrapper.locked_transaction(lock=lock):
pairs = await self.get_pairs(tree_id=tree_id, lock=False)
pairs = await self.get_keys_values(tree_id=tree_id, lock=False)
if len(pairs) == 0:
reference_node_hash = None
@@ -456,7 +456,7 @@ class DataStore:
if not was_empty:
# TODO: is there any way the db can enforce this?
pairs = await self.get_pairs(tree_id=tree_id, lock=False)
pairs = await self.get_keys_values(tree_id=tree_id, lock=False)
if any(key == node.key for node in pairs):
raise Exception(f"Key already present: {key.hex()}")
@@ -561,7 +561,7 @@ class DataStore:
async def get_node_by_key(self, key: bytes, tree_id: bytes32, *, lock: bool = True) -> TerminalNode:
async with self.db_wrapper.locked_transaction(lock=lock):
nodes = await self.get_pairs(tree_id=tree_id, lock=False)
nodes = await self.get_keys_values(tree_id=tree_id, lock=False)
for node in nodes:
if node.key == key:
+103 -16
View File
@@ -1,3 +1,4 @@
import dataclasses
from typing import Any, Callable, Dict, Optional
@@ -8,6 +9,7 @@ from chia.types.blockchain_format.sized_bytes import bytes32
from chia.util.byte_types import hexstr_to_bytes
# todo input assertions for all rpc's
from chia.util.streamable import recurse_jsonify
def process_change(change: Dict[str, Any]) -> Dict[str, Any]:
@@ -43,37 +45,55 @@ class DataLayerRpcApi:
def get_routes(self) -> Dict[str, Callable[[Any], Any]]:
return {
"/create_kv_store": self.create_kv_store,
"/update_kv_store": self.update_kv_store,
"/create_data_store": self.create_data_store,
"/batch_update": self.batch_update,
"/get_value": self.get_value,
"/get_pairs": self.get_pairs,
"/get_keys_values": self.get_keys_values,
"/get_ancestors": self.get_ancestors,
"/get_root": self.get_root,
"/get_roots": self.get_root,
"/delete_key": self.delete_key,
"/insert": self.insert,
"/create_wallet": self.create_wallet,
}
async def create_kv_store(self, request: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
async def create_data_store(self, request: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
if self.service is None:
raise Exception("Data layer not created")
value = await self.service.create_store()
return {"id": value.hex()}
async def get_value(self, request: Dict[str, Any]) -> Dict[str, Any]:
store_id = bytes32.from_hexstr(request["id"])
key = hexstr_to_bytes(request["key"])
if self.service is None:
raise Exception("Data layer not created")
value = await self.service.get_value(store_id=store_id, key=key)
return {"data": value.hex()}
hex = None
if value is not None:
hex = value.hex()
return {"value": hex}
async def get_pairs(self, request: Dict[str, Any]) -> Dict[str, Any]:
async def get_keys_values(self, request: Dict[str, Any]) -> Dict[str, Any]:
store_id = bytes32(hexstr_to_bytes(request["id"]))
value = await self.service.get_pairs(store_id)
# TODO: fix
return {"data": value.hex()} # type: ignore[attr-defined]
if self.service is None:
raise Exception("Data layer not created")
res = await self.service.get_keys_values(store_id)
json_nodes = []
for node in res:
json = recurse_jsonify(dataclasses.asdict(node)) # type: ignore[no-untyped-call]
json_nodes.append(json)
return {"keys_values": json_nodes}
async def get_ancestors(self, request: Dict[str, Any]) -> Dict[str, Any]:
store_id = bytes32(hexstr_to_bytes(request["id"]))
key = hexstr_to_bytes(request["key"])
# TODO: fix
value = await self.service.get_ancestors(key, store_id) # type:ignore[arg-type]
# TODO: fix
return {"data": value.hex()} # type: ignore[attr-defined]
node_hash = bytes32.from_hexstr(request["hash"])
if self.service is None:
raise Exception("Data layer not created")
value = await self.service.get_ancestors(node_hash, store_id)
return {"ancestors": value}
async def update_kv_store(self, request: Dict[str, Any]) -> None:
async def batch_update(self, request: Dict[str, Any]) -> Dict[str, Any]:
"""
rows_to_add a list of clvm objects as bytes to add to talbe
rows_to_remove a list of row hashes to remove
@@ -81,4 +101,71 @@ class DataLayerRpcApi:
changelist = [process_change(change) for change in request["changelist"]]
store_id = bytes32(hexstr_to_bytes(request["id"]))
# todo input checks
await self.service.insert(store_id, changelist)
if self.service is None:
raise Exception("Data layer not created")
await self.service.batch_update(store_id, changelist)
return {"tx_id": "id"}
async def insert(self, request: Dict[str, Any]) -> Dict[str, Any]:
"""
rows_to_add a list of clvm objects as bytes to add to talbe
rows_to_remove a list of row hashes to remove
"""
key = hexstr_to_bytes(request["key"])
value = hexstr_to_bytes(request["value"])
store_id = bytes32(hexstr_to_bytes(request["id"]))
# todo input checks
if self.service is None:
raise Exception("Data layer not created")
changelist = [{"action": "insert", "key": key.hex(), "value": value.hex()}]
txs = await self.service.batch_update(store_id, changelist)
return {"tx_id": "id"}
async def delete_key(self, request: Dict[str, Any]) -> Dict[str, Any]:
"""
rows_to_add a list of clvm objects as bytes to add to talbe
rows_to_remove a list of row hashes to remove
"""
key = hexstr_to_bytes(request["key"])
store_id = bytes32(hexstr_to_bytes(request["id"]))
# todo input checks
if self.service is None:
raise Exception("Data layer not created")
changelist = [{"action": "delete", "key": key.hex()}]
txs = await self.service.batch_update(store_id, changelist)
return {"tx_id": "id"}
async def get_root(self, request: Dict[str, Any]) -> Dict[str, Any]:
"""
rows_to_add a list of clvm objects as bytes to add to talbe
rows_to_remove a list of row hashes to remove
"""
store_id = bytes32(hexstr_to_bytes(request["id"]))
# todo input checks
if self.service is None:
raise Exception("Data layer not created")
res = await self.service.get_root(store_id)
return {"hash": res}
async def get_roots(self, request: Dict[str, Any]) -> Dict[str, Any]:
"""
rows_to_add a list of clvm objects as bytes to add to talbe
rows_to_remove a list of row hashes to remove
"""
store_ids = request["ids"]
# todo input checks
if self.service is None:
raise Exception("Data layer not created")
res = await self.service.get_roots(store_ids)
return {"hash": res}
async def create_wallet(self, request: Dict[str, Any]) -> Dict[str, Any]:
"""
rows_to_add a list of clvm objects as bytes to add to talbe
rows_to_remove a list of row hashes to remove
"""
# todo input checks
if self.service is None:
raise Exception("Data layer not created")
res = await self.service.create_wallet()
return {"tx_ids": res}
+31 -9
View File
@@ -1,24 +1,46 @@
from typing import Any, Dict
from typing import Any, Dict, List
from chia.rpc.rpc_client import RpcClient
from chia.types.blockchain_format.sized_bytes import bytes32
class DataLayerRpcClient(RpcClient):
async def create_kv_store(self) -> Dict[str, Any]:
response = await self.fetch("create_kv_store", {})
async def create_wallet(self) -> Dict[str, Any]:
response = await self.fetch("create_wallet", {})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def get_value(self, tree_id: bytes32, key: bytes) -> Dict[str, Any]:
response = await self.fetch("get_value", {"tree_id": tree_id.hex(), "key": key.hex()})
async def create_data_store(self) -> Dict[str, Any]:
response = await self.fetch("create_data_store", {})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def update_kv_store(self, tree_id: bytes32, changelist: Dict[str, str]) -> Dict[str, Any]:
response = await self.fetch("update_kv_store", {"tree_id": tree_id, "changelist": changelist})
async def get_value(self, store_id: bytes32, key: bytes) -> Dict[str, Any]:
response = await self.fetch("get_value", {"id": store_id.hex(), "key": key.hex()})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def get_tree_state(self, tree_id: bytes32) -> bytes32:
pass
async def update_data_store(self, store_id: bytes32, changelist: Dict[str, str]) -> Dict[str, Any]:
response = await self.fetch("batch_update", {"id": store_id.hex(), "changelist": changelist})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def get_keys_values(self, store_id: bytes32) -> Dict[str, Any]:
response = await self.fetch("get_keys_values", {"id": store_id.hex()})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def get_ancestors(self, store_id: bytes32, hash: bytes32) -> Dict[str, Any]:
response = await self.fetch("get_ancestors", {"id": store_id.hex(), "hash": hash})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def get_root(self, store_id: bytes32) -> Dict[str, Any]:
response = await self.fetch("get_root", {"id": store_id.hex()})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
async def get_roots(self, store_ids: List[bytes32]) -> Dict[str, Any]:
response = await self.fetch("get_roots", {"ids": store_ids})
# TODO: better hinting for .fetch() (probably a TypedDict)
return response # type: ignore[no-any-return]
+215 -129
View File
@@ -1,10 +1,11 @@
from typing import AsyncIterator, Dict, List, Tuple
import asyncio
from typing import AsyncIterator, Dict, List, Tuple, Any
import pytest
import aiosqlite
# flake8: noqa: F401
from chia.data_layer.data_layer import DataLayer
from chia.data_layer.data_store import DataStore
from chia.rpc.data_layer_rpc_api import DataLayerRpcApi
from chia.rpc.wallet_rpc_api import WalletRpcApi
from chia.server.server import ChiaServer
from chia.simulator.full_node_simulator import FullNodeSimulator
from chia.simulator.simulator_protocol import FarmNewBlockProtocol
@@ -12,27 +13,43 @@ from chia.types.blockchain_format.sized_bytes import bytes32
from chia.types.peer_info import PeerInfo
from chia.util.byte_types import hexstr_to_bytes
from chia.util.config import load_config
from chia.util.db_wrapper import DBWrapper
from chia.util.ints import uint16
from chia.wallet.transaction_record import TransactionRecord
from chia.wallet.wallet_node import WalletNode
from tests.core.data_layer.util import ChiaRoot
from tests.setup_nodes import setup_simulators_and_wallets, self_hostname
from tests.time_out_assert import time_out_assert
from tests.wallet.rl_wallet.test_rl_rpc import is_transaction_confirmed
pytestmark = pytest.mark.data_layer
nodes = Tuple[List[FullNodeSimulator], List[Tuple[WalletNode, ChiaServer]]]
async def init_data_layer(
full_node_api: Any, num_blocks: int, ph: bytes32, wallet_node: WalletNode, wallet_rpc_api: WalletRpcApi
) -> DataLayerRpcApi:
assert wallet_node.wallet_state_manager
data_layer = DataLayer(wallet_node.root_path, wallet_node)
data_rpc_api = DataLayerRpcApi(data_layer)
# res = await data_rpc_api.create_data_layer({"fee": 1})
# await asyncio.sleep(1)
# assert res["result"]
# tx0: TransactionRecord = res["result"][0]
# tx1: TransactionRecord = res["result"][1]
# await asyncio.sleep(1)
# for i in range(0, num_blocks):
# await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
# await time_out_assert(15, is_transaction_confirmed, True, tx0.wallet_id, wallet_rpc_api, tx0.name)
# await time_out_assert(15, is_transaction_confirmed, True, tx1.wallet_id, wallet_rpc_api, tx1.name)
return data_rpc_api
@pytest.fixture(scope="function")
async def one_wallet_node() -> AsyncIterator[nodes]:
async for _ in setup_simulators_and_wallets(1, 1, {}):
yield _
# TODO: fix this
@pytest.mark.xfail(reason="incomplete, needs caught up", strict=True)
@pytest.mark.asyncio
async def test_create_insert_get(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
root = chia_root.path
@@ -46,40 +63,42 @@ async def test_create_insert_get(chia_root: ChiaRoot, one_wallet_node: nodes) ->
assert wallet_node.wallet_state_manager is not None
wallet = wallet_node.wallet_state_manager.main_wallet
ph = await wallet.get_new_puzzlehash()
await server_2.start_client(PeerInfo(self_hostname, uint16(server_1._port)), None)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
print(f"confirmed balance is {await wallet.get_confirmed_balance()}")
print(f"unconfirmed balance is {await wallet.get_unconfirmed_balance()}")
data_layer = DataLayer(root_path=root, wallet_node=wallet_node)
async with aiosqlite.connect(data_layer.db_path) as connection:
data_layer.connection = connection
data_layer.db_wrapper = DBWrapper(data_layer.connection)
data_layer.data_store = await DataStore.create(data_layer.db_wrapper)
data_layer.initialized = True
rpc_api = DataLayerRpcApi(data_layer)
key = b"a"
value = b"\x00\x01"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key.hex(), "value": value.hex()}]
res = await rpc_api.create_kv_store()
store_id = bytes32(hexstr_to_bytes(res["id"]))
print(f"store id is {store_id}")
spendable = await wallet.get_spendable_balance()
print(f"spendable balance is {spendable}")
await rpc_api.update_kv_store({"id": store_id.hex(), "changelist": changelist})
res = await rpc_api.get_value({"id": store_id.hex(), "key": key.hex()})
assert hexstr_to_bytes(res["data"]) == value
changelist = [{"action": "delete", "key": key.hex()}]
await rpc_api.update_kv_store({"id": store_id.hex(), "changelist": changelist})
with pytest.raises(Exception):
await rpc_api.get_value({"id": store_id.hex(), "key": key.hex()})
wallet_rpc_api = WalletRpcApi(wallet_node)
data_rpc_api = await init_data_layer(full_node_api, num_blocks, ph, wallet_node, wallet_rpc_api)
key = b"a"
value = b"\x00\x01"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key.hex(), "value": value.hex()}]
res = await data_rpc_api.create_data_store()
await asyncio.sleep(1)
assert res is not None
store_id = bytes32(hexstr_to_bytes(res["id"]))
res = await data_rpc_api.batch_update({"id": store_id.hex(), "changelist": changelist})
update_tx_rec0 = res["tx_id"]
await asyncio.sleep(1)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec0.wallet_id, wallet_rpc_api, update_tx_rec0.name
)
res = await data_rpc_api.get_value({"id": store_id.hex(), "key": key.hex()})
assert hexstr_to_bytes(res["data"]) == value
changelist = [{"action": "delete", "key": key.hex()}]
res = await data_rpc_api.batch_update({"id": store_id.hex(), "changelist": changelist})
update_tx_rec1 = res["tx_id"]
await asyncio.sleep(1)
for i in range(0, num_blocks):
await asyncio.sleep(1)
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec1.wallet_id, wallet_rpc_api, update_tx_rec1.name
)
with pytest.raises(Exception):
val = await data_rpc_api.get_value({"id": store_id.hex(), "key": key.hex()})
# TODO: fix this
@pytest.mark.xfail(reason="incomplete, needs caught up", strict=True)
@pytest.mark.asyncio
async def test_create_double_insert(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
root = chia_root.path
@@ -93,47 +112,52 @@ async def test_create_double_insert(chia_root: ChiaRoot, one_wallet_node: nodes)
assert wallet_node.wallet_state_manager is not None
wallet = wallet_node.wallet_state_manager.main_wallet
ph = await wallet.get_new_puzzlehash()
await server_2.start_client(PeerInfo(self_hostname, uint16(server_1._port)), None)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
print(f"confirmed balance is {await wallet.get_confirmed_balance()}")
print(f"unconfirmed balance is {await wallet.get_unconfirmed_balance()}")
wallet_rpc_api = WalletRpcApi(wallet_node)
data_rpc_api = await init_data_layer(full_node_api, num_blocks, ph, wallet_node, wallet_rpc_api)
key1 = b"a"
value1 = b"\x01\x02"
key2 = b"b"
value2 = b"\x01\x23"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
res = await data_rpc_api.create_data_store()
store_id = bytes32(hexstr_to_bytes(res["id"]))
res = await data_rpc_api.batch_update({"id": store_id.hex(), "changelist": changelist})
update_tx_rec0 = res["tx_id"]
await asyncio.sleep(1)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec0.wallet_id, wallet_rpc_api, update_tx_rec0.name
)
data_layer = DataLayer(root_path=root, wallet_node=wallet_node)
async with aiosqlite.connect(data_layer.db_path) as connection:
data_layer.connection = connection
data_layer.db_wrapper = DBWrapper(data_layer.connection)
data_layer.data_store = await DataStore.create(data_layer.db_wrapper)
data_layer.initialized = True
rpc_api = DataLayerRpcApi(data_layer)
key1 = b"a"
value1 = b"\x01\x02"
key2 = b"b"
value2 = b"\x01\x23"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
res = await rpc_api.create_kv_store()
store_id = bytes32(hexstr_to_bytes(res["id"]))
await rpc_api.update_kv_store({"id": store_id.hex(), "changelist": changelist})
res = await rpc_api.get_value({"id": store_id.hex(), "key": key1.hex()})
assert hexstr_to_bytes(res["data"]) == value1
res = await data_rpc_api.get_value({"id": store_id.hex(), "key": key1.hex()})
assert hexstr_to_bytes(res["data"]) == value1
changelist = [{"action": "insert", "key": key2.hex(), "value": value2.hex()}]
await rpc_api.update_kv_store({"id": store_id.hex(), "changelist": changelist})
res = await rpc_api.get_value({"id": store_id.hex(), "key": key2.hex()})
assert hexstr_to_bytes(res["data"]) == value2
changelist = [{"action": "insert", "key": key2.hex(), "value": value2.hex()}]
res = await data_rpc_api.batch_update({"id": store_id.hex(), "changelist": changelist})
update_tx_rec1 = res["tx_id"]
await asyncio.sleep(1)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec0.wallet_id, wallet_rpc_api, update_tx_rec1.name
)
res = await data_rpc_api.get_value({"id": store_id.hex(), "key": key2.hex()})
assert hexstr_to_bytes(res["data"]) == value2
changelist = [{"action": "delete", "key": key1.hex()}]
await rpc_api.update_kv_store({"id": store_id.hex(), "changelist": changelist})
with pytest.raises(Exception):
await rpc_api.get_value({"id": store_id.hex(), "key": key1.hex()})
changelist = [{"action": "delete", "key": key1.hex()}]
await data_rpc_api.batch_update({"id": store_id.hex(), "changelist": changelist})
with pytest.raises(Exception):
val = await data_rpc_api.get_value({"id": store_id.hex(), "key": key1.hex()})
# TODO: fix this
@pytest.mark.xfail(reason="incomplete, needs caught up", strict=True)
@pytest.mark.asyncio
async def test_get_pairs(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
async def test_get_keys_values(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
root = chia_root.path
config = load_config(root, "config.yaml")
config["data_layer"]["database_path"] = "data_layer_test.sqlite"
@@ -145,45 +169,50 @@ async def test_get_pairs(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
assert wallet_node.wallet_state_manager is not None
wallet = wallet_node.wallet_state_manager.main_wallet
ph = await wallet.get_new_puzzlehash()
await server_2.start_client(PeerInfo(self_hostname, uint16(server_1._port)), None)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
print(f"confirmed balance is {await wallet.get_confirmed_balance()}")
print(f"unconfirmed balance is {await wallet.get_unconfirmed_balance()}")
data_layer = DataLayer(root_path=root, wallet_node=wallet_node)
async with aiosqlite.connect(data_layer.db_path) as connection:
data_layer.connection = connection
data_layer.db_wrapper = DBWrapper(data_layer.connection)
data_layer.data_store = await DataStore.create(data_layer.db_wrapper)
data_layer.initialized = True
rpc_api = DataLayerRpcApi(data_layer)
key1 = b"a"
value1 = b"\x01\x02"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
key2 = b"b"
value2 = b"\x03\x02"
changelist.append({"action": "insert", "key": key2.hex(), "value": value2.hex()})
key3 = b"c"
value3 = b"\x04\x05"
changelist.append({"action": "insert", "key": key3.hex(), "value": value3.hex()})
key4 = b"d"
value4 = b"\x06\x03"
changelist.append({"action": "insert", "key": key4.hex(), "value": value4.hex()})
key5 = b"e"
value5 = b"\x07\x01"
changelist.append({"action": "insert", "key": key5.hex(), "value": value5.hex()})
res = await rpc_api.create_kv_store()
tree_id = bytes32(hexstr_to_bytes(res["id"]))
await rpc_api.update_kv_store({"id": tree_id.hex(), "changelist": changelist})
await rpc_api.get_pairs({"id": tree_id.hex()})
# todo check values match
wallet_rpc_api = WalletRpcApi(wallet_node)
data_rpc_api = await init_data_layer(full_node_api, num_blocks, ph, wallet_node, wallet_rpc_api)
key1 = b"a"
value1 = b"\x01\x02"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
key2 = b"b"
value2 = b"\x03\x02"
changelist.append({"action": "insert", "key": key2.hex(), "value": value2.hex()})
key3 = b"c"
value3 = b"\x04\x05"
changelist.append({"action": "insert", "key": key3.hex(), "value": value3.hex()})
key4 = b"d"
value4 = b"\x06\x03"
changelist.append({"action": "insert", "key": key4.hex(), "value": value4.hex()})
key5 = b"e"
value5 = b"\x07\x01"
changelist.append({"action": "insert", "key": key5.hex(), "value": value5.hex()})
res = await data_rpc_api.create_data_store()
tree_id = bytes32(hexstr_to_bytes(res["id"]))
res = await data_rpc_api.batch_update({"id": tree_id.hex(), "changelist": changelist})
update_tx_rec0 = res["tx_id"]
await asyncio.sleep(1)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec0.wallet_id, wallet_rpc_api, update_tx_rec0.name
)
val = await data_rpc_api.get_keys_values({"id": tree_id.hex()})
dic = {}
for item in val["data"]:
dic[item.key] = item.value
assert dic[key1] == value1
assert dic[key2] == value2
assert dic[key3] == value3
assert dic[key4] == value4
assert dic[key5] == value5
# todo check values match
# TODO: fix this
@pytest.mark.xfail(reason="incomplete, needs caught up", strict=True)
@pytest.mark.asyncio
async def test_get_ancestors(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
root = chia_root.path
@@ -197,38 +226,95 @@ async def test_get_ancestors(chia_root: ChiaRoot, one_wallet_node: nodes) -> Non
assert wallet_node.wallet_state_manager
wallet = wallet_node.wallet_state_manager.main_wallet
ph = await wallet.get_new_puzzlehash()
await server_2.start_client(PeerInfo(self_hostname, uint16(server_1._port)), None)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
print(f"confirmed balance is {await wallet.get_confirmed_balance()}")
print(f"unconfirmed balance is {await wallet.get_unconfirmed_balance()}")
wallet_rpc_api = WalletRpcApi(wallet_node)
data_rpc_api = await init_data_layer(full_node_api, num_blocks, ph, wallet_node, wallet_rpc_api)
key1 = b"a"
value1 = b"\x01\x02"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
key2 = b"b"
value2 = b"\x03\x02"
changelist.append({"action": "insert", "key": key2.hex(), "value": value2.hex()})
key3 = b"c"
value3 = b"\x04\x05"
changelist.append({"action": "insert", "key": key3.hex(), "value": value3.hex()})
key4 = b"d"
value4 = b"\x06\x03"
changelist.append({"action": "insert", "key": key4.hex(), "value": value4.hex()})
key5 = b"e"
value5 = b"\x07\x01"
changelist.append({"action": "insert", "key": key5.hex(), "value": value5.hex()})
res = await data_rpc_api.create_data_store()
tree_id = bytes32(hexstr_to_bytes(res["id"]))
res = await data_rpc_api.batch_update({"id": tree_id.hex(), "changelist": changelist})
update_tx_rec0 = res["tx_id"]
await asyncio.sleep(1)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec0.wallet_id, wallet_rpc_api, update_tx_rec0.name
)
val = await data_rpc_api.get_keys_values({"id": tree_id.hex()})
assert val["data"]
val = await data_rpc_api.get_ancestors({"id": tree_id.hex(), "hash": val["data"][4].hash.hex()})
print(val)
# todo assert values
data_layer = DataLayer(root_path=root, wallet_node=wallet_node)
async with aiosqlite.connect(data_layer.db_path) as connection:
data_layer.connection = connection
data_layer.db_wrapper = DBWrapper(data_layer.connection)
data_layer.data_store = await DataStore.create(data_layer.db_wrapper)
data_layer.initialized = True
rpc_api = DataLayerRpcApi(data_layer)
key1 = b"a"
value1 = b"\x01\x02"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
key2 = b"b"
value2 = b"\x03\x02"
changelist.append({"action": "insert", "key": key2.hex(), "value": value2.hex()})
key3 = b"c"
value3 = b"\x04\x05"
changelist.append({"action": "insert", "key": key3.hex(), "value": value3.hex()})
key4 = b"d"
value4 = b"\x06\x03"
changelist.append({"action": "insert", "key": key4.hex(), "value": value4.hex()})
key5 = b"e"
value5 = b"\x07\x01"
changelist.append({"action": "insert", "key": key5.hex(), "value": value5.hex()})
res = await rpc_api.create_kv_store()
tree_id = bytes32(hexstr_to_bytes(res["id"]))
await rpc_api.update_kv_store({"id": tree_id.hex(), "changelist": changelist})
await rpc_api.get_ancestors({"id": tree_id.hex(), "key": key1.hex()})
# todo assert values
@pytest.mark.asyncio
async def test_get_roots(chia_root: ChiaRoot, one_wallet_node: nodes) -> None:
root = chia_root.path
config = load_config(root, "config.yaml")
config["data_layer"]["database_path"] = "data_layer_test.sqlite"
num_blocks = 5
full_nodes, wallets = one_wallet_node
full_node_api = full_nodes[0]
server_1 = full_node_api.full_node.server
wallet_node, server_2 = wallets[0]
assert wallet_node.wallet_state_manager
wallet = wallet_node.wallet_state_manager.main_wallet
ph = await wallet.get_new_puzzlehash()
await server_2.start_client(PeerInfo(self_hostname, uint16(server_1._port)), None)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
print(f"confirmed balance is {await wallet.get_confirmed_balance()}")
print(f"unconfirmed balance is {await wallet.get_unconfirmed_balance()}")
wallet_rpc_api = WalletRpcApi(wallet_node)
data_rpc_api = await init_data_layer(full_node_api, num_blocks, ph, wallet_node, wallet_rpc_api)
res = await data_rpc_api.create_data_store()
tree_id_1 = bytes32(hexstr_to_bytes(res["id"]))
res = await data_rpc_api.create_data_store()
tree_id_2 = bytes32(hexstr_to_bytes(res["id"]))
key1 = b"a"
value1 = b"\x01\x02"
changelist: List[Dict[str, str]] = [{"action": "insert", "key": key1.hex(), "value": value1.hex()}]
key2 = b"b"
value2 = b"\x03\x02"
changelist.append({"action": "insert", "key": key2.hex(), "value": value2.hex()})
key3 = b"c"
value3 = b"\x04\x05"
changelist.append({"action": "insert", "key": key3.hex(), "value": value3.hex()})
res = await data_rpc_api.batch_update({"id": tree_id_1.hex(), "changelist": changelist})
update_tx_rec0: TransactionRecord = res["tx_id"]
roots = await data_rpc_api.get_roots({"ids": [tree_id_1.hex(), tree_id_2.hex()]})
print(f"roots {roots}")
key4 = b"d"
value4 = b"\x06\x03"
changelist.append({"action": "insert", "key": key4.hex(), "value": value4.hex()})
key5 = b"e"
value5 = b"\x07\x01"
changelist.append({"action": "insert", "key": key5.hex(), "value": value5.hex()})
res = await data_rpc_api.batch_update({"id": tree_id_2.hex(), "changelist": changelist})
update_tx_rec1: TransactionRecord = res["tx_id"]
roots = await data_rpc_api.get_roots({"ids": [tree_id_1.hex(), tree_id_2.hex()]})
print(f"roots {roots}")
await asyncio.sleep(1)
for i in range(0, num_blocks):
await full_node_api.farm_new_transaction_block(FarmNewBlockProtocol(ph))
await time_out_assert(
15, is_transaction_confirmed, True, update_tx_rec1.wallet_id, wallet_rpc_api, update_tx_rec1.name
)
+2 -2
View File
@@ -241,14 +241,14 @@ async def test_get_pairs(
) -> None:
example = await create_example(data_store, tree_id)
pairs = await data_store.get_pairs(tree_id=tree_id)
pairs = await data_store.get_keys_values(tree_id=tree_id)
assert [node.hash for node in pairs] == example.terminal_nodes
@pytest.mark.asyncio
async def test_get_pairs_when_empty(data_store: DataStore, tree_id: bytes32) -> None:
pairs = await data_store.get_pairs(tree_id=tree_id)
pairs = await data_store.get_keys_values(tree_id=tree_id)
assert pairs == []