Near is a sharded blockchain system that hides its sharded nature from the end users, providing a single monolithic state abstraction for transaction execution. Currently, however, NEAR requires each step of inter-shard communication to be explicitly agreed upon by the underlying consensus algorithm. We argue that this is unnecessary.

Instead, we introduce SPICE, a new architecture of the NEAR that uses determinism inherent to a blockchain system to derive the transaction execution schedule, including inter-shard communication. This allows us to decouple transaction execution from ordering, relaxing the constraints imposed by the latter on the former and allowing for higher performance and more advanced transactions. Based on the same principles, we decouple the data availability of transaction payloads from ordering, reducing the chain itself to a lightweight component merely orchestrating other loosely coupled system components. This allows for short block times while maintaining high transaction throughput.

Introduction

A blockchain, at a high level, is a state machine replication system. It maintains its application state replicated across multiple nodes. Clients submit transactions (sometimes also called operations or requests) to the system which applies these transactions to the replicated application state while keeping all the copies of the state consistent. To this end, the blockchain system generally orders the transactions and applies (i.e., executes) them in the same order across all the copies of the state.

Especially if the system is large (consists of many nodes) and should sustain high transaction throughput (process many transactions per second), simply replicating all the required processing (ordering, execution, storage) across all nodes quickly becomes a scalability bottleneck.

Sharding is, at least conceptually, one of the most straightforward methods of improving scalability of blockchain systems and can be applied along several dimensions. A common way of applying sharding is to shard state storage. The blockchain state is partitioned into (usually disjoint) subsets called shards stored by different (usually also disjoint) sets of nodes. If execution is also sharded, different sets of nodes execute different sets of transactions. Some systems also employ the principles of sharding to ordering (sometimes also called sequencing) of transactions, introducing multiple more or less independent transaction sequences.

It is easy to see how the parallelism introduced by sharding can improve overall system throughput. However, increased throughput does not come for free. Sharding introduces significant challenges when it comes to coordinating the different shards. And without coordination, a fully sharded system would degrade to a set of multiple independent non-sharded systems.

Sharding in NEAR

NEAR is a blockchain system that applies all three of the above-mentioned dimensions of sharding: state storage, execution, and ordering. It partitions the state into subsets called shards, where each shard is only stored on a subset of the nodes. If a node $n$ stores the state of shard $s$, we say that $n$ is tracking $s$ or that $n$ belongs to $s$.

NEAR’s execution model allows for decomposing the execution of each transaction into parts we call receipts. A receipt is a unit of execution with the property of only accessing the state of a single shard. Each node only executes receipts accessing the shard it belongs to. If the execution of a receipt requires accessing state on a different shard, NEAR creates a new receipt describing the remote execution and sends it to the nodes tracking the other shard. When those nodes perform the necessary execution, they may respond with the execution result (also in the form of a receipt). Receipts can thus also be seen as messages exchanged by nodes from different shards used for composing the execution.

The ordering of NEAR transactions is also sharded to some extent. While NEAR produces only one chain of blocks, each block contains multiple independently ordered sequences (called chunks) of transactions, one chunk per shard. Each transaction targets a single shard where its execution starts. The execution may result in the creation of receipts which are included in later blocks alongside new transactions and executed in their respective destination shards.

Note that, given the above, the execution of a single transaction can span across multiple consecutive blocks and interleave with the execution of other transactions, making the transaction execution not linearizable. This is a price we pay for NEAR’s sharded nature.

The Problem

The main problem of the current design of NEAR is the interleaving of and the resulting strong coupling between ordering, execution, and data availability. The most important dependencies within the NEAR system are the following:

  • The system (like many others) requires the result of transaction execution to be known before a transaction is included in a block.

  • Receipts created during the execution of one block must be included in the next block.

  • A block can be created only after the previous one has been not just executed, but also verified (by means of multiple nodes re-executing its transactions).

  • A block does not directly contain the transaction payloads of its chunks. Instead, the block only has references (commitments) to that data. Therefore, for a new block to be created, the data it refers to must have been made available by replicating it on multiple nodes.

