///|
/// Lifecycle state for an application message held by the local store.
pub(all) enum AppStoreState {
Pending
Delivered
Acknowledged
DeadLetter
} derive(Eq, Debug)
///|
/// An immutable record containing message lifecycle data.
pub(all) struct AppStoreRecord {
message : AppMessage
state : AppStoreState
retries : Int
} derive(Eq, Debug)
///|
/// A bounded in-memory message store for demos, tests, and embedded workers.
pub(all) struct AppMessageStore {
records : Array[AppStoreRecord]
dead_letters : Array[AppMessage]
capacity : Int
max_retries : Int
} derive(Debug)
///|
/// Metrics for observing a store without exposing its internal arrays.
pub(all) struct AppStoreStats {
total : Int
pending : Int
delivered : Int
acknowledged : Int
dead_letters : Int
retries : Int
} derive(Eq, Debug)
///|
/// Create an empty bounded store.
pub fn app_store_new(capacity : Int, max_retries : Int) -> AppMessageStore {
{
records: [],
dead_letters: [],
capacity: if capacity < 1 {
1
} else {
capacity
},
max_retries: if max_retries < 1 {
1
} else {
max_retries
},
}
}
///|
fn store_copy_records(store : AppMessageStore) -> Array[AppStoreRecord] {
let records = []
for record in store.records {
records.push(record)
}
records
}
///|
fn store_copy_dead_letters(store : AppMessageStore) -> Array[AppMessage] {
let messages = []
for message in store.dead_letters {
messages.push(message)
}
messages
}
///|
fn store_record_index(store : AppMessageStore, id : String) -> Int? {
for i = 0; i < store.records.length(); i = i + 1 {
if store.records[i].message.id == id {
return Some(i)
}
}
None
}
///|
/// Insert a message; the Boolean is false for duplicate IDs or a full store.
pub fn app_store_append(
store : AppMessageStore,
message : AppMessage,
) -> (AppMessageStore, Bool) {
if store_record_index(store, message.id) is Some(_) ||
store.records.length() >= store.capacity {
return (store, false)
}
let records = store_copy_records(store)
records.push({ message, state: Pending, retries: 0 })
(
{
records,
dead_letters: store_copy_dead_letters(store),
capacity: store.capacity,
max_retries: store.max_retries,
},
true,
)
}
///|
/// Insert messages in order and return the accepted count.
pub fn app_store_append_batch(
store : AppMessageStore,
messages : Array[AppMessage],
) -> (AppMessageStore, Int) {
let mut current = store
let mut accepted = 0
for message in messages {
let (next, ok) = app_store_append(current, message)
current = next
if ok {
accepted = accepted + 1
}
}
(current, accepted)
}
///|
/// Find a message by ID, including messages awaiting acknowledgement.
pub fn app_store_get(store : AppMessageStore, id : String) -> AppMessage? {
match store_record_index(store, id) {
Some(index) => Some(store.records[index].message)
None => None
}
}
///|
/// Return a record with lifecycle metadata.
pub fn app_store_record(
store : AppMessageStore,
id : String,
) -> AppStoreRecord? {
match store_record_index(store, id) {
Some(index) => Some(store.records[index])
None => None
}
}
///|
fn store_update_state(
store : AppMessageStore,
id : String,
state : AppStoreState,
) -> AppMessageStore {
let records = store_copy_records(store)
match store_record_index(store, id) {
Some(index) =>
records[index] = {
message: records[index].message,
state,
retries: records[index].retries,
}
None => ()
}
{
records,
dead_letters: store_copy_dead_letters(store),
capacity: store.capacity,
max_retries: store.max_retries,
}
}
///|
/// Mark a message as delivered to a worker.
pub fn app_store_mark_delivered(
store : AppMessageStore,
id : String,
) -> AppMessageStore {
store_update_state(store, id, Delivered)
}
///|
/// Mark a message as successfully processed.
pub fn app_store_ack(store : AppMessageStore, id : String) -> AppMessageStore {
store_update_state(store, id, Acknowledged)
}
///|
/// Return whether a message has reached a terminal success state.
pub fn app_store_is_acknowledged(store : AppMessageStore, id : String) -> Bool {
match app_store_record(store, id) {
Some(record) => record.state == Acknowledged
None => false
}
}
///|
/// Record a retry; messages over the retry budget move to the dead-letter list.
pub fn app_store_record_retry(
store : AppMessageStore,
id : String,
) -> AppMessageStore {
match store_record_index(store, id) {
None => store
Some(index) => {
let records = store_copy_records(store)
let previous = records[index]
let retries = previous.retries + 1
let terminal = retries >= store.max_retries
records[index] = {
message: previous.message,
state: if terminal {
DeadLetter
} else {
Pending
},
retries,
}
let dead_letters = store_copy_dead_letters(store)
if terminal {
dead_letters.push(previous.message)
}
{
records,
dead_letters,
capacity: store.capacity,
max_retries: store.max_retries,
}
}
}
}
///|
/// Return the retry count for a message, or zero when it is unknown.
pub fn app_store_retry_count(store : AppMessageStore, id : String) -> Int {
match app_store_record(store, id) {
Some(record) => record.retries
None => 0
}
}
///|
/// Return the next pending message in insertion order.
pub fn app_store_next_pending(store : AppMessageStore) -> AppMessage? {
for record in store.records {
if record.state == Pending {
return Some(record.message)
}
}
None
}
///|
/// Remove a record after acknowledgement and preserve all other records.
pub fn app_store_remove_acknowledged(
store : AppMessageStore,
) -> AppMessageStore {
let records = []
for record in store.records {
if record.state != Acknowledged {
records.push(record)
}
}
{
records,
dead_letters: store_copy_dead_letters(store),
capacity: store.capacity,
max_retries: store.max_retries,
}
}
///|
/// Return all message IDs in insertion order.
pub fn app_store_ids(store : AppMessageStore) -> Array[String] {
let ids = []
for record in store.records {
ids.push(record.message.id)
}
ids
}
///|
/// Return all records that are ready for delivery.
pub fn app_store_pending(store : AppMessageStore) -> Array[AppStoreRecord] {
let result = []
for record in store.records {
if record.state == Pending {
result.push(record)
}
}
result
}
///|
/// Return the stored dead-letter messages.
pub fn app_store_dead_letters(store : AppMessageStore) -> Array[AppMessage] {
store_copy_dead_letters(store)
}
///|
/// Calculate lifecycle counters for dashboards and tests.
pub fn app_store_stats(store : AppMessageStore) -> AppStoreStats {
let mut pending = 0
let mut delivered = 0
let mut acknowledged = 0
let mut dead = 0
let mut retries = 0
for record in store.records {
retries = retries + record.retries
match record.state {
Pending => pending = pending + 1
Delivered => delivered = delivered + 1
Acknowledged => acknowledged = acknowledged + 1
DeadLetter => dead = dead + 1
}
}
{
total: store.records.length(),
pending,
delivered,
acknowledged,
dead_letters: dead + store.dead_letters.length(),
retries,
}
}
///|
/// Return the number of available capacity slots.
pub fn app_store_available_capacity(store : AppMessageStore) -> Int {
store.capacity - store.records.length()
}
///|
/// Return true when an ID is already known to the store.
pub fn app_store_contains(store : AppMessageStore, id : String) -> Bool {
store_record_index(store, id) is Some(_)
}
///|
/// Change the capacity only when it can hold the current records.
pub fn app_store_resize(
store : AppMessageStore,
capacity : Int,
) -> AppMessageStore {
let next_capacity = if capacity < store.records.length() {
store.records.length()
} else if capacity < 1 {
1
} else {
capacity
}
{
records: store_copy_records(store),
dead_letters: store_copy_dead_letters(store),
capacity: next_capacity,
max_retries: store.max_retries,
}
}
///|
/// Return a copy with a new retry budget for future failures.
pub fn app_store_set_retry_budget(
store : AppMessageStore,
max_retries : Int,
) -> AppMessageStore {
{
records: store_copy_records(store),
dead_letters: store_copy_dead_letters(store),
capacity: store.capacity,
max_retries: if max_retries < 1 {
1
} else {
max_retries
},
}
}
///|
/// Drain dead letters and return the cleared store plus the drained messages.
pub fn app_store_drain_dead_letters(
store : AppMessageStore,
) -> (AppMessageStore, Array[AppMessage]) {
let drained = store_copy_dead_letters(store)
let next = {
records: store_copy_records(store),
dead_letters: [],
capacity: store.capacity,
max_retries: store.max_retries,
}
(next, drained)
}
///|
/// Create a compact audit string for operational logs.
pub fn app_store_summary(store : AppMessageStore) -> String {
let stats = app_store_stats(store)
"total=\{stats.total} pending=\{stats.pending} delivered=\{stats.delivered} ack=\{stats.acknowledged} dead=\{stats.dead_letters} retries=\{stats.retries}"
}