Skip to content

Commit 5875168

Browse files
committed
Fix & improve tests
1 parent 94082ac commit 5875168

4 files changed

Lines changed: 168 additions & 15 deletions

File tree

test/membrane/buffer_metric/byte_size_test.exs

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -13,32 +13,31 @@ defmodule Membrane.Buffer.Metric.ByteSizeTest do
1313

1414
describe ".buffers_size/1" do
1515
test "should return size of all buffers" do
16-
size = ByteSize.buffers_size(@buffers)
17-
assert size == byte_size(@pay1) + byte_size(@pay2)
16+
assert ByteSize.buffers_size(@buffers) == {:ok, byte_size(@pay1) + byte_size(@pay2)}
1817
end
1918
end
2019

21-
describe ".split_buffers/2" do
20+
describe ".split_buffers/4" do
2221
test "when split position matches size of first buffer, extract only first buffer" do
23-
{buf, rest} = ByteSize.split_buffers(@buffers, byte_size(@pay1), nil)
22+
{buf, rest} = ByteSize.split_buffers(@buffers, byte_size(@pay1), nil, nil)
2423
assert buf == [@buf1]
2524
assert rest == [@buf2]
2625
end
2726

2827
test "when there is only one buffer where split position is greater than buffer size \
2928
returns the buffer and an empty list" do
30-
{buf, []} = ByteSize.split_buffers(@single_buffer, byte_size(@pay1) + 10, nil)
29+
{buf, []} = ByteSize.split_buffers(@single_buffer, byte_size(@pay1) + 10, nil, nil)
3130
assert buf == [@buf1]
3231
end
3332

3433
test "when there is only one buffer where split position is 0, it returns an empty \
3534
list and a list with the buffer" do
36-
{[], rest} = ByteSize.split_buffers(@single_buffer, 0, nil)
35+
{[], rest} = ByteSize.split_buffers(@single_buffer, 0, nil, nil)
3736
assert rest == [@buf1]
3837
end
3938

4039
test "when splitting is necessary it extracts the first buffer and splits the second into two" do
41-
{extracted, rest} = ByteSize.split_buffers(@buffers, byte_size(@pay1) + 1, nil)
40+
{extracted, rest} = ByteSize.split_buffers(@buffers, byte_size(@pay1) + 1, nil, nil)
4241
<<one_byte::binary-size(1), expected_rest::binary>> = @pay2
4342
assert extracted == [@buf1, %Membrane.Buffer{payload: one_byte}]
4443
assert rest == [%Membrane.Buffer{payload: expected_rest}]

test/membrane/buffer_metric/count_test.exs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,13 +11,13 @@ defmodule Membrane.Buffer.Metric.CountTest do
1111

1212
describe ".buffers_size/1" do
1313
test "should return count of all buffers" do
14-
assert Count.buffers_size(@buffers) == 2
14+
assert Count.buffers_size(@buffers) == {:ok, 2}
1515
end
1616
end
1717

18-
describe ".split_buffers/2" do
18+
describe ".split_buffers/4" do
1919
test "should return split buffers" do
20-
{extracted, rest} = Count.split_buffers(@buffers, @count)
20+
{extracted, rest} = Count.split_buffers(@buffers, @count, nil, nil)
2121

