pio_server.c 40.8 KB
Newer Older
Deike Kleberg's avatar
Deike Kleberg committed
1
2
/** @file ioServer.c
*/
3
4
5
6
#ifdef HAVE_CONFIG_H
#  include "config.h"
#endif

Deike Kleberg's avatar
Deike Kleberg committed
7
8
9
#include "pio_server.h"


10
#include <limits.h>
Deike Kleberg's avatar
Deike Kleberg committed
11
12
#include <stdlib.h>
#include <stdio.h>
13
14
15
16
17
18

#ifdef HAVE_PARALLEL_NC4
#include <core/ppm_combinatorics.h>
#include <core/ppm_rectilinear.h>
#include <ppm/ppm_uniform_partition.h>
#endif
19
#include <yaxt.h>
20

Deike Kleberg's avatar
Deike Kleberg committed
21
#include "cdi.h"
22
#include "cdipio.h"
23
#include "dmemory.h"
24
#include "namespace.h"
25
#include "taxis.h"
Deike Kleberg's avatar
Deike Kleberg committed
26
#include "pio.h"
Deike Kleberg's avatar
Deike Kleberg committed
27
#include "pio_comm.h"
28
#include "pio_interface.h"
Deike Kleberg's avatar
Deike Kleberg committed
29
#include "pio_rpc.h"
Deike Kleberg's avatar
Deike Kleberg committed
30
#include "pio_util.h"
31
#include "cdi_int.h"
32
33
34
#ifndef HAVE_NETCDF_PAR_H
#define MPI_INCLUDED
#endif
35
#include "pio_cdf_int.h"
36
#include "resource_handle.h"
37
#include "resource_unpack.h"
Thomas Jahns's avatar
Thomas Jahns committed
38
#include "stream_cdf.h"
Deike Kleberg's avatar
Deike Kleberg committed
39
#include "vlist_var.h"
40

41

42
extern void arrayDestroy ( void );
Deike Kleberg's avatar
Deike Kleberg committed
43

44
45
46
static struct
{
  size_t size;
Thomas Jahns's avatar
Thomas Jahns committed
47
  unsigned char *buffer;
48
  int dictSize;
49
50
} *rxWin = NULL;

Thomas Jahns's avatar
Thomas Jahns committed
51
static MPI_Win getWin = MPI_WIN_NULL;
Thomas Jahns's avatar
Thomas Jahns committed
52
static MPI_Group groupModel = MPI_GROUP_NULL;
Deike Kleberg's avatar
Deike Kleberg committed
53

54
55
56
57
58
#ifdef HAVE_PARALLEL_NC4
/* prime factorization of number of pio collectors */
static uint32_t *pioPrimes;
static int numPioPrimes;
#endif
Deike Kleberg's avatar
Deike Kleberg committed
59

Deike Kleberg's avatar
Deike Kleberg committed
60
61
/************************************************************************/

62
static
Deike Kleberg's avatar
Deike Kleberg committed
63
64
void serverWinCleanup ()
{
65
66
  if (getWin != MPI_WIN_NULL)
    xmpi(MPI_Win_free(&getWin));
67
68
  if (rxWin)
    {
69
      free(rxWin[0].buffer);
70
      free(rxWin);
Deike Kleberg's avatar
Deike Kleberg committed
71
    }
72

73
  xdebug("%s", "cleaned up mpi_win");
Deike Kleberg's avatar
Deike Kleberg committed
74
}
75

Deike Kleberg's avatar
Deike Kleberg committed
76
 /************************************************************************/
77

78
79
static size_t
collDefBufferSizes()
Deike Kleberg's avatar
Deike Kleberg committed
80
{
81
  int nstreams, * streamIndexList, streamNo, vlistID, nvars, varID, iorank;
82
83
  int modelID;
  size_t sumGetBufferSizes = 0;
84
  int rankGlob = commInqRankGlob ();
Deike Kleberg's avatar
Deike Kleberg committed
85
  int nProcsModel = commInqNProcsModel ();
86
  int root = commInqRootGlob ();
Deike Kleberg's avatar
Deike Kleberg committed
87

88
  xassert(rxWin != NULL);
Deike Kleberg's avatar
Deike Kleberg committed
89

Deike Kleberg's avatar
Deike Kleberg committed
90
  nstreams = reshCountType ( &streamOps );
Uwe Schulzweida's avatar
Uwe Schulzweida committed
91
  streamIndexList = (int*) xmalloc( nstreams * sizeof ( streamIndexList[0] ));
92
  reshGetResHListOfType ( nstreams, streamIndexList, &streamOps );
Deike Kleberg's avatar
Deike Kleberg committed
93
94
  for ( streamNo = 0; streamNo < nstreams; streamNo++ )
    {
95
      // space required for data
96
      vlistID = streamInqVlist ( streamIndexList[streamNo] );
Deike Kleberg's avatar
Deike Kleberg committed
97
98
99
100
      nvars = vlistNvars ( vlistID );
      for ( varID = 0; varID < nvars; varID++ )
        {
          iorank = vlistInqVarIOrank ( vlistID, varID );
Deike Kleberg's avatar
Deike Kleberg committed
101
          xassert ( iorank != CDI_UNDEFID );
Deike Kleberg's avatar
Deike Kleberg committed
102
103
          if ( iorank == rankGlob )
            {
Deike Kleberg's avatar
Deike Kleberg committed
104
              for ( modelID = 0; modelID < nProcsModel; modelID++ )
105
                {
106
107
108
109
110
111
112
113
                  int decoChunk;
                  {
                    int varSize = vlistInqVarSize(vlistID, varID);
                    int nProcsModel = commInqNProcsModel();
                    decoChunk =
                      (int)ceilf(cdiPIOpartInflate_
                                 * (varSize + nProcsModel - 1)/nProcsModel);
                  }
Deike Kleberg's avatar
Deike Kleberg committed
114
                  xassert ( decoChunk > 0 );
115
                  rxWin[modelID].size += decoChunk * sizeof (double)
116
117
118
119
                    /* re-align chunks to multiple of double size */
                    + sizeof (double) - 1
                    /* one header for data record, one for
                     * corresponding part descriptor*/
120
                    + 2 * sizeof (struct winHeaderEntry)
121
122
123
                    /* FIXME: heuristic for size of packed Xt_idxlist */
                    + sizeof (Xt_int) * decoChunk * 3;
                  rxWin[modelID].dictSize += 2;
124
                }
Deike Kleberg's avatar
Deike Kleberg committed
125
            }
126
        }
Deike Kleberg's avatar
Deike Kleberg committed
127
128
      // space required for the 3 function calls streamOpen, streamDefVlist, streamClose 
      // once per stream and timestep for all collprocs only on the modelproc root
129
      rxWin[root].size += numRPCFuncs * sizeof (struct winHeaderEntry)
130
131
132
133
        /* serialized filename */
        + MAXDATAFILENAME
        /* data part of streamDefTimestep */
        + (2 * CDI_MAX_NAME + sizeof (taxis_t));
134
      rxWin[root].dictSize += numRPCFuncs;
Deike Kleberg's avatar
Deike Kleberg committed
135
    }
136
  free ( streamIndexList );
Deike Kleberg's avatar
Deike Kleberg committed
137
138

  for ( modelID = 0; modelID < nProcsModel; modelID++ )
139
    {
140
      /* account for size header */
141
      rxWin[modelID].dictSize += 1;
142
      rxWin[modelID].size += sizeof (struct winHeaderEntry);
143
144
145
      rxWin[modelID].size = roundUpToMultiple(rxWin[modelID].size,
                                              PIO_WIN_ALIGN);
      sumGetBufferSizes += (size_t)rxWin[modelID].size;
146
    }
Deike Kleberg's avatar
Deike Kleberg committed
147
  xassert ( sumGetBufferSizes <= MAXWINBUFFERSIZE );
148
  return sumGetBufferSizes;
Deike Kleberg's avatar
Deike Kleberg committed
149
}
150

