StreamHandler.ts 12 KB

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