Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,8 @@ temp/
# Market and analytical data
# ------------------------------------------------------------
data/
!/src/ai4binance/data/
!/src/ai4binance/data/**
datasets/
market_data/
historical_data/
Expand Down
2 changes: 1 addition & 1 deletion config/operations/services.json
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,7 @@
"command": "python -m ai4binance.cli.futures_multitf daemon",
"mode": "RunFuturesMultiTf",
"required": false,
"enabled": false,
"enabled": true,
"health_file": "futures-multitf-health.json",
"lock_file": "futures-multitf.lock",
"max_health_age_seconds": 900,
Expand Down
29 changes: 29 additions & 0 deletions scripts/primary_local_reasoning_agent.ps1
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ $ErrorActionPreference = "Stop"
$root = (Resolve-Path -LiteralPath (Join-Path $PSScriptRoot "..")).Path
$taskName = "AI4BINANCE-Primary-Local-Reasoning"
$healthPath = Join-Path $root "runtime\state\primary-local-reasoning-health.json"
$lockPath = Join-Path $root "runtime\state\primary-local-reasoning.lock"
$serverScript = Join-Path $PSScriptRoot "start_llama_server.ps1"

function Test-PrimaryListener {
Expand All @@ -32,6 +33,7 @@ function Write-PrimaryHealth {
endpoint = "http://127.0.0.1:8080"
status = $Status
updated_at = [DateTimeOffset]::UtcNow.ToString("o")
pid = $PID
listener_pids = @($listeners | ForEach-Object { $_.OwningProcess } | Sort-Object -Unique)
blockers = @($Blockers)
execution_allowed = $false
Expand All @@ -41,6 +43,31 @@ function Write-PrimaryHealth {
} | ConvertTo-Json -Depth 4 | Set-Content -LiteralPath $healthPath -Encoding UTF8
}

function Write-PrimaryLock {
New-Item -ItemType Directory -Path (Split-Path $lockPath -Parent) -Force | Out-Null
$temporary = "$lockPath.$PID.tmp"
[System.IO.File]::WriteAllText(
$temporary,
[string]$PID,
[System.Text.Encoding]::ASCII
)
Move-Item -LiteralPath $temporary -Destination $lockPath -Force
}

function Remove-PrimaryLock {
if (-not (Test-Path -LiteralPath $lockPath -PathType Leaf)) {
return
}
try {
$ownerPid = [int](Get-Content -LiteralPath $lockPath -Raw -Encoding ASCII)
if ($ownerPid -eq $PID) {
Remove-Item -LiteralPath $lockPath -Force
}
}
catch {
}
}

function Start-PrimaryProvider {
try {
& $serverScript | Out-Null
Expand Down Expand Up @@ -84,12 +111,14 @@ $hasHandle = $false
try {
$hasHandle = $mutex.WaitOne(0)
if (-not $hasHandle) { exit 0 }
Write-PrimaryLock
while ($true) {
[void](Start-PrimaryProvider)
Start-Sleep -Seconds $PollSeconds
}
}
finally {
Remove-PrimaryLock
if ($hasHandle) { $mutex.ReleaseMutex() }
$mutex.Dispose()
}
51 changes: 38 additions & 13 deletions scripts/ykb_daily_report.ps1
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,7 @@ function Sync-LatestHumanReport {
Reason = "ykb_report_MISSING"
ReportId = ""
Status = ""
Blockers = @()
ObservedAt = ""
MarkdownPath = ""
LatestMarkdown = ""
Expand All @@ -124,6 +125,7 @@ function Sync-LatestHumanReport {
Reason = "YKB_HUMAN_REPORT_READY"
ReportId = [string]$payload.report_id
Status = [string]$payload.status
Blockers = @($payload.blockers)
ObservedAt = [string]$payload.observed_at
MarkdownPath = $markdownPath
LatestMarkdown = $markdownPath
Expand All @@ -135,6 +137,7 @@ function Sync-LatestHumanReport {
Reason = "YKB_HUMAN_REPORT_SYNC_FAILED"
ReportId = ""
Status = ""
Blockers = @()
ObservedAt = ""
MarkdownPath = ""
LatestMarkdown = ""
Expand Down Expand Up @@ -184,10 +187,18 @@ function Write-Health {
[AllowEmptyCollection()][string[]]$Blockers = @(),
[int]$RefreshExitCode = 0,
[int]$ReportExitCode = 0,
[string]$ReportStatus = "NOT_EVALUATED",
[AllowEmptyCollection()][string[]]$ReportBlockers = @(),
[string]$LastReportObservedAt = "",
[string]$LatestMarkdownPath = ""
)
New-Item -ItemType Directory -Path $stateDirectory -Force | Out-Null
[string[]]$normalizedBlockers = @(
$Blockers | Where-Object { -not [string]::IsNullOrWhiteSpace($_) }
)
[string[]]$normalizedReportBlockers = @(
$ReportBlockers | Where-Object { -not [string]::IsNullOrWhiteSpace($_) }
)
$payload = [ordered]@{
service = "ykb-report"
status = $Status
Expand All @@ -199,7 +210,9 @@ function Write-Health {
latest_report_observed_at = $LastReportObservedAt
refresh_exit_code = $RefreshExitCode
report_exit_code = $ReportExitCode
blockers = @($Blockers)
report_status = $ReportStatus
report_blockers = $normalizedReportBlockers
blockers = $normalizedBlockers
execution_allowed = $false
live_eligibility_status = "LIVE_ORDER_BLOCKED"
}
Expand Down Expand Up @@ -437,20 +450,26 @@ function Invoke-YkbRefreshAndReport {
}
$state = Read-LatestReportState
$humanReport = Sync-LatestHumanReport
$blockers = [System.Collections.Generic.List[string]]::new()
$operationalBlockers = [System.Collections.Generic.List[string]]::new()
if ($refreshExitCode -ne 0) {
$blockers.Add("RUNTIME_RESEARCH_REFRESH_WITH_BLOCKERS")
$operationalBlockers.Add("RUNTIME_RESEARCH_REFRESH_WITH_BLOCKERS")
}
if ($reportExitCode -ne 0) {
$blockers.Add("ykb_report_WITH_BLOCKERS")
if ($reportExitCode -notin @(0, 2)) {
$operationalBlockers.Add("YKB_REPORT_COMMAND_FAILED")
}
if (
$reportExitCode -eq 2 -and
$humanReport.Status -ne "RUNNING_WITH_BLOCKERS"
) {
$operationalBlockers.Add("YKB_REPORT_EXIT_STATUS_MISMATCH")
}
if (-not $state.Fresh) {
$blockers.Add([string]$state.Reason)
$operationalBlockers.Add([string]$state.Reason)
}
if (-not $humanReport.Synced) {
$blockers.Add([string]$humanReport.Reason)
$operationalBlockers.Add([string]$humanReport.Reason)
}
$status = if ($state.Fresh -and $blockers.Count -eq 0) {
$status = if ($state.Fresh -and $operationalBlockers.Count -eq 0) {
"READY"
}
elseif ($state.Fresh) {
Expand All @@ -461,9 +480,11 @@ function Invoke-YkbRefreshAndReport {
}
Write-Health `
-Status $status `
-Blockers $blockers.ToArray() `
-Blockers $operationalBlockers.ToArray() `
-RefreshExitCode $refreshExitCode `
-ReportExitCode $reportExitCode `
-ReportStatus ([string]$humanReport.Status) `
-ReportBlockers @($humanReport.Blockers) `
-LastReportObservedAt ([string]$state.ObservedAt) `
-LatestMarkdownPath ([string]$humanReport.LatestMarkdown)
if ($humanReport.Synced) {
Expand All @@ -474,7 +495,7 @@ function Invoke-YkbRefreshAndReport {
Write-Output "markdown_path=$($humanReport.MarkdownPath)"
Write-Output "latest_markdown_path=$($humanReport.LatestMarkdown)"
}
return $(if ($state.Fresh) { 0 } else { 1 })
return $(if ($status -eq "READY") { 0 } else { 1 })
}

function Start-YkbLoop {
Expand Down Expand Up @@ -616,13 +637,17 @@ if ($Mode -eq "RunLoop") {
$state = Read-LatestReportState
if ($state.Fresh -and -not $Force) {
$humanReport = Sync-LatestHumanReport
$status = if ($humanReport.Synced) { "READY" } else { "DEGRADED" }
$blockers = if ($humanReport.Synced) { @() } else { @($humanReport.Reason) }
Write-Health `
-Status "READY" `
-Blockers @() `
-Status $status `
-Blockers $blockers `
-ReportStatus ([string]$humanReport.Status) `
-ReportBlockers @($humanReport.Blockers) `
-LastReportObservedAt ([string]$state.ObservedAt) `
-LatestMarkdownPath ([string]$humanReport.LatestMarkdown)
Write-Output "ykb_report_FRESH"
exit 0
exit $(if ($humanReport.Synced) { 0 } else { 1 })
}

exit (Invoke-YkbRefreshAndReport)
Expand Down
157 changes: 101 additions & 56 deletions src/ai4binance/accounting/user_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
from urllib.parse import urlencode
from urllib.request import Request, urlopen

from websockets.exceptions import WebSocketException

from ai4binance.accounting.collectors import (
AccountingWebSocketCollector,
)
Expand Down Expand Up @@ -306,67 +308,28 @@ def __post_init__(self) -> None:
raise ValueError("WebSocket receive timeout is outside the safe range")

def collect_once(self, *, sync_run_id: str) -> WebSocketCollectionSummary:
accepted = 0
duplicates = 0
rejected = 0
blockers: list[str] = []
collector = AccountingWebSocketCollector(self.ledger)
deadline = time.monotonic() + self.collect_seconds
opened: list[UserDataStreamSession] = []
now = datetime.now(UTC)
opened, accepted, duplicates, rejected = self._open_sessions(
sync_run_id=sync_run_id,
received_at=datetime.now(UTC),
blockers=blockers,
)
try:
for session in self.sessions:
try:
session.open()
opened.append(session)
if self._append_subscription_event(session, sync_run_id, now):
accepted += 1
else:
duplicates += 1
except (OSError, RuntimeError, ValueError) as error:
rejected += 1
blockers.append(
f"{session.product_type.value}_WEBSOCKET_OPEN_FAILED:{type(error).__name__}"
)
while (
accepted + duplicates < self.event_limit and time.monotonic() < deadline
):
made_progress = False
for session in opened:
if accepted + duplicates >= self.event_limit:
break
try:
event = session.receive_event(
timeout_seconds=self.receive_timeout_seconds
)
except TimeoutError:
continue
except (OSError, RuntimeError, ValueError) as error:
blockers.append(
f"{session.product_type.value}_WEBSOCKET_EVENT_FAILED:{type(error).__name__}"
)
continue
result = collector.ingest_event(
event,
snapshot_id=sync_run_id,
sync_run_id=sync_run_id,
received_at=datetime.now(UTC),
)
accepted += result.accepted_count
duplicates += result.duplicate_count
rejected += result.rejected_count
blockers.extend(result.blockers)
made_progress = True
if not made_progress and not opened:
break
event_counts = self._collect_events(
collector=collector,
opened=opened,
sync_run_id=sync_run_id,
deadline=deadline,
accepted=accepted,
duplicates=duplicates,
rejected=rejected,
blockers=blockers,
)
accepted, duplicates, rejected = event_counts
finally:
for session in reversed(opened):
try:
session.close()
except (OSError, RuntimeError, ValueError):
blockers.append(
f"{session.product_type.value}_WEBSOCKET_CLOSE_FAILED"
)
self._close_sessions(opened, blockers)
return WebSocketCollectionSummary(
"COLLECTED" if not blockers and rejected == 0 else "DEGRADED",
accepted,
Expand All @@ -375,6 +338,88 @@ def collect_once(self, *, sync_run_id: str) -> WebSocketCollectionSummary:
tuple(dict.fromkeys(blockers)),
)

def _open_sessions(
self,
*,
sync_run_id: str,
received_at: datetime,
blockers: list[str],
) -> tuple[list[UserDataStreamSession], int, int, int]:
opened: list[UserDataStreamSession] = []
accepted = duplicates = rejected = 0
for session in self.sessions:
try:
session.open()
opened.append(session)
if self._append_subscription_event(session, sync_run_id, received_at):
accepted += 1
else:
duplicates += 1
except (OSError, RuntimeError, ValueError, WebSocketException) as error:
rejected += 1
blockers.append(
f"{session.product_type.value}_WEBSOCKET_OPEN_FAILED:{type(error).__name__}"
)
return opened, accepted, duplicates, rejected

def _collect_events(
self,
*,
collector: AccountingWebSocketCollector,
opened: list[UserDataStreamSession],
sync_run_id: str,
deadline: float,
accepted: int,
duplicates: int,
rejected: int,
blockers: list[str],
) -> tuple[int, int, int]:
while accepted + duplicates < self.event_limit and time.monotonic() < deadline:
made_progress = False
for session in opened:
if accepted + duplicates >= self.event_limit:
break
try:
event = session.receive_event(
timeout_seconds=self.receive_timeout_seconds
)
except TimeoutError:
continue
except (
OSError,
RuntimeError,
ValueError,
WebSocketException,
) as error:
blockers.append(
f"{session.product_type.value}_WEBSOCKET_EVENT_FAILED:{type(error).__name__}"
)
continue
result = collector.ingest_event(
event,
snapshot_id=sync_run_id,
sync_run_id=sync_run_id,
received_at=datetime.now(UTC),
)
accepted += result.accepted_count
duplicates += result.duplicate_count
rejected += result.rejected_count
blockers.extend(result.blockers)
made_progress = True
if not made_progress and not opened:
break
return accepted, duplicates, rejected

@staticmethod
def _close_sessions(
opened: list[UserDataStreamSession], blockers: list[str]
) -> None:
for session in reversed(opened):
try:
session.close()
except (OSError, RuntimeError, ValueError, WebSocketException):
blockers.append(f"{session.product_type.value}_WEBSOCKET_CLOSE_FAILED")

def _append_subscription_event(
self,
session: UserDataStreamSession,
Expand Down
Loading
Loading