2019-11-17 12:01:10 +00:00
|
|
|
package main
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
|
|
|
|
"github.com/ipfs/go-cid"
|
|
|
|
|
|
|
|
aapi "github.com/filecoin-project/lotus/api"
|
|
|
|
"github.com/filecoin-project/lotus/chain/types"
|
|
|
|
)
|
|
|
|
|
|
|
|
func subMpool(ctx context.Context, api aapi.FullNode, st *storage) {
|
|
|
|
sub, err := api.MpoolSub(ctx)
|
|
|
|
if err != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
for change := range sub {
|
|
|
|
if change.Type != aapi.MpoolAdd {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
|
2019-11-19 19:53:24 +00:00
|
|
|
log.Info("mpool message")
|
|
|
|
|
2019-11-17 12:01:10 +00:00
|
|
|
err := st.storeMessages(map[cid.Cid]*types.Message{
|
|
|
|
change.Message.Message.Cid(): &change.Message.Message,
|
|
|
|
})
|
|
|
|
if err != nil {
|
2019-12-11 22:17:44 +00:00
|
|
|
//log.Error(err)
|
2019-11-17 12:01:10 +00:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
|
|
|
|
if err := st.storeMpoolInclusion(change.Message.Message.Cid()); err != nil {
|
|
|
|
log.Error(err)
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|