These tight dependencies lead to an alternation between data availability, block production, execution, and verification. A new block is only produced if the contained transactions are provably available and the transactions and receipts in its predecessor have been executed and verified. Conversely, execution and verification of the new transactions can only happen when the new block has been produced. This can lead to stalling block production while executing and verifying transactions. Conversely, execution and verification may stall while a new block is being agreed upon. While pipelining can mitigate such effects to some extent, its impact is limited.

There is also variability in the overhead (computation time) of executing transactions assigned to different blocks, as well as in block production latency. Forcing execution and verification to happen lock-step with block production makes the system always progress at the pace of the momentarily slower one of these two processes even if they are pipelined.

The Solution: Decoupling Ordering, Execution, and Availability

The central concept of SPICE is having a very slim and simple mechanism for producing an agreed-upon sequence of blocks, providing just a very basic way of ordering events in the system. We aim for small blocks with little data inside, minimizing the amount of computation and communication necessary for producing the blocks. This enables us to produce blocks at a high rate and promotes decentralization by allowing the participation of nodes with limited hardware resources.

Our goal is moving most of the functionality required to run the blockchain system off the critical path of block production. This includes storing data, executing transactions, updating the application state, verifying the computation and certifying the execution results. The blocks themselves become merely a substrate that the other system components use for synchronization and committing to the outcomes they produce.

Decoupling the tasks of data dissemination, ordering, and execution allows them to run truly concurrently. Some inevitable dependence between these three tasks still remains, since all are necessary to advance the system’s state. However, the coupling becomes much looser, enabling the system to progress at the speed of the one that is slower on average, rather than the momentarily slower one, as has been the case so far.

Our design uses clean abstractions between the system components. Such a design, being desirable by itself from an engineering perspective, also makes NEAR future-proof. Advancements in the state of the art can more easily be integrated into our system, only modifying the concerned component. For example, we can replace the total-order broadcast (“consensus”) protocol by a more recent one if / when needed; we only need to modify the ordering component, without touching the rest of the system.

Overview

We now present the high-level structure of the system and ideas underlying its design. At the base of our system lies the Chain component that produces a sequence of blocks. Those blocks, however, do not contain user transactions at all, as is the case in many other blockchain systems. Instead, blocks only include bits of data that we call statements produced by various nodes participating in the protocol.

Statements can be seen as special transactions that only participating nodes (with adequate stake) can submit. The agreed-upon sequence of those statements is what defines NEAR’s core state. The core state only holds metadata about (references and commitments to) user transactions and the application state1, which are themselves partitioned into multiple shards. Statements are signed and, intuitively, express the participating nodes’ contributions to the protocol execution. For example, the meaning of a statement produced by a node can be “I verified that the outcome of executing user transactions X on top of the state with root Y leads to a state with root Z” or “I confirm having persistently stored the transaction data for shard X referenced in block Y”.

Different sets of nodes fulfill different functions of the NEAR system. All nodes locally replicate the core state and update it by submitting their statements. Note that the core state updates do not immediately happen locally, as all submitted statements need to be agreed upon before each node uses them to update their local copy of the core state. The Chain thus serves as a synchronization mechanism for different parts of the NEAR system.

Apart from the Chain, the two other main parts of the system are the Data Store and the Execution, respectively concerned with data availability and execution. The Data Store stores the payloads of user transactions and, through statements submitted to the chain, certifies that those payloads have been replicated in a way guaranteeing retrievability despite failures. The Execution Component is responsible for holding the application state and advancing it based on the transactions referenced on the Chain. It downloads the necessary transaction payloads from the Data Store and executes them. A high-level interaction between the different components is depicted below.

Image: Higl-level overview of SPICE

