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
14 changes: 13 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -211,7 +211,19 @@ assembly and coupling.
### Fixed

- Retained outputs and temporal dependency streams snapshot mutable values, so
later in-place model updates do not overwrite historical samples.
later in-place model updates do not overwrite historical samples. This also
covers the initial values of output requests, including selector reentry.
- Temporal history buffers grow when newly added consumers need longer windows.
`PreviousTimeStep` retains the preceding publication even when fractional
clocks or manual calls leave gaps between publications.
- `Many` inputs respect application and process filters for local producers.
Positional `Relation(...)` selectors keep separate bindings for each consumer
after objects are added.
- Newborn initializers can read fully initialized sources registered in the same
lifecycle event. Unrelated distributed outputs no longer hide newborn sources
from `Many` inputs.
- Mixed-type object IDs with the same text remain distinct during sorting and
identity-based lookup.
- Temporal inputs and `final_state` snapshots no longer alias nested mutable
producer values. Explicitly declared environment durations also survive
forcing-source remapping without duplicate fields or replacement.
Expand Down
117 changes: 77 additions & 40 deletions src/composite_model/compilation.jl
Original file line number Diff line number Diff line change
Expand Up @@ -2185,7 +2185,7 @@ function _many_binding_scope_anchor(
consumer_id::ObjectId,
)
selector_criteria = criteria(selector)
!isnothing(_criteria_get(selector_criteria, :relation, nothing)) &&
!isnothing(_criteria_value(selector_criteria, :relation, Relation)) &&
return (:consumer, consumer_id)
explicit_scope = _criteria_scope(selector_criteria)
scope = isnothing(explicit_scope) ?
Expand Down Expand Up @@ -2457,6 +2457,7 @@ function _append_added_many_sources!(
binding.application,
applications_by_id,
distributed_outputs,
applications_by_object,
)
isempty(new_source_ids) && return true

