1use crate::context::{ResourceContext, ServiceContext};
4use crate::event::EventBus;
5use crate::system::System;
6use async_trait::async_trait;
7use std::any::Any;
8use std::sync::Arc;
9
10use super::config::ResearchConfig;
11use super::events::*;
12use super::hook::ResearchHook;
13use super::research_projects::ResearchProjects;
14use super::state::ResearchState;
15use super::types::*;
16
17#[derive(Clone)]
32pub struct ResearchSystem {
33 hook: Arc<dyn ResearchHook>,
34}
35
36impl ResearchSystem {
37 pub fn new(hook: Arc<dyn ResearchHook>) -> Self {
39 Self { hook }
40 }
41
42 pub async fn process_events(
44 &mut self,
45 services: &ServiceContext,
46 resources: &mut ResourceContext,
47 ) {
48 self.process_queue_requests(services, resources).await;
49 self.process_start_requests(services, resources).await;
50 self.process_cancel_requests(services, resources).await;
51 self.process_progress_requests(services, resources).await;
52 self.process_complete_requests(services, resources).await;
53
54 self.auto_advance_progress(resources).await;
56 }
57
58 async fn process_queue_requests(
60 &mut self,
61 _services: &ServiceContext,
62 resources: &mut ResourceContext,
63 ) {
64 let requests = {
66 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
67 let reader = bus.reader::<ResearchQueueRequested>();
68 reader.iter().cloned().collect::<Vec<_>>()
69 } else {
70 Vec::new()
71 }
72 };
73
74 for request in requests {
75 let project = {
77 if let Some(projects) = resources.get::<ResearchProjects>().await {
78 match projects.get(&request.project_id) {
79 Some(p) => p.clone(),
80 None => continue,
81 }
82 } else {
83 continue;
84 }
85 };
86
87 {
89 let resources_ref = resources as &ResourceContext;
90 match self
91 .hook
92 .validate_prerequisites(&project, resources_ref)
93 .await
94 {
95 Ok(()) => {}
96 Err(_) => continue,
97 }
98 }
99
100 let cost = {
102 let resources_ref = resources as &ResourceContext;
103 match self
104 .hook
105 .calculate_research_cost(&project, resources_ref)
106 .await
107 {
108 Ok(c) => c,
109 Err(_) => continue,
110 }
111 };
112
113 {
115 if let Some(mut state) = resources.get_mut::<ResearchState>().await {
116 if state.queue(&request.project_id).is_err() {
117 continue;
118 }
119 } else {
120 continue;
121 }
122 }
123
124 {
126 let config = resources.get::<ResearchConfig>().await;
127 let state = resources.get_mut::<ResearchState>().await;
128
129 if let (Some(cfg), Some(mut st)) = (config, state) {
130 let max_slots = if cfg.allow_parallel_research {
131 cfg.max_parallel_slots
132 } else {
133 1
134 };
135 st.activate_next_queued(max_slots);
136 }
137 }
138
139 self.hook.on_research_queued(&project, resources).await;
141
142 let was_started = {
144 if let Some(state) = resources.get::<ResearchState>().await {
145 state.get_status(&request.project_id) == ResearchStatus::InProgress
146 } else {
147 false
148 }
149 };
150
151 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
153 bus.publish(ResearchQueuedEvent {
154 project_id: project.id.clone(),
155 project_name: project.name.clone(),
156 cost,
157 });
158
159 if was_started {
160 bus.publish(ResearchStartedEvent {
161 project_id: project.id.clone(),
162 project_name: project.name.clone(),
163 });
164 }
165 }
166 }
167 }
168
169 async fn process_start_requests(
171 &mut self,
172 _services: &ServiceContext,
173 resources: &mut ResourceContext,
174 ) {
175 let _requests = {
176 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
177 let reader = bus.reader::<ResearchStartRequested>();
178 reader.iter().cloned().collect::<Vec<_>>()
179 } else {
180 Vec::new()
181 }
182 };
183
184 }
186
187 async fn process_cancel_requests(
189 &mut self,
190 _services: &ServiceContext,
191 resources: &mut ResourceContext,
192 ) {
193 let requests = {
194 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
195 let reader = bus.reader::<ResearchCancelRequested>();
196 reader.iter().cloned().collect::<Vec<_>>()
197 } else {
198 Vec::new()
199 }
200 };
201
202 for request in requests {
203 let project = {
205 if let Some(projects) = resources.get::<ResearchProjects>().await {
206 match projects.get(&request.project_id) {
207 Some(p) => p.clone(),
208 None => continue,
209 }
210 } else {
211 continue;
212 }
213 };
214
215 {
217 if let Some(mut state) = resources.get_mut::<ResearchState>().await {
218 if state.cancel(&request.project_id).is_err() {
219 continue;
220 }
221 } else {
222 continue;
223 }
224 }
225
226 {
228 let config = resources.get::<ResearchConfig>().await;
229 let state = resources.get_mut::<ResearchState>().await;
230
231 if let (Some(cfg), Some(mut st)) = (config, state) {
232 let max_slots = if cfg.allow_parallel_research {
233 cfg.max_parallel_slots
234 } else {
235 1
236 };
237 st.activate_next_queued(max_slots);
238 }
239 }
240
241 self.hook.on_research_failed(&project, resources).await;
243
244 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
246 bus.publish(ResearchCancelledEvent {
247 project_id: project.id.clone(),
248 project_name: project.name.clone(),
249 });
250 }
251 }
252 }
253
254 async fn process_progress_requests(
256 &mut self,
257 _services: &ServiceContext,
258 resources: &mut ResourceContext,
259 ) {
260 let requests = {
261 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
262 let reader = bus.reader::<ResearchProgressRequested>();
263 reader.iter().cloned().collect::<Vec<_>>()
264 } else {
265 Vec::new()
266 }
267 };
268
269 for request in requests {
270 {
272 if let Some(mut state) = resources.get_mut::<ResearchState>().await {
273 state.add_progress(&request.project_id, request.amount);
274 } else {
275 continue;
276 }
277 }
278
279 let (completed, progress) = {
281 if let Some(state) = resources.get::<ResearchState>().await {
282 let prog = state.get_progress(&request.project_id);
283 (prog >= 1.0, prog)
284 } else {
285 continue;
286 }
287 };
288
289 if completed {
290 if let Some(mut state) = resources.get_mut::<ResearchState>().await {
292 let _ = state.complete(&request.project_id);
293 }
294
295 {
297 let config = resources.get::<ResearchConfig>().await;
298 let state = resources.get_mut::<ResearchState>().await;
299
300 if let (Some(cfg), Some(mut st)) = (config, state) {
301 let max_slots = if cfg.allow_parallel_research {
302 cfg.max_parallel_slots
303 } else {
304 1
305 };
306 st.activate_next_queued(max_slots);
307 }
308 }
309
310 let project = {
312 if let Some(projects) = resources.get::<ResearchProjects>().await {
313 projects.get(&request.project_id).cloned()
314 } else {
315 None
316 }
317 };
318
319 if let Some(proj) = project {
320 let result = ResearchResult {
321 project_id: proj.id.clone(),
322 success: true,
323 final_metrics: proj.metrics.clone(),
324 metadata: proj.metadata.clone(),
325 };
326
327 self.hook
329 .on_research_completed(&proj, &result, resources)
330 .await;
331
332 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
334 bus.publish(ResearchCompletedEvent {
335 project_id: proj.id.clone(),
336 project_name: proj.name.clone(),
337 result,
338 });
339 }
340 }
341 } else {
342 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
344 bus.publish(ResearchProgressUpdatedEvent {
345 project_id: request.project_id.clone(),
346 progress,
347 });
348 }
349 }
350 }
351 }
352
353 async fn process_complete_requests(
355 &mut self,
356 _services: &ServiceContext,
357 resources: &mut ResourceContext,
358 ) {
359 let requests = {
360 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
361 let reader = bus.reader::<ResearchCompleteRequested>();
362 reader.iter().cloned().collect::<Vec<_>>()
363 } else {
364 Vec::new()
365 }
366 };
367
368 for request in requests {
369 let project = {
371 if let Some(projects) = resources.get::<ResearchProjects>().await {
372 match projects.get(&request.project_id) {
373 Some(p) => p.clone(),
374 None => continue,
375 }
376 } else {
377 continue;
378 }
379 };
380
381 {
383 if let Some(mut state) = resources.get_mut::<ResearchState>().await {
384 if state.complete(&request.project_id).is_err() {
385 continue;
386 }
387 } else {
388 continue;
389 }
390 }
391
392 {
394 let config = resources.get::<ResearchConfig>().await;
395 let state = resources.get_mut::<ResearchState>().await;
396
397 if let (Some(cfg), Some(mut st)) = (config, state) {
398 let max_slots = if cfg.allow_parallel_research {
399 cfg.max_parallel_slots
400 } else {
401 1
402 };
403 st.activate_next_queued(max_slots);
404 }
405 }
406
407 let result = ResearchResult {
408 project_id: project.id.clone(),
409 success: true,
410 final_metrics: project.metrics.clone(),
411 metadata: project.metadata.clone(),
412 };
413
414 self.hook
416 .on_research_completed(&project, &result, resources)
417 .await;
418
419 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
421 bus.publish(ResearchCompletedEvent {
422 project_id: project.id.clone(),
423 project_name: project.name.clone(),
424 result,
425 });
426 }
427 }
428 }
429
430 async fn auto_advance_progress(&mut self, resources: &mut ResourceContext) {
432 let (enabled, progress_per_turn) = {
434 if let Some(config) = resources.get::<ResearchConfig>().await {
435 (config.auto_advance, config.base_progress_per_turn)
436 } else {
437 return;
438 }
439 };
440
441 if !enabled {
442 return;
443 }
444
445 let active_ids = {
447 if let Some(state) = resources.get::<ResearchState>().await {
448 state.active_projects()
449 } else {
450 return;
451 }
452 };
453
454 let mut completed = Vec::new();
455
456 {
458 if let Some(mut state) = resources.get_mut::<ResearchState>().await {
459 for id in &active_ids {
460 state.add_progress(id, progress_per_turn);
461
462 if state.get_progress(id) >= 1.0 {
463 let _ = state.complete(id);
464 completed.push(id.clone());
465 }
466 }
467 }
468 }
469
470 if !completed.is_empty() {
472 let config = resources.get::<ResearchConfig>().await;
473 let state = resources.get_mut::<ResearchState>().await;
474
475 if let (Some(cfg), Some(mut st)) = (config, state) {
476 let max_slots = if cfg.allow_parallel_research {
477 cfg.max_parallel_slots
478 } else {
479 1
480 };
481 st.activate_next_queued(max_slots);
482 }
483 }
484
485 for project_id in completed {
487 let project = {
488 if let Some(projects) = resources.get::<ResearchProjects>().await {
489 projects.get(&project_id).cloned()
490 } else {
491 None
492 }
493 };
494
495 if let Some(proj) = project {
496 let result = ResearchResult {
497 project_id: proj.id.clone(),
498 success: true,
499 final_metrics: proj.metrics.clone(),
500 metadata: proj.metadata.clone(),
501 };
502
503 self.hook
504 .on_research_completed(&proj, &result, resources)
505 .await;
506
507 if let Some(mut bus) = resources.get_mut::<EventBus>().await {
508 bus.publish(ResearchCompletedEvent {
509 project_id: proj.id.clone(),
510 project_name: proj.name.clone(),
511 result,
512 });
513 }
514 }
515 }
516 }
517}
518
519#[async_trait]
520impl System for ResearchSystem {
521 fn name(&self) -> &'static str {
522 "research_system"
523 }
524
525 fn as_any(&self) -> &dyn Any {
526 self
527 }
528
529 fn as_any_mut(&mut self) -> &mut dyn Any {
530 self
531 }
532}