apply method

void apply(
  1. QueueEvent event, {
  2. bool broadcast = true,
})

Folds a single event into the model. broadcast is false during history replay so the SSE tap only carries live events.

Implementation

void apply(QueueEvent event, {bool broadcast = true}) {
  final ts = event.timestamp;
  if (event.type.isTask && event.taskId != null) {
    final t = _task(event.taskId!, event.taskName ?? '?');
    if (event.taskName != null) t.name = event.taskName!;
    if (event.queue != null) t.queue = event.queue;
    if (event.workerId != null) t.workerId = event.workerId;
    if (ts.isAfter(t.updatedAt)) t.updatedAt = ts;
    final f = event.fields;
    if (f['retries'] is num) t.retries = (f['retries'] as num).toInt();
    if (f['max_retries'] is num) {
      t.maxRetries = (f['max_retries'] as num).toInt();
    }

    switch (event.type) {
      case QueueEventType.taskSent:
        // Idempotent: Rose may apply locally and also receive its own emit
        // back through the event source.
        if (t.sentAt == null) sent++;
        t.state = TaskState.pending;
        t.sentAt = ts;
        t.args = f['args'] ?? t.args;
        t.kwargs = f['kwargs'] ?? t.kwargs;
        if (f['eta'] is String) t.state = TaskState.scheduled;
        break;
      case QueueEventType.taskReceived:
        received++;
        t.receivedAt = ts;
        t.args ??= f['args'];
        t.kwargs ??= f['kwargs'];
        break;
      case QueueEventType.taskStarted:
        t.state = TaskState.started;
        t.startedAt = ts;
        break;
      case QueueEventType.taskSucceeded:
        succeeded++;
        t.state = TaskState.success;
        t.finishedAt = ts;
        t.result = f['result'];
        if (f['runtime_ms'] is num) {
          t.runtimeMs = (f['runtime_ms'] as num).toInt();
        }
        break;
      case QueueEventType.taskFailed:
        failed++;
        t.state = TaskState.failure;
        t.finishedAt = ts;
        t.error = f['error'] as String?;
        if (f['runtime_ms'] is num) {
          t.runtimeMs = (f['runtime_ms'] as num).toInt();
        }
        break;
      case QueueEventType.taskRetried:
        retried++;
        t.state = TaskState.retry;
        break;
      case QueueEventType.taskRateLimited:
        rateLimited++;
        t.state = TaskState.scheduled;
        break;
      case QueueEventType.taskRevoked:
        revoked++;
        t.state = TaskState.revoked;
        t.finishedAt = ts;
        t.error = f['reason'] as String?;
        break;
      default:
        break;
    }
  } else if (event.type.isWorker && event.workerId != null) {
    final w = _worker(event.workerId!);
    if (event.hostname != null) w.hostname = event.hostname;
    final f = event.fields;
    if (f['queues'] is List) {
      w.queues = (f['queues'] as List).map((e) => '$e').toList();
    }
    if (f['concurrency'] is num) {
      w.concurrency = (f['concurrency'] as num).toInt();
    }
    if (f['active'] is num) w.active = (f['active'] as num).toInt();
    if (f['succeeded'] is num) w.succeeded = (f['succeeded'] as num).toInt();
    if (f['failed'] is num) w.failed = (f['failed'] as num).toInt();
    if (f['retried'] is num) w.retried = (f['retried'] as num).toInt();
    if (ts.isAfter(w.lastHeartbeat)) w.lastHeartbeat = ts;
    w.online = event.type != QueueEventType.workerOffline;
  }

  if (broadcast && _tap.hasListener) _tap.add(event);
}