Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
175 changes: 175 additions & 0 deletions benchmark/keyed_reduce_compiler.jl
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
#!/usr/bin/env julia

# Reproducible compiler/allocation evidence for the warm sparse-key boundary.
import KernelAbstractions
import LocalMath
import TOML

struct KeyedCompilerNode end
struct KeyedCompilerEvaluator end
@inline (::KeyedCompilerEvaluator)(item::Int32, reads, parameters) =
(; delta = LocalMath.KeyedContribution(
(UInt32(isodd(item)), UInt32(item)), Int32(1)))

struct CollectCompilerEvaluator end
@inline (::CollectCompilerEvaluator)(item::Int32, reads, parameters) =
(; delta = LocalMath.CollectedValue(item))

struct KeyedCompilerSubtract end
@inline (::KeyedCompilerSubtract)(left::Int32, right::Int32) = left - right

function keyed_compiler_preparation(capacity::Int;
operation = +, retention = LocalMath.DropIdentityKeys())
source = LocalMath.Space(KeyedCompilerNode, 3)
key_type = Tuple{UInt32,UInt32}
collection = LocalMath.Collection(
LocalMath.KeyedValue{key_type,Int32}, capacity)
stage = LocalMath.Stage(source, NamedTuple(), (
LocalMath.Publication(collection,
LocalMath.KeyedReduce(key_type, Int32, operation;
seed = LocalMath.NewKeyIdentity(Int32(0)), retention);
value = :delta),),
LocalMath.Evaluator(KeyedCompilerEvaluator()), LocalMath.Control(),
LocalMath.SourceOrigin(:keyed_reduce_compiler, 1))
return LocalMath.prepare(LocalMath.LocalLaw(stage),
collection => LocalMath.Allocate();
backend = KernelAbstractions.CPU())
end

function collect_compiler_preparation(capacity::Int)
source = LocalMath.Space(KeyedCompilerNode, 3)
collection = LocalMath.Collection(Int32, capacity)
stage = LocalMath.Stage(source, NamedTuple(), (
LocalMath.Publication(collection,
LocalMath.Collect(Int32; maximum = 1); value = :delta),),
LocalMath.Evaluator(CollectCompilerEvaluator()), LocalMath.Control(),
LocalMath.SourceOrigin(:collect_compiler_control, 1))
return LocalMath.prepare(LocalMath.LocalLaw(stage),
collection => LocalMath.Allocate();
backend = KernelAbstractions.CPU())
end

function typed_metrics(callable, signature)
info, return_type = only(Base.code_typed_by_type(
Tuple{typeof(callable),signature.parameters...}; optimize = true))
calls = count(statement -> statement isa Expr &&
statement.head in (:call, :invoke), info.code)
any_indices = filter(index -> info.ssavaluetypes[index] === Any,
eachindex(info.code))
control_flow = statement -> statement isa Union{
Core.GotoNode,Core.GotoIfNot,Core.ReturnNode}
return Dict(
"statement_count" => length(info.code),
"call_count" => calls,
"any_ssa_count" => length(any_indices),
"any_control_flow_count" => count(
index -> control_flow(info.code[index]), any_indices),
"any_value_count" => count(
index -> !control_flow(info.code[index]), any_indices),
"return_type" => string(return_type),
)
end

function warm_public_execution_allocations(prepared)
wait(LocalMath.execute!(prepared))
return minimum(@allocated(wait(LocalMath.execute!(prepared))) for _ in 1:5)
end

