ADD Some headers names constants

This commit is contained in:
daniele.teti@gmail.com 2013-10-15 08:14:36 +00:00
parent 6688fc17a3
commit 41d42ff9b9
2 changed files with 56 additions and 77 deletions

View File

@ -1,6 +1,6 @@
// Stomp Client for Embarcadero Delphi & FreePascal
// Tested With ApacheMQ 5.2/5.3, Apache Apollo 1.2/1.6
// Copyright (c) 2009-2013 Daniele Teti
// Tested With ApacheMQ 5.2/5.3, Apache Apollo 1.2
// Copyright (c) 2009-2012 Daniele Teti
//
// Contributors:
// Daniel Gaspary: dgaspary@gmail.com
@ -102,7 +102,7 @@ type
procedure Receipt(const ReceiptID: string);
procedure Connect(Host: string = '127.0.0.1'; Port: Integer = DEFAULT_STOMP_PORT;
ClientID: string = ''; AcceptVersion: TStompAcceptProtocol = STOMP_Version_1_0);
procedure Disconnect(WithoutSendingFrame: boolean = false);
procedure Disconnect;
procedure Subscribe(QueueOrTopicName: string; Ack: TAckMode = amAuto;
Headers: IStompHeaders = nil);
procedure Unsubscribe(Queue: string);
@ -161,7 +161,7 @@ begin
Frame.SetCommand('ABORT');
Frame.GetHeaders.Add('transaction', TransactionIdentifier);
SendFrame(Frame);
FInTransaction := false;
FInTransaction := False;
FTransactions.Delete(FTransactions.IndexOf(TransactionIdentifier));
end
else
@ -175,7 +175,7 @@ var
begin
Frame := TStompFrame.Create;
Frame.SetCommand('ACK');
Frame.GetHeaders.Add('message-id', MessageID);
Frame.GetHeaders.Add(TStompHeaders.MESSAGE_ID, MessageID);
if TransactionIdentifier <> '' then
Frame.GetHeaders.Add('transaction', TransactionIdentifier);
SendFrame(Frame);
@ -200,6 +200,21 @@ begin
[TransactionIdentifier]);
end;
// procedure TStompClient.CheckReceipt(Frame: TStompFrame);
// var
// ReceiptID: string;
// begin
// if FEnableReceipts then
// begin
// ReceiptID := inttostr(GetTickCount);
// Frame.GetHeaders.Add('receipt', ReceiptID);
// SendFrame(Frame);
// Receipt(ReceiptID);
// end
// else
// SendFrame(Frame);
// end;
procedure TStompClient.CommitTransaction(const TransactionIdentifier: string);
var
Frame: IStompFrame;
@ -210,7 +225,7 @@ begin
Frame.SetCommand('COMMIT');
Frame.GetHeaders.Add('transaction', TransactionIdentifier);
SendFrame(Frame);
FInTransaction := false;
FInTransaction := False;
FTransactions.Delete(FTransactions.IndexOf(TransactionIdentifier));
end
else
@ -228,7 +243,7 @@ begin
{$IFDEF USESYNAPSE}
FSynapseConnected := false;
FSynapseConnected := False;
FSynapseTCP.Connect(Host, intToStr(Port));
FSynapseConnected := True;
@ -291,7 +306,7 @@ end;
constructor TStompClient.Create;
begin
inherited;
FInTransaction := false;
FInTransaction := False;
FSession := '';
FUserName := 'guest';
FPassword := 'guest';
@ -318,39 +333,28 @@ end;
destructor TStompClient.Destroy;
begin
try
Disconnect(false);
except
end;
Disconnect;
DeInit;
inherited;
end;
procedure TStompClient.Disconnect(WithoutSendingFrame: boolean);
procedure TStompClient.Disconnect;
var
Frame: IStompFrame;
begin
if Connected then
begin
if not WithoutSendingFrame then
begin
try
Frame := TStompFrame.Create;
Frame.SetCommand('DISCONNECT');
SendFrame(Frame);
except
// socket could be already dead
end;
end;
Frame := TStompFrame.Create;
Frame.SetCommand('DISCONNECT');
SendFrame(Frame);
{$IFDEF USESYNAPSE}
FSynapseTCP.CloseSocket;
FSynapseConnected := false;
FSynapseConnected := False;
{$ELSE}
FTCP.Socket.Close;
FTCP.Disconnect;
{$ENDIF}
@ -404,7 +408,7 @@ begin
if (Reason = HR_Error) and (FSynapseTCP.LastError <> WSAETIMEDOUT)
then
begin
FSynapseConnected := false;
FSynapseConnected := False;
end;
end;
@ -467,7 +471,7 @@ function TStompClient.Receive(ATimeout: Integer): IStompFrame;
s : string;
tout: boolean;
begin
tout := false;
tout := False;
Result := nil;
try
try
@ -526,14 +530,14 @@ function TStompClient.Receive(ATimeout: Integer): IStompFrame;
begin
// UTF8Encoding := TEncoding.UTF8;
UTF8Encoding := IndyTextEncoding_UTF8();
tout := false;
tout := False;
Result := nil;
try
sb := TStringBuilder.Create(1024 * 4);
try
FTCP.ReadTimeout := ATimeout;
try
FirstValidChar := false;
FirstValidChar := False;
FTCP.Socket.CheckForDataOnSource(1);
while True do
begin
@ -566,7 +570,7 @@ function TStompClient.Receive(ATimeout: Integer): IStompFrame;
begin
Result := StompUtils.CreateFrame(sb.toString + CHAR0);
if Result.GetCommand = 'ERROR' then
raise EStompDisconnectionError.Create(Result.GetHeaders.Value('message'));
raise EStomp.Create(Result.GetHeaders.Value('message'));
end;
finally
sb.Free;
@ -583,27 +587,16 @@ function TStompClient.Receive(ATimeout: Integer): IStompFrame;
begin
try
{$IFDEF USESYNAPSE}
{$IFDEF USESYNAPSE}
Result := InternalReceiveSynapse(ATimeout);
Result := InternalReceiveSynapse(ATimeout);
{$ELSE}
{$ELSE}
Result := InternalReceiveINDY(ATimeout);
Result := InternalReceiveINDY(ATimeout);
{$ENDIF}
except
on E: EStompDisconnectionError do
begin
Disconnect;
raise;
end;
on E: Exception do
raise;
end
{$ENDIF}
end;
@ -640,8 +633,6 @@ end;
procedure TStompClient.SendFrame(AFrame: IStompFrame);
begin
if not Connected then
raise EStomp.Create('StompClient not connected');
{$IFDEF USESYNAPSE}

