(event: StreamEndEvent)
| 7745 | } |
| 7746 | |
| 7747 | private async markWorkspaceTurnStreamEndDeferred(event: StreamEndEvent): Promise<void> { |
| 7748 | const metadata = this.getWorkspaceTurnMetadata(event); |
| 7749 | if (metadata == null) { |
| 7750 | return; |
| 7751 | } |
| 7752 | await this.workspaceTurnSettlementLocks.withLock(metadata.taskHandleId, async () => { |
| 7753 | const record = await this.taskHandleStore.getWorkspaceTurn( |
| 7754 | metadata.ownerWorkspaceId, |
| 7755 | metadata.taskHandleId |
| 7756 | ); |
| 7757 | if ( |
| 7758 | record == null || |
| 7759 | record.workspaceId !== event.workspaceId || |
| 7760 | record.turnId !== metadata.turnId || |
| 7761 | !this.isActiveWorkspaceTurn(record) || |
| 7762 | this.isDeferredWorkspaceTurnMessage(record, event.messageId) |
| 7763 | ) { |
| 7764 | return; |
| 7765 | } |
| 7766 | await this.taskHandleStore.upsertWorkspaceTurn({ |
| 7767 | ...record, |
| 7768 | updatedAt: getIsoNow(), |
| 7769 | deferredMessageIds: [...(record.deferredMessageIds ?? []), event.messageId], |
| 7770 | }); |
| 7771 | }); |
| 7772 | } |
| 7773 | |
| 7774 | private resolveWorkspaceTurnMuxMetadataForStreamEnd( |
| 7775 | event: StreamEndEvent |
no test coverage detected