function keyed_compiler_metrics(prepared)
launch = only(getfield(getfield(prepared, :runtime), :launches))
stage = getfield(launch, :stage)
validation = LocalMath._ProgramValidationTarget(
getfield(getfield(prepared, :runtime), :execution_gate), Int32(1))
signature = Tuple{typeof(stage),Tuple{},Int32,Tuple{},
typeof(getfield(launch, :guard)),typeof(validation)}
host = typed_metrics(LocalMath._execute_keyed_reduce_stage!, signature)
execution = getfield(stage, :execution)
plan = getfield(execution, :plan)
states = getfield(execution, :states)
emission = LocalMath.KeyedContribution(
(UInt32(1), UInt32(1)), Int32(1))
semantic_signature = Tuple{typeof(plan.emission),typeof(states.emission),
typeof(emission),Int32}
semantic = typed_metrics(LocalMath._keyed_reduce_materialize!,
semantic_signature)
allocated = warm_public_execution_allocations(prepared)
return Dict(
"host_orchestration" => host,
"emission_boundary" => semantic,
"sort_boundary" => typed_metrics(LocalMath._compacted_ordinal_less,
Tuple{typeof(plan.key_order),typeof(states.ordering),Int32,Int32}),
"fold_boundary" => typed_metrics(LocalMath._keyed_reduce_fold_segment!,
Tuple{typeof(plan.fold),typeof(states.fold),
typeof(states.ordering.order_a),Int32,Int32}),
"publish_boundary" => typed_metrics(LocalMath._keyed_reduce_publish_record!,
Tuple{typeof(plan.publication),typeof(states.publication),
typeof(execution.storage),Int32,Int32}),
"warm_public_execution_allocated_bytes" => allocated,
"host_any_classes" => [
"KernelAbstractions launch construction (_svec_ref and kwcall)",
"compacted scan/order launch calls whose host return is unused",
"host control-flow and return nodes represented as Any by CodeInfo",
],
"allocation_root_classes" => [
"shared ExecutionReceipt and validation-status grouping",
"KernelAbstractions Kernel, NDRange, keyword, and argument tuples per launch",
],
"prepared_stage_type" => string(typeof(stage)),
)
end

function collect_compiler_metrics(prepared)
launch = only(getfield(getfield(prepared, :runtime), :launches))
stage = getfield(launch, :stage)
validation = LocalMath._ProgramValidationTarget(
getfield(getfield(prepared, :runtime), :execution_gate), Int32(1))
signature = Tuple{typeof(stage),Tuple{},Int32,Tuple{},
typeof(getfield(launch, :guard)),typeof(validation)}
return typed_metrics(LocalMath._execute_collect_stage!, signature)
end

preparations = map(capacity -> keyed_compiler_preparation(capacity), (4, 8))
variants = (
keyed_compiler_preparation(4; operation = +,
retention = LocalMath.DropIdentityKeys()),
keyed_compiler_preparation(4; operation = +,
retention = LocalMath.RetainAllKeys()),
keyed_compiler_preparation(4; operation = KeyedCompilerSubtract(),
retention = LocalMath.DropIdentityKeys()),
keyed_compiler_preparation(4; operation = KeyedCompilerSubtract(),
retention = LocalMath.RetainAllKeys()),
)
metrics = keyed_compiler_metrics(first(preparations))
control = collect_compiler_preparation(4)
metrics["collect_host_orchestration"] = collect_compiler_metrics(control)
control_allocated = warm_public_execution_allocations(control)
metrics["collect_control_allocated_bytes"] = control_allocated
metrics["keyed_incremental_allocated_bytes"] =
metrics["warm_public_execution_allocated_bytes"] - control_allocated
metrics["capacity_specialization_count"] = length(unique(typeof(
getfield(only(getfield(getfield(prepared, :runtime), :launches)), :stage))
for prepared in preparations))
variant_parts = map(variants) do prepared
execution = getfield(getfield(only(getfield(
getfield(prepared, :runtime), :launches)), :stage), :execution)
(; plan = execution.plan, states = execution.states,
storage = execution.storage)
end
metrics["operation_retention_specializations"] = Dict(
"bounds" => length(unique(typeof(part.plan.bounds) for part in variant_parts)),
"emission" => length(unique(typeof(part.plan.emission) for part in variant_parts)),
"sort" => length(unique((typeof(part.plan.key_order),
typeof(part.states.ordering)) for part in variant_parts)),
"segment" => length(unique((typeof(part.plan.key_order),
typeof(part.states.segment)) for part in variant_parts)),
"fold" => length(unique((typeof(part.plan.fold),
typeof(part.states.fold)) for part in variant_parts)),
"finalize" => length(unique((typeof(part.plan.bounds),
typeof(part.states.final)) for part in variant_parts)),
"publish" => length(unique((typeof(part.plan.publication),
typeof(part.states.publication), typeof(part.storage))
for part in variant_parts)),
)
metrics["kaimon"] = "executed through a Kaimon persistent LocalMath project session; metrics are produced by these reproducible Base.code_typed_by_type probes"
TOML.print(stdout, Dict("keyed_reduce" => metrics); sorted = true)
println()
52 changes: 49 additions & 3 deletions docs/src/api/localmath.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,52 @@ Duplicate canonical identities fail validation before changing the previously
published count or records. The ordinary CPU and Metal collection-order tests
exercise these behaviors across partial workgroups with bounds checks enabled.

