StreamHandler.ts 8.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425
  1. import defaults from 'lodash/defaults';
  2. import {
  3. DataQueryRequest,
  4. FieldType,
  5. SeriesData,
  6. DataQueryResponse,
  7. DataQueryError,
  8. DataStreamObserver,
  9. DataStreamState,
  10. LoadingState,
  11. LogLevel,
  12. } from '@grafana/ui';
  13. import { TestDataQuery, StreamingQuery } from './types';
  14. export const defaultQuery: StreamingQuery = {
  15. type: 'signal',
  16. speed: 250, // ms
  17. spread: 3.5,
  18. noise: 2.2,
  19. };
  20. type StreamWorkers = {
  21. [key: string]: StreamWorker;
  22. };
  23. export class StreamHandler {
  24. workers: StreamWorkers = {};
  25. process(req: DataQueryRequest<TestDataQuery>, observer: DataStreamObserver): DataQueryResponse | undefined {
  26. let resp: DataQueryResponse;
  27. for (const query of req.targets) {
  28. if ('streaming_client' !== query.scenarioId) {
  29. continue;
  30. }
  31. if (!resp) {
  32. resp = { data: [] };
  33. }
  34. // set stream option defaults
  35. query.stream = defaults(query.stream, defaultQuery);
  36. // create stream key
  37. const key = req.dashboardId + '/' + req.panelId + '/' + query.refId;
  38. if (this.workers[key]) {
  39. const existing = this.workers[key];
  40. if (existing.update(query, req)) {
  41. continue;
  42. }
  43. existing.unsubscribe();
  44. delete this.workers[key];
  45. }
  46. const type = query.stream.type;
  47. if (type === 'signal') {
  48. this.workers[key] = new SignalWorker(key, query, req, observer);
  49. } else if (type === 'logs') {
  50. this.workers[key] = new LogsWorker(key, query, req, observer);
  51. } else {
  52. throw {
  53. message: 'Unknown Stream type: ' + type,
  54. refId: query.refId,
  55. } as DataQueryError;
  56. }
  57. }
  58. return resp;
  59. }
  60. }
  61. /**
  62. * Manages a single stream request
  63. */
  64. export class StreamWorker {
  65. query: StreamingQuery;
  66. stream: DataStreamState;
  67. observer: DataStreamObserver;
  68. last = -1;
  69. timeoutId = 0;
  70. constructor(key: string, query: TestDataQuery, request: DataQueryRequest, observer: DataStreamObserver) {
  71. this.stream = {
  72. key,
  73. state: LoadingState.Streaming,
  74. request,
  75. unsubscribe: this.unsubscribe,
  76. };
  77. this.query = query.stream;
  78. this.last = Date.now();
  79. this.observer = observer;
  80. console.log('Creating Test Stream: ', this);
  81. }
  82. unsubscribe = () => {
  83. this.observer = null;
  84. if (this.timeoutId) {
  85. clearTimeout(this.timeoutId);
  86. this.timeoutId = 0;
  87. }
  88. };
  89. update(query: TestDataQuery, request: DataQueryRequest): boolean {
  90. // Check if stream has been unsubscribed or query changed type
  91. if (this.observer === null || this.query.type !== query.stream.type) {
  92. return false;
  93. }
  94. this.query = query.stream;
  95. this.stream.request = request; // OK?
  96. console.log('Reuse Test Stream: ', this);
  97. return true;
  98. }
  99. appendRows(append: any[][]) {
  100. // Trim the maximum row count
  101. const { query, stream } = this;
  102. const maxRows = query.buffer ? query.buffer : stream.request.maxDataPoints;
  103. // Edit the first series
  104. const series = stream.series[0];
  105. let rows = series.rows.concat(append);
  106. const extra = maxRows - rows.length;
  107. if (extra < 0) {
  108. rows = rows.slice(extra * -1);
  109. }
  110. series.rows = rows;
  111. // Tell the event about only the rows that changed (it may want to process them)
  112. stream.delta = [{ ...series, rows: append }];
  113. // Broadcast the changes
  114. if (this.observer) {
  115. this.observer(stream);
  116. } else {
  117. console.log('StreamWorker working without any observer');
  118. }
  119. this.last = Date.now();
  120. }
  121. }
  122. export class SignalWorker extends StreamWorker {
  123. value: number;
  124. constructor(key: string, query: TestDataQuery, request: DataQueryRequest, observer: DataStreamObserver) {
  125. super(key, query, request, observer);
  126. setTimeout(() => {
  127. this.stream.series = [this.initBuffer(query.refId)];
  128. this.looper();
  129. }, 10);
  130. }
  131. nextRow = (time: number) => {
  132. const { spread, noise } = this.query;
  133. this.value += (Math.random() - 0.5) * spread;
  134. return [
  135. time,
  136. this.value, // Value
  137. this.value - Math.random() * noise, // MIN
  138. this.value + Math.random() * noise, // MAX
  139. ];
  140. };
  141. initBuffer(refId: string): SeriesData {
  142. const { speed, buffer } = this.query;
  143. const data = {
  144. fields: [
  145. { name: 'Time', type: FieldType.time },
  146. { name: 'Value', type: FieldType.number },
  147. { name: 'Min', type: FieldType.number },
  148. { name: 'Max', type: FieldType.number },
  149. ],
  150. rows: [],
  151. refId,
  152. name: 'Signal ' + refId,
  153. } as SeriesData;
  154. const request = this.stream.request;
  155. this.value = Math.random() * 100;
  156. const maxRows = buffer ? buffer : request.maxDataPoints;
  157. let time = Date.now() - maxRows * speed;
  158. for (let i = 0; i < maxRows; i++) {
  159. data.rows.push(this.nextRow(time));
  160. time += speed;
  161. }
  162. return data;
  163. }
  164. looper = () => {
  165. if (!this.observer) {
  166. const request = this.stream.request;
  167. const elapsed = request.startTime - Date.now();
  168. if (elapsed > 1000) {
  169. console.log('Stop looping');
  170. return;
  171. }
  172. }
  173. // Make sure it has a minimum speed
  174. const { query } = this;
  175. if (query.speed < 5) {
  176. query.speed = 5;
  177. }
  178. this.appendRows([this.nextRow(Date.now())]);
  179. this.timeoutId = window.setTimeout(this.looper, query.speed);
  180. };
  181. }
  182. export class LogsWorker extends StreamWorker {
  183. index = 0;
  184. constructor(key: string, query: TestDataQuery, request: DataQueryRequest, observer: DataStreamObserver) {
  185. super(key, query, request, observer);
  186. window.setTimeout(() => {
  187. this.stream.series = [this.initBuffer(query.refId)];
  188. this.looper();
  189. }, 10);
  190. }
  191. getRandomLogLevel(): LogLevel {
  192. const v = Math.random();
  193. if (v > 0.9) {
  194. return LogLevel.critical;
  195. }
  196. if (v > 0.8) {
  197. return LogLevel.error;
  198. }
  199. if (v > 0.7) {
  200. return LogLevel.warning;
  201. }
  202. if (v > 0.4) {
  203. return LogLevel.info;
  204. }
  205. if (v > 0.3) {
  206. return LogLevel.debug;
  207. }
  208. if (v > 0.1) {
  209. return LogLevel.trace;
  210. }
  211. return LogLevel.unknown;
  212. }
  213. getNextWord() {
  214. this.index = (this.index + Math.floor(Math.random() * 5)) % words.length;
  215. return words[this.index];
  216. }
  217. getRandomLine(length = 60) {
  218. let line = this.getNextWord();
  219. while (line.length < length) {
  220. line += ' ' + this.getNextWord();
  221. }
  222. return line;
  223. }
  224. nextRow = (time: number) => {
  225. return [time, '[' + this.getRandomLogLevel() + '] ' + this.getRandomLine()];
  226. };
  227. initBuffer(refId: string): SeriesData {
  228. const { speed, buffer } = this.query;
  229. const data = {
  230. fields: [{ name: 'Time', type: FieldType.time }, { name: 'Line', type: FieldType.string }],
  231. rows: [],
  232. refId,
  233. name: 'Logs ' + refId,
  234. } as SeriesData;
  235. const request = this.stream.request;
  236. const maxRows = buffer ? buffer : request.maxDataPoints;
  237. let time = Date.now() - maxRows * speed;
  238. for (let i = 0; i < maxRows; i++) {
  239. data.rows.push(this.nextRow(time));
  240. time += speed;
  241. }
  242. return data;
  243. }
  244. looper = () => {
  245. if (!this.observer) {
  246. const request = this.stream.request;
  247. const elapsed = request.startTime - Date.now();
  248. if (elapsed > 1000) {
  249. console.log('Stop looping');
  250. return;
  251. }
  252. }
  253. // Make sure it has a minimum speed
  254. const { query } = this;
  255. if (query.speed < 5) {
  256. query.speed = 5;
  257. }
  258. const variance = query.speed * 0.2 * (Math.random() - 0.5); // +-10%
  259. this.appendRows([this.nextRow(Date.now())]);
  260. this.timeoutId = window.setTimeout(this.looper, query.speed + variance);
  261. };
  262. }
  263. const words = [
  264. 'At',
  265. 'vero',
  266. 'eos',
  267. 'et',
  268. 'accusamus',
  269. 'et',
  270. 'iusto',
  271. 'odio',
  272. 'dignissimos',
  273. 'ducimus',
  274. 'qui',
  275. 'blanditiis',
  276. 'praesentium',
  277. 'voluptatum',
  278. 'deleniti',
  279. 'atque',
  280. 'corrupti',
  281. 'quos',
  282. 'dolores',
  283. 'et',
  284. 'quas',
  285. 'molestias',
  286. 'excepturi',
  287. 'sint',
  288. 'occaecati',
  289. 'cupiditate',
  290. 'non',
  291. 'provident',
  292. 'similique',
  293. 'sunt',
  294. 'in',
  295. 'culpa',
  296. 'qui',
  297. 'officia',
  298. 'deserunt',
  299. 'mollitia',
  300. 'animi',
  301. 'id',
  302. 'est',
  303. 'laborum',
  304. 'et',
  305. 'dolorum',
  306. 'fuga',
  307. 'Et',
  308. 'harum',
  309. 'quidem',
  310. 'rerum',
  311. 'facilis',
  312. 'est',
  313. 'et',
  314. 'expedita',
  315. 'distinctio',
  316. 'Nam',
  317. 'libero',
  318. 'tempore',
  319. 'cum',
  320. 'soluta',
  321. 'nobis',
  322. 'est',
  323. 'eligendi',
  324. 'optio',
  325. 'cumque',
  326. 'nihil',
  327. 'impedit',
  328. 'quo',
  329. 'minus',
  330. 'id',
  331. 'quod',
  332. 'maxime',
  333. 'placeat',
  334. 'facere',
  335. 'possimus',
  336. 'omnis',
  337. 'voluptas',
  338. 'assumenda',
  339. 'est',
  340. 'omnis',
  341. 'dolor',
  342. 'repellendus',
  343. 'Temporibus',
  344. 'autem',
  345. 'quibusdam',
  346. 'et',
  347. 'aut',
  348. 'officiis',
  349. 'debitis',
  350. 'aut',
  351. 'rerum',
  352. 'necessitatibus',
  353. 'saepe',
  354. 'eveniet',
  355. 'ut',
  356. 'et',
  357. 'voluptates',
  358. 'repudiandae',
  359. 'sint',
  360. 'et',
  361. 'molestiae',
  362. 'non',
  363. 'recusandae',
  364. 'Itaque',
  365. 'earum',
  366. 'rerum',
  367. 'hic',
  368. 'tenetur',
  369. 'a',
  370. 'sapiente',
  371. 'delectus',
  372. 'ut',
  373. 'aut',
  374. 'reiciendis',
  375. 'voluptatibus',
  376. 'maiores',
  377. 'alias',
  378. 'consequatur',
  379. 'aut',
  380. 'perferendis',
  381. 'doloribus',
  382. 'asperiores',
  383. 'repellat',
  384. ];