In the following, we provide more details on each of these three parts of the NEAR blockchain system.

Chain

The task of the Chain is to produce a continuously growing chain of blocks containing statements. The chain starts with a pre-defined genesis block known to all nodes. Each block has an integer height that is related to the block’s distance from the genesis block. We call an interval of heights an epoch. Blocks with heights in this interval are said to belong to the corresponding epoch. The concept of epochs facilitates reconfiguration of the system. The configuration (such as the set of participating nodes) is static during an epoch and only changes at epoch boundaries based on the agreed-upon system state.

To agree on the final sequence of blocks, NEAR uses the Doomslug Byzantine total-order broadcast (“consensus”) algorithm. We make the deliberate choice of hiding its details, since the concrete algorithm can be abstracted away and even replaced in the future without impacting the design of the rest of the system. It is only important that the system can receive inputs, order them (by including them in blocks), and deterministically process them using arbitrary processing logic.

The input the Chain receives are statements made by the participating nodes. Unlike in many other blockchain systems which directly accept user transactions from clients, the chain only accepts statements from nodes that participate in the protocol. Participation in the protocol can, in turn, be achieved by locking an adequate number of NEAR tokens as stake.

Core State

The core state is a small amount of data replicated on all nodes. It contains the most important information about the system itself, such as the set of nodes and their roles, the assignment of stake, or information about what application state has been certified and what transaction payloads have been made available and in what order. The core state can be deterministically derived solely from the sequence of statements included in the blocks produced by the Chain component.

We design the statement processing logic to be as simple as possible. The statements are signed and concern, for example, having stored some data or having verified a part of the execution.

A specific collection of statements constitutes a certificate. For example, if enough validator nodes submit a statement endorsing some particular execution result, the set of those statements constitutes a state certificate. If all these statements are included in blocks of the chain, the certificate becomes part of the core state and anyone observing it will consider the corresponding state certified.

The notion of certificate can be generalized to any data proving that some statement is true (such as data being available or execution result being valid). What exactly constitutes a certificate can (and should) be abstracted away and defined separately to allow the system to evolve. A certificate may take different forms. Natural examples include:

  • Quorum certificate: A collection of signed statements from a quorum of nodes, as described above
  • Threshold certificate: A threshold signature produced by a quorum of nodes
  • Succinct proof: A succinct proof (“ZK proof”) of the statement

In our design, we use quorum certificates for their simplicity, but they can be substituted for a more suitable kind in the future.

Ordering Transactions

The Chain orders transactions by including references to them in blocks. Designated nodes called chunk producers receive transactions from users. They take turns in proposing sequences (called chunks) of those transactions for inclusion in blocks. For each block, the Chain designates one chunk producer per shard. The system computes this assignment deterministically (pseudo-randomly) based on the amount of NEAR tokens each chunk producer staked. The assignment is part of the core state and is known to all nodes that download the blocks and are thus able to reconstruct the core state.

When a chunk producer assembles a chunk, it saves it in the Data Store and submits a statement containing a commitment to (hash of) the chunk. The assignment of chunk producers (and thus also of the chunks they produce) to blocks defines the order of the included transactions.

Note that, in general, the system has more than one shard, meaning that each block is assigned more than one chunk. Transactions in different chunks of the same block are not ordered with respect to each other. The way NEAR operates makes those transactions always commute and ordering them is thus not needed. Nevertheless, if one wanted to impose a total order on all NEAR’s transactions, it would suffice to define any deterministic ordering of transactions within the block.

Data Store

The Data Store is a very simple component that consists of a set of data owner nodes making sure that transaction payloads referenced on the chain are available (i.e., retrievable). Upon reception of a blob of data, a data owner persistently stores it and submits a signed statement to the Chain, confirming that the data is stored. The Chain then assembles these statements into an availability certificate for the stored data. The Data Store then serves the data to other nodes that need it.