2222
assert extracted == [@buf1]
2323
assert rest == [@buf2]
Lines changed: 143 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,143 @@
1+
defmodule Membrane.Buffer.Metric.TimestampTest do
2+
use ExUnit.Case, async: true
3+
4+
import ExUnit.CaptureLog
5+
6+
alias Membrane.Buffer
7+
alias Membrane.Buffer.Metric.Timestamp.{DTS, DTSOrPTS, PTS}
8+
9+
@t0 0 |> Membrane.Time.milliseconds()
10+
@t1 100 |> Membrane.Time.milliseconds()
11+
@t2 250 |> Membrane.Time.milliseconds()
12+
@t3 350 |> Membrane.Time.milliseconds()
13+
@t4 500 |> Membrane.Time.milliseconds()
14+
15+
defp buf(ts_field, ts), do: struct(%Buffer{payload: <<>>}, [{ts_field, ts}])
16+
17+
test ".init_manual_demand_size_value/0 returns -1 as the no-demand sentinel" do
18+
assert PTS.init_manual_demand_size_value() == -1
19+
assert DTS.init_manual_demand_size_value() == -1
20+
assert DTSOrPTS.init_manual_demand_size_value() == -1
21+
end
22+
23+
test ".reduce_demand/2 always returns the demand unchanged regardless of consumed size" do
24+
for module <- [PTS, DTS, DTSOrPTS] do
25+
assert module.reduce_demand(1_000, 5) == 1_000
26+
assert module.reduce_demand(1_000, nil) == 1_000
27+
end
28+
end
29+
30+
test ".buffers_size/1 returns {:error, :operation_not_supported}" do
31+
for module <- [PTS, DTS, DTSOrPTS] do
32+
assert module.buffers_size([]) == {:error, :operation_not_supported}
33+
assert module.buffers_size([%Buffer{payload: <<>>}]) == {:error, :operation_not_supported}
34+
end
35+
end
36+
37+
describe ".split_buffers/4" do
38+
for {module, name, ts_field} <- [{PTS, "PTS", :pts}, {DTS, "DTS", :dts}] do
39+
test "returns {[], buffers} when demand is the initial sentinel value (-1) for #{name}" do
40+
buffers = Enum.map([@t0, @t1, @t2], &buf(unquote(ts_field), &1))
41+
assert unquote(module).split_buffers(buffers, -1, nil, nil) == {[], buffers}
42+
end
43+
44+
test "uses first buffer's #{name} as offset when no buffers have been consumed yet" do
45+
buffers = Enum.map([@t0, @t1, @t2, @t3, @t4], &buf(unquote(ts_field), &1))
46+
# offset = @t0 = 0; demand = 300 ns → consume until ts - 0 >= 300 → stops at @t3 = 350
47+
{consumed, remaining} = unquote(module).split_buffers(buffers, 300, nil, nil)
48+
assert Enum.map(consumed, &Map.get(&1, unquote(ts_field))) == [@t0, @t1, @t2, @t3]
49+
assert Enum.map(remaining, &Map.get(&1, unquote(ts_field))) == [@t4]
50+
end
51+
52+
test "uses first_consumed_buffer's #{name} as offset when buffers have been consumed" do
53+
buffers = Enum.map([@t0, @t1, @t2, @t3, @t4], &buf(unquote(ts_field), &1))
54+
# first_consumed at @t0 → same offset as above, same split point
55+
first_consumed = buf(unquote(ts_field), @t0)
56+
last_consumed = buf(unquote(ts_field), @t1)
57+
58+
{consumed, remaining} =
59+
unquote(module).split_buffers(buffers, 300, first_consumed, last_consumed)
60+
61+
assert Enum.map(consumed, &Map.get(&1, unquote(ts_field))) == [@t0, @t1, @t2, @t3]
62+
assert Enum.map(remaining, &Map.get(&1, unquote(ts_field))) == [@t4]
63+
end
64+
65+
test "returns all buffers for #{name} when demand exceeds the available timestamp range" do
66+
buffers = Enum.map([@t0, @t1, @t2, @t3, @t4], &buf(unquote(ts_field), &1))
67+
{consumed, remaining} = unquote(module).split_buffers(buffers, 10_000, nil, nil)
68+
assert consumed == buffers
69+
assert remaining == []
70+
end
71+
72+
test "emits a warning and returns {[], buffers} for #{name} when elapsed duration already meets demand" do
73+
buffers = Enum.map([@t0, @t1, @t2, @t3, @t4], &buf(unquote(ts_field), &1))
74+
# last_consumed - first_consumed = @t4 - @t0 = 500 >= demand = 300
75+
first_consumed = buf(unquote(ts_field), @t0)
76+
last_consumed = buf(unquote(ts_field), @t4)
77+
78+
log =
79+
capture_log(fn ->
80+
assert unquote(module).split_buffers(buffers, 300, first_consumed, last_consumed) ==
81+
{[], buffers}
82+
end)
83+
84+
assert log =~ "warning"
85+
end
86+
end
87+
88+
test "DTSOrPTS prefers DTS when present, falls back to PTS" do
89+
# first buffer has dts=@t0 and pts=@t4: metric should use dts
90+
buf_with_dts = %Buffer{payload: <<>>, dts: @t0, pts: @t4}
91+
buf_pts_only = %Buffer{payload: <<>>, pts: @t3}
92+
93+
# offset = dts=@t0=0; buf_with_dts: 0-0=0 < 300 → include
94+
# buf_pts_only uses pts=@t3=350: 350-0=350 >= 300 → include, stop
95+
{consumed, remaining} = DTSOrPTS.split_buffers([buf_with_dts, buf_pts_only], 300, nil, nil)
96+
assert consumed == [buf_with_dts, buf_pts_only]
97+
assert remaining == []
98+
end
99+
end
100+
101+
describe ".generate_metric_specific_warnings/1" do
102+
test "returns :ok for an empty list" do
103+
assert PTS.generate_metric_specific_warnings([]) == :ok
104+
assert DTS.generate_metric_specific_warnings([]) == :ok
105+
assert DTSOrPTS.generate_metric_specific_warnings([]) == :ok
106+
end
107+
108+
for {module, name, ts_field} <- [{PTS, "PTS", :pts}, {DTS, "DTS", :dts}] do
109+
test "emits no warning for monotonically increasing #{name}s" do
110+
buffers = Enum.map([@t0, @t1, @t2, @t3], &buf(unquote(ts_field), &1))
111+
112+
log =
113+
capture_log(fn ->
114+
assert unquote(module).generate_metric_specific_warnings(buffers) == :ok
115+
end)
116+
117+
assert log == ""
118+
end
119+
120+
test "emits a warning for non-monotonic #{name}s" do
121+
buffers = Enum.map([@t0, @t3, @t1, @t4], &buf(unquote(ts_field), &1))
122+
123+
log =
124+
capture_log(fn ->
125+
assert unquote(module).generate_metric_specific_warnings(buffers) == :ok
126+
end)
127+
128+
assert log =~ "warning"
129+
end
130+
end
131+
132+
test "DTSOrPTS emits a warning for non-monotonic DTS-or-PTS values" do
133+
buffers = Enum.map([@t0, @t3, @t1, @t4], &%Buffer{payload: <<>>, dts: &1})
134+
135+
log =
136+
capture_log(fn ->
137+
assert DTSOrPTS.generate_metric_specific_warnings(buffers) == :ok
138+
end)
139+
140+
assert log =~ "warning"
141+
end
142+
end
143+
end

