Темы



Streamline: настраиваемые видеопроцессы с Cloudflare Stream и Workers

Коротко сгенерировано ИИ по тексту статьи
  • Cloudflare выпустила Streamline — площадку для создания конвейеров обработки видео на Workers, Containers и Durable Objects.
  • Медиадвижок в Container принимает RTMP, HLS или видео с веб-камеры и выдаёт RTMP-поток либо предпросмотр по WebSocket; сеанс работает независимо от запросов Worker.
  • API позволяет настраивать операции обработки, включая наложение изображений, фильтры и субтитры; в текущей реализации используется FFmpeg.

Почему это важно: Streamline позволяет создавать динамические аннотации для прямых трансляций и альтернативные версии видео со встроенными субтитрами.

Cloudflare Stream — это мощная платформа для трансляций , которая просто работает для многих наших клиентов. Но что, если вам понадобится выводить динамические аннотации поверх прямой трансляции или создать альтернативную версию размещённого видео со встроенными субтитрами? Для этого потребуется собственный конвейер обработки видео.

Сегодня мы выпускаем новую площадку для разработчиков — Streamline. Она показывает, как создать систему для предоставления таких персонализированных возможностей работы с видео на платформе Cloudflare для разработчиков. Мы расскажем, как Streamline использует Workers, Containers и несколько медиапротоколов для изменения видео и немедленной публикации результата в виде прямой трансляции или нового размещённого видео. Вы также сможете опробовать систему в своих проектах.

Для конвейера обработки нужна надёжная среда с длительным временем работы, способная выполнять специализированный скомпилированный код с предсказуемым объёмом памяти и вычислительных ресурсов. Обработка видеопотока может занимать минуты и часы, поэтому жизненный цикл медиапроцесса должен быть независим от запроса, который его запустил. Приложение должно иметь возможность запустить конвейер, передать ему входные данные, проверить его работу и остановить, не удерживая один запрос открытым всё это время.

В Cloudflare есть необходимые для этого базовые компоненты. Containers — это среды выполнения с длительным сроком работы, подходящие для обработки медиа. Durable Objects помогают управлять процессом. А Workers отлично подходят для передачи управляющих сигналов и мониторинга.

Для Streamline мы создали медиадвижок, работающий в Container и обрабатывающий медиа в реальном времени. Container управляется через Worker, который предоставляет агенту или пользователю средства управления, предварительного просмотра и тестирования. Обработка продолжается, даже если Worker отключается. Streamline построен из модульных компонентов, поэтому в будущем медиадвижок можно будет заменить специализированными продуктами для кодирования.

Архитектура

Развёртывание Streamline состоит из двух компонентов: медиадвижка, которая отвечает за ввод, вывод и обработку медиа, и управляющего приложения, которое создаёт, настраивает, отслеживает и останавливает медиасеансы.

Медиадвижок

Медиадвижок состоит из двух компонентов:

  • Контроллер. Это управляющая оболочка, написанная на Go. Она реализует HTTP-сервер, принимает входящие запросы и преобразует их в операции, которые может выполнить медиадвижок.
  • Процессор выполняет собственно обработку медиа. В текущей реализации используется FFmpeg, но это внутренняя деталь реализации, а не часть API, доступного пользователям.

Медиадвижок размещается в Container и отвечает за весь ввод, вывод и обработку медиа. Он может получать RTMPS-поток по сети с одного входа Stream Live и публиковать выходной RTMPS-поток на другой вход Stream Live. Он может получать HLS-манифест Cloudflare Stream и его сегменты, используя размещённые видео в качестве входных данных. Он может принимать видео от источника, указанного управляющим приложением, например веб-камеры. Он может передавать видео для предварительного просмотра по исходящему WebSocket в ретранслятор Durable Object. Приложение, которому нужен предпросмотр, может подключиться к этому ретранслятору через собственный WebSocket.

Приложение 

Приложение построено на Workers и может быть полнофункциональным браузерным приложением, агентом или встроенной системой. Оно состоит из следующих компонентов:

  • Пользовательский интерфейс (UI) включает клиентскую логику, управление идентификацией и политикой доступа. В этой статье в качестве конкретного примера рассматривается браузерное приложение, поэтому в него также входит интерфейс для браузера.
  • Оркестратор управляет сеансом, жизненным циклом Container и ретрансляцией видео для предпросмотра. Оркестратор реализован с помощью Durable Object.

