Extract timer runner runtime

master
dev 5 months ago
parent e9b1d7f8a4
commit 06bf5f44bc

@ -23,6 +23,7 @@ Estado al cierre:
- `arts/sess/engine-session.ts` delega la resolucion de revoke local/global/degradado en `session-revoke.ts`.
- `arts/perm/client.ts` delega claves/scope de cache en `client-keys.ts` y lectura de snapshot en `client-snapshot.ts`.
- `arts/stor/engine-storage.ts` delega lectura/escritura/validacion/migracion de entradas en `entry-runtime.ts`.
- `arts/timr/engine-timers.ts` delega armado nativo, ejecucion, finalizacion e intervalos en `timer-runner.ts`.
- `arts/conn/connection.ts` delega la programacion de reconnect y exhaustion en `connection-reconnect-runtime.ts`.
- `arts/conn/connection.ts` delega decode/routing de frames entrantes en `connection-message-router.ts`.
- `arts/conn/connection.ts` delega el intento open/auth/flush/join en `connection-connect.ts`.

@ -30,32 +30,22 @@ import {
ENGINE_METHOD_SCHEDULE,
ENGINE_METHOD_SCHEDULE_AT,
TIMER_DIAGNOSTIC_EVENTS,
TIMER_EVENT_COMPLETED,
TIMER_EVENT_DISPOSED,
TIMER_EVENT_FAILED,
TIMER_EVENT_RUNNING,
TIMER_EVENT_SCHEDULED,
TIMER_KIND_INTERVAL,
TIMER_KIND_TIMEOUT,
TIMER_STATUS_CANCELLED,
TIMER_STATUS_COMPLETED,
TIMER_STATUS_FAILED,
TIMER_STATUS_PENDING,
TIMER_STATUS_RUNNING
TIMER_KIND_TIMEOUT
} from './consts.ts';
import { createTimerDiagnostics, emitTimerDiagnostic } from './diagnostics.ts';
import { TimrDuplicateKeyError } from './errors.ts';
import { cancelTimerEntry, collectTimerKeysInScope } from './timer-cancel.ts';
import {
createInternalTimerEntry,
createTimerTaskContext,
getLiveTimerEntry,
hasTimerReachedMaxRuns,
timerEntrySnapshot,
type InternalTimerEntry
} from './timer-entry.ts';
import { createTimerEventBus } from './timer-events.ts';
import { createTimerHandle } from './timer-handle.ts';
import { createTimerRunner } from './timer-runner.ts';
import {
assertTimerDelay,
assertTimerEngineLive,
@ -94,121 +84,7 @@ export function createEngineTimers(options: EngineTimersOptions = {}): EngineTim
events.emit(event);
}
// ── Entry lifecycle ─────────────────────────────────────────────────────
/**
* Arm the native timeout for the given entry. Captures the tuple in
* the closure — `runEntry` will reject the callback if the entry has
* been cancelled, replaced or rescheduled in the meantime.
*/
function armEntry(entry: InternalTimerEntry, delayMs: number): void {
entry.version += 1;
const id = entry.id;
const key = entry.key;
const version = entry.version;
entry.native = clock.setTimeout(() => {
void runEntry(id, key, version);
}, delayMs);
}
async function runEntry(id: number, key: string, version: number): Promise<void> {
const entry = getLiveTimerEntry(entries, id, key, version);
if (entry === null) return;
if (entry.status !== TIMER_STATUS_PENDING) return;
const firedAt = clock.now();
entry.status = TIMER_STATUS_RUNNING;
entry.lastFiredAt = firedAt;
entry.runCount += 1;
entry.native = null;
const ctx = createTimerTaskContext(entry, firedAt);
emit({ type: TIMER_EVENT_RUNNING, entry: timerEntrySnapshot(entry) });
const isInterval = entry.kind === TIMER_KIND_INTERVAL;
const awaitTask = !isInterval || entry.awaitTask;
// Branch A — interval with awaitTask=false: arm the next tick BEFORE
// running so we honour cadence regardless of how long the task
// takes. Stale-callback guard still applies to next tick.
if (isInterval && !awaitTask && !hasTimerReachedMaxRuns(entry)) {
scheduleNextIntervalTick(entry);
}
try {
const maybe = entry.task(ctx);
if (awaitTask && maybe instanceof Promise) {
await maybe;
}
finishEntry(id, key, /* error */ null);
} catch (err) {
finishEntry(id, key, err);
}
}
/**
* Finalise the just-completed task. Identified by `(id, key)` only —
* NOT `version`, because in `awaitTask: false` mode the next tick's
* `armEntry()` already incremented `version` before the previous
* task settled. The version guard is for native-callback rejection,
* not for in-flight finalisation; here we only care that the entry
* is still the same logical instance and not cancelled/replaced.
*/
function finishEntry(id: number, key: string, error: unknown | null): void {
const entry = entries.get(key);
if (entry === undefined) return; // cancelled / replaced (deleted)
if (entry.id !== id) return; // replaced (new id)
if (entry.status === TIMER_STATUS_CANCELLED) return;
const now = clock.now();
const isInterval = entry.kind === TIMER_KIND_INTERVAL;
if (error === null) {
entry.lastCompletedAt = now;
entry.error = null;
emit({ type: TIMER_EVENT_COMPLETED, entry: timerEntrySnapshot(entry) });
} else {
entry.lastErrorAt = now;
entry.error = error;
emitTimerDiagnostic(diagnostics, TIMER_DIAGNOSTIC_EVENTS.TASK_FAILED, {
error,
key: entry.key
});
emit({ type: TIMER_EVENT_FAILED, entry: timerEntrySnapshot(entry), error });
}
if (isInterval) {
// `entry.status` may already be TIMER_STATUS_PENDING if branch A
// already armed the next tick; otherwise schedule it now.
if (entry.status === TIMER_STATUS_RUNNING) {
if (hasTimerReachedMaxRuns(entry)) {
entry.status = error === null ? TIMER_STATUS_COMPLETED : TIMER_STATUS_FAILED;
entries.delete(entry.key);
return;
}
scheduleNextIntervalTick(entry);
}
return;
}
// One-shot terminal state.
entry.status = error === null ? TIMER_STATUS_COMPLETED : TIMER_STATUS_FAILED;
if (entry.removeOnComplete) entries.delete(entry.key);
}
function scheduleNextIntervalTick(entry: InternalTimerEntry): void {
if (entry.intervalMs === null) return;
if (hasTimerReachedMaxRuns(entry)) {
entry.status = TIMER_STATUS_COMPLETED;
entries.delete(entry.key);
return;
}
const now = clock.now();
entry.scheduledAt = now;
entry.delayMs = entry.intervalMs;
entry.dueAt = now + entry.intervalMs;
entry.status = TIMER_STATUS_PENDING;
armEntry(entry, entry.intervalMs);
}
const runner = createTimerRunner({ entries, clock, diagnostics, emit });
// ── Cancellation ────────────────────────────────────────────────────────
@ -256,7 +132,7 @@ export function createEngineTimers(options: EngineTimersOptions = {}): EngineTim
nextId += 1;
entries.set(key, entry);
armEntry(entry, delayMs);
runner.armEntry(entry, delayMs);
emit({ type: TIMER_EVENT_SCHEDULED, entry: timerEntrySnapshot(entry) });
return makeHandle(entry);
}
@ -266,7 +142,7 @@ export function createEngineTimers(options: EngineTimersOptions = {}): EngineTim
entries,
clock,
cancelEntry,
armEntry,
armEntry: runner.armEntry,
emitScheduled: (live) => {
emit({ type: TIMER_EVENT_SCHEDULED, entry: timerEntrySnapshot(live) });
}

