Skip to content

Job Lifecycle Events and Hooks

DodaTech Updated 2026-06-28 6 min read

In this tutorial, you will learn about Job Lifecycle Events and Hooks. We cover key concepts, practical examples, and best practices to help you master this topic.

Track job lifecycle events from creation to completion with event hooks for enqueue, start, success, failure, retry, and completion stages in background processing.

What You Learn

You will learn how to implement lifecycle events, register event listeners, build event-driven workflows, and use events for auditing, metrics, and notifications.

Why It Matters

Lifecycle events decouple job processing from side effects. Instead of adding audit logging, metrics, and notifications to every job, you emit events and let listeners handle them.

Real-World Use

DodaTech's job system emits lifecycle events to: record audit trail in database, update real-time dashboard via Websocket, send Slack notification on failure, and collect Prometheus metrics. Each listener is independent.

Lifecycle Event Model

flowchart LR
    C[Created] --> E[Enqueued]
    E --> S[Started]
    S --> P[Processing]
    P --> SC[Success]
    P --> F[Failure]
    P --> R[Retry]
    SC --> D[Done]
    F --> D
    R --> E
    D --> A[Archived]

Event Emitter

import time
import threading

class EventEmitter:
    def __init__(self):
        self._listeners = {}

    def on(self, event, handler):
        if event not in self._listeners:
            self._listeners[event] = []
        self._listeners[event].append(handler)
        return handler

    def off(self, event, handler):
        if event in self._listeners:
            self._listeners[event] = [
                h for h in self._listeners[event] if h != handler
            ]

    def emit(self, event, **data):
        listeners = self._listeners.get(event, [])
        for handler in listeners:
            try:
                handler(data)
            except Exception as e:
                print(f"Listener error for {event}: {e}")

class LifecycleJob(EventEmitter):
    def __init__(self, name):
        super().__init__()
        self.name = name

    def enqueue(self, job_data):
        self.emit('enqueued', job_name=self.name, data=job_data, timestamp=time.time())

    def start(self):
        self.emit('started', job_name=self.name, timestamp=time.time())

    def success(self, result):
        self.emit('succeeded', job_name=self.name, result=result, timestamp=time.time())

    def failure(self, error):
        self.emit('failed', job_name=self.name, error=str(error), timestamp=time.time())

    def retry(self, attempt):
        self.emit('retried', job_name=self.name, attempt=attempt, timestamp=time.time())

job = LifecycleJob('email_sender')

@job.on('enqueued')
def log_enqueue(data):
    print(f"Log: Job enqueued - {data['job_name']}")

@job.on('started')
def metrics_start(data):
    print(f"Metrics: Job started - {data['job_name']}")

@job.on('succeeded')
def notify_success(data):
    print(f"Notify: Job completed - {data['job_name']}")

@job.on('failed')
def log_failure(data):
    print(f"Alert: Job FAILED - {data['job_name']}: {data['error']}")

job.enqueue({'to': 'user@example.com'})
job.start()
job.success({'sent': True})

Expected output:

Log: Job enqueued - email_sender
Metrics: Job started - email_sender
Notify: Job completed - email_sender

Lifecycle-Managed Worker

import time
import threading
import random

class LifecycleWorker:
    def __init__(self):
        self.listeners = {
            'before_enqueue': [],
            'after_enqueue': [],
            'before_process': [],
            'after_success': [],
            'after_failure': [],
            'after_retry': [],
        }

    def on(self, event, handler):
        if event in self.listeners:
            self.listeners[event].append(handler)

    def _trigger(self, event, context):
        for handler in self.listeners.get(event, []):
            try:
                handler(context)
            except Exception as e:
                print(f"Event handler error ({event}): {e}")

    def process(self, job_func, **context):
        self._trigger('before_enqueue', context)
        print(f"Enqueued: {context.get('name', 'job')}")
        self._trigger('after_enqueue', context)

        self._trigger('before_process', context)

        max_retries = context.get('max_retries', 3)
        for attempt in range(max_retries):
            try:
                print(f"  Attempt {attempt + 1}")
                result = job_func(context)
                context['result'] = result
                self._trigger('after_success', context)
                return result
            except Exception as e:
                context['error'] = str(e)
                context['attempt'] = attempt + 1
                if attempt < max_retries - 1:
                    self._trigger('after_retry', context)
                    time.sleep(0.5)
                else:
                    self._trigger('after_failure', context)
                    raise

worker = LifecycleWorker()

@worker.on('before_enqueue')
def audit_enqueue(ctx):
    print(f"  Audit: Enqueuing {ctx['name']}")

@worker.on('after_success')
def metrics_success(ctx):
    print(f"  Metrics: {ctx['name']} succeeded")

@worker.on('after_failure')
def alert_failure(ctx):
    print(f"  Alert: {ctx['name']} failed after {ctx.get('attempt', 0)} attempts")

@worker.on('after_retry')
def log_retry(ctx):
    print(f"  Log: Retrying {ctx['name']} (attempt {ctx['attempt']})")

def unreliable_task(ctx):
    if random.random() < 0.7:
        raise ValueError("Transient error")
    return {'processed': True}

try:
    worker.process(unreliable_task, name='data_sync', max_retries=3)
except ValueError:
    print("Job failed permanently")

Expected output:

  Audit: Enqueuing data_sync
Enqueued: data_sync
  Attempt 1
  Log: Retrying data_sync (attempt 1)
  Attempt 2
  Metrics: data_sync succeeded

Event Persistence and Replay

import json
import time
import sqlite3

