Ig0tU commited on
Commit
f64763a
·
1 Parent(s): c6a91aa

Self-healing fuzzy trail + /ui/trails endpoint

Browse files
Files changed (4) hide show
  1. .gitattributes +1 -0
  2. .gitignore +1 -0
  3. core/signal_registry.py +167 -25
  4. signalmesh_api.py +15 -0
.gitattributes CHANGED
@@ -33,3 +33,4 @@ saved_model/**/* filter=lfs diff=lfs merge=lfs -text
33
  *.zip filter=lfs diff=lfs merge=lfs -text
34
  *.zst filter=lfs diff=lfs merge=lfs -text
35
  *tfevents* filter=lfs diff=lfs merge=lfs -text
 
 
33
  *.zip filter=lfs diff=lfs merge=lfs -text
34
  *.zst filter=lfs diff=lfs merge=lfs -text
35
  *tfevents* filter=lfs diff=lfs merge=lfs -text
36
+ marketing/videos/*.mp4 filter=lfs diff=lfs merge=lfs -text
.gitignore CHANGED
@@ -4,3 +4,4 @@ __pycache__/
4
  .env
5
  *.log
6
  .DS_Store
 
 
4
  .env
5
  *.log
6
  .DS_Store
7
+ marketing/videos/
core/signal_registry.py CHANGED
@@ -1,54 +1,196 @@
1
- import json
2
- import os
3
  import time
4
- from typing import Dict, List, Any, Optional
 
5
  from app.logger import logger
6
 
 
7
  class SignalStream:
8
  def __init__(self, name: str, source_type: str, data: Any, metadata: Dict[str, Any]):
9
  self.name = name
10
- self.source_type = source_type # 'rss', 'mcp', 'context', 'log'
11
  self.data = data
12
  self.metadata = metadata
13
  self.timestamp = time.time()
14
 
 
15
  class SignalRegistry:
16
  """
17
  The Signal Mesh: A central hub where data sources broadcast signals.
18
- Agents (Antennae) tune into these based on their spatial/functional mapping.
 
 
 
 
 
 
 
 
19
  """
20
  _instance = None
 
 
21
 
22
  def __new__(cls):
23
  if cls._instance is None:
24
- cls._instance = super(SignalRegistry, cls).__new__(cls)
25
  cls._instance.streams: Dict[str, List[SignalStream]] = {}
 
 
26
  return cls._instance
27
 
28
- def broadcast(self, name: str, source_type: str, data: Any, metadata: Optional[Dict[str, Any]] = None):
29
- """Broadcast a signal into the mesh."""
 
 
30
  if name not in self.streams:
31
  self.streams[name] = []
32
-
33
- stream = SignalStream(name, source_type, data, metadata or {})
34
- self.streams[name].append(stream)
35
- # Keep last 100 signals per frequency
36
  self.streams[name] = self.streams[name][-100:]
37
  logger.info(f"📡 Signal Mesh: Broadcasting stream '{name}' ({source_type})")
38
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
39
  def tune_in(self, keywords: List[str]) -> List[Dict[str, Any]]:
40
- """Antennae use this to catch signals matching their frequency (keywords)."""
41
- caught_signals = []
42
- for name, buffer in self.streams.items():
43
- # Match signals to agent antennae keywords or name
44
- if any(key.lower() in name.lower() for key in keywords):
45
- for signal in buffer:
46
- caught_signals.append({
47
- "name": signal.name,
48
- "type": signal.source_type,
49
- "content": signal.data,
50
- "age": time.time() - signal.timestamp
51
- })
52
- return caught_signals
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
53
 
54
  signal_registry = SignalRegistry()
 
1
+ import re
 
2
  import time
3
+ from difflib import SequenceMatcher
4
+ from typing import Dict, List, Any, Optional, Tuple
5
  from app.logger import logger
6
 
7
+
8
  class SignalStream:
9
  def __init__(self, name: str, source_type: str, data: Any, metadata: Dict[str, Any]):
10
  self.name = name
11
+ self.source_type = source_type # 'rss', 'mcp', 'context', 'log'
12
  self.data = data
13
  self.metadata = metadata
14
  self.timestamp = time.time()
15
 
16
+
17
  class SignalRegistry:
18
  """
19
  The Signal Mesh: A central hub where data sources broadcast signals.
20
+ Agents (Antennae) tune in based on keywords.
21
+
22
+ Miss handling — self-healing fuzzy trail:
23
+ 1. Score unmapped keyword against all live frequencies (SequenceMatcher + token Jaccard)
24
+ 2. If best score >= TRAIL_THRESHOLD: acknowledge gap, write persistent trail
25
+ (keyword → [best_freq, ...token-overlap neighbors])
26
+ 3. Branch outward from best match to capture related frequencies
27
+ 4. Future calls with that keyword hit the trail as a fast-path — no re-scoring
28
+ 5. Gap metadata rides along in the returned signal so the caller can inspect it
29
  """
30
  _instance = None
31
+ TRAIL_THRESHOLD = 0.45 # min combined score to write a trail
32
+ BRANCH_THRESHOLD = 0.30 # min token-overlap ratio to include a neighbor
33
 
34
  def __new__(cls):
35
  if cls._instance is None:
36
+ cls._instance = super().__new__(cls)
37
  cls._instance.streams: Dict[str, List[SignalStream]] = {}
38
+ cls._instance.fuzzy_trails: Dict[str, List[str]] = {}
39
+ cls._instance.trail_gaps: Dict[str, Dict[str, float]] = {}
40
  return cls._instance
41
 
42
+ # ── Broadcast ─────────────────────────────────────────────────────────────
43
+
44
+ def broadcast(self, name: str, source_type: str, data: Any,
45
+ metadata: Optional[Dict[str, Any]] = None):
46
  if name not in self.streams:
47
  self.streams[name] = []
48
+ self.streams[name].append(SignalStream(name, source_type, data, metadata or {}))
 
 
 
49
  self.streams[name] = self.streams[name][-100:]
50
  logger.info(f"📡 Signal Mesh: Broadcasting stream '{name}' ({source_type})")
51
 
52
+ # ── Matching helpers ───────────────────────────────────────────────────────
53
+
54
+ @staticmethod
55
+ def _tokens(s: str) -> set:
56
+ return set(re.split(r'[_\-\s]+', s.lower())) - {''}
57
+
58
+ @staticmethod
59
+ def _similarity(a: str, b: str) -> float:
60
+ return SequenceMatcher(None, a.lower(), b.lower()).ratio()
61
+
62
+ def _fast_match(self, freq_name: str, kw: str) -> bool:
63
+ kl = kw.lower()
64
+ fl = freq_name.lower()
65
+ if kl in fl:
66
+ return True
67
+ if self._tokens(kl) & self._tokens(fl):
68
+ return True
69
+ # camelCase decomposition
70
+ camel = set(re.sub(r'([A-Z])', r' \1', kw).lower().split())
71
+ return bool(camel & self._tokens(fl))
72
+
73
+ def _score_all(self, keyword: str) -> List[Tuple[str, float]]:
74
+ kl = keyword.lower()
75
+ results = []
76
+ for name in self.streams:
77
+ tok_a = self._tokens(kl)
78
+ tok_b = self._tokens(name.lower())
79
+ union = tok_a | tok_b
80
+ jaccard = len(tok_a & tok_b) / len(union) if union else 0.0
81
+ score = 0.6 * self._similarity(kl, name) + 0.4 * jaccard
82
+ results.append((name, score))
83
+ return sorted(results, key=lambda x: x[1], reverse=True)
84
+
85
+ def _branch_from(self, anchor: str) -> List[str]:
86
+ anchor_tokens = self._tokens(anchor.lower())
87
+ neighbors = []
88
+ for name in self.streams:
89
+ if name == anchor:
90
+ continue
91
+ overlap = len(anchor_tokens & self._tokens(name.lower()))
92
+ if overlap and overlap / max(len(anchor_tokens), 1) >= self.BRANCH_THRESHOLD:
93
+ neighbors.append(name)
94
+ return neighbors
95
+
96
+ # ── Trail writer ───────────────────────────────────────────────────────────
97
+
98
+ def _write_trail(self, keyword: str, scores: List[Tuple[str, float]]) -> List[str]:
99
+ best_freq, best_score = scores[0]
100
+ if best_score < self.TRAIL_THRESHOLD:
101
+ logger.warning(
102
+ f"🕸 No trail for '{keyword}' — "
103
+ f"best match '{best_freq}' @ {best_score:.2f} (threshold {self.TRAIL_THRESHOLD})"
104
+ )
105
+ return []
106
+
107
+ branches = self._branch_from(best_freq)
108
+ trail = [best_freq] + branches
109
+ self.fuzzy_trails[keyword] = trail
110
+ self.trail_gaps[keyword] = {f: round(s, 3) for f, s in scores[:5]}
111
+ logger.info(
112
+ f"🕸 Trail written: '{keyword}' → {trail} "
113
+ f"(anchor '{best_freq}' @ {best_score:.2f}, {len(branches)} branch(es))"
114
+ )
115
+ return trail
116
+
117
+ # ── tune_in ────────────────────────────────────────────────────────────────
118
+
119
  def tune_in(self, keywords: List[str]) -> List[Dict[str, Any]]:
120
+ """
121
+ Match priority per keyword:
122
+ 1. Learned fuzzy trail fast-path
123
+ 2. Direct fast match (substring / token / camelCase)
124
+ 3. Miss → score all → write trail if above threshold return with _trail metadata
125
+ """
126
+ matched: Dict[str, Optional[str]] = {} # freq → trail_kw|None
127
+ gap_report: Dict[str, Dict] = {}
128
+
129
+ for kw in keywords:
130
+ # 1. Trail fast-path
131
+ if kw in self.fuzzy_trails:
132
+ for freq in self.fuzzy_trails[kw]:
133
+ if freq in self.streams:
134
+ matched.setdefault(freq, kw)
135
+ continue
136
+
137
+ # 2. Direct match
138
+ direct = False
139
+ for name in self.streams:
140
+ if self._fast_match(name, kw):
141
+ matched.setdefault(name, None)
142
+ direct = True
143
+ if direct:
144
+ continue
145
+
146
+ # 3. Miss → benchmark → write trail
147
+ scores = self._score_all(kw)
148
+ trail = self._write_trail(kw, scores)
149
+ for freq in trail:
150
+ if freq in self.streams:
151
+ matched.setdefault(freq, kw)
152
+ if trail:
153
+ gap_report[kw] = {
154
+ "bridged_to": trail[0],
155
+ "confidence": round(scores[0][1], 3),
156
+ "branches": trail[1:],
157
+ "top5": {f: round(s, 3) for f, s in scores[:5]},
158
+ }
159
+
160
+ # Collect signals
161
+ caught: List[Dict[str, Any]] = []
162
+ seen: set = set()
163
+ for freq_name, via_kw in matched.items():
164
+ for signal in self.streams.get(freq_name, []):
165
+ sig_id = (signal.name, signal.timestamp)
166
+ if sig_id in seen:
167
+ continue
168
+ seen.add(sig_id)
169
+ entry: Dict[str, Any] = {
170
+ "name": signal.name,
171
+ "type": signal.source_type,
172
+ "content": signal.data,
173
+ "age": time.time() - signal.timestamp,
174
+ }
175
+ if via_kw and via_kw in gap_report:
176
+ entry["_trail"] = gap_report[via_kw]
177
+ caught.append(entry)
178
+
179
+ return caught
180
+
181
+ # ── Introspection ──────────────────────────────────────────────────────────
182
+
183
+ def get_trails(self) -> Dict[str, Any]:
184
+ return {
185
+ kw: {"frequencies": freqs, "gap_scores": self.trail_gaps.get(kw, {})}
186
+ for kw, freqs in self.fuzzy_trails.items()
187
+ }
188
+
189
+ def clear_trail(self, keyword: str) -> bool:
190
+ removed = keyword in self.fuzzy_trails
191
+ self.fuzzy_trails.pop(keyword, None)
192
+ self.trail_gaps.pop(keyword, None)
193
+ return removed
194
+
195
 
196
  signal_registry = SignalRegistry()
signalmesh_api.py CHANGED
@@ -357,6 +357,21 @@ async def ui_hf(req: HfInjectReq): return await _do_hf_inject(req.repo_id, req.r
357
  @app.post("/ui/arxiv_inject")
358
  async def ui_arxiv(req: ArxivInjectReq): return await _do_arxiv_inject(req.arxiv_id)
359
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
360
  # ── Dashboard (v3 — clean bright build) ──────────────────────────────────────
361
  _HTML = r"""<!DOCTYPE html>
362
  <html lang="en">
 
357
  @app.post("/ui/arxiv_inject")
358
  async def ui_arxiv(req: ArxivInjectReq): return await _do_arxiv_inject(req.arxiv_id)
359
 
360
+ @app.get("/ui/trails")
361
+ def ui_trails():
362
+ """Return all learned fuzzy trails — keyword → frequencies + gap scores."""
363
+ return registry.get_trails()
364
+
365
+ @app.delete("/ui/trails/{keyword}")
366
+ def ui_clear_trail(keyword: str):
367
+ removed = registry.clear_trail(keyword)
368
+ return {"removed": removed, "keyword": keyword}
369
+
370
+ @app.get("/api/trails")
371
+ def api_trails(x_signalmesh_key: Optional[str] = Header(default=None)):
372
+ _check_key(x_signalmesh_key)
373
+ return registry.get_trails()
374
+
375
  # ── Dashboard (v3 — clean bright build) ──────────────────────────────────────
376
  _HTML = r"""<!DOCTYPE html>
377
  <html lang="en">