Во время разработки систему можно запускать локально. В этом случае контейнером служит локальный экземпляр Docker, а Durable Object не используется: пользователь только один, управляющему приложению не нужна авторизация для локального доступа, а видео для предпросмотра может подключаться напрямую к WebSocket на localhost.

Когда эти компоненты развёрнуты в Cloudflare, авторизованный пользователь или агент может открыть Worker и запустить новый сеанс. При необходимости запускается новый контейнер Streamline, его жизненный цикл управляется автоматически, предоставляется API для выполнения различных операций обработки видео, а входные данные из Cloudflare Stream и обработанный результат направляются обратно в Cloudflare Stream.

Пора подробно разобраться в технических аспектах работы системы.

Жизненный цикл контейнера и управление сеансами

Управляющее приложение Worker запускает длительный сеанс обработки медиа. После запуска сеанса приложение может безопасно отключаться и подключаться снова, пока Container продолжает обработку — до тех пор, пока управляющее приложение его не остановит. Мы также задаём максимальную длительность, чтобы сеанс в любом случае в итоге завершался и не мог работать бесконечно, даже без внешнего управления. Пока идёт сеанс обработки медиа, экземпляр контейнера недоступен для других приложений.

Cloudflare Container автоматически переходит в спящий режим, если в течение заданного интервала не получает входящих запросов. Однако в нашем случае работающий конвейер должен продолжать выполняться, даже если управляющее приложение отключится и запросы перестанут поступать. Такого поведения можно добиться, переопределив обратный вызов onActivityExpired() контейнера. Если время истечения срока ещё не наступило, мы продлеваем активность, в противном случае уничтожаем контейнер.

async onActivityExpired() {
  await this.withControlLock(async () => {
    const session = await this.getRelaySessionLocked()

    if (session?.expiresAt) {
      this.renewActivityTimeout()
      return
    }

    await this.destroy()
  })
}

API

HTTP-сервер, реализованный управляющей оболочкой на Go, и связанный с Container Durable Object вместе определяют низкоуровневый интерфейс системы. Однако мы хотели создать над ним абстракцию, чтобы система как можно меньше зависела от того, кто или что управляет сеансом, и от ненужных деталей внутренней реализации серверной части.

Для этого мы экспортируем из Streamline два пакета:

  • @cloudflare/streamline/client Предоставляет высокоуровневый API на основе сеансов.
  • @cloudflare/streamline/ Предоставляет базовый класс Durable Object, связанного с контейнером. Он направляет запросы API, реализует описанный ниже сервер ретрансляции предпросмотра и предоставляет точки расширения для настройки безопасности и политики доступа.

При удалённом развёртывании управляющий Worker должен импортировать @streamline/cloudflare и определить конкретный подкласс Durable Object, предоставляемого контейнером. Его можно использовать для логики и хранения данных, специфичных для приложения.

В локальном режиме, когда Durable Object отсутствует, во фронтенде определяется тонкий адаптерный слой. Он сохраняет API на основе сеансов, но подключается напрямую к локальному экземпляру Docker без контроля доступа и прочих ограничений.

В примере ниже показано, как управляющее приложение может использовать API для доступа к Streamline, подготовки сеанса и запуска конвейера обработки видео.

const streamline = createStreamline({ baseUrl: 'https://media.example' })
const session = await streamline.sessions.create()
const result = await session.start(config)

// At this point the pipeline is running, unless failure occurred.
console.log(result)

config — это JSON-объект, задающий конвейер обработки, который будет запущен; подробнее об этом говорится в следующих разделах.

В таблице ниже приведён полный список вызовов API.

Метод клиента

Функция

createStreamline()

Создаёт новый экземпляр Streamline.

streamline.sessions.create()

Создаёт новый сеанс обработки.

streamline.sessions.resume(id)

Повторно подключается к существующему сеансу.

session.start(config)

Запускает новый конвейер обработки.

session.ingest(chunk)

