Current section

Files

Jump to
viva_tensor src viva_tensor@optim@pool.erl
Raw

src/viva_tensor@optim@pool.erl

-module(viva_tensor@optim@pool).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/viva_tensor/optim/pool.gleam").
-export([parallel_map/2, similarity_search/3, top_k_similar/3, parallel_sum/1, benchmark_pool/0, main/0]).
-export_type([tensor_op/0, search_result/0, pid_/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(false).
-type tensor_op() :: {scale, float()} | normalize.
-type search_result() :: {search_result, integer(), float()}.
-type pid_() :: any().
-file("src/viva_tensor/optim/pool.gleam", 265).
?DOC(false).
-spec apply_op(viva_tensor@tensor:tensor(), tensor_op()) -> viva_tensor@tensor:tensor().
apply_op(T, Op) ->
case Op of
{scale, Factor} ->
viva_tensor@tensor:scale(T, Factor);
normalize ->
viva_tensor@tensor:normalize(T)
end.
-file("src/viva_tensor/optim/pool.gleam", 55).
?DOC(false).
-spec parallel_map(list(viva_tensor@tensor:tensor()), tensor_op()) -> list(viva_tensor@tensor:tensor()).
parallel_map(Tensors, Op) ->
Parent = erlang:self(),
Indexed = gleam@list:index_map(Tensors, fun(T, Idx) -> {Idx, T} end),
gleam@list:each(
Indexed,
fun(Pair) ->
{Idx@1, T@1} = Pair,
erlang:spawn(
fun() ->
Result = apply_op(T@1, Op),
viva_tensor_ffi:send_msg(Parent, {Idx@1, Result})
end
)
end
),
_pipe = viva_tensor_ffi:collect_n(erlang:length(Tensors)),
_pipe@1 = gleam@list:sort(
_pipe,
fun(A, B) ->
{Idx_a, _} = A,
{Idx_b, _} = B,
gleam@int:compare(Idx_a, Idx_b)
end
),
gleam@list:map(
_pipe@1,
fun(Pair@1) ->
{_, T@2} = Pair@1,
T@2
end
).
-file("src/viva_tensor/optim/pool.gleam", 85).
?DOC(false).
-spec similarity_search(
viva_tensor@tensor:tensor(),
list(viva_tensor@tensor:tensor()),
integer()
) -> list(search_result()).
similarity_search(Query, Documents, Chunk_size) ->
Parent = erlang:self(),
Chunks = gleam@list:sized_chunk(Documents, Chunk_size),
Num_chunks = erlang:length(Chunks),
gleam@list:index_map(
Chunks,
fun(Chunk, Chunk_idx) ->
erlang:spawn(
fun() ->
Start_idx = Chunk_idx * Chunk_size,
Results = gleam@list:index_map(
Chunk,
fun(Doc, Local_idx) ->
Similarity = case viva_tensor@tensor:dot(Query, Doc) of
{ok, Sim} ->
Sim;
{error, _} ->
+0.0
end,
{search_result, Start_idx + Local_idx, Similarity}
end
),
viva_tensor_ffi:send_msg(Parent, {Chunk_idx, Results})
end
)
end
),
_pipe = viva_tensor_ffi:collect_n(Num_chunks),
_pipe@1 = gleam@list:flat_map(
_pipe,
fun(Pair) ->
{_, Results@1} = Pair,
Results@1
end
),
gleam@list:sort(
_pipe@1,
fun(A, B) ->
gleam@float:compare(erlang:element(3, B), erlang:element(3, A))
end
).
-file("src/viva_tensor/optim/pool.gleam", 123).
?DOC(false).
-spec top_k_similar(
viva_tensor@tensor:tensor(),
list(viva_tensor@tensor:tensor()),
integer()
) -> list(search_result()).
top_k_similar(Query, Documents, K) ->
_pipe = similarity_search(Query, Documents, 100),
gleam@list:take(_pipe, K).
-file("src/viva_tensor/optim/pool.gleam", 133).
?DOC(false).
-spec parallel_sum(list(viva_tensor@tensor:tensor())) -> float().
parallel_sum(Tensors) ->
Parent = erlang:self(),
Chunks = gleam@list:sized_chunk(Tensors, 100),
gleam@list:each(
Chunks,
fun(Chunk) ->
erlang:spawn(
fun() ->
Partial = gleam@list:fold(
Chunk,
+0.0,
fun(Acc, T) -> Acc + viva_tensor@tensor:sum(T) end
),
viva_tensor_ffi:send_msg(Parent, Partial)
end
)
end
),
_pipe = viva_tensor_ffi:collect_n(erlang:length(Chunks)),
gleam@list:fold(_pipe, +0.0, fun(Acc@1, X) -> Acc@1 + X end).
-file("src/viva_tensor/optim/pool.gleam", 279).
?DOC(false).
-spec float_to_string(float()) -> binary().
float_to_string(F) ->
Rounded = erlang:float(erlang:round(F * 100.0)) / 100.0,
gleam_stdlib:float_to_string(Rounded).
-file("src/viva_tensor/optim/pool.gleam", 313).
?DOC(false).
-spec range_loop(integer(), integer(), list(integer())) -> list(integer()).
range_loop(From, To, Acc) ->
case From > To of
true ->
lists:reverse(Acc);
false ->
range_loop(From + 1, To, [From | Acc])
end.
-file("src/viva_tensor/optim/pool.gleam", 309).
?DOC(false).
-spec range_int(integer(), integer()) -> list(integer()).
range_int(From, To) ->
range_loop(From, To, []).
-file("src/viva_tensor/optim/pool.gleam", 272).
?DOC(false).
-spec safe_div(integer(), integer()) -> float().
safe_div(A, B) ->
case B > 0 of
true ->
case erlang:float(B) of
+0.0 -> +0.0;
-0.0 -> -0.0;
Gleam@denominator -> erlang:float(A) / Gleam@denominator
end;
false ->
1.0
end.
-file("src/viva_tensor/optim/pool.gleam", 156).
?DOC(false).
-spec benchmark_pool() -> nil.
benchmark_pool() ->
gleam_stdlib:println(
<<"╔══════════════════════════════════════════════════════════════════╗"/utf8>>
),
gleam_stdlib:println(
<<"║ TensorPool - Distributed Tensor Computing in Pure Gleam! ║"/utf8>>
),
gleam_stdlib:println(
<<"║ Something C/C++ CANNOT do with this simplicity! ║"/utf8>>
),
gleam_stdlib:println(
<<"╚══════════════════════════════════════════════════════════════════╝\n"/utf8>>
),
gleam_stdlib:println(<<"CONCEPT: GenServer + GPU"/utf8>>),
gleam_stdlib:println(
<<" - GenServer manages state and load balancing"/utf8>>
),
gleam_stdlib:println(
<<" - Workers are lightweight BEAM processes (~2KB each)"/utf8>>
),
gleam_stdlib:println(<<" - NIFs call cuBLAS/cuDNN for GPU compute"/utf8>>),
gleam_stdlib:println(
<<" - Zero data races, fault tolerant by design!\n"/utf8>>
),
gleam_stdlib:println(
<<"━━━ TEST 1: Parallel Map (Scale 1000 tensors x 512d) ━━━"/utf8>>
),
Tensors = begin
_pipe = range_int(1, 1000),
gleam@list:map(
_pipe,
fun(_) -> viva_tensor@tensor:random_uniform([512]) end
)
end,
{Seq_time, _} = timer:tc(
fun() ->
gleam@list:map(
Tensors,
fun(T) -> viva_tensor@tensor:scale(T, 2.0) end
)
end
),
{Par_time, _} = timer:tc(fun() -> parallel_map(Tensors, {scale, 2.0}) end),
Speedup1 = safe_div(Seq_time, Par_time),
gleam_stdlib:println(
<<<<" Sequential: "/utf8,
(erlang:integer_to_binary(Seq_time div 1000))/binary>>/binary,
"ms"/utf8>>
),
gleam_stdlib:println(
<<<<" Parallel: "/utf8,
(erlang:integer_to_binary(Par_time div 1000))/binary>>/binary,
"ms"/utf8>>
),
gleam_stdlib:println(
<<<<" Speedup: "/utf8, (float_to_string(Speedup1))/binary>>/binary,
"x"/utf8>>
),
gleam_stdlib:println(
<<"\n━━━ TEST 2: Similarity Search (10K docs x 512d) ━━━"/utf8>>
),
Query = viva_tensor@tensor:random_uniform([512]),
Docs = begin
_pipe@1 = range_int(1, 10000),
gleam@list:map(
_pipe@1,
fun(_) -> viva_tensor@tensor:random_uniform([512]) end
)
end,
{Search_time, Results} = timer:tc(
fun() -> top_k_similar(Query, Docs, 10) end
),
Throughput = case (erlang:float(Search_time) / 1000000.0) of
+0.0 -> +0.0;
-0.0 -> -0.0;
Gleam@denominator -> 10000.0 / Gleam@denominator
end,
gleam_stdlib:println(
<<<<" 10K docs searched in: "/utf8,
(erlang:integer_to_binary(Search_time div 1000))/binary>>/binary,
"ms"/utf8>>
),
gleam_stdlib:println(
<<<<" Throughput: "/utf8, (float_to_string(Throughput))/binary>>/binary,
" docs/sec"/utf8>>
),
gleam_stdlib:println(<<" Top 3 matches:"/utf8>>),
_pipe@2 = Results,
_pipe@3 = gleam@list:take(_pipe@2, 3),
gleam@list:each(
_pipe@3,
fun(R) ->
gleam_stdlib:println(
<<<<<<" Doc "/utf8,
(erlang:integer_to_binary(erlang:element(2, R)))/binary>>/binary,
": sim="/utf8>>/binary,
(float_to_string(erlang:element(3, R)))/binary>>
)
end
),
gleam_stdlib:println(
<<"\n━━━ TEST 3: Parallel Reduce (Sum 1000 tensors) ━━━"/utf8>>
),
{Reduce_time, Total} = timer:tc(fun() -> parallel_sum(Tensors) end),
gleam_stdlib:println(
<<<<" Time: "/utf8,
(erlang:integer_to_binary(Reduce_time div 1000))/binary>>/binary,
"ms"/utf8>>
),
gleam_stdlib:println(<<" Total: "/utf8, (float_to_string(Total))/binary>>),
gleam_stdlib:println(
<<"\n╔══════════════════════════════════════════════════════════════════╗"/utf8>>
),
gleam_stdlib:println(
<<"║ HOW THIS RUNS ON THE GPU: ║"/utf8>>
),
gleam_stdlib:println(
<<"║ ║"/utf8>>
),
gleam_stdlib:println(
<<"║ 1. GenServer receives tensor op request ║"/utf8>>
),
gleam_stdlib:println(
<<"║ 2. Spawns Worker (BEAM process ~2KB) ║"/utf8>>
),
gleam_stdlib:println(
<<"║ 3. Worker calls NIF (Rust/C) that accesses GPU ║"/utf8>>
),
gleam_stdlib:println(
<<"║ 4. NIF uses cuBLAS for matmul, cuDNN for conv ║"/utf8>>
),
gleam_stdlib:println(
<<"║ 5. Result returns to Worker, then to GenServer ║"/utf8>>
),
gleam_stdlib:println(
<<"║ ║"/utf8>>
),
gleam_stdlib:println(
<<"║ Advantage: thousands of ops in parallel, fault tolerant! ║"/utf8>>
),
gleam_stdlib:println(
<<"╚══════════════════════════════════════════════════════════════════╝"/utf8>>
).
-file("src/viva_tensor/optim/pool.gleam", 152).
?DOC(false).
-spec main() -> nil.
main() ->
benchmark_pool().