To prevent DoS attacks against data owners, we have them observe the core state which contains information about which nodes the protocol selected to be eligible to produce data. Data owners would not accept data from any other nodes. Even those nodes could be faulty, of course, and try to overwhelm data owners with maliciously crafted traffic. The amount of data the protocol requires to be stored from any single node is limited, however, making it possible for the data owners to easily detect fraudulent traffic.

Execution

The biggest change the SPICE architecture introduces concerns the execution. The main insight is that we move the execution of transactions off the critical path of block production. Instead, we implement a separate Execution component that is only loosely coupled to the chain through observing the core state and issuing statements. This allows us to produce blocks at a high rate, as a new block does not need to wait for the execution of the previous block to complete.

On a high level, the Execution Component observes the chain (in particular the core state), looking for commitments to chunks (lists of transactions) and corresponding availability certificates. When a commitment to an available chunk appears in the core state, Execution obtains the corresponding transaction payload from the Data Store, executes the transactions, and submits a proof that the transactions have been executed correctly (called a state certificate) to the Chain (by issuing corresponding statements).

The execution is performed by nodes called replicas. Replicas hold the application state and evaluate a state transition function. Each replica only holds the state of one shard and executes transactions and receipts associated with that shard. We explain how exactly this works a bit later. Each time a replica executes a chunk of transactions, it produces a state witness. This is a piece of data containing all the necessary information for anyone to verify that the transactions have been executed correctly. In practice, a state witness consists of all the inputs to the computation, including their Merkle inclusion proofs, along with the (hash of) the computation result.

State Certificates

NEAR periodically declares new versions of the blockchain state certified. We consider a particular version of the blockchain state certified if the core state contains a corresponding certificate. Similarly to availability certificates, state certificates are constructed and stored in the core state. While they can take any form in principle (such as a ZK proof), we use quorum certificates for simplicity (at least for now).

A quorum state certificate consists of signatures produced by nodes called validators. They are tasked with verifying the state witnesses created by replicas. Effectively, this means re-computing the new state, given a commitment to the old state and a state witness. Note that validators do not need to locally store the whole state of a shard. They only need a commitment to (Merkle root of) the shard state, which they obtain from the core state. We call this approach stateless validation.

If a validator obtains the same result as included in the state witness, it issues a statement confirming (i.e., endorsing) the execution result. Once enough validators perform the same execution and confirm a particular execution result, we consider the corresponding state certified.

Image: SPICE execution flow

Sharded Execution

Recall that, in NEAR, each user transaction is associated with a single shard. Only state within a shard can be modified atomically. A transaction can modify state residing in different shards, but NEAR executes these modifications asynchronously through sending transaction-like messages called receipts to other shards. Receipts, in turn, trigger the necessary state modifications within their respective shards.

The state of each shard advances in steps, separately from (but depending on) other shards. At each step, the execution component applies a state transition function to the shard state and a chunk (a list of transactions) to produce a new state for that shard. We refer to the application of the state transition function as execution.

In addition to transactions and the state, the execution also consumes (incoming) and produces (outgoing) receipts. When a receipt is produced at some execution step, it is consumed by its destination shard’s execution in the next step. In general, different nodes perform the execution for different shards. Receipts therefore need to be sent over the network from nodes that produce them in one step to nodes that consume them in the next step (unless source and destination shards are the same). NEAR employs a data dissemination subprotocol for this purpose that we will discuss in a future blog post. Receipts are the only means for shards to exchange information.

Lastly, the execution also depends on the core state, which may contain various parameters defining the details of how the execution must be performed (e.g., limits on how much execution should be performed in one step before postponing the rest to the next step).

To sum up, to produce a new version of the shard state, a replica needs 1) the core state, 2) the previous state of the shard, 3) a chunk of transactions, and 4) the incoming receipts from all other shards. The result of an execution step in a shard is 1) the new state of the shard, 2) outgoing receipts for other shards, and 3) a state witness for the validators who, after verifying it, produce statements certifying the execution result.

