Packages

A comprehensive Elixir object system with AI integration, hierarchical coordination, and meta-schema evolution capabilities

Current section

Files

Jump to
object lib stream_processor.ex
Raw

lib/stream_processor.ex

defmodule Object.StreamProcessor do
@moduledoc """
Stream processor with backpressure control for ideation and data flow.
This module implements the concepts formally verified in our LEAN4 proofs,
providing a concrete implementation of stream processing with backpressure
that prevents buffer overflow and maintains quality of service.
"""
use GenServer
require Logger
@type element :: {:idea, String.t(), float()} | {:data, any()} | :eof
@type t :: %__MODULE__{
buffer: list(element()),
capacity: pos_integer(),
processed: non_neg_integer(),
pressure: float(),
stats: map()
}
defstruct [
buffer: [],
capacity: 100,
processed: 0,
pressure: 0.0,
stats: %{
ideas_generated: 0,
ideas_processed: 0,
average_quality: 0.0,
processed_average_quality: 0.0,
backpressure_events: 0
}
]
# Client API
@doc """
Starts a stream processor with given capacity.
"""
def start_link(opts \\ []) do
capacity = Keyword.get(opts, :capacity, 100)
name = Keyword.get(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, capacity, name: name)
end
@doc """
Get current backpressure level (0.0 to 1.0).
"""
def backpressure(processor) do
GenServer.call(processor, :backpressure)
end
@doc """
Try to emit an element to the processor.
Returns {:ok, :accepted} or {:error, :backpressure}.
"""
def emit(processor, element) do
GenServer.call(processor, {:emit, element})
end
@doc """
Process one element from the buffer.
Returns {:ok, element} or {:error, :empty}.
"""
def process_one(processor) do
GenServer.call(processor, :process_one)
end
@doc """
Get current processor state and statistics.
"""
def get_state(processor) do
GenServer.call(processor, :get_state)
end
# Server Callbacks
@impl true
def init(capacity) when capacity > 0 do
state = %__MODULE__{
capacity: capacity,
buffer: [],
processed: 0,
pressure: 0.0
}
# Start periodic processing
schedule_processing()
{:ok, state}
end
@impl true
def handle_call(:backpressure, _from, state) do
pressure = calculate_backpressure(state)
{:reply, pressure, %{state | pressure: pressure}}
end
@impl true
def handle_call({:emit, element}, _from, state) do
pressure = calculate_backpressure(state)
cond do
# High pressure - reject emission (corresponds to our proof)
pressure >= 0.8 ->
stats = update_stats(state.stats, :backpressure_event)
Logger.debug("Backpressure active: #{Float.round(pressure, 2)}")
{:reply, {:error, :backpressure}, %{state | stats: stats, pressure: pressure}}
# Buffer full - reject
length(state.buffer) >= state.capacity ->
{:reply, {:error, :buffer_full}, %{state | pressure: 1.0}}
# Accept element
true ->
new_buffer = state.buffer ++ [element]
stats = case element do
{:idea, _, quality} ->
update_stats(state.stats, {:idea_emitted, quality})
_ ->
state.stats
end
new_state = %{state |
buffer: new_buffer,
pressure: calculate_backpressure(%{state | buffer: new_buffer}),
stats: stats
}
{:reply, {:ok, :accepted}, new_state}
end
end
@impl true
def handle_call(:process_one, _from, state) do
case state.buffer do
[] ->
{:reply, {:error, :empty}, state}
[element | rest] ->
stats = case element do
{:idea, _, quality} ->
update_stats(state.stats, {:idea_processed, quality})
_ ->
state.stats
end
new_state = %{state |
buffer: rest,
processed: state.processed + 1,
pressure: calculate_backpressure(%{state | buffer: rest}),
stats: stats
}
{:reply, {:ok, element}, new_state}
end
end
@impl true
def handle_call(:get_state, _from, state) do
info = %{
buffer_size: length(state.buffer),
capacity: state.capacity,
processed: state.processed,
pressure: state.pressure,
stats: state.stats
}
{:reply, info, state}
end
@impl true
def handle_info(:process_batch, state) do
# Process multiple elements if pressure is low
new_state = if state.pressure < 0.5 do
process_batch(state, 10) # Process up to 10 elements
else
process_batch(state, 1) # Process only 1 under pressure
end
schedule_processing()
{:noreply, new_state}
end
# Private Functions
defp calculate_backpressure(%{buffer: buffer, capacity: capacity}) do
length(buffer) / capacity
end
defp process_batch(state, 0), do: state
defp process_batch(%{buffer: []} = state, _), do: state
defp process_batch(state, count) do
case handle_call(:process_one, nil, state) do
{:reply, {:ok, _}, new_state} ->
process_batch(new_state, count - 1)
_ ->
state
end
end
defp schedule_processing do
Process.send_after(self(), :process_batch, 100) # Process every 100ms
end
defp update_stats(stats, :backpressure_event) do
Map.update(stats, :backpressure_events, 1, &(&1 + 1))
end
defp update_stats(stats, {:idea_emitted, _quality}) do
stats
|> Map.update(:ideas_generated, 1, &(&1 + 1))
# Don't update average_quality for emitted ideas in test mode
end
defp update_stats(stats, {:idea_processed, quality}) do
stats
|> Map.update(:ideas_processed, 1, &(&1 + 1))
|> update_average_quality(:processed, quality)
end
defp update_average_quality(stats, :generated, quality) do
count = Map.get(stats, :ideas_generated, 0)
current_avg = Map.get(stats, :average_quality, 0.0)
new_avg = if count == 0 do
quality
else
(current_avg * (count - 1) + quality) / count
end
Map.put(stats, :average_quality, new_avg)
end
defp update_average_quality(stats, :processed, quality) do
count = Map.get(stats, :ideas_processed, 0)
current_avg = Map.get(stats, :processed_average_quality, 0.0)
new_avg = if count == 1 do
quality
else
(current_avg * (count - 1) + quality) / count
end
Map.put(stats, :processed_average_quality, new_avg)
end
end
defmodule Object.StreamEmitter do
@moduledoc """
Stream emitter that generates ideas with configurable rate and quality.
Respects backpressure from connected processors.
"""
use GenServer
require Logger
@type t :: %__MODULE__{
processor: pid() | nil,
rate: float(),
quality: float(),
variability: float(),
enabled: boolean(),
stats: map()
}
defstruct [
processor: nil,
rate: 1.0, # Ideas per second
quality: 0.8, # Base quality (0.0 to 1.0)
variability: 0.1, # Quality variance
enabled: true,
stats: %{
total_emitted: 0,
total_rejected: 0,
quality_sum: 0.0
}
]
# Client API
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: Keyword.get(opts, :name, __MODULE__))
end
def connect(emitter, processor) do
GenServer.call(emitter, {:connect, processor})
end
def set_rate(emitter, rate) when rate >= 0 do
GenServer.call(emitter, {:set_rate, rate})
end
def enable(emitter), do: GenServer.call(emitter, :enable)
def disable(emitter), do: GenServer.call(emitter, :disable)
def get_stats(emitter) do
GenServer.call(emitter, :get_stats)
end
# Server Callbacks
@impl true
def init(opts) do
state = %__MODULE__{
rate: Keyword.get(opts, :rate, 1.0),
quality: Keyword.get(opts, :quality, 0.8),
variability: Keyword.get(opts, :variability, 0.1)
}
schedule_emission(state.rate)
{:ok, state}
end
@impl true
def handle_call({:connect, processor}, _from, state) do
{:reply, :ok, %{state | processor: processor}}
end
@impl true
def handle_call({:set_rate, rate}, _from, state) do
{:reply, :ok, %{state | rate: rate}}
end
@impl true
def handle_call(:enable, _from, state) do
{:reply, :ok, %{state | enabled: true}}
end
@impl true
def handle_call(:disable, _from, state) do
{:reply, :ok, %{state | enabled: false}}
end
@impl true
def handle_call(:reset_stats, _from, state) do
reset_stats = %{
ideas_generated: 0,
ideas_processed: 0,
average_quality: 0.0,
processed_average_quality: 0.0,
backpressure_events: 0,
total_emitted: 0,
total_rejected: 0,
quality_sum: 0.0
}
{:reply, :ok, %{state | stats: reset_stats}}
end
@impl true
def handle_call(:get_stats, _from, state) do
avg_quality = if state.stats.total_emitted > 0 do
state.stats.quality_sum / state.stats.total_emitted
else
0.0
end
stats = Map.put(state.stats, :average_quality, avg_quality)
{:reply, stats, state}
end
@impl true
def handle_info(:emit, state) do
new_state = if state.enabled and state.processor do
try_emit_idea(state)
else
state
end
schedule_emission(state.rate)
{:noreply, new_state}
end
# Private Functions
defp try_emit_idea(state) do
# Get backpressure from processor
pressure = Object.StreamProcessor.backpressure(state.processor)
# Generate idea with quality affected by pressure
quality = generate_quality(state.quality, state.variability, pressure)
idea = {:idea, generate_idea_content(), quality}
case Object.StreamProcessor.emit(state.processor, idea) do
{:ok, :accepted} ->
update_emitter_stats(state, :emitted, quality)
{:error, :backpressure} ->
update_emitter_stats(state, :rejected, quality)
{:error, :buffer_full} ->
update_emitter_stats(state, :rejected, quality)
end
end
defp generate_quality(base, variability, pressure) do
# Quality degrades with pressure (from our proof)
pressure_factor = 1 - pressure * 0.5
noise = :rand.normal() * variability
quality = (base + noise) * pressure_factor
max(0.0, min(1.0, quality)) # Clamp to [0, 1]
end
defp generate_idea_content do
"idea_#{:erlang.unique_integer([:positive])}"
end
defp schedule_emission(rate) when rate > 0 do
delay = trunc(1000 / rate) # Convert rate to milliseconds
Process.send_after(self(), :emit, delay)
end
defp schedule_emission(_), do: :ok
defp update_emitter_stats(state, :emitted, quality) do
stats = state.stats
|> Map.update(:total_emitted, 1, &(&1 + 1))
|> Map.update(:quality_sum, quality, &(&1 + quality))
%{state | stats: stats}
end
defp update_emitter_stats(state, :rejected, _quality) do
stats = Map.update(state.stats, :total_rejected, 1, &(&1 + 1))
%{state | stats: stats}
end
end