test/membrane/core/element/input_queue_test.exs

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -373,7 +373,15 @@ defmodule Membrane.Core.Element.InputQueueTest do
373373
end
374374

375375
defp prepare_input_queue(fields) do
376-
struct(InputQueue, [stalker_metrics: %{size: :atomics.new(1, [])}] ++ fields)
376+
outbound_metric = Keyword.get(fields, :outbound_metric)
377+
378+
defaults = [
379+
stalker_metrics: %{size: :atomics.new(1, [])},
380+
outbound_metric_demand_init_size:
381+
if(outbound_metric, do: outbound_metric.init_manual_demand_size_value(), else: 0)
382+
]
383+
384+
struct(InputQueue, defaults ++ fields)
377385
end
378386

379387
defp new_atomic_demand(),
@@ -390,10 +398,13 @@ defmodule Membrane.Core.Element.InputQueueTest do
390398
defp bufs_size(output, unit) do
391399
{_state, bufs} = output
392400

393-
Enum.flat_map(bufs, fn {:buffers, bufs_list, _inbound_metric_size, _outbound_metric_size} ->
394-
bufs_list
395-
end)
396-
|> Membrane.Buffer.Metric.from_unit(unit).buffers_size()
401+
{:ok, size} =
402+
Enum.flat_map(bufs, fn {:buffers, bufs_list, _inbound_metric_size, _outbound_metric_size} ->
403+
bufs_list
404+
end)
405+
|> Membrane.Buffer.Metric.from_unit(unit).buffers_size()
406+
407+
size
397408
end
398409

399410
defp bufs(n), do: Enum.to_list(1..n)

0 commit comments

Comments
 (0)