Oregami
Repositories/oxedyne/fe2o3

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

1use 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
19use oxedyne_fe2o3_core::channels::Recv;
20use oxedyne_fe2o3_jdat::{
21 prelude::*,
22 id::NumIdDat,
23};
24
25use 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.
37pub 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
59impl<
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
106impl<
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
119impl<
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}