You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Corda主流调用子流更新同一状态触发UnexpectedFlowEndException求助

问题:Corda主流调用子流更新同一LinearState时出现UnexpectedFlowEndException错误

我在Corda中编写了用于更新ContainerState的主流UpdateContainerStateFlow,尝试在该主流中调用另一个自定义子流UpdateContainerStateFlow2来更新同一ID的ContainerState属性(注:此为测试示例,实际可在主流内完成)。在构建主流交易前调用该子流后,出现如下错误:

错误日志

[WARN] 13:07:03,441 [Mock network] interceptors.DumpHistoryOnErrorInterceptor. - Flow [fb6a42fd-b9bf-45bb-ab98-2f2510235f91] error {fiber-id=10000008, flow-id=fb6a42fd-b9bf-45bb-ab98-2f2510235f91, invocation_id=38c29575-637a-459b-b8e9-81a960e93c7c, invocation_timestamp=2022-12-07T13:07:03.064Z, origin=O=Mock Company 1, L=London, C=GB, session_id=38c29575-637a-459b-b8e9-81a960e93c7c, session_timestamp=2022-12-07T13:07:03.064Z, thread-id=406, tx_id=EB1EA4EEA934C9C874461A6DDE35892057EBC472F405AF95241784CEDE3D6C92}

错误信息

net.corda.core.flows.UnexpectedFlowEndException: Counter-flow errored
at Received unexpected counter-flow exception from peer O=Mock Company 1, L=London, C=GB.() ~[?:?]

我疑惑子流是否应挂起主流并完成FinalityFlow的所有步骤,现附上相关代码,请求纠正理解误区并解决问题。

主流代码

@InitiatingFlow
@StartableByRPC
class UpdateContainerStateFlow(private val containerId: UUID,
val status: String, val comments: String) : FlowLogic<SignedTransaction>() {
    override val progressTracker = ProgressTracker()

    @Suspendable
    override fun call(): SignedTransaction {

        val queryCriteria = QueryCriteria.LinearStateQueryCriteria(linearId = listOf(UniqueIdentifier(id = containerId)))
        val currentContainer = serviceHub.vaultService.queryBy<ContainerState>(queryCriteria).states.single()

        val notary = serviceHub.networkMapCache.getNotary( CordaX500Name.parse("O=Notary,L=London,C=GB"))

        val updateContainer = currentContainer.state.data.copy(
           status = status,
           commentsFromStakeHolders = comments
        )

        val abc = subFlow(UpdateContainerStateFlow2(
            containerId = containerId,
            status = "subflowOne",
            comments = "test"
        ))

        val builder = TransactionBuilder(notary)
                .addCommand(ContainerContract.Commands.Update(), listOf(updateContainer.logisticCompany.owningKey, updateContainer.partnerCompany.owningKey))
                .addInputState(currentContainer)
                .addOutputState(updateContainer)

        // Step 4. Verify and sign it with our KeyPair.
        builder.verify(serviceHub)
        val ptx = serviceHub.signInitialTransaction(builder)

        val sessions = (updateContainer.participants - ourIdentity).map { initiateFlow(it) }.toSet()
        val stx = subFlow(CollectSignaturesFlow(ptx, sessions))

        // Step 7. Assuming no exceptions, we can now finalise the transaction
        return subFlow(FinalityFlow(stx, sessions))
    }
}

@InitiatedBy(UpdateContainerStateFlow::class)
class UpdateContainerStateFlowResponder(val counterpartySession: FlowSession) : FlowLogic<SignedTransaction>() {
    @Suspendable
    override fun call(): SignedTransaction {
        val signTransactionFlow = object : SignTransactionFlow(counterpartySession) {
            override fun checkTransaction(stx: SignedTransaction) = requireThat {
               //Addition checks
            }
        }
        val txId = subFlow(signTransactionFlow).id
        return subFlow(ReceiveFinalityFlow(counterpartySession, expectedTxId = txId))
    }
}

UpdateContainerStateFlow2()的实现与上述代码完全一致。


核心问题分析

  1. 子流角色定位错误:将UpdateContainerStateFlow2标记为@InitiatingFlow,作为子流调用时会尝试向参与者发起新会话,与主流已建立的会话冲突,触发异常。
  2. LinearState版本冲突:子流完成FinalityFlow后,目标State版本已升级,主流仍使用旧版本构建交易,会导致Vault状态校验失败。
  3. 子流设计误区:内部逻辑拆分的子流不应设计为完整的发起-响应式Flow,更不应单独提交FinalityFlow,否则会打断主流流程,造成节点状态不一致。

修正方案

方案一:将子流改为内部逻辑流(推荐)

如果只是拆分内部逻辑,将子流改为普通FlowLogic,去掉跨节点交互标记,仅负责状态预处理,不独立提交交易:

// 子流改为普通FlowLogic,仅处理状态更新逻辑
class UpdateContainerStateFlow2(private val containerId: UUID,
val status: String, val comments: String) : FlowLogic<ContainerState>() {
    @Suspendable
    override fun call(): ContainerState {
        val queryCriteria = QueryCriteria.LinearStateQueryCriteria(linearId = listOf(UniqueIdentifier(id = containerId)))
        val currentContainer = serviceHub.vaultService.queryBy<ContainerState>(queryCriteria).states.single()
        // 返回更新后的状态,由主流统一构建交易
        return currentContainer.state.data.copy(
            status = status,
            commentsFromStakeHolders = comments
        )
    }
}

主流中调整调用逻辑:

// 调用子流获取处理后的状态
val subUpdatedContainer = subFlow(UpdateContainerStateFlow2(
    containerId = containerId,
    status = "subflowOne",
    comments = "test"
))
// 使用子流返回的状态构建主流交易
val builder = TransactionBuilder(notary)
        .addCommand(ContainerContract.Commands.Update(), listOf(subUpdatedContainer.logisticCompany.owningKey, subUpdatedContainer.partnerCompany.owningKey))
        .addInputState(currentContainer)
        .addOutputState(subUpdatedContainer)

方案二:子流作为独立交易(仅特殊场景使用)

如果必须让子流提交独立交易,主流需在调用后重新查询最新版本的State:

// 调用子流完成独立交易
val abc = subFlow(UpdateContainerStateFlow2(
    containerId = containerId,
    status = "subflowOne",
    comments = "test"
))
// 重新查询最新版本的ContainerState
val updatedQueryCriteria = QueryCriteria.LinearStateQueryCriteria(linearId = listOf(UniqueIdentifier(id = containerId)))
val latestContainer = serviceHub.vaultService.queryBy<ContainerState>(updatedQueryCriteria).states.single()
// 基于最新状态构建主流交易
val updateContainer = latestContainer.state.data.copy(
    status = status,
    commentsFromStakeHolders = comments
)
val builder = TransactionBuilder(notary)
        .addCommand(ContainerContract.Commands.Update(), listOf(updateContainer.logisticCompany.owningKey, updateContainer.partnerCompany.owningKey))
        .addInputState(latestContainer)
        .addOutputState(updateContainer)

关键原则总结

  • 内部逻辑拆分的子流无需标记@InitiatingFlow,不独立提交交易,由主流统一处理跨节点交互。
  • 同一LinearState的连续更新尽量在单个交易中完成,避免跨交易版本冲突。
  • 禁止嵌套调用@InitiatingFlow,如需跨节点交互,由主流统一管理会话。

内容的提问来源于stack exchange,提问作者Benjamin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 06:05:21