Deike Kleberg's avatar
Deike Kleberg committed
151
 /************************************************************************/
152

153
154
155
static void
serverWinCreate(void)
{
Deike Kleberg's avatar
Deike Kleberg committed
156
  int ranks[1], modelID;
157
  MPI_Comm commCalc = commInqCommCalc ();
Deike Kleberg's avatar
Deike Kleberg committed
158
  MPI_Group groupCalc;
159
  int nProcsModel = commInqNProcsModel ();
160
161
162
  MPI_Info no_locks_info;
  xmpi(MPI_Info_create(&no_locks_info));
  xmpi(MPI_Info_set(no_locks_info, "no_locks", "true"));
Deike Kleberg's avatar
Deike Kleberg committed
163

164
  xmpi(MPI_Win_create(MPI_BOTTOM, 0, 1, no_locks_info, commCalc, &getWin));
Deike Kleberg's avatar
Deike Kleberg committed
165
166

  /* target group */
167
168
  ranks[0] = nProcsModel;
  xmpi ( MPI_Comm_group ( commCalc, &groupCalc ));
Deike Kleberg's avatar
Deike Kleberg committed
169
170
  xmpi ( MPI_Group_excl ( groupCalc, 1, ranks, &groupModel ));

171
  rxWin = xcalloc(nProcsModel, sizeof (rxWin[0]));
172
  size_t totalBufferSize = collDefBufferSizes();
Uwe Schulzweida's avatar
Uwe Schulzweida committed
173
  rxWin[0].buffer = (unsigned char*) xmalloc(totalBufferSize);
174
175
176
177
178
179
  size_t ofs = 0;
  for ( modelID = 1; modelID < nProcsModel; modelID++ )
    {
      ofs += rxWin[modelID - 1].size;
      rxWin[modelID].buffer = rxWin[0].buffer + ofs;
    }
Deike Kleberg's avatar
Deike Kleberg committed
180

181
182
  xmpi(MPI_Info_free(&no_locks_info));

183
  xdebug("%s", "created mpi_win, allocated getBuffer");
Deike Kleberg's avatar
Deike Kleberg committed
184
185
}

Deike Kleberg's avatar
Deike Kleberg committed
186
187
/************************************************************************/

188
static void
189
readFuncCall(struct winHeaderEntry *header)
Deike Kleberg's avatar
Deike Kleberg committed
190
191
{
  int root = commInqRootGlob ();
192
  int funcID = header->id;
193
  union funcArgs *funcArgs = &(header->specific.funcArgs);
Deike Kleberg's avatar
Deike Kleberg committed
194

195
  xassert(funcID >= MINFUNCID && funcID <= MAXFUNCID);
Deike Kleberg's avatar
Deike Kleberg committed
196
197
  switch ( funcID )
    {
198
199
    case STREAMCLOSE:
      {
200
        int streamID
201
          = namespaceAdaptKey2(funcArgs->streamChange.streamID);
202
203
204
205
        streamClose(streamID);
        xdebug("READ FUNCTION CALL FROM WIN:  %s, streamID=%d,"
               " closed stream",
               funcMap[(-1 - funcID)], streamID);
206
207
      }
      break;
Deike Kleberg's avatar
Deike Kleberg committed
208
    case STREAMOPEN:
209
      {
210
        size_t filenamesz = funcArgs->newFile.fnamelen;
Deike Kleberg's avatar
Deike Kleberg committed
211
        xassert ( filenamesz > 0 && filenamesz < MAXDATAFILENAME );
212
        const char *filename
213
          = (const char *)(rxWin[root].buffer + header->offset);
214
        xassert(filename[filenamesz] == '\0');
215
        int filetype = funcArgs->newFile.filetype;
216
        int streamID = streamOpenWrite(filename, filetype);
217
        xassert(streamID != CDI_ELIBNAVAIL);
218
219
        xdebug("READ FUNCTION CALL FROM WIN:  %s, filenamesz=%zu,"
               " filename=%s, filetype=%d, OPENED STREAM %d",
220
               funcMap[(-1 - funcID)], filenamesz, filename,
221
               filetype, streamID);
222
      }
223
      break;
224
225
    case STREAMDEFVLIST:
      {
226
        int streamID
227
228
          = namespaceAdaptKey2(funcArgs->streamChange.streamID);
        int vlistID = namespaceAdaptKey2(funcArgs->streamChange.vlistID);
229
230
231
232
        streamDefVlist(streamID, vlistID);
        xdebug("READ FUNCTION CALL FROM WIN:  %s, streamID=%d,"
               " vlistID=%d, called streamDefVlist ().",
               funcMap[(-1 - funcID)], streamID, vlistID);
233
234
      }
      break;
235
236
237
    case STREAMDEFTIMESTEP:
      {
        MPI_Comm commCalc = commInqCommCalc ();
238
        int streamID = funcArgs->streamNewTimestep.streamID;
239
        int originNamespace = namespaceResHDecode(streamID).nsp;
240
241
242
        streamID = namespaceAdaptKey2(streamID);
        int oldTaxisID
          = vlistInqTaxis(streamInqVlist(streamID));
243
        int position = header->offset;
244
245
        int changedTaxisID
          = taxisUnpack((char *)rxWin[root].buffer, (int)rxWin[root].size,
246
                        &position, originNamespace, &commCalc, 0);
247
248
249
250
        taxis_t *oldTaxisPtr = taxisPtr(oldTaxisID);
        taxis_t *changedTaxisPtr = taxisPtr(changedTaxisID);
        ptaxisCopy(oldTaxisPtr, changedTaxisPtr);
        taxisDestroy(changedTaxisID);
251
        streamDefTimestep(streamID, funcArgs->streamNewTimestep.tsID);
252
253
      }
      break;
Deike Kleberg's avatar
Deike Kleberg committed
254
    default:
255
      xabort ( "REMOTE FUNCTIONCALL NOT IMPLEMENTED!" );
Deike Kleberg's avatar
Deike Kleberg committed
256
257
258
259
260
    }
}

