-
Notifications
You must be signed in to change notification settings - Fork 3
Expand file tree
/
Copy pathMongoMessage.cs
More file actions
132 lines (103 loc) · 3.74 KB
/
Copy pathMongoMessage.cs
File metadata and controls
132 lines (103 loc) · 3.74 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
using System;
using System.Collections.Generic;
using Microsoft.AspNet.SignalR.Messaging;
using MongoDB.Bson;
using MongoDB.Bson.Serialization.Attributes;
namespace Signalr.MongoDb
{
[CollectionName("messagebus")]
public class MongoMessage
{
/*[BsonRepresentation(BsonType.ObjectId)]*/
[BsonId]
public ObjectId Id { get; set; }
[BsonElement("i")]
public int StreamIndex { get; set; }
[BsonElement("v")]
public byte[] Value { get; set; }
[BsonDefaultValue(0)]
[BsonElement("t")]
public byte Status { get; set; }
[BsonIgnore]
public DateTime Created
{
get
{
//if we retrieved the MongoMessage - then it has a ObjectId - so read the date/time
// return Id != null ? ObjectId.Parse(Id).CreationTime : DateTime.MinValue;
return Id.CreationTime;
}
}
public static byte[] ToBytes(IList<Message> messages)
{
if (messages == null)
{
throw new ArgumentNullException("messages");
}
/* using (var ms = new MemoryStream())
{
var binaryWriter = new BinaryWriter(ms);
var scaleoutMessage = new ScaleoutMessage(messages);
var buffer = scaleoutMessage.ToBytes();
binaryWriter.Write(buffer.Length);
binaryWriter.Write(buffer);
return ms.ToArray();
}*/
var _message = new ScaleoutMessage(messages);
return _message.ToBytes();
}
/* public static MongoMessage FromBytes(byte[] data)
{
using (var stream = new MemoryStream(data))
{
var message = new MongoMessage();
/* // read message id from memory stream until SPACE character
var messageIdBuilder = new StringBuilder();
do
{
// it is safe to read digits as bytes because they encoded by single byte in UTF-8
int charCode = stream.ReadByte();
if (charCode == -1)
{
throw new EndOfStreamException();
}
char c = (char)charCode;
if (c == ' ')
{
message.Id = ulong.Parse(messageIdBuilder.ToString(), CultureInfo.InvariantCulture);
messageIdBuilder = null;
}
else
{
messageIdBuilder.Append(c);
}
}
while (messageIdBuilder != null);#1#
var binaryReader = new BinaryReader(stream);
int count = binaryReader.ReadInt32();
byte[] buffer = binaryReader.ReadBytes(count);
message = ScaleoutMessage.FromBytes(buffer);
return message;
}
}*/
public static ScaleoutMessage FromBytes(MongoMessage msg)
{
return ScaleoutMessage.FromBytes(msg.Value);
}
public MongoMessage(IList<Message> messages)
: this(ToBytes(messages))
{
}
public MongoMessage(int streamIndex, IList<Message> messages)
: this(ToBytes(messages))
{
StreamIndex = streamIndex;
}
public MongoMessage(byte[] value)
{
/*Id = ObjectId.GenerateNewId().ToString();ConnectionId = connectionId;
EventKey = eventKey;*/
Value = value;
}
}
}