7a04f298d2
- update to latest telegram layer - remove some references to fields in tg.Entities that don't exist in the schema - originally added here: https://github.com/beeper/td/commit/820929062a2ba0104397bc01235ab58a9cff780e - referenced here - https://github.com/mautrix/telegramgo/commit/124f0967ed195b5a380c9bd02e170ada9710dde3 - https://github.com/mautrix/telegramgo/commit/4205047aab2e0639217148b5d125bfaab668bd8e
55 lines
897 B
Go
55 lines
897 B
Go
package downloader
|
|
|
|
import (
|
|
"context"
|
|
"io"
|
|
|
|
"github.com/go-faster/errors"
|
|
|
|
"go.mau.fi/mautrix-telegram/pkg/gotd/tdsync"
|
|
"go.mau.fi/mautrix-telegram/pkg/gotd/tg"
|
|
)
|
|
|
|
func (d *Downloader) stream(ctx context.Context, r *reader, w io.Writer) (tg.StorageFileTypeClass, error) {
|
|
var typ tg.StorageFileTypeClass
|
|
|
|
g := tdsync.NewCancellableGroup(ctx)
|
|
toWrite := make(chan block, 1)
|
|
|
|
stop := func(t tg.StorageFileTypeClass) {
|
|
typ = t
|
|
close(toWrite)
|
|
}
|
|
// Download loop
|
|
g.Go(func(ctx context.Context) error {
|
|
for {
|
|
b, err := r.Next(ctx)
|
|
if err != nil {
|
|
return errors.Wrap(err, "get file")
|
|
}
|
|
|
|
n := len(b.data)
|
|
if n < 1 {
|
|
stop(b.tag)
|
|
return nil
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case toWrite <- b:
|
|
}
|
|
|
|
if b.last() {
|
|
stop(b.tag)
|
|
return nil
|
|
}
|
|
}
|
|
})
|
|
|
|
// Write loop
|
|
g.Go(writeLoop(w, toWrite))
|
|
|
|
return typ, g.Wait()
|
|
}
|