fix(ludic.telemetry): nothing made per event - properties and the line written into kept buffers, the queue a ring of bytes made at the start

An event was a tree of values encoded into a new string, a dropped line was never given back, and
every batch and save joined the queue into another. Now: telemetry_props / telemetry_str / _int /
_bool write the properties as JSON into a buffer the state keeps; the line is written beside it,
its time worked out from the clock's seconds (TelemetryWorld.clock_s, was stamp); its bytes go
into one ring (ring_bytes, at most queue_max lines), the oldest overwritten past either; a batch
and the file are written into a third kept buffer and go as bytes (TelemetryTransport.send takes
the bytes and their length; Http.body_bytes, Fs.write_bytes). The facts are pooled.
tests/ring_test: ten thousand events past the cap in lines and in bytes, 0 bytes of heap; the line
exact JSON, escaped, 2000-02-29 right; the batch puts the id in; the file round-trips.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Orkun ÇAKILKAYA 2026-09-28 15:15:12 +03:00
parent fc0f3d3790
commit d1bea7f10e
14 changed files with 454 additions and 97 deletions

View file

@ -23,22 +23,31 @@ machine. With no salt, or where it cannot be read within ten seconds, it is rand
lives in the game's id file as `{"id": ...}`; any other key the game keeps there is left alone. lives in the game's id file as `{"id": ...}`; any other key the game keeps there is left alone.
An event queued before the id is known carries a mark that the batch replaces. An event queued before the id is known carries a mark that the batch replaces.
## Nothing is made per event
Ludic gives nothing back, so an event is never a tree of values or a string of its own. Its
properties are written as JSON into a kept buffer (`telemetry_props`), the line into another, and
the line's bytes into one ring made at the start (`ring_bytes`, at most `queue_max` lines): past
either, the oldest are overwritten. A batch and the queue file are written into a third kept buffer
and go out as bytes (`Http.body_bytes`, `Fs.write_bytes`). An event's time is worked out from the
clock's seconds, not formatted. `tests/ring_test` holds it: ten thousand events past the cap, 0 bytes.
## Config and ports (`bind`) ## Config and ports (`bind`)
```ludic ```ludic
property TelemetryConfig { property TelemetryConfig {
host, key, path ("/batch/"), lib # an empty host or key sends nothing, ever host, key, path ("/batch/"), lib # an empty host or key sends nothing, ever
queue_file, id_file, id_salt, id_scratch queue_file, id_file, id_salt, id_scratch
flush_ms (30000), flush_at (50), batch (200), queue_max (5000), backoff_max (6), save_ms (5000) flush_ms (30000), flush_at (50), batch (200), queue_max (5000), ring_bytes (1 MB), backoff_max (6), save_ms (5000)
} }
port TelemetryWorld { port TelemetryWorld {
can_send: fn() -> bool # may this run send at all - a test, a headless run (unbound: yes) can_send: fn() -> bool # may this run send at all - a test, a headless run (unbound: yes)
enabled: fn() -> bool # the player's switch (unbound: on) enabled: fn() -> bool # the player's switch (unbound: on)
now_ms: fn() -> int # a clock (unbound: the wall clock, by the second) now_ms: fn() -> int # a clock (unbound: the wall clock, by the second)
stamp: fn() -> string # an event's ISO time (unbound: now) clock_s: fn() -> int # the wall clock in seconds since 1970, for an event's time (unbound: now)
} }
port TelemetryTransport { # unbound: Http port TelemetryTransport { # unbound: Http
send: fn(string, string) -> int # url, body -> a handle, < 0 if it could not start send: fn(string, []byte, int) -> int # url, body, its length -> a handle, < 0 if it could not start
poll: fn(int) -> int # -1 pending, 0 failed, 1 taken poll: fn(int) -> int # -1 pending, 0 failed, 1 taken
drop: fn(int) -> void # abandon one drop: fn(int) -> void # abandon one
} }
@ -56,6 +65,7 @@ bind TelemetryWorld { can_send: fn my_can_send, enabled: fn my_switch }
| --- | --- | | --- | --- |
| `telemetry_config(c)` | before the start | | `telemetry_config(c)` | before the start |
| `telemetry_start()` | once, at boot: the id, a session, and the queue an earlier start left (or none, if off) | | `telemetry_start()` | once, at boot: the id, a session, and the queue an earlier start left (or none, if off) |
| `telemetry_props()`, `telemetry_str(p, k, s)`, `telemetry_int(p, k, n)`, `telemetry_bool(p, k, b)` | the next event's properties, written into a buffer the state keeps |
| `telemetry_event(name, props)` | queue one (props may be null); nothing is kept while off or while the run may not send | | `telemetry_event(name, props)` | queue one (props may be null); nothing is kept while off or while the run may not send |
| `telemetry_tick()` | every frame: the machine id's child, a batch when due (every `flush_ms`, or at `flush_at` events), its answer, the backoff, the queue on disk every `save_ms`; off forgets | | `telemetry_tick()` | every frame: the machine id's child, a batch when due (every `flush_ms`, or at `flush_at` events), its answer, the backoff, the queue on disk every `save_ms`; off forgets |
| `telemetry_save()` | the queue on disk now (a quit, a child about to start) | | `telemetry_save()` | the queue on disk now (a quit, a child about to start) |

View file

@ -0,0 +1,16 @@
# facts.ludic - a fact's record reused (the ludic.wildlife idiom): a batch sent or failed is a fact
# a fact record no drain holds any more, made (and kept) only when every one is held
function telemetry__fact_record(telemetry_st: mut TelemetryState) -> TelemetryFact {
if telemetry_st.telemetry__fnext[0] >= len(telemetry_st.telemetry__ffree) {
q_unheld(telemetry_st.telemetry__facts, telemetry_st.telemetry__fpool, telemetry_st.telemetry__ffree)
telemetry_st.telemetry__fnext[0] = 0
}
if telemetry_st.telemetry__fnext[0] < len(telemetry_st.telemetry__ffree) {
let f = telemetry_st.telemetry__ffree[telemetry_st.telemetry__fnext[0]]
telemetry_st.telemetry__fnext[0] += 1
return f
}
let f = new TelemetryFact
push(telemetry_st.telemetry__fpool, f)
return f
}

