Both. First it does a local fold on local inputs, then it does a global fold across all of the local fold results.
i.e. Assume we have three nodes and one key, "our_key", and values for "our_key" are distributed across the 3 nodes:
node1 : our_key : [ 1, 2, 3, 4, 5]
node2 : our_key : [ 6, 7, 8, 9, 10]
node3 : out_key : [11, 12, 13, 14, 15]
and our fold function ( which must be associative) is +. First we fold each node:
node1 : fold(+, [ 1, 2, 3, 4, 5]) = our_key : [15]
node2 : fold(+, [ 6, 7, 8, 9, 10]) = our_key : [40]
node3 : fold(+, [11, 12, 13, 14, 15]) = our_key : [65]
When each node is done with its local inputs, the key-value pairs are redistributed so that each key and its values are on a single node. In this case, assume our hash function moved "our_key" to node 2. Node 2 would now have:
node2 : our_key : [15, 40, 65]
And we would fold against that and get:
node2: our_key : 120
This is how we achieve massive scalability. The nodes can do all of their processing in parallel without any intercommunication except for the key-redistribution phase.
Two key things to point out are that HashFold does the folding as the data becomes available (we don't wait for all values to start folding... this reduces memory usage because we don't need to store the list of values, only the result of the fold as it progresses), and the framework described here is sufficient to perform any computation that MapReduce can do which should make the transition easier (as a lot of major data processing tasks are currently framed for MapReduce)