/************************************************************************/

261
262
263
264
265
static void
resizeVarGatherBuf(int vlistID, int varID, double **buf, int *bufSize)
{
  int size = vlistInqVarSize(vlistID, varID);
  if (size <= *bufSize) ; else
Uwe Schulzweida's avatar
Uwe Schulzweida committed
266
    *buf = (double*) xrealloc(*buf, (*bufSize = size) * sizeof (buf[0][0]));
267
268
269
270
271
272
273
}

static void
gatherArray(int root, int nProcsModel, int headerIdx,
            int vlistID,
            double *gatherBuf, int *nmiss)
{
274
275
  struct winHeaderEntry *winDict
    = (struct winHeaderEntry *)rxWin[root].buffer;
276
  int streamID = winDict[headerIdx].id;
277
  int varID = winDict[headerIdx].specific.dataRecord.varID;
278
  int varShape[3] = { 0, 0, 0 };
279
  cdiPioQueryVarDims(varShape, vlistID, varID);
280
281
282
283
284
  Xt_int varShapeXt[3];
  static const Xt_int origin[3] = { 0, 0, 0 };
  for (unsigned i = 0; i < 3; ++i)
    varShapeXt[i] = varShape[i];
  int varSize = varShape[0] * varShape[1] * varShape[2];
Uwe Schulzweida's avatar
Uwe Schulzweida committed
285
286
  struct Xt_offset_ext *partExts = (struct Xt_offset_ext*) xmalloc(nProcsModel * sizeof (partExts[0]));
  Xt_idxlist *part = (Xt_idxlist*) xmalloc(nProcsModel * sizeof (part[0]));
287
288
  MPI_Comm commCalc = commInqCommCalc();
  {
289
    int nmiss_ = 0;
290
291
292
    for (int modelID = 0; modelID < nProcsModel; modelID++)
      {
        struct dataRecord *dataHeader
293
294
295
296
          = &((struct winHeaderEntry *)
              rxWin[modelID].buffer)[headerIdx].specific.dataRecord;
        int position =
          ((struct winHeaderEntry *)rxWin[modelID].buffer)[headerIdx + 1].offset;
297
298
299
        xassert(namespaceAdaptKey2(((struct winHeaderEntry *)
                                    rxWin[modelID].buffer)[headerIdx].id)
                == streamID
300
                && dataHeader->varID == varID
301
302
                && ((struct winHeaderEntry *)
                    rxWin[modelID].buffer)[headerIdx + 1].id == PARTDESCMARKER
303
304
                && position > 0
                && ((size_t)position
305
                    >= sizeof (struct winHeaderEntry) * rxWin[modelID].dictSize)
306
307
308
309
                && ((size_t)position < rxWin[modelID].size));
        part[modelID] = xt_idxlist_unpack(rxWin[modelID].buffer,
                                          (int)rxWin[modelID].size,
                                          &position, commCalc);
310
        int partSize = xt_idxlist_get_num_indices(part[modelID]);
311
312
313
        size_t charOfs = (rxWin[modelID].buffer
                          + ((struct winHeaderEntry *)
                             rxWin[modelID].buffer)[headerIdx].offset)
314
315
316
317
          - rxWin[0].buffer;
        xassert(charOfs % sizeof (double) == 0
                && charOfs / sizeof (double) + partSize <= INT_MAX);
        int elemOfs = charOfs / sizeof (double);
318
319
320
        partExts[modelID].start = elemOfs;
        partExts[modelID].size = partSize;
        partExts[modelID].stride = 1;
321
322
323
324
325
        nmiss_ += dataHeader->nmiss;
      }
    *nmiss = nmiss_;
  }
  Xt_idxlist srcList = xt_idxlist_collection_new(part, nProcsModel);
326
  for (int modelID = 0; modelID < nProcsModel; modelID++)
327
328
329
330
331
332
333
334
    xt_idxlist_delete(part[modelID]);
  free(part);
  Xt_xmap gatherXmap;
  {
    Xt_idxlist dstList
      = xt_idxsection_new(0, 3, varShapeXt, varShapeXt, origin);
    struct Xt_com_list full = { .list = dstList, .rank = 0 };
    gatherXmap = xt_xmap_intersection_new(1, &full, 1, &full, srcList, dstList,
335
                                          MPI_COMM_SELF);
336
337
338
339
    xt_idxlist_delete(dstList);
  }
  xt_idxlist_delete(srcList);

340
  struct Xt_offset_ext gatherExt = { .start = 0, .size = varSize, .stride = 1 };
341
  Xt_redist gatherRedist
342
343
    = xt_redist_p2p_ext_new(gatherXmap, nProcsModel, partExts, 1, &gatherExt,
                            MPI_DOUBLE);
344
  xt_xmap_delete(gatherXmap);
345
  xt_redist_s_exchange1(gatherRedist, rxWin[0].buffer, gatherBuf);
346
  free(partExts);
347
  xt_redist_delete(gatherRedist);
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
}