View file

@ -4,8 +4,13 @@ module ludic_telemetry uses ludic_base
numbers float numbers float
import "ludic.base" import "ludic.base"
import "state.ludic" import "state.ludic"
import "facts.ludic"
import "ports.ludic" import "ports.ludic"
import "props.ludic"
import "stamp.ludic"
import "ring.ludic"
import "queue.ludic" import "queue.ludic"
import "load.ludic"
import "send.ludic" import "send.ludic"
import "ident.ludic" import "ident.ludic"
import "queries.ludic" import "queries.ludic"

View file

@ -0,0 +1,46 @@
# load.ludic - the unsent lines a start finds on disk, read once into the ring with their marks found
function telemetry__queue_load(telemetry_st: mut TelemetryState) -> void {
let f = telemetry__conf(telemetry_st).queue_file
if len(f) == 0 { return }
let b = Fs.read_bytes(f)
if b == null { return }
var start = 0
for i in 0 .. len(b) + 1 {
if i == len(b) or b[i] == 10 {
if i - start > 2 { tq_load_line(telemetry_st, b, start, i - start) }
start = i + 1
}
}
telemetry_st.telemetry__dirty = false
}
# a line read back (once, at a start): through the kept line buffer, its mark found again
function tq_load_line(telemetry_st: mut TelemetryState, b: []byte, at: int, n: int) -> void {
let sb = telemetry_st.telemetry__line
sb_clear(sb)
for i in 0 .. n { sb_byte(sb, b[at + i]) }
if tj_full(sb) { return }
tr_add(telemetry_st, sb.b, n, tq_find_mark(sb.b, n))
}
function tq_find_mark(b: []byte, n: int) -> int {
let m: string = TELEMETRY_MARK
let k = len(m)
for i in 0 .. n - k + 1 {
var same = true
for j in 0 .. k { if same and b[i + j] != m[j] { same = false } }
if same { return i }
}
return -1
}
function telemetry__read_obj(path: string) -> Val {
if len(path) == 0 { return null }
let s = Fs.read_text(path)
if s == null { return null }
let v = Json.parse(s)
if v == null { return null }
if Value.kind(v) != 6 {
Json.free(v)
return null
}
return v
}

View file

