lotus/chain/events/state/predicates.go

263 lines
8.3 KiB
Go
Raw Normal View History

2020-06-22 21:38:40 +00:00
package state
import (
"context"
"github.com/ipfs/go-cid"
cbor "github.com/ipfs/go-ipld-cbor"
2020-06-22 21:38:40 +00:00
"github.com/filecoin-project/go-address"
"github.com/filecoin-project/go-amt-ipld/v2"
2020-06-22 21:38:40 +00:00
"github.com/filecoin-project/specs-actors/actors/abi"
"github.com/filecoin-project/specs-actors/actors/builtin"
"github.com/filecoin-project/specs-actors/actors/builtin/market"
"github.com/filecoin-project/specs-actors/actors/builtin/miner"
"github.com/filecoin-project/specs-actors/actors/util/adt"
"github.com/filecoin-project/lotus/api/apibstore"
"github.com/filecoin-project/lotus/chain/types"
2020-06-22 21:38:40 +00:00
)
// UserData is the data returned from the DiffTipSetKeyFunc
2020-06-25 21:43:37 +00:00
type UserData interface{}
2020-06-26 19:36:48 +00:00
// ChainAPI abstracts out calls made by this class to external APIs
type ChainAPI interface {
2020-06-26 15:51:45 +00:00
apibstore.ChainIO
2020-06-25 21:43:37 +00:00
StateGetActor(ctx context.Context, actor address.Address, tsk types.TipSetKey) (*types.Actor, error)
}
2020-06-26 19:36:48 +00:00
// StatePredicates has common predicates for responding to state changes
2020-06-22 21:38:40 +00:00
type StatePredicates struct {
2020-06-26 19:36:48 +00:00
api ChainAPI
2020-06-25 21:43:37 +00:00
cst *cbor.BasicIpldStore
2020-06-22 21:38:40 +00:00
}
2020-06-26 19:36:48 +00:00
func NewStatePredicates(api ChainAPI) *StatePredicates {
2020-06-22 21:38:40 +00:00
return &StatePredicates{
2020-06-25 21:43:37 +00:00
api: api,
cst: cbor.NewCborStore(apibstore.NewAPIBlockstore(api)),
2020-06-22 21:38:40 +00:00
}
}
// DiffTipSetKeyFunc check if there's a change form oldState to newState, and returns
2020-06-26 18:59:23 +00:00
// - changed: was there a change
// - user: user-defined data representing the state change
// - err
type DiffTipSetKeyFunc func(ctx context.Context, oldState, newState types.TipSetKey) (changed bool, user UserData, err error)
2020-06-25 21:43:37 +00:00
type DiffActorStateFunc func(ctx context.Context, oldActorStateHead, newActorStateHead cid.Cid) (changed bool, user UserData, err error)
2020-06-22 21:38:40 +00:00
2020-06-26 19:36:48 +00:00
// OnActorStateChanged calls diffStateFunc when the state changes for the given actor
func (sp *StatePredicates) OnActorStateChanged(addr address.Address, diffStateFunc DiffActorStateFunc) DiffTipSetKeyFunc {
return func(ctx context.Context, oldState, newState types.TipSetKey) (changed bool, user UserData, err error) {
oldActor, err := sp.api.StateGetActor(ctx, addr, oldState)
2020-06-22 21:38:40 +00:00
if err != nil {
return false, nil, err
}
newActor, err := sp.api.StateGetActor(ctx, addr, newState)
2020-06-26 19:36:48 +00:00
if err != nil {
return false, nil, err
}
2020-06-22 21:38:40 +00:00
if oldActor.Head.Equals(newActor.Head) {
return false, nil, nil
}
return diffStateFunc(ctx, oldActor.Head, newActor.Head)
}
}
type DiffStorageMarketStateFunc func(ctx context.Context, oldState *market.State, newState *market.State) (changed bool, user UserData, err error)
2020-06-26 19:36:48 +00:00
// OnStorageMarketActorChanged calls diffStorageMarketState when the state changes for the market actor
func (sp *StatePredicates) OnStorageMarketActorChanged(diffStorageMarketState DiffStorageMarketStateFunc) DiffTipSetKeyFunc {
2020-06-22 21:38:40 +00:00
return sp.OnActorStateChanged(builtin.StorageMarketActorAddr, func(ctx context.Context, oldActorStateHead, newActorStateHead cid.Cid) (changed bool, user UserData, err error) {
var oldState market.State
2020-06-25 21:43:37 +00:00
if err := sp.cst.Get(ctx, oldActorStateHead, &oldState); err != nil {
2020-06-22 21:38:40 +00:00
return false, nil, err
}
var newState market.State
2020-06-25 21:43:37 +00:00
if err := sp.cst.Get(ctx, newActorStateHead, &newState); err != nil {
2020-06-22 21:38:40 +00:00
return false, nil, err
}
return diffStorageMarketState(ctx, &oldState, &newState)
})
}
type DiffDealStatesFunc func(ctx context.Context, oldDealStateRoot *amt.Root, newDealStateRoot *amt.Root) (changed bool, user UserData, err error)
2020-06-26 19:36:48 +00:00
// OnDealStateChanged calls diffDealStates when the market state changes
2020-06-22 21:38:40 +00:00
func (sp *StatePredicates) OnDealStateChanged(diffDealStates DiffDealStatesFunc) DiffStorageMarketStateFunc {
return func(ctx context.Context, oldState *market.State, newState *market.State) (changed bool, user UserData, err error) {
if oldState.States.Equals(newState.States) {
return false, nil, nil
}
2020-06-25 21:43:37 +00:00
oldRoot, err := amt.LoadAMT(ctx, sp.cst, oldState.States)
2020-06-22 21:38:40 +00:00
if err != nil {
return false, nil, err
}
2020-06-25 21:43:37 +00:00
newRoot, err := amt.LoadAMT(ctx, sp.cst, newState.States)
2020-06-22 21:38:40 +00:00
if err != nil {
return false, nil, err
}
2020-06-25 21:43:37 +00:00
2020-06-22 21:38:40 +00:00
return diffDealStates(ctx, oldRoot, newRoot)
}
}
2020-06-26 19:36:48 +00:00
// ChangedDeals is a set of changes to deal state
2020-06-26 15:51:45 +00:00
type ChangedDeals map[abi.DealID]DealStateChange
2020-06-26 19:36:48 +00:00
// DealStateChange is a change in deal state from -> to
2020-06-26 15:51:45 +00:00
type DealStateChange struct {
From market.DealState
2020-06-26 19:36:48 +00:00
To market.DealState
2020-06-26 15:51:45 +00:00
}
2020-06-26 19:36:48 +00:00
// DealStateChangedForIDs detects changes in the deal state AMT for the given deal IDs
2020-06-25 21:43:37 +00:00
func (sp *StatePredicates) DealStateChangedForIDs(dealIds []abi.DealID) DiffDealStatesFunc {
return func(ctx context.Context, oldDealStateRoot *amt.Root, newDealStateRoot *amt.Root) (changed bool, user UserData, err error) {
2020-06-26 15:51:45 +00:00
changedDeals := make(ChangedDeals)
2020-06-26 19:36:48 +00:00
for _, dealID := range dealIds {
2020-06-25 21:43:37 +00:00
var oldDeal, newDeal market.DealState
2020-06-26 19:36:48 +00:00
err := oldDealStateRoot.Get(ctx, uint64(dealID), &oldDeal)
2020-06-25 21:43:37 +00:00
if err != nil {
return false, nil, err
}
2020-06-26 19:36:48 +00:00
err = newDealStateRoot.Get(ctx, uint64(dealID), &newDeal)
2020-06-25 21:43:37 +00:00
if err != nil {
return false, nil, err
}
if oldDeal != newDeal {
2020-06-26 19:36:48 +00:00
changedDeals[dealID] = DealStateChange{oldDeal, newDeal}
2020-06-25 21:43:37 +00:00
}
2020-06-22 21:38:40 +00:00
}
2020-06-25 21:43:37 +00:00
if len(changedDeals) > 0 {
2020-06-26 15:51:45 +00:00
return true, changedDeals, nil
2020-06-22 21:38:40 +00:00
}
2020-06-25 21:43:37 +00:00
return false, nil, nil
2020-06-22 21:38:40 +00:00
}
}
type DiffMinerActorStateFunc func(ctx context.Context, oldState *miner.State, newState *miner.State) (changed bool, user UserData, err error)
func (sp *StatePredicates) OnMinerActorChange(minerAddr address.Address, diffMinerActorState DiffMinerActorStateFunc) DiffTipSetKeyFunc {
return sp.OnActorStateChanged(minerAddr, func(ctx context.Context, oldActorStateHead, newActorStateHead cid.Cid) (changed bool, user UserData, err error) {
var oldState miner.State
if err := sp.cst.Get(ctx, oldActorStateHead, &oldState); err != nil {
return false, nil, err
}
var newState miner.State
if err := sp.cst.Get(ctx, newActorStateHead, &newState); err != nil {
return false, nil, err
}
return diffMinerActorState(ctx, &oldState, &newState)
})
}
type MinerSectorChanges struct {
Added []miner.SectorOnChainInfo
Extended []SectorExtensions
Removed []miner.SectorOnChainInfo
}
type SectorExtensions struct {
From miner.SectorOnChainInfo
To miner.SectorOnChainInfo
}
func (sp *StatePredicates) OnMinerSectorChange() DiffMinerActorStateFunc {
return func(ctx context.Context, oldState, newState *miner.State) (changed bool, user UserData, err error) {
ctxStore := &contextStore{
ctx: context.TODO(),
cst: sp.cst,
}
sectorChanges := &MinerSectorChanges{
Added: []miner.SectorOnChainInfo{},
Extended: []SectorExtensions{},
Removed: []miner.SectorOnChainInfo{},
}
// no sector changes
if oldState.Sectors.Equals(newState.Sectors) {
return false, nil, nil
}
oldSectors, err := adt.AsArray(ctxStore, oldState.Sectors)
if err != nil {
return false, nil, err
}
newSectors, err := adt.AsArray(ctxStore, newState.Sectors)
if err != nil {
return false, nil, err
}
var osi miner.SectorOnChainInfo
// find all sectors that were extended or removed
if err := oldSectors.ForEach(&osi, func(i int64) error {
var nsi miner.SectorOnChainInfo
found, err := newSectors.Get(uint64(osi.Info.SectorNumber), &nsi)
if err != nil {
return err
}
if !found {
sectorChanges.Removed = append(sectorChanges.Removed, osi)
return nil
}
if nsi.Info.Expiration != osi.Info.Expiration {
sectorChanges.Extended = append(sectorChanges.Extended, SectorExtensions{
From: osi,
To: nsi,
})
}
// we don't update miners state filed with `newSectors.Root()` so this operation is safe.
if err := newSectors.Delete(uint64(osi.Info.SectorNumber)); err != nil {
return err
}
return nil
}); err != nil {
return false, nil, err
}
// all sectors that remain in newSectors are new
var nsi miner.SectorOnChainInfo
if err := newSectors.ForEach(&nsi, func(i int64) error {
sectorChanges.Added = append(sectorChanges.Added, nsi)
return nil
}); err != nil {
return false, nil, err
}
// nothing changed
if len(sectorChanges.Added)+len(sectorChanges.Extended)+len(sectorChanges.Removed) == 0 {
return false, nil, nil
}
return true, sectorChanges, nil
}
}
type contextStore struct {
ctx context.Context
cst *cbor.BasicIpldStore
}
func (cs *contextStore) Context() context.Context {
return cs.ctx
}
func (cs *contextStore) Get(ctx context.Context, c cid.Cid, out interface{}) error {
return cs.cst.Get(ctx, c, out)
}
func (cs *contextStore) Put(ctx context.Context, v interface{}) (cid.Cid, error) {
return cs.cst.Put(ctx, v)
}