@@ -8,16 +8,20 @@ import {
88 type ModelSelection ,
99 NodeId ,
1010 type OrchestrationV2Run ,
11+ type OrchestrationV2ThreadProjection ,
1112 ProjectId ,
1213 ProviderDriverKind ,
1314 ProviderInstanceId ,
1415 ProviderThreadId ,
16+ RunAttemptId ,
1517 RunId ,
18+ RuntimeRequestId ,
1619 ThreadId ,
1720} from "@t3tools/contracts" ;
1821import * as DateTime from "effect/DateTime" ;
1922import * as Effect from "effect/Effect" ;
2023import * as Layer from "effect/Layer" ;
24+ import * as Option from "effect/Option" ;
2125import * as Stream from "effect/Stream" ;
2226
2327import * as CheckpointStore from "../checkpointing/CheckpointStore.ts" ;
@@ -38,6 +42,7 @@ import * as Orchestrator from "./Orchestrator.ts";
3842import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts" ;
3943import { continueRestartedRun } from "./RestartContinuation.ts" ;
4044import * as RuntimeLayer from "./runtimeLayer.ts" ;
45+ import { makeSubagentChildThread } from "./SubagentProjection.ts" ;
4146import * as ProviderTurnStartServiceTestkit from "./ProviderTurnStartService.testkit.ts" ;
4247
4348const layerPlatformTest = Layer . merge (
@@ -474,6 +479,304 @@ it.layer(layerTest)("delegated completion delivery repairs", (it) => {
474479 } ) ,
475480 ) ;
476481
482+ it . effect ( "tells the parent once when its child blocks on a request" , ( ) =>
483+ Effect . gen ( function * ( ) {
484+ const orchestrator = yield * Orchestrator . OrchestratorV2 ;
485+ const sink = yield * EventSink . EventSinkV2 ;
486+ const now = yield * DateTime . now ;
487+ const threadId = ThreadId . make ( "thread:delegated-task-blocked" ) ;
488+ const runId = RunId . make ( "run:delegated-task-blocked" ) ;
489+ const rootNodeId = NodeId . make ( "node:delegated-task-blocked-root" ) ;
490+ const taskId = NodeId . make ( "node:delegated-task-blocked-task" ) ;
491+ const secondTaskId = NodeId . make ( "node:delegated-task-blocked-second" ) ;
492+ const childThreadId = ThreadId . make ( "thread:delegated-task-blocked-child" ) ;
493+ const secondChildThreadId = ThreadId . make ( "thread:delegated-task-blocked-second-child" ) ;
494+ yield * seedParentWithTerminalTask ( {
495+ threadId,
496+ runId,
497+ projectId : ProjectId . make ( "project:delegated-task-blocked" ) ,
498+ rootNodeId,
499+ taskId,
500+ deliveryState : "claimed" ,
501+ completionWake : "always" ,
502+ now,
503+ } ) ;
504+ const parent = yield * orchestrator . getThreadProjection ( threadId ) ;
505+ const parentRun = parent . runs . find ( ( run ) => run . id === runId ) ! ;
506+ const providerThreadId = parentRun . providerThreadId ! ;
507+ const attemptId = RunAttemptId . make ( "attempt:delegated-task-blocked" ) ;
508+ const runningTask = ( id : NodeId , child : ThreadId , title : string , taskRunId = runId ) => [
509+ {
510+ id : EventId . make ( `event:delegated-task-blocked:child:${ id } ` ) ,
511+ type : "thread.created" as const ,
512+ threadId : child ,
513+ driver,
514+ providerInstanceId : modelSelection . instanceId ,
515+ occurredAt : now ,
516+ payload : makeSubagentChildThread ( {
517+ parentThread : parent . thread ,
518+ childThreadId : child ,
519+ parentNodeId : id ,
520+ activeProviderThreadId : null ,
521+ providerInstanceId : modelSelection . instanceId ,
522+ modelSelection,
523+ title,
524+ now,
525+ createdBy : "agent" ,
526+ creationSource : "server" ,
527+ } ) ,
528+ } ,
529+ {
530+ id : EventId . make ( `event:delegated-task-blocked:task:${ id } ` ) ,
531+ type : "subagent.updated" as const ,
532+ threadId,
533+ runId : taskRunId ,
534+ nodeId : id ,
535+ driver,
536+ providerInstanceId : modelSelection . instanceId ,
537+ occurredAt : now ,
538+ payload : {
539+ ...parent . subagents [ 0 ] ! ,
540+ id,
541+ runId : taskRunId ,
542+ title,
543+ childThreadId : child ,
544+ status : "running" as const ,
545+ result : null ,
546+ completedAt : null ,
547+ completionDelivery : undefined ,
548+ } ,
549+ } ,
550+ ] ;
551+ // Two running children, and a parent run that Stop can interrupt.
552+ yield * sink . write ( {
553+ commandId : CommandId . make ( "command:delegated-task-blocked:seed" ) ,
554+ events : [
555+ ...runningTask ( taskId , childThreadId , "Load test lane" ) ,
556+ // Spawned by an earlier parent turn, so Stop's cohort barrier alone
557+ // would not reach its notices.
558+ ...runningTask (
559+ secondTaskId ,
560+ secondChildThreadId ,
561+ "Second lane" ,
562+ RunId . make ( "run:delegated-task-blocked-earlier" ) ,
563+ ) ,
564+ {
565+ id : EventId . make ( "event:delegated-task-blocked:root" ) ,
566+ type : "node.updated" ,
567+ threadId,
568+ runId,
569+ nodeId : rootNodeId ,
570+ occurredAt : now ,
571+ payload : {
572+ id : rootNodeId ,
573+ threadId,
574+ runId,
575+ parentNodeId : null ,
576+ rootNodeId,
577+ kind : "root_turn" ,
578+ status : "running" ,
579+ countsForRun : true ,
580+ providerThreadId,
581+ providerTurnId : null ,
582+ nativeItemRef : null ,
583+ runtimeRequestId : null ,
584+ checkpointScopeId : null ,
585+ startedAt : now ,
586+ completedAt : null ,
587+ } ,
588+ } ,
589+ {
590+ id : EventId . make ( "event:delegated-task-blocked:attempt" ) ,
591+ type : "run-attempt.updated" ,
592+ threadId,
593+ runId,
594+ nodeId : rootNodeId ,
595+ providerInstanceId : modelSelection . instanceId ,
596+ occurredAt : now ,
597+ payload : {
598+ id : attemptId ,
599+ runId,
600+ attemptOrdinal : 1 ,
601+ rootNodeId,
602+ providerInstanceId : modelSelection . instanceId ,
603+ providerThreadId,
604+ providerTurnId : null ,
605+ reason : "initial" ,
606+ status : "running" ,
607+ startedAt : now ,
608+ completedAt : null ,
609+ } ,
610+ } ,
611+ {
612+ id : EventId . make ( "event:delegated-task-blocked:active-attempt" ) ,
613+ type : "run.updated" ,
614+ threadId,
615+ runId,
616+ nodeId : rootNodeId ,
617+ providerInstanceId : modelSelection . instanceId ,
618+ occurredAt : now ,
619+ payload : { ...parentRun , activeAttemptId : attemptId } ,
620+ } ,
621+ ] ,
622+ } ) ;
623+
624+ // Records a pending request on a child. Returns the parent sequence to
625+ // wait from for the notice that the blocked-child reactor queues.
626+ const block = (
627+ child : ThreadId ,
628+ nodeId : NodeId ,
629+ id : string ,
630+ kind : "user_input" | "command" | "auth_refresh" ,
631+ attempt = 0 ,
632+ ) =>
633+ Effect . gen ( function * ( ) {
634+ const afterSequence = yield * sink . latestSequence ( { threadId } ) ;
635+ yield * sink . write ( {
636+ commandId : CommandId . make ( `command:delegated-task-blocked:${ id } :${ attempt } ` ) ,
637+ events : [
638+ {
639+ id : EventId . make ( `event:delegated-task-blocked:${ id } :${ attempt } ` ) ,
640+ type : "runtime-request.updated" ,
641+ threadId : child ,
642+ nodeId,
643+ occurredAt : now ,
644+ payload : {
645+ id : RuntimeRequestId . make ( id ) ,
646+ nodeId,
647+ providerTurnId : null ,
648+ nativeRequestRef : null ,
649+ kind,
650+ status : "pending" ,
651+ responseCapability : { type : "message" } ,
652+ createdAt : now ,
653+ resolvedAt : null ,
654+ } ,
655+ } ,
656+ ] ,
657+ } ) ;
658+ return afterSequence ;
659+ } ) ;
660+ const nextNotice = ( afterSequence : number ) =>
661+ sink . stream ( { threadId, afterSequence, eventType : "message.updated" } ) . pipe (
662+ Stream . filter (
663+ ( stored ) =>
664+ stored . event . type === "message.updated" &&
665+ stored . event . payload . notification ?. source . kind === "delegated_task" ,
666+ ) ,
667+ Stream . runHead ,
668+ Effect . map ( ( stored ) =>
669+ Option . isSome ( stored ) && stored . value . event . type === "message.updated"
670+ ? stored . value . event . payload
671+ : undefined ,
672+ ) ,
673+ ) ;
674+ const noticeRunStatus = (
675+ projection : OrchestrationV2ThreadProjection ,
676+ messageId : MessageId | undefined ,
677+ ) => projection . runs . find ( ( run ) => run . userMessageId === messageId ) ?. status ;
678+
679+ const question = yield * nextNotice (
680+ yield * block ( childThreadId , taskId , "request:question" , "user_input" ) ,
681+ ) ;
682+ assert . deepEqual ( question ?. notification , {
683+ source : { kind : "delegated_task" , taskIds : [ taskId ] } ,
684+ outcome : "updated" ,
685+ summary : "Load test lane is waiting for an answer to a question" ,
686+ } ) ;
687+ assert . include ( question ?. text ?? "" , "t3_pending_request_respond" ) ;
688+ assert . include ( question ?. text ?? "" , String ( childThreadId ) ) ;
689+
690+ // A repeated update for the same request must not queue a second notice.
691+ // The reactor handles events in order, so the approval notice below
692+ // arrives only after the repeat has been processed.
693+ yield * block ( childThreadId , taskId , "request:question" , "user_input" , 1 ) ;
694+ const approval = yield * nextNotice (
695+ yield * block ( childThreadId , taskId , "request:approval" , "command" ) ,
696+ ) ;
697+ assert . equal (
698+ approval ?. notification ?. summary ,
699+ "Load test lane is waiting for approval (command)" ,
700+ ) ;
701+ assert . include ( approval ?. text ?? "" , "Only the user can resolve this" ) ;
702+
703+ const signIn = yield * nextNotice (
704+ yield * block ( childThreadId , taskId , "request:sign-in" , "auth_refresh" ) ,
705+ ) ;
706+ assert . equal (
707+ signIn ?. notification ?. summary ,
708+ "Load test lane is waiting for the user to sign in again" ,
709+ ) ;
710+ assert . include ( signIn ?. text ?? "" , "Only the user can resolve this" ) ;
711+
712+ // A rejected earlier attempt must not swallow a later update's notice.
713+ yield * sink . commitRejectedCommand ( {
714+ commandId : CommandId . make ( "command:delegated-task-blocked:request:retry" ) ,
715+ threadId,
716+ commandType : "message.dispatch" ,
717+ rejectedAt : now ,
718+ error : "Thread has a pending merge-back transfer." ,
719+ } ) ;
720+ const retried = yield * nextNotice (
721+ yield * block ( childThreadId , taskId , "request:retry" , "user_input" ) ,
722+ ) ;
723+ assert . equal (
724+ retried ?. notification ?. summary ,
725+ "Load test lane is waiting for an answer to a question" ,
726+ ) ;
727+
728+ const blocked = yield * orchestrator . getThreadProjection ( threadId ) ;
729+ assert . equal (
730+ blocked . messages . filter (
731+ ( message ) =>
732+ message . notification ?. source . kind === "delegated_task" &&
733+ message . notification . source . taskIds . includes ( taskId ) ,
734+ ) . length ,
735+ 4 ,
736+ ) ;
737+ // The child is not finished: no result is published to the parent.
738+ assert . equal ( blocked . subagents . find ( ( row ) => row . id === taskId ) ?. status , "running" ) ;
739+ assert . isFalse ( blocked . contextTransfers . some ( ( row ) => row . type === "subagent_result" ) ) ;
740+ assert . equal ( noticeRunStatus ( blocked , question ?. id ) , "queued" ) ;
741+ assert . equal ( noticeRunStatus ( blocked , approval ?. id ) , "queued" ) ;
742+
743+ // Once the parent drops a task (task_cancel), its queued notices must
744+ // not start a parent turn.
745+ yield * orchestrator . dispatch ( {
746+ type : "delegated_task.completion-delivery.dispose" ,
747+ commandId : CommandId . make ( "command:delegated-task-blocked:dispose" ) ,
748+ parentThreadId : threadId ,
749+ taskId,
750+ } ) ;
751+ const disposed = yield * orchestrator . getThreadProjection ( threadId ) ;
752+ assert . equal ( noticeRunStatus ( disposed , question ?. id ) , "cancelled" ) ;
753+ assert . equal ( noticeRunStatus ( disposed , approval ?. id ) , "cancelled" ) ;
754+ assert . equal ( noticeRunStatus ( disposed , signIn ?. id ) , "cancelled" ) ;
755+ assert . equal ( noticeRunStatus ( disposed , retried ?. id ) , "cancelled" ) ;
756+
757+ // Stop must not let any queued notice start a new parent turn either.
758+ const second = yield * nextNotice (
759+ yield * block ( secondChildThreadId , secondTaskId , "request:second" , "user_input" ) ,
760+ ) ;
761+ assert . equal (
762+ noticeRunStatus ( yield * orchestrator . getThreadProjection ( threadId ) , second ?. id ) ,
763+ "queued" ,
764+ ) ;
765+ yield * orchestrator . dispatch ( {
766+ type : "run.interrupt" ,
767+ commandId : CommandId . make ( "command:delegated-task-blocked:stop-parent" ) ,
768+ threadId,
769+ runId,
770+ reason : "User stopped the parent." ,
771+ holdQueue : true ,
772+ } ) ;
773+ assert . equal (
774+ noticeRunStatus ( yield * orchestrator . getThreadProjection ( threadId ) , second ?. id ) ,
775+ "cancelled" ,
776+ ) ;
777+ } ) ,
778+ ) ;
779+
477780 it . effect ( "builds completion text and metadata from the same live cohort" , ( ) =>
478781 Effect . gen ( function * ( ) {
479782 const orchestrator = yield * Orchestrator . OrchestratorV2 ;
0 commit comments