wip(0.S3): the packages migrated by ludic migrate state packages - every package test green
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
parent
1e8b5b0523
commit
07505e7ef2
294 changed files with 14996 additions and 14675 deletions
|
|
@ -1,70 +1,69 @@
|
|||
# send.ludic - every frame: a batch when one is due, its answer when it lands, the wait doubled
|
||||
# when it failed, the queue on disk every few seconds. One request at a time; nothing waits on it.
|
||||
var telemetry__clock_on: bool = false
|
||||
|
||||
export function telemetry_tick() -> void {
|
||||
if not telemetry__ready { return }
|
||||
telemetry__id_tick()
|
||||
export function telemetry_tick(base_st: mut BaseState, telemetry_st: mut TelemetryState) -> void {
|
||||
if not telemetry_st.telemetry__ready { return }
|
||||
telemetry__id_tick(base_st, telemetry_st)
|
||||
if not telemetry__enabled() {
|
||||
if len(telemetry__q) > 0 or telemetry__h >= 0 { telemetry_forget() }
|
||||
if len(telemetry_st.telemetry__q) > 0 or telemetry_st.telemetry__h >= 0 { telemetry_forget(telemetry_st) }
|
||||
return
|
||||
}
|
||||
if not telemetry__can_send() { return }
|
||||
if not telemetry__can_send(telemetry_st) { return }
|
||||
let now = telemetry__now()
|
||||
let c = telemetry__conf()
|
||||
if not telemetry__clock_on {
|
||||
telemetry__clock_on = true
|
||||
telemetry__next_ms = now + c.flush_ms
|
||||
telemetry__saved_ms = now
|
||||
let c = telemetry__conf(telemetry_st)
|
||||
if not telemetry_st.telemetry__clock_on {
|
||||
telemetry_st.telemetry__clock_on = true
|
||||
telemetry_st.telemetry__next_ms = now + c.flush_ms
|
||||
telemetry_st.telemetry__saved_ms = now
|
||||
}
|
||||
if now - telemetry__saved_ms >= c.save_ms {
|
||||
telemetry_save()
|
||||
telemetry__saved_ms = now
|
||||
if now - telemetry_st.telemetry__saved_ms >= c.save_ms {
|
||||
telemetry_save(telemetry_st)
|
||||
telemetry_st.telemetry__saved_ms = now
|
||||
}
|
||||
if telemetry__h >= 0 {
|
||||
telemetry__answer(now)
|
||||
if telemetry_st.telemetry__h >= 0 {
|
||||
telemetry__answer(base_st, telemetry_st, now)
|
||||
return
|
||||
}
|
||||
if len(telemetry__q) == 0 { return }
|
||||
let early = telemetry__fails == 0 and len(telemetry__q) >= c.flush_at
|
||||
if early or now >= telemetry__next_ms {
|
||||
telemetry_save()
|
||||
telemetry__flush(now)
|
||||
if telemetry__fails == 0 { telemetry__next_ms = now + c.flush_ms }
|
||||
if len(telemetry_st.telemetry__q) == 0 { return }
|
||||
let early = telemetry_st.telemetry__fails == 0 and len(telemetry_st.telemetry__q) >= c.flush_at
|
||||
if early or now >= telemetry_st.telemetry__next_ms {
|
||||
telemetry_save(telemetry_st)
|
||||
telemetry__flush(base_st, telemetry_st, now)
|
||||
if telemetry_st.telemetry__fails == 0 { telemetry_st.telemetry__next_ms = now + c.flush_ms }
|
||||
}
|
||||
}
|
||||
|
||||
function telemetry__answer(now: int) -> void {
|
||||
let st = telemetry__poll(telemetry__h)
|
||||
function telemetry__answer(base_st: mut BaseState, telemetry_st: mut TelemetryState, now: int) -> void {
|
||||
let st = telemetry__poll(telemetry_st.telemetry__h)
|
||||
if st < 0 { return }
|
||||
telemetry__h = -1
|
||||
telemetry_st.telemetry__h = -1
|
||||
if st > 0 {
|
||||
telemetry__queue_drop(telemetry__sent_n)
|
||||
telemetry__fails = 0
|
||||
telemetry__next_ms = now + telemetry__conf().flush_ms
|
||||
telemetry__fact(TELEMETRY_SENT, telemetry__sent_n, st, 0)
|
||||
telemetry__queue_drop(telemetry_st, telemetry_st.telemetry__sent_n)
|
||||
telemetry_st.telemetry__fails = 0
|
||||
telemetry_st.telemetry__next_ms = now + telemetry__conf(telemetry_st).flush_ms
|
||||
telemetry__fact(base_st, telemetry_st, TELEMETRY_SENT, telemetry_st.telemetry__sent_n, st, 0)
|
||||
} else {
|
||||
telemetry__backoff(now)
|
||||
telemetry__fact(TELEMETRY_FAILED, telemetry__sent_n, st, telemetry__next_ms - now)
|
||||
telemetry__backoff(telemetry_st, now)
|
||||
telemetry__fact(base_st, telemetry_st, TELEMETRY_FAILED, telemetry_st.telemetry__sent_n, st, telemetry_st.telemetry__next_ms - now)
|
||||
}
|
||||
telemetry_save()
|
||||
telemetry_save(telemetry_st)
|
||||
}
|
||||
|
||||
function telemetry__flush(now: int) -> void {
|
||||
if telemetry__h >= 0 or len(telemetry__id) == 0 or len(telemetry__q) == 0 { return }
|
||||
let n = min(telemetry__conf().batch, len(telemetry__q))
|
||||
let h = telemetry__send(telemetry_batch_body(n))
|
||||
function telemetry__flush(base_st: mut BaseState, telemetry_st: mut TelemetryState, now: int) -> void {
|
||||
if telemetry_st.telemetry__h >= 0 or len(telemetry_st.telemetry__id) == 0 or len(telemetry_st.telemetry__q) == 0 { return }
|
||||
let n = min(telemetry__conf(telemetry_st).batch, len(telemetry_st.telemetry__q))
|
||||
let h = telemetry__send(telemetry_st, telemetry_batch_body(telemetry_st, n))
|
||||
if h < 0 {
|
||||
telemetry__backoff(now)
|
||||
telemetry__fact(TELEMETRY_FAILED, n, 0, telemetry__next_ms - now)
|
||||
telemetry__backoff(telemetry_st, now)
|
||||
telemetry__fact(base_st, telemetry_st, TELEMETRY_FAILED, n, 0, telemetry_st.telemetry__next_ms - now)
|
||||
return
|
||||
}
|
||||
telemetry__h = h
|
||||
telemetry__sent_n = n
|
||||
telemetry_st.telemetry__h = h
|
||||
telemetry_st.telemetry__sent_n = n
|
||||
}
|
||||
|
||||
# a failed send: wait twice as long next time, up to backoff_max doublings
|
||||
function telemetry__backoff(now: int) -> void {
|
||||
if telemetry__fails < telemetry__conf().backoff_max { telemetry__fails += 1 }
|
||||
telemetry__next_ms = now + (telemetry__conf().flush_ms << telemetry__fails)
|
||||
function telemetry__backoff(telemetry_st: mut TelemetryState, now: int) -> void {
|
||||
if telemetry_st.telemetry__fails < telemetry__conf(telemetry_st).backoff_max { telemetry_st.telemetry__fails += 1 }
|
||||
telemetry_st.telemetry__next_ms = now + (telemetry__conf(telemetry_st).flush_ms << telemetry_st.telemetry__fails)
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue