//=========================================================================
// Summary  : */
// Filename : TaskServer.cc
// Author   : */
// Project  : */
// Revision : 1
// Created  : 2000/08/19
// Modified : 2000/08/19
//=========================================================================
// Description :
//=========================================================================
#include <sys/kernel.h>
#include <stdio.h>
#include <unistd.h>
#include <errno.h>
#include <signal.h>
#include <sys/sched.h>
#include "TaskServer.h"
#include "Task.h"
#include "TaskInterface.h"
#include "Syslog.h"
#include "System.h"
#include "VehicleConfigurator.h"

TaskServer *TaskServer::_theTaskServer = 0;

TaskServer::TaskServer(const char *name, int priority, int sched_policy)
  : SharedObjectServer(name)
{
  Boolean debug = False;
  char errorBuf[100];

  if (_theTaskServer != 0) {
    // Oops, we've already created a Task! Only 1 allowed per process
    sprintf(errorBuf, "Task::Task() - instance (\"%s\") already exists!",
	    _theTaskServer->name());

    throw Exception(errorBuf);
  }


  // Set up the Scheduling
  struct sched_param sched_param ;
  pid_t pid = getpid();

  int result = sched_getparam( pid, &sched_param );
  if (priority > 0) sched_param.sched_priority = priority;

  if(result == -1)
    {
      sprintf(errorBuf,"Task::Task() - failed to get scheduler info");
      throw Exception(errorBuf);
    }

  result = sched_setscheduler( pid, sched_policy, &sched_param );
  if(result == -1)
    {
      sprintf(errorBuf,"Task::Task() - failed to Set scheduler info");
      throw Exception(errorBuf);
    }


  // Keep track of the most recently created TaskServer for cleanup
  // on signal/exit
  _theTaskServer = this;

  _theTaskServerName = new char[70];

  strcpy(_theTaskServerName,_theTaskServer->name());

  // Client subscription request
  addCallback(SubscribeCode, (CallbackPtr )TaskServer::addSubscriber);

  // Event notification
  addCallback(NotifyCode, (CallbackPtr )TaskServer::handleEvent);

  dprintf("TaskServer::TaskServer() - this=0x%x", this);

  //signal(SIGCHLD, TaskServer::signalHandler);
  signal(SIGCHLD, SIG_IGN);
  signal(SIGINT, TaskServer::signalHandler);
  signal(SIGTERM, TaskServer::signalHandler);
  signal(SIGQUIT, TaskServer::signalHandler);

  atexit(TaskServer::cleanup);
}



TaskServer::~TaskServer()
{
  Boolean debug = False;
  dprintf("TaskServer::~TaskServer()");

  delete [] _theTaskServerName;
  _theTaskServer = 0;
}



Boolean TaskServer::addSubscriber(Request *request, int nRequestBytes,
				  Reply **reply, int *nReplyBytes)
{
	Boolean debug = False;

	Subscription *subscription = (Subscription *)request;

	dprintf("TaskServer::addSubscriber() - subscriber=%d\n",
		subscription->subscriberPid);

	if (registerSubscriber(subscription->eventCode,
		subscription->subscriberPid, subscription->otherNode)) {
		_subscriptionReply.confirmed = True;
	}
	else
		_subscriptionReply.confirmed = False;

	*reply = &_subscriptionReply;
	*nReplyBytes = sizeof(SubscriptionReply);

	return True;

}


Boolean TaskServer::registerSubscriber(EventCode eventCode, pid_t subscriber, Boolean otherNode)
{
  Boolean debug = False;
  Boolean found = False;

  FwdEventTableSubscriber fwdSubscriber;
  fwdSubscriber.pid = subscriber;
  fwdSubscriber.otherNode = otherNode;

  // Add subscriber to list for this event
  FwdEventTableEntry *entry;
  for (int i = 0; i < _fwdEventTable.size(); i++) {

    _fwdEventTable.get(i, &entry);

    if (entry->eventCode == eventCode) {
      // Table contains entry for eventCode; add to subscriber list
      dprintf("TaskServer::registerSubscriber() - already have subscribers to "
	      "event %d; add this one\n", eventCode);

      entry->subscriber.add(&fwdSubscriber);
      found = True;
      break;
    }
  }

  if (!found) {
    // Event is not yet in subscriber table
    dprintf("TaskServer::registerSubscriber() - "
	    "add event/subscriber to table\n");
    entry = new FwdEventTableEntry();
    entry->eventCode = eventCode;
    // Add to subscriber list
    entry->subscriber.add(&fwdSubscriber);
    _fwdEventTable.add(&entry);
  }

  return True;
}