Expand Down Expand Up @@ -2541,6 +2542,9 @@ function _update_structural_many_sources!(
model::CompositeModel,
binding::CompiledModelInputBinding,
dirty_ids,
applications_by_object,
applications_by_id,
distributed_outputs,
)
binding.multiplicity == :many || return :fallback
binding.carrier_hint == :ref_vector || return :fallback
Expand All @@ -2559,6 +2563,14 @@ function _update_structural_many_sources!(
context=binding.consumer_id,
default_to_context=true,
default_scope=default_scope,
) && _matches_input_source_writer(
object_id,
binding.source_var,
binding.process,
binding.application,
applications_by_id,
distributed_outputs,
applications_by_object,
)
was_source == is_source && continue
source_reference = nothing
Expand Down Expand Up @@ -3602,6 +3614,9 @@ function _extend_compiled_scene(
model,
binding,
structural_dirty_ids,
applications_by_object,
applications_by_id,
distributed_outputs,
)
if structural_update != :fallback
processed_many_sources[binding.source_ids] = nothing
Expand Down Expand Up @@ -7147,6 +7162,7 @@ Base.@nospecializeinfer function _push_model_input_binding!(
application_filter,
applications_by_id,
distributed_outputs,
applications_by_object,
)
selector isa Many && sizehint!(source_ids, length(source_ids) + 1)
source_application_ids = if _selector_from_status(selector)
Expand Down Expand Up @@ -7328,63 +7344,84 @@ function _final_many_source_applications(
return isempty(canonical_ids) ? source_application_ids : canonical_ids
end

_filter_many_input_sources_by_writer!(
source_ids,
selector,
source_var,
process_filter,
application_filter,
applications_by_id,
::NoCompiledDistributedOutputs,
) = source_ids

function _filter_many_input_sources_by_writer!(
source_ids,
selector,
source_var::Symbol,
process_filter,
application_filter,
applications_by_id,
distributed_outputs::CompiledDistributedOutputs,
distributed_outputs,
applications_by_object,
)
selector isa Many || return source_ids
_selector_from_status(selector) && return source_ids
isnothing(process_filter) && isnothing(application_filter) &&
return source_ids
filter!(source_ids) do source_id
owners = get(
distributed_outputs.writer_ownership,
(source_id, source_var),
(),
return _matches_input_source_writer(
source_id,
source_var,
process_filter,
application_filter,
applications_by_id,
distributed_outputs,
applications_by_object,
)
owned = any(owners) do owner
source_application = get(
applications_by_id,
owner.application_id,
nothing,
)
isnothing(source_application) && return false
isnothing(process_filter) ||
source_application.process == process_filter || return false
isnothing(application_filter) ||
source_application.id == application_filter || return false
return true
end
owned && return true
# Manual callees (and stream-only local outputs) are not scheduled
# canonical owners. They remain valid explicitly selected producers;
# unrelated distributed outputs must not hide their local targets.
return any(values(applications_by_id)) do application
isnothing(process_filter) || application.process == process_filter || return false
isnothing(application_filter) || application.id == application_filter || return false
return _application_writes_object_variable(
NoCompiledDistributedOutputs(), application, source_id, source_var,
)
end
end
return source_ids
end

function _matches_input_source_writer(
source_id,
source_var,
process_filter,
application_filter,
applications_by_id,
distributed_outputs,
applications_by_object,
)
# This index also contains manual/stream-only writers and the targeted
# initializer overlay, whose newborns are not in compiled target_ids yet.
local_writer = any(get(applications_by_object, source_id, ())) do application
isnothing(process_filter) || application.process == process_filter || return false
isnothing(application_filter) || application.id == application_filter || return false
return haskey(_local_output_schema(application.spec), source_var)
end
local_writer && return true
return _has_matching_distributed_input_writer(
distributed_outputs,
source_id,
source_var,
process_filter,
application_filter,
applications_by_id,
)
end

_has_matching_distributed_input_writer(
::NoCompiledDistributedOutputs,
args...,
) = false

function _has_matching_distributed_input_writer(
distributed_outputs::CompiledDistributedOutputs,
source_id,
source_var,
process_filter,
application_filter,
applications_by_id,
)
owners = get(distributed_outputs.writer_ownership, (source_id, source_var), ())
return any(owners) do owner
application = get(applications_by_id, owner.application_id, nothing)
isnothing(application) && return false
isnothing(process_filter) || application.process == process_filter || return false
isnothing(application_filter) || application.id == application_filter || return false
return true
end
end

function _model_input_names(application::CompiledModelApplication)
return Symbol[Symbol(var) for var in keys(_input_schema(application.spec))]
end
Expand Down
69 changes: 61 additions & 8 deletions src/composite_model/runtime_outputs.jl
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,23 @@ function TemporalDependencyBuffer{T}(capacity::Integer) where {T}
)
end

function _grow_temporal_dependency_buffer!(
buffer::TemporalDependencyBuffer{T},
capacity::Int,
) where {T}
capacity <= length(buffer.times) && return buffer
times = Vector{Float64}(undef, capacity)
values = Vector{T}(undef, capacity)
for index in eachindex(buffer)
times[index], values[index] = buffer[index]
end
# Keep the buffer itself: compiled consumers and publishers share it.
buffer.times = times
buffer.values = values
buffer.first_slot = 1
return buffer
end

Base.IndexStyle(::Type{<:TemporalDependencyBuffer}) = IndexLinear()
Base.size(buffer::TemporalDependencyBuffer) = (buffer.sample_count,)

Expand Down Expand Up @@ -1520,7 +1537,16 @@ function _initialize_model_output_stream!(
sizehint_steps::Integer,
)
key = _model_stream_key(application.id, object_id, variable)
haskey(streams, key) && return streams
if haskey(streams, key)
stream = streams[key]
if stream isa TemporalDependencyBuffer
_grow_temporal_dependency_buffer!(
stream,
_model_dependency_capacity(retention, application.id, variable),
)
end
return streams
end
reference = _model_output_reference(
compiled,
application,
Expand Down Expand Up @@ -1710,7 +1736,10 @@ end
cutoff = output.dependency_horizon <= 0.0 ?
float(time) :
float(time) - output.dependency_horizon + 1.0
while !isempty(stream) &&
# The preceding publication is needed by PreviousTimeStep and linear
# extrapolation even when fractional clocks or manual calls leave a
# larger gap than the nominal cadence. Capacity still bounds storage.
while length(stream) > 2 &&
first(stream)[1] < cutoff - 1.0e-8
_temporal_dependency_popfirst!(stream)
end
Expand Down Expand Up @@ -5610,6 +5639,7 @@ mutable struct _TargetedTopologyRuntime{CS,MA} <:
ObjectId,
Vector{CompiledModelApplication},
}
indexed_addition_count::Int
manual_application_ids::MA
application_positions::Dict{Symbol,Int}
end
Expand Down Expand Up @@ -5644,6 +5674,7 @@ function _targeted_topology_runtime!(
Union{Nothing,_TargetedApplicationSet},
}(),
Dict{ObjectId,Vector{CompiledModelApplication}}(),
0,
compiled.scenario_plan.manual_application_ids,
Dict(
application_id => index
Expand Down Expand Up @@ -5698,17 +5729,14 @@ function _targeted_application_set!(
requested_ids,
)
key = Tuple(requested_ids)
return get!(runtime.application_sets, key) do
application_set = get!(runtime.application_sets, key) do
applications = _new_object_applications(
model,
runtime.compiled,
requested_ids,
)
isnothing(applications) && return nothing
output_applications, added_applications_by_object = applications
# Keep applications for every object targeted earlier in this same
# lifecycle delta. A later newborn can then bind an input to an earlier
# newborn without refreshing the whole scene at a mid-kernel barrier.
merge!(
runtime.added_applications_by_object,
added_applications_by_object,
Expand All @@ -5722,6 +5750,31 @@ function _targeted_application_set!(
false,
)
end
isnothing(application_set) && return nothing

# A fully initialized source can be registered without an explicit call.
# Include its application membership before resolving newborn inputs. The
# append-only delta cursor visits each addition once across a chain of calls;
# targets prepared above already have membership and need no extra lookup.
added = lifecycle_delta(model).added
unindexed_ids = nothing
for index in (runtime.indexed_addition_count + 1):length(added)
object_id = added[index].id
haskey(runtime.added_applications_by_object, object_id) && continue
isnothing(unindexed_ids) && (unindexed_ids = ObjectId[])
push!(unindexed_ids, object_id)
end
if !isnothing(unindexed_ids)
applications = _new_object_applications(
model,
runtime.compiled,
unindexed_ids,
)
isnothing(applications) && return nothing
merge!(runtime.added_applications_by_object, last(applications))
end
runtime.indexed_addition_count = length(added)
return application_set
end

function _targeted_callee_applications(
Expand Down Expand Up @@ -7225,7 +7278,7 @@ function _output_request_target(
OutputRequestMembership(
float(start_time),
nothing,
initial,
deepcopy(initial),
),
],
)
Expand Down Expand Up @@ -7381,7 +7434,7 @@ function _refresh_output_request_targets!(
OutputRequestMembership(
start_time,
nothing,
initial,
deepcopy(initial),
),
)
continue
Expand Down
9 changes: 7 additions & 2 deletions src/composite_model/selectors.jl
Original file line number Diff line number Diff line change
Expand Up @@ -478,8 +478,13 @@ end
function _object_id_isless(left::ObjectId, right::ObjectId)
left_value = left.value
right_value = right.value
if typeof(left_value) === typeof(right_value) &&
hasmethod(isless, Tuple{typeof(left_value),typeof(right_value)})
# Keep each identity type in one ordered group. Mixing natural ordering
# within a type with string ordering between types is not transitive, and
# distinct identities such as `1` and `Symbol("1")` otherwise compare equal.
if typeof(left_value) !== typeof(right_value)
return isless(string(typeof(left_value)), string(typeof(right_value)))
end
if hasmethod(isless, Tuple{typeof(left_value),typeof(right_value)})
return isless(left_value, right_value)
end
return isless(string(left_value), string(right_value))
Expand Down
Loading
Loading