From de97f778f9f39e163cf455394d1596a1c28c773d Mon Sep 17 00:00:00 2001 From: denis-tingaikin Date: Tue, 4 Feb 2025 04:54:44 +0300 Subject: [PATCH 1/2] add initial version of huly-stream Signed-off-by: denis-tingaikin --- .DS_Store | Bin 0 -> 6148 bytes .github/workflows/docker-push.yaml | 44 +++ .github/workflows/main.yaml | 73 +++++ .github/yamllint.yaml | 12 + .golangci.yaml | 160 ++++++++++ Dockerfile | 36 +++ LICENSE | 277 ++++++++++++++++++ README.md | 122 +++++++- cmd/huly-stream/main.go | 104 +++++++ go.mod | 44 +++ go.sum | 84 ++++++ internal/pkg/config/config.go | 51 ++++ internal/pkg/log/zap.go | 50 ++++ internal/pkg/manifest/hls.go | 125 ++++++++ internal/pkg/manifest/hls_test.go | 143 +++++++++ internal/pkg/pprof/pprof.go | 54 ++++ internal/pkg/sharedpipe/shared_pipe.go | 116 ++++++++ .../pkg/sharedpipe/shared_pipe_bench_test.go | 248 ++++++++++++++++ internal/pkg/transcoding/command.go | 163 +++++++++++ internal/pkg/transcoding/command_test.go | 62 ++++ internal/pkg/transcoding/limiter.go | 78 +++++ internal/pkg/transcoding/limiter_test.go | 88 ++++++ internal/pkg/transcoding/scheduler.go | 129 ++++++++ internal/pkg/transcoding/worker.go | 125 ++++++++ internal/pkg/uploader/datalake.go | 124 ++++++++ internal/pkg/uploader/options.go | 19 ++ internal/pkg/uploader/postpone.go | 40 +++ internal/pkg/uploader/postpone_test.go | 44 +++ internal/pkg/uploader/s3.go | 143 +++++++++ internal/pkg/uploader/uploader.go | 219 ++++++++++++++ 30 files changed, 2976 insertions(+), 1 deletion(-) create mode 100644 .DS_Store create mode 100644 .github/workflows/docker-push.yaml create mode 100644 .github/workflows/main.yaml create mode 100644 .github/yamllint.yaml create mode 100644 .golangci.yaml create mode 100644 Dockerfile create mode 100644 LICENSE create mode 100644 cmd/huly-stream/main.go create mode 100644 go.mod create mode 100644 go.sum create mode 100644 internal/pkg/config/config.go create mode 100644 internal/pkg/log/zap.go create mode 100644 internal/pkg/manifest/hls.go create mode 100644 internal/pkg/manifest/hls_test.go create mode 100644 internal/pkg/pprof/pprof.go create mode 100644 internal/pkg/sharedpipe/shared_pipe.go create mode 100644 internal/pkg/sharedpipe/shared_pipe_bench_test.go create mode 100644 internal/pkg/transcoding/command.go create mode 100644 internal/pkg/transcoding/command_test.go create mode 100644 internal/pkg/transcoding/limiter.go create mode 100644 internal/pkg/transcoding/limiter_test.go create mode 100644 internal/pkg/transcoding/scheduler.go create mode 100644 internal/pkg/transcoding/worker.go create mode 100644 internal/pkg/uploader/datalake.go create mode 100644 internal/pkg/uploader/options.go create mode 100644 internal/pkg/uploader/postpone.go create mode 100644 internal/pkg/uploader/postpone_test.go create mode 100644 internal/pkg/uploader/s3.go create mode 100644 internal/pkg/uploader/uploader.go diff --git a/.DS_Store b/.DS_Store new file mode 100644 index 0000000000000000000000000000000000000000..0923990b38bec54c6cb023fcb61ce6be5c8afe41 GIT binary patch literal 6148 zcmeHKyG{c^3>-s>NNG}1?k4~Z?J7#XARhn)5)DcsM5wRgyZE$>A3{1^C{m#)|UB-M9`4L$qUJ iv}10(9p6P!)-_-AycZ6ML1#YbMEwl7E;1=_Z3RxGZ4~VQ literal 0 HcmV?d00001 diff --git a/.github/workflows/docker-push.yaml b/.github/workflows/docker-push.yaml new file mode 100644 index 0000000000..0308ce654b --- /dev/null +++ b/.github/workflows/docker-push.yaml @@ -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 }} diff --git a/.github/workflows/main.yaml b/.github/workflows/main.yaml new file mode 100644 index 0000000000..a288973049 --- /dev/null +++ b/.github/workflows/main.yaml @@ -0,0 +1,73 @@ +--- +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 + - windows + - macos + 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 diff --git a/.github/yamllint.yaml b/.github/yamllint.yaml new file mode 100644 index 0000000000..7c94f8d8fe --- /dev/null +++ b/.github/yamllint.yaml @@ -0,0 +1,12 @@ +--- +extends: default + +yaml-files: + - '*.yaml' + - '*.yml' + +rules: + truthy: disable + line-length: disable + comments: + min-spaces-from-content: 1 diff --git a/.golangci.yaml b/.golangci.yaml new file mode 100644 index 0000000000..0099487e6f --- /dev/null +++ b/.golangci.yaml @@ -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 diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000000..0b4dca5d82 --- /dev/null +++ b/Dockerfile @@ -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"] \ No newline at end of file diff --git a/LICENSE b/LICENSE new file mode 100644 index 0000000000..e48e096345 --- /dev/null +++ b/LICENSE @@ -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. diff --git a/README.md b/README.md index 1b7f8dd774..8eb361d442 100644 --- a/README.md +++ b/README.md @@ -1 +1,121 @@ -# 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 +- **TUS Upload**: Resumable file uploads via TUS protocol. +- **s3 Upload**: Direct upload to Amazon S3. +- **datalake Upload**: Integration for data lake storage systems. +- **Supported Output Formats**: + - `aac` + - `hls` + +### 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. + + + +#### 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: " \ + --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! 🚀 \ No newline at end of file diff --git a/cmd/huly-stream/main.go b/cmd/huly-stream/main.go new file mode 100644 index 0000000000..2730cd8b35 --- /dev/null +++ b/cmd/huly-stream/main.go @@ -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 +} diff --git a/go.mod b/go.mod new file mode 100644 index 0000000000..89419cb728 --- /dev/null +++ b/go.mod @@ -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 +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000000..11e9c85d56 --- /dev/null +++ b/go.sum @@ -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= diff --git a/internal/pkg/config/config.go b/internal/pkg/config/config.go new file mode 100644 index 0000000000..1ef90f2192 --- /dev/null +++ b/internal/pkg/config/config.go @@ -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 +} diff --git a/internal/pkg/log/zap.go b/internal/pkg/log/zap.go new file mode 100644 index 0000000000..72e03ac379 --- /dev/null +++ b/internal/pkg/log/zap.go @@ -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) +} diff --git a/internal/pkg/manifest/hls.go b/internal/pkg/manifest/hls.go new file mode 100644 index 0000000000..712e97f072 --- /dev/null +++ b/internal/pkg/manifest/hls.go @@ -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 +} diff --git a/internal/pkg/manifest/hls_test.go b/internal/pkg/manifest/hls_test.go new file mode 100644 index 0000000000..15f3dbe022 --- /dev/null +++ b/internal/pkg/manifest/hls_test.go @@ -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) + } + }) + } +} diff --git a/internal/pkg/pprof/pprof.go b/internal/pkg/pprof/pprof.go new file mode 100644 index 0000000000..e7a5eb0302 --- /dev/null +++ b/internal/pkg/pprof/pprof.go @@ -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() +} diff --git a/internal/pkg/sharedpipe/shared_pipe.go b/internal/pkg/sharedpipe/shared_pipe.go new file mode 100644 index 0000000000..0a698b46c8 --- /dev/null +++ b/internal/pkg/sharedpipe/shared_pipe.go @@ -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) diff --git a/internal/pkg/sharedpipe/shared_pipe_bench_test.go b/internal/pkg/sharedpipe/shared_pipe_bench_test.go new file mode 100644 index 0000000000..3633e1fdf4 --- /dev/null +++ b/internal/pkg/sharedpipe/shared_pipe_bench_test.go @@ -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) + } + } +} diff --git a/internal/pkg/transcoding/command.go b/internal/pkg/transcoding/command.go new file mode 100644 index 0000000000..f43b5cf594 --- /dev/null +++ b/internal/pkg/transcoding/command.go @@ -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 +} diff --git a/internal/pkg/transcoding/command_test.go b/internal/pkg/transcoding/command_test.go new file mode 100644 index 0000000000..14a7db515d --- /dev/null +++ b/internal/pkg/transcoding/command_test.go @@ -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) + }) + } +} diff --git a/internal/pkg/transcoding/limiter.go b/internal/pkg/transcoding/limiter.go new file mode 100644 index 0000000000..f1fdf51d9f --- /dev/null +++ b/internal/pkg/transcoding/limiter.go @@ -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 +} diff --git a/internal/pkg/transcoding/limiter_test.go b/internal/pkg/transcoding/limiter_test.go new file mode 100644 index 0000000000..df7db82ba0 --- /dev/null +++ b/internal/pkg/transcoding/limiter_test.go @@ -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()) +} diff --git a/internal/pkg/transcoding/scheduler.go b/internal/pkg/transcoding/scheduler.go new file mode 100644 index 0000000000..25872a3047 --- /dev/null +++ b/internal/pkg/transcoding/scheduler.go @@ -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) +} diff --git a/internal/pkg/transcoding/worker.go b/internal/pkg/transcoding/worker.go new file mode 100644 index 0000000000..60baccc465 --- /dev/null +++ b/internal/pkg/transcoding/worker.go @@ -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 +} diff --git a/internal/pkg/uploader/datalake.go b/internal/pkg/uploader/datalake.go new file mode 100644 index 0000000000..6869cc9f99 --- /dev/null +++ b/internal/pkg/uploader/datalake.go @@ -0,0 +1,124 @@ +// 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" + "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 = filepath.Split(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 = filepath.Split(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 +} diff --git a/internal/pkg/uploader/options.go b/internal/pkg/uploader/options.go new file mode 100644 index 0000000000..076e2c70af --- /dev/null +++ b/internal/pkg/uploader/options.go @@ -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) diff --git a/internal/pkg/uploader/postpone.go b/internal/pkg/uploader/postpone.go new file mode 100644 index 0000000000..61b22b79b4 --- /dev/null +++ b/internal/pkg/uploader/postpone.go @@ -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) + } + }() +} diff --git a/internal/pkg/uploader/postpone_test.go b/internal/pkg/uploader/postpone_test.go new file mode 100644 index 0000000000..89522b4cec --- /dev/null +++ b/internal/pkg/uploader/postpone_test.go @@ -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()) +} diff --git a/internal/pkg/uploader/s3.go b/internal/pkg/uploader/s3.go new file mode 100644 index 0000000000..ba688618f9 --- /dev/null +++ b/internal/pkg/uploader/s3.go @@ -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 +} diff --git a/internal/pkg/uploader/uploader.go b/internal/pkg/uploader/uploader.go new file mode 100644 index 0000000000..86735191a4 --- /dev/null +++ b/internal/pkg/uploader/uploader.go @@ -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") + } +} From 1b8be96319b735d0196738772a6ff65d68f7a29e Mon Sep 17 00:00:00 2001 From: denis-tingaikin Date: Tue, 4 Feb 2025 05:15:28 +0300 Subject: [PATCH 2/2] apply review comments Signed-off-by: denis-tingaikin --- .github/workflows/main.yaml | 2 -- README.md | 16 +++++++++------- internal/pkg/uploader/datalake.go | 10 ++++++++-- 3 files changed, 17 insertions(+), 11 deletions(-) diff --git a/.github/workflows/main.yaml b/.github/workflows/main.yaml index a288973049..635804cf49 100644 --- a/.github/workflows/main.yaml +++ b/.github/workflows/main.yaml @@ -29,8 +29,6 @@ jobs: matrix: os: - ubuntu - - windows - - macos runs-on: ${{ matrix.os }}-latest steps: - name: Check out code diff --git a/README.md b/README.md index 8eb361d442..f977167dba 100644 --- a/README.md +++ b/README.md @@ -20,13 +20,15 @@ The Huly Stream high-performance HTTP-based transcoding service. Huly-stream is - `webm` ### Output Options -- **TUS Upload**: Resumable file uploads via TUS protocol. -- **s3 Upload**: Direct upload to Amazon S3. -- **datalake Upload**: Integration for data lake storage systems. - **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. @@ -74,13 +76,13 @@ STREAM_REMOVE_CONTENT_ON_UPLOAD True or False true STREAM_UPLOAD_RAW_CONTENT True or False false uploads content in raw quality to the endpoint if true ``` -### Metadata: +### Metadata -**resolutions** if passed, set the resolution for the output, for example, 'resolutions: 1920:1080, 1280:720.' +**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. +**token:** must be provided to be authorized in the Huly's datalake service. -**workspace** required for uploading content. +**workspace:** required for uploading content to the datalake storage. diff --git a/internal/pkg/uploader/datalake.go b/internal/pkg/uploader/datalake.go index 6869cc9f99..43ba535efd 100644 --- a/internal/pkg/uploader/datalake.go +++ b/internal/pkg/uploader/datalake.go @@ -18,6 +18,7 @@ import ( "context" "io" "mime/multipart" + "net/url" "os" "path/filepath" @@ -54,7 +55,7 @@ func (d *DatalakeStorage) UploadFile(ctx context.Context, fileName string) error _ = file.Close() }() - var _, objectKey = filepath.Split(fileName) + var objectKey = getObjectKey(fileName) var logger = log.FromContext(ctx).With(zap.String("datalake upload", d.workspace), zap.String("fileName", fileName)) logger.Debug("start uploading") @@ -101,7 +102,7 @@ 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 = filepath.Split(fileName) + var objectKey = getObjectKey(fileName) req := fasthttp.AcquireRequest() defer fasthttp.ReleaseRequest(req) @@ -122,3 +123,8 @@ func (d *DatalakeStorage) DeleteFile(ctx context.Context, fileName string) error return nil } + +func getObjectKey(s string) string { + var _, objectKey = filepath.Split(s) + return url.QueryEscape(objectKey) +}