Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 23 additions & 6 deletions guide/lua.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,15 @@ Each channel runs the script independently. It executes top-level initialization

## Events and payloads

A Source event has `type = "source"`, a timestamp and `payload`. A timer event has `type = "timer"`, a timestamp and no payload. Events and nested payload values are deeply read-only; modifying them is a sandbox violation.
| Field | Source event | Timer event |
| --- | --- | --- |
| `type` | `"source"` | `"timer"` |
| `timestamp` | Dispatch time, in integer milliseconds since Pipeline startup | Dispatch time, on the same timeline |
| `payload` | Read-only Source payload | Absent (`nil`) |
| `id` | Absent (`nil`) | Timer name (string), or `nil` for the anonymous timer |
| `eligibleAt` | Absent (`nil`) | Scheduled deadline, in integer milliseconds since Pipeline startup |

Events and nested payload values are deeply read-only; modifying them is a sandbox violation.

`event.timestamp` is a monotonic integer millisecond count since Pipeline startup, sampled when the channel dispatches the event. It is not Unix time. Use the wall-clock functions below for a real-world date or timestamp.

Expand Down Expand Up @@ -87,19 +95,28 @@ Registry lookup, Builder receiver/argument/type/range errors and timer argument

## Timers and state

Each VM has one one-shot timer slot:
Each VM has an anonymous one-shot timer slot plus independently named one-shot timers:

```lua
setTimeout(durationMs)
clearTimerTask()
setTimeout(durationMs, id)
clearTimeout()
clearTimeout(id)
hasTimeout()
hasTimeout(id)
```

The delay must be a nonnegative Lua integer. It starts at the call, on a monotonic timeline. Setting again replaces the pending timer; clearing an empty slot succeeds. `hasTimeout()` observes the current pending slot immediately. The slot is removed before its timer event enters `main`, so it reports false during that event unless another timer was scheduled.
The delay must be a nonnegative Lua integer. It starts at the call, on a monotonic timeline. Setting again replaces the pending timer with the same id; different ids remain independent. `clearTimeout` is idempotent, and `hasTimeout` observes the selected timer immediately. A timer is removed before its event enters `main`, so it reports false during that event unless another timer was scheduled. The id must be a non-empty UTF-8 string of at most 128 bytes. The maximum accepted integer is `9223372036854775807`; if conversion or deadline arithmetic cannot represent it, the catchable error is `setTimeout delay is out of range`.

These functions are allowed at top level and in `main`. A zero delay schedules eligibility on the next event-loop iteration. Execution may be later because the channel is busy. Periodic behavior explicitly schedules the next one-shot timer from the current timer event. Negative, non-integer, string, or missing delays raise `setTimeout delay must be a non-negative integer`; a delay that cannot be represented by the monotonic deadline or `eligibleAt` calculation raises `setTimeout delay is out of range`.

Omitting `id` or passing `nil` selects the anonymous timer. Invalid ids raise the catchable error `timer id must be a non-empty string of at most 128 bytes`.

`event.eligibleAt` is the timer's deadline in integer milliseconds since Pipeline startup, on the same monotonic timeline as `event.timestamp`. The difference `event.timestamp - event.eligibleAt` is its dispatch delay. Timers with the same deadline run in registration order. Among due timers and readable Source input, the earlier deadline or first observed Source readiness runs first; ties use their registration or observation order. Source readiness keeps its place until one record is consumed, so a timer that repeatedly reschedules itself with zero delay cannot indefinitely prevent readable Source input from running. Sink backpressure and lifecycle work can still delay either kind of event.

These functions are allowed at top level and in `main`. A zero delay schedules eligibility on the next event-loop iteration. Execution may be later because the channel is busy. Periodic behavior explicitly schedules the next one-shot timer from the current timer event. Invalid delays raise `setTimeout delay must be a non-negative integer`.
Named timers share the VM's batch completion behavior: the first successful `emit` in an event can complete all pending Source records in that VM. A timer id does not provide independent reliable completion for records with that key.

Lua state, Builders, snapshots and timers count toward the VM's memory allowance. Top-level initialization and each `main` invocation are bounded by the Runner's Lua CPU limit. State lives only in the current VM. A script/resource/sandbox failure invalidates the VM and its timers; rebuilding executes top-level code again. Relevant Document updates and process restarts also recreate state. Script replacement can briefly hold both old and new VMs, each with its own configured allowance.
Lua state, Builders, snapshots and timers count toward the VM's memory allowance. Each pending timer is charged 2 KiB plus three times its id length in UTF-8 bytes (zero id bytes for the anonymous timer). The fixed charge includes an allowance for the timer and its indexes; it is an estimate, not an exact allocation or process-memory measurement. There is no separate timer-count limit. Top-level initialization and each `main` invocation are bounded by the Runner's Lua CPU limit. State lives only in the current VM. A script/resource/sandbox failure invalidates the VM and its timers; rebuilding executes top-level code again. Relevant Document updates and process restarts also recreate state. Script replacement can briefly hold both old and new VMs, each with its own configured allowance.