View File

@ -1,6 +1,6 @@
// Stomp Client for Embarcadero Delphi & FreePascal
// Tested With ApacheMQ 5.2/5.3
// Copyright (c) 2009-2013 Daniele Teti
// Copyright (c) 2009-2009 Daniele Teti
//
// Contributors:
// Daniel Gaspary: dgaspary@gmail.com
@ -34,10 +34,6 @@ type
EStomp = class(Exception)
end;
EStompDisconnectionError = class(EStomp)
end;
TKeyValue = record
Key: string;
Value: string;
@ -76,7 +72,7 @@ type
procedure Receipt(const ReceiptID: string);
procedure Connect(Host: string = '127.0.0.1'; Port: Integer = 61613; ClientID: string = '';
AcceptVersion: TStompAcceptProtocol = STOMP_Version_1_0);
procedure Disconnect(WithoutSendingFrame: Boolean = false);
procedure Disconnect;
procedure Subscribe(QueueOrTopicName: string; Ack: TAckMode = amAuto;
Headers: IStompHeaders = nil);
procedure Unsubscribe(Queue: string);
@ -110,7 +106,12 @@ type
class function NewDurableSubscriptionHeader(const SubscriptionName: string): TKeyValue;
class function NewPersistentHeader(const Value: Boolean): TKeyValue;
class function NewReplyToHeader(const DestinationName: string): TKeyValue;
/// /////////////////////////////////////////////7
const
MESSAGE_ID: string = 'message-id';
TRANSACTION: string = 'transaction';
/// /
function Add(Key, Value: string): IStompHeaders; overload;
function Add(HeaderItem: TKeyValue): IStompHeaders; overload;
function Value(Key: string): string;
@ -533,26 +534,16 @@ var
frame : IStompFrame;
StopListen: Boolean;
begin
StopListen := false;
while (not terminated) and (not StopListen) do
StopListen := False;
while not terminated do
begin
try
if FStompClient.Receive(frame, 2000) then
if FStompClient.Receive(frame, 2000) then
begin
FStompClientListener.OnMessage(FStompClient, frame, StopListen);
if StopListen then
begin
try
FStompClientListener.OnMessage(FStompClient, frame, StopListen);
except
end;
end;
except
on e: EStompDisconnectionError do
begin
try
StopListen := true;
FStompClientListener.OnStopListen(FStompClient);
except
end;
FStompClientListener.OnStopListen(FStompClient);
StopListening;
end;
end;
end;
@ -566,11 +557,8 @@ end;
procedure TStompClientListener.StopListening;
begin
Terminate;
WaitFor;
try
FStompClientListener.OnStopListen(FStompClient);
except
end;
Free;
// WaitFor;
end;
function TStompClientListener._AddRef: Integer;