/** * LockedBlackboard - Atomic Commitment Layer for Multi-Agent Coordination * * This module provides file-system mutex locks to ensure atomic writes to the * swarm-blackboard.md, preventing split-brain scenarios when multiple agents * attempt concurrent updates. * * FEATURES: * - File-system mutexes (cross-platform) * - Atomic propose → validate → commit workflow * - Deadlock prevention with lock timeouts * - Split-brain detection and recovery * * @module LockedBlackboard * @version 1.0.0 * @license MIT */ import type { SecureAuditLogger } from '../security'; /** Conflict resolution strategy for concurrent writes to the same key. */ export type ConflictResolutionStrategy = 'first-commit-wins' | 'priority-wins'; /** Agent priority level (0=low, 1=normal, 2=high, 3=critical). */ export type AgentPriority = 0 | 1 | 2 | 3; /** Configuration options for LockedBlackboard. */ export interface LockedBlackboardOptions { /** How to resolve conflicts when multiple agents write to the same key. * - `'first-commit-wins'` (default): The first validated+committed change wins; later ones are aborted. * - `'priority-wins'`: Higher-priority changes preempt lower-priority pending/committed writes on the same key. * Equal-priority conflicts resolve in favor of the most recent proposal by arrival order (last-writer-wins). */ conflictResolution?: ConflictResolutionStrategy; /** Minimum milliseconds between consecutive write/commit operations (0 = no throttle) */ throttleMs?: number; /** * Environment name (e.g. `'dev'`, `'prod'`). When set, all data is scoped * to `//` keeping environments completely isolated. * Falls back to the `NETWORK_AI_ENV` environment variable when not provided. * The value is captured once at construction time — runtime changes to * `NETWORK_AI_ENV` after the instance is created have no effect. * * **Read isolation:** The mutex protects the `commit` step only. * Between `propose()` and `validate()`, other agents can read stale data. * For read-then-write safety, use `propose()` optimistically and handle * `CONFLICT` rejections from `validateChange()` by re-reading and re-proposing. */ env?: string; /** * When `true` (or when `NETWORK_AI_MINIMAL=1` env var is set), skip WAL * replay on startup and disable TTL sweep. Useful for CI and test * environments where fast startup is more important than crash recovery. */ disableWal?: boolean; } export interface BlackboardEntry { key: string; value: unknown; source_agent: string; timestamp: string; ttl: number | null; version: number; } export interface PendingChange { change_id: string; key: string; value: unknown; source_agent: string; proposed_at: string; ttl: number | null; status: 'pending' | 'validated' | 'committed' | 'aborted'; previous_hash: string | null; /** Agent priority (0=low, 1=normal, 2=high, 3=critical). Defaults to 0. */ priority: AgentPriority; validation?: { validated_at: string; validated_by: string; }; /** Set when this change was preempted by a higher-priority change. */ preempted_by?: string; } export interface LockInfo { locked: boolean; holder?: string; acquired_at?: string; timeout_at?: string; } export interface CommitResult { success: boolean; change_id: string; message: string; entry?: BlackboardEntry; } /** * Metadata about a blackboard entry — returned by {@link LockedBlackboard.readMetadata}. * Deliberately excludes the raw `value` so callers can inspect entry shape * without paying the cost of deserializing (or accidentally leaking) large values. */ export interface BlackboardEntryMetadata { /** Blackboard key. */ key: string; /** JavaScript `typeof` of the stored value. */ type: string; /** Approximate serialised byte size of the value. */ sizeBytes: number; /** Monotonically increasing write counter for this key. */ version: number; /** ISO-8601 timestamp of the last write. */ timestamp: string; /** TTL in milliseconds, or `null` for no expiry. */ ttl: number | null; } /** * Cross-platform file lock using lock files. * Works on Windows, Linux, and macOS. */ export declare class FileLock { private lockPath; private lockHolder; private lockFd; constructor(lockPath: string); private ensureDir; /** * Attempt to acquire the lock with timeout. * @param holderId Unique identifier for the lock holder * @param timeoutMs Maximum time to wait for lock (default: CONFIG.lockTimeoutMs) * @returns true if lock acquired, false if timeout */ acquire(holderId: string, timeoutMs?: number): boolean; /** * Release the lock if we hold it. * Verifies ownership before unlinking to prevent deleting a lock acquired by * another process after ours was force-released as stale. */ release(): boolean; /** * Force release a stale lock (use with caution). * Unconditional — only safe for corrupted lock files. */ forceRelease(): void; /** * Compare-and-delete a stale lock: only unlinks if the file still has the * exact same identity (acquired_at + pid) we observed when we decided it * was stale. Prevents deleting a valid lock that was freshly created by * another waiter that beat us to the cleanup. */ private forceReleaseStale; /** * Check current lock status. */ getStatus(): LockInfo; /** * Check if we hold the lock. */ isHeldByMe(): boolean; private sleep; } /** * LockedBlackboard - Thread-safe blackboard with atomic commits and audit trail. * * Every mutating operation (write, delete, commit) records an audit entry * capturing the lock holder, operation duration, and change details when * an optional {@link SecureAuditLogger} is provided. * * Usage: * ```typescript * const blackboard = new LockedBlackboard('./'); * * // Atomic write workflow * const changeId = blackboard.propose('task:123', { status: 'done' }, 'agent-1'); * const isValid = blackboard.validate(changeId, 'orchestrator'); * if (isValid) { * blackboard.commit(changeId); * } else { * blackboard.abort(changeId); * } * ``` */ export declare class LockedBlackboard { private basePath; private blackboardPath; private lockPath; private pendingDir; private lock; private cache; private pendingChanges; private auditLogger?; private conflictResolution; private paused; private throttleMs; private lastWriteTime; private walPath; private walOpCounter; private sweepTimer; private disableWal; constructor(basePath?: string, auditLoggerOrOptions?: SecureAuditLogger | LockedBlackboardOptions, options?: LockedBlackboardOptions); private initialize; private writeInitialBlackboard; private computeHash; private loadFromDisk; private loadPendingChanges; private cleanupOldPendingChanges; private persistToDisk; private isExpired; private savePendingChange; private archivePendingChange; /** * Pause all write and commit operations. * Read operations continue to work while paused. */ pause(): void; /** * Resume write and commit operations after a pause. */ resume(): void; /** * Check if the blackboard is currently paused. */ isPaused(): boolean; /** * Set the minimum interval between write/commit operations. * @param ms Milliseconds between writes (0 to disable throttling) */ setThrottle(ms: number): void; /** * Get the current throttle interval. */ getThrottle(): number; /** * Guard called before any mutating operation. * Throws if paused; enforces throttle delay. */ private enforceFlowControl; /** * Record the timestamp of a successful mutating operation. */ private recordWrite; /** * STEP 1: Propose a change (does NOT modify blackboard yet). * @param key Blackboard key to write * @param value Value to store * @param sourceAgent Agent proposing the change * @param ttl Optional time-to-live in seconds * @param priority Agent priority (0=low, 1=normal, 2=high, 3=critical). Defaults to 0. * @returns change_id for use in validate/commit/abort */ propose(key: string, value: unknown, sourceAgent: string, ttl?: number, priority?: AgentPriority): string; /** * Validate and clamp priority to the AgentPriority range. */ private validatePriority; /** * STEP 2: Validate a proposed change (check for conflicts). * * In `'first-commit-wins'` mode (default): fails if the key was modified since proposal. * In `'priority-wins'` mode: allows higher-priority changes to preempt lower-priority * pending/validated changes on the same key. * * @returns true if change can be safely committed */ validate(changeId: string, validatorAgent: string): boolean; /** * STEP 3a: Commit a validated change (applies to blackboard). * @returns CommitResult with success status */ commit(changeId: string): CommitResult; /** * STEP 3b: Abort a proposed/validated change. */ abort(changeId: string): boolean; /** * Find all pending/validated changes targeting the same key, excluding the given changeId. */ findConflictingPendingChanges(key: string, excludeChangeId: string): PendingChange[]; /** * Preempt a lower-priority change: abort it and emit an audit event. */ private preempt; /** * Get the priority of the last committed change to a key. * Falls back to 0 if unknown (legacy data without priority). */ private getLastCommittedPriority; /** * Get the current conflict resolution strategy. */ getConflictResolution(): ConflictResolutionStrategy; private persistToDiskInternal; /** * Read a value from the blackboard. */ read(key: string): BlackboardEntry | null; /** * Direct write with automatic locking (use propose/validate/commit for multi-agent safety). */ write(key: string, value: unknown, sourceAgent: string, ttl?: number): BlackboardEntry; /** * Delete a key from the blackboard. */ delete(key: string): boolean; /** * List all valid keys. */ listKeys(): string[]; /** * Return metadata for a single blackboard entry without exposing its value. * * Useful for orchestrators that need to inspect entry shape, age, or size * before deciding whether to read the full value. * * @param key - Blackboard key to query. * @returns Metadata object, or `null` if the key does not exist / has expired. */ readMetadata(key: string): BlackboardEntryMetadata | null; /** * Return metadata for all live (non-expired) blackboard entries. * * @returns Array of metadata objects — one per live key, in insertion order. */ listMetadata(): BlackboardEntryMetadata[]; /** @internal */ private _entryToMetadata; /** * Get full snapshot of blackboard state. */ getSnapshot(): Record; /** * List all pending changes. */ listPendingChanges(): PendingChange[]; /** * Get lock status. */ getLockStatus(): LockInfo; /** * Attach an audit logger at runtime (useful when the logger is created * after the blackboard, e.g., in the orchestrator constructor). */ setAuditLogger(logger: SecureAuditLogger): void; /** * Log an audit entry if an audit logger is attached. * Non-fatal: failures are swallowed so auditing never blocks operations. */ private audit; /** * Evict all expired entries from the in-memory cache and persist to disk. * Called automatically by `startSweep()` at the configured interval. * * @returns Number of entries evicted. */ purgeExpired(): number; /** * Start a background sweep timer that calls `purgeExpired()` periodically. * Safe to call multiple times — stops any existing timer first. * * The timer is unref'd so it does not prevent process exit. * * @param intervalMs Sweep interval in milliseconds (default 60 000 = 1 min). */ startSweep(intervalMs?: number): void; /** * Stop the background sweep timer started by `startSweep()`. * Safe to call even if no sweep is running. */ stopSweep(): void; /** * Append a WAL record for crash recovery. * Failures are logged but never propagate — WAL writes must not block ops. * @internal */ private appendToWAL; /** * Append a checkpoint record — signals that the matching op reached disk. * @internal */ private checkpointWAL; /** * Replay uncommitted WAL entries after a crash. * * Called automatically during construction (after `loadFromDisk()`). * Any WAL record without a matching `checkpoint` is replayed into the cache, * then the full state is persisted and the WAL is compacted. * * Malformed tail lines are silently skipped — partial writes at crash time * leave incomplete JSON that we must tolerate. */ replayWAL(): void; /** * Truncate the WAL file. * * Call after a full-state snapshot has been flushed to disk to prevent * unbounded WAL growth during long-running processes. */ compactWAL(): void; } export default LockedBlackboard; //# sourceMappingURL=locked-blackboard.d.ts.map