| 73 | } |
| 74 | |
| 75 | func newNetworkTaskStatusObserver( |
| 76 | networkRuntime networkRuntime, |
| 77 | tasks taskStore, |
| 78 | opts ...networkTaskStatusObserverOption, |
| 79 | ) *networkTaskStatusObserver { |
| 80 | prefs, ok := tasks.(storepkg.NetworkPreferenceStore) |
| 81 | if networkRuntime == nil || tasks == nil || !ok { |
| 82 | return nil |
| 83 | } |
| 84 | options := networkTaskStatusObserverOptions{ |
| 85 | logger: slog.Default(), |
| 86 | now: time.Now, |
| 87 | queueSize: aghconfig.DefaultTaskNetworkStatusQueueSize, |
| 88 | timeout: aghconfig.DefaultTaskNetworkStatusTimeout, |
| 89 | } |
| 90 | for _, opt := range opts { |
| 91 | if opt != nil { |
| 92 | opt(&options) |
| 93 | } |
| 94 | } |
| 95 | if options.logger == nil { |
| 96 | options.logger = slog.Default() |
| 97 | } |
| 98 | if options.now == nil { |
| 99 | options.now = time.Now |
| 100 | } |
| 101 | if options.queueSize <= 0 { |
| 102 | options.queueSize = aghconfig.DefaultTaskNetworkStatusQueueSize |
| 103 | } |
| 104 | if options.timeout <= 0 { |
| 105 | options.timeout = aghconfig.DefaultTaskNetworkStatusTimeout |
| 106 | } |
| 107 | ctx, cancel := context.WithCancel(context.Background()) |
| 108 | observer := &networkTaskStatusObserver{ |
| 109 | network: networkRuntime, |
| 110 | tasks: tasks, |
| 111 | prefs: prefs, |
| 112 | logger: options.logger, |
| 113 | now: options.now, |
| 114 | ctx: ctx, |
| 115 | cancel: cancel, |
| 116 | queue: make(chan taskpkg.EventRecord, options.queueSize), |
| 117 | timeout: options.timeout, |
| 118 | } |
| 119 | observer.start() |
| 120 | return observer |
| 121 | } |
| 122 | |
| 123 | func (o *networkTaskStatusObserver) OnTaskEvent(_ context.Context, record taskpkg.EventRecord) { |
| 124 | if o == nil || !networkTaskStatusEvent(record.Event.EventType) { |