struct xyzDims
{
  int sizes[3];
};

static inline int
xyzGridSize(struct xyzDims dims)
{
  return dims.sizes[0] * dims.sizes[1] * dims.sizes[2];
}

#ifdef HAVE_PARALLEL_NC4
static void
363
queryVarBounds(struct PPM_extent varShape[3], int vlistID, int varID)
364
{
365
366
  varShape[0].first = 0;
  varShape[1].first = 0;
367
  varShape[2].first = 0;
368
  int sizes[3];
369
  cdiPioQueryVarDims(sizes, vlistID, varID);
370
371
  for (unsigned i = 0; i < 3; ++i)
    varShape[i].size = sizes[i];
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
}

/* compute distribution of collectors such that number of collectors
 * <= number of variable grid cells in each dimension */
static struct xyzDims
varDimsCollGridMatch(const struct PPM_extent varDims[3])
{
  xassert(PPM_extents_size(3, varDims) >= commInqSizeColl());
  struct xyzDims collGrid = { { 1, 1, 1 } };
  /* because of storage order, dividing dimension 3 first is preferred */
  for (int i = 0; i < numPioPrimes; ++i)
    {
      for (int dim = 2; dim >=0; --dim)
        if (collGrid.sizes[dim] * pioPrimes[i] <= varDims[dim].size)
          {
            collGrid.sizes[dim] *= pioPrimes[i];
            goto nextPrime;
          }
      /* no position found, retrack */
      xabort("Not yet implemented back-tracking needed.");
      nextPrime:
      ;
    }
  return collGrid;
}

static void
myVarPart(struct PPM_extent varShape[3], struct xyzDims collGrid,
          struct PPM_extent myPart[3])
{
  int32_t myCollGridCoord[3];
  {
    struct PPM_extent collGridShape[3];
    for (int i = 0; i < 3; ++i)
      {
        collGridShape[i].first = 0;
        collGridShape[i].size = collGrid.sizes[i];
      }
    PPM_lidx2rlcoord_e(3, collGridShape, commInqRankColl(), myCollGridCoord);
    xdebug("my coord: (%d, %d, %d)", myCollGridCoord[0], myCollGridCoord[1],
           myCollGridCoord[2]);
  }
  PPM_uniform_partition_nd(3, varShape, collGrid.sizes,
                           myCollGridCoord, myPart);
}
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
#elif defined (HAVE_LIBNETCDF)
/* needed for writing when some files are only written to by a single process */
/* cdiOpenFileMap(fileID) gives the writer process */
int cdiPioSerialOpenFileMap(int streamID)
{
  return stream_to_pointer(streamID)->ownerRank;
}
/* for load-balancing purposes, count number of files per process */
/* cdiOpenFileCounts[rank] gives number of open files rank has to himself */
static int *cdiSerialOpenFileCount = NULL;
int cdiPioNextOpenRank()
{
  xassert(cdiSerialOpenFileCount != NULL);
  int commCollSize = commInqSizeColl();
  int minRank = 0, minOpenCount = cdiSerialOpenFileCount[0];
  for (int i = 1; i < commCollSize; ++i)
    if (cdiSerialOpenFileCount[i] < minOpenCount)
      {
        minOpenCount = cdiSerialOpenFileCount[i];
        minRank = i;
      }
  return minRank;
}

void cdiPioOpenFileOnRank(int rank)
{
  xassert(cdiSerialOpenFileCount != NULL
          && rank >= 0 && rank < commInqSizeColl());
  ++(cdiSerialOpenFileCount[rank]);
}


void cdiPioCloseFileOnRank(int rank)
{
  xassert(cdiSerialOpenFileCount != NULL
          && rank >= 0 && rank < commInqSizeColl());
  xassert(cdiSerialOpenFileCount[rank] > 0);
  --(cdiSerialOpenFileCount[rank]);
}

457
458
459
460
461
462
463
464
465
466
static void
cdiPioServerCdfDefVars(stream_t *streamptr)
{
  int rank, rankOpen;
  if (commInqIOMode() == PIO_NONE
      || ((rank = commInqRankColl())
          == (rankOpen = cdiPioSerialOpenFileMap(streamptr->self))))
    cdfDefVars(streamptr);
}

467
468
#endif

469
470
471
472
473
474
struct streamMapping {
  int streamID, filetype;
  int firstHeaderIdx, lastHeaderIdx;
  int numVars, *varMap;
};

475
476
477
478
479
480
struct streamMap
{
  struct streamMapping *entries;
  int numEntries;
};

Thomas Jahns's avatar
Thomas Jahns committed
481
482
483
484
485
486
487
488
static int
smCmpStreamID(const void *a_, const void *b_)
{
  const struct streamMapping *a = a_, *b = b_;
  int streamIDa = a->streamID, streamIDb = b->streamID;
  return (streamIDa > streamIDb) - (streamIDa < streamIDb);
}