Image: SPICE state transition

Summary

In the current version of NEAR, ordering transactions, storing their payloads, and executing them is all tightly coupled and performed in lock-step. At any point in time, the system performance is limited by the momentarily slowest of these three processes. While pipelining helps address this issue to some extent, its impact is limited.

We created a design that decouples these three processes and hides them behind clean abstractions. The Chain is reduced to a simple ordering service for transactions and system-internal events, allowing it to produce blocks independently of the other components. The remaining two components (Data Store and Execution) only interact with the chain by submitting statements to it and observing the core state constructed from those statements. The Data Store stores transaction payloads and produces availability certificates, while the Execution downloads those transactions and advances the application state accordingly.

The decoupling of the Data Store and the Execution from the Chain allows them to run completely asynchronously at network speed. They may temporarily “fall behind” the Chain during load spikes, meaning that while multiple new commitments to transactions already appear in the core state, they may not yet have been made available or executed. After the load spike, the Data Store and the Execution can “catch up” by making available and executing all the outstanding transactions.

Note that the Chain is the only component that effectively decides how the replicated state will evolve, as that state can be derived from the sequence of transactions committed to in the blocks. The Execution component only computes what this state evaluates to, but has no power to influence the result. In other words, the state is well-defined even before anybody knows what it is.

This blog post only presented the principles upon which the new version of the NEAR system is based. In further posts, we will look deeper into the challenges this approach presents and how we address them.

Bonus: Seeing NEAR Shards as Rollups

The sharded nature of NEAR with SPICE shares many similarities with L2 blockchain systems that run on top of a single main L1 chain. In particular, we share many concepts with rollups on Ethereum, each NEAR shard corresponding to a separate Ethereum rollup (L2) and the NEAR blockchain corresponding to Ethereum itself (L1). NEAR chunks correspond to L2 blocks that are constructed independently and the final global order is only established when those are posted to the L1.

In Ethereum rollups, the rollup state is determined by the sequence of L2 blocks associated with the L1 chain. Participants then periodically include commitments to the rollup state on the L1 chain (in L1 smart contracts). These commitments to the rollup state are then certified either explicitly (by submitting a validity proof in ZK rollups) or implicitly (by the absence of a fraud proof within a pre-defined part of the L1 chain in optimistic rollups).

SPICE has the same basic working principle, but the “rollup functionality” is built-in. In NEAR, each block generally contains a chunk from each shard, which is not necessarily the case for Ethereum rollups. The NEAR chain, in fact, only serves for ordering chunks (~L2 blocks) and state commitments, without the possibility of directly hosting general-purpose user smart contracts. That is, unlike Ethereum, NEAR does not allow any client code execution or even token transfers on its chain - everything must happen at NEAR’s L2s, i.e., shards.

The main difference between Ethereum rollups and NEAR shards is that while Ethereum rollups mostly exist at the application level (i.e., a rollup smart contract is no different from any other smart contract from Ethereum’s point of view), NEAR shards are an integral part of the system design. Conveniently, this allows us to make the communication between different NEAR shards part of the system specification (Ethereum rollups must interact without direct support of the Ethereum system). While this leads to more efficiency in NEAR (data can flow between shards directly, with minimal involvement of the L1), it has a flip side too: The design of NEAR (in particular the necessity of all shard’s receipts for the computation of every other shard’s state transition) still involves a tight coupling between shards, preventing the state of different shards from evolving independently.

The above, however, is not a fundamental limitation and NEAR’s design may be adapted to relax these constraints. Moreover, one could further generalize the design and allow it to become more similar to rollups, e.g., by enabling heterogeneous, dynamically reconfigurable shards. We postpone these considerations to future work.


$^1$We use the term application state for the replicated state the NEAR blockchain system maintains on behalf of the users. It consists of user accounts and is updated by executing transactions. (back to text)