class EventStore:
    def __init__(self, db_path=':memory:'):
        self.conn = sqlite3.connect(db_path)
        self.conn.execute('''
            CREATE TABLE IF NOT EXISTS lifecycle_events (
                id INTEGER PRIMARY KEY AUTOINCREMENT,
                job_id TEXT,
                event_type TEXT,
                data TEXT,
                created_at REAL
            )
        ''')
        self.conn.commit()
        self.listeners = {}

    def on(self, event_type, handler):
        if event_type not in self.listeners:
            self.listeners[event_type] = []
        self.listeners[event_type].append(handler)

    def record(self, job_id, event_type, data=None):
        self.conn.execute(
            'INSERT INTO lifecycle_events (job_id, event_type, data, created_at) VALUES (?, ?, ?, ?)',
            (job_id, event_type, json.dumps(data or {}), time.time())
        )
        self.conn.commit()

        for handler in self.listeners.get(event_type, []):
            try:
                handler(job_id, data)
            except Exception as e:
                print(f"Handler error: {e}")

    def get_events(self, job_id=None, event_type=None):
        query = 'SELECT * FROM lifecycle_events WHERE 1=1'
        params = []
        if job_id:
            query += ' AND job_id = ?'
            params.append(job_id)
        if event_type:
            query += ' AND event_type = ?'
            params.append(event_type)
        query += ' ORDER BY created_at'

        cursor = self.conn.execute(query, params)
        return [
            {
                'id': row[0], 'job_id': row[1],
                'event_type': row[2],
                'data': json.loads(row[3]),
                'created_at': row[4],
            }
            for row in cursor.fetchall()
        ]

store = EventStore()
store.on('failed', lambda jid, d: print(f"Alert: Job {jid} failed"))
store.on('completed', lambda jid, d: print(f"Log: Job {jid} completed"))

store.record('scan-001', 'started', {'file': 'doc.pdf'})
store.record('scan-001', 'completed', {'threats': 0})
store.record('scan-002', 'started', {'file': 'virus.exe'})
store.record('scan-002', 'failed', {'error': 'Timeout'})

events = store.get_events(job_id='scan-002')
print(f"Events for scan-002: {len(events)}")
for ev in events:
    print(f"  {ev['event_type']}: {ev['data']}")

Expected output:

Log: Job scan-001 completed
Alert: Job scan-002 failed
Events for scan-002: 2
  started: {'file': 'virus.exe'}
  failed: {'error': 'Timeout'}

Common Mistakes

1. Blocking Event Handlers

Slow event handlers block job processing. Run event handlers asynchronously or in separate threads.

2. Event Handler Exceptions

One failing handler should not break other handlers or the job itself. Wrap each handler in try/except.

3. Too Many Events

Firing an event for every minor state change creates noise. Define meaningful lifecycle stages: enqueued, started, succeeded, failed, retried.

4. No Event Ordering Guarantee

Asynchronous event handlers may Process events out of order. Include timestamps and sequence numbers for ordering.

5. Memory Leaks from Listeners

Listeners registered but never removed accumulate over time. Implement listener cleanup for short-lived jobs.

Practice Questions

1. What are lifecycle events in job processing?

Signals emitted at key stages of a job's life: enqueued, started, succeeded, failed, retried. Listeners react to these events.

2. How do events differ from middleware?

Middleware wraps execution and can modify behavior. Events are notifications after the fact. Events are for side effects, middleware for cross-cutting concerns.

3. Why persist lifecycle events?

For auditing, debugging, replay, and analytics. Persisted events provide a complete history of every job's journey through the system.

4. How do you prevent event handler failures from affecting jobs?

Wrap each handler in try/except. Log handler failures but do not propagate them to the job execution.

Challenge

Build a lifecycle event system for a payment processing queue: events for created, validated, processing, succeeded, failed, refunded. Persist events to database. Event handlers for audit log, metrics, Slack notification, and customer email.

FAQ

What is the difference between events and hooks?

Events are emitted after an action occurs. Hooks run before or after an action and can modify behavior. Events are passive, hooks are active.

Can lifecycle events trigger other jobs?

Yes. A 'succeeded' event can enqueue a follow-up job. This creates event-driven job chains without hardcoded dependencies.

How many events should a job emit?

4-6 core events: enqueued, started, succeeded, failed, retried. Add domain-specific events if needed for your use case.

Are lifecycle events async?

They should be. Fire-and-forget event emission prevents event handlers from blocking job execution.

How do I debug event handler issues?

Log every event emission and handler execution. Use a debug event handler that records all events for later inspection.

Mini Project: Lifecycle System

import time
import json

class JobLifecycle:
    def __init__(self):
        self.handlers = {}

    def on(self, event):
        def decorator(func):
            if event not in self.handlers:
                self.handlers[event] = []
            self.handlers[event].append(func)
            return func
        return decorator

    def emit(self, event, **data):
        for handler in self.handlers.get(event, []):
            handler(data)

    def run(self, job_func, **context):
        job_id = context.get('id', str(time.time()))
        self.emit('start', job_id=job_id, **context)
        try:
            result = job_func(context)
            self.emit('success', job_id=job_id, result=result)
            return result
        except Exception as e:
            self.emit('fail', job_id=job_id, error=str(e))
            raise

lc = JobLifecycle()

@lc.on('start')
def on_start(data):
    print(f"  [Event] Started: {data['job_id']}")

@lc.on('success')
def on_success(data):
    print(f"  [Event] Succeeded: {data['job_id']}")

@lc.on('fail')
def on_fail(data):
    print(f"  [Event] Failed: {data['job_id']}: {data['error']}")

lc.run(lambda ctx: time.sleep(0.1) or 'ok', id='job-001')

Expected output:

  [Event] Started: job-001
  [Event] Succeeded: job-001

What's Next

Now that you understand lifecycle events, explore job metrics with Prometheus for monitoring, then learn about job dashboard for visualization.

Built by the developers of DodaTech

Doda Browser, DodaZIP & Durga Antivirus Pro