### Incremental sparse keyed state

`KeyedReduce` updates one bounded keyed `Collection` without introducing a
second scheduler or storage authority. Prior records are intrinsic input state;
the required `seed` keyword supplies the initial value only for a key absent at
stage entry.

```julia
import KernelAbstractions
import LocalMath

struct Contact end
struct ContactDeltas end

@inline function (::ContactDeltas)(contact::Int32, reads, parameters)
owner = UInt32(isodd(contact) ? 1 : 2)
delta = isodd(contact) ? Int32(1) : Int32(-1)
return (; change = LocalMath.KeyedContribution(owner, delta))
end

contacts = LocalMath.Space(Contact, 4)
counts = LocalMath.Collection(LocalMath.KeyedValue{UInt32,Int32}, 8)
stage = LocalMath.Stage(contacts, NamedTuple(), (
LocalMath.Publication(counts,
LocalMath.KeyedReduce(UInt32, Int32, +;
maximum = 1,
seed = LocalMath.NewKeyIdentity(Int32(0)),
retention = LocalMath.DropIdentityKeys());
value = :change),),
LocalMath.Evaluator(ContactDeltas()), LocalMath.Control(),
LocalMath.SourceOrigin(:contact_counts, 1))
law = LocalMath.LocalLaw(stage)
prepared = LocalMath.prepare(law, counts => LocalMath.Allocate();
backend = KernelAbstractions.CPU())
wait(LocalMath.execute!(prepared))
records = LocalMath.storage(prepared, counts)
```

Every exact-key segment intrinsically folds an existing value first, then
participating tuple lanes in canonical `(source, lane)` order. Invalid prior counts, duplicate prior
keys, and final capacity overflow reject the whole publication, leaving its
records and logical count unchanged. `prepare` owns the bounded device
workspace; execution performs no device allocation. Public `execute!` and
`wait` still allocate shared host receipt/launch bookkeeping, which is tracked
separately rather than claimed as zero-allocation execution.

## Public surface