489
490
491
492
493
494
495
static inline int
inventorizeStream(struct streamMapping *streamMap, int numStreamIDs,
                  int *sizeStreamMap_, int streamID, int headerIdx)
{
  int sizeStreamMap = *sizeStreamMap_;
  if (numStreamIDs < sizeStreamMap) ; else
    {
Uwe Schulzweida's avatar
Uwe Schulzweida committed
496
497
498
      streamMap = (struct streamMapping*) xrealloc(streamMap,
                                                   (sizeStreamMap *= 2)
                                                   * sizeof (streamMap[0]));
499
500
501
502
      *sizeStreamMap_ = sizeStreamMap;
    }
  streamMap[numStreamIDs].streamID = streamID;
  streamMap[numStreamIDs].firstHeaderIdx = headerIdx;
503
  streamMap[numStreamIDs].lastHeaderIdx = headerIdx;
504
505
506
507
508
509
510
511
512
  streamMap[numStreamIDs].numVars = -1;
  int filetype = streamInqFiletype(streamID);
  streamMap[numStreamIDs].filetype = filetype;
  if (filetype == FILETYPE_NC || filetype == FILETYPE_NC2
      || filetype == FILETYPE_NC4)
    {
      int vlistID = streamInqVlist(streamID);
      int nvars = vlistNvars(vlistID);
      streamMap[numStreamIDs].numVars = nvars;
Uwe Schulzweida's avatar
Uwe Schulzweida committed
513
      streamMap[numStreamIDs].varMap = (int*) xmalloc(sizeof (streamMap[numStreamIDs].varMap[0]) * nvars);
514
515
516
517
518
519
      for (int i = 0; i < nvars; ++i)
        streamMap[numStreamIDs].varMap[i] = -1;
    }
  return numStreamIDs + 1;
}

520
521
522
523
524
525
526
527
528
529
static inline int
streamIsInList(struct streamMapping *streamMap, int numStreamIDs,
               int streamIDQuery)
{
  int p = 0;
  for (int i = 0; i < numStreamIDs; ++i)
    p |= streamMap[i].streamID == streamIDQuery;
  return p;
}

530
static struct streamMap
531
buildStreamMap(struct winHeaderEntry *winDict)
532
533
534
{
  int streamIDOld = CDI_UNDEFID;
  int oldStreamIdx = CDI_UNDEFID;
535
  int filetype = FILETYPE_UNDEF;
536
  int sizeStreamMap = 16;
Uwe Schulzweida's avatar
Uwe Schulzweida committed
537
  struct streamMapping *streamMap = (struct streamMapping *) xmalloc(sizeStreamMap * sizeof (streamMap[0]));
538
  int numDataEntries = winDict[0].specific.headerSize.numDataEntries;
539
  int numStreamIDs = 0;
540
  /* find streams written on this process */
541
542
543
  for (int headerIdx = 1; headerIdx < numDataEntries; headerIdx += 2)
    {
      int streamID
544
545
        = winDict[headerIdx].id
        = namespaceAdaptKey2(winDict[headerIdx].id);
546
547
548
549
550
551
552
553
554
555
556
      xassert(streamID > 0);
      if (streamID != streamIDOld)
        {
          for (int i = numStreamIDs - 1; i >= 0; --i)
            if ((streamIDOld = streamMap[i].streamID) == streamID)
              {
                oldStreamIdx = i;
                goto streamIDInventorized;
              }
          oldStreamIdx = numStreamIDs;
          streamIDOld = streamID;
557
558
          numStreamIDs = inventorizeStream(streamMap, numStreamIDs,
                                           &sizeStreamMap, streamID, headerIdx);
559
560
        }
      streamIDInventorized:
561
      filetype = streamMap[oldStreamIdx].filetype;
562
563
564
565
      streamMap[oldStreamIdx].lastHeaderIdx = headerIdx;
      if (filetype == FILETYPE_NC || filetype == FILETYPE_NC2
          || filetype == FILETYPE_NC4)
        {
566
          int varID = winDict[headerIdx].specific.dataRecord.varID;
567
568
569
          streamMap[oldStreamIdx].varMap[varID] = headerIdx;
        }
    }
570
571
572
573
  /* join with list of streams written to in total */
  {
    int *streamIDs, *streamIsWritten;
    int numTotalStreamIDs = streamSize();
Uwe Schulzweida's avatar
Uwe Schulzweida committed
574
    streamIDs = (int*) xmalloc(2 * sizeof (streamIDs[0]) * (size_t)numTotalStreamIDs);
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
    streamGetIndexList(numTotalStreamIDs, streamIDs);
    streamIsWritten = streamIDs + numTotalStreamIDs;
    for (int i = 0; i < numTotalStreamIDs; ++i)
      streamIsWritten[i] = streamIsInList(streamMap, numStreamIDs,
                                          streamIDs[i]);
    /* Find what streams are written to at all on any process */
    xmpi(MPI_Allreduce(MPI_IN_PLACE, streamIsWritten, numTotalStreamIDs,
                       MPI_INT, MPI_BOR, commInqCommColl()));
    /* append streams written to on other tasks to mapping */
    for (int i = 0; i < numTotalStreamIDs; ++i)
      if (streamIsWritten[i] && !streamIsInList(streamMap, numStreamIDs,
                                                streamIDs[i]))
        numStreamIDs = inventorizeStream(streamMap, numStreamIDs,
                                         &sizeStreamMap, streamIDs[i], -1);

    free(streamIDs);
  }
Thomas Jahns's avatar
Thomas Jahns committed
592
  /* sort written streams by streamID */
Uwe Schulzweida's avatar
Uwe Schulzweida committed
593
  streamMap = (struct streamMapping*) xrealloc(streamMap, sizeof (streamMap[0]) * numStreamIDs);
Thomas Jahns's avatar
Thomas Jahns committed
594
  qsort(streamMap, numStreamIDs, sizeof (streamMap[0]), smCmpStreamID);
595
596
597
  return (struct streamMap){ .entries = streamMap, .numEntries = numStreamIDs };
}

