Skip to content

Commit a84f586

Browse files
authored
Merge pull request #18 from mondain/feat/switch-stream-per-obj
added --stream-per-object flag
2 parents 6c1db80 + b4ef3b1 commit a84f586

7 files changed

Lines changed: 18 additions & 8 deletions

File tree

include/openmoq/publisher/cli_options.h

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ struct CliOptions {
4242
bool include_sap = false;
4343
bool include_msf_timeline = false;
4444
bool split_cmaf_chunks = true;
45+
bool stream_per_object = false;
4546
bool paced = false;
4647
bool loop = false;
4748
bool dump_plan = false;

include/openmoq/publisher/publisher_api.h

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ struct PublisherConfig {
2828
bool include_sap = false;
2929
bool include_msf_timeline = false;
3030
bool split_cmaf_chunks = true;
31+
bool live_stream_per_object = false;
3132
bool paced = false;
3233
bool loop = false;
3334
std::chrono::seconds subscriber_timeout = std::chrono::seconds(30);

include/openmoq/publisher/transport/moqt_session.h

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -67,11 +67,13 @@ class MoqtSession {
6767
TransportStatus publish(const openmoq::publisher::PublishPlan& plan);
6868
TransportStatus publish_live(std::istream& input,
6969
openmoq::publisher::DraftVersion draft_version,
70-
bool split_cmaf_chunks);
70+
bool split_cmaf_chunks,
71+
bool stream_per_object = false);
7172
TransportStatus publish_live(const LiveIngestOptions& ingest,
7273
std::istream* stdin_input,
7374
openmoq::publisher::DraftVersion draft_version,
74-
bool split_cmaf_chunks);
75+
bool split_cmaf_chunks,
76+
bool stream_per_object = false);
7577
TransportStatus publish_live_objects(const openmoq::publisher::LiveObjectSource& source,
7678
openmoq::publisher::DraftVersion draft_version);
7779
TransportStatus close(std::uint64_t application_error_code = 0);

src/cli_options.cpp

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,8 @@ CliOptions parse_cli_options(int argc, char** argv) {
178178
options.include_msf_timeline = true;
179179
} else if (argument == "--coalesce-cmaf-chunks") {
180180
options.split_cmaf_chunks = false;
181+
} else if (argument == "--stream-per-object") {
182+
options.stream_per_object = true;
181183
} else if (argument == "--timeout") {
182184
options.subscriber_timeout = parse_timeout(require_value("--timeout"));
183185
} else if (argument == "--paced") {
@@ -241,7 +243,7 @@ std::string build_usage(const char* argv0) {
241243
return std::string("Usage: ") + argv0 +
242244
" --input <mp4|-> [--live-source auto|stdin|srt] [--srt-config <path>]"
243245
" [--transport raw|webtransport] [--draft 14|16|17|18] [--namespace <value>] [--forward 0|1] [--timeout <seconds>]"
244-
" [--publish-catalog] [--sap] [--msf-timeline] [--coalesce-cmaf-chunks] [--paced] [--loop] [--dump-plan] [--emit-dir <dir>]"
246+
" [--publish-catalog] [--sap] [--msf-timeline] [--coalesce-cmaf-chunks] [--stream-per-object] [--paced] [--loop] [--dump-plan] [--emit-dir <dir>]"
245247
" [--endpoint host:port|moqt://host:port/path|https://host:port/path] [--alpn value] [--sni value]"
246248
" [--cert file] [--key file] [--ca file] [--insecure]";
247249
}

src/main.cpp

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ int main(int argc, char** argv) {
3030
.include_sap = options.include_sap,
3131
.include_msf_timeline = options.include_msf_timeline,
3232
.split_cmaf_chunks = options.split_cmaf_chunks,
33+
.live_stream_per_object = options.stream_per_object,
3334
.paced = options.paced,
3435
.loop = options.loop,
3536
.subscriber_timeout = options.subscriber_timeout,

src/publisher_api.cpp

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -256,7 +256,8 @@ transport::TransportStatus Publisher::publish_live(const LiveIngestConfig& inges
256256
status = active->session->publish_live(session_ingest,
257257
stdin_input,
258258
config_.draft_version,
259-
config_.split_cmaf_chunks);
259+
config_.split_cmaf_chunks,
260+
config_.live_stream_per_object);
260261
if (!status.ok) {
261262
const std::string error = "transport live publish failed: " + status.message;
262263
static_cast<void>(active->session->close(0));

src/transport/moqt_session.cpp

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2930,7 +2930,8 @@ TransportStatus MoqtSession::publish(const openmoq::publisher::PublishPlan& plan
29302930
TransportStatus MoqtSession::publish_live(const LiveIngestOptions& ingest,
29312931
std::istream* stdin_input,
29322932
openmoq::publisher::DraftVersion draft_version,
2933-
bool split_cmaf_chunks) {
2933+
bool split_cmaf_chunks,
2934+
bool stream_per_object) {
29342935
if (ingest.use_stdin && !ingest.srt_callers.empty()) {
29352936
return TransportStatus::failure("mixed stdin+SRT ingest is not supported; use either stdin or srt");
29362937
}
@@ -2941,7 +2942,7 @@ TransportStatus MoqtSession::publish_live(const LiveIngestOptions& ingest,
29412942
return TransportStatus::failure("live ingest requires at least one active source");
29422943
}
29432944
if (ingest.use_stdin) {
2944-
return publish_live(*stdin_input, draft_version, split_cmaf_chunks);
2945+
return publish_live(*stdin_input, draft_version, split_cmaf_chunks, stream_per_object);
29452946
}
29462947

29472948
// SRT-only path below.
@@ -3281,7 +3282,8 @@ TransportStatus MoqtSession::publish_live(const LiveIngestOptions& ingest,
32813282

32823283
TransportStatus MoqtSession::publish_live(std::istream& input,
32833284
openmoq::publisher::DraftVersion draft_version,
3284-
bool /*split_cmaf_chunks*/) {
3285+
bool /*split_cmaf_chunks*/,
3286+
bool stream_per_object) {
32853287
if (transport_.state() != ConnectionState::kConnected) {
32863288
return TransportStatus::failure("transport is not connected");
32873289
}
@@ -3590,7 +3592,7 @@ TransportStatus MoqtSession::publish_live(std::istream& input,
35903592
const std::span<const std::uint8_t> payload(fragment.payload.owned_bytes);
35913593
TransportStatus write_status = sender.serve(
35923594
transport_, draft_version, alias_it->second, send_seq,
3593-
object, true, true, payload);
3595+
object, true, stream_per_object, payload);
35943596
if (!write_status.ok) {
35953597
return write_status;
35963598
}

0 commit comments

Comments
 (0)