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
4 changes: 3 additions & 1 deletion .run/Desktop.run.xml
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
<component name="ProjectRunConfigurationManager">
<configuration default="false" name="Desktop" type="Application" factoryName="Application">
<option name="ALTERNATIVE_JRE_PATH" value="11" />
<option name="ALTERNATIVE_JRE_PATH_ENABLED" value="true" />
<option name="MAIN_CLASS_NAME" value="ai.flow.app.lwjgl3.Lwjgl3Launcher" />
<module name="flow-pilot.desktop.main" />
<option name="WORKING_DIRECTORY" value="$PROJECT_DIR$/assets" />
<method v="2" />
</configuration>
</component>
</component>
2 changes: 1 addition & 1 deletion cereal
30 changes: 30 additions & 0 deletions common/api/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
import os
import json

from common.params import Params

import urllib3

API_HOST = os.getenv('API_HOST', 'https://api.flowdrive.ai')

class Api():
def __init__(self):
self.params = Params()

self.http_client = urllib3.PoolManager()

def get_credentials(self):
# get userdata
self.email = self.params.get("Email")
self.token = self.params.get("Token")

# Get STS
r = self.http_client.request(
'POST',
f"{API_HOST}/auth/sts",
fields={'email': self.email, 'token': self.token}
)
print("status", r.status)

credentials = json.loads(r.data.decode('utf-8'))
return credentials
3 changes: 3 additions & 0 deletions common/params.cc
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,9 @@ std::unordered_map<std::string, uint32_t> keys = {
{"Offroad_TemperatureTooHigh", CLEAR_ON_MANAGER_START},
{"Offroad_UnofficialHardware", CLEAR_ON_MANAGER_START},
{"Offroad_UpdateFailed", CLEAR_ON_MANAGER_START},
{"IsOffroad", PERSISTENT},
{"Email", PERSISTENT},
{"Token", PERSISTENT},
};

lmdb::env Params::env = nullptr;
Expand Down
8 changes: 4 additions & 4 deletions common/realtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
import os
import time
import multiprocessing
from typing import Optional
from typing import Optional, List
from common.clock import sec_since_boot
from collections import deque
from selfdrive.swaglog import cloudlog
Expand Down Expand Up @@ -35,17 +35,17 @@ def set_realtime_priority(level: int) -> None:
cloudlog.info("Unable to set realtime priority")


def set_core_affinity(core: int) -> None:
def set_core_affinity(cores: List[int]) -> None:
try:
os.sched_setaffinity(0, [core,]) # type: ignore[attr-defined]
os.sched_setaffinity(0, cores) # type: ignore[attr-defined]
except:
cloudlog.info("Unable to set core affinity priority")


def config_realtime_process(core: int, priority: int) -> None:
gc.disable()
set_realtime_priority(priority)
set_core_affinity(core)
set_core_affinity([core])


class Ratekeeper:
Expand Down
2 changes: 1 addition & 1 deletion flowpilot_env.sh
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
ARCHNAME=$(arch)
ARCHNAME=$(uname -m)
SCRIPT=$(realpath "$0")
FLOWPILOT_DIR=$(dirname "$SCRIPT")

Expand Down
2 changes: 1 addition & 1 deletion launch_flowpilot.sh
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
set -e
source .env
source ./.env

# build changes
scons
Expand Down
59 changes: 38 additions & 21 deletions selfdrive/loggerd/uploader.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,14 +12,20 @@

from cereal import log
import cereal.messaging as messaging
#from common.api import Api
from common.api import Api
from common.params import Params
from common.realtime import set_core_affinity
from system.hardware import TICI
from selfdrive.loggerd.xattr_cache import getxattr, setxattr
from selfdrive.loggerd.config import ROOT
from selfdrive.swaglog import cloudlog

import boto3
from os.path import relpath
import os
import threading


NetworkType = log.DeviceState.NetworkType
UPLOAD_ATTR_NAME = 'user.upload'
UPLOAD_ATTR_VALUE = b'1'
Expand Down Expand Up @@ -57,7 +63,7 @@ def clear_locks(root):
class Uploader():
def __init__(self, dongle_id, root):
self.dongle_id = dongle_id
#self.api = Api(dongle_id)
self.api = Api()
self.root = root

