MCPcopy Create free account
hub / github.com/SkyworkAI/DeepResearchAgent / stream_insert

Method stream_insert

src/environment/alpacaentry/bars.py:197–253  ·  view source on GitHub ↗

Insert bars data from stream. Args: data: Raw bars data dictionary from Alpaca stream symbol: Symbol name asset_type: Optional asset type (not used for bars, kept for compatibility) Returns: True if successful, Fal

(self, data: Dict, symbol: str, asset_type=None)

Source from the content-addressed store, hash-verified

195 return db_data
196
197 async def stream_insert(self, data: Dict, symbol: str, asset_type=None) -> bool:
198 """Insert bars data from stream.
199
200 Args:
201 data: Raw bars data dictionary from Alpaca stream
202 symbol: Symbol name
203 asset_type: Optional asset type (not used for bars, kept for compatibility)
204
205 Returns:
206 True if successful, False otherwise
207 """
208 try:
209 # Ensure table exists
210 await self.ensure_table_exists(symbol)
211
212 # Prepare data for insertion
213 db_data = self._prepare_data_for_insert(data, symbol)
214
215 # Insert into database
216 base_name = self._sanitize_table_name(symbol)
217 table_name = f"{base_name}_bars"
218
219 insert_request = InsertRequest(
220 table_name=table_name,
221 data=db_data
222 )
223
224 # Debug: Log what we're trying to insert
225 logger.debug(f"| 🔍 Attempting to insert bars data for {symbol} into {table_name}: {db_data}")
226
227 result = await self.database_service.insert_data(insert_request)
228
229 if not result.success:
230 logger.error(f"| ❌ Failed to insert bars data for {symbol}: {result.message}")
231 logger.error(f"| ❌ Insert request: table={table_name}, data={db_data}")
232 return False
233
234 # Verify insertion by querying the table
235 verify_query = f"SELECT COUNT(*) as count FROM {table_name}"
236 verify_result = await self.database_service.execute_query(QueryRequest(query=verify_query))
237 if verify_result.success:
238 count = verify_result.extra.get("data", [{}])[0].get("count", 0)
239 logger.debug(f"| ✅ Bars data inserted for {symbol}. Total rows in {table_name}: {count}")
240 else:
241 logger.warning(f"| ⚠️ Insert succeeded but couldn't verify count for {symbol}")
242
243 # Update cache
244 self._update_cache(symbol, db_data)
245
246 # Calculate and store indicators
247 await self._calculate_and_store_indicators(symbol)
248
249 return True
250
251 except Exception as e:
252 logger.error(f"Error inserting bars data for {symbol}: {e}")
253 return False
254

Callers 1

_handle_dataMethod · 0.45

Calls 13

ensure_table_existsMethod · 0.95
_sanitize_table_nameMethod · 0.95
_update_cacheMethod · 0.95
InsertRequestClass · 0.90
QueryRequestClass · 0.90
debugMethod · 0.80
execute_queryMethod · 0.80
warningMethod · 0.80
insert_dataMethod · 0.45
errorMethod · 0.45

Tested by

no test coverage detected