Отправляет фрагмент видеоданных в режиме «веб-камера».

session.annotation(png)

Обновляет прозрачный слой аннотаций.

session.metrics()

Получает метрики текущего сеанса.

session.stop()

Останавливает обработку в текущем сеансе.

Определение и запуск конвейера обработки видео

session.start() создаёт и запускает конвейер обработки. Он принимает один аргумент — объект конфигурации JSON, определяющий выполняемые операции:

  • Входные данные
  • Операции
  • Выходные данные

В примере ниже запускается конвейер, который принимает в качестве входных данных трансляцию RTMP (протокол обмена сообщениями в реальном времени), например поток с входа Stream Live, получающего прямую трансляцию, накладывает изображение с прозрачностью и отправляет результат на RTMP-приёмник, например на другой вход Stream Live для записи или трансляции. Это позволяет приложению Worker в реальном времени создавать изменённую версию прямой трансляции.

const session = await streamline.sessions.create()

const result = await session.start({
  input: { type: 'rtmp', profile: 'primary-input' },
  pipeline: [
    {
      op: 'overlay',
      params: { image: '/app/assets/cf-logo.png', position: 'top-right' },
    },
    {
      op: 'encode',
      params: {
        codec: 'h264',
        preset: 'fast',
        bitrate: '1500k',
        resolution: '1280x720',
        fps: 30,
      },
    },
  ],
  output: { mode: 'rtmp', profile: 'primary-output' },
})

Ввод видео по запросу через HLS

Streamline также может принимать потоковое видео по HLS (HTTP Live Streaming), например видео, размещённое на Cloudflare Stream. В примере ниже показано, как приложение Worker может запустить конвейер, который получает видео из Stream, считывает встроенные субтитры для людей с нарушениями слуха, отображает их поверх видео и отправляет результат по RTMP — например, на вход Stream Live для трансляции или записи изменённой версии.

const streamVideoId = 'your-cloudflare-stream-video-id'
const session = await streamline.sessions.create()

await session.start({
  input: {
    type: 'hls',
    url: `https://videodelivery.net/${streamVideoId}/manifest/video.m3u8`,
  },
  pipeline: [
    { op: 'subtitle', params: { source: 'auto' } },
    {
      op: 'encode',
      params: {
        codec: 'h264',
        preset: 'fast',
        bitrate: '1500k',
        resolution: '1280x720',
        fps: 30,
      },
    },
  ],
  output: { mode: 'rtmp', profile: 'default' },
})

Передача видео в Streamline

Часто бывает полезно быстро проверить конвейер обработки, отправив видеоданные напрямую в Streamline, например с веб-камеры. Агенту или приложению на встроенном устройстве эта возможность тоже может пригодиться: например, чтобы отправлять записи с заводских камер на анализ ИИ или объединять потоки с нескольких камер в единое изображение.

В примере ниже создаётся конвейер, который ожидает входные данные от приложения Worker и формирует выходное видео для предпросмотра, доступное по WebSocket (подробнее о видео для предпросмотра через WebSocket рассказывается ниже). Конвейер применяет два фильтра и «аннотацию» — наложенное PNG-изображение, которое можно обновлять во время обработки, например для создания анимированной графики.

const session = await streamline.sessions.create()
const sessionId = session.id
if (!sessionId) throw new Error('Session creation returned no session ID')

const viewer = await openViewer('https://media.example', sessionId)

await session.start({
  input: { type: 'webcam' },
  pipeline: [
    { op: 'filter', params: { preset: 'brightness', amount: 0.1 } },
    { op: 'filter', params: { preset: 'flip' } },
    { op: 'overlay', params: { image: 'annotation', position: 'full' } },
    {
      op: 'encode',
      params: {
        codec: 'h264',
        preset: 'veryfast',
        bitrate: '1500k',
        resolution: '1280x720',
        fps: 30,
        gop: 60,
      },
    },
  ],
  output: { mode: 'websocket', format: 'fmp4' },
})

Приведённый выше фрагмент кода только запускает конвейер. Управляющий Worker пока не отправляет в Streamline медиаданные. Ниже мы рассмотрим функцию openViewer().

