MCPcopy Create free account
hub / github.com/QuantumNous/new-api / waitTaskPluginProtocol

Function waitTaskPluginProtocol

controller/plugin_protocol.go:639–829  ·  view source on GitHub ↗
(
	c *gin.Context,
	pinned pluginruntime.PinnedEndpoint,
	protocolRequest pluginruntime.ProtocolRequestContext,
	taskID string,
	machine *relay.PluginResponsesMachine,
	deps pluginProtocolBridgeDeps,
)

Source from the content-addressed store, hash-verified

637}
638
639func 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 }

Callers 1

serveTaskPluginProtocolFunction · 0.85

Calls 10

LogDebugFunction · 0.92
LogWarnFunction · 0.92
LogErrorFunction · 0.92
taskPluginDebugStatusFunction · 0.85
pluginProtocolTickDelayFunction · 0.85
DoneMethod · 0.80
StopMethod · 0.80

Tested by

no test coverage detected