|
|
|
@@ -19,6 +19,9 @@ import (
|
|
|
|
|
|
|
|
|
|
"github.com/google/uuid"
|
|
|
|
|
"github.com/ipfs/go-cid"
|
|
|
|
|
"github.com/ipfs/go-graphsync"
|
|
|
|
|
gsimpl "github.com/ipfs/go-graphsync/impl"
|
|
|
|
|
"github.com/ipfs/go-graphsync/peerstate"
|
|
|
|
|
"github.com/libp2p/go-libp2p-core/host"
|
|
|
|
|
"github.com/libp2p/go-libp2p-core/peer"
|
|
|
|
|
"go.uber.org/fx"
|
|
|
|
@@ -28,6 +31,7 @@ import (
|
|
|
|
|
|
|
|
|
|
"github.com/filecoin-project/go-address"
|
|
|
|
|
datatransfer "github.com/filecoin-project/go-data-transfer"
|
|
|
|
|
gst "github.com/filecoin-project/go-data-transfer/transport/graphsync"
|
|
|
|
|
"github.com/filecoin-project/go-fil-markets/piecestore"
|
|
|
|
|
"github.com/filecoin-project/go-fil-markets/retrievalmarket"
|
|
|
|
|
"github.com/filecoin-project/go-fil-markets/storagemarket"
|
|
|
|
@@ -70,6 +74,8 @@ type StorageMinerAPI struct {
|
|
|
|
|
RetrievalProvider retrievalmarket.RetrievalProvider `optional:"true"`
|
|
|
|
|
SectorAccessor retrievalmarket.SectorAccessor `optional:"true"`
|
|
|
|
|
DataTransfer dtypes.ProviderDataTransfer `optional:"true"`
|
|
|
|
|
StagingGraphsync dtypes.StagingGraphsync `optional:"true"`
|
|
|
|
|
Transport dtypes.ProviderTransport `optional:"true"`
|
|
|
|
|
DealPublisher *storageadapter.DealPublisher `optional:"true"`
|
|
|
|
|
SectorBlocks *sectorblocks.SectorBlocks `optional:"true"`
|
|
|
|
|
Host host.Host `optional:"true"`
|
|
|
|
@@ -553,6 +559,107 @@ func (sm *StorageMinerAPI) MarketDataTransferUpdates(ctx context.Context) (<-cha
|
|
|
|
|
return channels, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (sm *StorageMinerAPI) MarketDataTransferDiagnostics(ctx context.Context, mpid peer.ID) (*api.TransferDiagnostics, error) {
|
|
|
|
|
|
|
|
|
|
// gather information about active transport channels
|
|
|
|
|
transportChannels := sm.Transport.(*gst.Transport).ChannelsForPeer(mpid)
|
|
|
|
|
// gather information about graphsync state for peer
|
|
|
|
|
gsPeerState := sm.StagingGraphsync.(*gsimpl.GraphSync).PeerState(mpid)
|
|
|
|
|
|
|
|
|
|
sendingTransfers := sm.generateTransfers(ctx, transportChannels.SendingChannels, gsPeerState.IncomingState)
|
|
|
|
|
receivingTransfers := sm.generateTransfers(ctx, transportChannels.ReceivingChannels, gsPeerState.OutgoingState)
|
|
|
|
|
|
|
|
|
|
return &api.TransferDiagnostics{
|
|
|
|
|
SendingTransfers: sendingTransfers,
|
|
|
|
|
ReceivingTransfers: receivingTransfers,
|
|
|
|
|
}, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// generate transfers matches graphsync state and data transfer state for a given peer
|
|
|
|
|
// to produce detailed output on what's happening with a transfer
|
|
|
|
|
func (sm *StorageMinerAPI) generateTransfers(ctx context.Context,
|
|
|
|
|
transportChannels map[datatransfer.ChannelID]gst.ChannelGraphsyncRequests,
|
|
|
|
|
gsPeerState peerstate.PeerState) []*api.GraphSyncDataTransfer {
|
|
|
|
|
tc := &transferConverter{
|
|
|
|
|
matchedRequests: make(map[graphsync.RequestID]*api.GraphSyncDataTransfer),
|
|
|
|
|
gsDiagnostics: gsPeerState.Diagnostics(),
|
|
|
|
|
requestStates: gsPeerState.RequestStates,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// iterate through all operating data transfer transport channels
|
|
|
|
|
for channelID, channelRequests := range transportChannels {
|
|
|
|
|
originalState, err := sm.DataTransfer.ChannelState(ctx, channelID)
|
|
|
|
|
var baseDiagnostics []string
|
|
|
|
|
var channelState *api.DataTransferChannel
|
|
|
|
|
if err != nil {
|
|
|
|
|
baseDiagnostics = append(baseDiagnostics, fmt.Sprintf("Unable to lookup channel state: %s", err))
|
|
|
|
|
} else {
|
|
|
|
|
cs := api.NewDataTransferChannel(sm.Host.ID(), originalState)
|
|
|
|
|
channelState = &cs
|
|
|
|
|
}
|
|
|
|
|
// add the current request for this channel
|
|
|
|
|
tc.convertTransfer(&channelID, channelState, baseDiagnostics, channelRequests.Current, true)
|
|
|
|
|
for _, requestID := range channelRequests.Previous {
|
|
|
|
|
// add any previous requests that were cancelled for a restart
|
|
|
|
|
tc.convertTransfer(&channelID, channelState, baseDiagnostics, requestID, false)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// collect any graphsync data for channels we don't have any data transfer data for
|
|
|
|
|
tc.collectRemainingTransfers()
|
|
|
|
|
|
|
|
|
|
return tc.transfers
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type transferConverter struct {
|
|
|
|
|
matchedRequests map[graphsync.RequestID]*api.GraphSyncDataTransfer
|
|
|
|
|
transfers []*api.GraphSyncDataTransfer
|
|
|
|
|
gsDiagnostics map[graphsync.RequestID][]string
|
|
|
|
|
requestStates graphsync.RequestStates
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// convert transfer assembles transfer and diagnostic data for a given graphsync/data-transfer request
|
|
|
|
|
func (tc *transferConverter) convertTransfer(channelID *datatransfer.ChannelID, channelState *api.DataTransferChannel, baseDiagnostics []string,
|
|
|
|
|
requestID graphsync.RequestID, isCurrentChannelRequest bool) {
|
|
|
|
|
diagnostics := baseDiagnostics
|
|
|
|
|
state, hasState := tc.requestStates[requestID]
|
|
|
|
|
stateString := state.String()
|
|
|
|
|
if !hasState {
|
|
|
|
|
stateString = "no graphsync state found"
|
|
|
|
|
}
|
|
|
|
|
if channelID == nil {
|
|
|
|
|
diagnostics = append(diagnostics, fmt.Sprintf("No data transfer channel id for GraphSync request ID %d", requestID))
|
|
|
|
|
} else if isCurrentChannelRequest && !hasState {
|
|
|
|
|
diagnostics = append(diagnostics, fmt.Sprintf("No current request state for data transfer channel id %s", channelID))
|
|
|
|
|
} else if !isCurrentChannelRequest && hasState {
|
|
|
|
|
diagnostics = append(diagnostics, fmt.Sprintf("Graphsync request %d is a previous request on data transfer channel id %s that was restarted, but it is still running", requestID, channelID))
|
|
|
|
|
}
|
|
|
|
|
diagnostics = append(diagnostics, tc.gsDiagnostics[requestID]...)
|
|
|
|
|
transfer := &api.GraphSyncDataTransfer{
|
|
|
|
|
RequestID: requestID,
|
|
|
|
|
RequestState: stateString,
|
|
|
|
|
IsCurrentChannelRequest: isCurrentChannelRequest,
|
|
|
|
|
ChannelID: channelID,
|
|
|
|
|
ChannelState: channelState,
|
|
|
|
|
Diagnostics: diagnostics,
|
|
|
|
|
}
|
|
|
|
|
tc.transfers = append(tc.transfers, transfer)
|
|
|
|
|
tc.matchedRequests[requestID] = transfer
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (tc *transferConverter) collectRemainingTransfers() {
|
|
|
|
|
for requestID := range tc.requestStates {
|
|
|
|
|
if _, ok := tc.matchedRequests[requestID]; !ok {
|
|
|
|
|
tc.convertTransfer(nil, nil, nil, requestID, false)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
for requestID := range tc.gsDiagnostics {
|
|
|
|
|
if _, ok := tc.matchedRequests[requestID]; !ok {
|
|
|
|
|
tc.convertTransfer(nil, nil, nil, requestID, false)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (sm *StorageMinerAPI) MarketPendingDeals(ctx context.Context) (api.PendingDealInfo, error) {
|
|
|
|
|
return sm.DealPublisher.PendingDeals(), nil
|
|
|
|
|
}
|
|
|
|
|