Приложение Worker отправляет видеоданные в Streamline с помощью вызова session.ingest(). В примере ниже показано, как браузерное приложение может получать фрагменты с веб-камеры и пересылать их в Streamline.

const stream = await navigator.mediaDevices.getUserMedia({ video: true, audio: true })
const recorder = new MediaRecorder(stream, { mimeType: 'video/webm;codecs=vp8,opus' })
let uploadTail = Promise.resolve()

recorder.addEventListener('dataavailable', (event) => {
  if (event.data.size === 0) return
  uploadTail = uploadTail
    .then(() => session.ingest(event.data))
    .catch((error) => reportUploadFailure(error))
})

recorder.start(250)

Анимированный слой наложения

Слой аннотаций можно обновлять с помощью вызова session.annotation(). В примере ниже показано, как приложение Worker может делать снимок холста и отправлять его в Streamline. Это можно выполнять в цикле анимации, хотя на практике частота обновления может быть ограничена размером PNG-изображений для наложения, доступной пропускной способностью и вычислительной мощностью.

function canvasPng(canvas: HTMLCanvasElement): Promise<Blob> {
  return new Promise((resolve, reject) => {
    canvas.toBlob((blob) => {
      if (blob) resolve(blob)
      else reject(new Error('Canvas could not produce a PNG'))
    }, 'image/png')
  })
}

const png = await canvasPng(overlayCanvas)
await session.annotation(png)

Получение видео для предпросмотра из Streamline

Streamline также может формировать видео для предпросмотра, если указать output: { mode: 'websocket' }.

Streamline использует WebSocket для передачи видео предпросмотра в управляющее приложение с низкой задержкой: контейнер публикует фрагменты fMP4 в Durable Object, который пересылает их в выходной ретранслятор, доступный по WebSocket по URL /relay/view относительно источника приложения. Приложение должно подключиться к этому URL через WebSocket и затем будет получать видеоданные по мере их поступления из Streamline. В приведённом ниже фрагменте кода показано, как браузерное приложение может отображать видеопоток для предпросмотра.

function openViewer(origin: string, sessionId: string): Promise<WebSocket> {
  const url = new URL('/relay/view', origin)
  url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:'
  url.searchParams.set('session_id', sessionId)

  return new Promise((resolve, reject) => {
    const socket = new WebSocket(url)
    socket.binaryType = 'arraybuffer'
    socket.addEventListener('open', () => resolve(socket), { once: true })
    socket.addEventListener('error', () => reject(new Error('Preview relay failed')), { once: true })
  })
}

const sessionId = session.id
if (!sessionId) throw new Error('Session creation returned no session ID')

// Connect before session.start(), or the relay rejects the publisher.
const viewer = await openViewer('https://media.example', sessionId)

viewer.addEventListener('message', (event) => {
  if (typeof event.data === 'string') {
    if (event.data === '{"type":"eos"}') mediaSource.endOfStream()
    return
  }
  sourceBuffer.appendBuffer(new Uint8Array(event.data as ArrayBuffer))
})

В производственном проигрывателе MediaSource необходимо ставить фрагменты в очередь, пока SourceBuffer.updating имеет значение true. При локальной разработке браузер или другое управляющее приложение просто устанавливает WebSocket-соединение напрямую с локальным контейнером.

Поддерживаемые на данный момент операции

В объекте конфигурации, передаваемом в session.start() в приведённых выше примерах, pipeline — это массив операций из набора, поддерживаемого базовым медиадвижком. Порядок операций сейчас задаётся самим движком; порядок элементов в массиве значения не имеет. Ниже приведён список поддерживаемых на данный момент операций в том порядке, в котором они применяются.

Название операции

Функция

filter

Применяет фильтры, например размытие или настройку насыщенности.

overlay

Накладывает изображение, указанное по URL, либо двоичные данные PNG, переданные отдельно в вызове annotation().

subtitle

Встраивает субтитры в видео.

encode

Задаёт параметры кодирования выходных данных.

Безопасность

