= PgQ - queue for PostgreSQL =

== Queue creation ==

{{{
    pgq.create_queue(queue_name text)
}}}

Initialize event queue.

Returns 0 if event queue already exists, 1 otherwise.

== Producer ==

{{{
    pgq.insert_event(producer text, queue_name text, ref_id text, payload text)
}}}

Generate new event.  This should be called inside main tx - thus
rollbacked with it if needed.


== Consumer ==

{{{
    pgq.register_consumer(queue_name text, consumer_id text)
    pgq.register_consumer(queue_name text, consumer_id text, thread_timeout integer)
}}}

Attaches this consumer to particular event queue.

The `thread_timeout` specifies timeout for allocated batches in seconds.
If the batch is not finished in this time, the thread will be considered dead,
batch will be closed and events will be moved to retry queue.  The default
is 5 minutes.

Specifying 0 for thread_timeout disables dead thread reaping.
Then the consumer should make sure it's thread_id's do not change 
between runs, otherwise some batches may stay unprocessed.

More about it: ["../PgqNoDupes"].

Returns 0 if the consumer was already attached, 1 otherwise.

{{{
    pgq.unregister_consumer(queue_name text, consumer_id text)
}}}

Unregister and drop resources allocated to customer.


{{{
    pgq.next_batch(queue_name text, consumer_id text, thread_id text)
}}}

Allocates next batch of events to patricular consumer thread.

Returns batch id (int8), to be used in processing functions.  If no batches
are available, returns NULL.  That means that the ticker has not cut them yet.
This is the appropriate moment for consumer to sleep.

{{{
    pgq.fetch_batch_events(batch_id int8)
}}}

`pgq.fetch_batch_events()` returns a set of events in this batch.

There may be no events in the batch.  This is normal.  The batch must still be closed
with pgq.finish_batch().

Event fields: (event_id int8, ref_id text, payload text, producer text, creation_date, event_txid, retry_count)

{{{
    pgq.event_failed(batch_id int8, event_id int8, reason text)
}}}

Tag event as 'failed' - it will be stored, but not further processing is done.

{{{
    pgq.event_retry(batch_id int8, event_id int8, retry_seconds int4)
}}}

Tag event for 'retry' - after x seconds the event will be re-inserted
into main queue.

{{{
    pgq.finish_batch(batch_id int8, total_count int4, retry_count int4, failed_count int4)
}}}

Tag batch as finished.  Until this is not done, same thread will get
same batch again.  If the thread won't return, after some time the events
will be moved to retry queue.  Events not tagged 'failed' or 'retry'
will be assumed to processed successfully.

After calling finish_batch consumer cannot do any operations with events of that batch.
All operations must be done before.

== Failed queue operation ==

Events tagged as failed just stay on their queue.  Following
functions can be used to manage them.

{{{
    pgq.failed_event_list(queue_name, consumer)
    pgq.failed_event_list(queue_name, consumer, cnt, offset)
    pgq.failed_event_count(queue_name, consumer)
}}}

Get info about the queue.

Event fields: (event_id, ref_id, payload, producer, creation_date, retry_count, failure_date, failure_reason)

{{{
    pgq.failed_event_delete(queue_name, consumer, event_id)
    pgq.failed_event_retry(queue_name, consumer, event_id)
}}}

Remove an event from queue, or retry it.

== Info operations ==

{{{
    pgq.get_queue_list()
}}}

Get list of queues.

Result: (queue_name, rotation_delay, number_of_event_tables, event_timeout, ticker_lag)

{{{
    pgq.get_consumer_list()
    pgq.get_consumer_list(queue_name)
}}}

Get list of active consumers.

Result: (queue_name, consumer_name, thread_timeout, lag, last_seen)

{{{
    pgq.get_batch_info(batch_id)
}}}

Get info about batch.

Result fields: (queue_name, consumer_name, batch_start, batch_end, prev_tick_id, tick_id, lag)

== Notes ==

If the consumer is single-threaded, it can do all the batch in one transaction.
But if it wants to process multi-threaded, then it should call `pgq.next_batch`
in separate transaction.  Otherwise the processing will still happen in single
thread.  It could do even several batches in one TX, but it won't see any
new batches that way and also when the TX will be rollbacked, it needs to do
all the events again.

Producer can do several events in same transaction.

Consumer can also act as a producer, nothing should break.

Consumer '''must''' be able to process same event several times.

== Example ==

First, create event queue:

{{{
    select pgq.create_queue('LogEvent');
}}}

Then, producer side can do whenever it wishes:

{{{
    select pgq.insert_event('SampleProducer', 'LogEvent', '123', 'DataFor123');
}}}

First step for consumer is to register:

{{{
    select pgq.register_consumer('LogEvent', 'TestConsumer');
}}}

Then it can enter into consuming loop:

{{{
    begin;
    select pgq.next_batch('LogEvent', 'TestConsumer', 'thread0'); [into batch_id]
    commit;
}}}

That will reserve a batch of events for this thread.  If same thread calls 'next_batch'
again, without finalizing the batch, it will get same batch_id (unless thread_timeout
is reached and events moved to retry queue).

It is explicitly in separate transaction to let other threads grab their batches.
If consumer has only single thread it does not need to do it and can commit
only once per batch ().

To see the events in batch:

{{{
    select * from pgq.fetch_batch_events(batch_id);
}}}

That will give all events in batch.  The processing does not need to be happen
all in one transaction, framework can split the work how it wants.

If a events failed or needs to be tried again, framework can call:

{{{
    select pgq.event_retry(batch_id, event_id, 60);
    select pgq.event_failed(batch_id, event_id, 'Record deleted');
}}}

When all done, notify core about it:

{{{
    select pgq.finish_batch(batch_id, event_count, retry_count, failed_count);
}}}

'''NB:''': ''pgq.finish_batch'' assumes that all events not explicitly tagged
failed or retry are processed successfully.  This is done for  efficiency reasons
and is different how framework should act.  Framework should assume that
all event not tagged explicitly failed or done by user should be retried.
