apply method
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);
}