self.upload_thread = None
Expand Down Expand Up @@ -134,24 +140,27 @@ def next_file_to_upload(self):

def do_upload(self, key, fn):
try:
url_resp = self.api.get("v1.4/" + self.dongle_id + "/upload_url/", timeout=10, path=key, access_token=self.api.get_token())
if url_resp.status_code == 412:
self.last_resp = url_resp
return
credentials = self.api.get_credentials()

url_resp_json = json.loads(url_resp.text)
url = url_resp_json['url']
headers = url_resp_json['headers']
cloudlog.debug("upload_url v1.4 %s %s", url, str(headers))
access_key = credentials["access_key"]
secret_access_key = credentials["secret_access_key"]
session_token = credentials["session_token"]

if fake_upload:
cloudlog.debug(f"*** WARNING, THIS IS A FAKE UPLOAD TO {url} ***")
s3=boto3.client(
's3',
aws_access_key_id=access_key,
aws_secret_access_key=secret_access_key,
aws_session_token=session_token,
)

class FakeResponse():
def __init__(self):
self.status_code = 200
class FakeResponse():
def __init__(self):
self.status_code = 200

if fake_upload:
cloudlog.debug(f"*** WARNING, THIS IS A FAKE UPLOAD ***")
self.last_resp = FakeResponse()

else:
with open(fn, "rb") as f:
if key.endswith('.bz2') and not fn.endswith('.bz2'):
Expand All @@ -160,7 +169,14 @@ def __init__(self):
else:
data = f

self.last_resp = requests.put(url, data=data, headers=headers, timeout=10)
# api.get_credentials should populate api.email field, saving us a DB call
object_name = self.api.email.decode("utf-8") + "/" + relpath(fn, ROOT)
bucket = "fdusermedia"

self.last_resp = FakeResponse()
s3.upload_fileobj(data, bucket, object_name)
s3.Object(bucket, object_name).wait_until_exists()

except Exception as e:
self.last_exc = (e, traceback.format_exc())
raise
Expand Down Expand Up @@ -237,7 +253,6 @@ def uploader_fn(exit_event):

if dongle_id is None:
cloudlog.info("uploader missing dongle_id")
raise Exception("uploader can't start without dongle id")

if TICI and not Path("/data/media").is_mount():
cloudlog.warning("NVME not mounted")
Expand All @@ -250,11 +265,13 @@ def uploader_fn(exit_event):
while not exit_event.is_set():
sm.update(0)
offroad = params.get_bool("IsOffroad")
offroad = False
network_type = sm['deviceState'].networkType if not force_wifi else NetworkType.wifi
if network_type == NetworkType.none:
if allow_sleep:
time.sleep(60 if offroad else 5)
continue
# Start on any network type
# if network_type == NetworkType.none:
# if allow_sleep:
# time.sleep(60 if offroad else 5)
# continue

d = uploader.next_file_to_upload()
if d is None: # Nothing to upload
Expand Down
7 changes: 5 additions & 2 deletions selfdrive/ui/java/ai.flow.app/LoginScreen.java
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
import com.badlogic.gdx.utils.Align;
import com.badlogic.gdx.utils.viewport.FitViewport;


import java.net.HttpCookie;
import java.util.HashMap;
import java.util.List;
Expand Down Expand Up @@ -142,7 +143,7 @@ public void handleHttpResponse(Net.HttpResponse httpResponse) {
try {
String cookieHeader = httpResponse.getHeader("set-cookie");
List<HttpCookie> cookies = HttpCookie.parse(cookieHeader);
LoginSucceeded(cookies);
LoginSucceeded(email, token, cookies);

} catch (Exception exception) {
progressVal = 0;
Expand All @@ -162,8 +163,10 @@ public void cancelled() {
});
}

private void LoginSucceeded(List<HttpCookie> cookies) {
private void LoginSucceeded(String email, String token, List<HttpCookie> cookies) {
appContext.params.put("UserID", cookies.get(0).getValue());
appContext.params.put("Email", email);
appContext.params.put("Token", token);
progressVal++;
}

Expand Down