Ordinary authoring exports only the mathematical and execution vocabulary:
Expand All @@ -71,11 +117,11 @@ equation namespace:
|:--|:--|
| Lifecycle | `Plan`, `PreparedPlan`, `ExecutionReceipt`, `LocalMathValidationError`, `bind`, `plan`, `Allocate`, `Temporary`, `MutableRelationStorage`, `storage`, `inspect`, `compilation_report`, `execution_contract`, `lowering_identity` |
| Explicit laws | `Stage`, `Publication`, `Access`, `Control`, `SourceOrigin`, `Parameter`, `ParameterSchema`, `Evaluator`, `FieldPublication`, `CollectionPublication`, `FoldPublication`, `PublicationValue`, `sequence` |
| Collections | `CollectionAccess`, `CollectionCount`, `BoundedGroup`, `SourcePositionAccess`, `CompactedStorage`, `BoundedGroupView`, `one_group`, `group_by`, `source_order`, `canonical_by`, `persistent_source_position` |
| Publication laws | `Unique`, `Reduce`, `Resolve`, `Collect`, `OrderedFold`, `TotalCoverage`, `PartialCoverage`, `UnreachableEmpty`, `PreserveEmpty`, `FillEmpty`, `IdentitySeed`, `ExistingSeed`, `CanonicalLeftFold`, `RelaxedAtomic`, `ArgMin`, `ArgMax`, `CanonicalSourceLaneTie`, `TieMin`, `TieMax`, `RejectOverflow`, `EmptyCollection` |
| Collections | `CollectionAccess`, `CollectionCount`, `BoundedGroup`, `SourcePositionAccess`, `CompactedStorage`, `BoundedGroupView`, `KeyedValue`, `one_group`, `group_by`, `source_order`, `canonical_by`, `persistent_source_position` |
| Publication laws | `Unique`, `Reduce`, `Resolve`, `Collect`, `KeyedReduce`, `OrderedFold`, `TotalCoverage`, `PartialCoverage`, `UnreachableEmpty`, `PreserveEmpty`, `FillEmpty`, `IdentitySeed`, `ExistingSeed`, `NewKeyIdentity`, `RetainAllKeys`, `DropIdentityKeys`, `CanonicalLeftFold`, `RelaxedAtomic`, `ArgMin`, `ArgMax`, `CanonicalSourceLaneTie`, `TieMin`, `TieMax`, `RejectOverflow`, `EmptyCollection` |
| Ordered state | `FoldComponent`, `InitializedState`, `initialized_state`, `BoundedWrites`, `FoldStep` |
| Bounded scalar operations | `fold`, `BoundedFold`, `Where`, `RejectInvalid`, `SkipInvalid`, `FillInvalid`, `RejectEmpty`, `RelaxedAssociative`, `BoundedFoldOutcome`, `evaluate_bounded` |
| Evaluator outputs | `UniqueValue`, `ConditionalUniqueValue`, `RoutedUniqueValue`, `ConditionalRoutedUniqueValue`, `Contribution`, `RoutedContribution`, `ResolutionValue`, `RoutedResolutionValue`, `CollectedValue`, `GroupedCollectedValue`, `FoldValue` |
| Evaluator outputs | `UniqueValue`, `ConditionalUniqueValue`, `RoutedUniqueValue`, `ConditionalRoutedUniqueValue`, `Contribution`, `RoutedContribution`, `ResolutionValue`, `RoutedResolutionValue`, `CollectedValue`, `GroupedCollectedValue`, `KeyedContribution`, `FoldValue` |
| Advanced execution | `allocate_workspace`, `submission_capacity`, `ispending`, `success_gate` |

These qualified names are stable interfaces, not permission to access other
Expand Down
23 changes: 23 additions & 0 deletions spec/localmath.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ Publication laws define observable conflict behavior:
- reduction applies its declared operation and ordering law;
- resolution selects a score/payload pair with explicit tie and empty rules;
- collection publishes bounded records with explicit grouping and overflow;
- keyed reduction updates bounded sparse exact-key state by a canonical fold;
- ordered fold applies a bounded recurrence in its declared canonical order.

In authored `resolve_to` expressions, a noncanonical tie is an explicit
Expand Down Expand Up @@ -66,6 +67,28 @@ source-position lanes. Collection production distinguishes its global storage
capacity from the per-source `maximum` emission width. These forms lower to
the existing `CollectionAccess`, source-position, `Collect`, and control laws.

`KeyedReduce(K, V, operation; seed, ...)` is the narrow
incremental sparse-state companion to `Collect`. Its sole destination is a
`Collection{KeyedValue{K,V}}`. `K` is `Int32`, `UInt32`, or a bounded flat
tuple of those types. The stage-entry collection must contain unique keys.
Each exact-key segment folds its existing value first when present, then
participating `KeyedContribution`s in intrinsic `(source item, lane)` left-fold
order. It has no order selector, hash, registry, relaxed, or provider-specific
path. `DropIdentityKeys` frees capacity when the final
value equals the declared identity, while `RetainAllKeys` preserves it.
Invalid prior count, duplicate prior keys, invalid stage control, and final
capacity overflow reject the complete publication. Private workspace receives
all intermediate values; records and the device-resident count publish through
one validation gate, so failure leaves both unchanged. Workspace is
`O(capacity + source_count * maximum)` and exact-key sorting followed by
segmented folding is `O(n log n)` work plus linear scans.
Its bounded device workspace is allocated by `prepare`; execution performs no
device allocation. The public host `execute!`/`wait` path still allocates
receipt, launch, and event bookkeeping shared with other Stage executors, so
zero host allocation is not part of this contract. The reproducible compiler
benchmark reports both that shared baseline and the narrower keyed semantic
boundary rather than treating host orchestration `Any` values as device IR.