## Wall-clock time and dates

Expand Down
10 changes: 9 additions & 1 deletion src/lua/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ pub(super) enum ProcessEvent {
},
Timer {
timestamp: KernelTimestampMillis,
id: Option<Box<str>>,
eligible_at: KernelTimestampMillis,
},
}

Expand Down Expand Up @@ -167,9 +169,15 @@ pub(super) fn project(
project_source_message(lua, payload, Rc::clone(&fatal_fault))?,
)?;
}
ProcessEvent::Timer { timestamp } => {
ProcessEvent::Timer {
timestamp,
id,
eligible_at,
} => {
backing.raw_set("type", "timer")?;
backing.raw_set("timestamp", timestamp.0)?;
backing.raw_set("id", id.as_deref())?;
backing.raw_set("eligibleAt", eligible_at.0)?;
}
}
install_readonly_backing(lua, backing, fatal_fault)
Expand Down
123 changes: 123 additions & 0 deletions src/lua/memory.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

//! Shared native-memory budget for one thread-affine Lua VM.
//!
//! Sandbox initialization creates one budget and shares it with Payload and
//! timer owners. Each owner supplies its memory charge and releases it when
//! its storage is dropped. The budget reduces the Lua allocator's allowance
//! so native charges and Lua allocations share the same VM limit.
//!
//! The Lua reference is weak: reservations may outlive the interpreter during
//! teardown without keeping it alive. Failed reservations leave both the
//! charge and allocator allowance unchanged. This module owns no payload or
//! timer storage and does not define their accounting estimates.

use super::{LuaApiFailure, LuaApiResult, LuaVmFatalFault, record_fatal_fault};
use mlua::{Lua, WeakLua};
use std::cell::Cell;
use std::fmt;
use std::rc::Rc;

pub(super) struct LuaNativeMemoryBudget {
lua: WeakLua,
limit_bytes: usize,
used_bytes: Cell<usize>,
fatal_fault: Rc<Cell<Option<LuaVmFatalFault>>>,
}

impl LuaNativeMemoryBudget {
pub(super) fn new(
lua: &Lua,
limit_bytes: usize,
fatal_fault: Rc<Cell<Option<LuaVmFatalFault>>>,
) -> Self {
Self {
lua: lua.weak(),
limit_bytes,
used_bytes: Cell::new(0),
fatal_fault,
}
}

pub(super) fn used_bytes(&self) -> usize {
self.used_bytes.get()
}

pub(super) fn replace(&self, previous_bytes: usize, next_bytes: usize) -> LuaApiResult<()> {
let used_bytes = self.used_bytes.get();
let retained_bytes = used_bytes
.checked_sub(previous_bytes)
.ok_or(LuaApiFailure::InternalInvariantViolation)?;
let next_used_bytes = retained_bytes
.checked_add(next_bytes)
.ok_or(LuaApiFailure::MemoryExceeded)?;
let lua = self
.lua
.try_upgrade()
.ok_or(LuaApiFailure::InternalInvariantViolation)?;
let next_lua_limit = self
.limit_bytes
.checked_sub(next_used_bytes)
.ok_or(LuaApiFailure::MemoryExceeded)?;
if lua.used_memory() > next_lua_limit {
return Err(LuaApiFailure::MemoryExceeded);
}
lua.set_memory_limit(next_lua_limit)
.map_err(|_| LuaApiFailure::InternalInvariantViolation)?;
self.used_bytes.set(next_used_bytes);
Ok(())
}

pub(super) fn release(&self, bytes: usize) {
let Some(next_used_bytes) = self.used_bytes.get().checked_sub(bytes) else {
record_fatal_fault(
&self.fatal_fault,
LuaVmFatalFault::InternalInvariantViolation,
);
return;
};
self.used_bytes.set(next_used_bytes);
let Some(lua) = self.lua.try_upgrade() else {
return;
};
let Some(next_lua_limit) = self.limit_bytes.checked_sub(next_used_bytes) else {
record_fatal_fault(
&self.fatal_fault,
LuaVmFatalFault::InternalInvariantViolation,
);
return;
};
if lua.set_memory_limit(next_lua_limit).is_err() {
record_fatal_fault(
&self.fatal_fault,
LuaVmFatalFault::InternalInvariantViolation,
);
}
}
}
impl fmt::Debug for LuaNativeMemoryBudget {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("LuaNativeMemoryBudget")
.field("limit_bytes", &self.limit_bytes)
.field("used_bytes", &self.used_bytes.get())
.finish_non_exhaustive()
}
}
Loading