StreamHandler.ts 11 KB

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