598
599
600
601
602
603
604
605
606
607
static void
writeGribStream(struct winHeaderEntry *winDict, struct streamMapping *mapping,
                double **data_, int *currentDataBufSize, int root,
                int nProcsModel)
{
  int streamID = mapping->streamID;
  int headerIdx, lastHeaderIdx = mapping->lastHeaderIdx;
  int vlistID = streamInqVlist(streamID);
  if (lastHeaderIdx < 0)
    {
608
      /* write zero bytes to trigger synchronization code in fileWrite */
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
      cdiPioFileWrite(streamInqFileID(streamID), NULL, 0,
                      streamInqCurTimestepID(streamID));
    }
  else
    for (headerIdx = mapping->firstHeaderIdx;
         headerIdx <= lastHeaderIdx;
         headerIdx += 2)
      if (streamID == winDict[headerIdx].id)
        {
          int varID = winDict[headerIdx].specific.dataRecord.varID;
          int size = vlistInqVarSize(vlistID, varID);
          int nmiss;
          resizeVarGatherBuf(vlistID, varID, data_, currentDataBufSize);
          double *data = *data_;
          gatherArray(root, nProcsModel, headerIdx,
                      vlistID, data, &nmiss);
          streamWriteVar(streamID, varID, data, nmiss);
          if ( ddebug > 2 )
            {
              char text[1024];
              sprintf(text, "streamID=%d, var[%d], size=%d",
                      streamID, varID, size);
              xprintArray(text, data, size, DATATYPE_FLT);
            }
        }
}
635

636
637
638
639
640
641
642
#ifdef HAVE_NETCDF4
static void
buildWrittenVars(struct streamMapping *mapping, int **varIsWritten_,
                 int myCollRank, MPI_Comm collComm)
{
  int nvars = mapping->numVars;
  int *varMap = mapping->varMap;
Uwe Schulzweida's avatar
Uwe Schulzweida committed
643
  int *varIsWritten = *varIsWritten_ = (int*) xrealloc(*varIsWritten_, sizeof (*varIsWritten) * nvars);
644
645
646
647
648
649
650
  for (int varID = 0; varID < nvars; ++varID)
    varIsWritten[varID] = ((varMap[varID] != -1)
                           ?myCollRank+1 : 0);
  xmpi(MPI_Allreduce(MPI_IN_PLACE, varIsWritten, nvars,
                     MPI_INT, MPI_BOR, collComm));
}
#endif
651

