Merge branch 'lang/telemetry-ring' into lang/leaks2

This commit is contained in:
Orkun ÇAKILKAYA 2026-09-28 15:42:59 +03:00
commit 1cf132cb87
15 changed files with 543 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.
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`)
```ludic
property TelemetryConfig {
host, key, path ("/batch/"), lib # an empty host or key sends nothing, ever
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 {
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)
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
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
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_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_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) |

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,14 @@ module ludic_telemetry uses ludic_base
numbers float
import "ludic.base"
import "state.ludic"
import "facts.ludic"
import "ports.ludic"
import "props.ludic"
import "props_nest.ludic"
import "stamp.ludic"
import "ring.ludic"
import "queue.ludic"
import "load.ludic"
import "send.ludic"
import "ident.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
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 {
return telemetry_st.telemetry__cfg
@ -7,7 +10,7 @@ function telemetry__conf(telemetry_st: TelemetryState) -> TelemetryConfig {
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_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 {
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__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)
if h < 0 { return h }
Http.set(h, "Content-Type", "application/json")
Http.body(h, body)
Http.body_bytes(h, data_of(body), n)
Http.send(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__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) }
# a request abandoned (the player said no)

View file

@ -0,0 +1,72 @@
# 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 at the top
lv: []int = null # inside an object or a list: how many written at each level
d: int = 0 # how deep
}
# 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
p.d = 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 {
tp_next(p)
tj_str(p.sb, k)
sb_byte(p.sb, 58)
}
# a comma before all but the first at this level, and the count kept
function tp_next(p: TelemetryProps) -> void {
if p.d == 0 {
if p.n > 0 { sb_byte(p.sb, 44) }
p.n += 1
return
}
if p.lv[p.d] > 0 { sb_byte(p.sb, 44) }
p.lv[p.d] += 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

@ -0,0 +1,43 @@
# props_nest.ludic - an object or a list inside an event's properties ($set, a photograph's subjects, the
# party's names), written into the same kept buffer; and a signature of what was written, to tell a
# change without keeping a copy
export function telemetry_obj(p: TelemetryProps, k: string) -> void { tn_open(p, k, 123) }
export function telemetry_list(p: TelemetryProps, k: string) -> void { tn_open(p, k, 91) }
export function telemetry_obj_end(p: TelemetryProps) -> void { tn_close(p, 125) }
export function telemetry_list_end(p: TelemetryProps) -> void { tn_close(p, 93) }
function tn_open(p: TelemetryProps, k: string, c: int) -> void {
tp_key(p, k)
sb_byte(p.sb, c)
if p.lv == null { p.lv = words(8) }
if p.d < len(p.lv) - 1 { p.d += 1 }
p.lv[p.d] = 0
}
function tn_close(p: TelemetryProps, c: int) -> void {
sb_byte(p.sb, c)
if p.d > 0 { p.d -= 1 }
}
# a string in the list being written
export function telemetry_item_str(p: TelemetryProps, s: string) -> void {
tp_next(p)
tj_str(p.sb, s)
}
# a string item made of numbers the caller writes itself (sb_int, sb_byte - no quote or backslash)
export function telemetry_item_open(p: TelemetryProps) -> StrBuf {
tp_next(p)
sb_byte(p.sb, 34)
return p.sb
}
export function telemetry_item_close(p: TelemetryProps) -> void { sb_byte(p.sb, 34) }
# the same for a named string property; closed with telemetry_item_close
export function telemetry_str_open(p: TelemetryProps, k: string) -> StrBuf {
tp_key(p, k)
sb_byte(p.sb, 34)
return p.sb
}
# a number standing for what has been written: equal when the bytes are (FNV-1a, 31 bits)
export function telemetry_props_sig(p: TelemetryProps) -> int {
var h = 2166136261 & 2147483647
for i in 0 .. sb_len(p.sb) { h = ((h ^ p.sb.b[i]) * 16777619) & 2147483647 }
return h
}

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_in_flight(telemetry_st: TelemetryState) -> bool { return telemetry_st.telemetry__h >= 0 }
export function telemetry_queued(telemetry_st: TelemetryState) -> int {
if telemetry_st.telemetry__q == null { return 0 }
return len(telemetry_st.telemetry__q)
}
export function telemetry_queued(telemetry_st: TelemetryState) -> int { return telemetry_st.telemetry__n }
# 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)
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__id = ""
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__sent_n = 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
# whatever an earlier start left unsent - or nothing, if the player has it off
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__t0 = Time.now()
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", "")
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) }
@ -15,87 +15,80 @@ export function telemetry_start(telemetry_st: mut TelemetryState) -> void {
# 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) }
# an event: the properties (null for none) with the session and the library added, and the player
# as a mark put in when the batch goes, so an event queued before the id is known is not misfiled
export function telemetry_event(telemetry_st: mut TelemetryState, name: string, props: Val) -> void {
# an event: its properties (null for none) with the session and the library added, the player as a
# mark put in when the batch goes, so an event queued before the id is known is not misfiled. The
# 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 }
var p = props
if p == null { p = Value.object() }
Value.put(p, "$session_id", Value.str(telemetry_st.telemetry__session))
if len(telemetry__conf(telemetry_st).lib) > 0 { Value.put(p, "$lib", Value.str(telemetry__conf(telemetry_st).lib)) }
let e = Value.object()
Value.put(e, "event", Value.str(name))
Value.put(e, "distinct_id", Value.str(TELEMETRY_MARK))
Value.put(e, "properties", p)
Value.put(e, "timestamp", Value.str(telemetry__stamp()))
push(telemetry_st.telemetry__q, Json.encode(e))
telemetry_st.telemetry__dirty = true
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) }
let sb = telemetry_st.telemetry__line
sb_clear(sb)
sb_add(sb, "{\"event\":")
tj_str(sb, name)
sb_add(sb, ",\"distinct_id\":\"")
let mk = sb_len(sb)
sb_add(sb, TELEMETRY_MARK)
sb_add(sb, "\",\"properties\":{")
if props != null and props.n > 0 {
for i in 0 .. sb_len(props.sb) { sb_byte(sb, props.sb.b[i]) }
sb_byte(sb, 44)
}
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 {
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
let f = telemetry__conf(telemetry_st).queue_file
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) }
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
export function telemetry_forget(telemetry_st: mut TelemetryState) -> void {
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
let f = telemetry__conf(telemetry_st).queue_file
if len(f) > 0 and Fs.exists(f) { Fs.remove(f) }
}
function telemetry__queue_load(telemetry_st: mut TelemetryState) -> void {
telemetry_st.telemetry__q = new []string
let f = telemetry__conf(telemetry_st).queue_file
if len(f) == 0 { return }
let s = Fs.read_text(f)
if s == null { return }
let parts = Text.split(s, "\n")
for i in 0 .. len(parts) { if len(parts[i]) > 2 { push(telemetry_st.telemetry__q, parts[i]) } }
let max = telemetry__conf(telemetry_st).queue_max
if len(telemetry_st.telemetry__q) > max { telemetry__queue_drop(telemetry_st, len(telemetry_st.telemetry__q) - max) }
}
# 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
# the batch: the first n lines joined with the player id put in, into the kept buffer (tb_batch), and
# as text for a test
function tb_batch(telemetry_st: TelemetryState, n: int) -> void {
let out = telemetry_st.telemetry__out
sb_clear(out)
sb_add(out, "{\"api_key\":")
tj_str(out, telemetry__conf(telemetry_st).key)
sb_add(out, ",\"batch\":[")
for i in 0 .. n {
if i > 0 { sb_byte(out, 44) }
tr_line_to(telemetry_st, i, out, telemetry_st.telemetry__id)
}
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 }
telemetry__id_tick(telemetry_st)
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
}
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)
return
}
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 telemetry_st.telemetry__n == 0 { return }
let early = telemetry_st.telemetry__fails == 0 and telemetry_st.telemetry__n >= c.flush_at
if early or now >= telemetry_st.telemetry__next_ms {
telemetry_save(telemetry_st)
telemetry__flush(telemetry_st, now)
@ -38,7 +38,7 @@ function telemetry__answer(telemetry_st: mut TelemetryState, now: int) -> void {
if st < 0 { return }
telemetry_st.telemetry__h = -1
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__next_ms = now + telemetry__conf(telemetry_st).flush_ms
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 {
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 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, telemetry_st.telemetry__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 {
telemetry__backoff(telemetry_st, 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_at: int = 50 # or sooner, at this many events
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
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)
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
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
# answers -1 while pending, 0 on failure, 1 when the server took it; drop abandons one.
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
drop: fn(int) -> void = fn telemetry__http_drop
}
@ -56,8 +57,21 @@ export state TelemetryState {
telemetry__ready: bool = false
telemetry__id: string = ""
telemetry__session: string = ""
telemetry__q: []string = null # the queued lines, oldest first
telemetry__dirty: bool = false # telemetry__q differs from the file
telemetry__ring: []byte = null # the queued lines (ring.ludic), oldest first from head
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__sent_n: int = 0 # how many lines it carries
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
}
function telemetry__fact(telemetry_st: TelemetryState, what: int, count: int, status: int, retry: int) -> void {
let f = new TelemetryFact
function telemetry__fact(telemetry_st: mut TelemetryState, what: int, count: int, status: int, retry: int) -> void {
let f = telemetry__fact_record(telemetry_st)
f.what = what
f.count = count
f.status = status
@ -80,3 +94,9 @@ function telemetry__fact(telemetry_st: TelemetryState, what: int, count: int, st
q_push(telemetry_facts(telemetry_st), f)
}
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)
p.lv = words(8)
return p
}

View file

@ -0,0 +1,125 @@
# 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"))
}
test "an object and a list inside the properties are JSON, and the signature tells a change" (telemetry_st: mut TelemetryState) {
start(telemetry_st, 50, 65536)
let p = telemetry_props(telemetry_st)
telemetry_int(p, "grade", 3)
telemetry_list(p, "subjects")
telemetry_item_str(p, "elk")
telemetry_item_str(p, "deer")
telemetry_list_end(p)
telemetry_obj(p, "$set")
telemetry_str(p, "arch", "arm64")
telemetry_list(p, "looks")
let sb = telemetry_item_open(p)
sb_int(sb, 4)
sb_byte(sb, 44)
sb_int(sb, 2)
telemetry_item_close(p)
telemetry_list_end(p)
telemetry_obj_end(p)
let lk = telemetry_str_open(p, "look")
sb_int(lk, 7)
telemetry_item_close(p)
telemetry_bool(p, "last", false)
let sig = telemetry_props_sig(p)
telemetry_event(telemetry_st, "shot", p)
let a = telemetry_line(telemetry_st, 0)
expect(Text.contains(a, "{\"grade\":3,\"subjects\":[\"elk\",\"deer\"],\"$set\":{\"arch\":\"arm64\",\"looks\":[\"4,2\"]},\"look\":\"7\",\"last\":false,"))
expect(Json.parse(a) != null)
let q = telemetry_props(telemetry_st)
telemetry_int(q, "grade", 3)
expect(telemetry_props_sig(q) != sig)
}
}

View file

@ -18,18 +18,18 @@ program TelemetryTest {
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_can(telemetry_test_st: TelemetryTestState) -> bool { return telemetry_test_st.allowed }
function fake_stamp() -> string { return "2026-09-25T12:00:00Z" }
function fake_send(telemetry_test_st: mut TelemetryTestState, url: string, body: string) -> int {
function fake_clock() -> int { return 1790337600 }
function fake_send(telemetry_test_st: mut TelemetryTestState, url: string, body: []byte, n: int) -> int {
if telemetry_test_st.refuse { return -1 }
telemetry_test_st.sends += 1
telemetry_test_st.last_url = url
telemetry_test_st.last_body = body
telemetry_test_st.last_body = text_of(body, n)
return 7
}
function fake_poll(telemetry_test_st: TelemetryTestState, h: int) -> int { return telemetry_test_st.answer }
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 }
function qfile() -> string { return Os.temp_dir() + "/ludic-telemetry-test-queue.txt" }