@@ -873,20 +873,22 @@ func runDaemonRemoteStream(ctx context.Context, runner conversation.Runner, now
873873 writeStreamState ()
874874 deliveryOptions := newDaemonRemoteDeliveryOptions (registration , authToken , time.Time {})
875875 result := remotepkg .StreamWebSocketEvents (ctx , options , func (events []remotepkg.PollEvent ) error {
876- delivered , duplicates , ackEvents , ackSent , ackErrors , leaseEvents , leaseExpired , errorsOut := deliverDaemonRemoteEvents (ctx , runner , events , deliveryOptions )
876+ delivery := deliverDaemonRemoteEvents (ctx , runner , events , deliveryOptions )
877877 pumpState .RuntimeState = remotepkg .PumpRunning
878878 pumpState .EventCount += len (events )
879- pumpState .AckEventCount += ackEvents
880- pumpState .AckSentCount += ackSent
881- pumpState .AckErrorCount += ackErrors
882- pumpState .LeaseEventCount += leaseEvents
883- pumpState .LeaseExpiredCount += leaseExpired
884- pumpState .DeliveredCount += delivered
885- pumpState .DuplicateCount += duplicates
886- pumpState .ErrorCount += len (errorsOut )
879+ pumpState .AckEventCount += delivery .AckEvents
880+ pumpState .AckSentCount += delivery .AckSent
881+ pumpState .AckErrorCount += delivery .AckErrors
882+ pumpState .LeaseEventCount += delivery .LeaseEvents
883+ pumpState .LeaseExpiredCount += delivery .LeaseExpired
884+ pumpState .LeaseRenewSent += delivery .LeaseRenewSent
885+ pumpState .LeaseRenewErrors += delivery .LeaseRenewErrors
886+ pumpState .DeliveredCount += delivery .Delivered
887+ pumpState .DuplicateCount += delivery .Duplicates
888+ pumpState .ErrorCount += len (delivery .ErrorsOut )
887889 pumpState .LastPollAt = time .Now ().UTC ().Format (time .RFC3339Nano )
888- if len (errorsOut ) > 0 {
889- pumpState .LastError = fmt .Sprint (errorsOut [0 ]["error" ])
890+ if len (delivery . ErrorsOut ) > 0 {
891+ pumpState .LastError = fmt .Sprint (delivery . ErrorsOut [0 ]["error" ])
890892 }
891893 writeStreamState ()
892894 return nil
@@ -919,6 +921,8 @@ func runDaemonRemoteStream(ctx context.Context, runner conversation.Runner, now
919921 structured ["ack_error_count" ] = pumpState .AckErrorCount
920922 structured ["lease_event_count" ] = pumpState .LeaseEventCount
921923 structured ["lease_expired_count" ] = pumpState .LeaseExpiredCount
924+ structured ["lease_renew_sent_count" ] = pumpState .LeaseRenewSent
925+ structured ["lease_renew_error_count" ] = pumpState .LeaseRenewErrors
922926 structured ["stream_started_at" ] = pumpState .StreamStartedAt
923927 structured ["stream_ended_at" ] = pumpState .StreamEndedAt
924928 structured ["stream_stop_reason" ] = pumpState .StreamStopReason
@@ -1123,17 +1127,19 @@ func runDaemonRemotePoll(ctx context.Context, runner conversation.Runner, now ti
11231127 _ = remotepkg .WritePumpState (pumpPath , pumpState )
11241128 return contracts.ToolResult {Content : "Remote poll failed." , StructuredContent : structured }
11251129 }
1126- delivered , duplicates , ackEvents , ackSent , ackErrors , leaseEvents , leaseExpired , errorsOut := deliverDaemonRemoteEvents (ctx , runner , remoteFetch .Events , newDaemonRemoteDeliveryOptions (registration , authToken , now ))
1127- pumpState .AckEventCount = ackEvents
1128- pumpState .AckSentCount = ackSent
1129- pumpState .AckErrorCount = ackErrors
1130- pumpState .LeaseEventCount = leaseEvents
1131- pumpState .LeaseExpiredCount = leaseExpired
1132- pumpState .DeliveredCount = delivered
1133- pumpState .DuplicateCount = duplicates
1134- pumpState .ErrorCount = len (errorsOut )
1135- if len (errorsOut ) > 0 {
1136- pumpState .LastError = fmt .Sprint (errorsOut [0 ]["error" ])
1130+ delivery := deliverDaemonRemoteEvents (ctx , runner , remoteFetch .Events , newDaemonRemoteDeliveryOptions (registration , authToken , now ))
1131+ pumpState .AckEventCount = delivery .AckEvents
1132+ pumpState .AckSentCount = delivery .AckSent
1133+ pumpState .AckErrorCount = delivery .AckErrors
1134+ pumpState .LeaseEventCount = delivery .LeaseEvents
1135+ pumpState .LeaseExpiredCount = delivery .LeaseExpired
1136+ pumpState .LeaseRenewSent = delivery .LeaseRenewSent
1137+ pumpState .LeaseRenewErrors = delivery .LeaseRenewErrors
1138+ pumpState .DeliveredCount = delivery .Delivered
1139+ pumpState .DuplicateCount = delivery .Duplicates
1140+ pumpState .ErrorCount = len (delivery .ErrorsOut )
1141+ if len (delivery .ErrorsOut ) > 0 {
1142+ pumpState .LastError = fmt .Sprint (delivery .ErrorsOut [0 ]["error" ])
11371143 }
11381144 _ = remotepkg .WritePumpState (pumpPath , pumpState )
11391145 structured ["runtime_state" ] = pumpState .RuntimeState
@@ -1151,16 +1157,18 @@ func runDaemonRemotePoll(ctx context.Context, runner conversation.Runner, now ti
11511157 structured ["ack_error_count" ] = pumpState .AckErrorCount
11521158 structured ["lease_event_count" ] = pumpState .LeaseEventCount
11531159 structured ["lease_expired_count" ] = pumpState .LeaseExpiredCount
1160+ structured ["lease_renew_sent_count" ] = pumpState .LeaseRenewSent
1161+ structured ["lease_renew_error_count" ] = pumpState .LeaseRenewErrors
11541162 structured ["event_count" ] = pumpState .EventCount
11551163 structured ["delivered_count" ] = pumpState .DeliveredCount
11561164 structured ["duplicate_count" ] = pumpState .DuplicateCount
11571165 structured ["error_count" ] = pumpState .ErrorCount
1158- structured ["errors" ] = errorsOut
1166+ structured ["errors" ] = delivery . ErrorsOut
11591167 if remoteFetch .FallbackError != "" {
11601168 structured ["fallback_error" ] = remoteFetch .FallbackError
11611169 }
11621170 return contracts.ToolResult {
1163- Content : fmt .Sprintf ("Remote %s delivered %d event(s); %d duplicate(s); %d error(s)." , remoteFetch .Transport , delivered , duplicates , len (errorsOut )),
1171+ Content : fmt .Sprintf ("Remote %s delivered %d event(s); %d duplicate(s); %d error(s)." , remoteFetch .Transport , delivery . Delivered , delivery . Duplicates , len (delivery . ErrorsOut )),
11641172 StructuredContent : structured ,
11651173 }
11661174}
@@ -1181,17 +1189,20 @@ type daemonRemoteFetch struct {
11811189type daemonRemoteDeliveryOptions struct {
11821190 AuthToken string
11831191 AllowedOrigins []string
1192+ LeaseRenewURL string
11841193 Now time.Time
11851194}
11861195
11871196func newDaemonRemoteDeliveryOptions (registration remotepkg.RegistrationState , authToken string , now time.Time ) daemonRemoteDeliveryOptions {
11881197 return daemonRemoteDeliveryOptions {
11891198 AuthToken : authToken ,
11901199 AllowedOrigins : []string {
1200+ registration .RegistrationURL ,
11911201 registration .PollURL ,
11921202 registration .WebSocketURL ,
11931203 },
1194- Now : now ,
1204+ LeaseRenewURL : registration .LeaseRenewURL ,
1205+ Now : now ,
11951206 }
11961207}
11971208
@@ -1249,40 +1260,52 @@ func fetchDaemonRemoteEvents(ctx context.Context, registration remotepkg.Registr
12491260 }
12501261}
12511262
1252- func deliverDaemonRemoteEvents (ctx context.Context , runner conversation.Runner , events []remotepkg.PollEvent , options daemonRemoteDeliveryOptions ) (int , int , int , int , int , int , int , []map [string ]any ) {
1253- errorsOut := make ([]map [string ]any , 0 )
1254- delivered := 0
1255- duplicates := 0
1256- ackEvents := 0
1257- ackSent := 0
1258- ackErrors := 0
1259- leaseEvents := 0
1260- leaseExpired := 0
1263+ type daemonRemoteDeliveryResult struct {
1264+ Delivered int
1265+ Duplicates int
1266+ AckEvents int
1267+ AckSent int
1268+ AckErrors int
1269+ LeaseEvents int
1270+ LeaseExpired int
1271+ LeaseRenewSent int
1272+ LeaseRenewErrors int
1273+ ErrorsOut []map [string ]any
1274+ }
1275+
1276+ func deliverDaemonRemoteEvents (ctx context.Context , runner conversation.Runner , events []remotepkg.PollEvent , options daemonRemoteDeliveryOptions ) daemonRemoteDeliveryResult {
1277+ result := daemonRemoteDeliveryResult {ErrorsOut : make ([]map [string ]any , 0 )}
12611278 for _ , event := range events {
12621279 if event .AckURL != "" {
1263- ackEvents ++
1280+ result . AckEvents ++
12641281 }
12651282 if event .LeaseID != "" || event .LeaseExpiresAt != "" {
1266- leaseEvents ++
1283+ result . LeaseEvents ++
12671284 }
12681285 if daemonRemoteLeaseExpired (event , options .Now ) {
1269- leaseExpired ++
1286+ result . LeaseExpired ++
12701287 errorText := "remote event lease expired"
1271- errorsOut = append (errorsOut , map [string ]any {
1288+ result . ErrorsOut = append (result . ErrorsOut , map [string ]any {
12721289 "type" : "remote_lease" ,
12731290 "event_id" : event .EventID ,
12741291 "lease_id" : event .LeaseID ,
12751292 "lease_expires_at" : event .LeaseExpiresAt ,
12761293 "error" : errorText ,
12771294 })
12781295 sent , ackErr := acknowledgeDaemonRemoteEvent (ctx , event , options , "expired" , 0 , false , errorText )
1279- ackSent += sent
1296+ result . AckSent += sent
12801297 if ackErr != nil {
1281- ackErrors ++
1282- errorsOut = append (errorsOut , ackErr )
1298+ result . AckErrors ++
1299+ result . ErrorsOut = append (result . ErrorsOut , ackErr )
12831300 }
12841301 continue
12851302 }
1303+ sent , renewErr := renewDaemonRemoteLease (ctx , event , options )
1304+ result .LeaseRenewSent += sent
1305+ if renewErr != nil {
1306+ result .LeaseRenewErrors ++
1307+ result .ErrorsOut = append (result .ErrorsOut , renewErr )
1308+ }
12861309 input , err := json .Marshal (map [string ]any {
12871310 "team_id" : event .TeamID ,
12881311 "target" : event .Target ,
@@ -1292,16 +1315,16 @@ func deliverDaemonRemoteEvents(ctx context.Context, runner conversation.Runner,
12921315 "message" : event .Message ,
12931316 })
12941317 if err != nil {
1295- errorsOut = append (errorsOut , map [string ]any {"event_id" : event .EventID , "error" : err .Error ()})
1318+ result . ErrorsOut = append (result . ErrorsOut , map [string ]any {"event_id" : event .EventID , "error" : err .Error ()})
12961319 sent , ackErr := acknowledgeDaemonRemoteEvent (ctx , event , options , "failed" , 0 , false , err .Error ())
1297- ackSent += sent
1320+ result . AckSent += sent
12981321 if ackErr != nil {
1299- ackErrors ++
1300- errorsOut = append (errorsOut , ackErr )
1322+ result . AckErrors ++
1323+ result . ErrorsOut = append (result . ErrorsOut , ackErr )
13011324 }
13021325 continue
13031326 }
1304- result , err := tasktools .RunRemoteTrigger (tool.Context {
1327+ toolResult , err := tasktools .RunRemoteTrigger (tool.Context {
13051328 Context : ctx ,
13061329 WorkingDirectory : runner .WorkingDirectory ,
13071330 SessionID : runner .SessionID ,
@@ -1310,35 +1333,35 @@ func deliverDaemonRemoteEvents(ctx context.Context, runner conversation.Runner,
13101333 },
13111334 }, input , tool .NopProgressSink ())
13121335 if err != nil {
1313- errorsOut = append (errorsOut , map [string ]any {"event_id" : event .EventID , "team_id" : event .TeamID , "error" : err .Error ()})
1336+ result . ErrorsOut = append (result . ErrorsOut , map [string ]any {"event_id" : event .EventID , "team_id" : event .TeamID , "error" : err .Error ()})
13141337 sent , ackErr := acknowledgeDaemonRemoteEvent (ctx , event , options , "failed" , 0 , false , err .Error ())
1315- ackSent += sent
1338+ result . AckSent += sent
13161339 if ackErr != nil {
1317- ackErrors ++
1318- errorsOut = append (errorsOut , ackErr )
1340+ result . AckErrors ++
1341+ result . ErrorsOut = append (result . ErrorsOut , ackErr )
13191342 }
13201343 continue
13211344 }
1322- if duplicate , _ := result .StructuredContent ["duplicate" ].(bool ); duplicate {
1323- duplicates ++
1345+ if duplicate , _ := toolResult .StructuredContent ["duplicate" ].(bool ); duplicate {
1346+ result . Duplicates ++
13241347 sent , ackErr := acknowledgeDaemonRemoteEvent (ctx , event , options , "duplicate" , 0 , true , "" )
1325- ackSent += sent
1348+ result . AckSent += sent
13261349 if ackErr != nil {
1327- ackErrors ++
1328- errorsOut = append (errorsOut , ackErr )
1350+ result . AckErrors ++
1351+ result . ErrorsOut = append (result . ErrorsOut , ackErr )
13291352 }
13301353 continue
13311354 }
1332- sentCount := intMapValue (result .StructuredContent , "sent_count" )
1333- delivered += sentCount
1355+ sentCount := intMapValue (toolResult .StructuredContent , "sent_count" )
1356+ result . Delivered += sentCount
13341357 sent , ackErr := acknowledgeDaemonRemoteEvent (ctx , event , options , "delivered" , sentCount , false , "" )
1335- ackSent += sent
1358+ result . AckSent += sent
13361359 if ackErr != nil {
1337- ackErrors ++
1338- errorsOut = append (errorsOut , ackErr )
1360+ result . AckErrors ++
1361+ result . ErrorsOut = append (result . ErrorsOut , ackErr )
13391362 }
13401363 }
1341- return delivered , duplicates , ackEvents , ackSent , ackErrors , leaseEvents , leaseExpired , errorsOut
1364+ return result
13421365}
13431366
13441367func daemonRemoteLeaseExpired (event remotepkg.PollEvent , now time.Time ) bool {
@@ -1355,6 +1378,28 @@ func daemonRemoteLeaseExpired(event remotepkg.PollEvent, now time.Time) bool {
13551378 return ! expiresAt .After (now .UTC ())
13561379}
13571380
1381+ func renewDaemonRemoteLease (ctx context.Context , event remotepkg.PollEvent , options daemonRemoteDeliveryOptions ) (int , map [string ]any ) {
1382+ if strings .TrimSpace (event .LeaseID ) == "" || strings .TrimSpace (options .LeaseRenewURL ) == "" {
1383+ return 0 , nil
1384+ }
1385+ result := remotepkg .SendLeaseRenewal (ctx , remotepkg.LeaseRenewOptions {
1386+ LeaseRenewURL : options .LeaseRenewURL ,
1387+ AuthToken : options .AuthToken ,
1388+ EventID : event .EventID ,
1389+ LeaseID : event .LeaseID ,
1390+ AllowedOrigins : options .AllowedOrigins ,
1391+ })
1392+ if result .Error != "" {
1393+ return 0 , map [string ]any {
1394+ "type" : "remote_lease_renew" ,
1395+ "event_id" : event .EventID ,
1396+ "lease_id" : event .LeaseID ,
1397+ "error" : result .Error ,
1398+ }
1399+ }
1400+ return 1 , nil
1401+ }
1402+
13581403func acknowledgeDaemonRemoteEvent (ctx context.Context , event remotepkg.PollEvent , options daemonRemoteDeliveryOptions , status string , sentCount int , duplicate bool , errorText string ) (int , map [string ]any ) {
13591404 if event .AckURL == "" {
13601405 return 0 , nil
0 commit comments