@ -1,5 +1,8 @@
# ports.ludic - the config, and what the ports answer when the game leaves them unbound # ports.ludic - the config, and what the ports answer when the game leaves them unbound
export function telemetry_config(telemetry_st: mut TelemetryState, c: TelemetryConfig) -> void { telemetry_st.telemetry__cfg = c } export function telemetry_config(telemetry_st: mut TelemetryState, c: TelemetryConfig) -> void {
telemetry_st.telemetry__cfg = c
telemetry_st.telemetry__url = c.host + c.path
}
function telemetry__conf(telemetry_st: TelemetryState) -> TelemetryConfig { function telemetry__conf(telemetry_st: TelemetryState) -> TelemetryConfig {
return telemetry_st.telemetry__cfg return telemetry_st.telemetry__cfg
@ -7,7 +10,7 @@ function telemetry__conf(telemetry_st: TelemetryState) -> TelemetryConfig {
function telemetry__yes() -> bool { return true } function telemetry__yes() -> bool { return true }
function telemetry__wall_ms(telemetry_st: TelemetryState) -> int { return (Time.now() - telemetry_st.telemetry__t0) * 1000 } function telemetry__wall_ms(telemetry_st: TelemetryState) -> int { return (Time.now() - telemetry_st.telemetry__t0) * 1000 }
function telemetry__wall_stamp() -> string { return DateTime.format(Time.now(), "YYYY-MM-DDTHH:mm:ssZ") } function telemetry__wall_s() -> int { return Time.now() }
function telemetry__can_send(telemetry_st: TelemetryState) -> bool { function telemetry__can_send(telemetry_st: TelemetryState) -> bool {
let c = telemetry__conf(telemetry_st) let c = telemetry__conf(telemetry_st)
@ -17,13 +20,12 @@ function telemetry__can_send(telemetry_st: TelemetryState) -> bool {
function telemetry__enabled() -> bool { return TelemetryWorld.enabled() } function telemetry__enabled() -> bool { return TelemetryWorld.enabled() }
function telemetry__now() -> int { return TelemetryWorld.now_ms() } function telemetry__now() -> int { return TelemetryWorld.now_ms() }
function telemetry__stamp() -> string { return TelemetryWorld.stamp() }
function telemetry__http_send(url: string, body: string) -> int { function telemetry__http_send(url: string, body: []byte, n: int) -> int {
let h = Http.open("POST", url) let h = Http.open("POST", url)
if h < 0 { return h } if h < 0 { return h }
Http.set(h, "Content-Type", "application/json") Http.set(h, "Content-Type", "application/json")
Http.body(h, body) Http.body_bytes(h, data_of(body), n)
Http.send(h) Http.send(h)
return h return h
} }
@ -39,7 +41,7 @@ function telemetry__http_poll(h: int) -> int {
function telemetry__http_drop(h: int) -> void { Http.free(h) } function telemetry__http_drop(h: int) -> void { Http.free(h) }
function telemetry__send(telemetry_st: TelemetryState, body: string) -> int { return TelemetryTransport.send(telemetry__conf(telemetry_st).host + telemetry__conf(telemetry_st).path, body) } function telemetry__send(telemetry_st: TelemetryState, body: []byte, n: int) -> int { return TelemetryTransport.send(telemetry_st.telemetry__url, body, n) }
function telemetry__poll(h: int) -> int { return TelemetryTransport.poll(h) } function telemetry__poll(h: int) -> int { return TelemetryTransport.poll(h) }
# a request abandoned (the player said no) # a request abandoned (the player said no)

View file

@ -0,0 +1,60 @@
# props.ludic - an event's properties written as JSON into a buffer the state keeps, never a tree of
# values: a game says an event often, and Ludic gives nothing back
export const TELEMETRY_LINE: int = 2048 # the longest line an event may be; a longer one is dropped
export property TelemetryProps {
sb: StrBuf = null
n: int = 0 # properties written
}
# the writer, emptied, for the next event's properties
export function telemetry_props(telemetry_st: mut TelemetryState) -> TelemetryProps {
let p = telemetry_st.telemetry__props
sb_clear(p.sb)
p.n = 0
return p
}
export function telemetry_str(p: TelemetryProps, k: string, s: string) -> void {
tp_key(p, k)
tj_str(p.sb, s)
}
export function telemetry_int(p: TelemetryProps, k: string, v: int) -> void {
tp_key(p, k)
sb_int(p.sb, v)
}
export function telemetry_bool(p: TelemetryProps, k: string, b: bool) -> void {
tp_key(p, k)
if b { sb_add(p.sb, "true") } else { sb_add(p.sb, "false") }
}
function tp_key(p: TelemetryProps, k: string) -> void {
if p.n > 0 { sb_byte(p.sb, 44) }
tj_str(p.sb, k)
sb_byte(p.sb, 58)
p.n += 1
}
# a JSON string: quoted, a quote or a backslash escaped, a control byte as \u00XX
function tj_str(sb: StrBuf, s: string) -> void {
sb_byte(sb, 34)
if s != null {
let sp: string = s
for i in 0 .. len(sp) {
let c = sp[i]
if c == 34 or c == 92 {
sb_byte(sb, 92)
sb_byte(sb, c)
} else if c < 32 {
sb_add(sb, "\\u00")
sb_byte(sb, tj_hex(c / 16))
sb_byte(sb, tj_hex(c % 16))
} else { sb_byte(sb, c) }
}
}
sb_byte(sb, 34)
}
function tj_hex(v: int) -> int {
if v < 10 { return 48 + v }
return 87 + v
}
# a line that reached the end of its buffer was cut: it is dropped rather than sent broken
function tj_full(sb: StrBuf) -> bool { return sb_len(sb) >= len(sb.b) - 1 }

View file

@ -7,13 +7,14 @@ export function telemetry_fails(telemetry_st: TelemetryState) -> int { return te
export function telemetry_next_ms(telemetry_st: TelemetryState) -> int { return telemetry_st.telemetry__next_ms } export function telemetry_next_ms(telemetry_st: TelemetryState) -> int { return telemetry_st.telemetry__next_ms }
export function telemetry_in_flight(telemetry_st: TelemetryState) -> bool { return telemetry_st.telemetry__h >= 0 } export function telemetry_in_flight(telemetry_st: TelemetryState) -> bool { return telemetry_st.telemetry__h >= 0 }
export function telemetry_queued(telemetry_st: TelemetryState) -> int { export function telemetry_queued(telemetry_st: TelemetryState) -> int { return telemetry_st.telemetry__n }
if telemetry_st.telemetry__q == null { return 0 }
return len(telemetry_st.telemetry__q)
}
# a queued line as it will be sent, player mark and all # a queued line as it will be sent, player mark and all
export function telemetry_line(telemetry_st: TelemetryState, i: int) -> string { return telemetry_st.telemetry__q[i] } export function telemetry_line(telemetry_st: TelemetryState, i: int) -> string {
let sb = sb_new(TELEMETRY_LINE)
tr_line_to(telemetry_st, i, sb, null)
return text_of(sb.b, sb_len(sb))
}
# back to before telemetry_start, ports and config kept (tests; a start after a failed boot) # back to before telemetry_start, ports and config kept (tests; a start after a failed boot)
export function telemetry_reset(telemetry_st: mut TelemetryState) -> void { export function telemetry_reset(telemetry_st: mut TelemetryState) -> void {
@ -21,7 +22,7 @@ export function telemetry_reset(telemetry_st: mut TelemetryState) -> void {
telemetry_st.telemetry__ready = false telemetry_st.telemetry__ready = false
telemetry_st.telemetry__id = "" telemetry_st.telemetry__id = ""
telemetry_st.telemetry__session = "" telemetry_st.telemetry__session = ""
telemetry_st.telemetry__q = null if telemetry_st.telemetry__ring != null { tr_clear(telemetry_st) }
telemetry_st.telemetry__dirty = false telemetry_st.telemetry__dirty = false
telemetry_st.telemetry__sent_n = 0 telemetry_st.telemetry__sent_n = 0
telemetry_st.telemetry__next_ms = 0 telemetry_st.telemetry__next_ms = 0

View file

@ -1,4 +1,4 @@
# queue.ludic - a start, an event encoded once into a line, the queue on disk, and forgetting # queue.ludic - a start, an event written once into a line in the ring, the queue on disk, and forgetting
# once, at boot: who this is (from the id file, or begun from the machine), a new session, and # once, at boot: who this is (from the id file, or begun from the machine), a new session, and
# whatever an earlier start left unsent - or nothing, if the player has it off # whatever an earlier start left unsent - or nothing, if the player has it off
export function telemetry_start(telemetry_st: mut TelemetryState) -> void { export function telemetry_start(telemetry_st: mut TelemetryState) -> void {
@ -6,7 +6,7 @@ export function telemetry_start(telemetry_st: mut TelemetryState) -> void {
telemetry_st.telemetry__ready = true telemetry_st.telemetry__ready = true
telemetry_st.telemetry__t0 = Time.now() telemetry_st.telemetry__t0 = Time.now()
telemetry_st.telemetry__session = telemetry_random_hex(telemetry_st, 16) telemetry_st.telemetry__session = telemetry_random_hex(telemetry_st, 16)
telemetry_st.telemetry__q = new []string tr_init(telemetry_st)
telemetry_st.telemetry__id = sv_str(telemetry__read_obj(telemetry__conf(telemetry_st).id_file), "id", "") telemetry_st.telemetry__id = sv_str(telemetry__read_obj(telemetry__conf(telemetry_st).id_file), "id", "")
if len(telemetry_st.telemetry__id) == 0 { telemetry__id_begin(telemetry_st) } if len(telemetry_st.telemetry__id) == 0 { telemetry__id_begin(telemetry_st) }
if telemetry__enabled() { telemetry__queue_load(telemetry_st) } else { telemetry_forget(telemetry_st) } if telemetry__enabled() { telemetry__queue_load(telemetry_st) } else { telemetry_forget(telemetry_st) }
@ -15,87 +15,80 @@ export function telemetry_start(telemetry_st: mut TelemetryState) -> void {
# may an event be kept right now? # may an event be kept right now?
export function telemetry_on(telemetry_st: TelemetryState) -> bool { return telemetry_st.telemetry__ready and telemetry__enabled() and telemetry__can_send(telemetry_st) } export function telemetry_on(telemetry_st: TelemetryState) -> bool { return telemetry_st.telemetry__ready and telemetry__enabled() and telemetry__can_send(telemetry_st) }
# an event: the properties (null for none) with the session and the library added, and the player # an event: its properties (null for none) with the session and the library added, the player as a
# as a mark put in when the batch goes, so an event queued before the id is known is not misfiled # mark put in when the batch goes, so an event queued before the id is known is not misfiled. The
export function telemetry_event(telemetry_st: mut TelemetryState, name: string, props: Val) -> void { # line is written into a kept buffer and copied into the ring: nothing is made for it
export function telemetry_event(telemetry_st: mut TelemetryState, name: string, props: TelemetryProps) -> void {
if not telemetry_on(telemetry_st) { return } if not telemetry_on(telemetry_st) { return }
var p = props let sb = telemetry_st.telemetry__line
if p == null { p = Value.object() } sb_clear(sb)
Value.put(p, "$session_id", Value.str(telemetry_st.telemetry__session)) sb_add(sb, "{\"event\":")
if len(telemetry__conf(telemetry_st).lib) > 0 { Value.put(p, "$lib", Value.str(telemetry__conf(telemetry_st).lib)) } tj_str(sb, name)
let e = Value.object() sb_add(sb, ",\"distinct_id\":\"")
Value.put(e, "event", Value.str(name)) let mk = sb_len(sb)
Value.put(e, "distinct_id", Value.str(TELEMETRY_MARK)) sb_add(sb, TELEMETRY_MARK)
Value.put(e, "properties", p) sb_add(sb, "\",\"properties\":{")
Value.put(e, "timestamp", Value.str(telemetry__stamp())) if props != null and props.n > 0 {
push(telemetry_st.telemetry__q, Json.encode(e)) for i in 0 .. sb_len(props.sb) { sb_byte(sb, props.sb.b[i]) }
telemetry_st.telemetry__dirty = true sb_byte(sb, 44)
let c = telemetry__conf(telemetry_st) }
if len(telemetry_st.telemetry__q) > c.queue_max + c.batch { telemetry__queue_drop(telemetry_st, len(telemetry_st.telemetry__q) - c.queue_max) } sb_add(sb, "\"$session_id\":")
tj_str(sb, telemetry_st.telemetry__session)
if len(telemetry__conf(telemetry_st).lib) > 0 {
sb_add(sb, ",\"$lib\":")
tj_str(sb, telemetry__conf(telemetry_st).lib)
}
sb_add(sb, "},\"timestamp\":")
tj_stamp(sb, TelemetryWorld.clock_s())
sb_byte(sb, 125)
if tj_full(sb) or (props != null and tj_full(props.sb)) { return }
tr_add(telemetry_st, sb.b, sb_len(sb), mk)
} }
# the queue on disk now, if it changed (a quit, a crash about to happen, a child starting) # the queue on disk now, if it changed (a quit, a crash about to happen, a child starting): the lines,
# marks and all, a line each, written from the kept buffer
export function telemetry_save(telemetry_st: mut TelemetryState) -> void { export function telemetry_save(telemetry_st: mut TelemetryState) -> void {
if not telemetry_st.telemetry__dirty or telemetry_st.telemetry__q == null { return } if not telemetry_st.telemetry__dirty or telemetry_st.telemetry__ring == null { return }
telemetry_st.telemetry__dirty = false telemetry_st.telemetry__dirty = false
let f = telemetry__conf(telemetry_st).queue_file let f = telemetry__conf(telemetry_st).queue_file
if len(f) == 0 { return } if len(f) == 0 { return }
if len(telemetry_st.telemetry__q) == 0 { if telemetry_st.telemetry__n == 0 {
if Fs.exists(f) { Fs.remove(f) } if Fs.exists(f) { Fs.remove(f) }
return return
} }
Fs.write_text(f, Text.join(telemetry_st.telemetry__q, "\n") + "\n") let out = telemetry_st.telemetry__out
sb_clear(out)
for i in 0 .. telemetry_st.telemetry__n {
tr_line_to(telemetry_st, i, out, null)
sb_byte(out, 10)
}
Fs.write_bytes(f, out.b, sb_len(out))
} }
# the player turned it off: nothing more is sent, and nothing queued is kept, here or on disk # the player turned it off: nothing more is sent, and nothing queued is kept, here or on disk
export function telemetry_forget(telemetry_st: mut TelemetryState) -> void { export function telemetry_forget(telemetry_st: mut TelemetryState) -> void {
telemetry__drop_request(telemetry_st) telemetry__drop_request(telemetry_st)
telemetry_st.telemetry__q = new []string if telemetry_st.telemetry__ring != null { tr_clear(telemetry_st) }
telemetry_st.telemetry__dirty = false telemetry_st.telemetry__dirty = false
let f = telemetry__conf(telemetry_st).queue_file let f = telemetry__conf(telemetry_st).queue_file
if len(f) > 0 and Fs.exists(f) { Fs.remove(f) } if len(f) > 0 and Fs.exists(f) { Fs.remove(f) }
} }
function telemetry__queue_load(telemetry_st: mut TelemetryState) -> void { # the batch: the first n lines joined with the player id put in, into the kept buffer (tb_batch), and
telemetry_st.telemetry__q = new []string # as text for a test
let f = telemetry__conf(telemetry_st).queue_file function tb_batch(telemetry_st: TelemetryState, n: int) -> void {
if len(f) == 0 { return } let out = telemetry_st.telemetry__out
let s = Fs.read_text(f) sb_clear(out)
if s == null { return } sb_add(out, "{\"api_key\":")
let parts = Text.split(s, "\n") tj_str(out, telemetry__conf(telemetry_st).key)
for i in 0 .. len(parts) { if len(parts[i]) > 2 { push(telemetry_st.telemetry__q, parts[i]) } } sb_add(out, ",\"batch\":[")
let max = telemetry__conf(telemetry_st).queue_max for i in 0 .. n {
if len(telemetry_st.telemetry__q) > max { telemetry__queue_drop(telemetry_st, len(telemetry_st.telemetry__q) - max) } if i > 0 { sb_byte(out, 44) }
} tr_line_to(telemetry_st, i, out, telemetry_st.telemetry__id)
# drop the first n lines (sent, or the oldest past the cap)
function telemetry__queue_drop(telemetry_st: mut TelemetryState, n: int) -> void {
let rest = new []string
for i in n .. len(telemetry_st.telemetry__q) { push(rest, telemetry_st.telemetry__q[i]) }
telemetry_st.telemetry__q = rest
telemetry_st.telemetry__dirty = true
}
# the batch: the first n lines joined, with the player id put in - no parsing
export function telemetry_batch_body(telemetry_st: TelemetryState, n: int) -> string {
let part = new []string
for i in 0 .. n { push(part, telemetry_st.telemetry__q[i]) }
let head = Value.object()
Value.put(head, "api_key", Value.str(telemetry__conf(telemetry_st).key))
let h: string = Json.encode(head)
let body = h[0..len(h) - 1] + ",\"batch\":[" + Text.join(part, ",") + "]}"
return Text.replace(body, TELEMETRY_MARK, telemetry_st.telemetry__id)
}
function telemetry__read_obj(path: string) -> Val {
if len(path) == 0 { return null }
let s = Fs.read_text(path)
if s == null { return null }
let v = Json.parse(s)
if v == null { return null }
if Value.kind(v) != 6 {
Json.free(v)
return null
} }
return v sb_add(out, "]}")
}
export function telemetry_batch_body(telemetry_st: TelemetryState, n: int) -> string {
tb_batch(telemetry_st, Math.min(n, telemetry_st.telemetry__n))
return text_of(telemetry_st.telemetry__out.b, sb_len(telemetry_st.telemetry__out))
} }

View file

@ -0,0 +1,76 @@
# ring.ludic - the queued lines in one byte ring made at the start: a line's bytes, its length and where
# its player mark is, oldest first. Full (lines or bytes), the oldest are overwritten; nothing is made
# per event, per drop or per send
# the ring for this start, made once at its size and kept across starts of the same size
function tr_init(telemetry_st: mut TelemetryState) -> void {
let c = telemetry__conf(telemetry_st)
if telemetry_st.telemetry__ring == null or len(telemetry_st.telemetry__ring) != c.ring_bytes {
telemetry_st.telemetry__ring = buffer(c.ring_bytes)
telemetry_st.telemetry__out = sb_new(c.ring_bytes + c.batch * 128 + 1024)
}
if telemetry_st.telemetry__off == null or len(telemetry_st.telemetry__off) != c.queue_max {
telemetry_st.telemetry__off = words(c.queue_max)
telemetry_st.telemetry__len = words(c.queue_max)
telemetry_st.telemetry__mk = words(c.queue_max)
}
tr_clear(telemetry_st)
}
function tr_clear(telemetry_st: mut TelemetryState) -> void {
telemetry_st.telemetry__head = 0
telemetry_st.telemetry__n = 0
telemetry_st.telemetry__wr = 0
}
function tr_slot(telemetry_st: TelemetryState, i: int) -> int { return (telemetry_st.telemetry__head + i) % len(telemetry_st.telemetry__off) }
# the oldest k lines gone (sent, or overwritten)
function tr_drop(telemetry_st: mut TelemetryState, k: int) -> void {
let n = Math.min(k, telemetry_st.telemetry__n)
telemetry_st.telemetry__head = tr_slot(telemetry_st, n)
telemetry_st.telemetry__n -= n
if telemetry_st.telemetry__n == 0 { tr_clear(telemetry_st) }
telemetry_st.telemetry__dirty = true
}
# a line (the first `n` bytes of b) put at the ring's end, its mark `mk` bytes in (-1 none)
function tr_add(telemetry_st: mut TelemetryState, b: []byte, n: int, mk: int) -> void {
let ring = telemetry_st.telemetry__ring
if n <= 0 or n > len(ring) { return }
if telemetry_st.telemetry__n == len(telemetry_st.telemetry__off) { tr_drop(telemetry_st, 1) }
var at = telemetry_st.telemetry__wr
if at + n > len(ring) { at = 0 }
while telemetry_st.telemetry__n > 0 and tr_over(telemetry_st, at, n) { tr_drop(telemetry_st, 1) }
for i in 0 .. n { ring[at + i] = b[i] }
let s = tr_slot(telemetry_st, telemetry_st.telemetry__n)
telemetry_st.telemetry__off[s] = at
telemetry_st.telemetry__len[s] = n
telemetry_st.telemetry__mk[s] = mk
telemetry_st.telemetry__n += 1
telemetry_st.telemetry__wr = at + n
telemetry_st.telemetry__dirty = true
}
# would bytes at .. at + n overwrite the oldest line
function tr_over(telemetry_st: TelemetryState, at: int, n: int) -> bool {
let s = telemetry_st.telemetry__head
let o = telemetry_st.telemetry__off[s]
return o < at + n and at < o + telemetry_st.telemetry__len[s]
}
# line i (0 the oldest) onto sb, its mark written as `id` when id is not null
function tr_line_to(telemetry_st: TelemetryState, i: int, sb: StrBuf, id: string) -> void {
let s = tr_slot(telemetry_st, i)
let o = telemetry_st.telemetry__off[s]
let n = telemetry_st.telemetry__len[s]
let mk = telemetry_st.telemetry__mk[s]
let ring = telemetry_st.telemetry__ring
var k = 0
while k < n {
if id != null and k == mk {
sb_add(sb, id)
k += len(TELEMETRY_MARK)
} else {
sb_byte(sb, ring[o + k])
k += 1
}
}
}

View file

@ -5,7 +5,7 @@ export function telemetry_tick(telemetry_st: mut TelemetryState) -> void {
if not telemetry_st.telemetry__ready { return } if not telemetry_st.telemetry__ready { return }
telemetry__id_tick(telemetry_st) telemetry__id_tick(telemetry_st)
if not telemetry__enabled() { if not telemetry__enabled() {
if len(telemetry_st.telemetry__q) > 0 or telemetry_st.telemetry__h >= 0 { telemetry_forget(telemetry_st) } if telemetry_st.telemetry__n > 0 or telemetry_st.telemetry__h >= 0 { telemetry_forget(telemetry_st) }
return return
} }
if not telemetry__can_send(telemetry_st) { return } if not telemetry__can_send(telemetry_st) { return }
@ -24,8 +24,8 @@ export function telemetry_tick(telemetry_st: mut TelemetryState) -> void {
telemetry__answer(telemetry_st, now) telemetry__answer(telemetry_st, now)
return return
} }
if len(telemetry_st.telemetry__q) == 0 { return } if telemetry_st.telemetry__n == 0 { return }
let early = telemetry_st.telemetry__fails == 0 and len(telemetry_st.telemetry__q) >= c.flush_at let early = telemetry_st.telemetry__fails == 0 and telemetry_st.telemetry__n >= c.flush_at
if early or now >= telemetry_st.telemetry__next_ms { if early or now >= telemetry_st.telemetry__next_ms {
telemetry_save(telemetry_st) telemetry_save(telemetry_st)
telemetry__flush(telemetry_st, now) telemetry__flush(telemetry_st, now)
@ -38,7 +38,7 @@ function telemetry__answer(telemetry_st: mut TelemetryState, now: int) -> void {
if st < 0 { return } if st < 0 { return }
telemetry_st.telemetry__h = -1 telemetry_st.telemetry__h = -1
if st > 0 { if st > 0 {
telemetry__queue_drop(telemetry_st, telemetry_st.telemetry__sent_n) tr_drop(telemetry_st, telemetry_st.telemetry__sent_n)
telemetry_st.telemetry__fails = 0 telemetry_st.telemetry__fails = 0
telemetry_st.telemetry__next_ms = now + telemetry__conf(telemetry_st).flush_ms telemetry_st.telemetry__next_ms = now + telemetry__conf(telemetry_st).flush_ms
telemetry__fact(telemetry_st, TELEMETRY_SENT, telemetry_st.telemetry__sent_n, st, 0) telemetry__fact(telemetry_st, TELEMETRY_SENT, telemetry_st.telemetry__sent_n, st, 0)
@ -50,9 +50,10 @@ function telemetry__answer(telemetry_st: mut TelemetryState, now: int) -> void {
} }
function telemetry__flush(telemetry_st: mut TelemetryState, now: int) -> void { function telemetry__flush(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 } if telemetry_st.telemetry__h >= 0 or len(telemetry_st.telemetry__id) == 0 or telemetry_st.telemetry__n == 0 { return }
let n = min(telemetry__conf(telemetry_st).batch, len(telemetry_st.telemetry__q)) let n = min(telemetry__conf(telemetry_st).batch, telemetry_st.telemetry__n)
let h = telemetry__send(telemetry_st, telemetry_batch_body(telemetry_st, n)) tb_batch(telemetry_st, n)
let h = telemetry__send(telemetry_st, telemetry_st.telemetry__out.b, sb_len(telemetry_st.telemetry__out))
if h < 0 { if h < 0 {
telemetry__backoff(telemetry_st, now) telemetry__backoff(telemetry_st, now)
telemetry__fact(telemetry_st, TELEMETRY_FAILED, n, 0, telemetry_st.telemetry__next_ms - now) telemetry__fact(telemetry_st, TELEMETRY_FAILED, n, 0, telemetry_st.telemetry__next_ms - now)

View file

@ -0,0 +1,35 @@
# stamp.ludic - the ISO time an event happened, written from the wall clock's seconds with the date
# worked out in whole numbers (days from the civil calendar), not a string built by DateTime
function tj_stamp(sb: StrBuf, secs: int) -> void {
var days = secs / 86400
var rem = secs % 86400
if rem < 0 {
rem += 86400
days -= 1
}
let z = days + 719468
var era = z / 146097
if z < 0 { era = (z - 146096) / 146097 }
let doe = z - era * 146097
let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146096) / 365
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100)
let mp = (5 * doy + 2) / 153
let d = doy - (153 * mp + 2) / 5 + 1
var m = mp + 3
if mp >= 10 { m = mp - 9 }
var y = yoe + era * 400
if m <= 2 { y += 1 }
sb_byte(sb, 34)
sb_int(sb, y)
sb_byte(sb, 45)
sb_int2(sb, m)
sb_byte(sb, 45)
sb_int2(sb, d)
sb_byte(sb, 84)
sb_int2(sb, rem / 3600)
sb_byte(sb, 58)
sb_int2(sb, (rem / 60) % 60)
sb_byte(sb, 58)
sb_int2(sb, rem % 60)
sb_add(sb, "Z\"")
}

View file

@ -12,7 +12,8 @@ export property TelemetryConfig {
flush_ms: int = 30000 # a batch this often flush_ms: int = 30000 # a batch this often
flush_at: int = 50 # or sooner, at this many events flush_at: int = 50 # or sooner, at this many events
batch: int = 200 # events in one request batch: int = 200 # events in one request
queue_max: int = 5000 # the oldest go past this queue_max: int = 5000 # the oldest go past this, or past ring_bytes of lines
ring_bytes: int = 1048576 # the lines' ring (ring.ludic), made once at the start
backoff_max: int = 6 # the wait doubles at most this many times backoff_max: int = 6 # the wait doubles at most this many times
save_ms: int = 5000 # the queue on disk at most this far behind save_ms: int = 5000 # the queue on disk at most this far behind
} }
@ -22,13 +23,13 @@ export port TelemetryWorld {
can_send: fn() -> bool = fn telemetry__yes # may this run send at all (a test, a dev build: no) can_send: fn() -> bool = fn telemetry__yes # may this run send at all (a test, a dev build: no)
enabled: fn() -> bool = fn telemetry__yes # the player's switch; off forgets everything queued enabled: fn() -> bool = fn telemetry__yes # the player's switch; off forgets everything queued
now_ms: fn() -> int = fn telemetry__wall_ms # a clock in milliseconds now_ms: fn() -> int = fn telemetry__wall_ms # a clock in milliseconds
stamp: fn() -> string = fn telemetry__wall_stamp # the ISO time an event happened clock_s: fn() -> int = fn telemetry__wall_s # the wall clock in seconds since 1970, for an event's time
} }
# how a batch travels (unbound: Http). send returns a handle (< 0: could not start); poll # how a batch travels (unbound: Http). send returns a handle (< 0: could not start); poll
# answers -1 while pending, 0 on failure, 1 when the server took it; drop abandons one. # answers -1 while pending, 0 on failure, 1 when the server took it; drop abandons one.
export port TelemetryTransport { export port TelemetryTransport {
send: fn(string, string) -> int = fn telemetry__http_send send: fn(string, []byte, int) -> int = fn telemetry__http_send # (url, body, its length)
poll: fn(int) -> int = fn telemetry__http_poll poll: fn(int) -> int = fn telemetry__http_poll
drop: fn(int) -> void = fn telemetry__http_drop drop: fn(int) -> void = fn telemetry__http_drop
} }
@ -56,8 +57,21 @@ export state TelemetryState {
telemetry__ready: bool = false telemetry__ready: bool = false
telemetry__id: string = "" telemetry__id: string = ""
telemetry__session: string = "" telemetry__session: string = ""
telemetry__q: []string = null # the queued lines, oldest first telemetry__ring: []byte = null # the queued lines (ring.ludic), oldest first from head
telemetry__dirty: bool = false # telemetry__q differs from the file telemetry__off: []int = null
telemetry__len: []int = null
telemetry__mk: []int = null # where each line's player mark is, -1 none
telemetry__head: int = 0
telemetry__n: int = 0
telemetry__wr: int = 0
telemetry__line: StrBuf = sb_new(TELEMETRY_LINE) # an event's line, written here
telemetry__props: TelemetryProps = telemetry__props_new()
telemetry__out: StrBuf = null # a batch or the file, written here
telemetry__url: string = "" # host and path, joined once at the config
telemetry__fpool: []TelemetryFact = new []TelemetryFact # fact records, reused once no drain holds one
telemetry__ffree: []TelemetryFact = new []TelemetryFact
telemetry__fnext: []int = words(1)
telemetry__dirty: bool = false # the ring differs from the file
telemetry__h: int = -1 # the batch in flight telemetry__h: int = -1 # the batch in flight
telemetry__sent_n: int = 0 # how many lines it carries telemetry__sent_n: int = 0 # how many lines it carries
telemetry__next_ms: int = 0 # when the next send may go telemetry__next_ms: int = 0 # when the next send may go
@ -71,8 +85,8 @@ export function telemetry_facts(telemetry_st: TelemetryState) -> Queue<Telemetry
return telemetry_st.telemetry__facts return telemetry_st.telemetry__facts
} }
function telemetry__fact(telemetry_st: TelemetryState, what: int, count: int, status: int, retry: int) -> void { function telemetry__fact(telemetry_st: mut TelemetryState, what: int, count: int, status: int, retry: int) -> void {
let f = new TelemetryFact let f = telemetry__fact_record(telemetry_st)
f.what = what f.what = what
f.count = count f.count = count
f.status = status f.status = status
@ -80,3 +94,8 @@ function telemetry__fact(telemetry_st: TelemetryState, what: int, count: int, st
q_push(telemetry_facts(telemetry_st), f) q_push(telemetry_facts(telemetry_st), f)
} }
function telemetry__facts__new() -> Queue<TelemetryFact> { return queue_new("telemetry.facts") } function telemetry__facts__new() -> Queue<TelemetryFact> { return queue_new("telemetry.facts") }
function telemetry__props_new() -> TelemetryProps {
let p = new TelemetryProps
p.sb = sb_new(TELEMETRY_LINE)
return p
}

View file

@ -0,0 +1,93 @@
# ring_test.ludic - an event written once into a kept line and the queue a ring of bytes: ten thousand
# events, most past the cap in lines and in bytes, make nothing; the line is exact JSON with its time
# worked out from the seconds; the batch puts the id in; the file round-trips with its marks
import "ludic.telemetry"
import "ludic.base"
program RingTest {
numbers float
state RingTestState {
now: int = 1790337600
body: string = ""
}
function fk_now() -> int { return 0 }
function fk_on() -> bool { return true }
function fk_clock(ring_test_st: RingTestState) -> int { return ring_test_st.now }
function fk_send(url: string, body: []byte, n: int) -> int { return -1 }
function fk_poll(h: int) -> int { return -1 }
function fk_drop(h: int) -> void { }
bind TelemetryWorld { can_send: fn fk_on, enabled: fn fk_on, now_ms: fn fk_now, clock_s: fn fk_clock }
bind TelemetryTransport { send: fn fk_send, poll: fn fk_poll, drop: fn fk_drop }
function qfile() -> string { return Os.temp_dir() + "/ludic-telemetry-ring-queue.txt" }
function start(telemetry_st: mut TelemetryState, max: int, bytes: int) -> void {
telemetry_reset(telemetry_st)
if Fs.exists(qfile()) { Fs.remove(qfile()) }
Fs.write_text(Os.temp_dir() + "/ludic-telemetry-ring-id.json", "{\"id\":\"p-9\"}")
let c = new TelemetryConfig
c.host = "http://fake"
c.key = "k-9"
c.lib = "ring"
c.queue_file = qfile()
c.id_file = Os.temp_dir() + "/ludic-telemetry-ring-id.json"
c.queue_max = max
c.ring_bytes = bytes
c.batch = 3
telemetry_config(telemetry_st, c)
telemetry_start(telemetry_st)
}
function one(telemetry_st: mut TelemetryState, k: int) -> void {
let p = telemetry_props(telemetry_st)
telemetry_int(p, "n", k)
telemetry_str(p, "what", "a fish")
telemetry_bool(p, "dev", true)
telemetry_event(telemetry_st, "caught", p)
}
test "ten thousand events, past the cap in lines and in bytes, make nothing" (telemetry_st: mut TelemetryState) {
start(telemetry_st, 5000, 262144)
for k in 0 .. 200 { one(telemetry_st, k) }
let heap = Os.heap_bytes()
for k in 0 .. 10000 { one(telemetry_st, k) }
let grew = Os.heap_bytes() - heap
print(`ring: 10000 events, {telemetry_queued(telemetry_st)} held, heap {grew} bytes`)
expect(grew == long(0))
expect(Text.contains(telemetry_line(telemetry_st, telemetry_queued(telemetry_st) - 1), "\"n\":9999"))
start(telemetry_st, 40, 1048576)
let heap2 = Os.heap_bytes()
for k in 0 .. 10000 { one(telemetry_st, k) }
expect(Os.heap_bytes() - heap2 == long(0))
expect_eq(telemetry_queued(telemetry_st), 40)
expect(Text.contains(telemetry_line(telemetry_st, 0), "\"n\":9960"))
}
test "a line is exact JSON, escaped, its time from the seconds" (telemetry_st: mut TelemetryState, ring_test_st: mut RingTestState) {
start(telemetry_st, 50, 65536)
let p = telemetry_props(telemetry_st)
telemetry_str(p, "say", "a \"quote\"\\ and\na line")
telemetry_event(telemetry_st, "said", p)
ring_test_st.now = 951868799
telemetry_event(telemetry_st, "leap", null)
let a = telemetry_line(telemetry_st, 0)
expect(Text.contains(a, "\"properties\":{\"say\":\"a \\\"quote\\\"\\\\ and\\u000aa line\",\"$session_id\":"))
expect(Text.contains(a, "\"timestamp\":\"2026-09-25T12:00:00Z\"}"))
expect(Json.parse(a) != null)
expect(Text.contains(telemetry_line(telemetry_st, 1), "\"timestamp\":\"2000-02-29T23:59:59Z\""))
expect(Text.contains(telemetry_line(telemetry_st, 1), "\"properties\":{\"$session_id\""))
}
test "the batch puts the id in, and the file round-trips with its marks" (telemetry_st: mut TelemetryState) {
start(telemetry_st, 50, 65536)
for k in 0 .. 4 { one(telemetry_st, k) }
let body = telemetry_batch_body(telemetry_st, 3)
expect(Text.starts_with(body, "{\"api_key\":\"k-9\",\"batch\":[{"))
expect(Text.contains(body, "\"distinct_id\":\"p-9\""))
expect(not Text.contains(body, "\"n\":3"))
expect(Json.parse(body) != null)
telemetry_save(telemetry_st)
telemetry_reset(telemetry_st)
telemetry_start(telemetry_st)
expect_eq(telemetry_queued(telemetry_st), 4)
expect(Text.contains(telemetry_batch_body(telemetry_st, 4), "\"n\":3,"))
expect(not Text.contains(telemetry_batch_body(telemetry_st, 4), "TELEMETRY_PLAYER"))
}
}

View file

@ -18,18 +18,18 @@ program TelemetryTest {
function fake_now(telemetry_test_st: TelemetryTestState) -> int { return telemetry_test_st.clock } function fake_now(telemetry_test_st: TelemetryTestState) -> int { return telemetry_test_st.clock }
function fake_enabled(telemetry_test_st: TelemetryTestState) -> bool { return telemetry_test_st.on } function fake_enabled(telemetry_test_st: TelemetryTestState) -> bool { return telemetry_test_st.on }
function fake_can(telemetry_test_st: TelemetryTestState) -> bool { return telemetry_test_st.allowed } function fake_can(telemetry_test_st: TelemetryTestState) -> bool { return telemetry_test_st.allowed }
function fake_stamp() -> string { return "2026-09-25T12:00:00Z" } function fake_clock() -> int { return 1790337600 }
function fake_send(telemetry_test_st: mut TelemetryTestState, url: string, body: string) -> int { function fake_send(telemetry_test_st: mut TelemetryTestState, url: string, body: []byte, n: int) -> int {
if telemetry_test_st.refuse { return -1 } if telemetry_test_st.refuse { return -1 }
telemetry_test_st.sends += 1 telemetry_test_st.sends += 1
telemetry_test_st.last_url = url telemetry_test_st.last_url = url
telemetry_test_st.last_body = body telemetry_test_st.last_body = text_of(body, n)
return 7 return 7
} }
function fake_poll(telemetry_test_st: TelemetryTestState, h: int) -> int { return telemetry_test_st.answer } function fake_poll(telemetry_test_st: TelemetryTestState, h: int) -> int { return telemetry_test_st.answer }
function fake_drop(h: int) -> void { } function fake_drop(h: int) -> void { }
bind TelemetryWorld { can_send: fn fake_can, enabled: fn fake_enabled, now_ms: fn fake_now, stamp: fn fake_stamp } bind TelemetryWorld { can_send: fn fake_can, enabled: fn fake_enabled, now_ms: fn fake_now, clock_s: fn fake_clock }
bind TelemetryTransport { send: fn fake_send, poll: fn fake_poll, drop: fn fake_drop } bind TelemetryTransport { send: fn fake_send, poll: fn fake_poll, drop: fn fake_drop }
function qfile() -> string { return Os.temp_dir() + "/ludic-telemetry-test-queue.txt" } function qfile() -> string { return Os.temp_dir() + "/ludic-telemetry-test-queue.txt" }