Ordered recurrence is authored by declaring a total event order and every
evolving state component with its exact initial Field:

Expand Down
6 changes: 4 additions & 2 deletions src/LocalMath.jl
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ public sequence, allocate_workspace, submission_capacity, ispending, success_gat
public one_group, group_by, source_order, canonical_by
public persistent_source_position
public CompactedStorage, BoundedGroupView
public Unique, Reduce, Resolve, Collect, OrderedFold
public Unique, Reduce, Resolve, Collect, KeyedReduce, OrderedFold
public TotalCoverage, PartialCoverage, UnreachableEmpty, PreserveEmpty, FillEmpty
public IdentitySeed, ExistingSeed, CanonicalLeftFold, RelaxedAtomic
public ArgMin, ArgMax, CanonicalSourceLaneTie, TieMin, TieMax
Expand All @@ -46,7 +46,8 @@ public BoundedFoldOutcome, evaluate_bounded
public UniqueValue, ConditionalUniqueValue, RoutedUniqueValue
public ConditionalRoutedUniqueValue, Contribution, RoutedContribution
public ResolutionValue, RoutedResolutionValue, CollectedValue
public GroupedCollectedValue, FoldValue
public GroupedCollectedValue, KeyedValue, KeyedContribution, FoldValue
public NewKeyIdentity, RetainAllKeys, DropIdentityKeys
import Adapt
import Atomix
import KernelAbstractions
Expand Down Expand Up @@ -89,6 +90,7 @@ include("execution/ordered_fold_stage.jl")
include("execution/fixed_lane_support.jl")
include("execution/collect_physical_support.jl")
include("execution/collect_stage.jl")
include("execution/keyed_reduce_stage.jl")
include("execution/stage_program_kernelabstractions.jl")
include("execution/stage_program.jl")
include("execution/program_inspection.jl")
Expand Down
36 changes: 22 additions & 14 deletions src/bound_law.jl
Original file line number Diff line number Diff line change
Expand Up @@ -265,28 +265,36 @@ function _require_definite_field_initialization(law::LocalLaw, field::Field)
return nothing
end

function _collect_allocation_schema(law::LocalLaw, collection::Collection)
function _collection_allocation_schema(law::LocalLaw, collection::Collection)
schemas = Any[]
for stage in law.stages, publication in stage.publications
publication.law isa Collect || continue
publication.law isa Union{Collect,KeyedReduce} || continue
any(publication.components) do component
component isa CollectionPublication &&
_same_descriptor(component.collection, collection)
end || continue
law = publication.law
push!(schemas, (
grouped = _is_grouped(law.groups),
groups = Int(_compacted_group_count(law.groups)),
persistent_source_positions =
law.projection isa _PersistentSourcePosition,
source_position_count =
length(stage.source) * _publication_width(law),
))
publication_law = publication.law
if publication_law isa Collect
push!(schemas, (
grouped = _is_grouped(publication_law.groups),
groups = Int(_compacted_group_count(publication_law.groups)),
persistent_source_positions =
publication_law.projection isa _PersistentSourcePosition,
source_position_count =
length(stage.source) * _publication_width(publication_law),
))
else
push!(schemas, (
grouped = false, groups = 1,
persistent_source_positions = false,
source_position_count = 0,
))
end
end
isempty(schemas) && throw(LocalMathValidationError(
"an allocated Collection requires a producing Collect publication";
"an allocated Collection requires a producing collection publication";
stage = :bind, contract = :collection_allocation_producer,
expected = :collect_publication,
expected = :collection_publication,
actual = semantic_identity(collection),
))
all(schema -> schema == first(schemas), schemas) || throw(
Expand All @@ -306,7 +314,7 @@ function _collection_allocation(
stage = :bind, contract = :collection_allocation_initialization,
expected = :empty_collection, actual = request.initial,
))
schema = _collect_allocation_schema(law, collection)
schema = _collection_allocation_schema(law, collection)
capacity = Int(collection.capacity)
return CompactedStorage(
backend,
Expand Down
Loading