Skip to content
Merged
Show file tree
Hide file tree
Changes from 101 commits
Commits
Show all changes
124 commits
Select commit Hold shift + click to select a range
bb51ef7
Serverless
Mar 4, 2026
60933e3
TEMP
j0sh Mar 10, 2026
f495713
ai/scope: Take WebSocket URL and capacity from config file
j0sh Mar 17, 2026
e279b37
Create trickle streams on demand via websocket
j0sh Mar 18, 2026
b023f84
Temp LV2V hacks for Scope:
j0sh Mar 19, 2026
f26e354
Take websocket URL from params
j0sh Mar 19, 2026
977d660
Don't print event contents; they are large
j0sh Mar 20, 2026
75613d7
fixup! Take websocket URL from params
j0sh Mar 20, 2026
41e3e6c
Extend handshake to a `started` message; this would include auth
j0sh Mar 21, 2026
3e4134d
Canonicalize websocket URLs to lowercase
j0sh Mar 25, 2026
6bbce26
Retry event loop on errors
j0sh Mar 26, 2026
d541137
Close control/event channels on websocket termination.
j0sh Mar 26, 2026
c326b44
Handle websocket retries.
j0sh Mar 26, 2026
d38b897
trickle: Support publisher resets.
j0sh Mar 26, 2026
8da1c69
trickle: Add a /<channel>/next endpoint.
j0sh Mar 26, 2026
04ac4c5
Return a 401 for auth errors during the websocket handshake.
j0sh Mar 31, 2026
4f9535f
signer: Add an auth webhook callback.
j0sh Apr 7, 2026
5fc1345
signer: Cache webhook auth to skip redundant calls.
j0sh Apr 7, 2026
3ada704
ai/scope: Add Scope specific orchestrator endpoint.
j0sh Apr 7, 2026
4721c3b
Use control URL for base path extraction
j0sh Apr 8, 2026
92abb16
Use application-level ping-pong messages.
j0sh Apr 8, 2026
cccadb5
go fmt
j0sh Apr 8, 2026
0da6e61
Use full trickle channel URL as map key.
j0sh Apr 8, 2026
13a73ab
Handle interleaved pongs during handshake retries.
j0sh Apr 9, 2026
aa6c987
Revert "Handle interleaved pongs during handshake retries."
j0sh Apr 9, 2026
cacee23
Stop ping/pong during runner restarts.
j0sh Apr 9, 2026
fa35b2a
Wrap websocket
j0sh Apr 9, 2026
610ee3b
Configurable serverless timeout
j0sh Apr 10, 2026
1882716
trickle: Optimize fanout buffering
j0sh Apr 14, 2026
7b71171
feat(signer): restore /sign-byoc-job endpoint for BYOC inference
seanhanca Apr 11, 2026
656f8b6
Separate transport and business-logic auth errors
j0sh Apr 15, 2026
47c2a2b
ai/live: Add request headers for remote signer
j0sh Apr 15, 2026
3a07ed3
Return 400 on invalid ws url
j0sh Apr 15, 2026
efa5259
Add a notification + grace period for timeouts.
j0sh Apr 17, 2026
fb5dccb
Set ManifestID / DaydreamUserId from headers
j0sh Apr 19, 2026
18c28fa
trickle: Return 200 instead of 470 on empty segments.
j0sh Apr 23, 2026
114d610
Use clog instead of slog
j0sh Apr 23, 2026
4166548
Initial Live Runner implementation
j0sh May 1, 2026
cf1f8d3
Serverless worker as a synthetic LiveRunner
j0sh May 5, 2026
72eb806
Price handling tweaks, various simplifications
j0sh May 6, 2026
0b5a4c8
Rename to UseLiveRunners, enforce OrchSecret
j0sh May 6, 2026
47cb618
Move price conversion logic into runner itself
j0sh May 6, 2026
d3ac6fc
Add callback for trickle channel create / delete
j0sh May 6, 2026
b60016e
Add orchestrator URL to runner response
j0sh May 7, 2026
6d16641
Add live runner session endpoints
j0sh May 7, 2026
23e4839
Add create / delete trickle channel API to runner
j0sh May 7, 2026
57238ac
Separate creds for runner bootstrap, heartbeat, session
j0sh May 9, 2026
3c4ac32
Include runner ID with incoming requests
j0sh May 11, 2026
2d9d5a3
Override trickle base with configured url if necessary
j0sh May 11, 2026
405947c
Add runner registration from JSON config
j0sh May 12, 2026
a340df1
Add no-SSL mode via -httpAddr
j0sh May 12, 2026
889a44e
Runner single-shot mode
j0sh May 12, 2026
e78d14d
Support path-only static runner health URLs
j0sh May 12, 2026
0050d79
Offchain mode for live runners
j0sh May 12, 2026
9a4ad1b
Add single-shot alias
j0sh May 12, 2026
d1bec77
Update discovery price info
j0sh May 12, 2026
e1d994d
Add label-based routing
j0sh May 12, 2026
e8aeef5
Add Session-Control header for session control from internal runners
j0sh May 12, 2026
d74671e
Comments and wording tweaks
j0sh May 12, 2026
2720bce
Add runner payment support.
j0sh May 14, 2026
b3c73fa
Add payment refresh endpoint
j0sh May 14, 2026
ef1fb60
Add orch URL header to refresh message
j0sh May 14, 2026
2d79a35
Add payment challenge to /scope endpoint
j0sh May 14, 2026
e88988a
Remove Scope carveouts from LV2V
j0sh May 14, 2026
deaee30
Remove extraneous trickle benchmarks, fix test
j0sh May 14, 2026
9a59f26
Hide remote signer auth headers
j0sh May 14, 2026
cded6e3
Use manifest ID as session ID
j0sh May 15, 2026
860308f
Limit request payloads
j0sh May 15, 2026
0a6f1c9
Work in offchain mode with latest Scope
j0sh May 15, 2026
aa995c1
Add capacity to discovery
j0sh May 15, 2026
3cc3b4c
Add runners to signer discovery
j0sh May 18, 2026
5c77ba1
Granular per-runner locking
j0sh May 19, 2026
c3c93ef
Consolidate heartbeat checks
j0sh May 19, 2026
7fb1f62
Fixup tests
j0sh May 19, 2026
ffd643a
Add payment monitoring
j0sh May 20, 2026
41df986
Add r2o / o2r trickle channels
j0sh May 20, 2026
c2a374d
trickle: Don't autocreate local writes unless configured.
j0sh May 20, 2026
a358c9d
Add o2r keepalive
j0sh May 20, 2026
696d321
Remove r2o trickle channel
j0sh May 20, 2026
dfb5802
Add sessions to heartbeat
j0sh May 21, 2026
d2f3745
Add internal session-stop endpoint for live runner
j0sh May 26, 2026
6b4393d
Allow remote discovery from runners
j0sh May 27, 2026
e13ac6d
Split live runner callback URI from public service URI
j0sh May 28, 2026
bf8f95a
Fix live runner trickle public URL
j0sh May 29, 2026
d9ee753
Add non-grpc selfcheck using discovery
j0sh Jun 2, 2026
0de5e7c
Fix pricing
j0sh Jun 9, 2026
1f227c8
Fix serverless Scope payments with live runner
j0sh Jun 9, 2026
a7cbaef
Make header parsing more resilient
j0sh Jun 9, 2026
5ae964c
Fix payment refreshes
j0sh Jun 10, 2026
e8520f6
Fix legacy Scope clients.
j0sh Jun 10, 2026
c0e79cc
reject tickets with zero expiration block
j0sh Jun 10, 2026
e545b98
Add session-scoped payment endpoint
j0sh Jun 25, 2026
81f24cf
Add control URL to client session response
j0sh Jun 25, 2026
6033eb0
Merge branch 'master' into ja/live-runner
j0sh Jun 26, 2026
76c9c63
add serverless worker tests
j0sh Jun 26, 2026
64cf17c
fix discovery tests
j0sh Jun 26, 2026
e4f29f9
Merge branch 'master' into ja/live-runner
j0sh Jun 29, 2026
37f7bf6
fix(ai/live-runner): derive Session-Control header from liveRunnerAdd…
rickstaa Jun 29, 2026
3df2338
Merge branch 'master' into ja/live-runner
j0sh Jun 30, 2026
e645608
Set correct JSON mome type for O2R trickle channel.
j0sh Jun 30, 2026
aa81fd5
Fix mime type tests
j0sh Jul 1, 2026
4c2f623
feat(ai/runner): log live runner register and deregister events (#3943)
rickstaa Jul 9, 2026
25ade72
fix(remote-signer): keep retrying discovery while the cache is empty …
rickstaa Jul 9, 2026
c2db3cb
Add proxy URLs
j0sh Jul 13, 2026
0cff1bd
Stop filtering remote discovery runners by capability
j0sh Jul 15, 2026
004450d
Add 'live' pricing type
j0sh Jul 15, 2026
2a09807
s/Route/Routing/g
j0sh Jul 17, 2026
fe6c786
Add live runner docs
j0sh Jul 17, 2026
a03a0e8
Generalize live payment units
j0sh Jul 20, 2026
4a95c4f
Rename LV2V payment processor
j0sh Jul 20, 2026
a38177d
Add seconds-based live payment processor
j0sh Jul 20, 2026
253fef3
Minimize payment processor unit changes
j0sh Jul 20, 2026
a4dade6
Keep payment processor callback naming
j0sh Jul 20, 2026
884b277
Carry fractional live payment units
j0sh Jul 20, 2026
8273db9
Revert 7b71171d3f9b8dd94f8ecfd3edb97d018f2ed1b8
j0sh Jul 20, 2026
67dc4cd
Merge remote-tracking branch 'origin/master' into ja/live-runner
j0sh Jul 20, 2026
ca09ddf
Add metadata passthrough
j0sh Jul 21, 2026
b90c6b3
Support empty proxy target_url
j0sh Jul 21, 2026
7727bbe
Support single-shot proxy URLs
j0sh Jul 21, 2026
3800923
Fix proxy URL with port
j0sh Jul 21, 2026
802b103
Allow SDK and persistent default proxies
j0sh Jul 22, 2026
73142cd
early exit on invalid payment type
j0sh Jul 22, 2026
d258abb
Fix live payment signer tests
j0sh Jul 22, 2026
32c2352
Log on session open / close
j0sh Jul 22, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,397 changes: 1,397 additions & 0 deletions ai/runner/live_runner.go

Large diffs are not rendered by default.

1,744 changes: 1,744 additions & 0 deletions ai/runner/live_runner_test.go

Large diffs are not rendered by default.

975 changes: 975 additions & 0 deletions ai/worker/serverless_worker.go

Large diffs are not rendered by default.

414 changes: 414 additions & 0 deletions ai/worker/serverless_worker_test.go

Large diffs are not rendered by default.

63 changes: 63 additions & 0 deletions byoc/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package byoc
import (
"context"
"crypto/tls"
"encoding/binary"
"errors"
"math/big"
gonet "net"
Expand Down Expand Up @@ -279,3 +280,65 @@ type byocLiveRequestParams struct {
// when the write for the last segment started
lastSegmentTime time.Time
}

// Prevents cross-protocol signature replay.
const BYOCJobSigV1Prefix = "LP_BYOC_JOB_V1\x00\x00"

// BYOCJobSigningInput holds the fields that are bound into a BYOC job signature.
Comment thread
rickstaa marked this conversation as resolved.
Outdated
type BYOCJobSigningInput struct {
ID string
Capability string
Request string
Parameters string
TimeoutSeconds int
}

// FlattenBYOCJob produces a deterministic binary representation of a BYOC job
// for signing, similar to SegTranscodingMetadata.Flatten() used by LV2V.
//
// Wire format:
//
// version(16) || timeout(4,BE) || len(id)(4,BE) || id || len(cap)(4,BE) || cap
// || len(req)(4,BE) || req || len(params)(4,BE) || params
func FlattenBYOCJob(job *BYOCJobSigningInput) []byte {
idBytes := []byte(job.ID)
capBytes := []byte(job.Capability)
reqBytes := []byte(job.Request)
paramsBytes := []byte(job.Parameters)

size := 16 + 4 +
4 + len(idBytes) +
4 + len(capBytes) +
4 + len(reqBytes) +
4 + len(paramsBytes)

buf := make([]byte, size)
offset := 0

copy(buf[offset:], []byte(BYOCJobSigV1Prefix))
offset += 16

binary.BigEndian.PutUint32(buf[offset:], uint32(job.TimeoutSeconds))
offset += 4

binary.BigEndian.PutUint32(buf[offset:], uint32(len(idBytes)))
offset += 4
copy(buf[offset:], idBytes)
offset += len(idBytes)

binary.BigEndian.PutUint32(buf[offset:], uint32(len(capBytes)))
offset += 4
copy(buf[offset:], capBytes)
offset += len(capBytes)

binary.BigEndian.PutUint32(buf[offset:], uint32(len(reqBytes)))
offset += 4
copy(buf[offset:], reqBytes)
offset += len(reqBytes)

binary.BigEndian.PutUint32(buf[offset:], uint32(len(paramsBytes)))
offset += 4
copy(buf[offset:], paramsBytes)

return buf
}
19 changes: 18 additions & 1 deletion clog/clog.go
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,23 @@ func infof(ctx context.Context, lastErr bool, publicLog bool, format string, arg
// Example: Info(ctx, "hello", "key1", value1, "key2", value2)
// This will log: "hello key1=value1 key2=value2"
func Info(ctx context.Context, msg string, keyvals ...interface{}) {
msg = formatKeyvals(msg, keyvals...)
infof(ctx, false, false, msg)
}

// Warning logs a warning message with key-value pairs in a slog-like style.
func Warning(ctx context.Context, msg string, keyvals ...interface{}) {
msg = formatKeyvals(msg, keyvals...)
Warningf(ctx, "%s", msg)
}

// Error logs an error message with key-value pairs in a slog-like style.
func Error(ctx context.Context, msg string, keyvals ...interface{}) {
msg = formatKeyvals(msg, keyvals...)
Errorf(ctx, "%s", msg)
}

func formatKeyvals(msg string, keyvals ...interface{}) string {
if len(keyvals)%2 != 0 {
keyvals = append(keyvals[:len(keyvals)-1], "MISSING", keyvals[len(keyvals)-1])
}
Expand Down Expand Up @@ -225,7 +242,7 @@ func Info(ctx context.Context, msg string, keyvals ...interface{}) {
}
}

infof(ctx, false, false, sb.String())
return sb.String()
}

// V returns a Verbose instance for conditional logging at the specified level
Expand Down
4 changes: 4 additions & 0 deletions cmd/livepeer/starter/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,9 @@ func NewLivepeerConfig(fs *flag.FlagSet) LivepeerConfig {
// AI:
cfg.AIServiceRegistry = fs.Bool("aiServiceRegistry", *cfg.AIServiceRegistry, "Set to true to use an AI ServiceRegistry contract address")
cfg.AIWorker = fs.Bool("aiWorker", *cfg.AIWorker, "Set to true to run an AI worker")
cfg.AIServerless = fs.Bool("aiServerless", *cfg.AIServerless, "Set to true to use serverless AI worker")
cfg.UseLiveRunners = fs.Bool("useLiveRunners", *cfg.UseLiveRunners, "Set to true to use LiveRunner-backed runners for supported pipelines")
cfg.LiveRunnerConfig = fs.String("liveRunnerConfig", *cfg.LiveRunnerConfig, "Path to JSON config file for static LiveRunner registration")
cfg.AIModels = fs.String("aiModels", *cfg.AIModels, "Set models (pipeline:model_id) for AI worker to load upon initialization")
cfg.AIModelsDir = fs.String("aiModelsDir", *cfg.AIModelsDir, "Set directory where AI model weights are stored")
cfg.AIRunnerImage = fs.String("aiRunnerImage", *cfg.AIRunnerImage, "[Deprecated] Specify the base Docker image for the AI runner. Example: livepeer/ai-runner:0.0.1. Use -aiRunnerImageOverrides instead.")
Expand All @@ -67,6 +70,7 @@ func NewLivepeerConfig(fs *flag.FlagSet) LivepeerConfig {

// Live AI:
cfg.MediaMTXApiPassword = fs.String("mediaMTXApiPassword", "", "HTTP basic auth password for MediaMTX API requests")
cfg.LiveRunnerAddr = fs.String("liveRunnerAddr", *cfg.LiveRunnerAddr, "Base URL used by live runners for heartbeat, control-plane callbacks, and internal trickle channels. Must be a full URL such as http://go-livepeer:8935")
cfg.LiveAITrickleHostForRunner = fs.String("liveAITrickleHostForRunner", "", "Trickle Host used by AI Runner; It's used to overwrite the publicly available Trickle Host")
cfg.LiveAIAuthApiKey = fs.String("liveAIAuthApiKey", "", "API key to use for Live AI authentication requests")
cfg.LiveAIHeartbeatURL = fs.String("liveAIHeartbeatURL", "", "Base URL for Live AI heartbeat requests")
Expand Down
140 changes: 125 additions & 15 deletions cmd/livepeer/starter/starter.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import (
"github.com/ethereum/go-ethereum/params"
"github.com/ethereum/go-ethereum/rpc"
"github.com/golang/glog"
"github.com/livepeer/go-livepeer/ai/runner"
"github.com/livepeer/go-livepeer/ai/worker"
"github.com/livepeer/go-livepeer/build"
"github.com/livepeer/go-livepeer/common"
Expand Down Expand Up @@ -97,6 +98,9 @@ type LivepeerConfig struct {
Transcoder *bool
AIServiceRegistry *bool
AIWorker *bool
AIServerless *bool
UseLiveRunners *bool
LiveRunnerConfig *string
Gateway *bool
Broadcaster *bool
OrchSecret *string
Expand Down Expand Up @@ -164,6 +168,7 @@ type LivepeerConfig struct {
AuthWebhookURL *string
LiveAIAuthWebhookURL *string
LiveAITrickleHostForRunner *string
LiveRunnerAddr *string
OrchWebhookURL *string
OrchBlacklist *string
OrchMinLivepeerVersion *string
Expand Down Expand Up @@ -237,6 +242,9 @@ func DefaultLivepeerConfig() LivepeerConfig {
// AI:
defaultAIServiceRegistry := false
defaultAIWorker := false
defaultAIServerless := false
defaultUseLiveRunners := false
defaultLiveRunnerConfig := ""
defaultAIModels := ""
defaultAIModelsDir := ""
defaultAIRunnerImage := "livepeer/ai-runner:latest"
Expand All @@ -246,6 +254,7 @@ func DefaultLivepeerConfig() LivepeerConfig {
defaultAIMinRunnerVersion := "[]"
defaultAIRunnerImageOverrides := ""
defaultLiveAIAuthWebhookURL := ""
defaultLiveRunnerAddr := ""
defaultLivePaymentInterval := 5 * time.Second
defaultLiveOutSegmentTimeout := 0 * time.Second
defaultGatewayHost := ""
Expand Down Expand Up @@ -363,6 +372,9 @@ func DefaultLivepeerConfig() LivepeerConfig {
// AI:
AIServiceRegistry: &defaultAIServiceRegistry,
AIWorker: &defaultAIWorker,
AIServerless: &defaultAIServerless,
UseLiveRunners: &defaultUseLiveRunners,
LiveRunnerConfig: &defaultLiveRunnerConfig,
AIModels: &defaultAIModels,
AIModelsDir: &defaultAIModelsDir,
AIRunnerImage: &defaultAIRunnerImage,
Expand All @@ -372,6 +384,7 @@ func DefaultLivepeerConfig() LivepeerConfig {
AIMinRunnerVersion: &defaultAIMinRunnerVersion,
AIRunnerImageOverrides: &defaultAIRunnerImageOverrides,
LiveAIAuthWebhookURL: &defaultLiveAIAuthWebhookURL,
LiveRunnerAddr: &defaultLiveRunnerAddr,
LivePaymentInterval: &defaultLivePaymentInterval,
LiveOutSegmentTimeout: &defaultLiveOutSegmentTimeout,
GatewayHost: &defaultGatewayHost,
Expand Down Expand Up @@ -1319,6 +1332,15 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {

var aiCaps []core.Capability
capabilityConstraints := make(core.PerCapabilityConstraints)
var aiModelConfigs []core.AIModelConfig

if *cfg.AIModels != "" {
aiModelConfigs, err = core.ParseAIModelConfigs(*cfg.AIModels)
if err != nil {
glog.Errorf("Error parsing -aiModels: %v", err)
return
}
}

if *cfg.AIWorker {
gpus := []string{}
Expand Down Expand Up @@ -1378,10 +1400,28 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
}
}

n.AIWorker, err = worker.NewWorker(imageOverrides, *cfg.AIVerboseLogs, gpus, modelsDir, containerCreatorID)
if err != nil {
glog.Errorf("Error starting AI worker: %v", err)
return
if *cfg.AIServerless {
if len(aiModelConfigs) != 1 {
glog.Errorf("Serverless requires exactly one AI model config, got %d", len(aiModelConfigs))
return
}
config := aiModelConfigs[0]
if config.Pipeline != "live-video-to-video" || config.ModelID != "scope" {
glog.Errorf("Serverless only supports live-video-to-video/scope, got %s/%s", config.Pipeline, config.ModelID)
return
}

n.AIWorker, err = worker.NewServerlessWorker(strings.TrimSpace(config.URL), config.Capacity)
if err != nil {
glog.Errorf("Error starting Serverless AI worker: %v", err)
return
}
} else {
n.AIWorker, err = worker.NewWorker(imageOverrides, *cfg.AIVerboseLogs, gpus, modelsDir, containerCreatorID)
if err != nil {
glog.Errorf("Error starting AI worker: %v", err)
return
}
}

defer func() {
Expand All @@ -1397,20 +1437,15 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
}

if *cfg.AIModels != "" {
configs, err := core.ParseAIModelConfigs(*cfg.AIModels)
if err != nil {
glog.Errorf("Error parsing -aiModels: %v", err)
return
}

for _, config := range configs {
for _, config := range aiModelConfigs {
pipelineCap, err := core.PipelineToCapability(config.Pipeline)
if err != nil {
panic(fmt.Errorf("Pipeline is not valid capability: %v\n", config.Pipeline))
}
if *cfg.AIWorker {
modelConstraint := &core.ModelConstraint{Warm: config.Warm, Capacity: 1}
modelsCount := 1

if config.Capacity != 0 {
if config.URL == "" {
// Use multiple same configs if External Container is not used and capacity is set
Expand Down Expand Up @@ -1854,6 +1889,12 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
if cfg.LiveAITrickleHostForRunner != nil {
n.LiveAITrickleHostForRunner = *cfg.LiveAITrickleHostForRunner
}
if cfg.LiveRunnerAddr != nil && *cfg.LiveRunnerAddr != "" {
n.LiveRunnerAddr, err = parseLiveRunnerAddr(*cfg.LiveRunnerAddr)
if err != nil {
glog.Exitf("invalid -liveRunnerAddr: %v", err)
}
}
if cfg.LiveAICapRefreshModels != nil && *cfg.LiveAICapRefreshModels != "" {
glog.Warningf("The -liveAICapRefreshModels flag is deprecated, capacity is now available for all models, use -liveAICapReportInterval to set the interval for reporting capacity metrics")
}
Expand Down Expand Up @@ -1935,6 +1976,7 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
DiscoveryTimeout: *cfg.DiscoveryTimeout,
LiveAICapReportInterval: *cfg.LiveAICapReportInterval,
IgnoreCapacityCheck: true,
UseDiscoveryEndpoint: true,
}.New()
if err != nil {
exit("Could not create orchestrator pool with DB cache: %v", err)
Expand Down Expand Up @@ -1984,6 +2026,30 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {

orch := core.NewOrchestrator(s.LivepeerNode, timeWatcher)

if *cfg.UseLiveRunners || *cfg.LiveRunnerConfig != "" {
if n.OrchSecret == "" && *cfg.LiveRunnerConfig == "" {
glog.Exit("running with -useLiveRunners requires -orchSecret")
}
n.LiveRunnerManager = runner.NewLiveRunnerRegistry(runner.LiveRunnerRegistryConfig{
Host: liveRunnerHost{RunnerHost: orch, LivepeerNode: n},
Onchain: *cfg.Network != "offchain",
})
if n.OrchSecret == "" {
glog.Warning("No -orchSecret configured; dynamic LiveRunner heartbeat registration is disabled")
}
if *cfg.LiveRunnerConfig != "" {
configJSON, err := os.ReadFile(*cfg.LiveRunnerConfig)
if err != nil {
glog.Exitf("error reading -liveRunnerConfig: %v", err)
}
registration, err := n.LiveRunnerManager.(*runner.LiveRunnerRegistry).RegisterStaticRunnersJSON(configJSON)
if err != nil {
glog.Exitf("error registering -liveRunnerConfig: %v", err)
}
glog.Infof("Registered %d static live runners from %s", len(registration.Runners), *cfg.LiveRunnerConfig)
}
}

go func() {
err = server.StartTranscodeServer(orch, *cfg.HttpAddr, s.HTTPMux, n.WorkDir, n.TranscoderManager != nil, n.AIWorkerManager != nil, n)
if err != nil {
Expand All @@ -1997,10 +2063,20 @@ func StartLivepeer(ctx context.Context, cfg LivepeerConfig) {
// check whether or not the orchestrator is available
if *cfg.TestOrchAvail && doingWork {
time.Sleep(2 * time.Second)
orchAvail := server.CheckOrchestratorAvailability(orch)
var (
checkName string
orchAvail bool
)
if *cfg.UseLiveRunners || *cfg.LiveRunnerConfig != "" {
checkName = "discovery"
orchAvail = server.CheckOrchestratorDiscoveryAvailability(orch)
} else {
checkName = "grpc ping"
orchAvail = server.CheckOrchestratorAvailability(orch)
}
if !orchAvail {
// shut down orchestrator
glog.Infof("Orchestrator not available at %v (%v); shutting down", orch.ServiceURI(), *cfg.HttpAddr)
glog.Infof("Orchestrator not available at %v (%v) via %s check; shutting down", orch.ServiceURI(), *cfg.HttpAddr, checkName)
tc <- struct{}{}
}
}
Expand Down Expand Up @@ -2143,7 +2219,11 @@ func parseHeaderMap(raw string) map[string]string {
for _, header := range strings.Split(raw, ",") {
parts := strings.SplitN(header, ":", 2)
if len(parts) == 2 {
headers[parts[0]] = parts[1]
key := strings.TrimSpace(parts[0])
value := strings.TrimSpace(parts[1])
if key != "" {
headers[key] = value
}
}
}
return headers
Expand Down Expand Up @@ -2177,7 +2257,10 @@ func getServiceURI(n *core.LivepeerNode, serviceAddr string) (*url.URL, error) {
// special value to signal this node is not to be used for work
return url.Parse("")
}
return url.ParseRequestURI("https://" + serviceAddr)
if !strings.HasPrefix(serviceAddr, "http://") && !strings.HasPrefix(serviceAddr, "https://") {
serviceAddr = "https://" + serviceAddr
}
return url.ParseRequestURI(serviceAddr)
}

// Infer address
Expand Down Expand Up @@ -2504,3 +2587,30 @@ func exit(msg string, args ...any) {
glog.Errorf(msg, args...)
os.Exit(2)
}

type liveRunnerHost struct {
runner.RunnerHost
*core.LivepeerNode
}

func (h liveRunnerHost) LiveRunnerURI() *url.URL {
if h.LivepeerNode != nil && h.LivepeerNode.LiveRunnerAddr != nil {
v := *h.LivepeerNode.LiveRunnerAddr
return &v
}
return h.RunnerHost.ServiceURI()
}

func parseLiveRunnerAddr(addr string) (*url.URL, error) {
parsed, err := url.ParseRequestURI(addr)
if err != nil {
return nil, err
}
if !parsed.IsAbs() || parsed.Host == "" {
return nil, fmt.Errorf("must be an absolute URL")
}
if parsed.Scheme != "http" && parsed.Scheme != "https" {
return nil, fmt.Errorf("scheme must be http or https")
}
return parsed, nil
}
Loading
Loading