652
static void readGetBuffers()
Deike Kleberg's avatar
Deike Kleberg committed
653
{
654
  int nProcsModel = commInqNProcsModel ();
Deike Kleberg's avatar
Deike Kleberg committed
655
  int root        = commInqRootGlob ();
656
#ifdef HAVE_NETCDF4
657
  int myCollRank = commInqRankColl();
658
  MPI_Comm collComm = commInqCommColl();
659
#endif
660
  xdebug("%s", "START");
661

662
663
  struct winHeaderEntry *winDict
    = (struct winHeaderEntry *)rxWin[root].buffer;
664
  xassert(winDict[0].id == HEADERSIZEMARKER);
665
666
  {
    int dictSize = rxWin[root].dictSize,
667
      firstNonRPCEntry = dictSize - winDict[0].specific.headerSize.numRPCEntries - 1,
668
669
670
671
672
673
      headerIdx,
      numFuncCalls = 0;
    for (headerIdx = dictSize - 1;
         headerIdx > firstNonRPCEntry;
         --headerIdx)
      {
674
675
        xassert(winDict[headerIdx].id >= MINFUNCID
                && winDict[headerIdx].id <= MAXFUNCID);
676
        ++numFuncCalls;
677
        readFuncCall(winDict + headerIdx);
678
      }
679
    xassert(numFuncCalls == winDict[0].specific.headerSize.numRPCEntries);
680
  }
Thomas Jahns's avatar
Thomas Jahns committed
681
  /* build list of streams, data was transferred for */
682
  {
683
    struct streamMap map = buildStreamMap(winDict);
684
    double *data = NULL;
Thomas Jahns's avatar
Thomas Jahns committed
685
686
687
#ifdef HAVE_NETCDF4
    int *varIsWritten = NULL;
#endif
688
689
690
#if defined (HAVE_PARALLEL_NC4)
    double *writeBuf = NULL;
#endif
Thomas Jahns's avatar
Thomas Jahns committed
691
    int currentDataBufSize = 0;
692
    for (int streamIdx = 0; streamIdx < map.numEntries; ++streamIdx)
Thomas Jahns's avatar
Thomas Jahns committed
693
      {
694
        int streamID = map.entries[streamIdx].streamID;
Thomas Jahns's avatar
Thomas Jahns committed
695
        int vlistID = streamInqVlist(streamID);
696
        int filetype = map.entries[streamIdx].filetype;
Thomas Jahns's avatar
Thomas Jahns committed
697

698
        switch (filetype)
699
700
701
          {
          case FILETYPE_GRB:
          case FILETYPE_GRB2:
702
703
704
            writeGribStream(winDict, map.entries + streamIdx,
                            &data, &currentDataBufSize,
                            root, nProcsModel);
705
            break;
706
707
708
709
710
711
712
#ifdef HAVE_NETCDF4
          case FILETYPE_NC:
          case FILETYPE_NC2:
          case FILETYPE_NC4:
#ifdef HAVE_PARALLEL_NC4
            /* HAVE_PARALLE_NC4 implies having ScalES-PPM and yaxt */
            {
713
714
              int nvars = map.entries[streamIdx].numVars;
              int *varMap = map.entries[streamIdx].varMap;
715
716
              buildWrittenVars(map.entries + streamIdx, &varIsWritten,
                               myCollRank, collComm);
717
718
719
720
              for (int varID = 0; varID < nvars; ++varID)
                if (varIsWritten[varID])
                  {
                    struct PPM_extent varShape[3];
721
                    queryVarBounds(varShape, vlistID, varID);
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
                    struct xyzDims collGrid = varDimsCollGridMatch(varShape);
                    xdebug("writing varID %d with dimensions: "
                           "x=%d, y=%d, z=%d,\n"
                           "found distribution with dimensions:"
                           " x=%d, y=%d, z=%d.", varID,
                           varShape[0].size, varShape[1].size, varShape[2].size,
                           collGrid.sizes[0], collGrid.sizes[1],
                           collGrid.sizes[2]);
                    struct PPM_extent varChunk[3];
                    myVarPart(varShape, collGrid, varChunk);
                    int myChunk[3][2];
                    for (int i = 0; i < 3; ++i)
                      {
                        myChunk[i][0] = PPM_extent_start(varChunk[i]);
                        myChunk[i][1] = PPM_extent_end(varChunk[i]);
                      }
                    xdebug("Writing chunk { { %d, %d }, { %d, %d },"
                           " { %d, %d } }", myChunk[0][0], myChunk[0][1],
                           myChunk[1][0], myChunk[1][1], myChunk[2][0],
                           myChunk[2][1]);
                    Xt_int varSize[3];
                    for (int i = 0; i < 3; ++i)
                      varSize[2 - i] = varShape[i].size;
                    Xt_idxlist preRedistChunk, preWriteChunk;
                    /* prepare yaxt descriptor for current data
                       distribution after collect */
                    int nmiss;
                    if (varMap[varID] == -1)
                      {
                        preRedistChunk = xt_idxempty_new();
                        xdebug("%s", "I got none\n");
                      }
                    else
                      {
                        Xt_int preRedistStart[3] = { 0, 0, 0 };
                        preRedistChunk
                          = xt_idxsection_new(0, 3, varSize, varSize,
                                              preRedistStart);
                        resizeVarGatherBuf(vlistID, varID, &data,
                                           &currentDataBufSize);
                        int headerIdx = varMap[varID];
                        gatherArray(root, nProcsModel, headerIdx,
                                    vlistID, data, &nmiss);
                        xdebug("%s", "I got all\n");
                      }
                    MPI_Bcast(&nmiss, 1, MPI_INT, varIsWritten[varID] - 1,
                              collComm);
                    /* prepare yaxt descriptor for write chunk */
                    {
                      Xt_int preWriteChunkStart[3], preWriteChunkSize[3];
                      for (int i = 0; i < 3; ++i)
                        {
                          preWriteChunkStart[2 - i] = varChunk[i].first;
                          preWriteChunkSize[2 - i] = varChunk[i].size;
                        }
                      preWriteChunk = xt_idxsection_new(0, 3, varSize,
                                                        preWriteChunkSize,
                                                        preWriteChunkStart);
                    }
                    /* prepare redistribution */
                    {
                      Xt_xmap xmap = xt_xmap_all2all_new(preRedistChunk,
                                                         preWriteChunk,
                                                         collComm);
                      Xt_redist redist = xt_redist_p2p_new(xmap, MPI_DOUBLE);
                      xt_idxlist_delete(preRedistChunk);
                      xt_idxlist_delete(preWriteChunk);
                      xt_xmap_delete(xmap);
Uwe Schulzweida's avatar
Uwe Schulzweida committed
790
791
792
                      writeBuf = (double*) xrealloc(writeBuf,
                                                    sizeof (double)
                                                    * PPM_extents_size(3, varChunk));
793
                      xt_redist_s_exchange1(redist, data, writeBuf);
794
795
796
797
798
799
800
801
802
803
                      xt_redist_delete(redist);
                    }
                    /* write chunk */
                    streamWriteVarChunk(streamID, varID,
                                        (const int (*)[2])myChunk, writeBuf,
                                        nmiss);
                  }
            }
#else
            /* determine process which has stream open (writer) and
804
805
806
             * which has data for which variable (var owner)
             * three cases need to be distinguished */
            {
807
808
              int nvars = map.entries[streamIdx].numVars;
              int *varMap = map.entries[streamIdx].varMap;
809
810
              buildWrittenVars(map.entries + streamIdx, &varIsWritten,
                               myCollRank, collComm);
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
              int writerRank;
              if ((writerRank = cdiPioSerialOpenFileMap(streamID))
                  == myCollRank)
                {
                  for (int varID = 0; varID < nvars; ++varID)
                    if (varIsWritten[varID])
                      {
                        int nmiss;
                        int size = vlistInqVarSize(vlistID, varID);
                        resizeVarGatherBuf(vlistID, varID, &data,
                                           &currentDataBufSize);
                        int headerIdx = varMap[varID];
                        if (varIsWritten[varID] == myCollRank + 1)
                          {
                            /* this process has the full array and will
                             * write it */
                            xdebug("gathering varID=%d for direct writing",
                                   varID);
                            gatherArray(root, nProcsModel, headerIdx,
                                        vlistID, data, &nmiss);
                          }
                        else
                          {
                            /* another process has the array and will
                             * send it over */
                            MPI_Status stat;
                            xdebug("receiving varID=%d for writing from"
                                   " process %d",
                                   varID, varIsWritten[varID] - 1);
                            xmpiStat(MPI_Recv(&nmiss, 1, MPI_INT,
                                              varIsWritten[varID] - 1,
                                              COLLBUFNMISS,
                                              collComm, &stat), &stat);
                            xmpiStat(MPI_Recv(data, size, MPI_DOUBLE,
                                              varIsWritten[varID] - 1,
                                              COLLBUFTX,
                                              collComm, &stat), &stat);
                          }
                        streamWriteVar(streamID, varID, data, nmiss);
                      }
                }
              else
                for (int varID = 0; varID < nvars; ++varID)
                  if (varIsWritten[varID] == myCollRank + 1)
                    {
                      /* this process has the full array and another
                       * will write it */
                      int nmiss;
                      int size = vlistInqVarSize(vlistID, varID);
                      resizeVarGatherBuf(vlistID, varID, &data,
                                         &currentDataBufSize);
                      int headerIdx = varMap[varID];
                      gatherArray(root, nProcsModel, headerIdx,
                                  vlistID, data, &nmiss);
                      MPI_Request req;
                      MPI_Status stat;
                      xdebug("sending varID=%d for writing to"
                             " process %d",
                             varID, writerRank);
                      xmpi(MPI_Isend(&nmiss, 1, MPI_INT,
                                     writerRank, COLLBUFNMISS,
                                     collComm, &req));
                      xmpi(MPI_Send(data, size, MPI_DOUBLE,
                                    writerRank, COLLBUFTX,
                                    collComm));
                      xmpiStat(MPI_Wait(&req, &stat), &stat);
                    }
            }
879
880
881
#endif
            break;
#endif
882
883
884
          default:
            xabort("unhandled filetype in parallel I/O.");
          }
885
      }
Thomas Jahns's avatar
Thomas Jahns committed
886
887
#ifdef HAVE_NETCDF4
    free(varIsWritten);
Thomas Jahns's avatar
Thomas Jahns committed
888
889
890
#ifdef HAVE_PARALLEL_NC4
    free(writeBuf);
#endif
Thomas Jahns's avatar
Thomas Jahns committed
891
#endif
892
    free(map.entries);
Thomas Jahns's avatar
Thomas Jahns committed
893
    free(data);
894
  }
895
  xdebug("%s", "RETURN");
896
897
898
899
} 