Это Cloudflare, поэтому вопросы безопасности должны учитываться при проектировании системы, а не добавляться в самом конце. Необходимо гарантировать, что только авторизованные пользователи могут создавать сеансы или перехватывать управление существующими, а сеансы изолированы друг от друга. Ключи входных и выходных потоков Stream RTMPS нужно считать секретными и не раскрывать управляющему приложению. Также необходимо ограничить использование ресурсов.

Развёртывание владельца закрыто благодаря интеграции с Access для Workers. Владелец с настроенной учётной записью и другие разрешённые пользователи могут редактировать общие профили и запускать сеанс, пока единственный экземпляр простаивает. Перед приёмом запросов на управление Worker проверяет сеанс Access и связывает активный сеанс с подтверждённым субъектом. Одновременно может выполняться только один сеанс, а другой субъект не может остановить или заменить активный сеанс.

Ключи входов Stream Live хранятся в секретах Worker или в виде общих параметров, доступных только для записи, в хранилище Durable Object. Они никогда не возвращаются через API настроек и не сохраняются в браузере. Управляющее приложение указывает входы и выходы RTMPS по имени профиля. Worker находит профиль, прежде чем связаться с контейнером.

Для потока видео предпросмотра используются два учётных данных с разными назначениями. Сервисный токен Cloudflare Access подтверждает подлинность рабочей нагрузки контейнера перед конечной точкой публикации. Случайно сгенерированный токен для каждого сеанса разрешает публикацию только в текущий активный ретранслятор. Сервисный токен передаётся исходящим Worker контейнера и никогда не попадает в память контейнера. В первоначальном развёртывании временно используется обход Access для конкретного пути, при этом проверка токена сеанса сохраняется. После развёртывания и успешной проверки работоспособности обход заменяется на Service Auth.

Развёртывание владельца намеренно закрыто и направляется на единственный экземпляр. Это не модель безопасности для общедоступного многопользовательского сервиса.

Площадка для экспериментов и открытый исходный код

Мы хотим, чтобы вы попробовали Streamline и начали создавать свои решения! Поэтому вместе с этой статьёй мы публикуем исходный код системы и развёртываем общедоступную площадку для экспериментов.

Контейнер Streamline можно запустить локально или развернуть в своей учётной записи. Он предоставляет API Worker для управляющего приложения.

Также доступно примерное приложение Worker с веб-интерфейсом на Astro, демонстрирующее возможности Streamline на нескольких распространённых примерах: наложение изображений, декодирование субтитров, фильтры и «картинка в картинке». Функция проверки предоставляет метрики производительности и системную трассировку, которые могут пригодиться для отладки при разработке новых функций. Примерное приложение можно запустить на локальном сервере Astro или развернуть за Cloudflare Access, чтобы контролировать доступ к вашему экземпляру Streamline.

Оба репозитория опубликованы с открытым исходным кодом на GitHub Cloudflare:

Мы развернули общедоступную версию примерного приложения. При желании пользователь может развернуть её и у себя. В ней используется собственная конфигурация Access, отдельная идентичность контейнера для каждого подтверждённого пользователя, один активный сеанс на пользователя, глобальный контроль допуска, ограничения на параллельную работу, объём медиа и длительность сеанса. Пользователи не могут заменять сеансы друг друга.

Попробовать общедоступную площадку можно здесь:

Что дальше

Streamline показывает, как можно объединить существующие управляемые сервисы, например Stream, с низкоуровневыми базовыми компонентами и создавать гибко настраиваемые конвейеры обработки медиа. В текущей версии Streamline использует ресурсы CPU контейнера для обработки медиа, что становится узким местом при повышении качества или частоты кадров.

В дальнейшем нам интересно посмотреть, как мы вместе с сообществом разработчиков сможем расширить эту архитектуру: добавить поддержку конвейеров компьютерного зрения, аппаратно ускоренной обработки медиа и приложений реального времени на новых протоколах, таких как WebRTC и MoQ, а в конечном счёте — встроить средства кодирования и декодирования видео непосредственно в Workers.

Сегодня мы приглашаем вас ознакомиться с размещённой у нас демонстрацией Streamline и оценить возможности этих инструментов. Затем загляните в опубликованные нами исходные коды, чтобы увидеть, как легко развернуть Streamline в своей учётной записи и использовать его для создания собственных решений.

© Cloudflare Blog