2015-04-01 17:01:23 +02:00
|
|
|
unit MVCFramework.MessagingController;
|
2013-10-30 00:48:23 +01:00
|
|
|
|
|
|
|
interface
|
|
|
|
|
|
|
|
uses
|
|
|
|
MVCFramework,
|
|
|
|
StompClient,
|
|
|
|
StompTypes;
|
|
|
|
|
|
|
|
type
|
|
|
|
|
|
|
|
[MVCPath('/messages')]
|
|
|
|
TMVCBUSController = class(TMVCController)
|
2015-04-01 17:01:23 +02:00
|
|
|
strict protected
|
2013-10-30 00:48:23 +01:00
|
|
|
function GetUniqueDurableHeader(clientid, topicname: string): string;
|
2015-04-01 17:01:23 +02:00
|
|
|
procedure InternalSubscribeUserToTopics(clientid: string; Stomp: IStompClient);
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure InternalSubscribeUserToTopic(clientid: string; topicname: string;
|
|
|
|
StompClient: IStompClient);
|
|
|
|
|
|
|
|
procedure AddTopicToUserSubscriptions(const ATopic: string);
|
|
|
|
procedure RemoveTopicFromUserSubscriptions(const ATopic: string);
|
|
|
|
procedure OnBeforeAction(Context: TWebContext; const AActionNAme: string;
|
|
|
|
var Handled: Boolean); override;
|
|
|
|
|
|
|
|
public
|
2015-04-01 17:01:23 +02:00
|
|
|
[MVCHTTPMethod([httpPOST])]
|
|
|
|
[MVCPath('/clients/($clientid)')]
|
|
|
|
procedure SetClientID(CTX: TWebContext);
|
|
|
|
|
|
|
|
[MVCPath('/subscriptions/($topicorqueue)/($name)')]
|
|
|
|
[MVCHTTPMethod([httpPOST])]
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure SubscribeToTopic(CTX: TWebContext);
|
2015-04-01 17:01:23 +02:00
|
|
|
|
|
|
|
[MVCPath('/subscriptions/($topicorqueue)/($name)')]
|
|
|
|
[MVCHTTPMethod([httpDELETE])]
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure UnSubscribeFromTopic(CTX: TWebContext);
|
2015-04-01 17:01:23 +02:00
|
|
|
|
|
|
|
[MVCHTTPMethod([httpGET])]
|
|
|
|
[MVCPath]
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure ReceiveMessages(CTX: TWebContext);
|
|
|
|
|
|
|
|
[MVCHTTPMethod([httpPOST])]
|
2015-04-01 17:01:23 +02:00
|
|
|
[MVCPath('/($type)/($topicorqueue)')]
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure EnqueueMessage(CTX: TWebContext);
|
2015-04-01 17:01:23 +02:00
|
|
|
|
|
|
|
[MVCHTTPMethod([httpGET])]
|
|
|
|
[MVCPath('/subscriptions')]
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure CurrentlySubscribedTopics(CTX: TWebContext);
|
|
|
|
end;
|
|
|
|
|
|
|
|
implementation
|
|
|
|
|
|
|
|
{ TMVCBUSController }
|
|
|
|
|
|
|
|
uses
|
|
|
|
System.SysUtils,
|
|
|
|
MVCFramework.Commons,
|
|
|
|
System.DateUtils,
|
2014-09-05 12:47:40 +02:00
|
|
|
{$IF CompilerVersion < 27}
|
2013-10-30 00:48:23 +01:00
|
|
|
Data.DBXJSON,
|
2014-04-16 22:52:25 +02:00
|
|
|
{$ELSE}
|
|
|
|
System.JSON,
|
|
|
|
{$IFEND}
|
2013-10-30 00:48:23 +01:00
|
|
|
MVCFramework.Logger,
|
|
|
|
System.SyncObjs;
|
|
|
|
|
|
|
|
procedure TMVCBUSController.AddTopicToUserSubscriptions(const ATopic: string);
|
|
|
|
var
|
2013-11-09 14:22:11 +01:00
|
|
|
x: string;
|
2013-10-30 00:48:23 +01:00
|
|
|
topics: TArray<string>;
|
2013-11-09 14:22:11 +01:00
|
|
|
t: string;
|
|
|
|
ToAdd: Boolean;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
x := Session['__subscriptions'];
|
|
|
|
topics := x.Split([';']);
|
|
|
|
ToAdd := true;
|
|
|
|
for t in topics do
|
|
|
|
if t.Equals(ATopic) then
|
|
|
|
begin
|
|
|
|
ToAdd := False;
|
|
|
|
end;
|
|
|
|
if ToAdd then
|
|
|
|
begin
|
|
|
|
SetLength(topics, length(topics) + 1);
|
|
|
|
topics[length(topics) - 1] := ATopic;
|
|
|
|
Session['__subscriptions'] := string.Join(';', topics);
|
|
|
|
end;
|
|
|
|
end;
|
|
|
|
|
|
|
|
procedure TMVCBUSController.CurrentlySubscribedTopics(CTX: TWebContext);
|
|
|
|
begin
|
|
|
|
ContentType := TMVCMimeType.TEXT_PLAIN;
|
|
|
|
Render(Session['__subscriptions']);
|
|
|
|
end;
|
|
|
|
|
|
|
|
procedure TMVCBUSController.EnqueueMessage(CTX: TWebContext);
|
|
|
|
var
|
|
|
|
topicname: string;
|
2015-04-01 17:01:23 +02:00
|
|
|
queuetype: string;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
queuetype := CTX.Request.Params['type'].Trim.ToLower;
|
|
|
|
if (queuetype <> 'topic') and (queuetype <> 'queue') then
|
|
|
|
raise EMVCException.Create('Valid type are "queue" or "topic", got ' + queuetype);
|
|
|
|
|
|
|
|
topicname := CTX.Request.Params['topicorqueue'].Trim;
|
2013-10-30 00:48:23 +01:00
|
|
|
if topicname.IsEmpty then
|
|
|
|
raise EMVCException.Create('Invalid or empty topic');
|
2015-04-01 17:01:23 +02:00
|
|
|
if not CTX.Request.ThereIsRequestBody then
|
|
|
|
raise EMVCException.Create('Body request required');
|
|
|
|
EnqueueMessageOnTopicOrQueue(queuetype = 'queue', '/' + queuetype + '/' + topicname,
|
2013-10-30 00:48:23 +01:00
|
|
|
CTX.Request.BodyAsJSONObject.Clone as TJSONObject, true);
|
2015-04-01 17:01:23 +02:00
|
|
|
// EnqueueMessage('/queue/' + topicname, CTX.Request.BodyAsJSONObject.Clone as TJSONObject, true);
|
2013-10-30 00:48:23 +01:00
|
|
|
Render(200, 'Message sent to topic ' + topicname);
|
|
|
|
end;
|
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
function TMVCBUSController.GetUniqueDurableHeader(clientid, topicname: string): string;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
Result := clientid + '___' + topicname.Replace('/', '_', [rfReplaceAll]);
|
|
|
|
end;
|
|
|
|
|
|
|
|
procedure TMVCBUSController.ReceiveMessages(CTX: TWebContext);
|
|
|
|
var
|
2013-11-09 14:22:11 +01:00
|
|
|
Stomp: IStompClient;
|
2015-04-01 17:01:23 +02:00
|
|
|
LClientID: string;
|
2013-11-09 14:22:11 +01:00
|
|
|
frame: IStompFrame;
|
|
|
|
obj, res: TJSONObject;
|
2015-04-01 17:01:23 +02:00
|
|
|
LFrames: TArray<IStompFrame>;
|
2013-11-09 14:22:11 +01:00
|
|
|
arr: TJSONArray;
|
2015-04-01 17:01:23 +02:00
|
|
|
LLastReceivedMessageTS: TDateTime;
|
|
|
|
LTimeOut: Boolean;
|
2013-10-30 00:48:23 +01:00
|
|
|
const
|
|
|
|
|
2013-11-09 14:22:11 +01:00
|
|
|
{$IFDEF TEST}
|
2015-04-01 17:01:23 +02:00
|
|
|
RECEIVE_TIMEOUT = 5; // seconds
|
2013-10-30 00:48:23 +01:00
|
|
|
|
2013-11-09 14:22:11 +01:00
|
|
|
{$ELSE}
|
2013-10-30 00:48:23 +01:00
|
|
|
RECEIVE_TIMEOUT = 60 * 5; // 5 minutes
|
|
|
|
|
2013-11-09 14:22:11 +01:00
|
|
|
{$ENDIF}
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
LTimeOut := False;
|
|
|
|
LClientID := GetClientID;
|
|
|
|
Stomp := GetNewStompClient(LClientID);
|
2013-10-30 00:48:23 +01:00
|
|
|
try
|
2015-04-01 17:01:23 +02:00
|
|
|
InternalSubscribeUserToTopics(LClientID, Stomp);
|
2013-11-09 14:22:11 +01:00
|
|
|
// StartReceiving := now;
|
2013-10-30 00:48:23 +01:00
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
LLastReceivedMessageTS := now;
|
|
|
|
SetLength(LFrames, 0);
|
2013-10-30 00:48:23 +01:00
|
|
|
while not IsShuttingDown do
|
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
LTimeOut := False;
|
2013-10-30 00:48:23 +01:00
|
|
|
frame := nil;
|
2015-04-01 17:01:23 +02:00
|
|
|
Log('/messages receive');
|
|
|
|
Stomp.Receive(frame, 100);
|
2013-11-09 14:22:11 +01:00
|
|
|
if Assigned(frame) then
|
|
|
|
// get 10 messages at max, and then send them to client
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
LLastReceivedMessageTS := now;
|
|
|
|
SetLength(LFrames, length(LFrames) + 1);
|
|
|
|
LFrames[length(LFrames) - 1] := frame;
|
2013-10-30 00:48:23 +01:00
|
|
|
Stomp.Ack(frame.MessageID);
|
2015-04-01 17:01:23 +02:00
|
|
|
if length(LFrames) >= 10 then
|
2013-10-30 00:48:23 +01:00
|
|
|
break;
|
|
|
|
end
|
|
|
|
else
|
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
if (length(LFrames) > 0) then
|
2013-10-30 00:48:23 +01:00
|
|
|
break;
|
2015-04-01 17:01:23 +02:00
|
|
|
if SecondsBetween(now, LLastReceivedMessageTS) >= RECEIVE_TIMEOUT then
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
LTimeOut := true;
|
2013-10-30 00:48:23 +01:00
|
|
|
break;
|
|
|
|
end;
|
|
|
|
end;
|
|
|
|
end;
|
|
|
|
|
|
|
|
arr := TJSONArray.Create;
|
|
|
|
res := TJSONObject.Create(TJSONPair.Create('messages', arr));
|
2015-04-01 17:01:23 +02:00
|
|
|
for frame in LFrames do
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
if Assigned(frame) then
|
|
|
|
begin
|
|
|
|
obj := TJSONObject.ParseJSONValue(frame.GetBody) as TJSONObject;
|
|
|
|
if Assigned(obj) then
|
|
|
|
begin
|
|
|
|
arr.AddElement(obj);
|
|
|
|
end
|
|
|
|
else
|
|
|
|
begin
|
2013-11-09 14:22:11 +01:00
|
|
|
LogE(Format
|
|
|
|
('Not valid JSON object in topic requested by user %s. The raw message is "%s"',
|
2015-04-01 17:01:23 +02:00
|
|
|
[LClientID, frame.GetBody]));
|
2013-10-30 00:48:23 +01:00
|
|
|
end;
|
|
|
|
end;
|
|
|
|
end; // for in
|
|
|
|
res.AddPair('_timestamp', FormatDateTime('yyyy-mm-dd hh:nn:ss', now));
|
2015-04-01 17:01:23 +02:00
|
|
|
if LTimeOut then
|
|
|
|
begin
|
|
|
|
res.AddPair('_timeout', TJSONTrue.Create);
|
|
|
|
Render(http_status.RequestTimeout, res);
|
|
|
|
end
|
2013-10-30 00:48:23 +01:00
|
|
|
else
|
2015-04-01 17:01:23 +02:00
|
|
|
begin
|
|
|
|
res.AddPair('_timeout', TJSONFalse.Create);
|
|
|
|
Render(http_status.OK, res);
|
|
|
|
end;
|
2013-10-30 00:48:23 +01:00
|
|
|
|
|
|
|
finally
|
|
|
|
// Stomp.Disconnect;
|
|
|
|
end;
|
|
|
|
end;
|
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
procedure TMVCBUSController.RemoveTopicFromUserSubscriptions(const ATopic: string);
|
2013-10-30 00:48:23 +01:00
|
|
|
var
|
2013-11-09 14:22:11 +01:00
|
|
|
x: string;
|
2013-10-30 00:48:23 +01:00
|
|
|
topics, afterremovaltopics: TArray<string>;
|
2013-11-09 14:22:11 +01:00
|
|
|
IndexToRemove: Integer;
|
|
|
|
i: Integer;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
x := Session['__subscriptions'];
|
|
|
|
topics := x.Split([';']);
|
|
|
|
IndexToRemove := 0;
|
|
|
|
SetLength(afterremovaltopics, length(topics));
|
|
|
|
for i := 0 to length(topics) - 1 do
|
|
|
|
begin
|
|
|
|
if not topics[i].Equals(ATopic) then
|
|
|
|
begin
|
|
|
|
afterremovaltopics[IndexToRemove] := topics[i];
|
|
|
|
Inc(IndexToRemove);
|
|
|
|
end;
|
|
|
|
end;
|
|
|
|
if IndexToRemove <> length(ATopic) - 1 then
|
|
|
|
SetLength(afterremovaltopics, length(topics) - 1);
|
|
|
|
|
|
|
|
if length(afterremovaltopics) = 0 then
|
|
|
|
Session['__subscriptions'] := ''
|
|
|
|
else
|
|
|
|
Session['__subscriptions'] := string.Join(';', afterremovaltopics);
|
|
|
|
end;
|
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
procedure TMVCBUSController.SetClientID(CTX: TWebContext);
|
|
|
|
begin
|
|
|
|
Session[CLIENTID_KEY] := CTX.Request.Params['clientid'];
|
|
|
|
end;
|
|
|
|
|
2013-10-30 00:48:23 +01:00
|
|
|
procedure TMVCBUSController.SubscribeToTopic(CTX: TWebContext);
|
|
|
|
var
|
2015-04-01 17:01:23 +02:00
|
|
|
LStomp: IStompClient;
|
|
|
|
LClientID: string;
|
|
|
|
LTopicName: string;
|
|
|
|
LTopicOrQueue: string;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
LClientID := GetClientID;
|
|
|
|
LTopicName := CTX.Request.Params['name'].ToLower;
|
|
|
|
LTopicOrQueue := CTX.Request.Params['topicorqueue'].ToLower;
|
|
|
|
LStomp := GetNewStompClient(LClientID);
|
2013-10-30 00:48:23 +01:00
|
|
|
try
|
2015-04-01 17:01:23 +02:00
|
|
|
LTopicName := '/' + LTopicOrQueue + '/' + LTopicName;
|
|
|
|
InternalSubscribeUserToTopic(LClientID, LTopicName, LStomp);
|
|
|
|
Render(200, 'Subscription OK for ' + LTopicName);
|
2013-10-30 00:48:23 +01:00
|
|
|
finally
|
|
|
|
// Stomp.Disconnect;
|
|
|
|
end;
|
|
|
|
end;
|
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
procedure TMVCBUSController.InternalSubscribeUserToTopics(clientid: string; Stomp: IStompClient);
|
2013-10-30 00:48:23 +01:00
|
|
|
var
|
2013-11-09 14:22:11 +01:00
|
|
|
x, t: string;
|
2013-10-30 00:48:23 +01:00
|
|
|
topics: TArray<string>;
|
|
|
|
begin
|
|
|
|
x := Session['__subscriptions'];
|
|
|
|
topics := x.Split([';']);
|
|
|
|
for t in topics do
|
|
|
|
InternalSubscribeUserToTopic(clientid, t, Stomp);
|
|
|
|
end;
|
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
procedure TMVCBUSController.OnBeforeAction(Context: TWebContext; const AActionNAme: string;
|
|
|
|
var Handled: Boolean);
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
inherited;
|
2014-04-16 22:52:25 +02:00
|
|
|
if not StrToBool(Config['messaging']) then
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
Handled := true;
|
|
|
|
raise EMVCException.Create('Messaging extensions are not enabled');
|
|
|
|
end;
|
|
|
|
Handled := False;
|
|
|
|
end;
|
|
|
|
|
2015-04-01 17:01:23 +02:00
|
|
|
procedure TMVCBUSController.InternalSubscribeUserToTopic(clientid, topicname: string;
|
|
|
|
StompClient: IStompClient);
|
2013-10-30 00:48:23 +01:00
|
|
|
var
|
2015-04-01 17:01:23 +02:00
|
|
|
LDurSubHeader: string;
|
|
|
|
LHeaders: IStompHeaders;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
2015-04-01 17:01:23 +02:00
|
|
|
LHeaders := TStompHeaders.Create;
|
|
|
|
LDurSubHeader := GetUniqueDurableHeader(clientid, topicname);
|
|
|
|
LHeaders.Add(TStompHeaders.NewDurableSubscriptionHeader(LDurSubHeader));
|
|
|
|
|
|
|
|
if topicname.StartsWith('/topic') then
|
|
|
|
LHeaders.Add('id', clientid); //https://www.rabbitmq.com/stomp.html
|
|
|
|
|
|
|
|
StompClient.Subscribe(topicname, amClient, LHeaders);
|
|
|
|
LogE('SUBSCRIBE TO ' + clientid + '@' + topicname + ' dursubheader:' + LDurSubHeader);
|
2013-10-30 00:48:23 +01:00
|
|
|
AddTopicToUserSubscriptions(topicname);
|
|
|
|
end;
|
|
|
|
|
|
|
|
procedure TMVCBUSController.UnSubscribeFromTopic(CTX: TWebContext);
|
|
|
|
var
|
2013-11-09 14:22:11 +01:00
|
|
|
Stomp: IStompClient;
|
2013-10-30 00:48:23 +01:00
|
|
|
clientid: string;
|
2013-11-09 14:22:11 +01:00
|
|
|
thename: string;
|
|
|
|
s: string;
|
2013-10-30 00:48:23 +01:00
|
|
|
begin
|
|
|
|
clientid := GetClientID;
|
|
|
|
thename := CTX.Request.Params['name'].ToLower;
|
|
|
|
Stomp := GetNewStompClient(clientid);
|
2015-04-01 17:01:23 +02:00
|
|
|
s := '/queue/' + thename;
|
2013-10-30 00:48:23 +01:00
|
|
|
Stomp.Unsubscribe(s);
|
|
|
|
RemoveTopicFromUserSubscriptions(s);
|
|
|
|
Render(200, 'UnSubscription OK for ' + s);
|
|
|
|
end;
|
|
|
|
|
|
|
|
end.
|