From c6be194508ccc426c6e48517a68647360d19fa5f Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Mon, 10 Jun 2024 01:45:17 +0300 Subject: [PATCH] Implement experimental antietcd-based version of monitor --- mon/antietcd_adapter.js | 191 +++++++++++++++++++++++++++++++++ mon/etcd_adapter.js | 8 +- mon/mon-main.js | 2 +- mon/mon.js | 54 +++------- mon/package.json | 1 + mon/utils.js | 37 +++++++ mon/vitastor_persist_filter.js | 48 +++++++++ 7 files changed, 297 insertions(+), 44 deletions(-) create mode 100644 mon/antietcd_adapter.js create mode 100644 mon/utils.js create mode 100644 mon/vitastor_persist_filter.js diff --git a/mon/antietcd_adapter.js b/mon/antietcd_adapter.js new file mode 100644 index 00000000..c4dd2d6b --- /dev/null +++ b/mon/antietcd_adapter.js @@ -0,0 +1,191 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 (see README.md for details) + +const fs = require('fs'); + +const AntiEtcd = require('antietcd'); + +const vitastor_persist_filter = require('./vitastor_persist_filter.js'); +const { b64, local_ips } = require('./utils.js'); + +class AntiEtcdAdapter +{ + static async start_antietcd(config) + { + let antietcd; + if (config.use_antietcd) + { + let fileConfig = {}; + if (fs.existsSync(config.config_path||'/etc/vitastor/vitastor.conf')) + { + fileConfig = JSON.parse(fs.readFileSync(config.config_path||'/etc/vitastor/vitastor.conf', { encoding: 'utf-8' })); + } + let mergedConfig = { ...fileConfig, ...config }; + let cluster = mergedConfig.etcd_address; + if (!(cluster instanceof Array)) + cluster = cluster ? (''+(cluster||'')).split(/,+/) : []; + cluster = Object.keys(cluster.reduce((a, url) => + { + a[url.toLowerCase().replace(/^https?:\/\//, '').replace(/\/.*$/, '')] = true; + return a; + }, {})); + const cfg_port = mergedConfig.antietcd_port; + const is_local = local_ips(true).reduce((a, c) => { a[c] = true; return a; }, {}); + const selected = cluster.map(s => s.split(':', 2)).filter(ip => is_local[ip[0]] && (!cfg_port || ip[1] == cfg_port)); + if (selected.length > 1) + { + console.error('More than 1 etcd_address matches local IPs, please specify port'); + process.exit(1); + } + else if (selected.length == 1) + { + const antietcd_config = { + ip: selected[0][0], + port: selected[0][1], + data: mergedConfig.antietcd_data_file || ((mergedConfig.antietcd_data_dir || '/var/lib/vitastor') + '/mon_'+selected[0][1]+'.json.gz'), + persist_filter: vitastor_persist_filter(mergedConfig.etcd_prefix || '/vitastor'), + node_id: selected[0][0]+':'+selected[0][1], // node_id = ip:port + cluster: (cluster.length == 1 ? null : cluster), + cluster_key: (mergedConfig.etcd_prefix || '/vitastor'), + stale_read: 1, + }; + for (const key in config) + { + if (key.substr(0, 9) === 'antietcd_') + { + const noprefix = key.substr(9); + if (!(noprefix in antietcd_config) || noprefix == 'ip' || noprefix == 'cluster_key') + { + antietcd_config[noprefix] = config[key]; + } + } + } + antietcd = new AntiEtcd(antietcd_config); + await antietcd.start(); + } + else + { + console.log('Antietcd is enabled, but etcd_address does not contain local IPs, proceeding without it'); + } + } + return antietcd; + } + + constructor(mon, antietcd) + { + this.mon = mon; + this.antietcd = antietcd; + this.on_leader = []; + this.on_change = (st) => + { + if (st.state === 'leader') + { + for (const cb of this.on_leader) + { + cb(); + } + this.on_leader = []; + } + }; + this.antietcd.on('raftchange', this.on_change); + } + + parse_config(/*config*/) + { + } + + stop_watcher() + { + this.antietcd.off('raftchange', this.on_change); + const watch_id = this.watch_id; + if (watch_id) + { + this.watch_id = null; + this.antietcd.cancel_watch(watch_id).catch(console.error); + } + } + + async start_watcher() + { + if (this.watch_id) + { + await this.antietcd.cancel_watch(this.watch_id); + this.watch_id = null; + } + const watch_id = await this.antietcd.create_watch({ + key: b64(this.mon.config.etcd_prefix+'/'), + range_end: b64(this.mon.config.etcd_prefix+'0'), + start_revision: ''+this.mon.etcd_watch_revision, + watch_id: 1, + progress_notify: true, + }, (message) => + { + setImmediate(() => this.mon.on_message(message.result)); + }); + console.log('Successfully subscribed to antietcd revision '+this.antietcd.etctree.mod_revision); + this.watch_id = watch_id; + } + + async become_master() + { + if (!this.antietcd.raft) + { + console.log('Running in non-clustered mode'); + } + else + { + console.log('Waiting to become master'); + await new Promise(ok => this.on_leader.push(ok)); + } + const state = { ...this.mon.get_mon_state(), id: ''+this.mon.etcd_lease_id }; + await this.etcd_call('/kv/txn', { + success: [ { requestPut: { key: b64(this.mon.config.etcd_prefix+'/mon/master'), value: b64(JSON.stringify(state)), lease: ''+this.mon.etcd_lease_id } } ], + }, this.mon.config.etcd_start_timeout, 0); + if (this.antietcd.raft) + { + console.log('Became master'); + } + } + + async etcd_call(path, body, timeout, retries) + { + let retry = 0; + if (retries >= 0 && retries < 1) + { + retries = 1; + } + let prev = 0; + while (retries < 0 || retry < retries) + { + retry++; + if (this.mon.stopped) + { + throw new Error('Monitor instance is stopped'); + } + try + { + if (Date.now()-prev < timeout) + { + await new Promise(ok => setTimeout(ok, timeout-(Date.now()-prev))); + } + prev = Date.now(); + const res = await this.antietcd.api(path.replace(/^\/+/, '').replace(/\/+$/, '').replace(/\/+/g, '_'), body); + if (res.error) + { + console.error('Failed to query antietcd '+path+' (retry '+retry+'/'+retries+'): '+res.error); + } + else + { + return res; + } + } + catch (e) + { + console.error('Failed to query antietcd '+path+' (retry '+retry+'/'+retries+'): '+e.stack); + } + } + throw new Error('Failed to query antietcd ('+retries+' retries)'); + } +} + +module.exports = AntiEtcdAdapter; diff --git a/mon/etcd_adapter.js b/mon/etcd_adapter.js index 0c8d57ee..21d8dc72 100644 --- a/mon/etcd_adapter.js +++ b/mon/etcd_adapter.js @@ -3,6 +3,7 @@ const http = require('http'); const WebSocket = require('ws'); +const { b64, local_ips } = require('./utils.js'); const MON_STOPPED = 'Monitor instance is stopped'; @@ -23,7 +24,7 @@ class EtcdAdapter parse_etcd_addresses(addrs) { - const is_local_ip = this.mon.local_ips(true).reduce((a, c) => { a[c] = true; return a; }, {}); + const is_local_ip = local_ips(true).reduce((a, c) => { a[c] = true; return a; }, {}); this.etcd_local = []; this.etcd_urls = []; this.selected_etcd_url = null; @@ -348,9 +349,4 @@ function POST(url, body, timeout) }); } -function b64(str) -{ - return Buffer.from(str).toString('base64'); -} - module.exports = EtcdAdapter; diff --git a/mon/mon-main.js b/mon/mon-main.js index 1e15850c..79af1362 100755 --- a/mon/mon-main.js +++ b/mon/mon-main.js @@ -23,4 +23,4 @@ for (let i = 2; i < process.argv.length; i++) } } -Mon.run_forever(options); +Mon.run_forever(options).catch(console.error); diff --git a/mon/mon.js b/mon/mon.js index 7510c89c..6571b2b0 100644 --- a/mon/mon.js +++ b/mon/mon.js @@ -5,6 +5,7 @@ const { URL } = require('url'); const fs = require('fs'); const crypto = require('crypto'); const os = require('os'); +const AntiEtcdAdapter = require('./antietcd_adapter.js'); const EtcdAdapter = require('./etcd_adapter.js'); const { create_http_server } = require('./http_server.js'); const { export_prometheus_metrics } = require('./prometheus.js'); @@ -14,17 +15,23 @@ const { sum_op_stats, sum_object_counts, sum_inode_stats, serialize_bigints } = const stableStringify = require('./stable-stringify.js'); const { scale_pg_history } = require('./pg_utils.js'); const { get_osd_tree } = require('./osd_tree.js'); +const { b64, de64, local_ips } = require('./utils.js'); const { recheck_primary, save_new_pgs_txn, generate_pool_pgs } = require('./pg_gen.js'); class Mon { - static run_forever(config) + static async run_forever(config) { + let antietcd = await AntiEtcdAdapter.start_antietcd(config); let mon; const run = () => { console.log('Starting Monitor'); const my_mon = new Mon(config); + my_mon.etcd = antietcd + ? new AntiEtcdAdapter(my_mon, antietcd) + : new EtcdAdapter(my_mon); + my_mon.etcd.parse_config(my_mon.config); mon = my_mon; my_mon.on_die = () => { @@ -61,8 +68,6 @@ class Mon this.state = JSON.parse(JSON.stringify(etcd_tree)); this.prev_stats = { osd_stats: {}, osd_diff: {} }; this.recheck_pgs_active = false; - this.etcd = new EtcdAdapter(this); - this.etcd.parse_config(this.config); this.watcher_active = false; if (this.config.enable_prometheus || !('enable_prometheus' in this.config)) { @@ -177,8 +182,8 @@ class Mon this.etcd_watch_revision = BigInt(msg.header.revision)+BigInt(1); for (const e of msg.events||[]) { - this.parse_kv(e.kv); - const key = e.kv.key.substr(this.config.etcd_prefix.length); + const kv = this.parse_kv(e.kv); + const key = kv.key.substr(this.config.etcd_prefix.length); if (key.substr(0, 11) == '/osd/state/') { stats_changed = true; @@ -198,7 +203,7 @@ class Mon } if (this.config.verbose) { - console.log(JSON.stringify(e)); + console.log(JSON.stringify({ ...e, kv: kv || undefined })); } } if (pg_states_changed) @@ -282,7 +287,7 @@ class Mon get_mon_state() { - return { ip: this.local_ips(), hostname: os.hostname() }; + return { ip: local_ips(), hostname: os.hostname() }; } async get_lease() @@ -720,15 +725,16 @@ class Mon { if (!kv || !kv.key) { - return; + return kv; } + kv = { ...kv }; kv.key = de64(kv.key); kv.value = kv.value ? de64(kv.value) : null; let key = kv.key.substr(this.config.etcd_prefix.length+1); if (!etcd_allow.exec(key)) { console.log('Bad key in etcd: '+kv.key+' = '+kv.value); - return; + return kv; } try { @@ -737,7 +743,7 @@ class Mon catch (e) { console.log('Bad value in etcd: '+kv.key+' = '+kv.value); - return; + return kv; } let key_parts = key.split('/'); let cur = this.state; @@ -787,6 +793,7 @@ class Mon !this.state.osd.stats[osd_num] ? 0 : this.state.osd.stats[osd_num].time+this.config.osd_out_time ); } + return kv; } _die(err) @@ -796,33 +803,6 @@ class Mon this.on_stop().catch(console.error); this.on_die(); } - - local_ips(all) - { - const ips = []; - const ifaces = os.networkInterfaces(); - for (const ifname in ifaces) - { - for (const iface of ifaces[ifname]) - { - if (iface.family == 'IPv4' && !iface.internal || all) - { - ips.push(iface.address); - } - } - } - return ips; - } -} - -function b64(str) -{ - return Buffer.from(str).toString('base64'); -} - -function de64(str) -{ - return Buffer.from(str, 'base64').toString(); } function sha1hex(str) diff --git a/mon/package.json b/mon/package.json index 7b6aee5e..7861053f 100644 --- a/mon/package.json +++ b/mon/package.json @@ -9,6 +9,7 @@ "author": "Vitaliy Filippov", "license": "UNLICENSED", "dependencies": { + "antietcd": "^1.0.5", "sprintf-js": "^1.1.2", "ws": "^7.2.5" }, diff --git a/mon/utils.js b/mon/utils.js new file mode 100644 index 00000000..8a668532 --- /dev/null +++ b/mon/utils.js @@ -0,0 +1,37 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 (see README.md for details) + +const os = require('os'); + +function local_ips(all) +{ + const ips = []; + const ifaces = os.networkInterfaces(); + for (const ifname in ifaces) + { + for (const iface of ifaces[ifname]) + { + if (iface.family == 'IPv4' && !iface.internal || all) + { + ips.push(iface.address); + } + } + } + return ips; +} + +function b64(str) +{ + return Buffer.from(str).toString('base64'); +} + +function de64(str) +{ + return Buffer.from(str, 'base64').toString(); +} + +module.exports = { + b64, + de64, + local_ips, +}; diff --git a/mon/vitastor_persist_filter.js b/mon/vitastor_persist_filter.js new file mode 100644 index 00000000..810bccaf --- /dev/null +++ b/mon/vitastor_persist_filter.js @@ -0,0 +1,48 @@ +// AntiEtcd persistence filter for Vitastor +// (c) Vitaliy Filippov, 2024 +// License: Mozilla Public License 2.0 or Vitastor Network Public License 1.1 + +function vitastor_persist_filter(cfg) +{ + const prefix = cfg.vitastor_prefix || '/vitastor'; + return (key, value) => + { + if (key.substr(0, prefix.length+'/osd/stats/'.length) == prefix+'/osd/stats/') + { + if (value) + { + try + { + value = JSON.parse(value); + value = JSON.stringify({ + bitmap_granularity: value.bitmap_granularity || undefined, + data_block_size: value.data_block_size || undefined, + host: value.host || undefined, + immediate_commit: value.immediate_commit || undefined, + }); + } + catch (e) + { + console.error('invalid JSON in '+key+' = '+value+': '+e); + value = {}; + } + } + else + { + value = undefined; + } + return value; + } + else if (key.substr(0, prefix.length+'/osd/'.length) == prefix+'/osd/' || + key.substr(0, prefix.length+'/inode/stats/'.length) == prefix+'/inode/stats/' || + key.substr(0, prefix.length+'/pg/stats/'.length) == prefix+'/pg/stats/' || + key.substr(0, prefix.length+'/pool/stats/'.length) == prefix+'/pool/stats/' || + key == prefix+'/stats') + { + return undefined; + } + return value; + }; +} + +module.exports = vitastor_persist_filter;