@ -0,0 +1,152 @@
import {
TIMER_DIAGNOSTIC_EVENTS,
TIMER_EVENT_COMPLETED,
TIMER_EVENT_FAILED,
TIMER_EVENT_RUNNING,
TIMER_KIND_INTERVAL,
TIMER_STATUS_CANCELLED,
TIMER_STATUS_COMPLETED,
TIMER_STATUS_FAILED,
TIMER_STATUS_PENDING,
TIMER_STATUS_RUNNING
} from './consts.ts';
import { emitTimerDiagnostic, type TimerDiagnostics } from './diagnostics.ts';
import {
createTimerTaskContext,
getLiveTimerEntry,
hasTimerReachedMaxRuns,
timerEntrySnapshot,
type InternalTimerEntry
} from './timer-entry.ts';
import type { TimerClock, TimerEvent } from './types.ts';
interface TimerRunnerOptions {
readonly entries: Map<string, InternalTimerEntry>;
readonly clock: TimerClock;
readonly diagnostics: TimerDiagnostics;
emit(event: TimerEvent): void;
}
export interface TimerRunner {
armEntry(entry: InternalTimerEntry, delayMs: number): void;
}
export function createTimerRunner(options: TimerRunnerOptions): TimerRunner {
const { entries, clock, diagnostics, emit } = options;
/**
* Arm the native timeout for the given entry. Captures the tuple in
* the closure — `runEntry` will reject the callback if the entry has
* been cancelled, replaced or rescheduled in the meantime.
*/
function armEntry(entry: InternalTimerEntry, delayMs: number): void {
entry.version += 1;
const id = entry.id;
const key = entry.key;
const version = entry.version;
entry.native = clock.setTimeout(() => {
void runEntry(id, key, version);
}, delayMs);
}
async function runEntry(id: number, key: string, version: number): Promise<void> {
const entry = getLiveTimerEntry(entries, id, key, version);
if (entry === null) return;
if (entry.status !== TIMER_STATUS_PENDING) return;
const firedAt = clock.now();
entry.status = TIMER_STATUS_RUNNING;
entry.lastFiredAt = firedAt;
entry.runCount += 1;
entry.native = null;
const ctx = createTimerTaskContext(entry, firedAt);
emit({ type: TIMER_EVENT_RUNNING, entry: timerEntrySnapshot(entry) });
const isInterval = entry.kind === TIMER_KIND_INTERVAL;
const awaitTask = !isInterval || entry.awaitTask;
// Branch A — interval with awaitTask=false: arm the next tick BEFORE
// running so we honour cadence regardless of how long the task
// takes. Stale-callback guard still applies to next tick.
if (isInterval && !awaitTask && !hasTimerReachedMaxRuns(entry)) {
scheduleNextIntervalTick(entry);
}
try {
const maybe = entry.task(ctx);
if (awaitTask && maybe instanceof Promise) {
await maybe;
}
finishEntry(id, key, null);
} catch (err) {
finishEntry(id, key, err);
}
}
/**
* Finalise the just-completed task. Identified by `(id, key)` only —
* NOT `version`, because in `awaitTask: false` mode the next tick's
* `armEntry()` already incremented `version` before the previous
* task settled. The version guard is for native-callback rejection,
* not for in-flight finalisation; here we only care that the entry
* is still the same logical instance and not cancelled/replaced.
*/
function finishEntry(id: number, key: string, error: unknown | null): void {
const entry = entries.get(key);
if (entry === undefined) return; // cancelled / replaced (deleted)
if (entry.id !== id) return; // replaced (new id)
if (entry.status === TIMER_STATUS_CANCELLED) return;
const now = clock.now();
const isInterval = entry.kind === TIMER_KIND_INTERVAL;
if (error === null) {
entry.lastCompletedAt = now;
entry.error = null;
emit({ type: TIMER_EVENT_COMPLETED, entry: timerEntrySnapshot(entry) });
} else {
entry.lastErrorAt = now;
entry.error = error;
emitTimerDiagnostic(diagnostics, TIMER_DIAGNOSTIC_EVENTS.TASK_FAILED, {
error,
key: entry.key
});
emit({ type: TIMER_EVENT_FAILED, entry: timerEntrySnapshot(entry), error });
}
if (isInterval) {
// `entry.status` may already be TIMER_STATUS_PENDING if branch A
// already armed the next tick; otherwise schedule it now.
if (entry.status === TIMER_STATUS_RUNNING) {
if (hasTimerReachedMaxRuns(entry)) {
entry.status = error === null ? TIMER_STATUS_COMPLETED : TIMER_STATUS_FAILED;
entries.delete(entry.key);
return;
}
scheduleNextIntervalTick(entry);
}
return;
}
// One-shot terminal state.
entry.status = error === null ? TIMER_STATUS_COMPLETED : TIMER_STATUS_FAILED;
if (entry.removeOnComplete) entries.delete(entry.key);
}
function scheduleNextIntervalTick(entry: InternalTimerEntry): void {
if (entry.intervalMs === null) return;
if (hasTimerReachedMaxRuns(entry)) {
entry.status = TIMER_STATUS_COMPLETED;
entries.delete(entry.key);
return;
}
const now = clock.now();
entry.scheduledAt = now;
entry.delayMs = entry.intervalMs;
entry.dueAt = now + entry.intervalMs;
entry.status = TIMER_STATUS_PENDING;
armEntry(entry, entry.intervalMs);
}
return { armEntry };
}
Loading…
Cancel
Save

Powered by TurnKey Linux.