(
sessionId: string,
input: UserMessageInput,
options: TurnStartOptions = {},
)
| 646 | } |
| 647 | |
| 648 | async *startTurn( |
| 649 | sessionId: string, |
| 650 | input: UserMessageInput, |
| 651 | options: TurnStartOptions = {}, |
| 652 | ): AsyncIterable<SessionEvent> { |
| 653 | if (this.pendingContinuationSessions.has(sessionId)) { |
| 654 | throw new Error('Cannot start a turn while a runtime continuation is being claimed'); |
| 655 | } |
| 656 | const execution = this.takeExecutionClaim(sessionId, options.execution); |
| 657 | try { |
| 658 | await this.enterExecutionClaim(execution); |
| 659 | const header = await this.deps.store.readHeader(sessionId); |
| 660 | let workspaceIdentity: string | undefined; |
| 661 | if (this.deps.safeBoundaryResumeEnabled === true && this.deps.inspectContinuationSafety) { |
| 662 | try { |
| 663 | workspaceIdentity = (await this.deps.inspectContinuationSafety(sessionId)) |
| 664 | .workspaceIdentity; |
| 665 | } catch { |
| 666 | // A new turn remains usable without continuation metadata. Actual |
| 667 | // continuation claims inspect the same facts strictly below. |
| 668 | } |
| 669 | } |
| 670 | const run = new AgentRun({ |
| 671 | sessionId, |
| 672 | header, |
| 673 | userInput: input, |
| 674 | runId: options.runId, |
| 675 | userMessageId: options.userMessageId, |
| 676 | durability: options.durability, |
| 677 | store: this.deps.store, |
| 678 | runStore: this.deps.runStore, |
| 679 | runtimeEventStore: this.deps.runtimeEventStore, |
| 680 | ...(runtimeToolBoundaryProtocol(this.deps, header) |
| 681 | ? { toolBoundaryProtocol: runtimeToolBoundaryProtocol(this.deps, header) } |
| 682 | : {}), |
| 683 | repairRunRuntimeLedger: this.deps.repairRunRuntimeLedger, |
| 684 | newId: this.deps.newId, |
| 685 | now: this.deps.now, |
| 686 | ...(workspaceIdentity ? { workspaceIdentity } : {}), |
| 687 | hooks: { |
| 688 | reserveRun: async (targetSessionId, nextHeader, activeRun) => { |
| 689 | const active = await this.reserveParentRun( |
| 690 | targetSessionId, |
| 691 | nextHeader, |
| 692 | activeRun, |
| 693 | execution, |
| 694 | ); |
| 695 | this.reserveExecutionClaim(execution, active, activeRun); |
| 696 | return active; |
| 697 | }, |
| 698 | unregisterRun: (active, activeRun) => this.unregisterParentRun(active, activeRun), |
| 699 | updateHeader: (targetSessionId, patch) => this.updateHeader(targetSessionId, patch), |
| 700 | updateStatus: (targetSessionId, status, blockedReason, ts) => |
| 701 | this.updateStatus(targetSessionId, status, blockedReason, ts), |
| 702 | appendTurnState: (targetSessionId, turnId, status, lineage, options) => |
| 703 | this.appendTurnState(targetSessionId, turnId, status, lineage, options), |
| 704 | }, |
| 705 | }); |
nothing calls this directly
no test coverage detected