diff --git a/packages/ludic.telemetry/README.md b/packages/ludic.telemetry/README.md index 4c36483c..55320059 100644 --- a/packages/ludic.telemetry/README.md +++ b/packages/ludic.telemetry/README.md @@ -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) | diff --git a/packages/ludic.telemetry/facts.ludic b/packages/ludic.telemetry/facts.ludic new file mode 100644 index 00000000..a56dfe01 --- /dev/null +++ b/packages/ludic.telemetry/facts.ludic @@ -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 +} diff --git a/packages/ludic.telemetry/index.ludic b/packages/ludic.telemetry/index.ludic index b9b4956e..8539fd02 100644 --- a/packages/ludic.telemetry/index.ludic +++ b/packages/ludic.telemetry/index.ludic @@ -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" diff --git a/packages/ludic.telemetry/load.ludic b/packages/ludic.telemetry/load.ludic new file mode 100644 index 00000000..bb1415ee --- /dev/null +++ b/packages/ludic.telemetry/load.ludic @@ -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 +} diff --git a/packages/ludic.telemetry/ports.ludic b/packages/ludic.telemetry/ports.ludic index e8d2cb6c..96c80436 100644 --- a/packages/ludic.telemetry/ports.ludic +++ b/packages/ludic.telemetry/ports.ludic @@ -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) diff --git a/packages/ludic.telemetry/props.ludic b/packages/ludic.telemetry/props.ludic new file mode 100644 index 00000000..19c4599e --- /dev/null +++ b/packages/ludic.telemetry/props.ludic @@ -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 } diff --git a/packages/ludic.telemetry/props_nest.ludic b/packages/ludic.telemetry/props_nest.ludic new file mode 100644 index 00000000..bb539a4b --- /dev/null +++ b/packages/ludic.telemetry/props_nest.ludic @@ -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 +} diff --git a/packages/ludic.telemetry/queries.ludic b/packages/ludic.telemetry/queries.ludic index 641ba8f9..46e9aa4b 100644 --- a/packages/ludic.telemetry/queries.ludic +++ b/packages/ludic.telemetry/queries.ludic @@ -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 diff --git a/packages/ludic.telemetry/queue.ludic b/packages/ludic.telemetry/queue.ludic index 893720d7..56a297c5 100644 --- a/packages/ludic.telemetry/queue.ludic +++ b/packages/ludic.telemetry/queue.ludic @@ -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)) } diff --git a/packages/ludic.telemetry/ring.ludic b/packages/ludic.telemetry/ring.ludic new file mode 100644 index 00000000..809f3c07 --- /dev/null +++ b/packages/ludic.telemetry/ring.ludic @@ -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 + } + } +} diff --git a/packages/ludic.telemetry/send.ludic b/packages/ludic.telemetry/send.ludic index 654c574a..ab4dc84a 100644 --- a/packages/ludic.telemetry/send.ludic +++ b/packages/ludic.telemetry/send.ludic @@ -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) diff --git a/packages/ludic.telemetry/stamp.ludic b/packages/ludic.telemetry/stamp.ludic new file mode 100644 index 00000000..6702347b --- /dev/null +++ b/packages/ludic.telemetry/stamp.ludic @@ -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\"") +} diff --git a/packages/ludic.telemetry/state.ludic b/packages/ludic.telemetry/state.ludic index d8280f28..2c3ee351 100644 --- a/packages/ludic.telemetry/state.ludic +++ b/packages/ludic.telemetry/state.ludic @@ -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 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 { 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 +} diff --git a/packages/ludic.telemetry/tests/ring_test.ludic b/packages/ludic.telemetry/tests/ring_test.ludic new file mode 100644 index 00000000..ed079c1a --- /dev/null +++ b/packages/ludic.telemetry/tests/ring_test.ludic @@ -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) + } +} diff --git a/packages/ludic.telemetry/tests/telemetry_test.ludic b/packages/ludic.telemetry/tests/telemetry_test.ludic index 3c6edaa0..03879152 100644 --- a/packages/ludic.telemetry/tests/telemetry_test.ludic +++ b/packages/ludic.telemetry/tests/telemetry_test.ludic @@ -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" }