Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

ai/live: Simplify payments API around readers #3447

Draft
wants to merge 1 commit into
base: master
Choose a base branch
from
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
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
4 changes: 1 addition & 3 deletions server/ai_http.go
Original file line number Diff line number Diff line change
Expand Up @@ -214,11 +214,9 @@ func (h *lphttp) StartLiveVideoToVideo() http.Handler {
clog.Infof(ctx, "Error getting local trickle segment err=%v", err)
return
}
reader := segment.Reader
if paymentProcessor != nil {
reader = paymentProcessor.process(ctx, segment.Reader)
paymentProcessor.process(ctx, segment.Reader)
}
io.Copy(io.Discard, reader)
}
}()

Expand Down
2 changes: 1 addition & 1 deletion server/ai_live_video.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ func startTricklePublish(ctx context.Context, url *url.URL, params aiRequestPara
defer slowOrchChecker.EndSegment()
var r io.Reader = reader
if paymentProcessor != nil {
r = paymentProcessor.process(ctx, reader)
paymentProcessor.process(ctx, reader)
}

clog.V(8).Infof(ctx, "trickle publish writing data seq=%d", seq)
Expand Down
17 changes: 4 additions & 13 deletions server/live_payment_processor.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"time"

"github.com/livepeer/go-livepeer/clog"
"github.com/livepeer/go-livepeer/media"
"github.com/livepeer/lpms/ffmpeg"
)

Expand Down Expand Up @@ -103,20 +104,13 @@ func (p *LivePaymentProcessor) processOne(ctx context.Context, timestamp time.Ti
p.lastProcessedAt = timestamp
}

func (p *LivePaymentProcessor) process(ctx context.Context, reader io.Reader) io.Reader {
func (p *LivePaymentProcessor) process(ctx context.Context, reader media.CloneableReader) {
timestamp := time.Now()
if p.shouldSkip(timestamp) {
// We don't process every segment, because it's too compute-expensive
return reader
}

pipeReader, pipeWriter, err := os.Pipe()
if err != nil {
clog.InfofErr(ctx, "Error creating pipe", err)
return reader
return
}

resReader := io.TeeReader(reader, pipeWriter)
go func() {
select {
case p.processCh <- timestamp:
Expand All @@ -126,8 +120,7 @@ func (p *LivePaymentProcessor) process(ctx context.Context, reader io.Reader) io

// read the segment into the buffer, because the direct use of the reader causes Broken pipe
// it's probably related to different pace of reading by trickle and ffmpeg.GetCodecInfo()
defer pipeReader.Close()
segData, err := io.ReadAll(pipeReader)
segData, err := io.ReadAll(reader.Clone())
if err != nil {
clog.InfofErr(ctx, "Error reading segment data", err)
return
Expand All @@ -139,8 +132,6 @@ func (p *LivePaymentProcessor) process(ctx context.Context, reader io.Reader) io
// We process one segment at the time, no need to buffer them
}
}()

return resReader
}

func (p *LivePaymentProcessor) shouldSkip(timestamp time.Time) bool {
Expand Down
8 changes: 5 additions & 3 deletions trickle/local_subscriber.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,17 @@ package trickle

import (
"errors"
"io"
"log/slog"
"strconv"
"sync"

"github.com/livepeer/go-livepeer/media"
)

// local (in-memory) subscriber for trickle protocol

type TrickleData struct {
Reader io.Reader
Reader media.CloneableReader
Metadata map[string]string
}

Expand Down Expand Up @@ -44,7 +45,8 @@ func (c *TrickleLocalSubscriber) Read() (*TrickleData, error) {
return nil, errors.New("seq not found")
}
c.seq++
r, w := io.Pipe()
w := media.NewMediaWriter()
r := w.MakeReader()
go func() {
subscriber := &SegmentSubscriber{
segment: segment,
Expand Down
Loading