/************************************************************************/

Deike Kleberg's avatar
Deike Kleberg committed
900

Thomas Jahns's avatar
Thomas Jahns committed
901
902
static
void clearModelWinBuffer(int modelID)
Deike Kleberg's avatar
Deike Kleberg committed
903
904
905
{
  int nProcsModel = commInqNProcsModel ();

Deike Kleberg's avatar
Deike Kleberg committed
906
907
  xassert ( modelID                >= 0           &&
            modelID                 < nProcsModel &&
908
            rxWin != NULL && rxWin[modelID].buffer != NULL &&
909
910
            rxWin[modelID].size > 0 &&
            rxWin[modelID].size <= MAXWINBUFFERSIZE );
911
  memset(rxWin[modelID].buffer, 0, rxWin[modelID].size);
Deike Kleberg's avatar
Deike Kleberg committed
912
913
914
915
916
917
}


/************************************************************************/


918
static
919
void getTimeStepData()
Deike Kleberg's avatar
Deike Kleberg committed
920
{
921
  int modelID;
922
  char text[1024];
923
  int nProcsModel = commInqNProcsModel ();
Thomas Jahns's avatar
Thomas Jahns committed
924
925
  void *getWinBaseAddr;
  int attrFound;
926

927
  xdebug("%s", "START");
Deike Kleberg's avatar
Deike Kleberg committed
928

929
930
  for ( modelID = 0; modelID < nProcsModel; modelID++ )
    clearModelWinBuffer(modelID);
Deike Kleberg's avatar
Deike Kleberg committed
931
  // todo put in correct lbs and ubs
932
  xmpi(MPI_Win_start(groupModel, 0, getWin));
933
934
  xmpi(MPI_Win_get_attr(getWin, MPI_WIN_BASE, &getWinBaseAddr, &attrFound));
  xassert(attrFound);
Deike Kleberg's avatar
Deike Kleberg committed
935
936
  for ( modelID = 0; modelID < nProcsModel; modelID++ )
    {
937
      xdebug("modelID=%d, nProcsModel=%d, rxWin[%d].size=%zu,"
Thomas Jahns's avatar
Thomas Jahns committed
938
             " getWin=%p, sizeof(int)=%u",
939
             modelID, nProcsModel, modelID, rxWin[modelID].size,
Thomas Jahns's avatar
Thomas Jahns committed
940
             getWinBaseAddr, (unsigned)sizeof(int));
941
      /* FIXME: this needs to use MPI_PACK for portability */
942
943
944
      xmpi(MPI_Get(rxWin[modelID].buffer, rxWin[modelID].size,
                   MPI_UNSIGNED_CHAR, modelID, 0,
                   rxWin[modelID].size, MPI_UNSIGNED_CHAR, getWin));
Deike Kleberg's avatar
Deike Kleberg committed
945
    }
946
  xmpi ( MPI_Win_complete ( getWin ));
Deike Kleberg's avatar
Deike Kleberg committed
947

948
  if ( ddebug > 2 )
Deike Kleberg's avatar
Deike Kleberg committed
949
    for ( modelID = 0; modelID < nProcsModel; modelID++ )
950
      {
951
        sprintf(text, "rxWin[%d].size=%zu from PE%d rxWin[%d].buffer",
952
                modelID, rxWin[modelID].size, modelID, modelID);
953
        xprintArray(text, rxWin[modelID].buffer,
954
955
                    rxWin[modelID].size / sizeof (double),
                    DATATYPE_FLT);
956
      }
957
958
  readGetBuffers();

959
  xdebug("%s", "RETURN");
Deike Kleberg's avatar
Deike Kleberg committed
960
}
Deike Kleberg's avatar
Deike Kleberg committed
961
962
963

/************************************************************************/

964
965
966
967
968
969
970
971
972
973
974
975
#if defined (HAVE_LIBNETCDF) && ! defined (HAVE_PARALLEL_NC4)
static int
cdiPioStreamCDFOpenWrap(const char *filename, const char *filemode,
                        int filetype, stream_t *streamptr,
                        int recordBufIsToBeCreated)
{
  switch (filetype)
    {
    case FILETYPE_NC4:
    case FILETYPE_NC4C:
      {
        int rank, fileID;
Thomas Jahns's avatar
Thomas Jahns committed
976
977
        int ioMode = commInqIOMode();
        if (ioMode == PIO_NONE
978
979
980
981
            || commInqRankColl() == (rank = cdiPioNextOpenRank()))
          fileID = cdiStreamOpenDefaultDelegate(filename, filemode, filetype,
                                                streamptr,
                                                recordBufIsToBeCreated);
982
983
        else
          streamptr->filetype = filetype;
Thomas Jahns's avatar
Thomas Jahns committed
984
        if (ioMode != PIO_NONE)
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
          xmpi(MPI_Bcast(&fileID, 1, MPI_INT, rank, commInqCommColl()));
        streamptr->ownerRank = rank;
        return fileID;
      }
    default:
      return cdiStreamOpenDefaultDelegate(filename, filemode, filetype,
                                          streamptr, recordBufIsToBeCreated);
    }
}

static void
cdiPioStreamCDFCloseWrap(stream_t *streamptr, int recordBufIsToBeDeleted)
{
  int fileID   = streamptr->fileID;
  int filetype = streamptr->filetype;
  if ( fileID == CDI_UNDEFID )
For faster browsing, not all history is shown. View entire blame