UBERF-9192: Add initial version of huly-stream (#1)

This commit is contained in:
Denis Tingaikin
2025-02-07 15:08:16 +03:00
committed by GitHub
30 changed files with 2982 additions and 1 deletions
Vendored
BIN
View File
Binary file not shown.
+44
View File
@@ -0,0 +1,44 @@
---
name: Docker push
on:
workflow_dispatch:
inputs:
version:
description: 'Version tag for the image'
required: false
default: 'latest'
type: string
jobs:
push:
runs-on: ubuntu-latest
steps:
- name: "Checkout"
uses: actions/checkout@v4
- name: "Set up Docker Buildx"
uses: docker/setup-buildx-action@v1
- name: "Login to GitHub Container Registry"
uses: docker/login-action@v1
with:
username: hardcoreeng
password: ${{ secrets.DOCKER_ACCESS_TOKEN }}
- name: Docker meta
id: metaci
uses: docker/metadata-action@v3
with:
images: hardcoreeng/huly-stream
tags: |
type=ref,event=pr
type=sha,prefix=
- name: "Build and push"
uses: docker/build-push-action@v2
with:
file: Dockerfile
context: .
platforms: linux/amd64,linux/arm64
push: true
tags: ${{ inputs.version }}
+71
View File
@@ -0,0 +1,71 @@
---
name: ci
on:
push:
branches:
- main
- develop
pull_request:
concurrency:
group: 'main'
cancel-in-progress: true
jobs:
yamllint:
name: yamllint
runs-on: ubuntu-latest
steps:
- name: Check out code
uses: actions/checkout@v4
- name: yaml-lint
uses: ibiqlik/action-yamllint@v1
with:
config_file: .github/yamllint.yaml
strict: true
build-and-test:
strategy:
matrix:
os:
- ubuntu
runs-on: ${{ matrix.os }}-latest
steps:
- name: Check out code
uses: actions/checkout@v4
- name: Setup Go
uses: actions/setup-go@v5
with:
go-version: 1.23.5
- name: Build
run: go build -race ./...
- name: Test
run: go test -race ./...
checkgomod:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version: 1.23.5
- run: go mod tidy
- name: Check for changes in go.mod or go.sum
run: |
git diff --name-only --exit-code go.mod || ( echo "Run go mod tidy" && false )
golangci-lint:
name: golangci-lint
runs-on: ubuntu-latest
steps:
- name: Check out code into the Go module directory
uses: actions/checkout@v4
with:
fetch-depth: 0
- name: Setup Go
uses: actions/setup-go@v5
with:
go-version: 1.23.5
- name: golangci-lint
uses: golangci/golangci-lint-action@v4
with:
version: v1.60.3
args: --timeout 3m --verbose
+12
View File
@@ -0,0 +1,12 @@
---
extends: default
yaml-files:
- '*.yaml'
- '*.yml'
rules:
truthy: disable
line-length: disable
comments:
min-spaces-from-content: 1
+160
View File
@@ -0,0 +1,160 @@
---
run:
go: "1.23"
timeout: 2m
issues-exit-code: 1
tests: true
linters-settings:
goheader:
template: |-
Copyright © {{ mod-year-range }} Hardcore Engineering Inc.
Licensed under the Eclipse Public License, Version 2.0 (the "License");
you may not use this file except in compliance with the License. You may
obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
errcheck:
check-type-assertions: false
check-blank: false
govet:
enable:
- shadow
settings:
printf:
funcs:
- (github.com/sirupsen/logrus.FieldLogger).Infof
- (github.com/sirupsen/logrus.FieldLogger).Warnf
- (github.com/sirupsen/logrus.FieldLogger).Errorf
- (github.com/sirupsen/logrus.FieldLogger).Fatalf
revive:
confidence: 0.8
rules:
- name: exported
- name: blank-imports
- name: context-as-argument
- name: context-keys-type
- name: dot-imports
- name: error-return
- name: error-strings
- name: error-naming
- name: exported
- name: increment-decrement
- name: package-comments
- name: range
- name: receiver-naming
- name: time-naming
- name: unexported-return
- name: indent-error-flow
- name: errorf
- name: superfluous-else
- name: unreachable-code
goimports:
local-prefixes: github.com/networkservicemesh/sdk
gocyclo:
min-complexity: 15
dupl:
threshold: 150
funlen:
lines: 120
statements: 60
goconst:
min-len: 2
min-occurrences: 2
depguard:
rules:
main:
deny:
- pkg: "errors"
desc: "Please use \"github.com/pkg/errors\" instead of \"errors\" in go imports"
misspell:
locale: US
unparam:
check-exported: false
nakedret:
max-func-lines: 30
prealloc:
simple: true
range-loops: true
for-loops: false
gosec:
excludes:
- G115
- G204
- G301
- G302
- G306
gocritic:
enabled-checks:
- appendCombine
- boolExprSimplify
- builtinShadow
- commentedOutCode
- commentedOutImport
- docStub
- dupImport
- emptyFallthrough
- emptyStringTest
- equalFold
- evalOrder
- hexLiteral
- importShadow
- indexAlloc
- initClause
- methodExprCall
- nestingReduce
- nilValReturn
- octalLiteral
- paramTypeCombine
- rangeExprCopy
- rangeValCopy
- regexpPattern
- sloppyReassign
- stringXbytes
- typeAssertChain
- typeUnparen
- unlabelStmt
- unnamedResult
- unnecessaryBlock
- weakCond
- yodaStyleExpr
linters:
disable-all: true
enable:
- goheader
- bodyclose
- unused
- depguard
- dogsled
- dupl
- errcheck
- funlen
- gochecknoinits
- goconst
- gocritic
- gocyclo
- gofmt
- goimports
- revive
- gosec
- gosimple
- govet
- ineffassign
- misspell
- nakedret
- copyloopvar
- staticcheck
- stylecheck
- typecheck
- unconvert
- unparam
- whitespace
issues:
exclude-use-default: false
max-issues-per-linter: 0
max-same-issues: 0
+36
View File
@@ -0,0 +1,36 @@
# Copyright © 2025 Hardcore Engineering Inc.
#
# Licensed under the Eclipse Public License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License. You may
# obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
#
# See the License for the specific language governing permissions and
# limitations under the License.
FROM --platform=$BUILDPLATFORM golang:1.23.5 AS builder
ENV GO111MODULE=on
ENV CGO_ENABLED=0
ENV GOBIN=/bin
ARG BUILDARCH=amd64
COPY . ./
RUN set -xe && GOOS=$TARGETOS GOARCH=$TARGETARCH go build -o /go/bin/huly-stream ./cmd/huly-stream
FROM alpine
RUN set -xe && apk add --no-cache ffmpeg
RUN apk add --no-cache ca-certificates jq bash \
&& addgroup -g 1000 huly-stream \
&& adduser -u 1000 -G huly-stream -s /bin/sh -D huly-stream \
&& chown huly-stream:huly-stream /.
COPY --from=builder /go/bin/huly-stream /huly-stream
EXPOSE 1080
USER huly-stream
ENTRYPOINT ["/huly-stream"]
+277
View File
@@ -0,0 +1,277 @@
Eclipse Public License - v 2.0
THE ACCOMPANYING PROGRAM IS PROVIDED UNDER THE TERMS OF THIS ECLIPSE
PUBLIC LICENSE ("AGREEMENT"). ANY USE, REPRODUCTION OR DISTRIBUTION
OF THE PROGRAM CONSTITUTES RECIPIENT'S ACCEPTANCE OF THIS AGREEMENT.
1. DEFINITIONS
"Contribution" means:
a) in the case of the initial Contributor, the initial content
Distributed under this Agreement, and
b) in the case of each subsequent Contributor:
i) changes to the Program, and
ii) additions to the Program;
where such changes and/or additions to the Program originate from
and are Distributed by that particular Contributor. A Contribution
"originates" from a Contributor if it was added to the Program by
such Contributor itself or anyone acting on such Contributor's behalf.
Contributions do not include changes or additions to the Program that
are not Modified Works.
"Contributor" means any person or entity that Distributes the Program.
"Licensed Patents" mean patent claims licensable by a Contributor which
are necessarily infringed by the use or sale of its Contribution alone
or when combined with the Program.
"Program" means the Contributions Distributed in accordance with this
Agreement.
"Recipient" means anyone who receives the Program under this Agreement
or any Secondary License (as applicable), including Contributors.
"Derivative Works" shall mean any work, whether in Source Code or other
form, that is based on (or derived from) the Program and for which the
editorial revisions, annotations, elaborations, or other modifications
represent, as a whole, an original work of authorship.
"Modified Works" shall mean any work in Source Code or other form that
results from an addition to, deletion from, or modification of the
contents of the Program, including, for purposes of clarity any new file
in Source Code form that contains any contents of the Program. Modified
Works shall not include works that contain only declarations,
interfaces, types, classes, structures, or files of the Program solely
in each case in order to link to, bind by name, or subclass the Program
or Modified Works thereof.
"Distribute" means the acts of a) distributing or b) making available
in any manner that enables the transfer of a copy.
"Source Code" means the form of a Program preferred for making
modifications, including but not limited to software source code,
documentation source, and configuration files.
"Secondary License" means either the GNU General Public License,
Version 2.0, or any later versions of that license, including any
exceptions or additional permissions as identified by the initial
Contributor.
2. GRANT OF RIGHTS
a) Subject to the terms of this Agreement, each Contributor hereby
grants Recipient a non-exclusive, worldwide, royalty-free copyright
license to reproduce, prepare Derivative Works of, publicly display,
publicly perform, Distribute and sublicense the Contribution of such
Contributor, if any, and such Derivative Works.
b) Subject to the terms of this Agreement, each Contributor hereby
grants Recipient a non-exclusive, worldwide, royalty-free patent
license under Licensed Patents to make, use, sell, offer to sell,
import and otherwise transfer the Contribution of such Contributor,
if any, in Source Code or other form. This patent license shall
apply to the combination of the Contribution and the Program if, at
the time the Contribution is added by the Contributor, such addition
of the Contribution causes such combination to be covered by the
Licensed Patents. The patent license shall not apply to any other
combinations which include the Contribution. No hardware per se is
licensed hereunder.
c) Recipient understands that although each Contributor grants the
licenses to its Contributions set forth herein, no assurances are
provided by any Contributor that the Program does not infringe the
patent or other intellectual property rights of any other entity.
Each Contributor disclaims any liability to Recipient for claims
brought by any other entity based on infringement of intellectual
property rights or otherwise. As a condition to exercising the
rights and licenses granted hereunder, each Recipient hereby
assumes sole responsibility to secure any other intellectual
property rights needed, if any. For example, if a third party
patent license is required to allow Recipient to Distribute the
Program, it is Recipient's responsibility to acquire that license
before distributing the Program.
d) Each Contributor represents that to its knowledge it has
sufficient copyright rights in its Contribution, if any, to grant
the copyright license set forth in this Agreement.
e) Notwithstanding the terms of any Secondary License, no
Contributor makes additional grants to any Recipient (other than
those set forth in this Agreement) as a result of such Recipient's
receipt of the Program under the terms of a Secondary License
(if permitted under the terms of Section 3).
3. REQUIREMENTS
3.1 If a Contributor Distributes the Program in any form, then:
a) the Program must also be made available as Source Code, in
accordance with section 3.2, and the Contributor must accompany
the Program with a statement that the Source Code for the Program
is available under this Agreement, and informs Recipients how to
obtain it in a reasonable manner on or through a medium customarily
used for software exchange; and
b) the Contributor may Distribute the Program under a license
different than this Agreement, provided that such license:
i) effectively disclaims on behalf of all other Contributors all
warranties and conditions, express and implied, including
warranties or conditions of title and non-infringement, and
implied warranties or conditions of merchantability and fitness
for a particular purpose;
ii) effectively excludes on behalf of all other Contributors all
liability for damages, including direct, indirect, special,
incidental and consequential damages, such as lost profits;
iii) does not attempt to limit or alter the recipients' rights
in the Source Code under section 3.2; and
iv) requires any subsequent distribution of the Program by any
party to be under a license that satisfies the requirements
of this section 3.
3.2 When the Program is Distributed as Source Code:
a) it must be made available under this Agreement, or if the
Program (i) is combined with other material in a separate file or
files made available under a Secondary License, and (ii) the initial
Contributor attached to the Source Code the notice described in
Exhibit A of this Agreement, then the Program may be made available
under the terms of such Secondary Licenses, and
b) a copy of this Agreement must be included with each copy of
the Program.
3.3 Contributors may not remove or alter any copyright, patent,
trademark, attribution notices, disclaimers of warranty, or limitations
of liability ("notices") contained within the Program from any copy of
the Program which they Distribute, provided that Contributors may add
their own appropriate notices.
4. COMMERCIAL DISTRIBUTION
Commercial distributors of software may accept certain responsibilities
with respect to end users, business partners and the like. While this
license is intended to facilitate the commercial use of the Program,
the Contributor who includes the Program in a commercial product
offering should do so in a manner which does not create potential
liability for other Contributors. Therefore, if a Contributor includes
the Program in a commercial product offering, such Contributor
("Commercial Contributor") hereby agrees to defend and indemnify every
other Contributor ("Indemnified Contributor") against any losses,
damages and costs (collectively "Losses") arising from claims, lawsuits
and other legal actions brought by a third party against the Indemnified
Contributor to the extent caused by the acts or omissions of such
Commercial Contributor in connection with its distribution of the Program
in a commercial product offering. The obligations in this section do not
apply to any claims or Losses relating to any actual or alleged
intellectual property infringement. In order to qualify, an Indemnified
Contributor must: a) promptly notify the Commercial Contributor in
writing of such claim, and b) allow the Commercial Contributor to control,
and cooperate with the Commercial Contributor in, the defense and any
related settlement negotiations. The Indemnified Contributor may
participate in any such claim at its own expense.
For example, a Contributor might include the Program in a commercial
product offering, Product X. That Contributor is then a Commercial
Contributor. If that Commercial Contributor then makes performance
claims, or offers warranties related to Product X, those performance
claims and warranties are such Commercial Contributor's responsibility
alone. Under this section, the Commercial Contributor would have to
defend claims against the other Contributors related to those performance
claims and warranties, and if a court requires any other Contributor to
pay any damages as a result, the Commercial Contributor must pay
those damages.
5. NO WARRANTY
EXCEPT AS EXPRESSLY SET FORTH IN THIS AGREEMENT, AND TO THE EXTENT
PERMITTED BY APPLICABLE LAW, THE PROGRAM IS PROVIDED ON AN "AS IS"
BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, EITHER EXPRESS OR
IMPLIED INCLUDING, WITHOUT LIMITATION, ANY WARRANTIES OR CONDITIONS OF
TITLE, NON-INFRINGEMENT, MERCHANTABILITY OR FITNESS FOR A PARTICULAR
PURPOSE. Each Recipient is solely responsible for determining the
appropriateness of using and distributing the Program and assumes all
risks associated with its exercise of rights under this Agreement,
including but not limited to the risks and costs of program errors,
compliance with applicable laws, damage to or loss of data, programs
or equipment, and unavailability or interruption of operations.
6. DISCLAIMER OF LIABILITY
EXCEPT AS EXPRESSLY SET FORTH IN THIS AGREEMENT, AND TO THE EXTENT
PERMITTED BY APPLICABLE LAW, NEITHER RECIPIENT NOR ANY CONTRIBUTORS
SHALL HAVE ANY LIABILITY FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING WITHOUT LIMITATION LOST
PROFITS), HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN
CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE)
ARISING IN ANY WAY OUT OF THE USE OR DISTRIBUTION OF THE PROGRAM OR THE
EXERCISE OF ANY RIGHTS GRANTED HEREUNDER, EVEN IF ADVISED OF THE
POSSIBILITY OF SUCH DAMAGES.
7. GENERAL
If any provision of this Agreement is invalid or unenforceable under
applicable law, it shall not affect the validity or enforceability of
the remainder of the terms of this Agreement, and without further
action by the parties hereto, such provision shall be reformed to the
minimum extent necessary to make such provision valid and enforceable.
If Recipient institutes patent litigation against any entity
(including a cross-claim or counterclaim in a lawsuit) alleging that the
Program itself (excluding combinations of the Program with other software
or hardware) infringes such Recipient's patent(s), then such Recipient's
rights granted under Section 2(b) shall terminate as of the date such
litigation is filed.
All Recipient's rights under this Agreement shall terminate if it
fails to comply with any of the material terms or conditions of this
Agreement and does not cure such failure in a reasonable period of
time after becoming aware of such noncompliance. If all Recipient's
rights under this Agreement terminate, Recipient agrees to cease use
and distribution of the Program as soon as reasonably practicable.
However, Recipient's obligations under this Agreement and any licenses
granted by Recipient relating to the Program shall continue and survive.
Everyone is permitted to copy and distribute copies of this Agreement,
but in order to avoid inconsistency the Agreement is copyrighted and
may only be modified in the following manner. The Agreement Steward
reserves the right to publish new versions (including revisions) of
this Agreement from time to time. No one other than the Agreement
Steward has the right to modify this Agreement. The Eclipse Foundation
is the initial Agreement Steward. The Eclipse Foundation may assign the
responsibility to serve as the Agreement Steward to a suitable separate
entity. Each new version of the Agreement will be given a distinguishing
version number. The Program (including Contributions) may always be
Distributed subject to the version of the Agreement under which it was
received. In addition, after a new version of the Agreement is published,
Contributor may elect to Distribute the Program (including its
Contributions) under the new version.
Except as expressly stated in Sections 2(a) and 2(b) above, Recipient
receives no rights or licenses to the intellectual property of any
Contributor under this Agreement, whether expressly, by implication,
estoppel or otherwise. All rights in the Program not expressly granted
under this Agreement are reserved. Nothing in this Agreement is intended
to be enforceable by any entity that is not a Contributor or Recipient.
No third-party beneficiary rights are created under this Agreement.
Exhibit A - Form of Secondary Licenses Notice
"This Source Code may also be made available under the following
Secondary Licenses when the conditions for such availability set forth
in the Eclipse Public License, v. 2.0 are satisfied: {name license(s),
version(s), and exceptions or additional permissions here}."
Simply including a copy of this Agreement, including this Exhibit A
is not sufficient to license the Source Code under Secondary Licenses.
If it is not possible or desirable to put the notice in a particular
file, then You may include the notice in a location (such as a LICENSE
file in a relevant directory) where a recipient would be likely to
look for such a notice.
You may add additional accurate notices of copyright ownership.
+123 -1
View File
@@ -1 +1,123 @@
# huly-stream
# Huly Stream
[![X (formerly Twitter) Follow](https://img.shields.io/twitter/follow/huly_io?style=for-the-badge)](https://x.com/huly_io)
![GitHub License](https://img.shields.io/github/license/hcengineering/platform?style=for-the-badge)
## About
The Huly Stream high-performance HTTP-based transcoding service. Huly-stream is built around the **TUS protocol**, enabling reliable, resumable file uploads and downloads. Designed for seamless and consistent media processing,it supports advanced transcoding features with robust integration options.
---
## Features
### TUS Protocol Support
- **Resumable transcoding**: Leveraging the TUS protocol, Huly-stream ensures reliable and efficient stream processing.
### Input Support
- **Supported Input Formats**:
- `mp4`
- `webm`
### Output Options
- **Supported Output Formats**:
- `aac`
- `hls`
### Upload options
- **TUS Upload**: Resumable file uploads via TUS protocol.
- **s3 Upload**: Direct upload to S3 storage.
- **datalake Upload**: Upload to datalake storage.
### Key Functionalities
- **Live transcoing with minimal upload time**: Transcoding results are going to be avaible after stream completion.
- **Transcoding Cancelation**: Cancel or pause ongoing transcoding in real-time.
- **Transcoding Resumption**: Resume incomplete transcoding tasks efficiently.
---
## Installation
### Prerequisites
- [Go](https://golang.org/dl/) (v1.23+ recommended)
- [ffmpeg](https://www.ffmpeg.org/download.html) (ensure it’s installed and available in your system's PATH)
### Steps
1. Install dependencies:
```bash
go mod tidy
```
2. Build the service:
```bash
docker build . -t hcengineering/huly-stream:latest
```
---
## Configuraiton
### App env configuraiton
The following environment variables can be used:
```
KEY TYPE DEFAULT REQUIRED DESCRIPTION
STREAM_SECRET_TOKEN String secret token for authorize requests
STREAM_LOG_LEVEL String debug sets log level for the application
STREAM_PPROF_ENABLED True or False false starts profile server on localhost:6060 if true
STREAM_INSECURE True or False false ignores authorization check if true
STREAM_SERVE_URL String 0.0.0.0:1080 app listen url
STREAM_ENDPOINT_URL URL S3 or Datalake endpoint, example: s3://my-ip-address, datalake://my-ip-address
STREAM_MAX_CAPACITY Integer 6220800 represents the amount of maximum possible capacity for the transcoding. The default value is 1920 * 1080 * 3.
STREAM_MAX_THREADS Integer 4 means upper bound for the transcoing provider.
STREAM_OUTPUT_DIR String /tmp/transcoing/ path to the directory with transcoding result.
STREAM_REMOVE_CONTENT_ON_UPLOAD True or False true deletes all content when content delivered if true
STREAM_UPLOAD_RAW_CONTENT True or False false uploads content in raw quality to the endpoint if true
```
### Metadata
**resolutions:** if passed, set the resolution for the output, for example, 'resolutions: 1920:1080, 1280:720.'
**token:** must be provided to be authorized in the Huly's datalake service.
**workspace:** required for uploading content to the datalake storage.
#### S3 Env configuration
if you're working with S3 storage type, these envs must be provided:
**AWS_ACCESS_KEY_ID**
**AWS_SECRET_ACCESS_KEY**
## Usage
The service exposes an HTTP API. Below are some examples of how to interact with it.
### Upload a File for Transcoding via TUS
```bash
curl -X POST http://localhost:1080/transcoing \
-H "Tus-Resumable: 1.0.0" \
-H "Upload-Length: <file-size>" \
--data-binary @path/to/your/file.mp4
```
## Contributing
We welcome contributions! To get started:
1. Fork the repository.
2. Create a new branch for your feature or bug fix.
3. Submit a pull request describing your changes.
---
## License
This project is licensed under the [MIT License](LICENSE).
---
Enjoy seamless transcoding with huly-stream! 🚀
+104
View File
@@ -0,0 +1,104 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package main provides huly-stream entry point function
package main
import (
"context"
"net/http"
"os"
"os/signal"
"syscall"
"go.uber.org/zap"
"golang.org/x/exp/slog"
"github.com/huly-stream/internal/pkg/config"
"github.com/huly-stream/internal/pkg/log"
"github.com/huly-stream/internal/pkg/pprof"
"github.com/huly-stream/internal/pkg/transcoding"
tusd "github.com/tus/tusd/v2/pkg/handler"
)
const basePath = "/transcoding"
func main() {
var ctx, cancel = signal.NotifyContext(
context.Background(),
os.Interrupt,
syscall.SIGHUP,
syscall.SIGTERM,
syscall.SIGQUIT,
)
defer cancel()
ctx = log.WithLoggerFields(ctx)
var logger = log.FromContext(ctx)
var conf = must(config.FromEnv())
logger.Sugar().Debugf("provided config is %v", conf)
mustNoError(os.MkdirAll(conf.OutputDir, os.ModePerm))
if conf.PprofEnabled {
go pprof.ListenAndServe(ctx, "localhost:6060")
}
scheduler := transcoding.NewScheduler(ctx, conf)
tusComposer := tusd.NewStoreComposer()
tusComposer.UseCore(scheduler)
tusComposer.UseTerminater(scheduler)
tusComposer.UseConcater(scheduler)
tusComposer.UseLengthDeferrer(scheduler)
var handler = must(tusd.NewHandler(tusd.Config{
BasePath: basePath,
StoreComposer: tusComposer,
Logger: slog.New(slog.NewTextHandler(discardTextHandler{}, nil)),
}))
http.Handle("/transcoding/", http.StripPrefix("/transcoding/", handler))
http.Handle("/transcoding", http.StripPrefix("/transcoding", handler))
go func() {
// #nosec
var err = http.ListenAndServe(conf.ServeURL, nil)
if err != nil {
cancel()
logger.Debug("unable to listen", zap.Error(err))
}
}()
<-ctx.Done()
}
type discardTextHandler struct{}
func (discardTextHandler) Write([]byte) (int, error) {
return 0, nil
}
func mustNoError(err error) {
if err != nil {
panic(err.Error())
}
}
func must[T any](val T, err error) T {
mustNoError(err)
return val
}
+44
View File
@@ -0,0 +1,44 @@
module github.com/huly-stream
go 1.23.2
require (
github.com/aws/aws-sdk-go-v2 v1.32.3
github.com/aws/aws-sdk-go-v2/config v1.28.1
github.com/aws/aws-sdk-go-v2/credentials v1.17.42
github.com/aws/aws-sdk-go-v2/service/s3 v1.66.2
github.com/aws/smithy-go v1.22.0
github.com/fsnotify/fsnotify v1.8.0
github.com/google/uuid v1.6.0
github.com/kelseyhightower/envconfig v1.4.0
github.com/pkg/errors v0.9.1
github.com/stretchr/testify v1.9.0
github.com/tus/tusd/v2 v2.6.0
github.com/valyala/fasthttp v1.58.0
go.uber.org/zap v1.27.0
golang.org/x/exp v0.0.0-20230626212559-97b1e661b5df
)
require (
github.com/andybalholm/brotli v1.1.1 // indirect
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.6 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.18 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.22 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.22 // indirect
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1 // indirect
github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.22 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.0 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.3 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.3 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.3 // indirect
github.com/aws/aws-sdk-go-v2/service/sso v1.24.3 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.3 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.32.3 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/klauspost/compress v1.17.11 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/valyala/bytebufferpool v1.0.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
golang.org/x/sys v0.27.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
+84
View File
@@ -0,0 +1,84 @@
github.com/Acconut/go-httptest-recorder v1.0.0 h1:TAv2dfnqp/l+SUvIaMAUK4GeN4+wqb6KZsFFFTGhoJg=
github.com/Acconut/go-httptest-recorder v1.0.0/go.mod h1:CwQyhTH1kq/gLyWiRieo7c0uokpu3PXeyF/nZjUNtmM=
github.com/andybalholm/brotli v1.1.1 h1:PR2pgnyFznKEugtsUo0xLdDop5SKXd5Qf5ysW+7XdTA=
github.com/andybalholm/brotli v1.1.1/go.mod h1:05ib4cKhjx3OQYUY22hTVd34Bc8upXjOLL2rKwwZBoA=
github.com/aws/aws-sdk-go-v2 v1.32.3 h1:T0dRlFBKcdaUPGNtkBSwHZxrtis8CQU17UpNBZYd0wk=
github.com/aws/aws-sdk-go-v2 v1.32.3/go.mod h1:2SK5n0a2karNTv5tbP1SjsX0uhttou00v/HpXKM1ZUo=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.6 h1:pT3hpW0cOHRJx8Y0DfJUEQuqPild8jRGmSFmBgvydr0=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.6/go.mod h1:j/I2++U0xX+cr44QjHay4Cvxj6FUbnxrgmqN3H1jTZA=
github.com/aws/aws-sdk-go-v2/config v1.28.1 h1:oxIvOUXy8x0U3fR//0eq+RdCKimWI900+SV+10xsCBw=
github.com/aws/aws-sdk-go-v2/config v1.28.1/go.mod h1:bRQcttQJiARbd5JZxw6wG0yIK3eLeSCPdg6uqmmlIiI=
github.com/aws/aws-sdk-go-v2/credentials v1.17.42 h1:sBP0RPjBU4neGpIYyx8mkU2QqLPl5u9cmdTWVzIpHkM=
github.com/aws/aws-sdk-go-v2/credentials v1.17.42/go.mod h1:FwZBfU530dJ26rv9saAbxa9Ej3eF/AK0OAY86k13n4M=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.18 h1:68jFVtt3NulEzojFesM/WVarlFpCaXLKaBxDpzkQ9OQ=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.18/go.mod h1:Fjnn5jQVIo6VyedMc0/EhPpfNlPl7dHV916O6B+49aE=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.22 h1:Jw50LwEkVjuVzE1NzkhNKkBf9cRN7MtE1F/b2cOKTUM=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.22/go.mod h1:Y/SmAyPcOTmpeVaWSzSKiILfXTVJwrGmYZhcRbhWuEY=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.22 h1:981MHwBaRZM7+9QSR6XamDzF/o7ouUGxFzr+nVSIhrs=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.22/go.mod h1:1RA1+aBEfn+CAB/Mh0MB6LsdCYCnjZm7tKXtnk499ZQ=
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1 h1:VaRN3TlFdd6KxX1x3ILT5ynH6HvKgqdiXoTxAF4HQcQ=
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1/go.mod h1:FbtygfRFze9usAadmnGJNc8KsP346kEe+y2/oyhGAGc=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.22 h1:yV+hCAHZZYJQcwAaszoBNwLbPItHvApxT0kVIw6jRgs=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.22/go.mod h1:kbR1TL8llqB1eGnVbybcA4/wgScxdylOdyAd51yxPdw=
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.0 h1:TToQNkvGguu209puTojY/ozlqy2d/SFNcoLIqTFi42g=
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.0/go.mod h1:0jp+ltwkf+SwG2fm/PKo8t4y8pJSgOCO4D8Lz3k0aHQ=
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.3 h1:kT6BcZsmMtNkP/iYMcRG+mIEA/IbeiUimXtGmqF39y0=
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.3/go.mod h1:Z8uGua2k4PPaGOYn66pK02rhMrot3Xk3tpBuUFPomZU=
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.3 h1:qcxX0JYlgWH3hpPUnd6U0ikcl6LLA9sLkXE2w1fpMvY=
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.3/go.mod h1:cLSNEmI45soc+Ef8K/L+8sEA3A3pYFEYf5B5UI+6bH4=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.3 h1:ZC7Y/XgKUxwqcdhO5LE8P6oGP1eh6xlQReWNKfhvJno=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.3/go.mod h1:WqfO7M9l9yUAw0HcHaikwRd/H6gzYdz7vjejCA5e2oY=
github.com/aws/aws-sdk-go-v2/service/s3 v1.66.2 h1:p9TNFL8bFUMd+38YIpTAXpoxyz0MxC7FlbFEH4P4E1U=
github.com/aws/aws-sdk-go-v2/service/s3 v1.66.2/go.mod h1:fNjyo0Coen9QTwQLWeV6WO2Nytwiu+cCcWaTdKCAqqE=
github.com/aws/aws-sdk-go-v2/service/sso v1.24.3 h1:UTpsIf0loCIWEbrqdLb+0RxnTXfWh2vhw4nQmFi4nPc=
github.com/aws/aws-sdk-go-v2/service/sso v1.24.3/go.mod h1:FZ9j3PFHHAR+w0BSEjK955w5YD2UwB/l/H0yAK3MJvI=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.3 h1:2YCmIXv3tmiItw0LlYf6v7gEHebLY45kBEnPezbUKyU=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.3/go.mod h1:u19stRyNPxGhj6dRm+Cdgu6N75qnbW7+QN0q0dsAk58=
github.com/aws/aws-sdk-go-v2/service/sts v1.32.3 h1:wVnQ6tigGsRqSWDEEyH6lSAJ9OyFUsSnbaUWChuSGzs=
github.com/aws/aws-sdk-go-v2/service/sts v1.32.3/go.mod h1:VZa9yTFyj4o10YGsmDO4gbQJUvvhY72fhumT8W4LqsE=
github.com/aws/smithy-go v1.22.0 h1:uunKnWlcoL3zO7q+gG2Pk53joueEOsnNB28QdMsmiMM=
github.com/aws/smithy-go v1.22.0/go.mod h1:irrKGvNn1InZwb2d7fkIRNucdfwR8R+Ts3wxYa/cJHg=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/fsnotify/fsnotify v1.8.0 h1:dAwr6QBTBZIkG8roQaJjGof0pp0EeF+tNV7YBP3F/8M=
github.com/fsnotify/fsnotify v1.8.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0=
github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc=
github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/kelseyhightower/envconfig v1.4.0 h1:Im6hONhd3pLkfDFsbRgu68RDNkGF1r3dvMUtDTo2cv8=
github.com/kelseyhightower/envconfig v1.4.0/go.mod h1:cccZRl6mQpaq41TPp5QxidR+Sa3axMbJDNb//FQX6Gg=
github.com/klauspost/compress v1.17.11 h1:In6xLpyWOi1+C7tXUUWv2ot1QvBjxevKAaI6IXrJmUc=
github.com/klauspost/compress v1.17.11/go.mod h1:pMDklpSncoRMuLFrf1W9Ss9KT+0rH90U12bZKk7uwG0=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg=
github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
github.com/tus/tusd/v2 v2.6.0 h1:Je243QDKnFTvm/WkLH2bd1oQ+7trolrflRWyuI0PdWI=
github.com/tus/tusd/v2 v2.6.0/go.mod h1:1Eb1lBoSRBfYJ/mQfFVjyw8ZdNMdBqW17vgQKl3Ah9g=
github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw=
github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc=
github.com/valyala/fasthttp v1.58.0 h1:GGB2dWxSbEprU9j0iMJHgdKYJVDyjrOwF9RE59PbRuE=
github.com/valyala/fasthttp v1.58.0/go.mod h1:SYXvHHaFp7QZHGKSHmoMipInhrI5StHrhDTYVEjK/Kw=
github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU=
github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0=
go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y=
go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8=
go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E=
golang.org/x/exp v0.0.0-20230626212559-97b1e661b5df h1:UA2aFVmmsIlefxMk29Dp2juaUSth8Pyn3Tq5Y5mJGME=
golang.org/x/exp v0.0.0-20230626212559-97b1e661b5df/go.mod h1:FXUEEKJgO7OQYeo8N01OfiKP8RXMtf6e8aTskBGqWdc=
golang.org/x/net v0.31.0 h1:68CPQngjLL0r2AlUKiSxtQFKvzRVbnzLwMUn5SzcLHo=
golang.org/x/net v0.31.0/go.mod h1:P4fl1q7dY2hnZFxEk4pPSkDHF+QqjitcnDjUQyMM+pM=
golang.org/x/sys v0.27.0 h1:wBqf8DvsY9Y/2P8gAfPDEYNuS30J4lPHJxXSb/nJZ+s=
golang.org/x/sys v0.27.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/text v0.20.0 h1:gK/Kv2otX8gz+wn7Rmb3vT96ZwuoxnQlY+HlJVj7Qug=
golang.org/x/text v0.20.0/go.mod h1:D4IsuqiFMhST5bX19pQ9ikHC2GsaKyk/oF+pn3ducp4=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+51
View File
@@ -0,0 +1,51 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package config provides configuration for the application
package config
import (
"net/url"
"github.com/kelseyhightower/envconfig"
)
// Config represents configuration for the huly-stream application.
type Config struct {
SecretToken string `split_words:"true" desc:"secret token for authorize requests"`
LogLevel string `split_words:"true" default:"debug" desc:"sets log level for the application"`
PprofEnabled bool `default:"false" split_words:"true" desc:"starts profile server on localhost:6060 if true"`
Insecure bool `default:"false" desc:"ignores authorization check if true"`
ServeURL string `split_words:"true" desc:"app listen url" default:"0.0.0.0:1080"`
EndpointURL *url.URL `split_words:"true" desc:"S3 or Datalake endpoint, example: s3://my-ip-address, datalake://my-ip-address"`
MaxCapacity int64 `split_words:"true" default:"6220800" desc:"represents the amount of maximum possible capacity for the transcoding. The default value is 1920 * 1080 * 3."`
MaxThreads int `split_words:"true" default:"4" desc:"means upper bound for the transcoing provider."`
OutputDir string `split_words:"true" default:"/tmp/transcoing/" desc:"path to the directory with transcoding result."`
RemoveContentOnUpload bool `split_words:"true" default:"true" desc:"deletes all content when content delivered if true"`
UploadRawContent bool `split_words:"true" default:"false" desc:"uploads content in raw quality to the endpoint if true"`
}
// FromEnv creates new Config from env
func FromEnv() (*Config, error) {
var result Config
if err := envconfig.Usage("stream", &result); err != nil {
return nil, err
}
if err := envconfig.Process("stream", &result); err != nil {
return nil, err
}
return &result, nil
}
+50
View File
@@ -0,0 +1,50 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package log provides simple api for using inherited logging
package log
import (
"context"
"go.uber.org/zap"
)
type contextKey struct{}
// WithLoggerFields createsa new context with zap.Logger and passed fields
func WithLoggerFields(ctx context.Context, fields ...zap.Field) context.Context {
var logger = FromContext(ctx)
if logger == nil {
var err error
logger, err = zap.NewDevelopment()
if err != nil {
panic(err.Error())
}
logger.Info("zap logger was initialized")
go func() {
<-ctx.Done()
_ = logger.Sync()
}()
}
return context.WithValue(ctx, contextKey{}, logger.With(fields...))
}
// FromContext returns zap.Logger from the context
func FromContext(ctx context.Context) *zap.Logger {
var val = ctx.Value(contextKey{})
if val == nil {
return nil
}
return val.(*zap.Logger)
}
+125
View File
@@ -0,0 +1,125 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package manifest provides data types for manifest based media files.
package manifest
import (
"bufio"
"fmt"
"strconv"
"strings"
)
// HLSManifest represents an HLS manifest file
// with metadata about the playlist and its segments.
type HLSManifest struct {
Version int
TargetDuration int
SequenceNumber int
Segments []Segment
EndList bool
}
// Segment represents a media segment in the HLS manifest.
type Segment struct {
URI string
Duration float64
Title string
}
// ToM3U8 serializes the HLSManifest to an M3U8 file format.
func (m *HLSManifest) ToM3U8() string {
var builder strings.Builder
builder.WriteString("#EXTM3U\n")
builder.WriteString(fmt.Sprintf("#EXT-X-VERSION:%d\n", m.Version))
builder.WriteString(fmt.Sprintf("#EXT-X-TARGETDURATION:%d\n", m.TargetDuration))
builder.WriteString(fmt.Sprintf("#EXT-X-MEDIA-SEQUENCE:%d\n", m.SequenceNumber))
for _, segment := range m.Segments {
if segment.Title != "" {
builder.WriteString(fmt.Sprintf("#EXTINF:%.2f,%s\n", segment.Duration, segment.Title))
} else {
builder.WriteString(fmt.Sprintf("#EXTINF:%.2f,\n", segment.Duration))
}
builder.WriteString(fmt.Sprintf("%s\n", segment.URI))
}
if m.EndList {
builder.WriteString("#EXT-X-ENDLIST\n")
}
return builder.String()
}
// FromM3U8 converts raw input to the hls master file
// nolint
func FromM3U8(data string) (*HLSManifest, error) {
scanner := bufio.NewScanner(strings.NewReader(data))
manifest := &HLSManifest{}
var currentSegment *Segment
for scanner.Scan() {
line := strings.TrimSpace(scanner.Text())
if line == "" {
continue
}
if strings.HasPrefix(line, "#EXTM3U") {
continue
}
if strings.HasPrefix(line, "#EXT-X-VERSION:") {
version, err := strconv.Atoi(strings.TrimPrefix(line, "#EXT-X-VERSION:"))
if err != nil {
return nil, err
}
manifest.Version = version
} else if strings.HasPrefix(line, "#EXT-X-TARGETDURATION:") {
targetDuration, err := strconv.Atoi(strings.TrimPrefix(line, "#EXT-X-TARGETDURATION:"))
if err != nil {
return nil, err
}
manifest.TargetDuration = targetDuration
} else if strings.HasPrefix(line, "#EXT-X-MEDIA-SEQUENCE:") {
sequenceNumber, err := strconv.Atoi(strings.TrimPrefix(line, "#EXT-X-MEDIA-SEQUENCE:"))
if err != nil {
return nil, err
}
manifest.SequenceNumber = sequenceNumber
} else if strings.HasPrefix(line, "#EXTINF:") {
parts := strings.SplitN(strings.TrimPrefix(line, "#EXTINF:"), ",", 2)
duration, err := strconv.ParseFloat(parts[0], 64)
if err != nil {
return nil, err
}
title := ""
if len(parts) > 1 {
title = parts[1]
}
currentSegment = &Segment{Duration: duration, Title: title}
} else if strings.HasPrefix(line, "#EXT-X-ENDLIST") {
manifest.EndList = true
} else if currentSegment != nil {
currentSegment.URI = line
manifest.Segments = append(manifest.Segments, *currentSegment)
currentSegment = nil
}
}
if err := scanner.Err(); err != nil {
return nil, err
}
return manifest, nil
}
+143
View File
@@ -0,0 +1,143 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package manifest_test
import (
"testing"
"github.com/huly-stream/internal/pkg/manifest"
"github.com/stretchr/testify/assert"
)
func TestToM3U8(t *testing.T) {
tests := []struct {
name string
manifest manifest.HLSManifest
expected string
}{
{
name: "simple manifest",
manifest: manifest.HLSManifest{
Version: 3,
TargetDuration: 10,
SequenceNumber: 1,
Segments: []manifest.Segment{
{URI: "segment1.ts", Duration: 9.5, Title: "Segment 1"},
{URI: "segment2.ts", Duration: 9.0, Title: "Segment 2"},
},
EndList: true,
},
expected: `#EXTM3U
#EXT-X-VERSION:3
#EXT-X-TARGETDURATION:10
#EXT-X-MEDIA-SEQUENCE:1
#EXTINF:9.50,Segment 1
segment1.ts
#EXTINF:9.00,Segment 2
segment2.ts
#EXT-X-ENDLIST
`,
},
{
name: "empty manifest",
manifest: manifest.HLSManifest{
Version: 3,
TargetDuration: 10,
SequenceNumber: 1,
Segments: []manifest.Segment{},
EndList: false,
},
expected: `#EXTM3U
#EXT-X-VERSION:3
#EXT-X-TARGETDURATION:10
#EXT-X-MEDIA-SEQUENCE:1
`,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
actual := tt.manifest.ToM3U8()
assert.Equal(t, tt.expected, actual)
})
}
}
func TestFromM3U8(t *testing.T) {
tests := []struct {
name string
data string
expected manifest.HLSManifest
err bool
}{
{
name: "valid manifest",
data: `#EXTM3U
#EXT-X-VERSION:3
#EXT-X-TARGETDURATION:10
#EXT-X-MEDIA-SEQUENCE:1
#EXTINF:9.50,Segment 1
segment1.ts
#EXTINF:9.00,Segment 2
segment2.ts
#EXT-X-ENDLIST
`,
expected: manifest.HLSManifest{
Version: 3,
TargetDuration: 10,
SequenceNumber: 1,
Segments: []manifest.Segment{
{URI: "segment1.ts", Duration: 9.5, Title: "Segment 1"},
{URI: "segment2.ts", Duration: 9.0, Title: "Segment 2"},
},
EndList: true,
},
err: false,
},
// {
// name: "missing target duration",
// data: `#EXTM3U
// #EXT-X-VERSION:3
// #EXT-X-MEDIA-SEQUENCE:1
// #EXTINF:9.50,Segment 1
// segment1.ts
// `,
// expected: manifest.HLSManifest{},
// err: true,
// },
{
name: "empty file",
data: "",
expected: manifest.HLSManifest{
Version: 0,
TargetDuration: 0,
SequenceNumber: 0,
EndList: false,
},
err: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
actual, err := manifest.FromM3U8(tt.data)
if tt.err {
assert.Error(t, err)
} else {
assert.NoError(t, err)
assert.Equal(t, tt.expected, *actual)
}
})
}
}
+54
View File
@@ -0,0 +1,54 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package pprof provides all-in configured pprof server for debugging purposes
package pprof
import (
"context"
"net/http"
"net/http/pprof"
"time"
"github.com/huly-stream/internal/pkg/log"
"go.uber.org/zap"
)
// ListenAndServe - configures pprof http handlers
func ListenAndServe(ctx context.Context, listenOn string) {
log.FromContext(ctx).Debug("Profiler is enabled", zap.String("listening on", listenOn))
mux := http.NewServeMux()
mux.HandleFunc("/debug/pprof/", pprof.Index)
mux.HandleFunc("/debug/pprof/cmdline", pprof.Cmdline)
mux.HandleFunc("/debug/pprof/profile", pprof.Profile)
mux.HandleFunc("/debug/pprof/symbol", pprof.Symbol)
mux.HandleFunc("/debug/pprof/trace", pprof.Trace)
mux.Handle("/debug/pprof/allocs", pprof.Handler("allocs"))
mux.Handle("/debug/pprof/block", pprof.Handler("block"))
mux.Handle("/debug/pprof/goroutine", pprof.Handler("goroutine"))
mux.Handle("/debug/pprof/heap", pprof.Handler("heap"))
mux.Handle("/debug/pprof/mutex", pprof.Handler("mutex"))
mux.Handle("/debug/pprof/threadcreate", pprof.Handler("threadcreate"))
server := &http.Server{
Addr: listenOn,
Handler: mux,
ReadTimeout: 10 * time.Second,
WriteTimeout: 10 * time.Second,
}
if err := server.ListenAndServe(); err != nil {
log.FromContext(ctx).Debug("Failed to start profiler", zap.Error(err))
}
<-ctx.Done()
_ = server.Close()
}
+116
View File
@@ -0,0 +1,116 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package sharedpipe provided a shared pipe that can be used when one writer can be shared between a couple of readers.
package sharedpipe
import (
"io"
"sync"
"sync/atomic"
)
// Chunk represents a chunk of raw data for the readers
type Chunk struct {
content []byte
Next atomic.Pointer[Chunk]
ready chan struct{}
done chan struct{}
}
// NewWriter creates a new shared pipe Writer
func NewWriter() *Writer {
var res = &Writer{
done: make(chan struct{}),
}
res.curr.Store(&Chunk{done: res.done, ready: make(chan struct{})})
return res
}
// Transpile creates a new Reader
func (w *Writer) Transpile() *Reader {
var res = &Reader{
curr: w.curr.Load(),
done: make(chan struct{}),
}
return res
}
// Writer represents a shared pipe writer
type Writer struct {
curr atomic.Pointer[Chunk]
done chan struct{}
once sync.Once
}
// Close closes the pipe for all readers
func (w *Writer) Close() error {
w.once.Do(func() { close(w.done) })
return nil
}
func (w *Writer) Write(p []byte) (n int, err error) {
var completePrevious = w.curr.Load().ready
var curr = w.curr.Load()
curr.Next.Store(&Chunk{content: p, ready: make(chan struct{}), done: w.done})
w.curr.Store(curr.Next.Load())
close(completePrevious)
return len(p), nil
}
// Reader is reader from shared pipe, imlements io.Reader interface
type Reader struct {
curr *Chunk
once sync.Once
done chan struct{}
pos int
}
// Close closes reader stream
func (s *Reader) Close() error {
s.once.Do(func() {
close(s.done)
})
return nil
}
func (s *Reader) Read(in []byte) (n int, err error) {
var curr = s.curr
for i := 0; i < len(in); {
for s.pos >= len(curr.content) {
select {
case <-curr.done:
curr = curr.Next.Load()
if curr == nil {
_ = s.Close()
return i, io.EOF
}
case <-s.done:
return i, io.ErrClosedPipe
case <-curr.ready:
curr = curr.Next.Load()
}
s.pos = 0
s.curr = curr
}
var n = copy(in[i:], curr.content[s.pos:])
s.pos += n
i += n
}
return len(in), nil
}
var _ io.Closer = (*Writer)(nil)
var _ io.Closer = (*Reader)(nil)
var _ io.Reader = (*Reader)(nil)
var _ io.Writer = (*Writer)(nil)
@@ -0,0 +1,248 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
//
package sharedpipe
import (
"bytes"
"io"
"math/rand"
"strings"
"sync"
"testing"
"time"
"github.com/stretchr/testify/require"
)
const sendMessageSize = 8 * 1000 * 1000
const readerCount = 10
func TestStability(t *testing.T) {
for range 10 {
testStability(t)
}
}
// #nosec
func testStability(t *testing.T) {
var writer = NewWriter()
var readers []io.Reader
for range 1000 {
readers = append(readers, writer.Transpile())
}
var buff [4]byte
for range rand.Intn(1000) {
_, _ = writer.Write([]byte("ping"))
for i := range rand.Intn(10) {
_, _ = readers[i].Read(buff[:])
}
}
_ = writer.Close()
for _, r := range readers {
_, err := io.ReadAll(r)
require.NoError(t, err)
}
}
func TestBasicWriteRead(t *testing.T) {
writer := NewWriter()
defer func() { _ = writer.Close() }()
reader := writer.Transpile()
data := []byte("Hello, World!")
n, err := writer.Write(data)
if err != nil {
t.Fatalf("Unexpected error on Write: %v", err)
}
if n != len(data) {
t.Fatalf("Expected to write %d bytes, wrote %d", len(data), n)
}
readBuf := make([]byte, len(data))
n, err = reader.Read(readBuf)
if err != nil && err != io.EOF {
t.Fatalf("Unexpected error on Read: %v", err)
}
if !bytes.Equal(readBuf[:n], data) {
t.Fatalf("Expected to read %q, got %q", data, readBuf[:n])
}
}
func TestConcurrentWriteRead(t *testing.T) {
writer := NewWriter()
defer func() { _ = writer.Close() }()
reader := writer.Transpile()
var wg sync.WaitGroup
data := []byte("Hello, Concurrent World!")
readBuf := make([]byte, len(data))
wg.Add(2)
go func() {
defer wg.Done()
_, err := writer.Write(data)
if err != nil {
t.Errorf("Unexpected error on Write: %v", err)
}
}()
// Reader goroutine
go func() {
defer wg.Done()
n, err := reader.Read(readBuf)
if err != nil && err != io.EOF {
t.Errorf("Unexpected error on Read: %v", err)
}
if n != len(data) {
t.Errorf("Expected to read %d bytes, read %d", len(data), n)
}
if !bytes.Equal(readBuf[:n], data) {
t.Errorf("Expected to read %q, got %q", data, readBuf[:n])
}
}()
wg.Wait()
}
func TestWriterClose(t *testing.T) {
writer := NewWriter()
reader := writer.Transpile()
_, _ = writer.Write([]byte("Hello"))
_ = writer.Close()
readBuf := make([]byte, 5)
n, err := reader.Read(readBuf)
if err != nil && err != io.EOF {
t.Fatalf("Unexpected error on Read: %v", err)
}
if n != 5 {
t.Fatalf("Expected to read 5 bytes, read %d", n)
}
n, err = reader.Read(readBuf)
if err != io.EOF {
t.Fatalf("Expected EOF after writer close, got %v", err)
}
if n != 0 {
t.Fatalf("Expected to read 0 bytes after EOF, read %d", n)
}
}
// Test reading from an empty writer.
func TestReadFromEmptyWriter(t *testing.T) {
writer := NewWriter()
_ = writer.Close()
reader := writer.Transpile()
readBuf := make([]byte, 5)
_, err := reader.Read(readBuf)
require.ErrorIs(t, io.EOF, err)
}
func Test_PipeWait(t *testing.T) {
var writer = NewWriter()
var reader = writer.Transpile()
var buff [4]byte
go func() {
time.Sleep(time.Second / 10)
_, _ = writer.Write([]byte("test"))
}()
_, _ = reader.Read(buff[:])
require.Equal(t, "test", string(buff[:]))
}
func Test_Consistent(t *testing.T) {
var writer = NewWriter()
var readers []io.Reader
for range readerCount {
readers = append(readers, writer.Transpile())
}
_, _ = writer.Write([]byte("Hello"))
_, _ = writer.Write([]byte(" "))
_, _ = writer.Write([]byte("World!"))
_ = writer.Close()
var res strings.Builder
for i := range readerCount {
for {
var buffer = make([]byte, 2)
_, err := readers[i].Read(buffer)
if err == io.EOF {
break
}
_, _ = res.WriteString(string(buffer))
}
require.Equal(t, "Hello World!", res.String())
res.Reset()
}
}
// Benchmark_DefaultPipe-8 (4 b) 61956 19177 ns/op 48 B/op 1 allocs/op
// Benchmark_DefaultPipe-8 (8 mb) 22 49187741 ns/op 118 B/op 1 allocs/op
func Benchmark_DefaultPipe(b *testing.B) {
var data [sendMessageSize]byte
var buffer = make([]byte, len(data))
var readers []io.Reader
var writers []io.Writer
for i := 0; i < readerCount; i++ {
r, w := io.Pipe()
readers = append(readers, r)
writers = append(writers, w)
}
b.ReportAllocs()
b.ResetTimer()
for range b.N {
go func() {
for i := 0; i < readerCount; i++ {
_, _ = writers[i].Write(data[:])
}
}()
for i := 0; i < readerCount; i++ {
_, _ = readers[i].Read(buffer)
}
}
}
// Benchmark_SharedPipe-8 (4 b) 161847 8131 ns/op 160 B/op 2 allocs/op
// Benchmark_SharedPipe-8 (8 mb) 69 15880031 ns/op 160 B/op 2 allocs/op
func Benchmark_SharedPipe(b *testing.B) {
var data [sendMessageSize]byte
var buffer = make([]byte, len(data))
var writer = NewWriter()
var readers []io.Reader
for range readerCount {
readers = append(readers, writer.Transpile())
}
b.ReportAllocs()
b.ResetTimer()
for range b.N {
_, _ = writer.Write(data[:])
for i := 0; i < readerCount; i++ {
_, _ = readers[i].Read(buffer)
}
}
}
+163
View File
@@ -0,0 +1,163 @@
//
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
//
package transcoding
import (
"context"
"fmt"
"io"
"os"
"os/exec"
"path/filepath"
"sort"
"strconv"
"strings"
"github.com/pkg/errors"
"github.com/huly-stream/internal/pkg/log"
"go.uber.org/zap"
)
// Options represents configuration for the ffmpeg command
type Options struct {
OuputDir string
Resolutions []string
Threads int
UploadID string
}
func measure(options *Options) int64 {
var res int64
for _, resolution := range options.Resolutions {
var w, h int
var parts = strings.Split(resolution, ":")
if len(parts) > 1 {
w, _ = strconv.Atoi(parts[0])
w = max(w, 320)
h, _ = strconv.Atoi(parts[1])
h = max(h, 240)
res += int64(w) * int64(h)
}
}
return max(res, 320*240)
}
func newFfmpegCommand(ctx context.Context, in io.Reader, options *Options) (*exec.Cmd, error) {
if options == nil {
return nil, errors.New("options should not be nil")
}
if ctx == nil {
return nil, errors.New("ctx should not be nil")
}
var logger = log.FromContext(ctx).With(zap.String("func", "NewFFMpegCommand"))
var args []string
if options.Resolutions == nil {
logger.Debug("resolutions were not provided, building audio command...")
args = BuildAudioCommand(options)
} else {
logger.Debug("building video command...")
args = BuildVideoCommand(options)
}
logger.Debug("prepared command: ", zap.Strings("args", args))
var result = exec.CommandContext(ctx, "ffmpeg", args...)
result.Stderr = os.Stdout
result.Stdout = os.Stdout
result.Stdin = in
return result, nil
}
func buildCommonComamnd(opts *Options) []string {
return []string{
"-threads", fmt.Sprint(opts.Threads),
"-i", "pipe:0",
}
}
// BuildAudioCommand returns flags for getting the audio from the input
func BuildAudioCommand(opts *Options) []string {
var commonPart = buildCommonComamnd(opts)
return append(commonPart,
"-vn", "-acodec",
"copy", filepath.Join(opts.OuputDir, opts.UploadID),
)
}
// BuildVideoCommand returns flags for ffmpeg for video transcoding
func BuildVideoCommand(opts *Options) []string {
var result = buildCommonComamnd(opts)
for _, res := range opts.Resolutions {
var prefix string
var w, h int
var parts = strings.Split(res, ":")
if len(parts) > 1 {
w, _ = strconv.Atoi(parts[0])
h, _ = strconv.Atoi(parts[1])
}
w = max(w, 640)
h = max(h, 480)
prefix = ResolutionFromPixels(w * h)
result = append(result,
"-vf", fmt.Sprintf("scale=%d:%d", w, h),
"-c:v",
"libx264",
"-preset", "veryfast",
"-crf", "23",
"-g", "60",
"-hls_time", "5",
"-hls_list_size", "0",
"-hls_segment_filename", filepath.Join(opts.OuputDir, fmt.Sprintf("%s_%s_%s.ts", opts.UploadID, "%03d", prefix)),
filepath.Join(opts.OuputDir, fmt.Sprintf("%s_%s_master.m3u8", opts.UploadID, prefix)))
}
return result
}
var resolutions = []struct {
pixels int
label string
}{
{pixels: 640 * 480, label: "320p"},
{pixels: 1280 * 720, label: "480p"},
{pixels: 1920 * 1080, label: "720p"},
{pixels: 2560 * 1440, label: "1k"},
{pixels: 3840 * 2160, label: "2k"},
{pixels: 5120 * 2160, label: "4k"},
{pixels: 7680 * 4320, label: "5k"},
}
// ResolutionFromPixels converts pixel count to short string
func ResolutionFromPixels(pixels int) string {
idx := sort.Search(len(resolutions), func(i int) bool {
return pixels < resolutions[i].pixels
})
if idx == len(resolutions) {
return "8k"
}
return resolutions[idx].label
}
+62
View File
@@ -0,0 +1,62 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package transcoding_test
import (
"runtime"
"strings"
"testing"
"github.com/huly-stream/internal/pkg/transcoding"
"github.com/stretchr/testify/require"
)
func Test_BuildVideoCommand_Basic(t *testing.T) {
if runtime.GOOS == "windows" {
t.Skip()
}
var simpleHlsCommand = transcoding.BuildVideoCommand(&transcoding.Options{
OuputDir: "test",
UploadID: "1",
Threads: 4,
Resolutions: []string{"1280:720"},
})
const expected = `-threads 4 -i pipe:0 -vf scale=1280:720 -c:v libx264 -preset veryfast -crf 23 -g 60 -hls_time 5 -hls_list_size 0 -hls_segment_filename test/1_%03d_720p.ts test/1_720p_master.m3u8`
require.Contains(t, expected, strings.Join(simpleHlsCommand, " "))
}
func TestResolutionFromPixels(t *testing.T) {
tests := []struct {
pixels int
expected string
}{
{pixels: 320 * 240, expected: "320p"},
{pixels: 640 * 480, expected: "480p"},
{pixels: 1280 * 720, expected: "720p"},
{pixels: 1920 * 1080, expected: "1k"},
{pixels: 2560 * 1440, expected: "2k"},
{pixels: 3840 * 2160, expected: "4k"},
{pixels: 5120 * 2160, expected: "5k"},
{pixels: 9000 * 4000, expected: "8k"},
}
for _, tt := range tests {
t.Run(tt.expected, func(t *testing.T) {
result := transcoding.ResolutionFromPixels(tt.pixels)
require.Equal(t, tt.expected, result, "ResolutionFromPixels(%d)", tt.pixels)
})
}
}
+78
View File
@@ -0,0 +1,78 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package transcoding
import "sync/atomic"
// Limiter is a simple CAS data structure for managing resources.
type Limiter struct {
capacity int64
maxCapacity int64
}
// NewLimiter creates a new limiter with the given initial capacity.
func NewLimiter(capacity int64) *Limiter {
return &Limiter{
capacity: capacity,
maxCapacity: capacity,
}
}
// TryConsume attempts to consume the specified amount of capacity.
// Returns true if successful, false otherwise.
func (l *Limiter) TryConsume(amount int64) bool {
if amount <= 0 {
return false
}
for {
current := atomic.LoadInt64(&l.capacity)
if current < amount {
return false
}
updated := current - amount
if atomic.CompareAndSwapInt64(&l.capacity, current, updated) {
return true
}
}
}
// ReturnCapacity adds the specified amount back to the limiter's capacity.
// Does not exceed the maximum capacity.
func (l *Limiter) ReturnCapacity(amount int64) {
if amount <= 0 {
return
}
for {
current := atomic.LoadInt64(&l.capacity)
updated := current + amount
if updated > l.maxCapacity {
updated = l.maxCapacity
}
if atomic.CompareAndSwapInt64(&l.capacity, current, updated) {
break
}
}
}
// GetCapacity retrieves the current capacity for debugging or monitoring purposes.
func (l *Limiter) GetCapacity() int64 {
return atomic.LoadInt64(&l.capacity)
}
// GetMaxCapacity retrieves the maximum capacity.
func (l *Limiter) GetMaxCapacity() int64 {
return l.maxCapacity
}
+88
View File
@@ -0,0 +1,88 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package transcoding_test
import (
"sync"
"sync/atomic"
"testing"
"github.com/huly-stream/internal/pkg/transcoding"
"github.com/stretchr/testify/require"
)
func TestLimiter(t *testing.T) {
limiter := transcoding.NewLimiter(10)
t.Run("Initial capacity", func(t *testing.T) {
require.Equal(t, int64(10), limiter.GetCapacity())
})
t.Run("Successful consume", func(t *testing.T) {
success := limiter.TryConsume(5)
require.True(t, success)
require.Equal(t, int64(5), limiter.GetCapacity())
})
t.Run("Failed consume", func(t *testing.T) {
success := limiter.TryConsume(10)
require.False(t, success)
require.Equal(t, int64(5), limiter.GetCapacity())
})
t.Run("Return capacity", func(t *testing.T) {
limiter.ReturnCapacity(3)
require.Equal(t, int64(8), limiter.GetCapacity())
})
t.Run("Exceeding max capacity", func(t *testing.T) {
limiter.ReturnCapacity(10)
require.Equal(t, int64(10), limiter.GetCapacity())
})
}
func TestLimiterConcurrency(t *testing.T) {
limiter := transcoding.NewLimiter(10)
var wg sync.WaitGroup
for i := 0; i < 10; i++ {
wg.Add(1)
go func() {
defer wg.Done()
limiter.TryConsume(2)
}()
}
wg.Wait()
require.LessOrEqual(t, limiter.GetCapacity(), int64(0))
}
func TestLimiterCAS(t *testing.T) {
limiter := transcoding.NewLimiter(10)
var successful int64
var wg sync.WaitGroup
for i := 0; i < 1000; i++ {
wg.Add(1)
go func() {
defer wg.Done()
if limiter.TryConsume(1) {
atomic.AddInt64(&successful, 1)
}
}()
}
wg.Wait()
require.Equal(t, int64(10), successful)
require.Equal(t, int64(0), limiter.GetCapacity())
}
+129
View File
@@ -0,0 +1,129 @@
//
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
//
package transcoding
import (
"context"
"strings"
"sync"
"github.com/pkg/errors"
"github.com/google/uuid"
"github.com/huly-stream/internal/pkg/config"
"github.com/huly-stream/internal/pkg/log"
"github.com/huly-stream/internal/pkg/sharedpipe"
"github.com/huly-stream/internal/pkg/uploader"
"github.com/tus/tusd/v2/pkg/handler"
"go.uber.org/zap"
)
// Scheduler represents manager for worker. It creates a new worker for clients and manages its life cycle.
type Scheduler struct {
conf *config.Config
limiter *Limiter
mainContext context.Context
logger *zap.Logger
workers sync.Map
}
// NewScheduler creates a new scheduler for transcode operations.
func NewScheduler(ctx context.Context, c *config.Config) *Scheduler {
return &Scheduler{
conf: c,
limiter: NewLimiter(c.MaxCapacity),
mainContext: ctx,
logger: log.FromContext(ctx).With(zap.String("Scheduler", c.OutputDir)),
}
}
// NewUpload creates a new worker with passed parameters
func (s *Scheduler) NewUpload(ctx context.Context, info handler.FileInfo) (handler.Upload, error) {
if info.ID == "" {
info.ID = uuid.NewString()
}
s.logger.Debug("NewUpload", zap.String("ID", info.ID))
var result = &Worker{
done: make(chan struct{}),
writer: sharedpipe.NewWriter(),
info: info,
logger: log.FromContext(s.mainContext).With(zap.String("Worker", info.ID)),
}
var resolutions = strings.Split(info.MetaData["resolutions"], ",")
var commandOptions = Options{
OuputDir: s.conf.OutputDir,
Threads: s.conf.MaxThreads,
UploadID: info.ID,
Resolutions: resolutions,
}
result.cost = measure(&commandOptions)
if !s.limiter.TryConsume(result.cost) {
s.logger.Error("run out of resources")
return nil, errors.New("run out of resources")
}
if s.conf.EndpointURL != nil {
s.logger.Debug("found endpoint url in the config, starting uploader...")
var contentUploader, err = uploader.New(s.mainContext, *s.conf, info.ID, info.MetaData)
if err != nil {
return nil, err
}
result.contentUploader = contentUploader
go func() {
var serverErr = result.contentUploader.Serve()
result.logger.Debug("content uploader has finished", zap.Error(serverErr))
}()
}
s.workers.Store(result.info.ID, result)
s.logger.Sugar().Debugf("New Upload: info %v", result.info)
if err := result.start(s.mainContext, &commandOptions); err != nil {
return nil, err
}
return result, nil
}
// GetUpload returns current a worker based on upload id
func (s *Scheduler) GetUpload(ctx context.Context, id string) (upload handler.Upload, err error) {
if v, ok := s.workers.Load(id); ok {
s.logger.Debug("GetUpload: found worker by id", zap.String("id", id))
return v.(*Worker), nil
}
s.logger.Debug("GetUpload: worker not found", zap.String("id", id))
return nil, errors.New("bad id")
}
// AsTerminatableUpload returns tusd handler.TerminatableUpload
func (s *Scheduler) AsTerminatableUpload(upload handler.Upload) handler.TerminatableUpload {
var worker = upload.(*Worker)
s.logger.Debug("AsTerminatableUpload, trying to return capacity", zap.Int64("cost", worker.cost))
s.limiter.ReturnCapacity(worker.cost)
return worker
}
// AsLengthDeclarableUpload returns tusd handler.LengthDeclarableUpload
func (s *Scheduler) AsLengthDeclarableUpload(upload handler.Upload) handler.LengthDeclarableUpload {
s.logger.Debug("AsLengthDeclarableUpload")
return upload.(*Worker)
}
+125
View File
@@ -0,0 +1,125 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package transcoding provides objects and functions for video trnascoding
package transcoding
import (
"context"
"io"
"github.com/pkg/errors"
"github.com/huly-stream/internal/pkg/sharedpipe"
"github.com/huly-stream/internal/pkg/uploader"
"github.com/tus/tusd/v2/pkg/handler"
"go.uber.org/zap"
)
// Worker manages client's input and transcodes it based on the passsed configuration
type Worker struct {
contentUploader uploader.Uploader
logger *zap.Logger
info handler.FileInfo
writer *sharedpipe.Writer
reader *sharedpipe.Reader
cost int64
done chan struct{}
}
// WriteChunk calls when client sends a chunk of raw data
func (w *Worker) WriteChunk(ctx context.Context, _ int64, src io.Reader) (int64, error) {
w.logger.Debug("Write Chunk start", zap.Int64("offset", w.info.Offset))
var bytes, err = io.ReadAll(src)
_, _ = w.writer.Write(bytes)
var n = int64(len(bytes))
w.info.Offset += n
w.logger.Debug("Write Chunk end", zap.Int64("offset", w.info.Offset), zap.Error(err))
return n, err
}
// DeclareLength sets length of the video input
func (w *Worker) DeclareLength(ctx context.Context, length int64) error {
w.info.Size = length
w.info.SizeIsDeferred = false
w.logger.Debug("DeclareLength", zap.Int64("size", length), zap.Bool("SizeIsDeferred", w.info.SizeIsDeferred))
return nil
}
// GetInfo returns info about transcoing status
func (w *Worker) GetInfo(ctx context.Context) (handler.FileInfo, error) {
w.logger.Debug("GetInfo is executed")
return w.info, nil
}
// GetReader returns worker's bytes stream
func (w *Worker) GetReader(ctx context.Context) (io.ReadCloser, error) {
w.logger.Debug("GetReader is executed, creating current reader...")
return w.reader, nil
}
// Terminate calls when upload has failed
func (w *Worker) Terminate(ctx context.Context) error {
w.logger.Debug("Terminating...")
if w.contentUploader != nil {
go func() {
<-w.done
w.contentUploader.Rollback()
}()
}
return w.writer.Close()
}
// ConcatUploads calls when upload resumed after fail
func (w *Worker) ConcatUploads(ctx context.Context, partialUploads []handler.Upload) error {
w.logger.Debug("ConcatUploads was executed, it's not implemented")
//
// TODO: load raw source from the Buckup bucket, terminate all workers with same ID and start process again.
//
return errors.New("not implemented")
}
// FinishUpload calls when upload finished without errors on the client side
func (w *Worker) FinishUpload(ctx context.Context) error {
w.logger.Debug("finishing upload...")
if w.contentUploader != nil {
go func() {
<-w.done
w.contentUploader.Terminate()
}()
}
return w.writer.Close()
}
// AsConcatableUpload returns tusd handler.ConcatableUpload
func (s *Scheduler) AsConcatableUpload(upload handler.Upload) handler.ConcatableUpload {
s.logger.Debug("AsConcatableUpload is executed")
return upload.(*Worker)
}
func (w *Worker) start(ctx context.Context, options *Options) error {
w.reader = w.writer.Transpile()
var cmd, err = newFfmpegCommand(ctx, w.reader, options)
if err != nil {
return err
}
go func() {
defer close(w.done)
if runErr := cmd.Run(); runErr != nil {
w.logger.Error("transoding provider is exited with error", zap.Error(err))
} else {
w.logger.Debug("transoding provider has finished without errors")
}
}()
return nil
}
+130
View File
@@ -0,0 +1,130 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package uploader
import (
"bytes"
"context"
"io"
"mime/multipart"
"net/url"
"os"
"path/filepath"
"github.com/huly-stream/internal/pkg/log"
"github.com/pkg/errors"
"github.com/valyala/fasthttp"
"go.uber.org/zap"
)
// DatalakeStorage represents datalake storage
type DatalakeStorage struct {
baseURL string
workspace string
token string
}
// NewDatalakeStorage creates a new datalake client
func NewDatalakeStorage(baseURL, workspace, token string) Storage {
return &DatalakeStorage{
baseURL: "https://" + baseURL,
token: token,
workspace: workspace,
}
}
// UploadFile uploads file to the datalake
func (d *DatalakeStorage) UploadFile(ctx context.Context, fileName string) error {
// #nosec
file, err := os.Open(fileName)
if err != nil {
return err
}
defer func() {
_ = file.Close()
}()
var objectKey = getObjectKey(fileName)
var logger = log.FromContext(ctx).With(zap.String("datalake upload", d.workspace), zap.String("fileName", fileName))
logger.Debug("start uploading")
body := &bytes.Buffer{}
writer := multipart.NewWriter(body)
part, err := writer.CreateFormFile("file", objectKey)
if err != nil {
return errors.Wrapf(err, "failed to create form file")
}
_, err = io.Copy(part, file)
if err != nil {
return errors.Wrapf(err, "failed to copy file data")
}
_ = writer.Close()
req := fasthttp.AcquireRequest()
defer fasthttp.ReleaseRequest(req)
res := fasthttp.AcquireResponse()
defer fasthttp.ReleaseResponse(res)
req.SetRequestURI(d.baseURL + "/upload/form-data/" + d.workspace)
req.Header.SetMethod(fasthttp.MethodPost)
req.Header.Add("Authorization", "Bearer "+d.token)
req.Header.SetContentType(writer.FormDataContentType())
req.SetBody(body.Bytes())
client := fasthttp.Client{}
if err := client.Do(req, res); err != nil {
return errors.Wrapf(err, "upload failed")
}
logger.Debug("file uploaded")
return nil
}
// DeleteFile deletes file from the datalake
func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error {
var logger = log.FromContext(ctx).With(zap.String("datalake delete", d.workspace), zap.String("fileName", fileName))
logger.Debug("start deleting")
var objectKey = getObjectKey(fileName)
req := fasthttp.AcquireRequest()
defer fasthttp.ReleaseRequest(req)
res := fasthttp.AcquireResponse()
defer fasthttp.ReleaseResponse(res)
req.SetRequestURI(d.baseURL + "/blob/" + d.workspace + "/" + objectKey)
req.Header.SetMethod(fasthttp.MethodDelete)
req.Header.Add("Authorization", "Bearer "+d.token)
client := fasthttp.Client{}
if err := client.Do(req, res); err != nil {
return errors.Wrapf(err, "delete failed")
}
logger.Debug("file deleted")
return nil
}
func getObjectKey(s string) string {
var _, objectKey = filepath.Split(s)
return url.QueryEscape(objectKey)
}
+19
View File
@@ -0,0 +1,19 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package uploader
type options struct{}
// Option provides option for storages
type Option func(*options)
+40
View File
@@ -0,0 +1,40 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package uploader
import (
"context"
"time"
)
func (u *uploader) postpone(id string, action func()) {
var ctx, cancel = context.WithCancel(context.Background())
var startCh = time.After(u.postponeDuration)
if v, ok := u.contexts.Load(id); ok {
(*v.(*context.CancelFunc))()
}
u.contexts.Store(id, &cancel)
go func() {
defer cancel()
select {
case <-ctx.Done():
return
case <-startCh:
action()
u.contexts.CompareAndDelete(id, &cancel)
}
}()
}
+44
View File
@@ -0,0 +1,44 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package uploader
import (
"sync/atomic"
"testing"
"time"
"github.com/stretchr/testify/require"
)
func Test_Postpone(t *testing.T) {
var u = uploader{
postponeDuration: time.Second / 4,
}
var counter atomic.Int32
u.postpone("1", func() { counter.Add(1) })
time.Sleep(time.Second / 8)
u.postpone("1", func() { counter.Add(1) })
time.Sleep(time.Second / 2)
require.Equal(t, int32(1), counter.Load())
time.Sleep(time.Second / 2)
require.Equal(t, int32(1), counter.Load())
}
func Test_WithoutPostpone(t *testing.T) {
var counter atomic.Int32
var u uploader
u.postpone("1", func() { counter.Add(1) })
time.Sleep(time.Second / 10)
require.Equal(t, int32(1), counter.Load())
}
+143
View File
@@ -0,0 +1,143 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
package uploader
import (
"context"
"fmt"
"github.com/pkg/errors"
"os"
"path/filepath"
"strings"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/smithy-go"
"github.com/huly-stream/internal/pkg/log"
"go.uber.org/zap"
)
// S3Storage represents S3 storage
type S3Storage struct {
client *s3.Client
bucketName string
}
// NewS3 creates a new S3 storage
func NewS3(ctx context.Context, endpoint string) Storage {
var accessKeyID = os.Getenv("AWS_ACCESS_KEY_ID")
var accessKeySecret = os.Getenv("AWS_SECRET_ACCESS_KEY")
var bucketName = os.Getenv("AWS_BUCKET_NAME")
cfg, err := config.LoadDefaultConfig(ctx,
config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKeyID, accessKeySecret, "")),
config.WithRegion("auto"),
)
if err != nil {
panic(err.Error())
}
var s3Client = s3.NewFromConfig(cfg, func(o *s3.Options) {
endpoint = "https://" + endpoint
o.BaseEndpoint = &endpoint
})
return &S3Storage{
client: s3Client,
bucketName: bucketName,
}
}
func getContentType(objectKey string) string {
if strings.HasSuffix(objectKey, ".txt") {
return "txt"
}
if strings.HasSuffix(objectKey, ".ts") {
return "video/mp2t"
}
if strings.HasSuffix(objectKey, ".m3u8") {
return "application/x-mpegurl"
}
return "application/octet-stream"
}
// DeleteFile deletes file from the s3 storage
func (u *S3Storage) DeleteFile(ctx context.Context, fileName string) error {
var _, objectKey = filepath.Split(fileName)
var logger = log.FromContext(ctx).With(zap.String("s3 delete", u.bucketName), zap.String("fileName", fileName))
logger.Debug("start deleting")
input := &s3.DeleteObjectInput{
Bucket: aws.String(u.bucketName),
Key: aws.String(objectKey),
}
_, err := u.client.DeleteObject(ctx, input)
if err != nil {
return fmt.Errorf("failed to delete file from S3: %w", err)
}
logger.Debug("file deleted")
return nil
}
// UploadFile uploads file to the s3 storage
func (u *S3Storage) UploadFile(ctx context.Context, fileName string) error {
var _, objectKey = filepath.Split(fileName)
var logger = log.FromContext(ctx).With(zap.String("s3 upload", u.bucketName), zap.String("fileName", fileName))
logger.Debug("start upload file")
// #nosec
var file, err = os.Open(fileName)
if err != nil {
logger.Error("can not open file", zap.Error(err))
return err
}
defer func() {
_ = file.Close()
}()
_, err = u.client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(u.bucketName),
Key: aws.String(objectKey),
Body: file,
ContentType: aws.String(getContentType(objectKey)),
})
if err != nil {
var apiErr smithy.APIError
if errors.As(err, &apiErr) && apiErr.ErrorCode() == "EntityTooLarge" {
logger.Error("Error while uploading object. The object is too large." +
"To upload objects larger than 5GB, use the S3 console (160GB max)" +
"or the multipart upload API (5TB max).")
} else {
logger.Error("Couldn't upload file", zap.Error(err))
}
return apiErr
}
err = s3.NewObjectExistsWaiter(u.client).Wait(
ctx, &s3.HeadObjectInput{Bucket: aws.String(u.bucketName), Key: aws.String(objectKey)}, time.Minute)
if err != nil {
logger.Debug("Failed attempt to wait for object to exist.")
}
logger.Debug("file has uploaded")
return err
}
+219
View File
@@ -0,0 +1,219 @@
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
// Package uploader provides objects and functions to work with uploading and monitoring files
package uploader
import (
"context"
"net/url"
"os"
"strings"
"sync"
"time"
"github.com/pkg/errors"
"github.com/fsnotify/fsnotify"
"github.com/huly-stream/internal/pkg/config"
"github.com/huly-stream/internal/pkg/log"
"go.uber.org/zap"
)
type uploader struct {
ctx context.Context
cancel context.CancelFunc
baseDir string
uploadID string
masterFiles sync.Map
postponeDuration time.Duration
sentFiles sync.Map
storage Storage
contexts sync.Map
retryCount int
removeLocalContentOnUpload bool
eventBufferCount uint
isMasterFileFunc func(s string) bool
}
// Rollback deletes all delivered files and also deletes all local content by uploadID
func (u *uploader) Rollback() {
log.FromContext(u.ctx).Debug("cancel")
defer u.cancel()
u.sentFiles.Range(func(key, value any) bool {
log.FromContext(u.ctx).Debug("deleting remote file", zap.String("key", key.(string)))
for range u.retryCount {
var err = u.storage.DeleteFile(u.ctx, key.(string))
if err == nil {
break
}
log.FromContext(u.ctx).Debug("can not delete file", zap.Error(err))
}
return true
})
if !u.removeLocalContentOnUpload {
return
}
u.sentFiles.Range(func(key, value any) bool {
log.FromContext(u.ctx).Debug("deleting local file", zap.String("key", key.(string)))
_ = os.Remove(key.(string))
return true
})
}
func (u *uploader) Terminate() {
log.FromContext(u.ctx).Debug("terminate")
defer u.cancel()
u.masterFiles.Range(func(key, value any) bool {
log.FromContext(u.ctx).Debug("uploading master file", zap.String("key", key.(string)))
for range u.retryCount {
var uploadErr = u.storage.UploadFile(u.ctx, key.(string))
if uploadErr == nil {
break
}
log.FromContext(u.ctx).Debug("can not upload file", zap.Error(uploadErr))
}
return true
})
if !u.removeLocalContentOnUpload {
return
}
u.masterFiles.Range(func(key, value any) bool {
log.FromContext(u.ctx).Debug("deleting local master file", zap.String("key", key.(string)))
_ = os.Remove(key.(string))
return true
})
u.sentFiles.Range(func(key, value any) bool {
log.FromContext(u.ctx).Debug("deleting local file", zap.String("key", key.(string)))
_ = os.Remove(key.(string))
return true
})
}
func (u *uploader) Serve() error {
var logger = log.FromContext(u.ctx)
logger = logger.With(zap.String("uploader", u.uploadID), zap.String("dir", u.baseDir))
var watcher, err = fsnotify.NewBufferedWatcher(u.eventBufferCount)
if err != nil {
logger.Error("can not start watcher")
return err
}
if err := watcher.Add(u.baseDir); err != nil {
return err
}
defer func() {
_ = watcher.Close()
}()
logger.Debug("uploader initialized and started to watch")
for {
select {
case <-u.ctx.Done():
logger.Debug("done")
return u.ctx.Err()
case event, ok := <-watcher.Events:
if !strings.Contains(event.Name, u.uploadID) {
continue
}
if !ok {
return u.ctx.Err()
}
if u.isMasterFileFunc(event.Name) {
u.masterFiles.Store(event.Name, struct{}{})
logger.Debug("found master file", zap.String("eventName", event.Name))
continue
}
u.postpone(event.Name, func() {
logger.Debug("started to upload", zap.String("eventName", event.Name))
for range u.retryCount {
var uploadErr = u.storage.UploadFile(u.ctx, event.Name)
if uploadErr == nil {
break
}
logger.Error("can not upload file", zap.Error(uploadErr))
}
logger.Debug("added to sentFiles", zap.String("eventName", event.Name))
u.sentFiles.Store(event.Name, struct{}{})
})
case err, ok := <-watcher.Errors:
if !ok {
return u.ctx.Err()
}
logger.Error("get an error", zap.Error(err))
}
}
}
// Uploader manages content delivering
type Uploader interface {
Terminate()
Rollback()
Serve() error
}
// Storage represents file-based storage
type Storage interface {
UploadFile(ctx context.Context, fileName string) error
DeleteFile(ctx context.Context, fileName string) error
}
// New creates a new instance of Uplaoder
func New(ctx context.Context, conf config.Config, uploadID string, metadata map[string]string) (Uploader, error) {
var uploaderCtx, uploaderCancel = context.WithCancel(ctx)
var storage Storage
var err error
if conf.EndpointURL != nil {
storage, err = NewStorageByURL(ctx, conf.EndpointURL, metadata)
if err != nil {
uploaderCancel()
return nil, err
}
}
return &uploader{
ctx: uploaderCtx,
cancel: uploaderCancel,
uploadID: uploadID,
removeLocalContentOnUpload: conf.RemoveContentOnUpload,
postponeDuration: time.Second * 2,
storage: storage,
retryCount: 5,
baseDir: conf.OutputDir,
eventBufferCount: 100,
isMasterFileFunc: func(s string) bool {
return strings.HasSuffix(s, "m3u8")
},
}, nil
}
// NewStorageByURL creates a new storage basd on the type from the url scheme, for example "datalake://my-datalake-endpoint"
func NewStorageByURL(ctx context.Context, u *url.URL, headers map[string]string) (Storage, error) {
switch u.Scheme {
case "tus":
return nil, errors.New("not imlemented yet")
case "datalake":
if headers["workspace"] == "" {
return nil, errors.New("missed workspace in the client's metadata")
}
if headers["token"] == "" {
return nil, errors.New("missed auth token in the client's metadata")
}
return NewDatalakeStorage(u.Hostname(), headers["workspace"], headers["token"]), nil
case "s3":
return NewS3(ctx, u.Hostname()), nil
default:
return nil, errors.New("unknown scheme")
}
}