oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/bot_config.rs
8.1 KiB, 68 runs
created by r1870400018:739, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | use crate::{ |
| 2 | prelude::*, |
| 3 | base::index::ZoneInd, |
| 4 | bots::base::bot_deps::*, |
| 5 | comm::{ |
| 6 | channels::{ |
| 7 | BotChannels, |
| 8 | ChannelPool, |
| 9 | }, |
| 10 | response::Responder, |
| 11 | }, |
| 12 | file::zdir::{ |
| 13 | ZoneDir, |
| 14 | ZoneDirStr, |
| 15 | }, |
| 16 | format_zone_dir, |
| 17 | }; |
| 18 | |
| 19 | use oxedyne_fe2o3_core::channels::Recv; |
| 20 | use oxedyne_fe2o3_jdat::{ |
| 21 | prelude::*, |
| 22 | id::NumIdDat, |
| 23 | }; |
| 24 | |
| 25 | use std::{ |
| 26 | collections::BTreeMap, |
| 27 | path::PathBuf, |
| 28 | str::FromStr, |
| 29 | sync::Arc, |
| 30 | time::{ |
| 31 | Duration, |
| 32 | SystemTime, |
| 33 | }, |
| 34 | }; |
| 35 | |
| 36 | /// Watches the config file. |
| 37 | pub struct ConfigBot< |
| 38 | const UIDL: usize, |
| 39 | UID: NumIdDat<UIDL>, |
| 40 | ENC: Encrypter, |
| 41 | KH: Hasher, |
| 42 | PR: Hasher, |
| 43 | CS: Checksummer, |
| 44 | >{ |
| 45 | // Bot |
| 46 | sem: Semaphore, |
| 47 | errc: Arc<Mutex<usize>>, |
| 48 | log_stream_id: String, |
| 49 | // Comms |
| 50 | chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 51 | // API |
| 52 | api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>, |
| 53 | // State |
| 54 | inited: bool, |
| 55 | last: SystemTime, // Last check of file |
| 56 | sleep: Duration, |
| 57 | } |
| 58 | |
| 59 | impl< |
| 60 | const UIDL: usize, |
| 61 | UID: NumIdDat<UIDL> + 'static, |
| 62 | ENC: Encrypter + 'static, |
| 63 | KH: Hasher + 'static, |
| 64 | PR: Hasher, |
| 65 | CS: Checksummer, |
| 66 | > |
| 67 | Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for ConfigBot<UIDL, UID, ENC, KH, PR, CS> |
| 68 | { |
| 69 | bot_methods!(); |
| 70 | |
| 71 | fn go(&mut self) { |
| 72 | |
| 73 | sync_log::set_stream(self.log_stream_id()); |
| 74 | |
| 75 | if self.no_init() { return; } |
| 76 | info!(sync_log::stream(), "{}: Checking config file for changes every {:?}.", self.ozid(), self.sleep); |
| 77 | self.now_listening(); |
| 78 | loop { |
| 79 | if self.listen().must_end() { break; } |
| 80 | //let result = self.check_file(); |
| 81 | //self.result(result); |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | fn listen(&mut self) -> LoopBreak { |
| 86 | match self.chan_in().recv_timeout(self.sleep) { |
| 87 | Recv::Empty => (), |
| 88 | Recv::Result(Err(e)) => self.err_cannot_receive(err!(e, |
| 89 | "{}: Waiting for message.", self.ozid(); |
| 90 | IO, Channel)), |
| 91 | Recv::Result(Ok(OzoneMsg::ZoneInitTrigger(resp))) => { |
| 92 | // Each zone asked answers for itself. The one this failed on, and any after it, |
| 93 | // never will, so the supervisor waiting on them is told why here. |
| 94 | if let Err(e) = self.zone_init(&resp) { |
| 95 | self.error(e.clone()); |
| 96 | self.respond(Err(e), &resp); |
| 97 | } |
| 98 | }, |
| 99 | // ... listen here for custom messages |
| 100 | Recv::Result(Ok(msg)) => return self.listen_more(msg), |
| 101 | } |
| 102 | LoopBreak(false) |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | impl< |
| 107 | const UIDL: usize, |
| 108 | UID: NumIdDat<UIDL> + 'static, |
| 109 | ENC: Encrypter + 'static, |
| 110 | KH: Hasher + 'static, |
| 111 | PR: Hasher, |
| 112 | CS: Checksummer, |
| 113 | > |
| 114 | OzoneBot<UIDL, UID, ENC, KH, PR, CS> for ConfigBot<UIDL, UID, ENC, KH, PR, CS> |
| 115 | { |
| 116 | ozonebot_methods!(); |
| 117 | } |
| 118 | |
| 119 | impl< |
| 120 | const UIDL: usize, |
| 121 | UID: NumIdDat<UIDL> + 'static, |
| 122 | ENC: Encrypter + 'static, |
| 123 | KH: Hasher + 'static, |
| 124 | PR: Hasher, |
| 125 | CS: Checksummer, |
| 126 | > |
| 127 | ConfigBot<UIDL, UID, ENC, KH, PR, CS> |
| 128 | { |
| 129 | pub fn new( |
| 130 | args: BotInitArgs<UIDL, UID, ENC, KH, PR, CS>, |
| 131 | ) |
| 132 | -> Self |
| 133 | { |
| 134 | Self { |
| 135 | // Bot |
| 136 | sem: args.sem, |
| 137 | errc: Arc::new(Mutex::new(0)), |
| 138 | log_stream_id: args.log_stream_id, |
| 139 | // Comms |
| 140 | chan_in: args.chan_in, |
| 141 | // API |
| 142 | api: args.api, |
| 143 | // State |
| 144 | inited: false, |
| 145 | last: SystemTime::now(), |
| 146 | sleep: Duration::from_secs(constant::CONFIGWATCHER_CHECK_INTERVAL_SECS), |
| 147 | } |
| 148 | } |
| 149 | |
| 150 | fn zbots(&self) -> &ChannelPool<UIDL, UID, ENC, KH> { &self.chans().all_zbots() } |
| 151 | |
| 152 | pub fn zone_init(&mut self, resp: &Responder<UIDL, UID, ENC, KH>) -> Outcome<()> { |
| 153 | let zcfg = self.cfg().zone_config(); |
| 154 | let default = ZoneDir { |
| 155 | dir: self.cfg().zone_root(self.db_root()), |
| 156 | max_size: constant::DEFAULT_MAX_ZONE_DIR_BYTES, |
| 157 | }; |
| 158 | let zdirs = res!(self.process_zone_overrides()); |
| 159 | let nz = self.cfg().num_zones; |
| 160 | for z in 0..nz { |
| 161 | let zind = ZoneInd::new(z); |
| 162 | let mut zdir = match zdirs.get(&z) { |
| 163 | None => default.clone(), |
| 164 | Some(zd) => zd.clone(), |
| 165 | }; |
| 166 | zdir.dir.push(fmt!(format_zone_dir!(), z+1)); |
| 167 | if !zdir.dir.exists() { |
| 168 | res!(std::fs::create_dir(&zdir.dir)); |
| 169 | } |
| 170 | let bot = res!(self.zbots().get_bot(z as usize)); |
| 171 | if let Err(e) = bot.send(OzoneMsg::ZoneInit(zdir, zcfg.clone(), resp.clone())) { |
| 172 | return Err(err!(e, |
| 173 | "{}: Sending init data to zone {:?}", self.ozid(), zind; |
| 174 | Init, IO, Channel)); |
| 175 | } |
| 176 | } |
| 177 | info!(sync_log::stream(), "{}: Initialisation request sent to all {} zones.", self.ozid(), nz); |
| 178 | Ok(()) |
| 179 | } |
| 180 | |
| 181 | /// The configuration may contain a map `BTreeMap<Dat, Dat>` of zones settings that |
| 182 | /// override the default. This method turns the map into a map `BTreeMap<u16, ZoneDir>`. |
| 183 | pub fn process_zone_overrides(&self) -> Outcome<BTreeMap<u16, ZoneDir>> { |
| 184 | let zones = self.cfg().zone_overrides(); |
| 185 | if self.cfg().num_zones() < zones.len() { |
| 186 | return Err(err!( |
| 187 | "{}: There are {} zone overrides but only {} zones in the configuration.", |
| 188 | self.ozid(), zones.len(), self.cfg().num_zones; |
| 189 | Invalid, Input)); |
| 190 | } |
| 191 | let mut zdirs = BTreeMap::new(); |
| 192 | for (zdat, zmapdat) in zones.iter() { |
| 193 | match zdat { |
| 194 | Dat::U16(z) => { |
| 195 | if *z == 0 { |
| 196 | return Err(err!( |
| 197 | "{}: The configured zone {} override cannot index from 0, \ |
| 198 | zones are indexed from 1.", self.ozid(), z; |
| 199 | Invalid, Input)); |
| 200 | } |
| 201 | match zmapdat { |
| 202 | Dat::Map(zmap) => { |
| 203 | // Duplicate keys should be detected in from_datmap |
| 204 | let zdstr: ZoneDirStr = res!(ZoneDirStr::from_datmap(zmap.clone())); |
| 205 | let zdir = { |
| 206 | let mut raw_path = if zdstr.dir.len() == 0 { |
| 207 | PathBuf::from(self.db_root()) |
| 208 | } else { |
| 209 | res!(PathBuf::from_str(&zdstr.dir)) |
| 210 | }; |
| 211 | if !raw_path.is_absolute() { |
| 212 | // Relative paths are relative to the db_root |
| 213 | let mut tmp = PathBuf::from(self.db_root()); |
| 214 | tmp.push(&raw_path); |
| 215 | raw_path = tmp; |
| 216 | } |
| 217 | let canonical_path_str = |
| 218 | res!(OzoneConfig::canonicalize_path(raw_path)); |
| 219 | let dir = self.cfg().zone_root(&Path::new(&canonical_path_str)); |
| 220 | res!(std::fs::create_dir_all(&dir)); |
| 221 | ZoneDir { |
| 222 | dir, |
| 223 | max_size: zdstr.max_size, |
| 224 | } |
| 225 | }; |
| 226 | zdirs.insert(*z - 1, zdir); |
| 227 | }, |
| 228 | _ => return Err(err!( |
| 229 | "{}: The configured zone {} override '{:?}' is not a Dat::Map.", |
| 230 | self.ozid(), z, zmapdat; |
| 231 | Invalid, Input)), |
| 232 | } |
| 233 | }, |
| 234 | _ => return Err(err!( |
| 235 | "{}: The zone number in the configured zone override map, '{:?}', |
| 236 | is not a Dat::U16.", self.ozid(), zdat; |
| 237 | Invalid, Input)), |
| 238 | } |
| 239 | } |
| 240 | Ok(zdirs) |
| 241 | } |
| 242 | } |