( c *gin.Context, pinned pluginruntime.PinnedEndpoint, protocolRequest pluginruntime.ProtocolRequestContext, taskID string, machine *relay.PluginResponsesMachine, deps pluginProtocolBridgeDeps, )
| 637 | } |
| 638 | |
| 639 | func waitTaskPluginProtocol( |
| 640 | c *gin.Context, |
| 641 | pinned pluginruntime.PinnedEndpoint, |
| 642 | protocolRequest pluginruntime.ProtocolRequestContext, |
| 643 | taskID string, |
| 644 | machine *relay.PluginResponsesMachine, |
| 645 | deps pluginProtocolBridgeDeps, |
| 646 | ) { |
| 647 | generation := pinned.Generation.Number |
| 648 | pluginKey := pinned.Plugin.Meta.Key |
| 649 | logger.LogDebug( |
| 650 | c, |
| 651 | "task_plugin subsystem=protocol event=observation_start generation=%d plugin=%q mode=nonstream public_task_id=%q timeout_ms=%d tick_ms=%d", |
| 652 | generation, |
| 653 | pluginKey, |
| 654 | taskID, |
| 655 | deps.observationTimeout.Milliseconds(), |
| 656 | deps.tickInterval.Milliseconds(), |
| 657 | ) |
| 658 | observationContext, cancelObservation := context.WithTimeout(c.Request.Context(), deps.observationTimeout) |
| 659 | defer cancelObservation() |
| 660 | tickNumber := uint64(0) |
| 661 | lastStatus := "" |
| 662 | for { |
| 663 | loadStarted := deps.now() |
| 664 | loadContext, cancelLoad := context.WithTimeout(observationContext, deps.loadTimeout) |
| 665 | task, exists, err := deps.loadTask( |
| 666 | loadContext, |
| 667 | common.GetContextKeyInt(c, constant.ContextKeyUserId), |
| 668 | constant.TaskPlatform(pinned.Plugin.Meta.Key), |
| 669 | taskID, |
| 670 | ) |
| 671 | loadContextErr := loadContext.Err() |
| 672 | cancelLoad() |
| 673 | loadElapsed := deps.now().Sub(loadStarted) |
| 674 | loadOverloaded := errors.Is(loadContextErr, context.DeadlineExceeded) && |
| 675 | observationContext.Err() == nil && |
| 676 | c.Request.Context().Err() == nil |
| 677 | if loadOverloaded { |
| 678 | logger.LogWarn(c, fmt.Sprintf( |
| 679 | "task protocol database observation overloaded; plugin=%s task=%s", |
| 680 | pinned.Plugin.Meta.Key, |
| 681 | taskID, |
| 682 | )) |
| 683 | logger.LogDebug( |
| 684 | c, |
| 685 | "task_plugin subsystem=protocol event=observation_tick generation=%d plugin=%q mode=nonstream tick=%d load_ms=%d overloaded=true", |
| 686 | generation, |
| 687 | pluginKey, |
| 688 | tickNumber, |
| 689 | loadElapsed.Milliseconds(), |
| 690 | ) |
| 691 | } else if err != nil || !exists || task == nil { |
| 692 | if errors.Is(observationContext.Err(), context.DeadlineExceeded) { |
| 693 | logger.LogDebug(c, "task_plugin subsystem=protocol event=observation_timeout generation=%d plugin=%q mode=nonstream last_status=%q", generation, pluginKey, taskPluginDebugStatus(lastStatus)) |
| 694 | writeTaskPluginProtocolTimeoutResponse(c, machine, lastStatus) |
| 695 | return |
| 696 | } |
no test coverage detected