void TaskServer::subscribe(TaskInterface *client,
			   EventCode eventCode,
			   EventCallback callback)
{
  /* IMPORTANT! Note: This function is copied, byte-for-byte, in EventService::subscribe().
     If you change something here, change it in EventService::subscribe(), too! */

  int publisherNode;
  pid_t localProxyPid=-1, remoteProxyPid=-1;
  Boolean debug = False;

  dprintf("TaskServer::subscribe() - pid %d subscribing to server %s, event code %d.",
  	getpid(), client->serverName(), eventCode);

  // First, let's find out if we've already subscribed to this event. It's not likely, but if
  //  there's some piece of code that just subscribes to something on another node willy-nilly,
  //  we'll end up with tons of remote proxies unless we check. Here goes...
  for (int i = 0; i < _eventTable.size(); i++) {
	  EventTableEntry *entry;
	  _eventTable.get(i, &entry);
	  if ( ( entry->_eventCode == eventCode ) && ( entry->_callback == callback )
		&& ( !strcmp( entry->_client->serverName(), client->serverName() ) ) )
	  {
		  // Good news, we've already subscribed!
		  dprintf("EventService::subscribe() - already subscribed to event %d from %s.\n",
			eventCode, client->serverName() );
		  return;
	  }
  }

  publisherNode = VehicleConfigurator::getNodeFromIFName( (char *)client->serverName() );

  if ( publisherNode != getnid() )
  {
	  dprintf("TaskServer::subscribe() - publisher is on another node.");

      TaskServer::Event *event = new TaskServer::Event();

      event->eventCode = eventCode;
      event->serverPid = VehicleConfigurator::getServerPidFromIFName( (char *)client->serverName() );

      // The node stored in the event is the node the event comes from. Hence, this isn't getnid(), it's publisherNode.
      event->node = publisherNode;

      // Ok, so this is confusing. The event object that we're creating here is put into memory.
      // This memory is read when the event is triggered. So, when this proxy is triggered
      event->forward = False;

	  if ( ( localProxyPid = qnx_proxy_attach(0, event, sizeof(TaskServer::Event), -1) ) == -1 )
	  {
		  Syslog::write("TaskServer::subscribe() - couldn't attach local proxy - errno %d. Not subscribing.", errno);
		  delete event;
		  return;
	  }
	  delete event;

      if ( ( remoteProxyPid = qnx_proxy_rem_attach(publisherNode, localProxyPid) ) == -1 )
      {
          Syslog::write("TaskServer::subscribe() - couldn't attach remote proxy - errno %d. Not subscribing.\n", errno);
          qnx_proxy_detach(localProxyPid);
	      return;
	  }
	  dprintf("TaskServer::subscribe() - local proxy = %d, remote = %d.", localProxyPid, remoteProxyPid);
  }

  dprintf("TaskServer::subscribe() - about to call TaskInterface::subscribe(%d, %d).", eventCode, remoteProxyPid);
  client->subscribe(eventCode, remoteProxyPid);

  EventTableEntry *entry = new EventTableEntry(client, eventCode, callback, localProxyPid, remoteProxyPid);
  _eventTable.add(&entry);
}


Boolean TaskServer::handleEvent(Request *request, int nRequestBytes,
				Reply **reply, int *nReplyBytes)
{
  Boolean debug = False;

  Event *event = (Event *)request;

  dprintf("TaskServer::handleEvent() - "
	  "clientPid=%d, reqCode=%d, eventCode=%d, pid=%d, forward=%d node=%d",
	  clientPid(), event->code, event->eventCode, event->serverPid,
	  event->forward, event->node);

  if (event->forward) {

    // Forward event to subscribers
    forwardEvent(event);
  }
  else {
    // Event's destination is here (i.e. server subscribed to it,
    // through a contained TaskInterface)

    // If the event is on the same node as us, then we did a qnx_proxy_attach in forwardEvent
    // and we need to detach it to free memory.
    if ( event->node == getnid() )
    {
		dprintf("TaskServer::handleEvent() - detaching proxy %d.", clientPid() );
        if (qnx_proxy_detach(clientPid()) == -1) {
            Syslog::write("TaskServer::handleEvent() - %s", strerror(errno));
        }
    }

    invokeEventCallback(event);
  }
  *reply = 0;
  *nReplyBytes = 0;

  // Don't reply to proxy
  return False;
}


