当Client接收到用户的存储交易,创建一个/fil/storage/mk/1.0.1协议的流,再通过流发送存储交易,去处理这个协议的正是HandleDealStream方法。这个方法直接调用自身的receiveDeal方法进行处理。receiveDeal方法处理如下:
从流中读取存储提案Proposal对象。
proposal,err:=s.ReadDealProposal()这里的流对象是dealStream对象(storagemarket/network/deal_stream.go),这个对象对原始流对象进行了封装。
获取 ipld node 对象。
proposalNd,err:=cborutil.AsIpld(proposal.DealProposal)生成矿工交易对象。
deal:=&storagemarket.MinerDeal{ Client:s.RemotePeer(), Miner:p.net.ID(), ClientDealProposal:*proposal.DealProposal, ProposalCid:proposalNd.Cid(), State:storagemarket.StorageDealUnknown, Ref:proposal.Piece, }调用 fsm 状态组的Begin的方法,生成一个状态机,并开始跟踪矿工交易对象。
err=p.deals.Begin(proposalNd.Cid(),deal)保存流对象到连接管理器中。
err=p.conns.AddStream(proposalNd.Cid(),s)发送事件到 fsm 状态组,从而开始对交易对象进行处理。
returnp.deals.Send(proposalNd.Cid(),storagemarket.ProviderEventOpen)当处理机收到ProviderEventOpen状态事件时,因为初始状态为默认值 0,即StorageDealUnknown,事件处理器对象经过内部处理找到对应的目的状态为StorageDealValidating,从而调用其处理函数ValidateDealProposal函数进行处理。
1、`ValidateDealProposal` 函数这个函数用来验证交易提案对象。
调用 Lotus Provider 适配器对象的GetChainHead方法,获取区块链顶部 tipset key 和其高度。
tok,height,err:=environment.Node().GetChainHead(ctx.Context())iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("nodeerrorgettingmostrecentstateid:%w",err)) }
验证客户发送的交易提案对象。如果验证不通过,则发送拒绝事件。
iferr:=providerutils.VerifyProposal(ctx.Context(),deal.ClientDealProposal,tok,environment.Node().VerifySignature);err!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("verifyingStorageDealProposal:%w",err)) }检查交易提案中指定的矿工地址是否正确。如果不正确,则发送拒绝事件。
proposal:=deal.Proposalifproposal.Provider!=environment.Address(){ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("incorrectproviderfordeal")) }
检查交易指定的高度是否正确。如果不正确,则发送拒绝事件。
ifheight>proposal.StartEpoch-environment.DealAcceptanceBuffer(){ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("dealstartepochistoosoonordealalreadyexpired")) }检查费用是否OK,如果不OK,则发送拒绝事件。
minPrice:=big.Div(big.Mul(environment.Ask().Price,abi.NewTokenAmount(int64(proposal.PieceSize))),abi.NewTokenAmount(1<<30)) ifproposal.StoragePricePerEpoch.LessThan(minPrice){ returnctx.Trigger(storagemarket.ProviderEventDealRejected, xerrors.Errorf("storagepriceperepochlessthanaskingprice:%s<%s",proposal.StoragePricePerEpoch,minPrice)) }检查交易的大小是否匹配。如果不匹配,则发送拒绝事件。
ifproposal.PieceSize<environment.Ask().MinPieceSize{ returnctx.Trigger(storagemarket.ProviderEventDealRejected, xerrors.Errorf("piecesizelessthanminimumrequiredsize:%d<%d",proposal.PieceSize,environment.Ask().MinPieceSize)) }ifproposal.PieceSize>environment.Ask().MaxPieceSize{ returnctx.Trigger(storagemarket.ProviderEventDealRejected, xerrors.Errorf("piecesizemorethanmaximumallowedsize:%d>%d",proposal.PieceSize,environment.Ask().MaxPieceSize)) }
获取客户的资金。
clientMarketBalance,err:=environment.Node().GetBalance(ctx.Context(),proposal.Client,tok) iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("nodeerrorgettingclientmarketbalancefailed:%w",err)) }如果客户可用资金小于总的交易费用,则发送拒绝事件。
ifclientMarketBalance.Available.LessThan(proposal.TotalStorageFee()){ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.New("clientMarketBalance.Availabletoosmall")) }如果交易是验证过的,则进行验证。
fsm 上下文对象的Trigger方法,发送事件。
returnctx.Trigger(storagemarket.ProviderEventDealDeciding)当状态机收到这个事件后,经过事件处理器把状态从StorageDealUnknown修改为StorageDealAcceptWait,从而调用其处理函数DecideOnProposal确定是否接收交易。
2、`DecideOnProposal` 函数这个函数用来决定接受或拒绝交易。
调用环境对象的RunCustomDecisionLogic方法,运行自定义逻辑来验证是不接收客户交易。
accept,reason,err:=environment.RunCustomDecisionLogic(ctx.Context(),deal)iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("customdealdecisionlogicfailed:%w",err)) }
如果不接收,则发送拒绝事件。
if!accept{ returnctx.Trigger(storagemarket.ProviderEventDealRejected,fmt.Errorf(reason)) }调用环境对象的SendSignedResponse方法,发送签名的响应给客户端。
err=environment.SendSignedResponse(ctx.Context(),&network.Response{ State:storagemarket.StorageDealWaitingForData, Proposal:deal.ProposalCid, })iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventSendResponseFailed,err) }
这个方法找到对应的流,然后对响应进行签名,生成签名的响应对象,最后通过流发送响应。
断开与客户端的连接。
iferr:=environment.Disconnect(deal.ProposalCid);err!=nil{ log.Warnf("closingclientconnection:%+v",err) }调用 fsm 上下文对象的Trigger方法,发送一个事件。
returnctx.Trigger(storagemarket.ProviderEventDataRequested)当状态机收到这个事件后,经过事件处理器把状态从StorageDealAcceptWait修改为StorageDealWaitingForData,因为没有指定的处理函数,从而不会调用函数进行处理,一直等待数据传输过程发送事件。
当数据开始传输时,数据传输组件发送ProviderEventDataTransferInitiated事件,经过事件处理器把状态从StorageDealWaitingForData修改为StorageDealTransferring,因为没有指定的处理函数,从而不会调用函数进行处理,一直等待数据传输过程发送事件。
当数据传输完成时,数据传输组件发送ProviderEventDataTransferCompleted事件,经过事件处理器把状态从StorageDealTransferring修改为StorageDealVerifyData,从而调用其处理函数VerifyData验证数据。
3、`VerifyData` 函数这个函数验证接受到的数据与交易提案中的 pieceCID 相匹配。
VerifyData函数流程如下:
调用环境对象的GeneratePieceCommitmentToFile方法,生成碎片的 CID 、碎片所在目录和元数据目录。
pieceCid,piecePath,metadataPath,err:=environment.GeneratePieceCommitmentToFile(deal.Ref.Root,shared.AllSelector())GeneratePieceCommitmentToFile方法内容如下:
调用文件存储对象的CreateTemp方法,创建一个临时文件。
f,err:=pio.store.CreateTemp()生成一个清理函数。
cleanup:=func(){ f.Close() _=pio.store.Delete(f.Path()) }从底层存储对象中获取指定 CID 的内容,然后写入指定文件。
err=pio.carIO.WriteCar(context.Background(),pio.bs,payloadCid,selector,f,userOnNewCarBlocks...)获取文件大小,即碎片大小。
pieceSize:=uint64(f.Size())定位到文件开头位置。
_,err=f.Seek(0,io.SeekStart)使用文件内容生成碎片 ID。
commitment,paddedSize,err:=GeneratePieceCommitment(rt,f,pieceSize)关闭文件。
_=f.Close()返回碎片 CID 和文件路径。
returncommitment,f.Path(),paddedSize,nil如果矿工设置了universalRetrievalEnabled标志,则直接调用GeneratePieceCommitmentWithMetadata函数进行处理。
ifp.p.universalRetrievalEnabled{ returnproviderutils.GeneratePieceCommitmentWithMetadata(p.p.fs,p.p.pio.GeneratePieceCommitmentToFile,p.p.proofType,payloadCid,selector) }universalRetrievalEnabled标志如果为真,则存储矿工会跟踪碎片中的所有 CID,因此对于所有 CID 都可以被检索,而不仅是 Root CID。
否则,调用 piece IO 对象的GeneratePieceCommitmentToFile方法进行处理。
pieceCid,piecePath,_,err:=p.p.pio.GeneratePieceCommitmentToFile(p.p.proofType,payloadCid,selector)payloadCid表示根 Root CID。
piece IO 对象的GeneratePieceCommitmentToFile方法处理如下:
返回碎片 CID 和碎片路径。
returnpieceCid,piecePath,filestore.Path(""),err验证生成的碎片 CID 和矿工交易中交易提案的碎片 CID是否一致。如果不一致,则发送拒绝事件。
ifpieceCid!=deal.Proposal.PieceCID{ returnctx.Trigger(storagemarket.ProviderEventDealRejected,xerrors.Errorf("proposalCommPdoesn'tmatchcalculatedCommP")) }3. 调用 fsm 上下文对象的Trigger方法,发送一个事件。returnctx.Trigger(storagemarket.ProviderEventVerifiedData,piecePath,metadataPath)当状态机收到这个事件后,经过事件处理器把状态从`StorageDealVerifyData`修改为`StorageDealEnsureProviderFunds`,从而调用其处理函数`EnsureProviderFunds`确定是否接收交易。同时,在调用处理函数之前,通过`Action`函数,修改矿工交易对象的`PiecePath`和`MetadataPath`两个属性。
4、`EnsureProviderFunds` 函数这个函数用来确定矿工有足够的资金来处理当前交易。
获取 Lotus Provider 适配器。
node:=environment.Node()获取区块链顶部 tipset 对应的 key 和高度。
tok,_,err:=node.GetChainHead(ctx.Context())iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventNodeErrored,xerrors.Errorf("acquiringchainhead:%w",err)) }
获取矿工的 worker 地址。
waddr,err:=node.GetMinerWorkerAddress(ctx.Context(),deal.Proposal.Provider,tok)iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventNodeErrored,xerrors.Errorf("lookingupminerworker:%w",err)) }
调用 Lotus Provider 适配器的EnsureFunds方法,确保矿工有足够的资金来处理当前交易。
mcid,err:=node.EnsureFunds(ctx.Context(),deal.Proposal.Provider,waddr,deal.Proposal.ProviderCollateral,tok)iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventNodeErrored,xerrors.Errorf("ensuringfunds:%w",err)) }
如果返回的mcid是空的,那么意味着已经实时确认,则调用 fsm 上下文对象的Trigger方法,发送一个事件。
ifmcid==cid.Undef{ returnctx.Trigger(storagemarket.ProviderEventFunded) }否则,调用 fsm 上下文对象的Trigger方法,发送另一个事件。
returnctx.Trigger(storagemarket.ProviderEventFundingInitiated,mcid)当状态机收到这个事件后,经过事件处理器把状态从StorageDealEnsureProviderFunds修改为StorageDealProviderFunding,从而调用其处理函数WaitForFunding等待产一步的消息上链。同时,在调用处理函数之前,通过Action函数,修改矿工交易对象的PublishCid属性。
5、`WaitForFunding` 函数这个函数用来等待消息上链。消息上链之后,调用 fsm 上下文对象的Trigger方法,发送一个事件。
函数内容如下:
node:=environment.Node()returnnode.WaitForMessage(ctx.Context(),*deal.AddFundsCid,func(codeexitcode.ExitCode,bytes[]byte,errerror)error{ iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventNodeErrored,xerrors.Errorf("AddFundserrored:%w",err)) } ifcode!=exitcode.Ok{ returnctx.Trigger(storagemarket.ProviderEventNodeErrored,xerrors.Errorf("AddFundsexitcode:%s",code.String())) } returnctx.Trigger(storagemarket.ProviderEventFunded) })
当状态机收到ProviderEventFunded这个事件后,经过事件处理器把状态从StorageDealProviderFunding修改为StorageDealPublish,从而调用其处理函数PublishDeal把交易信息上链。同时,在调用处理函数之前,通过Action函数,修改矿工交易对象的PublishCid属性。
6、`PublishDeal` 函数这个函数主要用来提交交易信息上链。
生成矿工交易对象。
smDeal:=storagemarket.MinerDeal{ Client:deal.Client, ClientDealProposal:deal.ClientDealProposal, ProposalCid:deal.ProposalCid, State:deal.State, Ref:deal.Ref, }调用 Lotus Provider 适配器对象的PublishDeals把交易信息上链。
mcid,err:=environment.Node().PublishDeals(ctx.Context(),smDeal) iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventNodeErrored,xerrors.Errorf("publishingdeal:%w",err)) }调用 fsm 上下文对象的Trigger方法,发送事件。
returnctx.Trigger(storagemarket.ProviderEventDealPublishInitiated,mcid)当状态机收到这个事件后,经过事件处理器把状态从StorageDealPublish修改为StorageDealPublishing,从而调用其处理函数WaitForPublish等待交易信息上链。
7、`WaitForPublish` 函数这个函数用来等待交易信息上链,然后给客户端发送响应,然后断开与客户端的连接。最后调用 fsm 上下文对象的Trigger方法,通过事件处理生成一个事件对象,然后发送事件对象到状态机。此处生成的事件对象名称为ProviderEventDealPublished。
当状态机收到这个事件后,经过事件处理器把状态从StorageDealPublishing修改为StorageDealStaged,从而调用其处理函数HandoffDeal开始扇区密封处理。同时,在调用处理函数之前,通过Action函数,修改矿工交易对象的ConnectionClosed和DealID属性。
returnenvironment.Node().WaitForMessage(ctx.Context(),*deal.PublishCid,func(codeexitcode.ExitCode,retBytes[]byte,errerror)error{ iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealPublishError,xerrors.Errorf("PublishStorageDealserrored:%w",err)) } ifcode!=exitcode.Ok{ returnctx.Trigger(storagemarket.ProviderEventDealPublishError,xerrors.Errorf("PublishStorageDealsexitcode:%s",code.String())) } varretvalmarket.PublishStorageDealsReturn err=retval.UnmarshalCBOR(bytes.NewReader(retBytes)) iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealPublishError,xerrors.Errorf("PublishStorageDealserrorunmarshallingresult:%w",err)) }returnctx.Trigger(storagemarket.ProviderEventDealPublished,retval.IDs[0]) })
8、`HandoffDeal` 函数这个函数调用 miner 的Provide适配器的
使用碎片路径生成文件对象。
file,err:=environment.FileStore().Open(deal.PiecePath)iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventFileStoreErrored,xerrors.Errorf("readingpieceatpath%s:%w",deal.PiecePath,err)) }
使用碎片文件流生成碎片流。
paddedReader,paddedSize:=padreader.New(file,uint64(file.Size()))调用 Lotus Provider 适配器对象的OnDealComplete方法,通知交易已经完成,从而把碎片加入某个扇区中。
err=environment.Node().OnDealComplete( ctx.Context(), storagemarket.MinerDeal{ Client:deal.Client, ClientDealProposal:deal.ClientDealProposal, ProposalCid:deal.ProposalCid, State:deal.State, Ref:deal.Ref, DealID:deal.DealID, FastRetrieval:deal.FastRetrieval, PiecePath:filestore.Path(environment.FileStore().Filename(deal.PiecePath)), }, paddedSize, paddedReader, )iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealHandoffFailed,err) }
调用 fsm 上下文对象的Trigger方法,发送事件。
returnctx.Trigger(storagemarket.ProviderEventDealHandedOff)当状态机收到这个事件后,经过事件处理器把状态从StorageDealStaged修改为StorageDealSealing,从而调用其处理函数VerifyDealActivated等待扇区密封结果。
9、`VerifyDealActivated` 函数生成回调函数。
cb:=func(errerror){ iferr!=nil{ _=ctx.Trigger(storagemarket.ProviderEventDealActivationFailed,err) }else{ _=ctx.Trigger(storagemarket.ProviderEventDealActivated) } }当 Lotus Provider 适配器对象检查到交易对象变化时会调用这个回调函数,从而发送相应的事件。
当状态机收到这个事件后,经过事件处理器把状态从StorageDealSealing修改为StorageDealActive,从而调用其处理函数RecordPieceInfo记录相关信息。
调用 Lotus Provider 适配器对象的OnDealSectorCommitted方法,等待扇区被提交。
err:=environment.Node().OnDealSectorCommitted(ctx.Context(),deal.Proposal.Provider,deal.DealID,cb)iferr!=nil{ returnctx.Trigger(storagemarket.ProviderEventDealActivationFailed,err) }
返回空。
returnnil9、`RecordPieceInfo` 函数这个函数主要记录相关信息。
最后调用 fsm 上下文对象的Trigger方法,通过事件处理生成一个事件对象,然后发送事件对象到状态机。此处生成的事件对象名称为ProviderEventDealCompleted。
当状态机收到这个事件后,经过事件处理器把状态从StorageDealActive修改为StorageDealCompleted,最终结束状态机处理。
这里会删除碎片的临时文件。
Filecoin宣布推出FVM漏洞賞金計劃
Filecoin宣布Network Indexer已正式發布
Filecoin推出存儲提供商指導資助計劃
Filecoin獨立協議Lotus發布v1.15.0
Filecoin網絡今晚將進行V15 OhSnap網絡升級
CoinList宣布推出低抵押Filecoin借貸計劃
Filecoin 區塊鏈瀏覽器 Filfox 推出多鏈加密貨幣錢包 Fox Wallet,支持鏈上數據查看以及礦工信息查詢等功能
三分鐘了解去中心化失能開關協議 Sarcophagus
Filecoin 將分階段推出 EVM 兼容的 Filecoin 虛擬機,首階段版本將于 2021 年第 4 季度上線
以跨鏈漫游、時間分片為地基,Chainge 「自金融」基礎設施服務商的愿景如何實現?
三分鐘了解 Numbers Protocol:Web3 去中心化圖片網絡
Protocol Labs 產品負責人:存儲在 Filecoin 上的 NFT 數據量達 54TB,數量超 700 萬
讀懂 Arweave 如何利用博弈設計實現永久網絡存儲
2021 萬向區塊鏈黑客馬拉松收官,一覽獲獎的 14 個項目
Vitalik Buterin、Juan Benet 與 Dominic Williams 探討公鏈技術創新之路
加拿大科技公司 ChainSafe 推出基于 Rust 語言的 Filecoin 客戶端 Forest
去中心化圖片網絡 Numbers Protocol 完成 600 萬美元私募輪和種子輪融資,Protocol Labs、幣安等參投
火星云礦將在 12 月 31 日前清退中國大陸挖礦資產和服務
CoinList 最新一期種子選手速覽,DeFi 和 NFT 仍占主導地位
專訪 Stratos 創始人:如何為區塊行業構建去中心化數據基礎設施?



