Skip to content
Closed
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
8 changes: 4 additions & 4 deletions .github/workflows/CI.yml
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ jobs:
test:
name: Julia ${{ matrix.version }} - ${{ matrix.os }} (${{ matrix.arch }})
runs-on: ${{ matrix.os }}
timeout-minutes: 150
timeout-minutes: 180
if: ${{ !contains(github.event.head_commit.message, '[skip tests]') }}
env:
JULIA_NUM_THREADS: '1'
Expand All @@ -32,7 +32,7 @@ jobs:
- {version: '1.12', os: ubuntu-latest, arch: x64}
- {version: '1', os: ubuntu-latest, arch: x64}
- {version: '1', os: windows-latest, arch: x64}
- {version: '1', os: macos-15, arch: x64}
- {version: '1', os: macos-26, arch: arm64}
- {version: 'nightly', os: ubuntu-latest, arch: x64}
steps:
- uses: actions/checkout@v4
Expand All @@ -55,7 +55,7 @@ jobs:
test-multithreaded:
name: Julia ${{ matrix.version }} - ${{ matrix.os }} (${{ matrix.arch }}) (multithreaded)
runs-on: ${{ matrix.os }}
timeout-minutes: 90
timeout-minutes: 120
if: ${{ !contains(github.event.head_commit.message, '[skip tests]') }}
strategy:
fail-fast: false
Expand All @@ -66,7 +66,7 @@ jobs:
- {version: '1.12', os: ubuntu-latest, arch: x64}
- {version: '1', os: ubuntu-latest, arch: x64}
- {version: '1', os: windows-latest, arch: x64}
- {version: '1', os: macos-15, arch: x64}
- {version: '1', os: macos-26, arch: arm64}
- {version: 'nightly', os: ubuntu-latest, arch: x64}
steps:
- uses: actions/checkout@v4
Expand Down
6 changes: 4 additions & 2 deletions docs/src/darray.md
Original file line number Diff line number Diff line change
Expand Up @@ -705,8 +705,10 @@ From `LinearAlgebra`:
- `mul!` (In-place Matrix-Matrix and Matrix-Vector multiply)
- `cholesky`/`cholesky!` (In-place/Out-of-place Cholesky factorization)
- `lu`/`lu!` (In-place/Out-of-place LU factorization (`NoPivot` and `RowMaximum`))
- `\`/`ldiv!` (In-place/Out-of-place Linear solving with LU and Cholesky factorizations)
- `inv` (Out-of-place matrix inversion)
- `qr`/`qr!` (In-place/Out-of-place QR factorization)
- `\`/`ldiv!` (In-place/Out-of-place Linear solving with LU, Cholesky, QR, and SVD factorizations)
- `inv` (Out-of-place matrix inversion, including via SVD)
- `svd`/`svd!`/`svdvals!` (In-place/Out-of-place Singular Value Decomposition)

From `AbstractFFTs`:
- `fft`/`fft!`
Expand Down
65 changes: 56 additions & 9 deletions lib/TimespanLogging/src/core.jl
Original file line number Diff line number Diff line change
Expand Up @@ -206,32 +206,53 @@ function Base.setindex!(ml::MultiEventLog, c, name::Symbol)
ml.consumers[name] = c
end

function get_state(ml::MultiEventLog)
lock(event_log_lock) do
mls = get!(()->MultiEventLogState(), MultiEventLogState_PLS, ml.uid)
max_length = reduce(max, map(length, values(mls.consumer_logs)); init=0)
# Resolve (and lazily initialize) the process-local state for `ml`. The caller
# MUST already hold `event_log_lock`. Split out from `get_state` so the hot
# `write_event` path takes the lock exactly once (was twice: get_state + write).
function _get_state_locked(ml::MultiEventLog)
mls = get!(()->MultiEventLogState(), MultiEventLogState_PLS, ml.uid)
# Fast path: the consumer set is stable, so there is nothing to initialize.
# Consumers are only ever added (never removed -- see FIXME), so a length
# mismatch is a sufficient and cheap "needs init" test, and it lets us skip
# the per-event `map(length, ...)` allocation in steady state.
if length(mls.consumers) != length(ml.consumers)
max_length = 0
for v in values(mls.consumer_logs)
l = length(v)
if l > max_length
max_length = l
end
end
for name in keys(ml.consumers)
if !haskey(mls.consumers, name)
mls.consumers[name] = init_similar(ml.consumers[name])
mls.consumer_logs[name] = Vector{Any}(fill(nothing, max_length))
end
end
end
if length(mls.aggregators) != length(ml.aggregators)
for name in keys(ml.aggregators)
if !haskey(mls.aggregators, name)
mls.aggregators[name] = init_similar(ml.aggregators[name])
end
end
# FIXME: Remove deleted consumers and aggregators
mls
end
# FIXME: Remove deleted consumers and aggregators
return mls
end

function get_state(ml::MultiEventLog)
lock(event_log_lock) do
_get_state_locked(ml)
end
end

"Creates a copy of `x` with the same configuration, but fresh/empty data."
init_similar(x) = x

function write_event(ml::MultiEventLog, event::Event)
mls = get_state(ml)
lock(event_log_lock) do
mls = _get_state_locked(ml)
for name in keys(mls.consumers)
cevent = try
mls.consumers[name](event)
Expand Down Expand Up @@ -278,16 +299,30 @@ end

empty_prof() = ProfilerResult(UInt[], Dict{UInt64, Vector{Base.StackTraces.StackFrame}}(), UInt[])

# Shared, immutable-in-practice empty profiler result. Every non-profiled event
# carries *this* object instead of allocating a fresh `ProfilerResult` (2 vectors
# + a dict) per event. It is only ever read (e.g. `mix_samples` does a read-only
# `vcat`), never mutated, on the non-profiling path.
const EMPTY_PROF = empty_prof()
const _NO_TASKS = Task[]

# Profiling is opt-in and rare; this flag lets the (very hot) `timespan_finish`
# and per-compute-task `prof_task_put!` paths skip `prof_lock` entirely when
# profiling is off. `Dagger.enable_logging!(profile=true)` sets it.
const PROFILE_TASKS = Ref{Bool}(false)

const prof_refcount = Ref{Threads.Atomic{Int}}(Threads.Atomic{Int}(0))
const prof_lock = Threads.ReentrantLock()
const prof_tasks = IdDict{Any, Vector{Task}}()

function prof_task_put!(id, task::Task=Base.current_task())
PROFILE_TASKS[] || return
lock(prof_lock) do
push!(get!(()->Task[], prof_tasks, id), task)
end
end
function prof_tasks_take!(id)
PROFILE_TASKS[] || return _NO_TASKS
lock(prof_lock) do
if haskey(prof_tasks, id)
pop!(prof_tasks, id)
Expand Down Expand Up @@ -317,7 +352,9 @@ function timespan_start(ctx, category::Symbol, @nospecialize(id), @nospecialize(
Profile.start_timer()
end
end
ev = Event(:start, category, id, tl, time_ns(), gc_num(), empty_prof())
# Start events never carry profiler samples (those are gathered at finish),
# so always reuse the shared empty result rather than allocating.
ev = Event(:start, category, id, tl, time_ns(), gc_num(), EMPTY_PROF)
write_event(sink, ev)
nothing
end
Expand All @@ -330,12 +367,22 @@ categorized by `category`, and uniquely identified by `id`; these two must be
the same as previously passed to `timespan_start`. `tl` is the "timeline" of
the event, which is just an arbitrary payload attached to the event.
"""
function timespan_finish(ctx, category::Symbol, @nospecialize(id), @nospecialize(tl); tasks=prof_tasks_take!(id))
function timespan_finish(ctx, category::Symbol, @nospecialize(id), @nospecialize(tl); tasks=nothing)
sink = log_sink(ctx)
isa(sink, NoOpLog) && return
do_profile = profile(ctx, category, id, tl)
time = time_ns()
gcn = gc_num()
if !do_profile
# Hot path: no profiling. Skip `prof_lock`, `Profile.fetch`, and the
# per-event `ProfilerResult` allocation by reusing the shared empty one.
ev = Event(:finish, category, id, tl, time, gcn, EMPTY_PROF)
write_event(sink, ev)
return nothing
end
if tasks === nothing
tasks = prof_tasks_take!(id)
end
prof = UInt[]
lidict = Dict{UInt64, Vector{Base.StackTraces.StackFrame}}()
GC.@preserve tasks begin
Expand Down
2 changes: 2 additions & 0 deletions src/Dagger.jl
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,7 @@ include("datadeps/chunkview.jl")
include("datadeps/remainders.jl")
include("datadeps/scheduling.jl")
include("datadeps/queue.jl")
include("datadeps/hierarchical.jl")

# Stencils
include("utils/haloarray.jl")
Expand Down Expand Up @@ -129,6 +130,7 @@ include("array/cholesky.jl")
include("array/trsm.jl")
include("array/lu.jl")
include("array/qr.jl")
include("array/svd.jl")

# GPU
include("gpu.jl")
Expand Down
5 changes: 5 additions & 0 deletions src/array/darray.jl
Original file line number Diff line number Diff line change
Expand Up @@ -321,6 +321,11 @@ function Base.isequal(x::ArrayOp, y::ArrayOp)
x === y
end

aliasing(x::DArray) =
throw(ConcurrencyViolationError("DArray aliasing may be mixed and unstable"))
memory_space(x::DArray) =
throw(ConcurrencyViolationError("DArray memory spaces may be mixed and unstable"))

Base.similar(D::DArray{T,N} where T, ::Type{S}, dims::Dims{N}) where {S,N} =
DArray{S,N}(undef, D.partitioning, dims)
Base.similar(D::DArray{T,N1} where T, ::Type{S}, dims::Dims{N2}) where {S,N1,N2} =
Expand Down
23 changes: 17 additions & 6 deletions src/array/mul.jl
Original file line number Diff line number Diff line change
Expand Up @@ -487,11 +487,22 @@ function gemv_dagger!(
alpha = T(_alpha)
beta = T(_beta)

if Ant != Bmt
throw(DimensionMismatch(lazy"A has number of blocks ($Amt,$Ant) but B has number of blocks ($Bmt)"))
end
if Amt != Cmt
throw(DimensionMismatch(lazy"A has number of blocks ($Amt,$Ant) but C has number of blocks ($Cmt)"))
# For op(A)*x: when A is not transposed, x matches A's column-blocks and
# C matches A's row-blocks; when A is [conj-]transposed the roles swap.
if transA == 'N'
if Ant != Bmt
throw(DimensionMismatch(lazy"A has number of blocks ($Amt,$Ant) but B has number of blocks ($Bmt)"))
end
if Amt != Cmt
throw(DimensionMismatch(lazy"A has number of blocks ($Amt,$Ant) but C has number of blocks ($Cmt)"))
end
else
if Amt != Bmt
throw(DimensionMismatch(lazy"A' has number of blocks ($Ant,$Amt) but B has number of blocks ($Bmt)"))
end
if Ant != Cmt
throw(DimensionMismatch(lazy"A' has number of blocks ($Ant,$Amt) but C has number of blocks ($Cmt)"))
end
end

Dagger.spawn_datadeps() do
Expand All @@ -510,7 +521,7 @@ function gemv_dagger!(
)
end
else
# A: [Conj]Trans
# A: [Conj]Trans — C's blocks index A's column-blocks
for k in range(1, Amt)
mzone = k == 1 ? beta : T(1.0)
Dagger.@spawn BLAS.gemv!(
Expand Down
Loading
Loading