void TaskServer::triggerEvent(EventCode eventCode)
{
  // Create event and forward to subscribers
  Event *event = new Event();
  event->eventCode = eventCode;
  event->node = getnid();

  forwardEvent(event);

  delete event;
}


void TaskServer::forwardEvent(Event *event)
{
  Boolean debug = False;
  pid_t proxy;

  // Tag event with server pid
  event->serverPid = getpid();

  // Sending to subscribers, which is final destination
  event->forward = False;

  Boolean found = False;

  // Find list of subscribers for this event
  for (int i = 0; i < _fwdEventTable.size(); i++) {

    FwdEventTableEntry *entry;

    _fwdEventTable.get(i, &entry);

    if (entry->eventCode == event->eventCode) {

      found = True;

      // Notify all clients who subscribed to this event
      for (int j = 0; j < entry->subscriber.size(); j++) {
	pid_t processId;
	FwdEventTableSubscriber fwdSubscriber;
	entry->subscriber.get(j, &fwdSubscriber);

	processId = fwdSubscriber.pid;

	if ( fwdSubscriber.otherNode )
	    proxy = fwdSubscriber.pid;
	else
	    proxy = qnx_proxy_attach(processId, event, sizeof(Event), -1);

	if (proxy == -1) {
	     if( errno == ESRCH ) {
		  // This error will occur when the subscribing process
		  // no longer exists (i.e. test programs, etc).  The
		  // efficient thing to do would be to remove the subscriber
		  // from the subscribe list.  This is not efficient.
		  //
		  // ...
		  ;
	     } else {
		  Syslog::write("TaskServer::forwardEvent() - %s", strerror(errno));
	     }
	}
	else {
	  dprintf("TaskServer::forwardEvent() - "
		  "trigger proxy %d, for pid %d",
		  proxy, processId);

	  Trigger(proxy);
	}
      }

      return;
    }
  }

  if (!found) {
	  dprintf("TaskServer::forwardEvent(), %s - "
		  "no subscribers found for event code %d",
		  name(), event->eventCode);
  }
}


void TaskServer::invokeEventCallback(Event *event)
{
  Boolean debug = False;

  // Determine callback which corresponds to this notification
  Boolean found = False;
  for (int i = 0; i < _eventTable.size(); i++) {
    EventTableEntry *entry;
    _eventTable.get(i, &entry);

    //dprintf("TaskServer::invokeEventCallback() - "
	//    "entry eventCode=%d, serverPid=%d\n",
	//    entry->_eventCode, VehicleConfigurator::getServerPidFromIFName((char *)entry->_client->serverName()) );

    if (event->eventCode == entry->_eventCode &&
	event->serverPid == VehicleConfigurator::getServerPidFromIFName( (char *)entry->_client->serverName() ) ) {

      // Invoke callback
      dprintf("TaskServer::invokeEventCallback() - found callback\n");
      callMemberFunction(this, entry->_callback)(entry->_client,
						 event->eventCode);
      found = True;
    }
  }
  if (!found) {
    Syslog::write("TaskServer::invokeEventCallback(), %s - unknown event, "
		  "code=%d, pid=%d\n",
		  name(), event->eventCode, event->serverPid);
  }
}



void TaskServer::run()
{

  SharedObjectServer::run();
}


void TaskServer::signalHandler(int sigNo)
{
  Boolean debug = False;

  if (!_theTaskServer) {
    // Doesn't exist
    dprintf("TaskServer::signalHandler() - _theTaskServer doesn't exist");
    return;
  }

  dprintf("TaskServer::signalHandler() for task %s - exit()",
	  _theTaskServer->name());

  // Invoke exit(), which will invoke cleanup()
  exit(1);
}


void TaskServer::cleanup()
{
  Boolean debug = False;

  char serverName[100];

  serverName[0] = 0x00;

  //strcpy(serverName, name());

  dprintf("TaskServer::cleanup() - delete _theTaskServer (%s)", serverName);

  delete _theTaskServer;

  dprintf("TaskServer::cleanup() - done with %s", serverName);
  return;
}


size_t TaskServer::maxRequestBytes()
{
  return sizeof(TaskServer::Event);
}

const char* TaskServer::getSerializedInterfaceDescription(int iMethod)
{
  // By default, TaskServer has no IDL calls.
  return 0;
}

