-
Notifications
You must be signed in to change notification settings - Fork 163
/
starter.go
executable file
·1676 lines (1479 loc) · 50.9 KB
/
starter.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
package starter
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"io/ioutil"
"math"
"math/big"
"net"
"net/http"
"net/url"
"os"
"os/user"
"path"
"path/filepath"
"strconv"
"strings"
"time"
ethcommon "github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/ethclient"
"github.com/ethereum/go-ethereum/params"
"github.com/ethereum/go-ethereum/rpc"
"github.com/golang/glog"
"github.com/livepeer/ai-worker/worker"
"github.com/livepeer/go-livepeer/build"
"github.com/livepeer/go-livepeer/common"
"github.com/livepeer/go-livepeer/core"
"github.com/livepeer/go-livepeer/discovery"
"github.com/livepeer/go-livepeer/eth"
"github.com/livepeer/go-livepeer/eth/blockwatch"
"github.com/livepeer/go-livepeer/eth/watchers"
lpmon "github.com/livepeer/go-livepeer/monitor"
"github.com/livepeer/go-livepeer/pm"
"github.com/livepeer/go-livepeer/server"
"github.com/livepeer/go-livepeer/verification"
"github.com/livepeer/go-tools/drivers"
"github.com/livepeer/livepeer-data/pkg/event"
"github.com/livepeer/lpms/ffmpeg"
)
var (
// The timeout for ETH RPC calls
ethRPCTimeout = 20 * time.Second
// The maximum blocks for the block watcher to retain
blockWatcherRetentionLimit = 20
// Estimate of the gas required to redeem a PM ticket on L1 Ethereum
redeemGasL1 = 350000
// Estimate of the gas required to redeem a PM ticket on L2 Arbitrum
redeemGasL2 = 1200000
// The multiplier on the transaction cost to use for PM ticket faceValue
txCostMultiplier = 100
// The interval at which to clean up cached max float values for PM senders and balances per stream
cleanupInterval = 10 * time.Minute
// The time to live for cached max float values for PM senders (else they will be cleaned up) in seconds
smTTL = 172800 // 2 days
aiWorkerContainerImageID = "livepeer/ai-runner:latest"
aiWorkerContainerStopTimeout = 5 * time.Second
)
const (
BroadcasterRpcPort = "9935"
BroadcasterCliPort = "5935"
BroadcasterRtmpPort = "1935"
OrchestratorRpcPort = "8935"
OrchestratorCliPort = "7935"
TranscoderCliPort = "6935"
RefreshPerfScoreInterval = 10 * time.Minute
)
type LivepeerConfig struct {
Network *string
RtmpAddr *string
CliAddr *string
HttpAddr *string
ServiceAddr *string
OrchAddr *string
VerifierURL *string
EthController *string
VerifierPath *string
LocalVerify *bool
HttpIngest *bool
Orchestrator *bool
Transcoder *bool
AIWorker *bool
Broadcaster *bool
OrchSecret *string
TranscodingOptions *string
AIModels *string
MaxAttempts *int
SelectRandWeight *float64
SelectStakeWeight *float64
SelectPriceWeight *float64
SelectPriceExpFactor *float64
OrchPerfStatsURL *string
Region *string
MaxPricePerUnit *int
MinPerfScore *float64
MaxSessions *string
CurrentManifest *bool
Nvidia *string
Netint *string
TestTranscoder *bool
EthAcctAddr *string
EthPassword *string
EthKeystorePath *string
EthOrchAddr *string
EthUrl *string
TxTimeout *time.Duration
MaxTxReplacements *int
GasLimit *int
MinGasPrice *int64
MaxGasPrice *int
InitializeRound *bool
TicketEV *string
MaxFaceValue *string
MaxTicketEV *string
MaxTotalEV *string
DepositMultiplier *int
PricePerUnit *int
PixelsPerUnit *int
AutoAdjustPrice *bool
PricePerBroadcaster *string
BlockPollingInterval *int
Redeemer *bool
RedeemerAddr *string
Reward *bool
Monitor *bool
MetricsPerStream *bool
MetricsExposeClientIP *bool
MetadataQueueUri *string
MetadataAmqpExchange *string
MetadataPublishTimeout *time.Duration
Datadir *string
AIModelsDir *string
Objectstore *string
Recordstore *string
FVfailGsBucket *string
FVfailGsKey *string
AuthWebhookURL *string
OrchWebhookURL *string
OrchBlacklist *string
TestOrchAvail *bool
}
// DefaultLivepeerConfig creates LivepeerConfig exactly the same as when no flags are passed to the livepeer process.
func DefaultLivepeerConfig() LivepeerConfig {
// Network & Addresses:
defaultNetwork := "offchain"
defaultRtmpAddr := ""
defaultCliAddr := ""
defaultHttpAddr := ""
defaultServiceAddr := ""
defaultOrchAddr := ""
defaultVerifierURL := ""
defaultVerifierPath := ""
// Transcoding:
defaultOrchestrator := false
defaultTranscoder := false
defaultBroadcaster := false
defaultOrchSecret := ""
defaultTranscodingOptions := "P240p30fps16x9,P360p30fps16x9"
defaultMaxAttempts := 3
defaultSelectRandWeight := 0.3
defaultSelectStakeWeight := 0.7
defaultSelectPriceWeight := 0.0
defaultSelectPriceExpFactor := 100.0
defaultMaxSessions := strconv.Itoa(10)
defaultOrchPerfStatsURL := ""
defaultRegion := ""
defaultMinPerfScore := 0.0
defaultCurrentManifest := false
defaultNvidia := ""
defaultNetint := ""
defaultTestTranscoder := true
// AI:
defaultAIWorker := false
defaultAIModels := ""
defaultAIModelsDir := ""
// Onchain:
defaultEthAcctAddr := ""
defaultEthPassword := ""
defaultEthKeystorePath := ""
defaultEthOrchAddr := ""
defaultEthUrl := ""
defaultTxTimeout := 5 * time.Minute
defaultMaxTxReplacements := 1
defaultGasLimit := 0
defaultMaxGasPrice := 0
defaultEthController := ""
defaultInitializeRound := false
defaultTicketEV := "8000000000"
defaultMaxFaceValue := "0"
defaultMaxTicketEV := "3000000000000"
defaultMaxTotalEV := "20000000000000"
defaultDepositMultiplier := 1
defaultMaxPricePerUnit := 0
defaultPixelsPerUnit := 1
defaultAutoAdjustPrice := true
defaultPricePerBroadcaster := ""
defaultBlockPollingInterval := 5
defaultRedeemer := false
defaultRedeemerAddr := ""
defaultMonitor := false
defaultMetricsPerStream := false
defaultMetricsExposeClientIP := false
defaultMetadataQueueUri := ""
defaultMetadataAmqpExchange := "lp_golivepeer_metadata"
defaultMetadataPublishTimeout := 1 * time.Second
// Ingest:
defaultHttpIngest := true
// Verification:
defaultLocalVerify := true
// Storage:
defaultDatadir := ""
defaultObjectstore := ""
defaultRecordstore := ""
// Fast Verification GS bucket:
defaultFVfailGsBucket := ""
defaultFVfailGsKey := ""
// API
defaultAuthWebhookURL := ""
defaultOrchWebhookURL := ""
// Flags
defaultTestOrchAvail := true
return LivepeerConfig{
// Network & Addresses:
Network: &defaultNetwork,
RtmpAddr: &defaultRtmpAddr,
CliAddr: &defaultCliAddr,
HttpAddr: &defaultHttpAddr,
ServiceAddr: &defaultServiceAddr,
OrchAddr: &defaultOrchAddr,
VerifierURL: &defaultVerifierURL,
VerifierPath: &defaultVerifierPath,
// Transcoding:
Orchestrator: &defaultOrchestrator,
Transcoder: &defaultTranscoder,
Broadcaster: &defaultBroadcaster,
OrchSecret: &defaultOrchSecret,
TranscodingOptions: &defaultTranscodingOptions,
MaxAttempts: &defaultMaxAttempts,
SelectRandWeight: &defaultSelectRandWeight,
SelectStakeWeight: &defaultSelectStakeWeight,
SelectPriceWeight: &defaultSelectPriceWeight,
SelectPriceExpFactor: &defaultSelectPriceExpFactor,
MaxSessions: &defaultMaxSessions,
OrchPerfStatsURL: &defaultOrchPerfStatsURL,
Region: &defaultRegion,
MinPerfScore: &defaultMinPerfScore,
CurrentManifest: &defaultCurrentManifest,
Nvidia: &defaultNvidia,
Netint: &defaultNetint,
TestTranscoder: &defaultTestTranscoder,
// AI:
AIWorker: &defaultAIWorker,
AIModels: &defaultAIModels,
AIModelsDir: &defaultAIModelsDir,
// Onchain:
EthAcctAddr: &defaultEthAcctAddr,
EthPassword: &defaultEthPassword,
EthKeystorePath: &defaultEthKeystorePath,
EthOrchAddr: &defaultEthOrchAddr,
EthUrl: &defaultEthUrl,
TxTimeout: &defaultTxTimeout,
MaxTxReplacements: &defaultMaxTxReplacements,
GasLimit: &defaultGasLimit,
MaxGasPrice: &defaultMaxGasPrice,
EthController: &defaultEthController,
InitializeRound: &defaultInitializeRound,
TicketEV: &defaultTicketEV,
MaxFaceValue: &defaultMaxFaceValue,
MaxTicketEV: &defaultMaxTicketEV,
MaxTotalEV: &defaultMaxTotalEV,
DepositMultiplier: &defaultDepositMultiplier,
MaxPricePerUnit: &defaultMaxPricePerUnit,
PixelsPerUnit: &defaultPixelsPerUnit,
AutoAdjustPrice: &defaultAutoAdjustPrice,
PricePerBroadcaster: &defaultPricePerBroadcaster,
BlockPollingInterval: &defaultBlockPollingInterval,
Redeemer: &defaultRedeemer,
RedeemerAddr: &defaultRedeemerAddr,
Monitor: &defaultMonitor,
MetricsPerStream: &defaultMetricsPerStream,
MetricsExposeClientIP: &defaultMetricsExposeClientIP,
MetadataQueueUri: &defaultMetadataQueueUri,
MetadataAmqpExchange: &defaultMetadataAmqpExchange,
MetadataPublishTimeout: &defaultMetadataPublishTimeout,
// Ingest:
HttpIngest: &defaultHttpIngest,
// Verification:
LocalVerify: &defaultLocalVerify,
// Storage:
Datadir: &defaultDatadir,
Objectstore: &defaultObjectstore,
Recordstore: &defaultRecordstore,
// Fast Verification GS bucket:
FVfailGsBucket: &defaultFVfailGsBucket,
FVfailGsKey: &defaultFVfailGsKey,
// API
AuthWebhookURL: &defaultAuthWebhookURL,
OrchWebhookURL: &defaultOrchWebhookURL,
// Flags
TestOrchAvail: &defaultTestOrchAvail,
}
}
func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
if *cfg.MaxSessions == "auto" && *cfg.Orchestrator {
if *cfg.Transcoder {
glog.Exit("-maxSessions 'auto' cannot be used when both -orchestrator and -transcoder are specified")
}
core.MaxSessions = 0
} else {
intMaxSessions, err := strconv.Atoi(*cfg.MaxSessions)
if err != nil || intMaxSessions <= 0 {
glog.Exit("-maxSessions must be 'auto' or greater than zero")
}
core.MaxSessions = intMaxSessions
}
if *cfg.Netint != "" && *cfg.Nvidia != "" {
glog.Exit("both -netint and -nvidia arguments specified, this is not supported")
}
blockPollingTime := time.Duration(*cfg.BlockPollingInterval) * time.Second
type NetworkConfig struct {
ethController string
minGasPrice int64
redeemGas int
}
configOptions := map[string]*NetworkConfig{
"rinkeby": {
ethController: "0x9a9827455911a858E55f07911904fACC0D66027E",
redeemGas: redeemGasL1,
},
"arbitrum-one-rinkeby": {
ethController: "0x9ceC649179e2C7Ab91688271bcD09fb707b3E574",
redeemGas: redeemGasL2,
},
"mainnet": {
ethController: "0xf96d54e490317c557a967abfa5d6e33006be69b3",
minGasPrice: int64(params.GWei),
redeemGas: redeemGasL1,
},
"arbitrum-one-mainnet": {
ethController: "0xD8E8328501E9645d16Cf49539efC04f734606ee4",
redeemGas: redeemGasL2,
},
}
if *cfg.Network == "rinkeby" || *cfg.Network == "arbitrum-one-rinkeby" {
glog.Warning("The Rinkeby/ArbRinkeby networks are deprecated in favor of the Goerli/ArbGoerli networks which will be launched in January 2023.")
}
// If multiple orchAddr specified, ensure other necessary flags present and clean up list
orchURLs := parseOrchAddrs(*cfg.OrchAddr)
// Setting config options based on specified network
var redeemGas int
minGasPrice := int64(0)
if cfg.MinGasPrice != nil {
minGasPrice = *cfg.MinGasPrice
}
if netw, ok := configOptions[*cfg.Network]; ok {
if *cfg.EthController == "" {
*cfg.EthController = netw.ethController
}
if cfg.MinGasPrice == nil {
minGasPrice = netw.minGasPrice
}
redeemGas = netw.redeemGas
glog.Infof("***Livepeer is running on the %v network: %v***", *cfg.Network, *cfg.EthController)
} else {
redeemGas = redeemGasL1
glog.Infof("***Livepeer is running on the %v network***", *cfg.Network)
}
if *cfg.Datadir == "" {
homedir := os.Getenv("HOME")
if homedir == "" {
usr, err := user.Current()
if err != nil {
exit("Cannot find current user: %v", err)
}
homedir = usr.HomeDir
}
*cfg.Datadir = filepath.Join(homedir, ".lpData", *cfg.Network)
}
//Make sure datadir is present
if _, err := os.Stat(*cfg.Datadir); os.IsNotExist(err) {
glog.Infof("Creating data dir: %v", *cfg.Datadir)
if err = os.MkdirAll(*cfg.Datadir, 0755); err != nil {
glog.Errorf("Error creating datadir: %v", err)
}
}
//Set Gs bucket for fast verification fail case
if *cfg.FVfailGsBucket != "" && *cfg.FVfailGsKey != "" {
drivers.SetCreds(*cfg.FVfailGsBucket, *cfg.FVfailGsKey)
}
//Set up DB
dbh, err := common.InitDB(*cfg.Datadir + "/lpdb.sqlite3")
if err != nil {
glog.Errorf("Error opening DB: %v", err)
return
}
defer dbh.Close()
n, err := core.NewLivepeerNode(nil, *cfg.Datadir, dbh)
if err != nil {
glog.Errorf("Error creating livepeer node: %v", err)
}
if *cfg.OrchSecret != "" {
n.OrchSecret, _ = common.ReadFromFile(*cfg.OrchSecret)
}
var transcoderCaps []core.Capability
if *cfg.Transcoder {
core.WorkDir = *cfg.Datadir
accel := ffmpeg.Software
var devicesStr string
if *cfg.Nvidia != "" {
accel = ffmpeg.Nvidia
devicesStr = *cfg.Nvidia
}
if *cfg.Netint != "" {
accel = ffmpeg.Netint
devicesStr = *cfg.Netint
}
if accel != ffmpeg.Software {
accelName := ffmpeg.AccelerationNameLookup[accel]
tf, err := core.GetTranscoderFactoryByAccel(accel)
if err != nil {
exit("Error unsupported acceleration: %v", err)
}
// Get a list of device ids
devices, err := common.ParseAccelDevices(devicesStr, accel)
glog.Infof("%v devices: %v", accelName, devices)
if err != nil {
exit("Error while parsing '-%v %v' flag: %v", strings.ToLower(accelName), devices, err)
}
glog.Infof("Transcoding on these %v devices: %v", accelName, devices)
// Test transcoding with specified device
if *cfg.TestTranscoder {
transcoderCaps, err = core.TestTranscoderCapabilities(devices, tf)
if err != nil {
glog.Exit(err)
}
} else {
// no capability test was run, assume default capabilities
transcoderCaps = append(transcoderCaps, core.DefaultCapabilities()...)
}
// Initialize LB transcoder
n.Transcoder = core.NewLoadBalancingTranscoder(devices, tf)
} else {
// for local software mode, enable all capabilities
transcoderCaps = append(core.DefaultCapabilities(), core.OptionalCapabilities()...)
n.Transcoder = core.NewLocalTranscoder(*cfg.Datadir)
}
}
var aiCaps []core.Capability
constraints := make(map[core.Capability]*core.Constraints)
if *cfg.AIWorker {
gpus := []string{}
if *cfg.Nvidia != "" {
var err error
gpus, err = common.ParseAccelDevices(*cfg.Nvidia, ffmpeg.Nvidia)
if err != nil {
glog.Errorf("Error parsing -nvidia for devices: %v", err)
return
}
}
modelsDir := *cfg.AIModelsDir
if modelsDir == "" {
var err error
modelsDir, err = filepath.Abs(path.Join(*cfg.Datadir, "models"))
if err != nil {
glog.Error("Error creating absolute path for models dir: %v", modelsDir)
return
}
}
if err := os.MkdirAll(modelsDir, 0755); err != nil {
glog.Error("Error creating models dir %v", modelsDir)
return
}
n.AIWorker, err = worker.NewWorker(aiWorkerContainerImageID, gpus, modelsDir)
if err != nil {
glog.Errorf("Error starting AI worker: %v", err)
return
}
if *cfg.AIModels != "" {
configs, err := core.ParseAIModelConfigs(*cfg.AIModels)
if err != nil {
glog.Error("Error parsing -aiModels: %v", err)
return
}
for _, config := range configs {
modelConstraint := &core.ModelConstraint{Warm: config.Warm}
// If the config contains a URL we call Warm() anyway because AIWorker will just register
// the endpoint for an external container
if config.Warm || config.URL != "" {
endpoint := worker.RunnerEndpoint{URL: config.URL, Token: config.Token}
if err := n.AIWorker.Warm(ctx, config.Pipeline, config.ModelID, endpoint, config.OptimizationFlags); err != nil {
glog.Errorf("Error AI worker warming %v container: %v", config.Pipeline, err)
return
}
}
// Show warning if people set OptimizationFlags but not Warm.
if len(config.OptimizationFlags) > 0 && !config.Warm {
glog.Warningf("Model %v has 'optimization_flags' set without 'warm'. Optimization flags are currently only used for warm containers.", config.ModelID)
}
switch config.Pipeline {
case "text-to-image":
_, ok := constraints[core.Capability_TextToImage]
if !ok {
aiCaps = append(aiCaps, core.Capability_TextToImage)
constraints[core.Capability_TextToImage] = &core.Constraints{
Models: make(map[string]*core.ModelConstraint),
}
}
constraints[core.Capability_TextToImage].Models[config.ModelID] = modelConstraint
n.SetBasePriceForCap("default", core.Capability_TextToImage, config.ModelID, big.NewRat(config.PricePerUnit, config.PixelsPerUnit))
case "image-to-image":
_, ok := constraints[core.Capability_ImageToImage]
if !ok {
aiCaps = append(aiCaps, core.Capability_ImageToImage)
constraints[core.Capability_ImageToImage] = &core.Constraints{
Models: make(map[string]*core.ModelConstraint),
}
}
constraints[core.Capability_ImageToImage].Models[config.ModelID] = modelConstraint
n.SetBasePriceForCap("default", core.Capability_ImageToImage, config.ModelID, big.NewRat(config.PricePerUnit, config.PixelsPerUnit))
case "image-to-video":
_, ok := constraints[core.Capability_ImageToVideo]
if !ok {
aiCaps = append(aiCaps, core.Capability_ImageToVideo)
constraints[core.Capability_ImageToVideo] = &core.Constraints{
Models: make(map[string]*core.ModelConstraint),
}
}
constraints[core.Capability_ImageToVideo].Models[config.ModelID] = modelConstraint
n.SetBasePriceForCap("default", core.Capability_ImageToVideo, config.ModelID, big.NewRat(config.PricePerUnit, config.PixelsPerUnit))
}
}
}
defer func() {
ctx, cancel := context.WithTimeout(context.Background(), aiWorkerContainerStopTimeout)
defer cancel()
if err := n.AIWorker.Stop(ctx); err != nil {
glog.Errorf("Error stopping AI worker containers: %v", err)
return
}
glog.Infof("Stopped AI worker containers")
}()
}
if *cfg.Redeemer {
n.NodeType = core.RedeemerNode
} else if *cfg.Orchestrator {
n.NodeType = core.OrchestratorNode
if !*cfg.Transcoder {
n.TranscoderManager = core.NewRemoteTranscoderManager()
n.Transcoder = n.TranscoderManager
}
} else if *cfg.Transcoder {
n.NodeType = core.TranscoderNode
} else if *cfg.Broadcaster {
n.NodeType = core.BroadcasterNode
} else if (cfg.Reward == nil || !*cfg.Reward) && !*cfg.InitializeRound {
exit("No services enabled; must be at least one of -broadcaster, -transcoder, -orchestrator, -redeemer, -reward or -initializeRound")
}
lpmon.NodeID = *cfg.EthAcctAddr
if lpmon.NodeID != "" {
lpmon.NodeID += "-"
}
hn, _ := os.Hostname()
lpmon.NodeID += hn
if *cfg.Monitor {
if *cfg.MetricsExposeClientIP {
*cfg.MetricsPerStream = true
}
lpmon.Enabled = true
lpmon.PerStreamMetrics = *cfg.MetricsPerStream
lpmon.ExposeClientIP = *cfg.MetricsExposeClientIP
nodeType := lpmon.Default
switch n.NodeType {
case core.BroadcasterNode:
nodeType = lpmon.Broadcaster
case core.OrchestratorNode:
nodeType = lpmon.Orchestrator
case core.TranscoderNode:
nodeType = lpmon.Transcoder
case core.RedeemerNode:
nodeType = lpmon.Redeemer
}
lpmon.InitCensus(nodeType, core.LivepeerVersion)
}
watcherErr := make(chan error)
serviceErr := make(chan error)
var timeWatcher *watchers.TimeWatcher
if *cfg.Network == "offchain" {
glog.Infof("***Livepeer is in off-chain mode***")
if err := checkOrStoreChainID(dbh, big.NewInt(0)); err != nil {
glog.Error(err)
return
}
} else {
n.SelectionAlgorithm, err = createSelectionAlgorithm(cfg)
if err != nil {
exit("Incorrect parameters for selection algorithm, err=%v", err)
}
var keystoreDir = filepath.Join(*cfg.Datadir, "keystore")
keystoreInfo, err := parseEthKeystorePath(*cfg.EthKeystorePath)
if err == nil {
if keystoreInfo.path != "" {
keystoreDir = keystoreInfo.path
} else if (keystoreInfo.address != ethcommon.Address{}) {
ethKeystoreAddr := keystoreInfo.address.Hex()
ethAcctAddr := ethcommon.HexToAddress(*cfg.EthAcctAddr).Hex()
if (ethAcctAddr == ethcommon.Address{}.Hex()) || ethKeystoreAddr == ethAcctAddr {
*cfg.EthAcctAddr = ethKeystoreAddr
} else {
glog.Exit("-ethKeystorePath and -ethAcctAddr were both provided, but ethAcctAddr does not match the address found in keystore")
}
}
} else {
glog.Exit(fmt.Errorf(err.Error()))
}
//Get the Eth client connection information
if *cfg.EthUrl == "" {
glog.Exit("Need to specify an Ethereum node JSON-RPC URL using -ethUrl")
}
//Set up eth client
backend, err := ethclient.Dial(*cfg.EthUrl)
if err != nil {
glog.Errorf("Failed to connect to Ethereum client: %v", err)
return
}
chainID, err := backend.ChainID(ctx)
if err != nil {
glog.Errorf("failed to get chain ID from remote ethereum node: %v", err)
return
}
if !build.ChainSupported(chainID.Int64()) {
glog.Errorf("node does not support chainID = %v right now", chainID)
return
}
if err := checkOrStoreChainID(dbh, chainID); err != nil {
glog.Error(err)
return
}
var bigMaxGasPrice *big.Int
if *cfg.MaxGasPrice > 0 {
bigMaxGasPrice = big.NewInt(int64(*cfg.MaxGasPrice))
}
gpm := eth.NewGasPriceMonitor(backend, blockPollingTime, big.NewInt(minGasPrice), bigMaxGasPrice)
// Start gas price monitor
_, err = gpm.Start(ctx)
if err != nil {
glog.Errorf("Error starting gas price monitor: %v", err)
return
}
defer gpm.Stop()
am, err := eth.NewAccountManager(ethcommon.HexToAddress(*cfg.EthAcctAddr), keystoreDir, chainID, *cfg.EthPassword)
if err != nil {
glog.Errorf("Error creating Ethereum account manager: %v", err)
return
}
if err := am.Unlock(*cfg.EthPassword); err != nil {
glog.Errorf("Error unlocking Ethereum account: %v", err)
return
}
tm := eth.NewTransactionManager(backend, gpm, am, *cfg.TxTimeout, *cfg.MaxTxReplacements)
go tm.Start()
defer tm.Stop()
ethCfg := eth.LivepeerEthClientConfig{
AccountManager: am,
ControllerAddr: ethcommon.HexToAddress(*cfg.EthController),
EthClient: backend,
GasPriceMonitor: gpm,
TransactionManager: tm,
Signer: types.LatestSignerForChainID(chainID),
CheckTxTimeout: time.Duration(int64(*cfg.TxTimeout) * int64(*cfg.MaxTxReplacements+1)),
}
client, err := eth.NewClient(ethCfg)
if err != nil {
glog.Errorf("Failed to create Livepeer Ethereum client: %v", err)
return
}
if err := client.SetGasInfo(uint64(*cfg.GasLimit)); err != nil {
glog.Errorf("Failed to set gas info on Livepeer Ethereum Client: %v", err)
return
}
if err := client.SetMaxGasPrice(bigMaxGasPrice); err != nil {
glog.Errorf("Failed to set max gas price: %v", err)
return
}
n.Eth = client
addrMap := n.Eth.ContractAddresses()
// Initialize block watcher that will emit logs used by event watchers
blockWatcherClient, err := blockwatch.NewRPCClient(*cfg.EthUrl, ethRPCTimeout)
if err != nil {
glog.Errorf("Failed to setup blockwatch client: %v", err)
return
}
topics := watchers.FilterTopics()
blockWatcherCfg := blockwatch.Config{
Store: n.Database,
PollingInterval: blockPollingTime,
StartBlockDepth: rpc.LatestBlockNumber,
BlockRetentionLimit: blockWatcherRetentionLimit,
WithLogs: true,
Topics: topics,
Client: blockWatcherClient,
}
// Wait until all event watchers have been initialized before starting the block watcher
blockWatcher := blockwatch.New(blockWatcherCfg)
timeWatcher, err = watchers.NewTimeWatcher(addrMap["RoundsManager"], blockWatcher, n.Eth)
if err != nil {
glog.Errorf("Failed to setup roundswatcher: %v", err)
return
}
timeWatcherErr := make(chan error, 1)
go func() {
if err := timeWatcher.Watch(); err != nil {
timeWatcherErr <- fmt.Errorf("roundswatcher failed to start watching for events: %v", err)
}
}()
defer timeWatcher.Stop()
// Initialize unbonding watcher to update the DB with latest state of the node's unbonding locks
unbondingWatcher, err := watchers.NewUnbondingWatcher(n.Eth.Account().Address, addrMap["BondingManager"], blockWatcher, n.Database)
if err != nil {
glog.Errorf("Failed to setup unbonding watcher: %v", err)
return
}
// Start unbonding watcher (logs will not be received until the block watcher is started)
go unbondingWatcher.Watch()
defer unbondingWatcher.Stop()
senderWatcher, err := watchers.NewSenderWatcher(addrMap["TicketBroker"], blockWatcher, n.Eth, timeWatcher)
if err != nil {
glog.Errorf("Failed to setup senderwatcher: %v", err)
return
}
go senderWatcher.Watch()
defer senderWatcher.Stop()
orchWatcher, err := watchers.NewOrchestratorWatcher(addrMap["BondingManager"], blockWatcher, dbh, n.Eth, timeWatcher)
if err != nil {
glog.Errorf("Failed to setup orchestrator watcher: %v", err)
return
}
go orchWatcher.Watch()
defer orchWatcher.Stop()
serviceRegistryWatcher, err := watchers.NewServiceRegistryWatcher(addrMap["ServiceRegistry"], blockWatcher, dbh, n.Eth)
if err != nil {
glog.Errorf("Failed to set up service registry watcher: %v", err)
return
}
go serviceRegistryWatcher.Watch()
defer serviceRegistryWatcher.Stop()
n.Balances = core.NewAddressBalances(cleanupInterval)
defer n.Balances.StopCleanup()
// By default the ticket recipient is the node's address
// If the address of an on-chain registered orchestrator is provided, then it should be specified as the ticket recipient
recipientAddr := n.Eth.Account().Address
if *cfg.EthOrchAddr != "" {
recipientAddr = ethcommon.HexToAddress(*cfg.EthOrchAddr)
}
smCfg := &pm.LocalSenderMonitorConfig{
Claimant: recipientAddr,
CleanupInterval: cleanupInterval,
TTL: smTTL,
RedeemGas: redeemGas,
SuggestGasPrice: client.Backend().SuggestGasPrice,
RPCTimeout: ethRPCTimeout,
}
if *cfg.Orchestrator {
// Set price per pixel base info
if *cfg.PixelsPerUnit <= 0 {
// Can't divide by 0
panic(fmt.Errorf("-pixelsPerUnit must be > 0, provided %d", *cfg.PixelsPerUnit))
}
if cfg.PricePerUnit == nil {
// Prevent orchestrators from unknowingly providing free transcoding
panic(fmt.Errorf("-pricePerUnit must be set"))
}
if *cfg.PricePerUnit < 0 {
panic(fmt.Errorf("-pricePerUnit must be >= 0, provided %d", *cfg.PricePerUnit))
}
n.SetBasePrice("default", big.NewRat(int64(*cfg.PricePerUnit), int64(*cfg.PixelsPerUnit)))
glog.Infof("Price: %d wei for %d pixels\n ", *cfg.PricePerUnit, *cfg.PixelsPerUnit)
if *cfg.PricePerBroadcaster != "" {
ppb := getBroadcasterPrices(*cfg.PricePerBroadcaster)
for _, p := range ppb {
price := big.NewRat(p.PricePerUnit, p.PixelsPerUnit)
n.SetBasePrice(p.EthAddress, price)
glog.Infof("Price: %v set for broadcaster %v", price.RatString(), p.EthAddress)
}
}
n.AutoSessionLimit = *cfg.MaxSessions == "auto"
n.AutoAdjustPrice = *cfg.AutoAdjustPrice
ev, _ := new(big.Int).SetString(*cfg.TicketEV, 10)
if ev == nil {
glog.Errorf("-ticketEV must be a valid integer, but %v provided. Restart the node with a different valid value for -ticketEV", *cfg.TicketEV)
return
}
if ev.Cmp(big.NewInt(0)) < 0 {
glog.Errorf("-ticketEV must be greater than 0, but %v provided. Restart the node with a different valid value for -ticketEV", *cfg.TicketEV)
return
}
if err := setupOrchestrator(n, recipientAddr); err != nil {
glog.Errorf("Error setting up orchestrator: %v", err)
return
}
sigVerifier := &pm.DefaultSigVerifier{}
validator := pm.NewValidator(sigVerifier, timeWatcher)
var sm pm.SenderMonitor
if *cfg.RedeemerAddr != "" {
*cfg.RedeemerAddr = defaultAddr(*cfg.RedeemerAddr, "127.0.0.1", OrchestratorRpcPort)
rc, err := server.NewRedeemerClient(*cfg.RedeemerAddr, senderWatcher, timeWatcher)
if err != nil {
glog.Error("Unable to start redeemer client: ", err)
return
}
sm = rc
} else {
sm = pm.NewSenderMonitor(smCfg, n.Eth, senderWatcher, timeWatcher, n.Database)
}
// Start sender monitor
sm.Start()
defer sm.Stop()
tcfg := pm.TicketParamsConfig{
EV: ev,
RedeemGas: redeemGas,
TxCostMultiplier: txCostMultiplier,
}
n.Recipient, err = pm.NewRecipient(
recipientAddr,
n.Eth,
validator,
gpm,
sm,
timeWatcher,
tcfg,
)
if err != nil {
glog.Errorf("Error setting up PM recipient: %v", err)
return
}
mfv, _ := new(big.Int).SetString(*cfg.MaxFaceValue, 10)
if mfv == nil {
panic(fmt.Errorf("-maxFaceValue must be a valid integer, but %v provided. Restart the node with a different valid value for -maxFaceValue", *cfg.MaxFaceValue))
return
} else {
n.SetMaxFaceValue(mfv)
}
}
if n.NodeType == core.BroadcasterNode {
maxEV, _ := new(big.Rat).SetString(*cfg.MaxTicketEV)
if maxEV == nil {
panic(fmt.Errorf("-maxTicketEV must be a valid rational number, but %v provided. Restart the node with a valid value for -maxTicketEV", *cfg.MaxTicketEV))
}
if maxEV.Cmp(big.NewRat(0, 1)) < 0 {
panic(fmt.Errorf("-maxTicketEV must not be negative, but %v provided. Restart the node with a valid value for -maxTicketEV", *cfg.MaxTicketEV))
}
maxTotalEV, _ := new(big.Rat).SetString(*cfg.MaxTotalEV)
if maxTotalEV.Cmp(big.NewRat(0, 1)) < 0 {
panic(fmt.Errorf("-maxTotalEV must not be negative, but %v provided. Restart the node with a valid value for -maxTotalEV", *cfg.MaxTotalEV))
}
if *cfg.DepositMultiplier <= 0 {
panic(fmt.Errorf("-depositMultiplier must be greater than 0, but %v provided. Restart the node with a valid value for -depositMultiplier", *cfg.DepositMultiplier))
}
// Fetch and cache broadcaster on-chain info
info, err := senderWatcher.GetSenderInfo(n.Eth.Account().Address)
if err != nil {
glog.Error("Failed to get broadcaster on-chain info: ", err)
return
}
glog.Info("Broadcaster Deposit: ", eth.FormatUnits(info.Deposit, "ETH"))
glog.Info("Broadcaster Reserve: ", eth.FormatUnits(info.Reserve.FundsRemaining, "ETH"))
n.Sender = pm.NewSender(n.Eth, timeWatcher, senderWatcher, maxEV, maxTotalEV, *cfg.DepositMultiplier)
if *cfg.PixelsPerUnit <= 0 {
// Can't divide by 0
panic(fmt.Errorf("The amount of pixels per unit must be greater than 0, provided %d instead\n", *cfg.PixelsPerUnit))
}
if *cfg.MaxPricePerUnit > 0 {
server.BroadcastCfg.SetMaxPrice(big.NewRat(int64(*cfg.MaxPricePerUnit), int64(*cfg.PixelsPerUnit)))
} else {
glog.Infof("Maximum transcoding price per pixel is not greater than 0: %v, broadcaster is currently set to accept ANY price.\n", *cfg.MaxPricePerUnit)
glog.Infoln("To update the broadcaster's maximum acceptable transcoding price per pixel, use the CLI or restart the broadcaster with the appropriate 'maxPricePerUnit' and 'pixelsPerUnit' values")
}
}
if n.NodeType == core.RedeemerNode {
if err := setupOrchestrator(n, recipientAddr); err != nil {
glog.Errorf("Error setting up orchestrator: %v", err)
return
}